From 539d089611d79426f71ded16509e516fa58090a7 Mon Sep 17 00:00:00 2001 From: Lee <7932644+strahe@users.noreply.github.com> Date: Thu, 24 Sep 2026 18:30:47 +0800 Subject: [PATCH] fix(storage): recover Store uploads and scope dashboard health --- docs/en/operations/health-metrics.md | 2 + docs/en/operations/troubleshooting.md | 4 +- docs/en/operations/upgrade-recovery.md | 2 +- docs/en/reference/cli-api.md | 2 +- docs/zh/operations/health-metrics.md | 2 + docs/zh/operations/troubleshooting.md | 4 +- docs/zh/operations/upgrade-recovery.md | 2 +- docs/zh/reference/cli-api.md | 2 +- internal/admin/api_overview.go | 80 ++- internal/admin/api_overview_test.go | 52 +- internal/db/repository/cache_eviction_repo.go | 23 + .../repository/cache_store_pin_repo_test.go | 105 ++++ internal/db/repository/interfaces.go | 3 + internal/db/repository/observability_repo.go | 41 ++ .../db/repository/observability_repo_test.go | 151 +++++ internal/db/repository/storage_commit_repo.go | 47 +- internal/db/repository/task_repo.go | 33 + internal/storagecommit/advancer_test.go | 105 +++- internal/synapse/readiness.go | 2 +- internal/synapse/readiness_test.go | 17 + .../storage_checkpoint_recovery_test.go | 586 ++++++++++++++++++ internal/worker/storage_task_handlers.go | 246 +++++++- internal/worker/task_handlers_test.go | 42 +- ui/src/lib/overview.ts | 16 + ui/src/routes/-root-content.ts | 21 +- ui/src/routes/__root.tsx | 6 +- ui/src/routes/index.tsx | 7 +- ui/src/routes/tasks.tsx | 2 +- ui/test/overview.test.ts | 37 +- ui/test/root-content.test.ts | 26 +- 30 files changed, 1537 insertions(+), 131 deletions(-) create mode 100644 internal/db/repository/cache_store_pin_repo_test.go create mode 100644 internal/worker/storage_checkpoint_recovery_test.go diff --git a/docs/en/operations/health-metrics.md b/docs/en/operations/health-metrics.md index 1631125..3145796 100644 --- a/docs/en/operations/health-metrics.md +++ b/docs/en/operations/health-metrics.md @@ -38,6 +38,8 @@ Failed check: ## Background Task Activity +The dashboard overview counts this node's data sets and their providers. The Providers page shows the full observed provider list. + SynapS3 reports an unhealthy task processor when it stops reporting activity for longer than its configured health window. This detects stalled background storage even when no upload is active. Check background task state: diff --git a/docs/en/operations/troubleshooting.md b/docs/en/operations/troubleshooting.md index be2bc22..f4c9088 100644 --- a/docs/en/operations/troubleshooting.md +++ b/docs/en/operations/troubleshooting.md @@ -132,11 +132,11 @@ Retry only after RPC connectivity, storage provider availability, wallet funds, synaps3 admin task retry 42 ``` -The API indicates whether Retry is available for each failed task. Provider replacement work is recovered from **Details** → **Storage** → **Data Sets**. A wallet operation can be recovered from Tasks only when no broadcast started; an uncertain broadcast remains non-retryable. An uncertain Store offers **Check again**, which observes the provider without uploading again. For remote-copy removal that remains unconfirmed after 24 hours, **Recover** checks the chain first and may submit another paid request if the copy is still present and not queued; the earlier request may still succeed. If the data set is no longer active, the chain cannot confirm this individual removal, so the task keeps checking instead of claiming the copy is gone. Use **Dismiss** or `synaps3 admin task acknowledge ` only after reviewing the failure; acknowledged tasks remain available for the configured retention period before cleanup. When failures have piled up, **Dismiss all** on the Tasks page clears the ones the current Operation filter selects, and `synaps3 admin task acknowledge --type --yes` does the same from the CLI; failures recorded after you confirm stay in the list. +The API indicates whether Retry is available for each failed task. Provider replacement work is recovered from **Details** → **Storage** → **Data Sets**. A wallet operation can be recovered from Tasks only when no broadcast started; an uncertain broadcast remains non-retryable. A failed Store offers **Retry upload**; it checks the provider first and resends the piece only if missing. For remote-copy removal that remains unconfirmed after 24 hours, **Recover** checks the chain first and may submit another paid request if the copy is still present and not queued; the earlier request may still succeed. If the data set is no longer active, the chain cannot confirm this individual removal, so the task keeps checking instead of claiming the copy is gone. Use **Dismiss** or `synaps3 admin task acknowledge ` only after reviewing the failure; acknowledged tasks remain available for the configured retention period before cleanup. When failures have piled up, **Dismiss all** on the Tasks page clears the ones the current Operation filter selects, and `synaps3 admin task acknowledge --type --yes` does the same from the CLI; failures recorded after you confirm stay in the list. ## Provider or RPC Issues -Check provider health and Filecoin readiness in the dashboard, or inspect the Admin API: +Check this node's storage health on Overview and Filecoin readiness in Settings, or inspect the Admin API: ```bash curl -u admin http://127.0.0.1:9090/api/v1/filecoin/readiness diff --git a/docs/en/operations/upgrade-recovery.md b/docs/en/operations/upgrade-recovery.md index 0713d6f..1204759 100644 --- a/docs/en/operations/upgrade-recovery.md +++ b/docs/en/operations/upgrade-recovery.md @@ -67,7 +67,7 @@ After a restart, unfinished work becomes eligible to continue automatically. - Retry a failed task only when the dashboard or API marks it retryable. - Recover provider replacements from **Details** → **Storage** → **Data Sets**. - A wallet operation can be retried from Tasks only when no broadcast started. An uncertain broadcast remains non-retryable. -- For an uncertain Store, **Check again** checks the provider without uploading the object again. +- **Retry upload** checks whether the provider has the piece, then uploads it again if missing. A repeat can use more bandwidth or open another upload session. - `status=failed` lists unacknowledged failures. Use `status=dismissed` to list acknowledged failures. - Review unresolved storage confirmations with `synaps3 admin storage-confirmation list`. diff --git a/docs/en/reference/cli-api.md b/docs/en/reference/cli-api.md index c9f7664..194185f 100644 --- a/docs/en/reference/cli-api.md +++ b/docs/en/reference/cli-api.md @@ -101,7 +101,7 @@ Admin global flags must appear after `admin` and before the subcommand: Task listing supports `--type`, `--status`, `--limit`, and ID-based `--cursor`. Valid status filters are `pending`, `running`, `completed`, `failed`, `cancelled`, and `dismissed`. Pending work is presented as queued, scheduled, or waiting; `failed` returns unacknowledged failures and `dismissed` returns acknowledged failures. -`synaps3 admin task retry` recovers only tasks whose response says they are retryable. Provider replacement recovery remains under **Details** → **Storage** → **Data Sets**. A wallet operation can be retried only when its broadcast never started; an operation with an uncertain broadcast remains non-retryable. For an uncertain Store, the dashboard labels Retry as **Check again**: this checks the provider without uploading again. For remote-copy removal that remains unconfirmed after 24 hours, Recover checks whether the copy is gone or already queued for removal; otherwise it may submit another paid request, even though the earlier request may still succeed. Use `synaps3 admin task acknowledge ` to dismiss a failed task after reviewing its outcome; acknowledgement starts its retention period, after which it may be cleaned up. To clear a backlog, run it without an ID and confirm with `--yes`: `--type` limits it to one operation, and `--before` sets an RFC 3339 cutoff that defaults to now, so failures recorded later stay visible. +`synaps3 admin task retry` recovers only tasks whose response says they are retryable. Provider replacement recovery remains under **Details** → **Storage** → **Data Sets**. A wallet operation can be retried only when its broadcast never started; an operation with an uncertain broadcast remains non-retryable. For a failed Store, **Retry upload** checks the provider first and sends the piece again only if missing. For remote-copy removal that remains unconfirmed after 24 hours, Recover checks whether the copy is gone or already queued for removal; otherwise it may submit another paid request, even though the earlier request may still succeed. Use `synaps3 admin task acknowledge ` to dismiss a failed task after reviewing its outcome; acknowledgement starts its retention period, after which it may be cleaned up. To clear a backlog, run it without an ID and confirm with `--yes`: `--type` limits it to one operation, and `--before` sets an RFC 3339 cutoff that defaults to now, so failures recorded later stay visible. `synaps3 admin storage-confirmation list` shows storage confirmations that need review. Verify the piece CID, provider, current attempt ID, attempted time, and any available transaction evidence before running `storage-confirmation release --attempt-id --yes`; release only if you accept that the provider may already store the piece and resubmission may create duplicate paid storage. A stale attempt ID is refused. diff --git a/docs/zh/operations/health-metrics.md b/docs/zh/operations/health-metrics.md index 50e171c..497a291 100644 --- a/docs/zh/operations/health-metrics.md +++ b/docs/zh/operations/health-metrics.md @@ -38,6 +38,8 @@ curl http://127.0.0.1:9090/healthz ## 后台任务活动 +仪表盘概览只统计本节点的数据集及其存储提供方;Providers 页面仍显示完整的观测名单。 + 如果后台任务处理长时间没有在配置的健康窗口内报告活动,SynapS3 会将其标记为不健康。即使没有正在上传的对象,这项检查也能发现后台存储已经停滞。 检查后台任务状态: diff --git a/docs/zh/operations/troubleshooting.md b/docs/zh/operations/troubleshooting.md index 24e32f3..4d77221 100644 --- a/docs/zh/operations/troubleshooting.md +++ b/docs/zh/operations/troubleshooting.md @@ -132,11 +132,11 @@ synaps3 admin task list --status failed --limit 100 synaps3 admin task retry 42 ``` -API 会标明每个失败任务是否可使用 Retry。存储提供方替换从 **Details** → **Storage** → **Data Sets** 恢复。只有尚未发出广播的钱包操作可以从 Tasks 恢复;广播结果不确定时仍不可重试。Store 结果不确定时会提供 **Check again**,它只观察存储提供方,不会再次上传。远端副本删除超过 24 小时仍无法确认时,**Recover** 会先检查链上状态;若副本仍存在且未排队删除,可能再次提交付费请求,而先前的请求仍可能成功。若数据集在链上不再活跃,链上无法据此确认这份副本已删除,任务会继续核查,不会将其记为已删除。只有在核对失败结果后才使用 **Dismiss** 或 `synaps3 admin task acknowledge `;确认后的任务会继续保留配置的时长,再由后台清理。失败任务积压时,任务页的 **Dismiss all** 会处理当前操作类型筛选下的失败任务,命令行对应 `synaps3 admin task acknowledge --type <操作> --yes`;在你确认之后才记录的失败仍会留在列表里。 +API 会标明每个失败任务是否可使用 Retry。存储提供方替换从 **Details** → **Storage** → **Data Sets** 恢复。只有尚未发出广播的钱包操作可以从 Tasks 恢复;广播结果不确定时仍不可重试。Store 失败时可使用 **Retry upload**;它先检查提供方,确认分片缺失才重新上传。远端副本删除超过 24 小时仍无法确认时,**Recover** 会先检查链上状态;若副本仍存在且未排队删除,可能再次提交付费请求,而先前的请求仍可能成功。若数据集在链上不再活跃,链上无法据此确认这份副本已删除,任务会继续核查,不会将其记为已删除。只有在核对失败结果后才使用 **Dismiss** 或 `synaps3 admin task acknowledge `;确认后的任务会继续保留配置的时长,再由后台清理。失败任务积压时,任务页的 **Dismiss all** 会处理当前操作类型筛选下的失败任务,命令行对应 `synaps3 admin task acknowledge --type <操作> --yes`;在你确认之后才记录的失败仍会留在列表里。 ## 存储提供方或 RPC 问题 -在仪表盘查看存储提供方健康状态和 Filecoin readiness,或检查 Admin API: +在 Overview 查看本节点存储健康状态,在 Settings 查看 Filecoin readiness,或检查 Admin API: ```bash curl -u admin http://127.0.0.1:9090/api/v1/filecoin/readiness diff --git a/docs/zh/operations/upgrade-recovery.md b/docs/zh/operations/upgrade-recovery.md index 5a8154a..842ff75 100644 --- a/docs/zh/operations/upgrade-recovery.md +++ b/docs/zh/operations/upgrade-recovery.md @@ -67,7 +67,7 @@ SynapS3 不会修改不兼容的数据库。废弃的 `worker.upload`、`worker. - 只重试仪表盘或 API 标记为可重试的失败任务。 - 从 **Details** → **Storage** → **Data Sets** 恢复存储提供方替换。 - 钱包操作只有在尚未发出广播时才能从 Tasks 重试;广播结果不确定时仍不可重试。 -- Store 结果不确定时,**Check again** 只查询存储提供方,不会重新上传对象。 +- **Retry upload** 会先检查存储提供方是否已有分片;确认缺失后才重新上传。重传可能增加带宽用量或开启另一次上传会话。 - `status=failed` 只列出尚未确认的失败;使用 `status=dismissed` 查看已确认的失败。 - 使用 `synaps3 admin storage-confirmation list` 核对尚未解决的存储确认。 diff --git a/docs/zh/reference/cli-api.md b/docs/zh/reference/cli-api.md index ad52feb..410426f 100644 --- a/docs/zh/reference/cli-api.md +++ b/docs/zh/reference/cli-api.md @@ -101,7 +101,7 @@ Admin 全局 flags 必须放在 `admin` 之后、子命令之前: 列出后台任务时支持 `--type`、`--status`、`--limit` 和基于任务 ID 的 `--cursor`。有效的状态过滤值为 `pending`、`running`、`completed`、`failed`、`cancelled` 和 `dismissed`。pending 工作会显示为 queued、scheduled 或 waiting;`failed` 返回尚未确认的失败,`dismissed` 返回已确认的失败。 -`synaps3 admin task retry` 只恢复响应中标记为可重试的失败任务。存储提供方替换仍在 **Details** → **Storage** → **Data Sets** 中恢复。只有尚未发出广播的钱包操作可以重试;广播结果不确定时仍不可重试。Store 结果不确定时,dashboard 会把 Retry 显示为 **Check again**:该操作只查询存储提供方,不会重新上传。远端副本删除超过 24 小时仍无法确认时,**Recover** 会先检查副本是否已删除或已排队删除;若仍存在且未排队,可能再次提交付费删除请求,而先前的请求仍可能成功。核对失败结果后,可用 `synaps3 admin task acknowledge ` 将任务标记为已处理;确认后开始计算保留期,到期后可能被清理。需要清理积压时,不带 ID 运行并用 `--yes` 确认:`--type` 限定某一种操作,`--before` 指定 RFC 3339 截止时刻(默认为当前时间),该时刻之后记录的失败仍然可见。 +`synaps3 admin task retry` 只恢复响应中标记为可重试的失败任务。存储提供方替换仍在 **Details** → **Storage** → **Data Sets** 中恢复。只有尚未发出广播的钱包操作可以重试;广播结果不确定时仍不可重试。Store 失败后,**Retry upload** 会先检查提供方;确认分片缺失才再次上传。远端副本删除超过 24 小时仍无法确认时,**Recover** 会先检查副本是否已删除或已排队删除;若仍存在且未排队,可能再次提交付费删除请求,而先前的请求仍可能成功。核对失败结果后,可用 `synaps3 admin task acknowledge ` 将任务标记为已处理;确认后开始计算保留期,到期后可能被清理。需要清理积压时,不带 ID 运行并用 `--yes` 确认:`--type` 限定某一种操作,`--before` 指定 RFC 3339 截止时刻(默认为当前时间),该时刻之后记录的失败仍然可见。 `synaps3 admin storage-confirmation list` 会显示需要核对的存储确认。核对存储提供方、transaction 和当前 attempt 后,使用 `storage-confirmation release --attempt-id --yes` 表示确认存储提供方可能已经接受该 piece,并允许正常恢复流程再次提交。过期的 attempt ID 会被拒绝。 diff --git a/internal/admin/api_overview.go b/internal/admin/api_overview.go index dc88840..5866f9c 100644 --- a/internal/admin/api_overview.go +++ b/internal/admin/api_overview.go @@ -184,38 +184,76 @@ func (s *Server) filecoinStorageHealthOverview(ctx context.Context) filecoinStor return health } - providers, err := s.observability.ListProviderObservations(ctx, observability.ListOptions{Limit: 1}) + dataSetRows, providerStates, dataSetStates, providerCheckedAt, dataSetCheckedAt, err := s.repos.Observability.OverviewStorageStates(ctx) if err != nil { if s.logger != nil { - s.logger.Warn("overview: failed to load provider observability summary", "error", err) + s.logger.Warn("overview: failed to load storage observability summary", "error", err) } - health.PartialErrors["observability_providers"] = "provider health query failed" + health.PartialErrors["observability"] = "storage health query failed" health.Level = observability.WorstSignalLevel(health.Level, observability.SignalWarning) - } else { - health.Providers = &filecoinStorageHealthObservationOverview{ - Summary: providers.Summary, - SummarySignal: providers.SummarySignal, - } - health.Level = observability.WorstSignalLevel(health.Level, providers.SummarySignal.Level) + return health } - - dataSets, err := s.observability.ListDataSetObservations(ctx, observability.ListOptions{Limit: 1}) - if err != nil { - if s.logger != nil { - s.logger.Warn("overview: failed to load data set observability summary", "error", err) + providerByID := make(map[string]observability.ProviderState, len(providerStates)) + for _, state := range providerStates { + providerByID[state.ProviderID.String()] = state + providerCheckedAt = olderOverviewObservation(providerCheckedAt, state.LastCheckedAt) + } + dataSetByID := make(map[int64]observability.DataSetState, len(dataSetStates)) + for _, state := range dataSetStates { + dataSetByID[state.LocalDataSetID] = state + dataSetCheckedAt = olderOverviewObservation(dataSetCheckedAt, state.LastCheckedAt) + } + providers := observability.Summary{} + dataSets := observability.Summary{Total: len(dataSetRows)} + seenProviders := make(map[string]bool, len(dataSetRows)) + for _, row := range dataSetRows { + if state, ok := dataSetByID[row.ID]; ok { + addStorageHealthStatus(&dataSets, state.Status) + } else { + dataSets.Unknown++ } - health.PartialErrors["observability_data_sets"] = "data set health query failed" - health.Level = observability.WorstSignalLevel(health.Level, observability.SignalWarning) - } else { - health.DataSets = &filecoinStorageHealthObservationOverview{ - Summary: dataSets.Summary, - SummarySignal: dataSets.SummarySignal, + id := row.ProviderID.String() + if seenProviders[id] { + continue + } + seenProviders[id] = true + providers.Total++ + if state, ok := providerByID[id]; ok { + addStorageHealthStatus(&providers, state.Status) + } else { + providers.Unknown++ } - health.Level = observability.WorstSignalLevel(health.Level, dataSets.SummarySignal.Level) } + now := time.Now().UTC() + interval := s.observability.RefreshInterval() + providerSignal := observability.DefaultAttentionSummarySignal(providers, providerCheckedAt, interval, now) + dataSetSignal := observability.DefaultAttentionSummarySignal(dataSets, dataSetCheckedAt, interval, now) + health.Providers = &filecoinStorageHealthObservationOverview{Summary: providers, SummarySignal: providerSignal} + health.DataSets = &filecoinStorageHealthObservationOverview{Summary: dataSets, SummarySignal: dataSetSignal} + health.Level = observability.WorstSignalLevel(providerSignal.Level, dataSetSignal.Level) return health } +func olderOverviewObservation(current *time.Time, observed time.Time) *time.Time { + if observed.IsZero() || (current != nil && !observed.Before(*current)) { + return current + } + return &observed +} + +func addStorageHealthStatus(summary *observability.Summary, status observability.Status) { + switch status { + case observability.StatusAvailable: + summary.Available++ + case observability.StatusDegraded: + summary.Degraded++ + case observability.StatusUnavailable: + summary.Unavailable++ + default: + summary.Unknown++ + } +} + func taskPipelineOverviewRows(counts []repository.TaskPipelineCount) []taskPipelineOverview { rows := make([]taskPipelineOverview, 0) index := make(map[string]int) diff --git a/internal/admin/api_overview_test.go b/internal/admin/api_overview_test.go index ee9a810..2145929 100644 --- a/internal/admin/api_overview_test.go +++ b/internal/admin/api_overview_test.go @@ -15,13 +15,41 @@ import ( "github.com/strahe/synaps3/internal/observability" taskengine "github.com/strahe/synaps3/internal/task" "github.com/strahe/synaps3/internal/testutil" + "github.com/strahe/synaps3/internal/types" "github.com/uptrace/bun" ) +type overviewHealthRepository struct { + repository.ObservabilityRepository + dataSets []model.StorageDataSet + providers []observability.ProviderState + states []observability.DataSetState + checkedAt time.Time + err error +} + +func (r *overviewHealthRepository) OverviewStorageStates(context.Context) ([]model.StorageDataSet, []observability.ProviderState, []observability.DataSetState, *time.Time, *time.Time, error) { + return r.dataSets, r.providers, r.states, &r.checkedAt, &r.checkedAt, r.err +} + +func overviewHealthRepo(base repository.ObservabilityRepository, providerStatuses, dataSetStatuses []observability.Status, checkedAt time.Time) repository.ObservabilityRepository { + r := &overviewHealthRepository{ObservabilityRepository: base, checkedAt: checkedAt} + for i, status := range providerStatuses { + r.providers = append(r.providers, observability.ProviderState{ProviderID: types.NewOnChainID(uint64(i + 1)), Status: status}) + } + for i, status := range dataSetStatuses { + providerID := types.NewOnChainID(uint64(i%len(providerStatuses) + 1)) + r.dataSets = append(r.dataSets, model.StorageDataSet{ID: int64(i + 1), ProviderID: providerID}) + r.states = append(r.states, observability.DataSetState{LocalDataSetID: int64(i + 1), Status: status}) + } + return r +} + func TestAPIOverviewFilecoinStorageHealthUsesObservabilitySummaries(t *testing.T) { db := testutil.NewTestDB(t) repos := repository.NewRepositories(db) - checkedAt := time.Date(2026, 5, 18, 12, 0, 0, 0, time.UTC) + checkedAt := time.Now().UTC() + repos.Observability = overviewHealthRepo(repos.Observability, []observability.Status{observability.StatusAvailable, observability.StatusAvailable}, []observability.Status{observability.StatusAvailable, observability.StatusAvailable, observability.StatusAvailable}, checkedAt) srv := newTestServer(":0", db, &stubCache{rootDir: t.TempDir()}, 100, repos, nil, nil, config.DefaultFilecoinCopies, testLogger()). WithObservability(observability.NewService(observability.ServiceOptions{ Store: &observabilityStateStore{ @@ -56,7 +84,8 @@ func TestAPIOverviewFilecoinStorageHealthUsesObservabilitySummaries(t *testing.T func TestAPIOverviewFilecoinStorageHealthWarnsForObservabilitySignalsWithoutReinterpretingSummary(t *testing.T) { db := testutil.NewTestDB(t) repos := repository.NewRepositories(db) - checkedAt := time.Date(2026, 5, 18, 12, 0, 0, 0, time.UTC) + checkedAt := time.Now().UTC() + repos.Observability = overviewHealthRepo(repos.Observability, []observability.Status{observability.StatusDegraded}, []observability.Status{observability.StatusUnknown}, checkedAt) srv := newTestServer(":0", db, &stubCache{rootDir: t.TempDir()}, 100, repos, nil, nil, config.DefaultFilecoinCopies, testLogger()). WithObservability(observability.NewService(observability.ServiceOptions{ Store: &observabilityStateStore{ @@ -79,10 +108,11 @@ func TestAPIOverviewFilecoinStorageHealthWarnsForObservabilitySignalsWithoutRein } } -func TestAPIOverviewFilecoinStorageHealthRollsUpBlockingObservabilitySignal(t *testing.T) { +func TestAPIOverviewKeepsBlockingLevelWhenObservationsAreStale(t *testing.T) { db := testutil.NewTestDB(t) repos := repository.NewRepositories(db) - checkedAt := time.Date(2026, 5, 18, 12, 0, 0, 0, time.UTC) + checkedAt := time.Now().UTC().Add(-10 * time.Minute) + repos.Observability = overviewHealthRepo(repos.Observability, []observability.Status{observability.StatusUnavailable}, []observability.Status{observability.StatusAvailable}, checkedAt) srv := newTestServer(":0", db, &stubCache{rootDir: t.TempDir()}, 100, repos, nil, nil, config.DefaultFilecoinCopies, testLogger()). WithObservability(observability.NewService(observability.ServiceOptions{ Store: &observabilityStateStore{ @@ -103,6 +133,9 @@ func TestAPIOverviewFilecoinStorageHealthRollsUpBlockingObservabilitySignal(t *t if body.FilecoinStorageHealth.Level != observability.SignalBlocking { t.Fatalf("filecoin storage health level = %s, want blocking", body.FilecoinStorageHealth.Level) } + if !body.FilecoinStorageHealth.Providers.SummarySignal.Freshness.Stale { + t.Fatal("provider observation should be stale") + } } func TestAPIOverviewFilecoinStorageHealthHandlesMissingObservability(t *testing.T) { @@ -125,6 +158,7 @@ func TestAPIOverviewFilecoinStorageHealthHandlesMissingObservability(t *testing. func TestAPIOverviewFilecoinStorageHealthHandlesObservabilityQueryFailures(t *testing.T) { db := testutil.NewTestDB(t) repos := repository.NewRepositories(db) + repos.Observability = &overviewHealthRepository{ObservabilityRepository: repos.Observability, err: errors.New("provider rpc failed with sensitive detail")} srv := newTestServer(":0", db, &stubCache{rootDir: t.TempDir()}, 100, repos, nil, nil, config.DefaultFilecoinCopies, testLogger()). WithObservability(observability.NewService(observability.ServiceOptions{ Store: &observabilityStateStore{ @@ -138,11 +172,8 @@ func TestAPIOverviewFilecoinStorageHealthHandlesObservabilityQueryFailures(t *te if body.FilecoinStorageHealth.Level != observability.SignalWarning { t.Fatalf("filecoin storage health level = %s, want warning", body.FilecoinStorageHealth.Level) } - if got := body.FilecoinStorageHealth.PartialErrors["observability_providers"]; got != "provider health query failed" { - t.Fatalf("provider partial error = %q, want sanitized query failure", got) - } - if got := body.FilecoinStorageHealth.PartialErrors["observability_data_sets"]; got != "data set health query failed" { - t.Fatalf("data set partial error = %q, want sanitized query failure", got) + if got := body.FilecoinStorageHealth.PartialErrors["observability"]; got != "storage health query failed" { + t.Fatalf("storage partial error = %q, want sanitized query failure", got) } } @@ -150,7 +181,8 @@ func TestAPIOverviewFilecoinStorageHealthIgnoresTaskPressure(t *testing.T) { db := testutil.NewTestDB(t) repos := repository.NewRepositories(db) taskService := newAdminTestTaskService(t, repos) - checkedAt := time.Date(2026, 5, 18, 12, 0, 0, 0, time.UTC) + checkedAt := time.Now().UTC() + repos.Observability = overviewHealthRepo(repos.Observability, []observability.Status{observability.StatusAvailable}, []observability.Status{observability.StatusAvailable}, checkedAt) overviewSeedTask(t, taskService, repos, model.TaskTypeStorageStore, "running", model.TaskStatusRunning) overviewSeedTask(t, taskService, repos, model.TaskTypeStorageStore, "failed", model.TaskStatusFailed) srv := newTestServer(":0", db, &stubCache{rootDir: t.TempDir()}, 100, repos, nil, nil, config.DefaultFilecoinCopies, testLogger()). diff --git a/internal/db/repository/cache_eviction_repo.go b/internal/db/repository/cache_eviction_repo.go index 0679db7..8824d5d 100644 --- a/internal/db/repository/cache_eviction_repo.go +++ b/internal/db/repository/cache_eviction_repo.go @@ -195,6 +195,7 @@ func (r *BunCacheEvictionRepo) ListLRUCandidates(ctx context.Context, limit int) Where("storage_content.content_size > 0"). Where("object_cache.cache_accessed_at IS NOT NULL"). Where(minimumDurabilityMetSQL("storage_content", "durability_bucket")). + Where(noUnfinishedStoreCacheDependencySQL("storage_content.id")). Where(noUnfinishedReplacementCacheDependencySQL("storage_content.id"), storagereplacement.ItemStatusPending, storagereplacement.ItemStatusAttention, @@ -270,6 +271,17 @@ func (r *BunCacheEvictionRepo) AuthorizeDeletion( if pendingReplacement > 0 { return cacheeviction.ErrNoLongerEligible } + var pendingStore int + if err := db.NewRaw(`SELECT COUNT(*) FROM storage_copies AS copy + WHERE copy.content_id = ? AND copy.status IN ('pending', 'piece_ready', 'committing') + AND copy.transfer_method IN ('ingress', 'cache_restore') + AND EXISTS (SELECT 1 FROM object_versions AS version WHERE version.content_id = copy.content_id)`, contentID). + Scan(ctx, &pendingStore); err != nil { + return err + } + if pendingStore > 0 { + return cacheeviction.ErrNoLongerEligible + } bucket, err := lockBucketByID(ctx, db, content.BucketID) if err != nil { return err @@ -476,6 +488,7 @@ func nextBucketDurabilityCandidate(ctx context.Context, db bun.IDB, bucketID int Where("cache_entry.in_cache = ?", true). Where("cache_entry.cache_active_task_id IS NULL"). Where(minimumDurabilityMetSQL("storage_content", "durability_bucket")). + Where(noUnfinishedStoreCacheDependencySQL("storage_content.id")). Where(noUnfinishedReplacementCacheDependencySQL("storage_content.id"), storagereplacement.ItemStatusPending, storagereplacement.ItemStatusAttention, @@ -516,6 +529,16 @@ func noUnfinishedReplacementCacheDependencySQL(contentIDExpr string) string { )`, contentIDExpr) } +func noUnfinishedStoreCacheDependencySQL(contentIDExpr string) string { + return fmt.Sprintf(`NOT EXISTS ( + SELECT 1 FROM storage_copies AS store_copy + WHERE store_copy.content_id = %s + AND store_copy.status IN ('pending', 'piece_ready', 'committing') + AND store_copy.transfer_method IN ('ingress', 'cache_restore') + AND EXISTS (SELECT 1 FROM object_versions AS version WHERE version.content_id = store_copy.content_id) + )`, contentIDExpr) +} + func cacheAccessTime(entry *model.ObjectCache) time.Time { if entry.CacheAccessedAt != nil { return *entry.CacheAccessedAt diff --git a/internal/db/repository/cache_store_pin_repo_test.go b/internal/db/repository/cache_store_pin_repo_test.go new file mode 100644 index 0000000..f0e558e --- /dev/null +++ b/internal/db/repository/cache_store_pin_repo_test.go @@ -0,0 +1,105 @@ +package repository_test + +import ( + "errors" + "fmt" + "testing" + "time" + + "github.com/strahe/synaps3/internal/cacheeviction" + "github.com/strahe/synaps3/internal/db/repository" + "github.com/strahe/synaps3/internal/model" + "github.com/strahe/synaps3/internal/testutil" +) + +func TestUnfinishedStoreProtectsCacheAfterMinimumDurability(t *testing.T) { + db := testDB(t) + repos := repository.NewRepositories(db) + bucket := seedBucket(t, db, "unfinished-store-cache") + if _, err := db.NewUpdate().Model((*model.Bucket)(nil)). + Set("default_copies = 2").Set("minimum_durable_copies = 1").Where("id = ?", bucket.ID).Exec(t.Context()); err != nil { + t.Fatal(err) + } + content, err := repos.Contents.EnsureContent(t.Context(), repository.EnsureContentInput{ + BucketID: bucket.ID, ContentSize: 128, Checksum: testutil.StorageChecksum("unfinished-store-cache"), RequestedCopies: 2, + }) + if err != nil { + t.Fatal(err) + } + version := &model.ObjectVersion{ + VersionID: model.NewVersionID(), BucketID: bucket.ID, + Key: "pinned.bin", ContentID: &content.ID, Size: 128, ETag: "pinned", ContentType: "application/octet-stream", + } + if _, err := repos.Objects.CreateVersionAndSetCurrent(t.Context(), version); err != nil { + t.Fatal(err) + } + if err := repos.Objects.SetVersionCachePresence(t.Context(), version.VersionID, true); err != nil { + t.Fatal(err) + } + bindings := make([]*model.StorageDataSet, 2) + for i := range bindings { + providerID := onChainID(t, fmt.Sprint(101+i)) + binding, err := repos.Contents.EnsureDataSetBinding(t.Context(), repository.EnsureDataSetBindingInput{ + BucketID: bucket.ID, ProviderID: providerID, CopyIndex: i, CreatedByContentID: content.ID, + }) + if err != nil { + t.Fatal(err) + } + dataSetID := onChainID(t, fmt.Sprint(301+i)) + clientID := onChainID(t, fmt.Sprint(501+i)) + if err := repos.Contents.MarkDataSetReady(t.Context(), repository.MarkDataSetReadyInput{ + ID: binding.ID, ContentID: content.ID, DataSetID: dataSetID, ClientDataSetID: &clientID, + }); err != nil { + t.Fatal(err) + } + bindings[i] = binding + } + if err := repos.Contents.CreateUploadCopiesForBindings(t.Context(), content.ID, []repository.UploadCopyBindingInput{ + {StorageDataSetID: bindings[0].ID, CopyIndex: 0, ProviderID: bindings[0].ProviderID, TransferMethod: model.StorageCopyTransferMethodPeerPull}, + {StorageDataSetID: bindings[1].ID, CopyIndex: 1, ProviderID: bindings[1].ProviderID, TransferMethod: model.StorageCopyTransferMethodIngress}, + }); err != nil { + t.Fatal(err) + } + copies, err := repos.Contents.ListCopies(t.Context(), content.ID) + if err != nil || len(copies) != 2 { + t.Fatalf("copies = %#v, err=%v", copies, err) + } + firstPieceID := onChainID(t, "71") + testutil.CommitStorageCopy(t, db, repos, repository.MarkUploadCopyCommittedInput{ + StorageCopyID: copies[0].ID, ContentID: content.ID, CopyIndex: 0, + PieceCID: "bafk2bzacecpinnedsource", PieceID: &firstPieceID, RetrievalURL: "https://source.example/piece", + }) + if candidates, err := repos.CacheEvictions.ListLRUCandidates(t.Context(), 10); err != nil || len(candidates) != 0 { + t.Fatalf("LRU candidates with unfinished Store = %#v, err=%v", candidates, err) + } + durabilityTask := enqueueAndClaimTask(t, repos, "unfinished-store-durability", time.Minute) + durabilityGeneration, err := repos.CacheEvictions.NextDurabilityGeneration(t.Context(), bucket.ID) + if err != nil { + t.Fatal(err) + } + if err := repos.CacheEvictions.BindDurabilityTask(t.Context(), bucket.ID, durabilityGeneration, durabilityTask.ID); err != nil { + t.Fatal(err) + } + if candidate, err := repos.CacheEvictions.NextBucketDurabilityCandidate(t.Context(), bucket.ID, durabilityGeneration, durabilityTask.ID); err != nil || candidate != nil { + t.Fatalf("capacity cleanup candidate with unfinished Store = %#v, err=%v", candidate, err) + } + evict := enqueueAndClaimTask(t, repos, "unfinished-store-eviction", time.Minute) + reservation, err := repos.CacheEvictions.PrepareEviction(t.Context(), content.ID) + if err != nil { + t.Fatal(err) + } + if err := repos.CacheEvictions.BindEvictionTask(t.Context(), content.ID, reservation.Generation, evict.ID); err != nil { + t.Fatal(err) + } + if _, err := repos.CacheEvictions.AuthorizeDeletion(t.Context(), content.ID, reservation.Generation, evict.ID, nil); !errors.Is(err, cacheeviction.ErrNoLongerEligible) { + t.Fatalf("final deletion authorization = %v, want ineligible", err) + } + secondPieceID := onChainID(t, "72") + testutil.CommitStorageCopy(t, db, repos, repository.MarkUploadCopyCommittedInput{ + StorageCopyID: copies[1].ID, ContentID: content.ID, CopyIndex: 1, + PieceCID: "bafk2bzacecpinnedsource", PieceID: &secondPieceID, RetrievalURL: "https://target.example/piece", + }) + if _, err := repos.CacheEvictions.AuthorizeDeletion(t.Context(), content.ID, reservation.Generation, evict.ID, nil); err != nil { + t.Fatalf("final deletion authorization after Store commit: %v", err) + } +} diff --git a/internal/db/repository/interfaces.go b/internal/db/repository/interfaces.go index 3e036ef..81a2ed3 100644 --- a/internal/db/repository/interfaces.go +++ b/internal/db/repository/interfaces.go @@ -683,9 +683,11 @@ type TaskRepository interface { Enqueue(ctx context.Context, task *model.Task) (*model.Task, bool, error) GetByID(ctx context.Context, id int64) (*model.Task, error) GetByIdentity(ctx context.Context, taskType model.TaskType, idempotencyKey string) (*model.Task, error) + PreviousStoreCheckpoints(ctx context.Context, copyID, taskID int64) ([]model.Task, error) ClaimNext(ctx context.Context, leaseDuration time.Duration) (*model.Task, error) RenewLease(ctx context.Context, id, generation int64, leaseDuration time.Duration) (time.Time, error) WriteCheckpoint(ctx context.Context, id, generation int64, checkpoint json.RawMessage) error + ConsumeStoreRetry(ctx context.Context, id, generation int64) error ValidateClaim(ctx context.Context, id, generation int64) error Settle(ctx context.Context, id, generation int64, transition TaskTransition) error ShortenLease(ctx context.Context, id, generation int64, duration time.Duration) error @@ -755,6 +757,7 @@ type WalletOperationRepository interface { } type ObservabilityRepository interface { + OverviewStorageStates(ctx context.Context) ([]model.StorageDataSet, []observability.ProviderState, []observability.DataSetState, *time.Time, *time.Time, error) ReplaceProviderStates(ctx context.Context, checkedAt time.Time, states []observability.ProviderState) error ListProviderStates(ctx context.Context, opts observability.ListOptions) (observability.ProviderStatePage, error) ReplaceDataSetStates(ctx context.Context, checkedAt time.Time, states []observability.DataSetState) error diff --git a/internal/db/repository/observability_repo.go b/internal/db/repository/observability_repo.go index f41e6bc..1860abe 100644 --- a/internal/db/repository/observability_repo.go +++ b/internal/db/repository/observability_repo.go @@ -3,12 +3,53 @@ package repository import ( "context" "database/sql" + "fmt" "time" + "github.com/strahe/synaps3/internal/model" "github.com/strahe/synaps3/internal/observability" "github.com/uptrace/bun" ) +// OverviewStorageStates scopes health to local dependencies without changing global observations. +func (r *BunObservabilityRepo) OverviewStorageStates(ctx context.Context) ([]model.StorageDataSet, []observability.ProviderState, []observability.DataSetState, *time.Time, *time.Time, error) { + var dataSets []model.StorageDataSet + err := r.db.NewSelect().Model(&dataSets).Where(overviewDataSetDependencySQL("storage_data_set")).Scan(ctx) + if err != nil { + return nil, nil, nil, nil, nil, err + } + var providers []observability.ProviderState + var observedDataSets []observability.DataSetState + if err := r.db.NewSelect().Model(&providers).ModelTableExpr("observability_provider_states AS provider_state"). + Where(`EXISTS (SELECT 1 FROM storage_data_sets AS scoped_data_set + WHERE scoped_data_set.provider_id = provider_state.provider_id AND ` + overviewDataSetDependencySQL("scoped_data_set") + `)`).Scan(ctx); err != nil { + return nil, nil, nil, nil, nil, err + } + if err := r.db.NewSelect().Model(&observedDataSets). + Where(`EXISTS (SELECT 1 FROM storage_data_sets AS scoped_data_set + WHERE scoped_data_set.id = observability_data_set_state.local_data_set_id AND ` + overviewDataSetDependencySQL("scoped_data_set") + `)`).Scan(ctx); err != nil { + return nil, nil, nil, nil, nil, err + } + providerCheckedAt, err := r.collectionLastCheckedAt(ctx, observability.CollectionProviders) + if err != nil { + return nil, nil, nil, nil, nil, err + } + dataSetCheckedAt, err := r.collectionLastCheckedAt(ctx, observability.CollectionDataSets) + return dataSets, providers, observedDataSets, providerCheckedAt, dataSetCheckedAt, err +} + +func overviewDataSetDependencySQL(alias string) string { + return fmt.Sprintf(`((%[1]s.is_current AND %[1]s.status <> 'retired') + OR EXISTS (SELECT 1 FROM storage_replacements AS replacement + WHERE replacement.status NOT IN ('completed', 'superseded') + AND (replacement.source_data_set_id = %[1]s.id OR replacement.target_data_set_id = %[1]s.id)) + OR (%[1]s.status IN ('ready', 'draining') AND EXISTS ( + SELECT 1 FROM storage_copies AS copy + JOIN object_versions AS version ON version.content_id = copy.content_id + WHERE copy.storage_data_set_id = %[1]s.id AND %[2]s)))`, + alias, readableCommittedCopyPredicateSQL("copy", alias)) +} + const ( defaultObservabilityListLimit = 100 maxObservabilityListLimit = 500 diff --git a/internal/db/repository/observability_repo_test.go b/internal/db/repository/observability_repo_test.go index 8ea2d9c..3606265 100644 --- a/internal/db/repository/observability_repo_test.go +++ b/internal/db/repository/observability_repo_test.go @@ -8,6 +8,8 @@ import ( "github.com/strahe/synaps3/internal/db/repository" "github.com/strahe/synaps3/internal/model" "github.com/strahe/synaps3/internal/observability" + "github.com/strahe/synaps3/internal/storagereplacement" + "github.com/strahe/synaps3/internal/testutil" "github.com/uptrace/bun" ) @@ -91,6 +93,155 @@ func TestObservabilityRepoReplacesProviderStatesAndSummarizes(t *testing.T) { } } +func TestOverviewStorageStatesUsesLocalDependenciesWithoutFilteringGlobalObservations(t *testing.T) { + db := testDB(t) + repos := repository.NewRepositories(db) + bucket := seedBucket(t, db, "overview-local") + current := seedStorageDataSet(t, db, bucket.ID, "101", "1001", model.StorageDataSetStatusReady) + retiredBucket := seedBucket(t, db, "overview-retired") + retired := seedStorageDataSet(t, db, retiredBucket.ID, "202", "2002", model.StorageDataSetStatusReady) + if _, err := db.NewUpdate().Model((*model.StorageDataSet)(nil)).Set("is_current = ?", false).Set("status = ?", model.StorageDataSetStatusRetired). + Where("id = ?", retired.ID).Exec(t.Context()); err != nil { + t.Fatal(err) + } + checkedAt := time.Now().UTC() + if err := repos.Observability.ReplaceProviderStates(t.Context(), checkedAt, []observability.ProviderState{ + {ProviderID: current.ProviderID, Status: observability.StatusAvailable}, + {ProviderID: retired.ProviderID, Status: observability.StatusUnavailable}, + }); err != nil { + t.Fatal(err) + } + if err := repos.Observability.ReplaceDataSetStates(t.Context(), checkedAt, []observability.DataSetState{ + { + LocalDataSetID: current.ID, BucketID: bucket.ID, CopyIndex: current.CopyIndex, + ProviderID: current.ProviderID, Status: observability.StatusAvailable, + }, + { + LocalDataSetID: retired.ID, BucketID: retiredBucket.ID, CopyIndex: retired.CopyIndex, + ProviderID: retired.ProviderID, Status: observability.StatusUnavailable, + }, + }); err != nil { + t.Fatal(err) + } + dataSets, providers, states, providerCheckedAt, dataSetCheckedAt, err := repos.Observability.OverviewStorageStates(t.Context()) + if err != nil { + t.Fatal(err) + } + if len(dataSets) != 1 || dataSets[0].ID != current.ID || len(providers) != 1 || providers[0].ProviderID.String() != current.ProviderID.String() || len(states) != 1 || states[0].LocalDataSetID != current.ID { + t.Fatalf("scoped states = data sets:%#v providers:%#v observations:%#v", dataSets, providers, states) + } + if providerCheckedAt == nil || dataSetCheckedAt == nil { + t.Fatal("scoped overview lost collection freshness") + } + global, err := repos.Observability.ListProviderStates(t.Context(), observability.ListOptions{}) + if err != nil || global.Summary.Total != 2 { + t.Fatalf("global provider observations = %#v, err=%v", global, err) + } +} + +func TestOverviewStorageStatesIncludesBothSidesOfUnfinishedReplacement(t *testing.T) { + db := testDB(t) + repos := repository.NewRepositories(db) + bucket := seedBucket(t, db, "overview-replacement") + source, err := repos.Contents.EnsureDataSetBinding(t.Context(), repository.EnsureDataSetBindingInput{ + BucketID: bucket.ID, ProviderID: onChainID(t, "301"), CopyIndex: 0, + }) + if err != nil { + t.Fatal(err) + } + replacement, created, err := repos.Replacements.Authorize(t.Context(), repository.AuthorizeReplacementInput{ + BucketID: bucket.ID, SourceDataSetID: source.ID, + SelectionMode: storagereplacement.SelectionModeManual, + TargetProviderID: onChainID(t, "302"), ClientRequestID: "overview-replacement", + }) + if err != nil || !created { + t.Fatalf("replacement = %#v, created=%v, err=%v", replacement, created, err) + } + sets, _, _, _, _, err := repos.Observability.OverviewStorageStates(t.Context()) + if err != nil || len(sets) != 2 { + t.Fatalf("unfinished replacement data sets = %#v, err=%v", sets, err) + } + seen := map[int64]bool{} + for _, row := range sets { + seen[row.ID] = true + } + if !seen[source.ID] || !seen[replacement.TargetDataSetID] { + t.Fatalf("replacement source and target missing: %#v", sets) + } + if _, err := db.NewUpdate().Model((*storagereplacement.Replacement)(nil)). + Set("status = ?", storagereplacement.StatusCompleted).Where("id = ?", replacement.ID).Exec(t.Context()); err != nil { + t.Fatal(err) + } + sets, _, _, _, _, err = repos.Observability.OverviewStorageStates(t.Context()) + if err != nil || len(sets) != 1 || sets[0].ID != source.ID { + t.Fatalf("completed replacement without readable old copy = %#v, err=%v", sets, err) + } +} + +func TestOverviewStorageStatesKeepsReadableOlderGenerationOnlyWhileReferenced(t *testing.T) { + db := testDB(t) + repos := repository.NewRepositories(db) + bucket := seedBucket(t, db, "overview-old-readable") + content, err := repos.Contents.EnsureContent(t.Context(), repository.EnsureContentInput{ + BucketID: bucket.ID, ContentSize: 128, + Checksum: testutil.StorageChecksum("overview-old-readable"), RequestedCopies: 1, + }) + if err != nil { + t.Fatal(err) + } + version := &model.ObjectVersion{ + VersionID: model.NewVersionID(), BucketID: bucket.ID, + Key: "old.bin", ContentID: &content.ID, Size: 128, ETag: "old", ContentType: "application/octet-stream", + } + if _, err := repos.Objects.CreateVersionAndSetCurrent(t.Context(), version); err != nil { + t.Fatal(err) + } + old, err := repos.Contents.EnsureDataSetBinding(t.Context(), repository.EnsureDataSetBindingInput{ + BucketID: bucket.ID, ProviderID: onChainID(t, "401"), CopyIndex: 0, + }) + if err != nil { + t.Fatal(err) + } + dataSetID := onChainID(t, "1401") + clientID := onChainID(t, "2401") + if err := repos.Contents.MarkDataSetReady(t.Context(), repository.MarkDataSetReadyInput{ + ID: old.ID, ContentID: content.ID, DataSetID: dataSetID, ClientDataSetID: &clientID, + }); err != nil { + t.Fatal(err) + } + if err := repos.Contents.CreateUploadCopiesForBindings(t.Context(), content.ID, []repository.UploadCopyBindingInput{{ + StorageDataSetID: old.ID, CopyIndex: 0, ProviderID: old.ProviderID, + TransferMethod: model.StorageCopyTransferMethodIngress, + }}); err != nil { + t.Fatal(err) + } + copies, err := repos.Contents.ListCopies(t.Context(), content.ID) + if err != nil || len(copies) != 1 { + t.Fatalf("copies = %#v, err=%v", copies, err) + } + pieceID := onChainID(t, "51") + testutil.CommitStorageCopy(t, db, repos, repository.MarkUploadCopyCommittedInput{ + StorageCopyID: copies[0].ID, ContentID: content.ID, CopyIndex: 0, + PieceCID: "bafk2bzacecoverviewold", PieceID: &pieceID, RetrievalURL: "https://old.example/piece", + }) + if _, err := db.NewUpdate().Model((*model.StorageDataSet)(nil)).Set("is_current = ?", false). + Where("id = ?", old.ID).Exec(t.Context()); err != nil { + t.Fatal(err) + } + sets, _, _, _, _, err := repos.Observability.OverviewStorageStates(t.Context()) + if err != nil || len(sets) != 1 || sets[0].ID != old.ID { + t.Fatalf("readable old data set = %#v, err=%v", sets, err) + } + if _, err := db.NewUpdate().Model((*model.StorageDataSet)(nil)).Set("status = ?", model.StorageDataSetStatusRetired). + Where("id = ?", old.ID).Exec(t.Context()); err != nil { + t.Fatal(err) + } + sets, _, _, _, _, err = repos.Observability.OverviewStorageStates(t.Context()) + if err != nil || len(sets) != 0 { + t.Fatalf("retired old data set = %#v, err=%v", sets, err) + } +} + func TestObservabilityRepoProviderOrderAndPaginatedSummary(t *testing.T) { ctx := context.Background() db := testDB(t) diff --git a/internal/db/repository/storage_commit_repo.go b/internal/db/repository/storage_commit_repo.go index 4a472b0..7cd0a9b 100644 --- a/internal/db/repository/storage_commit_repo.go +++ b/internal/db/repository/storage_commit_repo.go @@ -498,6 +498,18 @@ func (r *BunStorageContentRepo) ReleaseCommitAttention(ctx context.Context, inpu } return err } + if initial.ActiveTaskID != nil { + res, err := db.NewUpdate().Model((*model.Task)(nil)). + Set("updated_at = updated_at"). + Where("id = ? AND type = ?", *initial.ActiveTaskID, model.TaskTypeStorageCommit). + Exec(ctx) + if err != nil { + return fmt.Errorf("locking storage commit task: %w", err) + } + if rows, _ := res.RowsAffected(); rows != 1 { + return ErrConflict + } + } identity := storagecommit.CopyIdentity{ StorageCopyID: initial.ID, ContentID: initial.ContentID, @@ -528,10 +540,20 @@ func (r *BunStorageContentRepo) ReleaseCommitAttention(ctx context.Context, inpu if err := projectResolvedCommitAttempt(ctx, db, copyID, now, true, true, nil); err != nil { return fmt.Errorf("releasing storage confirmation attention: %w", err) } - if err := wakeCommitFIFOHead(ctx, db, initial.StorageDataSetID); err != nil { + current := new(model.StorageCopy) + if err := db.NewSelect().Model(current).Where("id = ?", copyID).Scan(ctx); err != nil { + return err + } + if current.ContentID != initial.ContentID || current.StorageDataSetID != initial.StorageDataSetID || + current.CopyIndex != initial.CopyIndex || current.WorkGeneration != initial.WorkGeneration || + (current.ActiveTaskID == nil) != (initial.ActiveTaskID == nil) || + (current.ActiveTaskID != nil && *current.ActiveTaskID != *initial.ActiveTaskID) { + return ErrConflict + } + if err := resumeCommitTaskAfterAttentionRelease(ctx, db, initial.ActiveTaskID); err != nil { return err } - return resumeCommitTaskAfterAttentionRelease(ctx, db, initial.ActiveTaskID) + return wakeCommitFIFOHead(ctx, db, initial.StorageDataSetID) }) } @@ -549,6 +571,27 @@ func resumeCommitTaskAfterAttentionRelease(ctx context.Context, db bun.IDB, task } switch taskRow.Status { case model.TaskStatusPending, model.TaskStatusRunning: + now := time.Now() + res, err := db.NewUpdate().Model((*model.Task)(nil)). + Set("status = ?", model.TaskStatusPending). + Set("resume_mode = ?", model.TaskResumeModeRecover). + Set("available_at = ?", now). + Set("retry_count = 0"). + Set("claimed_at = NULL"). + Set("lease_until = NULL"). + Set("wait_reason = NULL"). + Set("failure_reason = NULL"). + Set("last_error = NULL"). + Set("status_message = NULL"). + Set("updated_at = ?", now). + Where("id = ? AND type = ? AND status = ? AND claim_generation = ?", taskRow.ID, model.TaskTypeStorageCommit, taskRow.Status, taskRow.ClaimGeneration). + Exec(ctx) + if err != nil { + return fmt.Errorf("resuming released storage commit task: %w", err) + } + if rows, _ := res.RowsAffected(); rows != 1 { + return ErrConflict + } return nil case model.TaskStatusFailed: if err := tasks.RetryFailed(ctx, taskRow.ID); err != nil { diff --git a/internal/db/repository/task_repo.go b/internal/db/repository/task_repo.go index 71e5814..b1e7aaa 100644 --- a/internal/db/repository/task_repo.go +++ b/internal/db/repository/task_repo.go @@ -191,6 +191,22 @@ func (r *BunTaskRepo) GetByIdentity(ctx context.Context, taskType model.TaskType return task, nil } +func (r *BunTaskRepo) PreviousStoreCheckpoints(ctx context.Context, copyID, taskID int64) ([]model.Task, error) { + if copyID < 1 || taskID < 1 { + return nil, ErrInvalidInput + } + var tasks []model.Task + err := withTaskPayload(r.db.NewSelect().Model(&tasks)). + Where("task.type = ? AND task.status = ?", model.TaskTypeStorageStore, model.TaskStatusFailed). + Where("task.subject_type = ? AND task.subject_key = ?", "storage_copy", fmt.Sprint(copyID)). + Where("task.id <> ? AND task_payload.checkpoint_json IS NOT NULL", taskID). + OrderExpr("task.id DESC").Scan(ctx) + if err != nil { + return nil, fmt.Errorf("selecting previous Store checkpoints: %w", err) + } + return tasks, nil +} + func (r *BunTaskRepo) ClaimNext(ctx context.Context, leaseDuration time.Duration) (*model.Task, error) { if leaseDuration <= 0 { return nil, fmt.Errorf("lease duration must be positive: %w", ErrInvalidInput) @@ -315,6 +331,22 @@ func (r *BunTaskRepo) WriteCheckpoint(ctx context.Context, id, generation int64, }) } +func (r *BunTaskRepo) ConsumeStoreRetry(ctx context.Context, id, generation int64) error { + now := time.Now() + result, err := r.db.NewUpdate().Model((*model.Task)(nil)). + Set("retry_count = retry_count + 1").Set("updated_at = ?", now). + Where("id = ? AND type = ? AND status = ?", id, model.TaskTypeStorageStore, model.TaskStatusRunning). + Where("claim_generation = ? AND lease_until > ?", generation, now). + Where("retry_limit IS NULL OR retry_count < retry_limit").Exec(ctx) + if err != nil { + return fmt.Errorf("consuming Store retry: %w", err) + } + if rows, _ := result.RowsAffected(); rows != 1 { + return ErrConflict + } + return nil +} + func (r *BunTaskRepo) ValidateClaim(ctx context.Context, id, generation int64) error { now := time.Now() var found int64 @@ -471,6 +503,7 @@ func (r *BunTaskRepo) RetryFailed(ctx context.Context, id int64) error { Set("resume_mode = ?", model.TaskResumeModeRecover). Set("available_at = ?", now). Set("retry_count = 0"). + Set("retry_limit = CASE WHEN type = ? AND retry_limit = 0 THEN 1 ELSE retry_limit END", model.TaskTypeStorageStore). Set("failure_reason = NULL"). Set("last_error = NULL"). Set("status_message = NULL"). diff --git a/internal/storagecommit/advancer_test.go b/internal/storagecommit/advancer_test.go index 172a37b..6c09f5e 100644 --- a/internal/storagecommit/advancer_test.go +++ b/internal/storagecommit/advancer_test.go @@ -4,6 +4,8 @@ import ( "context" "errors" "fmt" + "strings" + "sync" "testing" "time" @@ -894,8 +896,8 @@ func TestAdvancerAttemptOnlyUsesPieceStatusAsDiagnosticEvidence(t *testing.T) { } } -func TestReleaseCommitAttentionResumesFencedFailedTask(t *testing.T) { - db := testutil.NewTestDB(t) +func seedCommitAttentionTask(t *testing.T, db *bun.DB) (*repository.Repositories, model.StorageCopy, *model.Task, string) { + t.Helper() repos := repository.NewRepositories(db) _, copies, _ := seedAdvancerCopies(t, db, 1) copyRow := copies[0] @@ -935,6 +937,11 @@ func TestReleaseCommitAttentionResumesFencedFailedTask(t *testing.T) { if err != nil || claimed == nil || claimed.ID != taskRow.ID { t.Fatalf("claim commit task = %#v err=%v", claimed, err) } + return repos, copyRow, claimed, attemptID +} + +func TestReleaseCommitAttentionResumesFencedFailedTask(t *testing.T) { + repos, copyRow, claimed, attemptID := seedCommitAttentionTask(t, testutil.NewTestDB(t)) reason := string(storagecommit.AttentionAttemptOnlyAmbiguous) message := "storage registration requires attention" if err := repos.Tasks.Settle(t.Context(), claimed.ID, claimed.ClaimGeneration, repository.TaskTransition{ @@ -949,18 +956,108 @@ func TestReleaseCommitAttentionResumesFencedFailedTask(t *testing.T) { }); err != nil { t.Fatalf("release commit attention: %v", err) } - resumed, err := repos.Tasks.GetByID(t.Context(), taskRow.ID) + resumed, err := repos.Tasks.GetByID(t.Context(), claimed.ID) if err != nil || resumed == nil || resumed.Status != model.TaskStatusPending || resumed.ResumeMode != model.TaskResumeModeRecover || resumed.RetryCount != 0 { t.Fatalf("resumed task = %#v err=%v", resumed, err) } persisted := loadAdvancerCopy(t, repos, copyRow.ID) if persisted.Status != model.StorageCopyStatusPieceReady || persisted.CommitAttemptID != nil || - persisted.ActiveTaskID == nil || *persisted.ActiveTaskID != taskRow.ID { + persisted.ActiveTaskID == nil || *persisted.ActiveTaskID != claimed.ID { t.Fatalf("released copy = %#v, want piece-ready copy fenced to resumed task", persisted) } } +func TestReleaseCommitAttentionFencesRunningTask(t *testing.T) { + repos, copyRow, claimed, attemptID := seedCommitAttentionTask(t, testutil.NewTestDB(t)) + if err := repos.Contents.ReleaseCommitAttention(t.Context(), storagecommit.ManualReleaseInput{ + CopyID: copyRow.ID, ExpectedAttemptID: attemptID, AcknowledgePossibleDuplicate: true, + }); err != nil { + t.Fatalf("release commit attention: %v", err) + } + if err := repos.Tasks.Settle(t.Context(), claimed.ID, claimed.ClaimGeneration, repository.TaskTransition{ + Status: model.TaskStatusFailed, ResumeMode: model.TaskResumeModeRecover, + }); !errors.Is(err, repository.ErrTaskLeaseLost) { + t.Fatalf("old task settlement = %v, want lost lease", err) + } + resumed, err := repos.Tasks.GetByID(t.Context(), claimed.ID) + if err != nil || resumed == nil || resumed.Status != model.TaskStatusPending || + resumed.ResumeMode != model.TaskResumeModeRecover || resumed.LeaseUntil != nil || resumed.ClaimedAt != nil { + t.Fatalf("resumed task = %#v err=%v", resumed, err) + } + persisted := loadAdvancerCopy(t, repos, copyRow.ID) + if persisted.Status != model.StorageCopyStatusPieceReady || persisted.CommitAttemptID != nil || + persisted.ActiveTaskID == nil || *persisted.ActiveTaskID != claimed.ID { + t.Fatalf("released copy = %#v", persisted) + } + next, err := repos.Tasks.ClaimNext(t.Context(), time.Minute) + if err != nil || next == nil || next.ID != claimed.ID || next.ClaimGeneration <= claimed.ClaimGeneration || + next.ResumeMode != model.TaskResumeModeRecover { + t.Fatalf("next claim = %#v err=%v", next, err) + } +} + +type commitReleaseTaskLockSignal struct { + once sync.Once + started chan struct{} +} + +type commitReleaseContextKey struct{} + +func (h *commitReleaseTaskLockSignal) BeforeQuery(ctx context.Context, event *bun.QueryEvent) context.Context { + query := strings.ToLower(event.Query) + if ctx.Value(commitReleaseContextKey{}) != nil && strings.Contains(query, "update") && + strings.Contains(query, "tasks") && strings.Contains(query, "updated_at = updated_at") { + h.once.Do(func() { close(h.started) }) + } + return ctx +} + +func (*commitReleaseTaskLockSignal) AfterQuery(context.Context, *bun.QueryEvent) {} + +func TestPostgresReleaseCommitAttentionLocksTaskBeforeStorage(t *testing.T) { + db := testutil.NewTestPostgresDB(t) + repos, copyRow, claimed, attemptID := seedCommitAttentionTask(t, db) + ctx, cancel := context.WithTimeout(t.Context(), 5*time.Second) + defer cancel() + tx, err := db.BeginTx(ctx, nil) + if err != nil { + t.Fatalf("begin settlement: %v", err) + } + defer func() { _ = tx.Rollback() }() + if err := repository.NewRepositories(tx).Tasks.ValidateClaim(ctx, claimed.ID, claimed.ClaimGeneration); err != nil { + t.Fatalf("lock settling task: %v", err) + } + signal := &commitReleaseTaskLockSignal{started: make(chan struct{})} + db.AddQueryHook(signal) + released := make(chan error, 1) + go func() { + releaseCtx := context.WithValue(ctx, commitReleaseContextKey{}, true) + released <- repos.Contents.ReleaseCommitAttention(releaseCtx, storagecommit.ManualReleaseInput{ + CopyID: copyRow.ID, ExpectedAttemptID: attemptID, AcknowledgePossibleDuplicate: true, + }) + }() + select { + case <-signal.started: + case <-ctx.Done(): + t.Fatalf("release did not try to lock the task first: %v", ctx.Err()) + } + if _, err := tx.NewRaw("UPDATE storage_contents SET updated_at = updated_at WHERE id = ?", copyRow.ContentID).Exec(ctx); err != nil { + t.Fatalf("settlement storage lock: %v", err) + } + if err := tx.Commit(); err != nil { + t.Fatalf("commit settlement: %v", err) + } + select { + case err := <-released: + if err != nil { + t.Fatalf("release after settlement lock: %v", err) + } + case <-ctx.Done(): + t.Fatalf("release deadlocked with settlement: %v", ctx.Err()) + } +} + func TestAdvancerCanceledObservationDoesNotWriteAttention(t *testing.T) { db := testutil.NewTestDB(t) repos := repository.NewRepositories(db) diff --git a/internal/synapse/readiness.go b/internal/synapse/readiness.go index 8994c75..a865bbf 100644 --- a/internal/synapse/readiness.go +++ b/internal/synapse/readiness.go @@ -322,7 +322,7 @@ func (c *ReadinessChecker) checkStorage( if err != nil { c.addPartialError(result, "storage_info", err) } - if info == nil { + if err != nil || info == nil { result.unknown("providers", "Approved storage providers could not be checked.") addStorageDependencyUnknowns(result, "Storage cost estimate could not be calculated.", "Payment funding could not be checked.", "FWSS approval could not be checked.") return diff --git a/internal/synapse/readiness_test.go b/internal/synapse/readiness_test.go index cd90f58..96170f2 100644 --- a/internal/synapse/readiness_test.go +++ b/internal/synapse/readiness_test.go @@ -200,6 +200,23 @@ func TestReadinessCheckerStorageErrorsAreUnknownAndSanitized(t *testing.T) { } } +func TestReadinessCheckerDoesNotBlockOnPartialProviderInventory(t *testing.T) { + cfg := readyReadinessConfig() + cfg.DefaultCopies = 3 + client := readyReadinessClient(2) + client.storage.infoErr = errors.New("provider lookup failed") + checker := NewReadinessChecker(cfg, client) + + got := checker.CheckRuntime(context.Background()) + + requireResultStatus(t, got, ReadinessStatusUnknown) + requireCheckStatus(t, got.Checks, "providers", ReadinessStatusUnknown) + requireCheckStatus(t, got.Checks, "fwss_approval", ReadinessStatusUnknown) + if len(client.storage.costRefs) != 0 { + t.Fatalf("cost refs = %d, want none after incomplete provider inventory", len(client.storage.costRefs)) + } +} + func TestReadinessCheckerDoesNotWarnForCancelledRequests(t *testing.T) { cfg := readyReadinessConfig() client := readyReadinessClient(1) diff --git a/internal/worker/storage_checkpoint_recovery_test.go b/internal/worker/storage_checkpoint_recovery_test.go new file mode 100644 index 0000000..e002132 --- /dev/null +++ b/internal/worker/storage_checkpoint_recovery_test.go @@ -0,0 +1,586 @@ +package worker_test + +import ( + "context" + "encoding/json" + "errors" + "io" + "log/slog" + "os" + "strings" + "sync/atomic" + "testing" + "time" + + "github.com/ipfs/go-cid" + "github.com/strahe/synaps3/internal/cache" + "github.com/strahe/synaps3/internal/db/repository" + "github.com/strahe/synaps3/internal/model" + "github.com/strahe/synaps3/internal/synapse" + taskengine "github.com/strahe/synaps3/internal/task" + "github.com/strahe/synaps3/internal/testutil" + "github.com/strahe/synaps3/internal/worker" + "github.com/strahe/synapse-go/piece" + "github.com/strahe/synapse-go/storage" + sdktypes "github.com/strahe/synapse-go/types" +) + +func storeRecoveryFixture(t *testing.T, retries int, payload string, parked synapse.ParkedPieceChecker, cacheStore cache.Cache, target *testutil.MockStorageTarget) (handlerTestRuntime, seededCopyPipeline, cid.Cid) { + t.Helper() + storageClient := &testutil.MockStorageClient{} + runtime := newHandlerTestRuntime(t, handlerRuntimeOptions{ + cache: cacheStore, storage: storageClient, parkedPieces: parked, maxRetries: &retries, + policy: cache.EvictionPolicyNone, + register: func(handlers *worker.TaskHandlers, registry *taskengine.Registry) error { + return handlers.RegisterStorage(registry) + }, + }) + pipeline := seedCopyPipeline(t, runtime, model.StorageCopyStatusPending) + version := &model.ObjectVersion{ + VersionID: model.NewVersionID(), BucketID: pipeline.upload.BucketID, + Key: "store-recovery.bin", ContentID: &pipeline.upload.ID, + Size: 128, ETag: "store-recovery", ContentType: "application/octet-stream", + } + if _, err := runtime.repos.Objects.CreateVersionAndSetCurrent(t.Context(), version); err != nil { + t.Fatal(err) + } + info, err := piece.Calculate(strings.NewReader(payload)) + if err != nil { + t.Fatal(err) + } + if _, err := runtime.db.NewUpdate().Model((*model.StorageCopy)(nil)). + Set("transfer_method = ?", model.StorageCopyTransferMethodCacheRestore). + Where("id = ?", pipeline.target.ID).Exec(t.Context()); err != nil { + t.Fatal(err) + } + if _, err := runtime.db.NewUpdate().Model((*model.StorageContent)(nil)). + Set("piece_cid = ?", info.CIDv2.String()).Where("id = ?", pipeline.upload.ID).Exec(t.Context()); err != nil { + t.Fatal(err) + } + target.ProviderIDValue = pipeline.targetSet.ProviderID.SDK() + targetDataSetID := pipeline.targetSet.DataSetID.SDK() + target.DataSetIDValue = &targetDataSetID + target.ClientDataSetIDValue = pipeline.targetClient + storageClient.OpenDataSetTargetFunc = func(context.Context, sdktypes.BigInt, storage.NewDataSetContextOptions) (synapse.DataSetTarget, error) { + return target, nil + } + return runtime, pipeline, info.CIDv2 +} + +func seedStoreCheckpoint(t *testing.T, runtime handlerTestRuntime, taskID, copyID int64, pieceCID cid.Cid, providerURL string, attemptedAt time.Time) { + t.Helper() + checkpoint, err := json.Marshal(map[string]any{ + "attempted_at": attemptedAt.UTC(), "intended_piece_cid": pieceCID.String(), + "provider_service_url": providerURL, "ingress_attempt": 1, + }) + if err != nil { + t.Fatal(err) + } + if _, err := runtime.db.NewRaw(`UPDATE task_payloads SET checkpoint_json = ? WHERE task_id = ?`, checkpoint, taskID).Exec(t.Context()); err != nil { + t.Fatal(err) + } + if _, err := runtime.db.NewRaw(`UPDATE tasks SET resume_mode = 'recover' WHERE id = ?`, taskID).Exec(t.Context()); err != nil { + t.Fatal(err) + } + if _, err := runtime.db.NewRaw(`UPDATE storage_copies SET ingress_store_attempt = 1 WHERE id = ?`, copyID).Exec(t.Context()); err != nil { + t.Fatal(err) + } +} + +func TestStoreProcessingAndQueryFailureLeaveRetryableCheckpoint(t *testing.T) { + for _, tt := range []struct { + name string + state synapse.ParkedPieceState + err error + reason string + }{ + {name: "processing", state: synapse.ParkedPieceProcessing, reason: "store_processing_timeout"}, + {name: "query failure", err: errors.New("provider unavailable"), reason: "store_check_failed"}, + } { + t.Run(tt.name, func(t *testing.T) { + payload := strings.Repeat("k", 128) + target := &testutil.MockStorageTarget{ServiceURLValue: "https://store.example"} + parked := parkedPieceCheckerFunc(func(context.Context, string, cid.Cid) (synapse.ParkedPieceState, error) { + return tt.state, tt.err + }) + runtime, pipeline, pieceCID := storeRecoveryFixture(t, 2, payload, parked, nil, target) + taskRow := bindCopyTask(t, runtime, pipeline.target, model.TaskTypeStorageStore) + seedStoreCheckpoint(t, runtime, taskRow.ID, pipeline.target.ID, pieceCID, target.ServiceURL(), time.Now().Add(-31*time.Minute)) + cancel, done := runHandlerEngine(t, runtime) + defer stopHandlerEngine(t, cancel, done) + failed := waitForTask(t, runtime.repos, taskRow.ID, func(task *model.Task) bool { return task.Status == model.TaskStatusFailed }) + if failed.FailureReason == nil || *failed.FailureReason != tt.reason || !runtime.service.Retryable(failed) || len(failed.Checkpoint) == 0 { + t.Fatalf("failed Store task = %#v", failed) + } + copyRow, err := runtime.repos.Contents.GetUploadCopyByID(t.Context(), pipeline.target.ID) + if err != nil || copyRow.Status != model.StorageCopyStatusPending || copyRow.ActiveTaskID == nil || *copyRow.ActiveTaskID != taskRow.ID { + t.Fatalf("Store copy after timeout = %#v, err=%v", copyRow, err) + } + }) + } +} + +func TestStoreManualRetryAllowsOneUploadWhenAutomaticRetriesDisabled(t *testing.T) { + payload := strings.Repeat("r", 128) + var stores atomic.Int64 + cacheStore := &testutil.MockCache{GetFunc: func(context.Context, string, string) (io.ReadCloser, *cache.ObjectInfo, error) { + return io.NopCloser(strings.NewReader(payload)), &cache.ObjectInfo{Size: 128}, nil + }} + parked := parkedPieceCheckerFunc(func(context.Context, string, cid.Cid) (synapse.ParkedPieceState, error) { + return synapse.ParkedPieceMissing, nil + }) + target := &testutil.MockStorageTarget{ServiceURLValue: "https://store.example", StoreFunc: func(context.Context, io.Reader, *storage.StoreOptions) (*storage.StoreResult, error) { + stores.Add(1) + return nil, errors.New("connection closed after upload") + }} + runtime, pipeline, _ := storeRecoveryFixture(t, 0, payload, parked, cacheStore, target) + taskRow := bindCopyTask(t, runtime, pipeline.target, model.TaskTypeStorageStore) + cancel, done := runHandlerEngine(t, runtime) + defer stopHandlerEngine(t, cancel, done) + waitForTask(t, runtime.repos, taskRow.ID, func(task *model.Task) bool { + return task.Status == model.TaskStatusPending && len(task.Checkpoint) > 0 + }) + if _, err := runtime.db.NewRaw(`UPDATE tasks SET available_at = ? WHERE id = ?`, time.Now().Add(-time.Second), taskRow.ID).Exec(t.Context()); err != nil { + t.Fatal(err) + } + failed := waitForTask(t, runtime.repos, taskRow.ID, func(task *model.Task) bool { return task.Status == model.TaskStatusFailed }) + if failed.FailureReason == nil || *failed.FailureReason != "store_retry_limit" || !runtime.service.Retryable(failed) || stores.Load() != 1 { + t.Fatalf("zero-retry Store = task:%#v uploads:%d", failed, stores.Load()) + } + if err := runtime.service.Retry(t.Context(), taskRow.ID); err != nil { + t.Fatal(err) + } + retrying := waitForTask(t, runtime.repos, taskRow.ID, func(task *model.Task) bool { + return task.Status == model.TaskStatusPending && task.RetryCount == 1 && stores.Load() == 2 + }) + if retrying.RetryLimit == nil || *retrying.RetryLimit != 1 { + t.Fatalf("manual Store retry limit = %v, want 1", retrying.RetryLimit) + } +} + +func TestStoreStopsAfterConfiguredAutomaticRetransmissions(t *testing.T) { + payload := strings.Repeat("l", 128) + var stores atomic.Int64 + cacheStore := &testutil.MockCache{GetFunc: func(context.Context, string, string) (io.ReadCloser, *cache.ObjectInfo, error) { + return io.NopCloser(strings.NewReader(payload)), &cache.ObjectInfo{Size: 128}, nil + }} + parked := parkedPieceCheckerFunc(func(context.Context, string, cid.Cid) (synapse.ParkedPieceState, error) { + return synapse.ParkedPieceMissing, nil + }) + target := &testutil.MockStorageTarget{ServiceURLValue: "https://store.example", StoreFunc: func(context.Context, io.Reader, *storage.StoreOptions) (*storage.StoreResult, error) { + stores.Add(1) + return nil, errors.New("connection closed after upload") + }} + runtime, pipeline, _ := storeRecoveryFixture(t, 1, payload, parked, cacheStore, target) + taskRow := bindCopyTask(t, runtime, pipeline.target, model.TaskTypeStorageStore) + cancel, done := runHandlerEngine(t, runtime) + defer stopHandlerEngine(t, cancel, done) + waitForTask(t, runtime.repos, taskRow.ID, func(task *model.Task) bool { + return task.Status == model.TaskStatusPending && stores.Load() == 1 + }) + if _, err := runtime.db.NewRaw(`UPDATE tasks SET available_at = ? WHERE id = ?`, time.Now().Add(-time.Second), taskRow.ID).Exec(t.Context()); err != nil { + t.Fatal(err) + } + waitForTask(t, runtime.repos, taskRow.ID, func(task *model.Task) bool { + return task.Status == model.TaskStatusPending && stores.Load() == 2 && task.RetryCount == 1 + }) + if _, err := runtime.db.NewRaw(`UPDATE tasks SET available_at = ? WHERE id = ?`, time.Now().Add(-time.Second), taskRow.ID).Exec(t.Context()); err != nil { + t.Fatal(err) + } + failed := waitForTask(t, runtime.repos, taskRow.ID, func(task *model.Task) bool { return task.Status == model.TaskStatusFailed }) + if failed.FailureReason == nil || *failed.FailureReason != "store_retry_limit" || !runtime.service.Retryable(failed) || stores.Load() != 2 { + t.Fatalf("Store retry limit = task:%#v uploads:%d", failed, stores.Load()) + } + copyRow, err := runtime.repos.Contents.GetUploadCopyByID(t.Context(), pipeline.target.ID) + if err != nil || copyRow.ActiveTaskID == nil || *copyRow.ActiveTaskID != taskRow.ID { + t.Fatalf("Store copy lost its retryable task: %#v, err:%v", copyRow, err) + } +} + +func TestNewStoreTaskAdoptsPreviousCheckpoint(t *testing.T) { + payload := strings.Repeat("a", 128) + var storeCalls atomic.Int64 + parked := parkedPieceCheckerFunc(func(context.Context, string, cid.Cid) (synapse.ParkedPieceState, error) { + return synapse.ParkedPieceReady, nil + }) + target := &testutil.MockStorageTarget{ + ServiceURLValue: "https://store.example", + StoreFunc: func(context.Context, io.Reader, *storage.StoreOptions) (*storage.StoreResult, error) { + storeCalls.Add(1) + return nil, errors.New("unexpected upload") + }, + PresignForCommitFunc: func(context.Context, []storage.PieceInput) ([]byte, error) { return []byte{0xaa}, nil }, + } + runtime, pipeline, pieceCID := storeRecoveryFixture(t, 2, payload, parked, nil, target) + old := bindCopyTask(t, runtime, pipeline.target, model.TaskTypeStorageStore) + seedStoreCheckpoint(t, runtime, old.ID, pipeline.target.ID, pieceCID, target.ServiceURL(), time.Now()) + if _, err := runtime.db.NewRaw(`UPDATE tasks SET status = 'failed', finished_at = ?, resume_mode = 'recover' WHERE id = ?`, time.Now(), old.ID).Exec(t.Context()); err != nil { + t.Fatal(err) + } + if _, err := runtime.db.NewRaw(`UPDATE storage_copies SET status = 'failed', active_task_id = NULL WHERE id = ?`, pipeline.target.ID).Exec(t.Context()); err != nil { + t.Fatal(err) + } + if err := runtime.repos.Contents.ReopenFailedUploadCopy(t.Context(), pipeline.target.ID); err != nil { + t.Fatal(err) + } + newTask := bindCopyTask(t, runtime, pipeline.target, model.TaskTypeStorageStore) + cancel, done := runHandlerEngine(t, runtime) + defer stopHandlerEngine(t, cancel, done) + completed := waitForTask(t, runtime.repos, newTask.ID, func(task *model.Task) bool { return task.Status == model.TaskStatusCompleted }) + if len(completed.Checkpoint) == 0 { + t.Fatal("new Store task did not retain the previous checkpoint") + } + copyRow, err := runtime.repos.Contents.GetUploadCopyByID(t.Context(), pipeline.target.ID) + if err != nil || copyRow.Status == model.StorageCopyStatusPending || copyRow.Status == model.StorageCopyStatusFailed || storeCalls.Load() != 0 { + t.Fatalf("adopted Store copy = %#v, uploads:%d, err:%v", copyRow, storeCalls.Load(), err) + } +} + +func TestNewStoreTaskWithoutCheckpointNeedsReadableCache(t *testing.T) { + payload := strings.Repeat("m", 128) + cacheStore := &testutil.MockCache{GetFunc: func(context.Context, string, string) (io.ReadCloser, *cache.ObjectInfo, error) { + return nil, nil, os.ErrNotExist + }} + target := &testutil.MockStorageTarget{ServiceURLValue: "https://store.example"} + runtime, pipeline, _ := storeRecoveryFixture(t, 2, payload, nil, cacheStore, target) + taskRow := bindCopyTask(t, runtime, pipeline.target, model.TaskTypeStorageStore) + if _, err := runtime.db.NewRaw(`UPDATE storage_copies SET ingress_store_attempt = 1 WHERE id = ?`, pipeline.target.ID).Exec(t.Context()); err != nil { + t.Fatal(err) + } + cancel, done := runHandlerEngine(t, runtime) + defer stopHandlerEngine(t, cancel, done) + failed := waitForTask(t, runtime.repos, taskRow.ID, func(task *model.Task) bool { return task.Status == model.TaskStatusFailed }) + if failed.FailureReason == nil || *failed.FailureReason != "store_cache_missing" || len(failed.Checkpoint) != 0 { + t.Fatalf("uncached Store task = %#v", failed) + } +} + +func TestNewStoreTaskWithoutCheckpointChecksProviderBeforeUploading(t *testing.T) { + payload := strings.Repeat("n", 128) + var stores, checks atomic.Int64 + cacheStore := &testutil.MockCache{GetFunc: func(context.Context, string, string) (io.ReadCloser, *cache.ObjectInfo, error) { + return io.NopCloser(strings.NewReader(payload)), &cache.ObjectInfo{Size: 128}, nil + }} + parked := parkedPieceCheckerFunc(func(context.Context, string, cid.Cid) (synapse.ParkedPieceState, error) { + checks.Add(1) + return synapse.ParkedPieceMissing, nil + }) + target := &testutil.MockStorageTarget{ + ServiceURLValue: "https://store.example", + StoreFunc: func(_ context.Context, _ io.Reader, options *storage.StoreOptions) (*storage.StoreResult, error) { + stores.Add(1) + return &storage.StoreResult{PieceCID: options.PieceCID, Size: 128}, nil + }, + PresignForCommitFunc: func(context.Context, []storage.PieceInput) ([]byte, error) { return []byte{0xaa}, nil }, + } + runtime, pipeline, _ := storeRecoveryFixture(t, 2, payload, parked, cacheStore, target) + taskRow := bindCopyTask(t, runtime, pipeline.target, model.TaskTypeStorageStore) + if _, err := runtime.db.NewRaw(`UPDATE storage_copies SET ingress_store_attempt = 1 WHERE id = ?`, pipeline.target.ID).Exec(t.Context()); err != nil { + t.Fatal(err) + } + cancel, done := runHandlerEngine(t, runtime) + defer stopHandlerEngine(t, cancel, done) + completed := waitForTask(t, runtime.repos, taskRow.ID, func(task *model.Task) bool { return task.Status == model.TaskStatusCompleted }) + if stores.Load() != 1 || checks.Load() == 0 || completed.RetryCount != 1 || len(completed.Checkpoint) == 0 { + t.Fatalf("Store without old checkpoint = task:%#v uploads:%d checks:%d", completed, stores.Load(), checks.Load()) + } +} + +func TestNewStoreTaskKeepsRetryAfterReadyPiecePresignFailure(t *testing.T) { + payload := strings.Repeat("s", 128) + var presigns, stores atomic.Int64 + cacheStore := &testutil.MockCache{GetFunc: func(context.Context, string, string) (io.ReadCloser, *cache.ObjectInfo, error) { + return io.NopCloser(strings.NewReader(payload)), &cache.ObjectInfo{Size: 128}, nil + }} + parked := parkedPieceCheckerFunc(func(context.Context, string, cid.Cid) (synapse.ParkedPieceState, error) { + return synapse.ParkedPieceReady, nil + }) + target := &testutil.MockStorageTarget{ + ServiceURLValue: "https://store.example", + StoreFunc: func(context.Context, io.Reader, *storage.StoreOptions) (*storage.StoreResult, error) { + stores.Add(1) + return nil, errors.New("unexpected upload") + }, + PresignForCommitFunc: func(context.Context, []storage.PieceInput) ([]byte, error) { + if presigns.Add(1) == 1 { + return nil, errors.New("temporary presign failure") + } + return []byte{0xaa}, nil + }, + } + runtime, pipeline, _ := storeRecoveryFixture(t, 2, payload, parked, cacheStore, target) + taskRow := bindCopyTask(t, runtime, pipeline.target, model.TaskTypeStorageStore) + if _, err := runtime.db.NewRaw(`UPDATE storage_copies SET ingress_store_attempt = 1 WHERE id = ?`, pipeline.target.ID).Exec(t.Context()); err != nil { + t.Fatal(err) + } + cancel, done := runHandlerEngine(t, runtime) + defer stopHandlerEngine(t, cancel, done) + failed := waitForTask(t, runtime.repos, taskRow.ID, func(task *model.Task) bool { return task.Status == model.TaskStatusFailed }) + var checkpoint struct { + IngressAttempt int `json:"ingress_attempt"` + } + if failed.FailureReason == nil || *failed.FailureReason != "commit_presign_failed" || + !runtime.service.Retryable(failed) || json.Unmarshal(failed.Checkpoint, &checkpoint) != nil || + checkpoint.IngressAttempt != 1 || failed.RetryCount != 0 || stores.Load() != 0 { + t.Fatalf("ready piece after presign failure = task:%#v checkpoint:%#v uploads:%d", failed, checkpoint, stores.Load()) + } + copyRow, err := runtime.repos.Contents.GetUploadCopyByID(t.Context(), pipeline.target.ID) + if err != nil || copyRow.Status != model.StorageCopyStatusPending || copyRow.ActiveTaskID == nil || *copyRow.ActiveTaskID != taskRow.ID { + t.Fatalf("ready piece copy after presign failure = %#v, err:%v", copyRow, err) + } + if err := runtime.service.Retry(t.Context(), taskRow.ID); err != nil { + t.Fatal(err) + } + waitForTask(t, runtime.repos, taskRow.ID, func(task *model.Task) bool { return task.Status == model.TaskStatusCompleted }) + if presigns.Load() != 2 || stores.Load() != 0 { + t.Fatalf("manual retry = presigns:%d uploads:%d, want 2/0", presigns.Load(), stores.Load()) + } +} + +func TestNewStoreTaskWaitsForUnavailableProviderAfterEarlierUpload(t *testing.T) { + payload := strings.Repeat("u", 128) + var providerReady atomic.Bool + cacheStore := &testutil.MockCache{GetFunc: func(context.Context, string, string) (io.ReadCloser, *cache.ObjectInfo, error) { + return io.NopCloser(strings.NewReader(payload)), &cache.ObjectInfo{Size: 128}, nil + }} + parked := parkedPieceCheckerFunc(func(context.Context, string, cid.Cid) (synapse.ParkedPieceState, error) { + return synapse.ParkedPieceReady, nil + }) + target := &testutil.MockStorageTarget{ + ServiceURLValue: "https://store.example", + PresignForCommitFunc: func(context.Context, []storage.PieceInput) ([]byte, error) { return []byte{0xaa}, nil }, + } + runtime, pipeline, _ := storeRecoveryFixture(t, 0, payload, parked, cacheStore, target) + runtime.storage.OpenDataSetTargetFunc = func(context.Context, sdktypes.BigInt, storage.NewDataSetContextOptions) (synapse.DataSetTarget, error) { + if !providerReady.Load() { + return nil, storage.ErrDataSetUnavailable + } + return target, nil + } + taskRow := bindCopyTask(t, runtime, pipeline.target, model.TaskTypeStorageStore) + if _, err := runtime.db.NewRaw(`UPDATE storage_copies SET ingress_store_attempt = 1 WHERE id = ?`, pipeline.target.ID).Exec(t.Context()); err != nil { + t.Fatal(err) + } + cancel, done := runHandlerEngine(t, runtime) + defer stopHandlerEngine(t, cancel, done) + waiting := waitForTask(t, runtime.repos, taskRow.ID, func(task *model.Task) bool { + return task.Status == model.TaskStatusPending && task.WaitReason != nil && *task.WaitReason == "provider" + }) + if waiting.RetryCount != 0 || len(waiting.Checkpoint) != 0 { + t.Fatalf("unavailable provider consumed retry or checkpointed Store: %#v", waiting) + } + providerReady.Store(true) + if _, err := runtime.db.NewRaw(`UPDATE tasks SET available_at = ? WHERE id = ?`, time.Now().Add(-time.Second), taskRow.ID).Exec(t.Context()); err != nil { + t.Fatal(err) + } + waitForTask(t, runtime.repos, taskRow.ID, func(task *model.Task) bool { return task.Status == model.TaskStatusCompleted }) +} + +func TestStoreDoesNotUseOldProviderPieceForChangedTarget(t *testing.T) { + payload := strings.Repeat("p", 128) + var stores, oldChecks, currentChecks atomic.Int64 + cacheStore := &testutil.MockCache{GetFunc: func(context.Context, string, string) (io.ReadCloser, *cache.ObjectInfo, error) { + return io.NopCloser(strings.NewReader(payload)), &cache.ObjectInfo{Size: 128}, nil + }} + parked := parkedPieceCheckerFunc(func(_ context.Context, serviceURL string, _ cid.Cid) (synapse.ParkedPieceState, error) { + if serviceURL == "https://old.example" { + oldChecks.Add(1) + return synapse.ParkedPieceReady, nil + } + currentChecks.Add(1) + return synapse.ParkedPieceMissing, nil + }) + target := &testutil.MockStorageTarget{ + ServiceURLValue: "https://current.example", + StoreFunc: func(_ context.Context, _ io.Reader, options *storage.StoreOptions) (*storage.StoreResult, error) { + stores.Add(1) + return &storage.StoreResult{PieceCID: options.PieceCID, Size: 128}, nil + }, + PresignForCommitFunc: func(context.Context, []storage.PieceInput) ([]byte, error) { return []byte{0xaa}, nil }, + } + runtime, pipeline, pieceCID := storeRecoveryFixture(t, 2, payload, parked, cacheStore, target) + taskRow := bindCopyTask(t, runtime, pipeline.target, model.TaskTypeStorageStore) + seedStoreCheckpoint(t, runtime, taskRow.ID, pipeline.target.ID, pieceCID, "https://old.example", time.Now()) + cancel, done := runHandlerEngine(t, runtime) + defer stopHandlerEngine(t, cancel, done) + completed := waitForTask(t, runtime.repos, taskRow.ID, func(task *model.Task) bool { return task.Status == model.TaskStatusCompleted }) + copyRow, err := runtime.repos.Contents.GetUploadCopyByID(t.Context(), pipeline.target.ID) + if err != nil || copyRow.Status == model.StorageCopyStatusPending || copyRow.Status == model.StorageCopyStatusFailed || stores.Load() != 1 || oldChecks.Load() == 0 || currentChecks.Load() == 0 || completed.RetryCount != 1 { + t.Fatalf("changed provider Store = copy:%#v task:%#v uploads:%d old checks:%d current checks:%d err:%v", copyRow, completed, stores.Load(), oldChecks.Load(), currentChecks.Load(), err) + } +} + +func TestStoreRestartBeforeRetryCheckpointRechecksOldPiece(t *testing.T) { + payload := strings.Repeat("b", 128) + var stores atomic.Int64 + cacheStore := &testutil.MockCache{GetFunc: func(context.Context, string, string) (io.ReadCloser, *cache.ObjectInfo, error) { + return io.NopCloser(strings.NewReader(payload)), &cache.ObjectInfo{Size: 128}, nil + }} + parked := parkedPieceCheckerFunc(func(context.Context, string, cid.Cid) (synapse.ParkedPieceState, error) { + return synapse.ParkedPieceMissing, nil + }) + target := &testutil.MockStorageTarget{ServiceURLValue: "https://store.example", StoreFunc: func(context.Context, io.Reader, *storage.StoreOptions) (*storage.StoreResult, error) { + stores.Add(1) + return nil, errors.New("transfer interrupted") + }} + runtime, pipeline, pieceCID := storeRecoveryFixture(t, 2, payload, parked, cacheStore, target) + taskRow := bindCopyTask(t, runtime, pipeline.target, model.TaskTypeStorageStore) + seedStoreCheckpoint(t, runtime, taskRow.ID, pipeline.target.ID, pieceCID, target.ServiceURL(), time.Now()) + limitedRepos := *runtime.repos + limitedRepos.Tasks = &limitedClaimRepository{TaskRepository: runtime.repos.Tasks, maximum: 1} + firstEngine, err := taskengine.NewEngine(taskengine.EngineConfig{ + Concurrency: 1, PollInterval: 5 * time.Millisecond, LeaseDuration: 300 * time.Millisecond, + Retention: time.Hour, ProviderMutationConcurrency: 4, DestructiveMutationConcurrency: 2, + }, &limitedRepos, runtime.registry, slog.Default()) + if err != nil { + t.Fatal(err) + } + cancel, done := runEngine(t, firstEngine) + before := waitForTask(t, runtime.repos, taskRow.ID, func(task *model.Task) bool { + return task.Status == model.TaskStatusPending && task.ResumeMode == model.TaskResumeModeExecute + }) + stopHandlerEngine(t, cancel, done) + var checkpoint struct { + IngressAttempt int `json:"ingress_attempt"` + } + if err := json.Unmarshal(before.Checkpoint, &checkpoint); err != nil || checkpoint.IngressAttempt != 1 || stores.Load() != 0 { + t.Fatalf("before new checkpoint = attempt:%d uploads:%d err:%v", checkpoint.IngressAttempt, stores.Load(), err) + } + cancel, done = runHandlerEngine(t, runtime) + defer stopHandlerEngine(t, cancel, done) + after := waitForTask(t, runtime.repos, taskRow.ID, func(task *model.Task) bool { + return task.Status == model.TaskStatusPending && task.ResumeMode == model.TaskResumeModeRecover && stores.Load() == 1 + }) + if err := json.Unmarshal(after.Checkpoint, &checkpoint); err != nil || checkpoint.IngressAttempt != 2 || after.RetryCount != 1 { + t.Fatalf("after new checkpoint = attempt:%d retry:%d err:%v", checkpoint.IngressAttempt, after.RetryCount, err) + } +} + +func TestStoreRestartAfterRetryCheckpointOnlyQueriesProvider(t *testing.T) { + payload := strings.Repeat("c", 128) + var stores atomic.Int64 + var parkedState atomic.Value + parkedState.Store(synapse.ParkedPieceMissing) + enteredStore := make(chan struct{}) + cacheStore := &testutil.MockCache{GetFunc: func(context.Context, string, string) (io.ReadCloser, *cache.ObjectInfo, error) { + return io.NopCloser(strings.NewReader(payload)), &cache.ObjectInfo{Size: 128}, nil + }} + parked := parkedPieceCheckerFunc(func(context.Context, string, cid.Cid) (synapse.ParkedPieceState, error) { + return parkedState.Load().(synapse.ParkedPieceState), nil + }) + target := &testutil.MockStorageTarget{ + ServiceURLValue: "https://store.example", + StoreFunc: func(ctx context.Context, _ io.Reader, _ *storage.StoreOptions) (*storage.StoreResult, error) { + if stores.Add(1) == 1 { + close(enteredStore) + } + <-ctx.Done() + return nil, ctx.Err() + }, + PresignForCommitFunc: func(context.Context, []storage.PieceInput) ([]byte, error) { return []byte{0xaa}, nil }, + } + runtime, pipeline, pieceCID := storeRecoveryFixture(t, 2, payload, parked, cacheStore, target) + taskRow := bindCopyTask(t, runtime, pipeline.target, model.TaskTypeStorageStore) + seedStoreCheckpoint(t, runtime, taskRow.ID, pipeline.target.ID, pieceCID, target.ServiceURL(), time.Now()) + cancel, done := runHandlerEngine(t, runtime) + select { + case <-enteredStore: + case <-time.After(3 * time.Second): + stopHandlerEngine(t, cancel, done) + t.Fatal("Store did not start after the new checkpoint") + } + inFlight, err := runtime.repos.Tasks.GetByID(t.Context(), taskRow.ID) + if err != nil { + t.Fatal(err) + } + var checkpoint struct { + IngressAttempt int `json:"ingress_attempt"` + } + if err := json.Unmarshal(inFlight.Checkpoint, &checkpoint); err != nil || checkpoint.IngressAttempt != 2 || inFlight.RetryCount != 1 { + t.Fatalf("in-flight checkpoint = attempt:%d retry:%d err:%v", checkpoint.IngressAttempt, inFlight.RetryCount, err) + } + stopHandlerEngine(t, cancel, done) + parkedState.Store(synapse.ParkedPieceReady) + time.Sleep(350 * time.Millisecond) + restarted, err := taskengine.NewEngine(taskengine.EngineConfig{ + Concurrency: 1, PollInterval: 5 * time.Millisecond, LeaseDuration: 300 * time.Millisecond, + Retention: time.Hour, ProviderMutationConcurrency: 4, DestructiveMutationConcurrency: 2, + }, runtime.repos, runtime.registry, slog.Default()) + if err != nil { + t.Fatal(err) + } + cancel, done = runEngine(t, restarted) + defer stopHandlerEngine(t, cancel, done) + waitForTask(t, runtime.repos, taskRow.ID, func(task *model.Task) bool { return task.Status == model.TaskStatusCompleted }) + if stores.Load() != 1 { + t.Fatalf("Store calls after restart = %d, want 1", stores.Load()) + } +} + +func TestPeerPullWaitsUntilStoreCopyIsReadable(t *testing.T) { + runtime := newHandlerTestRuntime(t, handlerRuntimeOptions{ + register: func(handlers *worker.TaskHandlers, registry *taskengine.Registry) error { + return handlers.RegisterStorage(registry) + }, + }) + bucket := &model.Bucket{Name: "peer-source-wait", Status: model.BucketStatusActive, DefaultCopies: 2, MinimumDurableCopies: 1} + if err := runtime.repos.Buckets.Create(t.Context(), bucket); err != nil { + t.Fatal(err) + } + content, err := runtime.repos.Contents.EnsureContent(t.Context(), repository.EnsureContentInput{ + BucketID: bucket.ID, ContentSize: 128, Checksum: testutil.StorageChecksum("peer-source-wait"), RequestedCopies: 2, + }) + if err != nil { + t.Fatal(err) + } + version := &model.ObjectVersion{ + VersionID: model.NewVersionID(), BucketID: bucket.ID, Key: "peer.bin", + ContentID: &content.ID, Size: 128, ETag: "peer", ContentType: "application/octet-stream", + } + if _, err := runtime.repos.Objects.CreateVersionAndSetCurrent(t.Context(), version); err != nil { + t.Fatal(err) + } + bindings := make([]*model.StorageDataSet, 2) + for i := range bindings { + providerID := testOnChainID(t, int64(61001+i)) + binding, err := runtime.repos.Contents.EnsureDataSetBinding(t.Context(), repository.EnsureDataSetBindingInput{ + BucketID: bucket.ID, ProviderID: providerID, CopyIndex: i, CreatedByContentID: content.ID, + }) + if err != nil { + t.Fatal(err) + } + dataSetID := testOnChainID(t, int64(62001+i)) + clientID := testOnChainID(t, int64(63001+i)) + if err := runtime.repos.Contents.MarkDataSetReady(t.Context(), repository.MarkDataSetReadyInput{ + ID: binding.ID, ContentID: content.ID, DataSetID: dataSetID, ClientDataSetID: &clientID, + }); err != nil { + t.Fatal(err) + } + bindings[i] = binding + } + if err := runtime.repos.Contents.CreateUploadCopiesForBindings(t.Context(), content.ID, []repository.UploadCopyBindingInput{ + {StorageDataSetID: bindings[0].ID, CopyIndex: 0, ProviderID: bindings[0].ProviderID, TransferMethod: model.StorageCopyTransferMethodIngress}, + {StorageDataSetID: bindings[1].ID, CopyIndex: 1, ProviderID: bindings[1].ProviderID, TransferMethod: model.StorageCopyTransferMethodPeerPull}, + }); err != nil { + t.Fatal(err) + } + copies, err := runtime.repos.Contents.ListCopies(t.Context(), content.ID) + if err != nil || len(copies) != 2 { + t.Fatalf("copies = %#v, err:%v", copies, err) + } + peerTask := bindCopyTask(t, runtime, &copies[1], model.TaskTypeStorageTransferPlan) + cancel, done := runHandlerEngine(t, runtime) + defer stopHandlerEngine(t, cancel, done) + waitForTask(t, runtime.repos, peerTask.ID, func(task *model.Task) bool { + return task.Status == model.TaskStatusPending && task.WaitReason != nil && *task.WaitReason == "source" + }) + pieceID := testOnChainID(t, 64001) + testutil.CommitStorageCopy(t, runtime.db, runtime.repos, repository.MarkUploadCopyCommittedInput{ + StorageCopyID: copies[0].ID, ContentID: content.ID, CopyIndex: 0, + PieceCID: testPieceCID(t, "peer-source-wait").String(), PieceID: &pieceID, RetrievalURL: "https://source.example/piece", + }) + if _, err := runtime.db.NewRaw(`UPDATE tasks SET available_at = ? WHERE id = ?`, time.Now().Add(-time.Second), peerTask.ID).Exec(t.Context()); err != nil { + t.Fatal(err) + } + waitForTask(t, runtime.repos, peerTask.ID, func(task *model.Task) bool { return task.Status == model.TaskStatusCompleted }) +} diff --git a/internal/worker/storage_task_handlers.go b/internal/worker/storage_task_handlers.go index db2a122..5f26d0d 100644 --- a/internal/worker/storage_task_handlers.go +++ b/internal/worker/storage_task_handlers.go @@ -4,6 +4,7 @@ import ( "context" "crypto/rand" "encoding/hex" + "encoding/json" "errors" "fmt" "math/big" @@ -1015,7 +1016,9 @@ func (h *TaskHandlers) storeHandler() taskengine.Handler { switch *task.FailureReason { case "store_not_started": return len(task.Checkpoint) == 0 - case "store_outcome_unknown", "copy_owner_missing", "copy_context_failed": + case "store_retry_limit", "store_cache_missing", "store_checkpoint_write_failed", "copy_authorization_failed": + return true + case "store_outcome_unknown", "store_processing_timeout", "store_check_failed", "store_result_invalid", "store_cache_read_failed", "store_identity_mismatch", "commit_presign_failed", "copy_owner_missing", "copy_context_failed": return len(task.Checkpoint) > 0 default: return false @@ -1049,16 +1052,73 @@ func (h *TaskHandlers) runStore(ctx context.Context, execution taskengine.Execut } _, target, content, bucket, err := h.copyContext(ctx, copyRow) if err != nil { + if !hasCheckpoint && copyRow.IngressStoreAttempt > 0 { + if synapse.IsProviderUnavailable(err) || errors.Is(err, storage.ErrDataSetUnavailable) { + return h.copyContextFailure(execution, input, copyRow, err, false) + } + return h.retryStoreNotStarted(execution, err) + } return h.copyContextFailure(execution, input, copyRow, err, !hasCheckpoint) } - if hasCheckpoint { + if hasCheckpoint && !mayStore { return h.recoverStore(ctx, execution, input, copyRow, target, checkpoint) } if !mayStore { return taskengine.Suspend(model.TaskResumeModeExecute, 0, "safe_to_execute", "Storage transfer is ready", nil) } + if hasCheckpoint { + if err := validateStoreCheckpoint(checkpoint, copyRow.ContentSize); err != nil { + return taskengine.Fail(err, "invalid_checkpoint", nil) + } + pieceCID, err := cid.Parse(checkpoint.IntendedPieceCID) + if err != nil { + return taskengine.Fail(err, "invalid_checkpoint", nil) + } + state, checkErr := h.storePieceState(ctx, target, checkpoint, pieceCID) + if checkErr != nil || state == synapse.ParkedPieceProcessing { + return h.waitForStoreStatus(checkpoint, state, checkErr) + } + if state == synapse.ParkedPieceReady { + return h.finishPieceTransfer(ctx, execution, input, copyRow, target, pieceCID) + } + if execution.RetryWillFail() { + return taskengine.Fail(errors.New("storage transfer retry limit reached"), "store_retry_limit", nil) + } + } + if !hasCheckpoint && copyRow.IngressStoreAttempt > 0 { + previous, err := h.deps.Repositories.Tasks.PreviousStoreCheckpoints(ctx, copyRow.ID, execution.ID()) + if err != nil { + return h.retryStoreNotStarted(execution, err) + } + for _, old := range previous { + var oldInput storagepipeline.CopyGenerationInput + var candidate storeCheckpoint + if err := json.Unmarshal(old.Checkpoint, &candidate); err != nil || candidate.IngressAttempt != copyRow.IngressStoreAttempt { + continue + } + if err := json.Unmarshal(old.Input, &oldInput); err != nil || oldInput.CopyID != copyRow.ID || oldInput.Generation >= input.Generation { + return taskengine.Fail(errors.New("previous Store task identity conflicts with this copy"), "store_checkpoint_conflict", nil) + } + if err := validateStoreCheckpoint(candidate, copyRow.ContentSize); err != nil { + return taskengine.Fail(err, "invalid_checkpoint", nil) + } + if hasCheckpoint && checkpoint != candidate { + return taskengine.Fail(errors.New("previous Store checkpoints disagree"), "store_checkpoint_conflict", nil) + } + checkpoint, hasCheckpoint = candidate, true + } + if hasCheckpoint { + if err := execution.WriteCheckpoint(ctx, checkpoint); err != nil { + return h.retryStoreNotStarted(execution, err) + } + return taskengine.Suspend(model.TaskResumeModeRecover, 0, "provider_confirmation", "Checking storage transfer", nil) + } + } if h.deps.Cache == nil || h.deps.CacheGate == nil { - return h.failCopyTask(execution, input, copyRow, errors.New("cache reader is unavailable"), "dependency_unavailable") + if hasCheckpoint { + return taskengine.Fail(errors.New("cache reader is unavailable"), "store_cache_read_failed", nil) + } + return h.retryStoreNotStarted(execution, errors.New("cache reader is unavailable")) } cacheKey := model.ContentCacheKey(content.ID) releaseCache := h.deps.CacheGate.HoldRead(cacheKey) @@ -1067,13 +1127,16 @@ func (h *TaskHandlers) runStore(ctx context.Context, execution taskengine.Execut // that has to wait for a slot never hashes the bytes first. var outcome taskengine.Result err = execution.WithResource(ctx, taskengine.ResourceProviderMutation, func(ctx context.Context) error { - outcome = h.storeWithProviderSlot(ctx, execution, input, copyRow, target, content, bucket, cacheKey) + outcome = h.storeWithProviderSlot(ctx, execution, input, copyRow, target, content, bucket, cacheKey, checkpoint, hasCheckpoint) return nil }) if errors.Is(err, taskengine.ErrResourceBusy) { return taskengine.ResourceWait("Waiting for other storage operations to finish") } if err != nil { + if hasCheckpoint { + return taskengine.Fail(err, "store_checkpoint_write_failed", nil) + } return h.retryStoreNotStarted(execution, err) } return outcome @@ -1090,38 +1153,90 @@ func (h *TaskHandlers) storeWithProviderSlot( content *model.StorageContent, bucket *model.Bucket, cacheKey string, + previous storeCheckpoint, + hasCheckpoint bool, ) taskengine.Result { + if hasCheckpoint { + pieceCID, err := cid.Parse(previous.IntendedPieceCID) + if err != nil { + return taskengine.Fail(err, "invalid_checkpoint", nil) + } + state, checkErr := h.storePieceState(ctx, target, previous, pieceCID) + if checkErr != nil || state == synapse.ParkedPieceProcessing { + return h.waitForStoreStatus(previous, state, checkErr) + } + if state == synapse.ParkedPieceReady { + return h.finishPieceTransfer(ctx, execution, input, copyRow, target, pieceCID) + } + if execution.RetryWillFail() { + return taskengine.Fail(errors.New("storage transfer retry limit reached"), "store_retry_limit", nil) + } + } calculateReader, _, err := h.deps.Cache.Get(ctx, bucket.Name, cacheKey) if err != nil { if os.IsNotExist(err) { - if copyRow.TransferMethod == model.StorageCopyTransferMethodCacheRestore { - return taskengine.Fail(errors.New("stored content migration cannot read its local cache"), "migration_cache_missing", nil) + if hasCheckpoint || copyRow.IngressStoreAttempt > 0 || copyRow.TransferMethod == model.StorageCopyTransferMethodCacheRestore { + return taskengine.Fail(errors.New("storage transfer needs local bytes that are no longer cached"), "store_cache_missing", nil) } return taskengine.Suspend(model.TaskResumeModeExecute, storageDependencyWait, "source", "Waiting for retained cache data", nil) } - return h.retryCopyTask(execution, input, copyRow, err, "cache_open_failed") + return h.storeCacheFailure(execution, hasCheckpoint, err) } pieceInfo, calculateErr := piece.Calculate(calculateReader) closeErr := calculateReader.Close() if calculateErr != nil { - return h.retryCopyTask(execution, input, copyRow, calculateErr, "store_identity_failed") + return h.storeCacheFailure(execution, hasCheckpoint, calculateErr) } if closeErr != nil { - return h.retryCopyTask(execution, input, copyRow, closeErr, "cache_close_failed") + return h.storeCacheFailure(execution, hasCheckpoint, closeErr) } if !pieceInfo.CIDv2.Defined() || content.ContentSize < 0 || pieceInfo.RawSize != uint64(content.ContentSize) { err := fmt.Errorf("calculated storage identity has size %d, expected %d", pieceInfo.RawSize, content.ContentSize) + if hasCheckpoint { + return taskengine.Fail(err, "store_identity_mismatch", nil) + } return h.failCopyTask(execution, input, copyRow, err, "store_identity_mismatch") } + if hasCheckpoint && pieceInfo.CIDv2.String() != previous.IntendedPieceCID { + return taskengine.Fail(errors.New("cached bytes differ from the checkpointed piece"), "store_identity_mismatch", nil) + } + if !hasCheckpoint && copyRow.IngressStoreAttempt > 0 { + if h.deps.ParkedPieces == nil { + return h.retryStoreNotStarted(execution, errors.New("storage provider status checker is unavailable")) + } + state, checkErr := h.deps.ParkedPieces.FindParkedPiece(ctx, target.ServiceURL(), pieceInfo.CIDv2) + if checkErr != nil || state == synapse.ParkedPieceProcessing { + if checkErr == nil { + checkErr = errors.New("storage provider is still processing the piece") + } + return h.retryStoreNotStarted(execution, checkErr) + } + if state == synapse.ParkedPieceReady { + checkpoint := storeCheckpoint{ + AttemptedAt: time.Now().UTC(), IntendedPieceCID: pieceInfo.CIDv2.String(), + ProviderServiceURL: target.ServiceURL(), IngressAttempt: copyRow.IngressStoreAttempt, + } + if err := execution.WriteCheckpoint(ctx, checkpoint); err != nil { + return taskengine.Fail(err, "store_checkpoint_write_failed", nil) + } + return h.finishPieceTransfer(ctx, execution, input, copyRow, target, pieceInfo.CIDv2) + } + if state != synapse.ParkedPieceMissing { + return h.retryStoreNotStarted(execution, fmt.Errorf("unexpected storage provider state %q", state)) + } + } + if (hasCheckpoint || copyRow.IngressStoreAttempt > 0) && execution.RetryWillFail() { + return taskengine.Fail(errors.New("storage transfer retry limit reached"), "store_retry_limit", nil) + } storeReader, _, err := h.deps.Cache.Get(ctx, bucket.Name, cacheKey) if err != nil { if os.IsNotExist(err) { - if copyRow.TransferMethod == model.StorageCopyTransferMethodCacheRestore { - return taskengine.Fail(errors.New("stored content migration cannot read its local cache"), "migration_cache_missing", nil) + if hasCheckpoint || copyRow.IngressStoreAttempt > 0 || copyRow.TransferMethod == model.StorageCopyTransferMethodCacheRestore { + return taskengine.Fail(errors.New("storage transfer needs local bytes that are no longer cached"), "store_cache_missing", nil) } return taskengine.Suspend(model.TaskResumeModeExecute, storageDependencyWait, "source", "Waiting for retained cache data", nil) } - return h.retryCopyTask(execution, input, copyRow, err, "cache_open_failed") + return h.storeCacheFailure(execution, hasCheckpoint, err) } defer func() { _ = storeReader.Close() }() checkpoint := storeCheckpoint{ @@ -1132,6 +1247,11 @@ func (h *TaskHandlers) storeWithProviderSlot( if copyRow.TransferMethod == model.StorageCopyTransferMethodIngress || copyRow.TransferMethod == model.StorageCopyTransferMethodCacheRestore { checkpoint.IngressAttempt = copyRow.IngressStoreAttempt + 1 checkpointSettlement = func(ctx context.Context, repos *repository.Repositories) error { + if hasCheckpoint || copyRow.IngressStoreAttempt > 0 { + if err := repos.Tasks.ConsumeStoreRetry(ctx, execution.ID(), execution.ClaimGeneration()); err != nil { + return err + } + } _, err := repos.Contents.BeginIngressStoreProgress(ctx, repository.BeginIngressStoreProgressInput{ CopyID: copyRow.ID, Generation: input.Generation, TaskID: execution.ID(), Attempt: checkpoint.IngressAttempt, }) @@ -1155,13 +1275,16 @@ func (h *TaskHandlers) storeWithProviderSlot( defer progress.Close() if err != nil { if !attempted { + if hasCheckpoint || copyRow.IngressStoreAttempt > 0 { + return taskengine.Fail(err, "store_checkpoint_write_failed", nil) + } return h.retryStoreNotStarted(execution, err) } return taskengine.Suspend(model.TaskResumeModeRecover, storagePollInterval, "provider_confirmation", "Checking storage transfer", nil) } if stored == nil || !stored.PieceCID.Equals(pieceInfo.CIDv2) || stored.Size != content.ContentSize { err := errors.New("storage provider returned a mismatched piece identity or size") - return h.failCopyTask(execution, input, copyRow, err, "store_result_invalid") + return taskengine.Fail(err, "store_result_invalid", nil) } progress.Flush(content.ContentSize, true) return h.finishPieceTransfer(ctx, execution, input, copyRow, target, stored.PieceCID) @@ -1175,34 +1298,85 @@ func (h *TaskHandlers) recoverStore( target synapse.DataSetTarget, checkpoint storeCheckpoint, ) taskengine.Result { - if checkpoint.AttemptedAt.IsZero() || checkpoint.IntendedPieceCID == "" || checkpoint.ProviderServiceURL == "" { - return taskengine.Fail(errors.New("storage transfer checkpoint is incomplete"), "invalid_checkpoint", nil) + if err := validateStoreCheckpoint(checkpoint, copyRow.ContentSize); err != nil { + return taskengine.Fail(err, "invalid_checkpoint", nil) } pieceCID, err := cid.Parse(checkpoint.IntendedPieceCID) if err != nil { return taskengine.Fail(err, "invalid_checkpoint", nil) } - pieceInfo, err := piece.ParseV2(pieceCID) - if err != nil || copyRow.ContentSize < 0 || pieceInfo.RawSize != uint64(copyRow.ContentSize) { - if err == nil { - err = fmt.Errorf("checkpointed storage identity has size %d, expected %d", pieceInfo.RawSize, copyRow.ContentSize) + state, findErr := h.storePieceState(ctx, target, checkpoint, pieceCID) + if findErr == nil && state == synapse.ParkedPieceReady { + return h.finishPieceTransfer(ctx, execution, input, copyRow, target, pieceCID) + } + if findErr == nil && state == synapse.ParkedPieceMissing { + if execution.RetryWillFail() { + return taskengine.Fail(errors.New("storage transfer retry limit reached"), "store_retry_limit", nil) } - return taskengine.Fail(err, "invalid_checkpoint", nil) + return taskengine.Suspend(model.TaskResumeModeExecute, 0, "provider_confirmation", "Retrying storage transfer", nil) } - if h.deps.ParkedPieces == nil { - return taskengine.Suspend(model.TaskResumeModeRecover, storageDependencyWait, "provider", "Waiting for storage provider", nil) + return h.waitForStoreStatus(checkpoint, state, findErr) +} + +func validateStoreCheckpoint(checkpoint storeCheckpoint, contentSize int64) error { + if checkpoint.AttemptedAt.IsZero() || checkpoint.IntendedPieceCID == "" || checkpoint.ProviderServiceURL == "" { + return errors.New("storage transfer checkpoint is incomplete") } - state, findErr := h.deps.ParkedPieces.FindParkedPiece(ctx, checkpoint.ProviderServiceURL, pieceCID) - if findErr == nil && state == synapse.ParkedPieceReady { - return h.finishPieceTransfer(ctx, execution, input, copyRow, target, pieceCID) + pieceCID, err := cid.Parse(checkpoint.IntendedPieceCID) + if err != nil { + return err } - if time.Since(checkpoint.AttemptedAt) >= storeAttentionAfter { - if findErr == nil { - findErr = fmt.Errorf("storage transfer remains %s", state) + pieceInfo, err := piece.ParseV2(pieceCID) + if err != nil { + return err + } + if contentSize < 0 || pieceInfo.RawSize != uint64(contentSize) { + return fmt.Errorf("checkpointed storage identity has size %d, expected %d", pieceInfo.RawSize, contentSize) + } + return nil +} + +func (h *TaskHandlers) storePieceState(ctx context.Context, target synapse.DataSetTarget, checkpoint storeCheckpoint, pieceCID cid.Cid) (synapse.ParkedPieceState, error) { + if h.deps.ParkedPieces == nil { + return "", errors.New("storage provider status checker is unavailable") + } + oldState, oldErr := h.deps.ParkedPieces.FindParkedPiece(ctx, checkpoint.ProviderServiceURL, pieceCID) + if checkpoint.ProviderServiceURL == target.ServiceURL() { + if oldErr == nil && oldState != synapse.ParkedPieceMissing && oldState != synapse.ParkedPieceProcessing && oldState != synapse.ParkedPieceReady { + return "", fmt.Errorf("unexpected storage provider state %q", oldState) } - return taskengine.Fail(findErr, "store_outcome_unknown", nil) + return oldState, oldErr + } + currentState, currentErr := h.deps.ParkedPieces.FindParkedPiece(ctx, target.ServiceURL(), pieceCID) + if currentErr == nil && currentState != synapse.ParkedPieceMissing && currentState != synapse.ParkedPieceProcessing && currentState != synapse.ParkedPieceReady { + return "", fmt.Errorf("unexpected storage provider state %q", currentState) + } + if currentErr == nil && currentState == synapse.ParkedPieceReady { + return currentState, nil + } + if currentErr != nil { + return "", currentErr + } + if oldErr != nil { + return "", oldErr + } + if oldState != synapse.ParkedPieceMissing && oldState != synapse.ParkedPieceProcessing && oldState != synapse.ParkedPieceReady { + return "", fmt.Errorf("unexpected storage provider state %q", oldState) } - return taskengine.Suspend(model.TaskResumeModeRecover, storagePollInterval, "provider_confirmation", "Checking storage transfer", nil) + if currentState == synapse.ParkedPieceMissing && (oldState == synapse.ParkedPieceReady || oldState == synapse.ParkedPieceMissing) { + return synapse.ParkedPieceMissing, nil + } + return synapse.ParkedPieceProcessing, nil +} + +func (h *TaskHandlers) waitForStoreStatus(checkpoint storeCheckpoint, state synapse.ParkedPieceState, checkErr error) taskengine.Result { + if time.Since(checkpoint.AttemptedAt) < storeAttentionAfter { + return taskengine.Suspend(model.TaskResumeModeRecover, storagePollInterval, "provider_confirmation", "Checking storage transfer", nil) + } + if checkErr != nil { + return taskengine.Fail(checkErr, "store_check_failed", nil) + } + return taskengine.Fail(fmt.Errorf("storage provider still reports %s", state), "store_processing_timeout", nil) } func (h *TaskHandlers) retryStoreNotStarted(execution taskengine.Execution, err error) taskengine.Result { @@ -1212,6 +1386,13 @@ func (h *TaskHandlers) retryStoreNotStarted(execution taskengine.Execution, err return retryTask(err, "store_not_started") } +func (h *TaskHandlers) storeCacheFailure(execution taskengine.Execution, hasCheckpoint bool, err error) taskengine.Result { + if hasCheckpoint { + return taskengine.Fail(err, "store_cache_read_failed", nil) + } + return h.retryStoreNotStarted(execution, err) +} + func (h *TaskHandlers) pullHandler() taskengine.Handler { definition := copyDefinition(model.TaskTypeStoragePull, h.retryLimit()) return taskHandler{ @@ -1537,6 +1718,9 @@ func (h *TaskHandlers) authorizeCopyTask( return input, nil, true, taskengine.Cancel("Storage work was superseded", nil) } if err != nil { + if execution.Type() == model.TaskTypeStorageStore && execution.RetryWillFail() { + return input, nil, true, taskengine.Fail(err, "copy_authorization_failed", nil) + } if execution.RetryWillFail() { message := err.Error() return input, nil, true, taskengine.Fail(err, "copy_authorization_failed", func(ctx context.Context, repos *repository.Repositories) error { @@ -1717,7 +1901,7 @@ func (h *TaskHandlers) finishPieceTransfer( ) taskengine.Result { extra, err := target.PresignForCommit(ctx, []storage.PieceInput{{PieceCID: pieceCID}}) if err != nil { - return h.retryCopyTask(execution, input, copyRow, err, "commit_presign_failed") + return taskengine.Fail(err, "commit_presign_failed", nil) } return h.finishPieceTransferWithExtra(execution, input, copyRow, target, pieceCID, hex.EncodeToString(extra), "") } diff --git a/internal/worker/task_handlers_test.go b/internal/worker/task_handlers_test.go index 0f14e30..b4314d4 100644 --- a/internal/worker/task_handlers_test.go +++ b/internal/worker/task_handlers_test.go @@ -473,11 +473,7 @@ func (p *recordingWorkerEvents) Publish(topic string, payload map[string]any) { } } -// TestTerminalStoreFailureSettlesCopyAndContent checks that a store failure -// settles everything the transfer owns: the copy is failed and unbound, ingress -// progress stays on the copy that produced it, and the version's derived -// position follows the copies to failed without a stored column. -func TestTerminalStoreFailureSettlesCopyAndContent(t *testing.T) { +func TestStoreResultMismatchRetainsCopyAndCheckpoint(t *testing.T) { events := &recordingWorkerEvents{events: make(chan recordedWorkerEvent, 8)} cacheStore := &testutil.MockCache{GetFunc: func(context.Context, string, string) (io.ReadCloser, *cache.ObjectInfo, error) { return io.NopCloser(strings.NewReader(strings.Repeat("s", 128))), &cache.ObjectInfo{Size: 128}, nil @@ -572,11 +568,11 @@ func TestTerminalStoreFailureSettlesCopyAndContent(t *testing.T) { failedTask := waitForTask(t, runtime.repos, taskRow.ID, func(task *model.Task) bool { return task.Status == model.TaskStatusFailed }) - if runtime.service.Retryable(failedTask) { - t.Fatal("terminal copy task is unexpectedly retryable") + if !runtime.service.Retryable(failedTask) || len(failedTask.Checkpoint) == 0 { + t.Fatal("Store result failure lost its recovery checkpoint") } copyRow, err := runtime.repos.Contents.GetUploadCopyByID(ctx, copies[0].ID) - if err != nil || copyRow.Status != model.StorageCopyStatusFailed || copyRow.ActiveTaskID != nil { + if err != nil || copyRow.Status != model.StorageCopyStatusPending || copyRow.ActiveTaskID == nil || *copyRow.ActiveTaskID != taskRow.ID { t.Fatalf("terminal copy = %#v, err=%v", copyRow, err) } // Ingress progress belongs to the transfer, so it survives on the copy. @@ -584,7 +580,7 @@ func TestTerminalStoreFailureSettlesCopyAndContent(t *testing.T) { t.Fatalf("terminal copy progress = attempt:%d bytes:%d at:%v", copyRow.IngressStoreAttempt, copyRow.IngressBytesTransferred, copyRow.ProgressUpdatedAt) } storedVersion, err := runtime.repos.Objects.GetVersionByID(ctx, version.VersionID) - if err != nil || storedVersion.State != model.ObjectStateFailed { + if err != nil || storedVersion.State != model.ObjectStateUploading { t.Fatalf("terminal version = %#v, err=%v", storedVersion, err) } } @@ -2948,8 +2944,8 @@ func TestStoreManualRetryRequiresUnsettledRecoveryEvidence(t *testing.T) { {name: "checkpointed owner missing", reason: "copy_owner_missing", checkpoint: checkpoint, want: true}, {name: "context failure before checkpoint", reason: "copy_context_failed"}, {name: "invalid checkpoint", reason: "invalid_checkpoint", checkpoint: checkpoint}, - {name: "settled store result failure", reason: "store_result_invalid", checkpoint: checkpoint}, - {name: "settled presign failure", reason: "commit_presign_failed", checkpoint: checkpoint}, + {name: "store result failure", reason: "store_result_invalid", checkpoint: checkpoint, want: true}, + {name: "presign failure", reason: "commit_presign_failed", checkpoint: checkpoint, want: true}, } for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { @@ -2965,7 +2961,7 @@ func TestStoreManualRetryRequiresUnsettledRecoveryEvidence(t *testing.T) { } } -func TestStoreUnknownRecoveryOnlyRechecksParkedPiece(t *testing.T) { +func TestStoreRecoveryRetriesMissingPieceAfterCheckpoint(t *testing.T) { payload := strings.Repeat("parked-store", 12) payload = payload[:128] info, err := piece.Calculate(strings.NewReader(payload)) @@ -3018,6 +3014,10 @@ func TestStoreUnknownRecoveryOnlyRechecksParkedPiece(t *testing.T) { }, }) pipeline := seedCopyPipeline(t, runtime, model.StorageCopyStatusPending) + if _, err := runtime.db.NewUpdate().Model((*model.StorageCopy)(nil)).Set("transfer_method = ?", model.StorageCopyTransferMethodCacheRestore). + Where("id = ?", pipeline.target.ID).Exec(t.Context()); err != nil { + t.Fatalf("set cache restore transfer: %v", err) + } version := &model.ObjectVersion{ VersionID: model.NewVersionID(), BucketID: pipeline.upload.BucketID, Key: "parked-store.bin", ContentID: &pipeline.upload.ID, Size: pipeline.upload.ContentSize, ETag: "parked-store", @@ -3039,7 +3039,7 @@ func TestStoreUnknownRecoveryOnlyRechecksParkedPiece(t *testing.T) { } taskRow := bindCopyTask(t, runtime, pipeline.target, model.TaskTypeStorageStore) limitedRepos := *runtime.repos - limited := &limitedClaimRepository{TaskRepository: runtime.repos.Tasks, maximum: 5} + limited := &limitedClaimRepository{TaskRepository: runtime.repos.Tasks, maximum: 6} limitedRepos.Tasks = limited runtime.repos.Tasks = limited engine, err := taskengine.NewEngine(taskengine.EngineConfig{ @@ -3093,25 +3093,25 @@ func TestStoreUnknownRecoveryOnlyRechecksParkedPiece(t *testing.T) { if _, err := runtime.db.NewRaw(`UPDATE tasks SET available_at = ? WHERE id = ?`, time.Now().Add(-time.Second), taskRow.ID).Exec(t.Context()); err != nil { t.Fatalf("wake store recovery: %v", err) } - failed := waitForTask(t, runtime.repos, taskRow.ID, func(task *model.Task) bool { - return task.Status == model.TaskStatusFailed + retrying := waitForTask(t, runtime.repos, taskRow.ID, func(task *model.Task) bool { + return limited.claims.Load() == 5 && task.Status == model.TaskStatusPending && task.ResumeMode == model.TaskResumeModeRecover }) - if failed.FailureReason == nil || *failed.FailureReason != "store_outcome_unknown" || !runtime.service.Retryable(failed) { - t.Fatalf("unknown store task = %#v", failed) + if retrying.RetryCount != 1 || storeCalls.Load() != 2 { + t.Fatalf("retransmitted store = retry count:%d calls:%d, want 1/2", retrying.RetryCount, storeCalls.Load()) } copyRow, err := runtime.repos.Contents.GetUploadCopyByID(t.Context(), pipeline.target.ID) if err != nil || copyRow.Status != model.StorageCopyStatusPending || copyRow.ActiveTaskID == nil || *copyRow.ActiveTaskID != taskRow.ID { t.Fatalf("unknown store copy = %#v, err=%v", copyRow, err) } parkedState.Store(synapse.ParkedPieceReady) - if err := runtime.service.Retry(t.Context(), taskRow.ID); err != nil { - t.Fatalf("retry unknown store: %v", err) + if _, err := runtime.db.NewRaw(`UPDATE tasks SET available_at = ? WHERE id = ?`, time.Now().Add(-time.Second), taskRow.ID).Exec(t.Context()); err != nil { + t.Fatalf("wake retransmitted store recovery: %v", err) } waitForTask(t, runtime.repos, taskRow.ID, func(task *model.Task) bool { return task.Status == model.TaskStatusCompleted }) - if storeCalls.Load() != 1 || cacheOpens.Load() != 2 || parkedCalls.Load() != 4 { - t.Fatalf("store recovery calls = store:%d cache:%d parked:%d, want 1/2/4", storeCalls.Load(), cacheOpens.Load(), parkedCalls.Load()) + if storeCalls.Load() != 2 || cacheOpens.Load() != 4 || parkedCalls.Load() < 5 { + t.Fatalf("store recovery calls = store:%d cache:%d parked:%d", storeCalls.Load(), cacheOpens.Load(), parkedCalls.Load()) } } diff --git a/ui/src/lib/overview.ts b/ui/src/lib/overview.ts index fe5bd2d..a126e1c 100644 --- a/ui/src/lib/overview.ts +++ b/ui/src/lib/overview.ts @@ -242,6 +242,8 @@ export function filecoinStorageHealthReadyPercent(available: number, total: numb } export function filecoinStorageHealthStatusLabel(health: OverviewData['filecoin_storage_health']) { + if (health.providers?.summary_signal.freshness.stale || health.data_sets?.summary_signal.freshness.stale) + return hasFreshBlockingStorageHealthSignal(health) ? 'Blocked' : 'Health data outdated' if (Object.keys(health.partial_errors ?? {}).length > 0) return filecoinStorageHealthLevelLabel(health.level) if (health.level === 'blocking') return filecoinStorageHealthLevelLabel(health.level) if (hasFilecoinStorageHealthIssue(health.data_sets)) return 'Degraded' @@ -251,6 +253,20 @@ export function filecoinStorageHealthStatusLabel(health: OverviewData['filecoin_ return filecoinStorageHealthLevelLabel(health.level) } +export function filecoinStorageHealthDisplayLevel( + health: OverviewData['filecoin_storage_health'] +): FilecoinStorageHealthLevel { + if (!health.providers?.summary_signal.freshness.stale && !health.data_sets?.summary_signal.freshness.stale) + return health.level + return hasFreshBlockingStorageHealthSignal(health) ? 'blocking' : 'warning' +} + +function hasFreshBlockingStorageHealthSignal(health: OverviewData['filecoin_storage_health']) { + return [health.providers, health.data_sets].some( + (section) => section?.summary_signal.level === 'blocking' && !section.summary_signal.freshness.stale + ) +} + export function filecoinStorageHealthPartialErrorRows(partialErrors: Record) { const labels: Record = { observability: 'Observability unavailable', diff --git a/ui/src/routes/-root-content.ts b/ui/src/routes/-root-content.ts index 5b8dbd5..7bff50a 100644 --- a/ui/src/routes/-root-content.ts +++ b/ui/src/routes/-root-content.ts @@ -1,5 +1,5 @@ import type { FilecoinReadinessCheck, FilecoinReadinessData, FilecoinReadinessStatus } from '../api/client.ts' -import { filecoinReadinessSummary, importantFilecoinReadinessChecks } from '../lib/filecoin-readiness.ts' +import { importantFilecoinReadinessChecks } from '../lib/filecoin-readiness.ts' export type RootSettingsState = { mode: 'ready' | 'setup' @@ -15,9 +15,12 @@ export type GlobalFilecoinReadinessAlertState = summary: string status: FilecoinReadinessStatus primaryCheck?: FilecoinReadinessCheck + details?: FilecoinReadinessData failed: boolean } +const provisioningCheckIds = new Set(['providers', 'storage_cost', 'payment_funding', 'fwss_approval']) + export function rootUsesSetupShell(settings: RootSettingsState | undefined) { return settings?.mode === 'setup' || settings?.runtime_available === false } @@ -62,25 +65,17 @@ export function globalFilecoinReadinessAlertState({ if (!data) return { show: false } - const visibleChecks = importantFilecoinReadinessChecks(data.checks, dismissedCheckIds) + const checks = data.checks.filter((check) => !provisioningCheckIds.has(check.id)) + const visibleChecks = importantFilecoinReadinessChecks(checks, dismissedCheckIds) const primaryCheck = visibleChecks[0] if (primaryCheck?.status === 'blocked' || primaryCheck?.status === 'unknown') { return { show: true, title: filecoinReadinessAlertTitle(primaryCheck.status), - summary: filecoinReadinessSummary(data, dismissedCheckIds), + summary: primaryCheck.message, status: primaryCheck.status, primaryCheck, - failed: false, - } - } - - if (data.checks.length === 0 && (data.status === 'blocked' || data.status === 'unknown')) { - return { - show: true, - title: filecoinReadinessAlertTitle(data.status), - summary: filecoinReadinessSummary(data, dismissedCheckIds), - status: data.status, + details: { ...data, status: primaryCheck.status, checks }, failed: false, } } diff --git a/ui/src/routes/__root.tsx b/ui/src/routes/__root.tsx index 51388e6..aa0fef3 100644 --- a/ui/src/routes/__root.tsx +++ b/ui/src/routes/__root.tsx @@ -319,8 +319,8 @@ function GlobalFilecoinReadinessAlert({ enabled }: { enabled: boolean }) { type="button" variant="outline" size="sm" - disabled={!data} - onClick={() => data && setDetailsOpen(true)} + disabled={!alert.details} + onClick={() => alert.details && setDetailsOpen(true)} > Details @@ -331,7 +331,7 @@ function GlobalFilecoinReadinessAlert({ enabled }: { enabled: boolean }) { @@ -172,7 +173,7 @@ function StorageHealthCard({ health }: { health: OverviewData['filecoin_storage_
- Storage copies + Data sets {formatOptionalNumber(dataSets.available)} / {formatOptionalNumber(dataSets.total)} ready @@ -180,7 +181,7 @@ function StorageHealthCard({ health }: { health: OverviewData['filecoin_storage_
)} - {task.failure_reason === 'store_outcome_unknown' ? 'Check again' : 'Recover'} + {task.type === 'storage_store' ? 'Retry upload' : 'Recover'} )} {task.acknowledgeable && ( diff --git a/ui/test/overview.test.ts b/ui/test/overview.test.ts index c90fc19..6a0c14e 100644 --- a/ui/test/overview.test.ts +++ b/ui/test/overview.test.ts @@ -4,6 +4,7 @@ import type { ObservabilityFreshness, ObservabilitySummary, OverviewData } from import { attentionDisplayRows, filecoinStorageHealthCheckedLabel, + filecoinStorageHealthDisplayLevel, filecoinStorageHealthFreshnessLabel, filecoinStorageHealthLevelLabel, filecoinStorageHealthLevelStyle, @@ -217,7 +218,7 @@ test('filecoin storage health status label prefers concrete storage and provider ) assert.equal( filecoinStorageHealthStatusLabel( - filecoinStorageHealthFixture({ providers: { total: 2, available: 1, unavailable: 1 } }) + filecoinStorageHealthFixture({ providers: { total: 2, available: 1, degraded: 1 } }) ), 'Provider degraded' ) @@ -229,6 +230,36 @@ test('filecoin storage health status label prefers concrete storage and provider ) }) +test('outdated health data takes priority over old provider failures', () => { + const health = filecoinStorageHealthFixture( + { level: 'blocking', providers: { total: 1, unavailable: 1 } }, + { providerFreshness: { stale: true, warnings: ['stale_state'], last_checked_at: '2026-01-01T00:00:00Z' } } + ) + assert.equal(filecoinStorageHealthStatusLabel(health), 'Health data outdated') + assert.equal(filecoinStorageHealthDisplayLevel(health), 'warning') +}) + +test('fresh storage failures remain blocking when the other health group is outdated', () => { + for (const staleGroup of ['providers', 'data_sets'] as const) { + const health = filecoinStorageHealthFixture( + { + level: 'blocking', + providers: staleGroup === 'data_sets' ? { total: 1, unavailable: 1 } : { total: 1, available: 1 }, + dataSets: staleGroup === 'providers' ? { total: 1, unavailable: 1 } : { total: 1, available: 1 }, + }, + { + providerFreshness: + staleGroup === 'providers' ? { stale: true, warnings: ['stale_state'] } : { stale: false, warnings: [] }, + dataSetFreshness: + staleGroup === 'data_sets' ? { stale: true, warnings: ['stale_state'] } : { stale: false, warnings: [] }, + } + ) + assert.equal(filecoinStorageHealthStatusLabel(health), 'Blocked') + assert.equal(filecoinStorageHealthDisplayLevel(health), 'blocking') + assert.equal(filecoinStorageHealthCheckedLabel(health), 'Stale') + } +}) + test('filecoin storage health partial error rows use fixed labels', () => { assert.deepEqual(filecoinStorageHealthPartialErrorRows({ observability_providers: 'health query failed' }), [ { key: 'observability_providers', label: 'Provider summary unavailable', message: 'health query failed' }, @@ -260,14 +291,14 @@ function filecoinStorageHealthFixture( providers: { summary: filecoinStorageHealthSummaryFixture(providers), summary_signal: { - level: providers.unavailable || providers.degraded || providers.unknown ? 'warning' : 'ok', + level: providers.unavailable ? 'blocking' : providers.degraded || providers.unknown ? 'warning' : 'ok', freshness: options.providerFreshness ?? freshness, }, }, data_sets: { summary: filecoinStorageHealthSummaryFixture(dataSets), summary_signal: { - level: dataSets.unavailable || dataSets.degraded || dataSets.unknown ? 'warning' : 'ok', + level: dataSets.unavailable ? 'blocking' : dataSets.degraded || dataSets.unknown ? 'warning' : 'ok', freshness: options.dataSetFreshness ?? freshness, }, }, diff --git a/ui/test/root-content.test.ts b/ui/test/root-content.test.ts index 1e84027..9ab8454 100644 --- a/ui/test/root-content.test.ts +++ b/ui/test/root-content.test.ts @@ -69,7 +69,7 @@ test('global filecoin readiness alert shows blocked checks first', () => { const alert = globalFilecoinReadinessAlertState({ enabled: true, data: readinessData('blocked', [ - { id: 'storage_cost', status: 'unknown', message: 'Storage cost could not be estimated.' }, + { id: 'providers', status: 'blocked', message: 'Only 0 approved providers are available.' }, { id: 'wallet_fil_gas', status: 'blocked', message: 'FIL gas balance is empty.' }, ]), }) @@ -79,6 +79,10 @@ test('global filecoin readiness alert shows blocked checks first', () => { assert.equal(alert.status, 'blocked') assert.equal(alert.summary, 'FIL gas balance is empty.') assert.equal(alert.primaryCheck?.id, 'wallet_fil_gas') + assert.deepEqual( + alert.details?.checks.map((check) => check.id), + ['wallet_fil_gas'] + ) }) test('global filecoin readiness alert shows unknown checks without blocked checks', () => { @@ -86,27 +90,29 @@ test('global filecoin readiness alert shows unknown checks without blocked check enabled: true, data: readinessData('warning', [ { id: 'payment_runway', status: 'warning', message: 'Payment runway is low.' }, - { id: 'storage_cost', status: 'unknown', message: 'Storage cost could not be estimated.' }, + { id: 'payment_account', status: 'unknown', message: 'Payment account could not be checked.' }, ]), }) assert.equal(alert.show, true) assert.equal(alert.title, 'Filecoin readiness unknown') assert.equal(alert.status, 'unknown') - assert.equal(alert.summary, 'Storage cost could not be estimated.') - assert.equal(alert.primaryCheck?.id, 'storage_cost') + assert.equal(alert.summary, 'Payment account could not be checked.') + assert.equal(alert.primaryCheck?.id, 'payment_account') }) -test('global filecoin readiness alert falls back to aggregate unknown when checks are empty', () => { +test('global filecoin readiness alert ignores provider selection and its dependent checks', () => { const alert = globalFilecoinReadinessAlertState({ enabled: true, - data: readinessData('unknown', []), + data: readinessData('blocked', [ + { id: 'providers', status: 'blocked', message: 'Only 0 approved providers are available.' }, + { id: 'storage_cost', status: 'unknown', message: 'Storage cost could not be estimated.' }, + { id: 'payment_funding', status: 'unknown', message: 'Payment funding could not be checked.' }, + { id: 'fwss_approval', status: 'unknown', message: 'FWSS approval could not be checked.' }, + ]), }) - assert.equal(alert.show, true) - assert.equal(alert.title, 'Filecoin readiness unknown') - assert.equal(alert.status, 'unknown') - assert.equal(alert.summary, 'Unknown') + assert.equal(alert.show, false) }) test('global filecoin readiness alert shows query failures and respects disabled state', () => {