diff --git a/docs/changelog/index.md b/docs/changelog/index.md index 9dc94e9d..7a83c8af 100644 --- a/docs/changelog/index.md +++ b/docs/changelog/index.md @@ -23,7 +23,7 @@ sidebar: false ### 新增 -- Pages 项目新增持久部署源,可配置 Remote URL 或公开 GitHub Release,并支持手动检查、同步发布、来源状态查看与同一 Release 资源替换确认;部署历史会保留安全的来源快照。 +- Pages 项目新增持久部署源,可配置 Remote URL 或公开 GitHub Release,并支持手动检查、同步发布、来源状态查看与同一 Release 资源替换确认;GitHub latest 来源可按设定间隔自动检查并发布更新,部署历史会保留安全的来源快照。 - WAF 规则编排新增「UA 检查」节点:可要求携带 User-Agent、按浏览器/操作系统白名单(且/或)匹配,并优先屏蔽常见爬虫、非正常 UA(不含爬虫)与自定义正则 UA。 ### 改进 @@ -33,7 +33,7 @@ sidebar: false ### 修复 -- 修复 Pages 部署包路径校验、归档展开限额、历史版本裁剪、代理路由绑定与 Agent 下载过程中的安全和一致性问题;大包改为流式处理,部署入口、旧版目录切换、保留版本及上传记录在并发场景下更加可靠。 +- 修复 Pages 部署包路径校验、归档展开限额、历史版本裁剪、代理路由绑定与 Agent 下载过程中的安全和一致性问题;大包改为流式处理,部署入口、旧版目录切换、保留版本及上传记录在并发场景下更加可靠,异常中断遗留的部署包也会被安全补偿清理。 ## [v3.4.0] - 2026-07-19 diff --git a/internal/apps/openflare/pages/errs.go b/internal/apps/openflare/pages/errs.go index 64fc9a25..0d2493ef 100644 --- a/internal/apps/openflare/pages/errs.go +++ b/internal/apps/openflare/pages/errs.go @@ -61,6 +61,7 @@ const ( errPagesSourceActionInvalid = "pages 部署源任务参数无效" errPagesSourceActionStale = "pages 部署源配置已变化,本次任务已跳过" errPagesSourceLeaseLost = "pages 部署源任务执行权已失效" + errPagesSourceLeaseExpired = "上次 pages 部署源任务租约已过期" errPagesSourceSyncFailed = "pages 部署源同步失败" errPagesSourceTaskDispatchFailed = "pages 部署源任务入队失败" errPagesSourceInternal = "pages 部署源操作失败,请稍后重试" diff --git a/internal/apps/openflare/pages/github_source.go b/internal/apps/openflare/pages/github_source.go index d2e8fa66..12db7291 100644 --- a/internal/apps/openflare/pages/github_source.go +++ b/internal/apps/openflare/pages/github_source.go @@ -44,6 +44,7 @@ type githubSourceConfig struct { Selector string Tag string AssetName string + AutoUpdate bool CheckInterval int SourceIdentity string } @@ -56,9 +57,6 @@ func validateGitHubSourceInput(input SourceUpdateInput) error { strings.TrimSpace(input.RemoteNetworkPolicy) != "" { return errors.New(errPagesSourceGitHubFields) } - if input.AutoUpdateEnabled { - return errors.New(errPagesSourceAutoNotAvailable) - } if _, err := normalizeGitHubRepositoryURL(input.RepositoryURL); err != nil { return err } @@ -83,7 +81,7 @@ func validateGitHubSourceInput(input SourceUpdateInput) error { return errors.New(errPagesSourceCheckInterval) } case githubReleaseSelectorTag: - if !validGitHubReleaseTagConfig(input.ReleaseTag) || input.CheckIntervalMinutes != 0 { + if !validGitHubReleaseTagConfig(input.ReleaseTag) || input.AutoUpdateEnabled || input.CheckIntervalMinutes != 0 { return errors.New(errPagesSourceSelectorInvalid) } default: @@ -110,11 +108,17 @@ func buildGitHubSourceConfig(input SourceUpdateInput) (githubSourceConfig, error if selector == githubReleaseSelectorLatest && interval == 0 { interval = defaultCheckInterval } + autoUpdate := input.AutoUpdateEnabled + if selector == githubReleaseSelectorTag { + autoUpdate = false + interval = 0 + } return githubSourceConfig{ Repository: repository, Selector: selector, Tag: tag, AssetName: assetName, + AutoUpdate: autoUpdate, CheckInterval: interval, SourceIdentity: buildGitHubSourceIdentity(repository, selector, tag, assetName), }, nil @@ -251,7 +255,7 @@ func createGitHubSourceTx(tx *gorm.DB, projectID uint, config githubSourceConfig ReleaseSelector: config.Selector, ReleaseTag: config.Tag, AssetName: config.AssetName, - AutoUpdateEnabled: false, + AutoUpdateEnabled: config.AutoUpdate, CheckIntervalMinutes: config.CheckInterval, ConfigVersion: 1, SourceIdentity: config.SourceIdentity, @@ -276,7 +280,7 @@ func githubSourceUpdates(config githubSourceConfig, version int) map[string]any "release_selector": config.Selector, "release_tag": config.Tag, "asset_name": config.AssetName, - sourceColumnAutoUpdateEnabled: false, + sourceColumnAutoUpdateEnabled: config.AutoUpdate, "check_interval_minutes": config.CheckInterval, sourceColumnConfigVersion: version, "source_identity": config.SourceIdentity, @@ -287,7 +291,7 @@ func githubSourceConfigChanged(existing *model.PagesProjectSource, config github return existing.SourceType != PagesSourceTypeGitHubRelease || existing.RemoteURL != "" || existing.RemoteNetworkPolicy != "" || existing.GitHubRepository != config.Repository || existing.ReleaseSelector != config.Selector || existing.ReleaseTag != config.Tag || - existing.AssetName != config.AssetName || existing.AutoUpdateEnabled || + existing.AssetName != config.AssetName || existing.AutoUpdateEnabled != config.AutoUpdate || existing.CheckIntervalMinutes != config.CheckInterval } diff --git a/internal/apps/openflare/pages/github_source_action.go b/internal/apps/openflare/pages/github_source_action.go index d4fc1e89..6aaea778 100644 --- a/internal/apps/openflare/pages/github_source_action.go +++ b/internal/apps/openflare/pages/github_source_action.go @@ -28,9 +28,10 @@ const githubSourceDetailProvider = "github" var githubDigestPattern = regexp.MustCompile(`^sha256:[0-9a-f]{64}$`) type githubSourceProviderDomainError struct { - message string - permanent bool - retryAt *time.Time + message string + permanent bool + retryAt *time.Time + statusCode int } func (domainError *githubSourceProviderDomainError) Error() string { @@ -56,9 +57,12 @@ type githubSourceTarget struct { } type githubCheckTaskResult struct { - Message string - Detail string - Stale bool + Message string + Detail string + Revision string + Status string + RetryAt *time.Time + Stale bool } type preparedGitHubSource struct { @@ -99,13 +103,21 @@ func checkGitHubSource( return nil, domainErr } if result.NotModified { - if err := finishGitHubCheckNotModified(ctx, snapshot, result); err != nil { + revision, status, err := finishGitHubCheckNotModified(ctx, snapshot, result) + if err != nil { if errors.Is(err, errSourceFinalFence) { return &githubCheckTaskResult{Message: errPagesSourceActionStale, Stale: true}, nil } return nil, err } - return &githubCheckTaskResult{Message: "GitHub Release 检查完成,内容未变化"}, nil + detail, _ := json.Marshal(map[string]string{"revision": revision, pagesDeploymentColumnStatus: status}) + return &githubCheckTaskResult{ + Message: "GitHub Release 检查完成,内容未变化", + Detail: string(detail), + Revision: revision, + Status: status, + RetryAt: result.RetryAt, + }, nil } target, err := buildGitHubSourceTarget(result.Release, result.Asset, result.RetryAt) if err != nil { @@ -136,7 +148,10 @@ func checkGitHubSource( case pagesSourceStatusAttention: message = "检测到同一 Release 的资源被替换,需要确认" } - return &githubCheckTaskResult{Message: message, Detail: string(detail)}, nil + return &githubCheckTaskResult{ + Message: message, Detail: string(detail), Revision: target.Revision, + Status: status, RetryAt: result.RetryAt, + }, nil } func buildGitHubSourceTarget( @@ -183,17 +198,22 @@ func finishGitHubCheckNotModified( ctx context.Context, snapshot *sourceExecutionSnapshot, result githubrelease.ResolveResult, -) error { - return db.DB(ctx).Transaction(func(tx *gorm.DB) error { +) (string, string, error) { + var revision string + var status string + err := db.DB(ctx).Transaction(func(tx *gorm.DB) error { runtime, now, err := lockOwnedSourceRuntime(tx, snapshot) if err != nil { return err } + revision = runtime.LastSeenRevision + status = normalizedSourceRuntimeStatus(runtime) updates := githubCheckTerminalUpdates(snapshot, now, result.RetryAt) updates["etag"] = result.ETag - updates[sourceRuntimeColumnSyncStatus] = normalizedSourceRuntimeStatus(runtime) + updates[sourceRuntimeColumnSyncStatus] = status return tx.Model(runtime).Updates(updates).Error }) + return revision, status, err } func finishGitHubCheckTarget( @@ -230,7 +250,7 @@ func githubCheckTerminalUpdates( sourceRuntimeColumnLeaseToken: "", sourceRuntimeColumnLeaseExpiresAt: nil, } - updates["next_check_at"] = nextCheckAfterGitHubResponse(snapshot, now, retryAt) + updates[sourceRuntimeColumnNextCheckAt] = nextCheckAfterGitHubResponse(snapshot, now, retryAt) return updates } @@ -244,7 +264,7 @@ func nextCheckAfterGitHubResponse( } next := nextGitHubCheckAt(now, snapshot.SourceID, snapshot.CheckIntervalMinutes) if retryAt != nil && retryAt.After(next) { - next = retryAt.UTC() + next = retryAt.In(now.Location()) } return &next } @@ -275,7 +295,7 @@ func failGitHubCheckLease( now := time.Now() next := now.Add(initialCheckRetryDelay) if retryAt.After(next) { - next = retryAt.UTC() + next = retryAt.In(now.Location()) } updates := map[string]any{ sourceRuntimeColumnSyncStatus: pagesSourceStatusFailed, @@ -285,9 +305,9 @@ func failGitHubCheckLease( sourceRuntimeColumnLeaseExpiresAt: nil, } if snapshot.ReleaseSelector == githubReleaseSelectorLatest { - updates["next_check_at"] = &next + updates[sourceRuntimeColumnNextCheckAt] = &next } else { - updates["next_check_at"] = nil + updates[sourceRuntimeColumnNextCheckAt] = nil } result := db.DB(ctx).Model(&model.PagesProjectSourceRuntime{}). Where("source_id = ? AND lease_token = ? AND lease_expires_at > ?", snapshot.SourceID, snapshot.LeaseToken, now). @@ -332,16 +352,20 @@ func preflightGitHubSyncConfirmation(ctx context.Context, sourceID uint, confirm return nil } -func syncGitHubSource( +func syncGitHubSourceWithTrigger( ctx context.Context, snapshot *sourceExecutionSnapshot, actor string, targetRevision string, confirmedRevision string, + triggerType string, ) (outcome *sourceSyncOutcome, resultErr error) { if snapshot == nil || snapshot.SourceType != PagesSourceTypeGitHubRelease || !validPagesSourceActor(actor) { return nil, errors.New(errPagesSourceActionInvalid) } + if !validSourceDeploymentTrigger(triggerType) { + return nil, errors.New(errPagesSourceActionInvalid) + } defer func() { resultErr = finalizeGitHubSyncFailure(ctx, snapshot, resultErr) }() @@ -383,7 +407,7 @@ func syncGitHubSource( if !renewed { return &sourceSyncOutcome{Stale: true}, nil } - return activatePreparedGitHubSource(ctx, snapshot, actor, prepared) + return activatePreparedGitHubSource(ctx, snapshot, actor, triggerType, prepared) } func finalizeGitHubSyncFailure( @@ -429,12 +453,13 @@ func activatePreparedGitHubSource( ctx context.Context, snapshot *sourceExecutionSnapshot, actor string, + triggerType string, prepared *preparedGitHubSource, ) (*sourceSyncOutcome, error) { task.AppendLog(ctx, "[activate] 正在原子切换 GitHub Release 部署") - deployment, reused, referenced, err := commitSourceDeployment( + deployment, reused, referenced, err := commitSourceDeploymentWithTrigger( ctx, snapshot, prepared.target.Revision, prepared.download.SHA256, - prepared.target.Detail, prepared.target.DetailJSON, actor, prepared.manifest, + prepared.target.Detail, prepared.target.DetailJSON, actor, triggerType, prepared.manifest, prepared.ingestState.Result, prepared.ingestState.HasIngest, prepared.target.RetryAt, ) prepared.ingestState.Referenced = referenced @@ -634,7 +659,7 @@ func releaseGitHubSyncWithoutActivation( if expedite && snapshot.ReleaseSelector == githubReleaseSelectorLatest { next := now.Add(initialCheckRetryDelay) if retryAt != nil && retryAt.After(next) { - next = retryAt.UTC() + next = retryAt.In(now.Location()) } nextCheckAt = &next } @@ -644,7 +669,7 @@ func releaseGitHubSyncWithoutActivation( sourceRuntimeColumnSyncStatus: status, sourceRuntimeColumnLastError: lastError, sourceRuntimeColumnLastCheckedAt: &now, - "next_check_at": nextCheckAt, + sourceRuntimeColumnNextCheckAt: nextCheckAt, sourceRuntimeColumnLeaseToken: "", sourceRuntimeColumnLeaseExpiresAt: nil, } @@ -689,32 +714,37 @@ func safeGitHubSourceError(err error) string { func githubSourceDomainError(err error) error { message := errPagesSourceSyncFailed + statusCode := 0 + var providerError *githubrelease.Error + if errors.As(err, &providerError) { + statusCode = providerError.StatusCode + } retryAt, hasRetryAt := githubrelease.RetryAt(err) var retryDeadline *time.Time if hasRetryAt { retryDeadline = &retryAt } if err == nil { - return &githubSourceProviderDomainError{message: message, permanent: false} + return &githubSourceProviderDomainError{message: message, permanent: false, statusCode: statusCode} } if githubrelease.IsDigestError(err) { message = errPagesSourceDigestMismatch - return &githubSourceProviderDomainError{message: message, permanent: true} + return &githubSourceProviderDomainError{message: message, permanent: true, statusCode: statusCode} } if githubrelease.IsNotFound(err) { message = errPagesSourceReleaseNotFound - return &githubSourceProviderDomainError{message: message, permanent: true} + return &githubSourceProviderDomainError{message: message, permanent: true, statusCode: statusCode} } if errors.Is(err, githubrelease.ErrAssetTooLarge) { message = errPagesPackageURLTooLarge - return &githubSourceProviderDomainError{message: message, permanent: true} + return &githubSourceProviderDomainError{message: message, permanent: true, statusCode: statusCode} } if errors.Is(err, githubrelease.ErrEmptyAsset) { message = errPagesPackageEmpty - return &githubSourceProviderDomainError{message: message, permanent: true} + return &githubSourceProviderDomainError{message: message, permanent: true, statusCode: statusCode} } return &githubSourceProviderDomainError{ - message: message, permanent: !githubrelease.IsRetryable(err), retryAt: retryDeadline, + message: message, permanent: !githubrelease.IsRetryable(err), retryAt: retryDeadline, statusCode: statusCode, } } diff --git a/internal/apps/openflare/pages/github_source_test.go b/internal/apps/openflare/pages/github_source_test.go index 6d823319..61b5e340 100644 --- a/internal/apps/openflare/pages/github_source_test.go +++ b/internal/apps/openflare/pages/github_source_test.go @@ -190,7 +190,7 @@ func TestGitHubSourceSaveSurvivesInitialCheckDispatchFailure(t *testing.T) { } } -func TestGitHubSourceRejectsUnsafeOrPhaseThreeFields(t *testing.T) { +func TestGitHubSourceRejectsUnsafeOrModeIncompatibleFields(t *testing.T) { tests := []SourceUpdateInput{ {SourceType: PagesSourceTypeGitHubRelease, RepositoryURL: "http://github.com/a/b"}, {SourceType: PagesSourceTypeGitHubRelease, RepositoryURL: "https://github.com/a%20b/repo"}, @@ -202,8 +202,8 @@ func TestGitHubSourceRejectsUnsafeOrPhaseThreeFields(t *testing.T) { {SourceType: PagesSourceTypeGitHubRelease, RepositoryURL: "https://github.com/a/b", AssetName: "dist\n.zip"}, {SourceType: PagesSourceTypeGitHubRelease, RepositoryURL: "https://github.com/a/b", AssetName: "dist\u202e.zip"}, {SourceType: PagesSourceTypeGitHubRelease, RepositoryURL: "https://github.com/a/b", AssetName: "dir/dist.zip"}, - {SourceType: PagesSourceTypeGitHubRelease, RepositoryURL: "https://github.com/a/b", AutoUpdateEnabled: true}, {SourceType: PagesSourceTypeGitHubRelease, RepositoryURL: "https://github.com/a/b", ReleaseSelector: "tag", ReleaseTag: "v1", CheckIntervalMinutes: 60}, + {SourceType: PagesSourceTypeGitHubRelease, RepositoryURL: "https://github.com/a/b", ReleaseSelector: "tag", ReleaseTag: "v1", AutoUpdateEnabled: true}, {SourceType: PagesSourceTypeGitHubRelease, RepositoryURL: "https://github.com/a/b", ReleaseSelector: "tag", ReleaseTag: " v1"}, {SourceType: PagesSourceTypeGitHubRelease, RepositoryURL: "https://github.com/a/b", ReleaseSelector: "tag", ReleaseTag: "v1\n"}, {SourceType: PagesSourceTypeGitHubRelease, RepositoryURL: "https://github.com/a/b", ReleaseSelector: "tag", ReleaseTag: "v1\u2028draft"}, @@ -779,6 +779,7 @@ func TestSourceActionPayloadSeparatesSystemTargetAndUserConfirmation(t *testing. revision := strings.Repeat("a", 64) invalid := []SourceActionPayload{ {SourceID: 1, ConfigVersion: 1, Action: sourceActionSync, Actor: "user:1", TargetRevision: revision}, + {SourceID: 1, ConfigVersion: 1, Action: sourceActionSync, Actor: pagesSourceCreatedBySystem, TriggerType: pagesSourceTriggerManualSync}, {SourceID: 1, ConfigVersion: 1, Action: sourceActionSync, Actor: pagesSourceCreatedBySystem, ConfirmedRevision: revision}, {SourceID: 1, ConfigVersion: 1, Action: sourceActionSync, Actor: pagesSourceCreatedBySystem, TargetRevision: revision, ConfirmedRevision: revision}, } @@ -789,8 +790,8 @@ func TestSourceActionPayloadSeparatesSystemTargetAndUserConfirmation(t *testing. } } valid := []SourceActionPayload{ - {SourceID: 1, ConfigVersion: 1, Action: sourceActionSync, Actor: pagesSourceCreatedBySystem, TargetRevision: revision}, - {SourceID: 1, ConfigVersion: 1, Action: sourceActionSync, Actor: "user:1", ConfirmedRevision: revision}, + {SourceID: 1, ConfigVersion: 1, Action: sourceActionSync, Actor: pagesSourceCreatedBySystem, TriggerType: pagesSourceTriggerScheduledAutoUpdate, TargetRevision: revision}, + {SourceID: 1, ConfigVersion: 1, Action: sourceActionSync, Actor: "user:1", TriggerType: pagesSourceTriggerManualSync, ConfirmedRevision: revision}, } for _, payload := range valid { raw, _ := json.Marshal(payload) diff --git a/internal/apps/openflare/pages/source_manual_test_helpers_test.go b/internal/apps/openflare/pages/source_manual_test_helpers_test.go new file mode 100644 index 00000000..d54fc991 --- /dev/null +++ b/internal/apps/openflare/pages/source_manual_test_helpers_test.go @@ -0,0 +1,51 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package pages + +import ( + "context" + "time" + + "github.com/Rain-kl/Wavelet/internal/apps/upload" + "github.com/Rain-kl/Wavelet/internal/model" +) + +func syncRemoteSource( + ctx context.Context, + snapshot *sourceExecutionSnapshot, + actor string, +) (*sourceSyncOutcome, error) { + return syncRemoteSourceWithTrigger(ctx, snapshot, actor, pagesSourceTriggerManualSync) +} + +func syncGitHubSource( + ctx context.Context, + snapshot *sourceExecutionSnapshot, + actor string, + targetRevision string, + confirmedRevision string, +) (*sourceSyncOutcome, error) { + return syncGitHubSourceWithTrigger( + ctx, snapshot, actor, targetRevision, confirmedRevision, pagesSourceTriggerManualSync, + ) +} + +func commitSourceDeployment( + ctx context.Context, + snapshot *sourceExecutionSnapshot, + revision string, + packageChecksum string, + detail sourceDetail, + detailJSON string, + actor string, + manifest *deploymentManifest, + ingestResult upload.IngestResult, + hasIngest bool, + nextCheckNotBefore *time.Time, +) (*model.PagesDeployment, bool, bool, error) { + return commitSourceDeploymentWithTrigger( + ctx, snapshot, revision, packageChecksum, detail, detailJSON, actor, + pagesSourceTriggerManualSync, manifest, ingestResult, hasIngest, nextCheckNotBefore, + ) +} diff --git a/internal/apps/openflare/pages/source_runtime.go b/internal/apps/openflare/pages/source_runtime.go index 98801661..fdefb536 100644 --- a/internal/apps/openflare/pages/source_runtime.go +++ b/internal/apps/openflare/pages/source_runtime.go @@ -29,6 +29,7 @@ const ( sourceRuntimeColumnSyncStatus = "sync_status" sourceRuntimeColumnLastError = "last_error" sourceRuntimeColumnLastCheckedAt = "last_checked_at" + sourceRuntimeColumnNextCheckAt = "next_check_at" sourceRuntimeColumnLeaseToken = "lease_token" sourceRuntimeColumnLeaseExpiresAt = "lease_expires_at" pagesDeploymentColumnStatus = "status" @@ -58,6 +59,7 @@ type sourceExecutionSnapshot struct { ReleaseSelector string ReleaseTag string AssetName string + AutoUpdateEnabled bool CheckIntervalMinutes int ETag string LastSeenRevision string @@ -179,6 +181,7 @@ func loadSourceExecutionSnapshot( ReleaseSelector: source.ReleaseSelector, ReleaseTag: source.ReleaseTag, AssetName: source.AssetName, + AutoUpdateEnabled: source.AutoUpdateEnabled, CheckIntervalMinutes: source.CheckIntervalMinutes, ETag: runtime.ETag, LastSeenRevision: runtime.LastSeenRevision, @@ -300,6 +303,41 @@ func sourceLeaseIsBusy(ctx context.Context, sourceID uint) (bool, error) { return runtime.LeaseExpiresAt != nil && runtime.LeaseExpiresAt.After(time.Now()), nil } +// recoverExpiredSourceLease clears one exact expired lease owner. Matching the +// token, observed expiry and status prevents a scanner from overwriting a +// worker that renewed or was replaced after the candidate query. +func recoverExpiredSourceLease( + ctx context.Context, + sourceID uint, + token string, + expiresAt time.Time, + status string, + now time.Time, + nextCheckAt *time.Time, +) (bool, error) { + if sourceID == 0 || token == "" || + (status != pagesSourceStatusChecking && status != pagesSourceStatusSyncing) { + return false, nil + } + updates := map[string]any{ + sourceRuntimeColumnSyncStatus: pagesSourceStatusFailed, + sourceRuntimeColumnLastError: errPagesSourceLeaseExpired, + sourceRuntimeColumnLeaseToken: "", + sourceRuntimeColumnLeaseExpiresAt: nil, + sourceRuntimeColumnNextCheckAt: nextCheckAt, + } + result := db.DB(ctx).Model(&model.PagesProjectSourceRuntime{}). + Where("source_id = ?", sourceID). + Where("lease_token = ?", token). + Where("lease_expires_at = ? AND lease_expires_at <= ?", expiresAt, now). + Where("sync_status = ?", status). + Updates(updates) + if result.Error != nil { + return false, result.Error + } + return result.RowsAffected == 1, nil +} + // fenceAndNormalizeRuntime invalidates in-flight work while preserving safe // seen/applied cursors. The caller must already hold the source row lock. func fenceAndNormalizeRuntime(tx *gorm.DB, sourceID uint) error { diff --git a/internal/apps/openflare/pages/source_scanner.go b/internal/apps/openflare/pages/source_scanner.go new file mode 100644 index 00000000..29fc0dc8 --- /dev/null +++ b/internal/apps/openflare/pages/source_scanner.go @@ -0,0 +1,478 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package pages + +import ( + "bytes" + "context" + "encoding/json" + "errors" + "fmt" + "time" + + "github.com/Rain-kl/Wavelet/internal/db" + "github.com/Rain-kl/Wavelet/internal/model" + "github.com/Rain-kl/Wavelet/internal/task" + "github.com/Rain-kl/Wavelet/pkg/logger" + "gorm.io/gorm" +) + +const ( + // PagesSourceScanTask is the private Asynq task type for the periodic scanner. + PagesSourceScanTask = "openflare:pages_source_scan" + // TaskTypePagesSourceScan is the internal task meta type seeded in w_schedules. + TaskTypePagesSourceScan = "of_pages_source_scan" + + pagesSourceScanBatchSize = 20 +) + +// PagesSourceScanMeta is available to the scheduler registry but hidden from +// generic Admin task dispatch and schedule mutation APIs. +var PagesSourceScanMeta = task.TaskMeta{ + Type: TaskTypePagesSourceScan, + AsynqTask: PagesSourceScanTask, + Name: "OpenFlare Pages 部署源扫描", + Description: "补偿孤儿部署包、恢复过期执行权并串行检查到期的 GitHub latest 部署源", + SupportsTime: false, + MaxRetry: 0, + Queue: task.QueueDefault, + Retryable: false, + InternalOnly: true, +} + +type pagesSourceScanPayload struct{} + +type pagesSourceScanSummary struct { + ExpiredCandidates int `json:"expired_candidates"` + RecoveredLeases int `json:"recovered_leases"` + OrphanCleanup PagesOrphanCleanupSummary `json:"orphan_cleanup"` + DueSources int `json:"due_sources"` + SelectedSources int `json:"selected_sources"` + CheckedSources int `json:"checked_sources"` + UpdatesFound int `json:"updates_found"` + AttentionSources int `json:"attention_sources"` + DispatchedSyncs int `json:"dispatched_syncs"` + FailedDispatches int `json:"failed_dispatches"` + BusySources int `json:"busy_sources"` + StaleSources int `json:"stale_sources"` + FailedSources int `json:"failed_sources"` + Backlog int `json:"backlog"` + ProviderBackoffs []pagesSourceProviderBackoff `json:"provider_backoffs,omitempty"` +} + +type pagesSourceProviderBackoff struct { + SourceID uint `json:"source_id"` + StatusCode int `json:"status_code"` + RetryAt string `json:"retry_at"` +} + +type expiredSourceLeaseCandidate struct { + SourceID uint + LeaseToken string + LeaseExpiresAt time.Time + SyncStatus string + SourceType string + ReleaseSelector string +} + +type dueGitHubSourceCandidate struct { + SourceID uint + ConfigVersion int +} + +var ( + pagesSourceScanNow = time.Now + reconcilePagesSourceOrphans = ReconcilePagesOrphanUploads + dispatchPagesSourceAutoSync = func( + ctx context.Context, + source model.PagesProjectSource, + targetRevision string, + ) (*SourceActionReceipt, error) { + return dispatchSourceActionSnapshotWithTrigger( + ctx, + source, + sourceActionSync, + pagesSourceCreatedBySystem, + pagesSourceTriggerScheduledAutoUpdate, + targetRevision, + "", + "system", + ) + } +) + +// SourceScanHandler serializes provider checks inside one scheduled task. A +// source-level lease still permits overlapping scanner executions safely. +type SourceScanHandler struct{} + +// ValidatePayload accepts only an empty object; the scanner has no user input. +func (handler *SourceScanHandler) ValidatePayload(payload []byte) ([]byte, error) { + if len(bytes.TrimSpace(payload)) == 0 { + payload = []byte("{}") + } + var input pagesSourceScanPayload + decoder := json.NewDecoder(bytes.NewReader(payload)) + decoder.DisallowUnknownFields() + if err := decoder.Decode(&input); err != nil { + return nil, errors.New(errPagesSourceActionInvalid) + } + if err := ensureJSONEOF(decoder); err != nil { + return nil, errors.New(errPagesSourceActionInvalid) + } + return []byte("{}"), nil +} + +// Execute recovers expired leases and checks at most 20 due latest sources in +// stable order. Provider and dispatch failures are isolated per source. +func (handler *SourceScanHandler) Execute(ctx context.Context, payload []byte) (*task.TaskResult, error) { + if _, err := handler.ValidatePayload(payload); err != nil { + return nil, task.PermanentError(errPagesSourceActionInvalid) + } + + now := pagesSourceScanNow() + summary := pagesSourceScanSummary{} + if err := recoverExpiredPagesSourceLeases(ctx, now, &summary); err != nil { + return nil, err + } + orphanSummary, err := reconcilePagesSourceOrphans(ctx, now) + if err != nil { + return nil, err + } + summary.OrphanCleanup = orphanSummary + task.AppendLog( + ctx, + "[cleanup] orphan 候选=%d,已补偿=%d,仍被引用=%d,lease busy=%d,非法 marker=%d,跳过=%d,失败=%d", + orphanSummary.Candidates, + orphanSummary.Reconciled, + orphanSummary.Referenced, + orphanSummary.LeaseBusy, + orphanSummary.InvalidMarker, + orphanSummary.Skipped, + orphanSummary.Failed, + ) + if err := scanDueGitHubSources(ctx, now, &summary); err != nil { + return nil, err + } + + detail, err := json.Marshal(summary) + if err != nil { + return nil, err + } + message := fmt.Sprintf( + "Pages 部署源扫描完成:恢复 %d 个租约,补偿 %d 个孤儿记录,检查 %d 个来源,投递 %d 个自动更新,积压 %d 个", + summary.RecoveredLeases, + summary.OrphanCleanup.Reconciled, + summary.CheckedSources, + summary.DispatchedSyncs, + summary.Backlog, + ) + return &task.TaskResult{Message: message, Detail: string(detail)}, nil +} + +func recoverExpiredPagesSourceLeases( + ctx context.Context, + now time.Time, + summary *pagesSourceScanSummary, +) error { + var candidates []expiredSourceLeaseCandidate + err := db.DB(ctx). + Table("of_pages_project_source_runtime AS runtime"). + Select(`runtime.source_id, runtime.lease_token, runtime.lease_expires_at, + runtime.sync_status, source.source_type, source.release_selector`). + Joins("JOIN of_pages_project_sources AS source ON source.id = runtime.source_id"). + Where("runtime.lease_token <> ''"). + Where("runtime.lease_expires_at IS NOT NULL AND runtime.lease_expires_at <= ?", now). + Where("runtime.sync_status IN ?", []string{pagesSourceStatusChecking, pagesSourceStatusSyncing}). + Order("runtime.source_id ASC"). + Scan(&candidates).Error + if err != nil { + return err + } + summary.ExpiredCandidates = len(candidates) + for _, candidate := range candidates { + var nextCheckAt *time.Time + if candidate.SourceType == PagesSourceTypeGitHubRelease && + candidate.ReleaseSelector == githubReleaseSelectorLatest { + next := nextGitHubCheckAt(now, candidate.SourceID, minimumCheckInterval) + nextCheckAt = &next + } + recovered, recoverErr := recoverExpiredSourceLease( + ctx, + candidate.SourceID, + candidate.LeaseToken, + candidate.LeaseExpiresAt, + candidate.SyncStatus, + now, + nextCheckAt, + ) + if recoverErr != nil { + summary.FailedSources++ + logger.WarnF( + ctx, + "[PagesSourceScan] recover expired lease failed: source_id=%d error=%v", + candidate.SourceID, + recoverErr, + ) + continue + } + if recovered { + summary.RecoveredLeases++ + task.AppendLog(ctx, "[recover] 已恢复过期来源租约:source_id=%d", candidate.SourceID) + } + } + return nil +} + +func scanDueGitHubSources( + ctx context.Context, + now time.Time, + summary *pagesSourceScanSummary, +) error { + dueQuery := func() *gorm.DB { + return db.DB(ctx). + Table("of_pages_project_source_runtime AS runtime"). + Joins("JOIN of_pages_project_sources AS source ON source.id = runtime.source_id"). + Where("source.source_type = ?", PagesSourceTypeGitHubRelease). + Where("source.release_selector = ?", githubReleaseSelectorLatest). + Where("runtime.next_check_at IS NOT NULL AND runtime.next_check_at <= ?", now) + } + var dueCount int64 + if err := dueQuery().Count(&dueCount).Error; err != nil { + return err + } + summary.DueSources = int(dueCount) + + var candidates []dueGitHubSourceCandidate + if err := dueQuery(). + Select("source.id AS source_id, source.config_version"). + Order("runtime.next_check_at ASC"). + Order("source.id ASC"). + Limit(pagesSourceScanBatchSize). + Scan(&candidates).Error; err != nil { + return err + } + summary.SelectedSources = len(candidates) + task.AppendLog( + ctx, + "[scan] 到期来源=%d,本批=%d", + summary.DueSources, + summary.SelectedSources, + ) + + for _, candidate := range candidates { + scanOneDueGitHubSource(ctx, candidate, summary) + } + var remainingDue int64 + if err := dueQuery().Count(&remainingDue).Error; err != nil { + return err + } + summary.Backlog = int(remainingDue) + task.AppendLog(ctx, "[scan] 本批处理后仍到期来源=%d", summary.Backlog) + return nil +} + +func scanOneDueGitHubSource( + ctx context.Context, + candidate dueGitHubSourceCandidate, + summary *pagesSourceScanSummary, +) { + snapshot, outcome, err := acquireSourceLease( + ctx, + candidate.SourceID, + candidate.ConfigVersion, + sourceActionCheck, + ) + if err != nil { + summary.FailedSources++ + logger.WarnF(ctx, "[PagesSourceScan] acquire check lease failed: source_id=%d error=%v", candidate.SourceID, err) + return + } + switch outcome { + case sourceLeaseBusy: + summary.BusySources++ + task.AppendLog(ctx, "[check] 来源正在执行其它任务,跳过:source_id=%d", candidate.SourceID) + return + case sourceLeaseStale: + summary.StaleSources++ + return + } + if snapshot == nil || snapshot.SourceType != PagesSourceTypeGitHubRelease || + snapshot.ReleaseSelector != githubReleaseSelectorLatest { + summary.StaleSources++ + if snapshot != nil { + if finalizeErr := failSourceLease(ctx, snapshot, errPagesSourceActionStale); finalizeErr != nil { + logger.WarnF( + ctx, + "[PagesSourceScan] finalize stale source failed: source_id=%d error=%v", + snapshot.SourceID, + finalizeErr, + ) + } + } + return + } + + checkResult, checkErr := checkGitHubSource(ctx, snapshot) + if checkErr != nil { + summary.FailedSources++ + recordPagesSourceProviderBackoff(ctx, candidate.SourceID, checkErr, summary) + logger.WarnF( + ctx, + "[PagesSourceScan] source check failed: source_id=%d error=%s", + candidate.SourceID, + safeGitHubSourceError(checkErr), + ) + return + } + if checkResult == nil || checkResult.Stale { + summary.StaleSources++ + return + } + handleCheckedGitHubSource(ctx, snapshot, checkResult, summary) +} + +func recordPagesSourceProviderBackoff( + ctx context.Context, + sourceID uint, + checkErr error, + summary *pagesSourceScanSummary, +) { + var domainError *githubSourceProviderDomainError + if !errors.As(checkErr, &domainError) || + (domainError.statusCode != 403 && domainError.statusCode != 429) { + return + } + + retryAt := domainError.retryAt + var runtime model.PagesProjectSourceRuntime + if err := db.DB(ctx). + Select("next_check_at"). + Where("source_id = ?", sourceID). + First(&runtime).Error; err != nil { + logger.WarnF(ctx, "[PagesSourceScan] load provider backoff deadline failed: source_id=%d error=%v", sourceID, err) + } else if runtime.NextCheckAt != nil { + retryAt = runtime.NextCheckAt + } + + retryAtText := "unknown" + if retryAt != nil { + retryAtText = retryAt.UTC().Format(time.RFC3339) + } + summary.ProviderBackoffs = append(summary.ProviderBackoffs, pagesSourceProviderBackoff{ + SourceID: sourceID, StatusCode: domainError.statusCode, RetryAt: retryAtText, + }) + task.AppendLog( + ctx, + "[check] GitHub provider 退避:source_id=%d status=%d retry_at=%s", + sourceID, + domainError.statusCode, + retryAtText, + ) +} + +func handleCheckedGitHubSource( + ctx context.Context, + snapshot *sourceExecutionSnapshot, + checkResult *githubCheckTaskResult, + summary *pagesSourceScanSummary, +) { + summary.CheckedSources++ + switch checkResult.Status { + case pagesSourceStatusUpdateAvailable: + summary.UpdatesFound++ + case pagesSourceStatusAttention: + summary.AttentionSources++ + } + if !snapshot.AutoUpdateEnabled || checkResult.Status != pagesSourceStatusUpdateAvailable || + !validOptionalSourceRevision(checkResult.Revision) || checkResult.Revision == "" { + return + } + + source := model.PagesProjectSource{ + ID: snapshot.SourceID, + ProjectID: snapshot.ProjectID, + ConfigVersion: snapshot.SourceConfigVersion, + } + receipt, dispatchErr := dispatchPagesSourceAutoSync(ctx, source, checkResult.Revision) + if dispatchErr == nil { + summary.DispatchedSyncs++ + if receipt != nil { + task.AppendLog( + ctx, + "[dispatch] 已投递自动更新:source_id=%d execution_id=%s revision=%s", + snapshot.SourceID, + receipt.ExecutionID, + checkResult.Revision, + ) + } + return + } + + summary.FailedSources++ + summary.FailedDispatches++ + logger.WarnF( + ctx, + "[PagesSourceScan] dispatch auto sync failed: source_id=%d revision=%s error=%v", + snapshot.SourceID, + checkResult.Revision, + dispatchErr, + ) + updated, recordErr := recordPagesSourceAutoDispatchFailure( + ctx, + snapshot, + checkResult.Revision, + checkResult.RetryAt, + ) + if recordErr != nil { + logger.WarnF( + ctx, + "[PagesSourceScan] record auto sync dispatch failure failed: source_id=%d error=%v", + snapshot.SourceID, + recordErr, + ) + } else if !updated { + summary.StaleSources++ + } +} + +func recordPagesSourceAutoDispatchFailure( + ctx context.Context, + snapshot *sourceExecutionSnapshot, + revision string, + retryAt *time.Time, +) (bool, error) { + if snapshot == nil || revision == "" { + return false, nil + } + now := pagesSourceScanNow() + next := nextGitHubCheckAt(now, snapshot.SourceID, minimumCheckInterval) + if retryAt != nil && retryAt.After(next) { + next = retryAt.In(now.Location()) + } + result := db.DB(ctx).Model(&model.PagesProjectSourceRuntime{}). + Where("source_id = ?", snapshot.SourceID). + Where("sync_status = ? AND last_seen_revision = ?", pagesSourceStatusUpdateAvailable, revision). + Where("lease_expires_at IS NULL OR lease_expires_at <= ?", now). + Where(`EXISTS ( + SELECT 1 FROM of_pages_project_sources AS source + WHERE source.id = ? AND source.config_version = ? + AND source.source_type = ? AND source.release_selector = ? + AND source.auto_update_enabled = ? + )`, + snapshot.SourceID, + snapshot.SourceConfigVersion, + PagesSourceTypeGitHubRelease, + githubReleaseSelectorLatest, + true, + ). + Updates(map[string]any{ + sourceRuntimeColumnSyncStatus: pagesSourceStatusUpdateAvailable, + sourceRuntimeColumnLastError: errPagesSourceTaskDispatchFailed, + sourceRuntimeColumnNextCheckAt: &next, + }) + if result.Error != nil { + return false, result.Error + } + return result.RowsAffected == 1, nil +} diff --git a/internal/apps/openflare/pages/source_scanner_test.go b/internal/apps/openflare/pages/source_scanner_test.go new file mode 100644 index 00000000..2d7c2394 --- /dev/null +++ b/internal/apps/openflare/pages/source_scanner_test.go @@ -0,0 +1,511 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package pages + +import ( + "context" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "errors" + "fmt" + "os" + "path/filepath" + "strings" + "testing" + "time" + + "github.com/Rain-kl/Wavelet/internal/db" + "github.com/Rain-kl/Wavelet/internal/integration/githubrelease" + "github.com/Rain-kl/Wavelet/internal/model" + "gorm.io/gorm" +) + +type scannerDispatchedSync struct { + SourceID uint + Revision string +} + +func TestGitHubLatestAutoConfigPreservesIdentityAndRuntimeCursor(t *testing.T) { + ctx := setupPagesSourceTest(t) + project := mustCreatePagesSourceProject(t, ctx, "scanner-auto-config") + source, _ := mustConfigureGitHubSourceWithoutDispatch(t, ctx, project.ID, SourceUpdateInput{ + SourceType: PagesSourceTypeGitHubRelease, + RepositoryURL: "https://github.com/scanner/auto-config", + AutoUpdateEnabled: false, + CheckIntervalMinutes: 60, + }) + identity := source.SourceIdentity + seenRevision := strings.Repeat("a", sourceRevisionHexLength) + appliedRevision := strings.Repeat("b", sourceRevisionHexLength) + future := time.Now().Add(time.Hour) + if err := db.DB(ctx).Model(&model.PagesProjectSourceRuntime{}). + Where("source_id = ?", source.ID). + Updates(map[string]any{ + "etag": `"cursor-etag"`, + "last_seen_revision": seenRevision, + "last_seen_detail": `{"provider":"github","release_id":"2","asset_id":"2","tag":"v2","asset_name":"dist.zip"}`, + "last_applied_revision": appliedRevision, + "last_applied_detail": `{"provider":"github","release_id":"1","asset_id":"1","tag":"v1","asset_name":"dist.zip"}`, + "sync_status": pagesSourceStatusSyncing, + "last_error": "old error", + "lease_token": "in-flight", + "lease_expires_at": &future, + }).Error; err != nil { + t.Fatalf("seed runtime error = %v, want nil", err) + } + + input := SourceUpdateInput{ + SourceType: PagesSourceTypeGitHubRelease, + RepositoryURL: "https://github.com/scanner/auto-config", + AutoUpdateEnabled: true, + CheckIntervalMinutes: 15, + } + if err := validateGitHubSourceInput(input); err != nil { + t.Fatalf("validateGitHubSourceInput(auto latest) error = %v, want nil", err) + } + if err := db.DB(ctx).Transaction(func(tx *gorm.DB) error { + changed, err := updateGitHubSourceTx(tx, project.ID, input) + if err == nil && !changed { + return errors.New("auto config update was treated as no-op") + } + return err + }); err != nil { + t.Fatalf("updateGitHubSourceTx(auto latest) error = %v, want nil", err) + } + + updated, runtime := mustLoadPagesSource(t, ctx, project.ID) + if updated.SourceIdentity != identity || updated.ConfigVersion != source.ConfigVersion+1 { + t.Errorf( + "updated source = identity:%q version:%d, want identity:%q version:%d", + updated.SourceIdentity, + updated.ConfigVersion, + identity, + source.ConfigVersion+1, + ) + } + if !updated.AutoUpdateEnabled || updated.CheckIntervalMinutes != 15 { + t.Errorf("updated auto config = enabled:%t interval:%d, want true/15", updated.AutoUpdateEnabled, updated.CheckIntervalMinutes) + } + if runtime.ETag != `"cursor-etag"` || runtime.LastSeenRevision != seenRevision || + runtime.LastAppliedRevision != appliedRevision { + t.Errorf("runtime cursor changed after auto-only update: %+v", runtime) + } + if runtime.LeaseToken != "" || runtime.LeaseExpiresAt != nil || runtime.LastError != "" || + runtime.SyncStatus != pagesSourceStatusUpdateAvailable || runtime.NextCheckAt == nil { + t.Errorf( + "runtime fence = token:%q expiry:%v error:%q status:%q next:%v", + runtime.LeaseToken, + runtime.LeaseExpiresAt, + runtime.LastError, + runtime.SyncStatus, + runtime.NextCheckAt, + ) + } + + tagConfig, err := buildGitHubSourceConfig(SourceUpdateInput{ + SourceType: PagesSourceTypeGitHubRelease, + RepositoryURL: "https://github.com/scanner/auto-config", + ReleaseSelector: githubReleaseSelectorTag, + ReleaseTag: "v1", + AutoUpdateEnabled: true, + CheckIntervalMinutes: 60, + }) + if err != nil { + t.Fatalf("buildGitHubSourceConfig(tag) error = %v, want nil", err) + } + if tagConfig.AutoUpdate || tagConfig.CheckInterval != 0 { + t.Errorf("tag config auto/interval = %t/%d, want false/0", tagConfig.AutoUpdate, tagConfig.CheckInterval) + } +} + +func TestRecoverExpiredPagesSourceLeaseUsesExactCASAndStableJitter(t *testing.T) { + ctx := setupPagesSourceTest(t) + project := mustCreatePagesSourceProject(t, ctx, "scanner-expired-lease") + source, _ := mustConfigureGitHubSourceWithoutDispatch(t, ctx, project.ID, SourceUpdateInput{ + SourceType: PagesSourceTypeGitHubRelease, + RepositoryURL: "https://github.com/scanner/expired-lease", + }) + now := time.Date(2026, 7, 19, 12, 0, 0, 0, time.UTC) + usePagesSourceScannerClock(t, now) + expiredAt := now.Add(-time.Minute) + if err := db.DB(ctx).Model(&model.PagesProjectSourceRuntime{}). + Where("source_id = ?", source.ID). + Updates(map[string]any{ + "sync_status": pagesSourceStatusChecking, + "lease_token": "expired-owner", + "lease_expires_at": &expiredAt, + "next_check_at": &expiredAt, + }).Error; err != nil { + t.Fatalf("seed expired lease error = %v, want nil", err) + } + + summary := pagesSourceScanSummary{} + if err := recoverExpiredPagesSourceLeases(ctx, now, &summary); err != nil { + t.Fatalf("recoverExpiredPagesSourceLeases() error = %v, want nil", err) + } + if summary.ExpiredCandidates != 1 || summary.RecoveredLeases != 1 || summary.FailedSources != 0 { + t.Errorf("recovery summary = %+v, want one recovered lease", summary) + } + _, runtime := mustLoadPagesSource(t, ctx, project.ID) + wantNext := nextGitHubCheckAt(now, source.ID, minimumCheckInterval) + if runtime.SyncStatus != pagesSourceStatusFailed || runtime.LastError != errPagesSourceLeaseExpired || + runtime.LeaseToken != "" || runtime.LeaseExpiresAt != nil || runtime.NextCheckAt == nil || + runtime.NextCheckAt.Sub(wantNext) != 0 { + t.Errorf( + "recovered runtime = status:%q error:%q token:%q expiry:%v next:%v, want next %v", + runtime.SyncStatus, + runtime.LastError, + runtime.LeaseToken, + runtime.LeaseExpiresAt, + runtime.NextCheckAt, + wantNext, + ) + } + + renewedExpiry := now.Add(time.Minute) + if err := db.DB(ctx).Model(&model.PagesProjectSourceRuntime{}). + Where("source_id = ?", source.ID). + Updates(map[string]any{ + "sync_status": pagesSourceStatusSyncing, + "lease_token": "renewed-owner", + "lease_expires_at": &renewedExpiry, + }).Error; err != nil { + t.Fatalf("seed renewed lease error = %v, want nil", err) + } + recovered, err := recoverExpiredSourceLease( + ctx, + source.ID, + "renewed-owner", + expiredAt, + pagesSourceStatusSyncing, + now, + &wantNext, + ) + if err != nil || recovered { + t.Fatalf("recoverExpiredSourceLease(stale expiry) = %t, %v; want false, nil", recovered, err) + } + _, runtime = mustLoadPagesSource(t, ctx, project.ID) + if runtime.LeaseToken != "renewed-owner" || runtime.LeaseExpiresAt == nil || + runtime.LeaseExpiresAt.Sub(renewedExpiry) != 0 || runtime.SyncStatus != pagesSourceStatusSyncing { + t.Errorf("stale recovery overwrote renewed lease: %+v", runtime) + } +} + +func TestScheduledAutoSyncPersistsExplicitDeploymentTrigger(t *testing.T) { + ctx := setupPagesSourceSyncTest(t) + project := mustCreatePagesSourceProject(t, ctx, "scanner-scheduled-trigger") + source, _ := mustConfigureGitHubSourceWithoutDispatch(t, ctx, project.ID, SourceUpdateInput{ + SourceType: PagesSourceTypeGitHubRelease, + RepositoryURL: "https://github.com/scanner/scheduled-trigger", + AutoUpdateEnabled: true, + }) + packageBytes := testPagesZip(t, map[string]string{"index.html": "scheduled-v1"}) + packageHash := sha256.Sum256(packageBytes) + release := githubrelease.Release{ID: "scheduled-release", Tag: "v1"} + asset := githubrelease.Asset{ + ID: "scheduled-asset", Name: defaultGitHubAssetName, State: "uploaded", + UpdatedAt: time.Date(2026, 7, 19, 12, 30, 0, 0, time.UTC), + } + target, err := buildGitHubSourceTarget(release, asset, nil) + if err != nil { + t.Fatalf("buildGitHubSourceTarget() error = %v, want nil", err) + } + useFakeGitHubReleaseClient(t, &fakeGitHubReleaseClient{ + resolve: func(context.Context, githubrelease.ResolveRequest) (githubrelease.ResolveResult, error) { + return githubrelease.ResolveResult{Release: release, Asset: asset}, nil + }, + download: func(context.Context, githubrelease.DownloadRequest) (*githubrelease.DownloadResult, error) { + path := filepath.Join(t.TempDir(), "scheduled.zip") + if err := os.WriteFile(path, packageBytes, 0o600); err != nil { + t.Fatalf("os.WriteFile(scheduled package) error = %v", err) + } + return &githubrelease.DownloadResult{ + Path: path, Size: int64(len(packageBytes)), SHA256: hex.EncodeToString(packageHash[:]), + }, nil + }, + }) + snapshot, outcome, err := acquireSourceLease(ctx, source.ID, source.ConfigVersion, sourceActionSync) + if err != nil || outcome != sourceLeaseAcquired { + t.Fatalf("acquireSourceLease() = %+v, %q, %v; want acquired", snapshot, outcome, err) + } + synced, err := syncGitHubSourceWithTrigger( + ctx, + snapshot, + pagesSourceCreatedBySystem, + target.Revision, + "", + pagesSourceTriggerScheduledAutoUpdate, + ) + if err != nil || synced == nil || synced.Deployment == nil || synced.Stale { + t.Fatalf("syncGitHubSourceWithTrigger() = %+v, %v; want active deployment", synced, err) + } + deployment, err := model.GetPagesDeploymentByID(ctx, synced.Deployment.ID) + if err != nil { + t.Fatalf("GetPagesDeploymentByID(%d) error = %v", synced.Deployment.ID, err) + } + if deployment.TriggerType != pagesSourceTriggerScheduledAutoUpdate || + deployment.CreatedBy != pagesSourceCreatedBySystem { + t.Errorf( + "scheduled provenance = trigger:%q actor:%q, want %q/%q", + deployment.TriggerType, + deployment.CreatedBy, + pagesSourceTriggerScheduledAutoUpdate, + pagesSourceCreatedBySystem, + ) + } +} + +func TestPagesSourceScannerSerialBatchIsolation304AndActualBacklog(t *testing.T) { + ctx := setupPagesSourceTest(t) + now := time.Now().Truncate(time.Second) + usePagesSourceScannerClock(t, now) + dueAt := now.Add(-time.Hour) + + type fixture struct { + source *model.PagesProjectSource + runtime *model.PagesProjectSourceRuntime + repository string + } + fixtures := make([]fixture, 0, 22) + byRepository := make(map[string]int, 22) + for index := 1; index <= 22; index++ { + project := mustCreatePagesSourceProject(t, ctx, fmt.Sprintf("scanner-batch-%02d", index)) + repository := fmt.Sprintf("scanner/source-%02d", index) + source, runtime := mustConfigureGitHubSourceWithoutDispatch(t, ctx, project.ID, SourceUpdateInput{ + SourceType: PagesSourceTypeGitHubRelease, + RepositoryURL: "https://github.com/" + repository, + AutoUpdateEnabled: index != 4, + CheckIntervalMinutes: 60, + }) + if err := db.DB(ctx).Model(&model.PagesProjectSourceRuntime{}). + Where("source_id = ?", source.ID). + Update("next_check_at", &dueAt).Error; err != nil { + t.Fatalf("mark source %d due error = %v, want nil", source.ID, err) + } + fixtures = append(fixtures, fixture{source: source, runtime: runtime, repository: repository}) + byRepository[repository] = index + } + + busyUntil := now.Add(time.Hour) + if err := db.DB(ctx).Model(&model.PagesProjectSourceRuntime{}). + Where("source_id = ?", fixtures[0].source.ID). + Updates(map[string]any{ + "sync_status": pagesSourceStatusChecking, + "lease_token": "busy-owner", + "lease_expires_at": &busyUntil, + }).Error; err != nil { + t.Fatalf("seed busy source error = %v, want nil", err) + } + + stored304Revision := strings.Repeat("3", sourceRevisionHexLength) + if err := db.DB(ctx).Model(&model.PagesProjectSourceRuntime{}). + Where("source_id = ?", fixtures[2].source.ID). + Updates(map[string]any{ + "etag": `"stored-etag"`, + "last_seen_revision": stored304Revision, + "last_seen_detail": `{"provider":"github","release_id":"release-3","asset_id":"3","tag":"v3","asset_name":"dist.zip"}`, + "sync_status": pagesSourceStatusUpdateAvailable, + }).Error; err != nil { + t.Fatalf("seed 304 cursor error = %v, want nil", err) + } + + if err := db.DB(ctx).Model(&model.PagesProjectSourceRuntime{}). + Where("source_id = ?", fixtures[4].source.ID). + Updates(map[string]any{ + "last_applied_revision": strings.Repeat("a", sourceRevisionHexLength), + "last_applied_detail": `{"provider":"github","release_id":"shared-release","asset_id":"old","tag":"v5","asset_name":"dist.zip"}`, + }).Error; err != nil { + t.Fatalf("seed replacement cursor error = %v, want nil", err) + } + + retryAt := now.Add(2 * time.Hour) + calledRepositories := make([]string, 0, pagesSourceScanBatchSize) + useFakeGitHubReleaseClient(t, &fakeGitHubReleaseClient{ + resolve: func(_ context.Context, request githubrelease.ResolveRequest) (githubrelease.ResolveResult, error) { + index := byRepository[request.Repository] + calledRepositories = append(calledRepositories, request.Repository) + switch index { + case 2: + return githubrelease.ResolveResult{}, &githubrelease.Error{ + Kind: githubrelease.ErrMetadata, StatusCode: 429, RetryAt: &retryAt, + } + case 3: + if request.ETag != `"stored-etag"` { + t.Errorf("304 source ETag = %q, want stored ETag", request.ETag) + } + return githubrelease.ResolveResult{NotModified: true, ETag: request.ETag}, nil + default: + releaseID := fmt.Sprintf("release-%d", index) + if index == 5 { + releaseID = "shared-release" + } + result := githubrelease.ResolveResult{ + ETag: fmt.Sprintf(`"etag-%d"`, index), + Release: githubrelease.Release{ID: releaseID, Tag: fmt.Sprintf("v%d", index)}, + Asset: githubrelease.Asset{ + ID: fmt.Sprintf("asset-%d", index), + Name: defaultGitHubAssetName, + State: "uploaded", + UpdatedAt: now.Add(time.Duration(index) * time.Minute), + }, + } + if index == 6 { + result.RetryAt = &retryAt + } + return result, nil + } + }, + download: func(context.Context, githubrelease.DownloadRequest) (*githubrelease.DownloadResult, error) { + t.Fatal("scanner downloaded an asset; want check-only behavior") + return nil, nil + }, + }) + + dispatched := make([]scannerDispatchedSync, 0, pagesSourceScanBatchSize) + previousDispatch := dispatchPagesSourceAutoSync + dispatchPagesSourceAutoSync = func( + _ context.Context, + source model.PagesProjectSource, + revision string, + ) (*SourceActionReceipt, error) { + if source.ID == fixtures[5].source.ID { + return nil, errors.New("injected dispatch failure") + } + dispatched = append(dispatched, scannerDispatchedSync{SourceID: source.ID, Revision: revision}) + return &SourceActionReceipt{ExecutionID: fmt.Sprintf("%d", source.ID), Action: sourceActionSync}, nil + } + t.Cleanup(func() { dispatchPagesSourceAutoSync = previousDispatch }) + + result, err := (&SourceScanHandler{}).Execute(ctx, []byte("{}")) + if err != nil { + t.Fatalf("SourceScanHandler.Execute() error = %v, want nil", err) + } + var summary pagesSourceScanSummary + if err := json.Unmarshal([]byte(result.Detail), &summary); err != nil { + t.Fatalf("json.Unmarshal(scan detail) error = %v, want nil", err) + } + if summary.DueSources != 22 || summary.SelectedSources != pagesSourceScanBatchSize || + summary.CheckedSources != 18 || summary.UpdatesFound != 17 || summary.AttentionSources != 1 || + summary.DispatchedSyncs != 15 || summary.FailedDispatches != 1 || + summary.BusySources != 1 || summary.FailedSources != 2 || + summary.Backlog != 3 { + t.Errorf("scan summary = %+v, want due=22 selected=20 checked=18 updates=17 attention=1 dispatched=15 dispatch_failed=1 busy=1 failed=2 backlog=3", summary) + } + if len(summary.ProviderBackoffs) != 1 || + summary.ProviderBackoffs[0].SourceID != fixtures[1].source.ID || + summary.ProviderBackoffs[0].StatusCode != 429 || + summary.ProviderBackoffs[0].RetryAt != retryAt.UTC().Format(time.RFC3339) { + t.Errorf("scan provider backoffs = %+v, want source=%d status=429 retry_at=%s", summary.ProviderBackoffs, fixtures[1].source.ID, retryAt.UTC().Format(time.RFC3339)) + } + + if len(calledRepositories) != 19 { + t.Fatalf("Resolve calls = %d, want 19 (one busy source in selected batch)", len(calledRepositories)) + } + for index, repository := range calledRepositories { + want := fixtures[index+1].repository + if repository != want { + t.Fatalf("Resolve order[%d] = %q, want %q", index, repository, want) + } + } + + if !containsDispatchedSource(dispatched, fixtures[2].source.ID, stored304Revision) { + t.Errorf("304 stored revision was not dispatched: %+v", dispatched) + } + if containsDispatchedSource(dispatched, fixtures[3].source.ID, "") { + t.Errorf("auto=false source was dispatched: %+v", dispatched) + } + if containsDispatchedSource(dispatched, fixtures[4].source.ID, "") { + t.Errorf("attention source was dispatched: %+v", dispatched) + } + + _, dispatchFailedRuntime := mustLoadPagesSource(t, ctx, fixtures[5].source.ProjectID) + if dispatchFailedRuntime.SyncStatus != pagesSourceStatusUpdateAvailable || + dispatchFailedRuntime.LastError != errPagesSourceTaskDispatchFailed || + dispatchFailedRuntime.NextCheckAt == nil || dispatchFailedRuntime.NextCheckAt.Before(retryAt) || + dispatchFailedRuntime.LastSeenRevision == "" { + t.Errorf( + "dispatch failure runtime = status:%q error:%q next:%v seen:%q, want preserved update and provider deadline >= %v", + dispatchFailedRuntime.SyncStatus, + dispatchFailedRuntime.LastError, + dispatchFailedRuntime.NextCheckAt, + dispatchFailedRuntime.LastSeenRevision, + retryAt, + ) + } +} + +func TestPagesSourceScannerIncludesOrphanCleanupSummary(t *testing.T) { + ctx := setupPagesSourceTest(t) + now := time.Now().Truncate(time.Second) + usePagesSourceScannerClock(t, now) + + previousReconcile := reconcilePagesSourceOrphans + reconcilePagesSourceOrphans = func( + _ context.Context, + gotNow time.Time, + ) (PagesOrphanCleanupSummary, error) { + if !gotNow.Equal(now) { + t.Errorf("orphan cleanup now = %v, want %v", gotNow, now) + } + return PagesOrphanCleanupSummary{ + Candidates: 7, + Reconciled: 1, + Referenced: 2, + LeaseBusy: 1, + InvalidMarker: 1, + Skipped: 1, + Failed: 1, + }, nil + } + t.Cleanup(func() { reconcilePagesSourceOrphans = previousReconcile }) + + result, err := (&SourceScanHandler{}).Execute(ctx, []byte("{}")) + if err != nil { + t.Fatalf("SourceScanHandler.Execute() error = %v, want nil", err) + } + var summary pagesSourceScanSummary + if err := json.Unmarshal([]byte(result.Detail), &summary); err != nil { + t.Fatalf("json.Unmarshal(scan detail) error = %v, want nil", err) + } + if summary.OrphanCleanup.Candidates != 7 || summary.OrphanCleanup.Reconciled != 1 || + summary.OrphanCleanup.Referenced != 2 || summary.OrphanCleanup.LeaseBusy != 1 || + summary.OrphanCleanup.InvalidMarker != 1 || summary.OrphanCleanup.Skipped != 1 || + summary.OrphanCleanup.Failed != 1 { + t.Errorf("orphan cleanup summary = %+v, want injected result", summary.OrphanCleanup) + } +} + +func TestPagesSourceScanPayloadAndMetaAreInternalOnly(t *testing.T) { + handler := &SourceScanHandler{} + if normalized, err := handler.ValidatePayload(nil); err != nil || string(normalized) != "{}" { + t.Errorf("ValidatePayload(nil) = %s, %v; want {}, nil", normalized, err) + } + if _, err := handler.ValidatePayload([]byte(`{"unexpected":true}`)); err == nil { + t.Error("ValidatePayload(unknown field) error = nil, want non-nil") + } + if !PagesSourceScanMeta.InternalOnly || PagesSourceScanMeta.Type != TaskTypePagesSourceScan || + PagesSourceScanMeta.AsynqTask != PagesSourceScanTask || PagesSourceScanMeta.MaxRetry != 0 { + t.Errorf("PagesSourceScanMeta = %+v, want internal bounded scheduled scanner", PagesSourceScanMeta) + } + if PagesSourceScanMeta.SupportsTime { + t.Error("PagesSourceScanMeta.SupportsTime = true, want empty scanner payload") + } +} + +func usePagesSourceScannerClock(t *testing.T, now time.Time) { + t.Helper() + previous := pagesSourceScanNow + pagesSourceScanNow = func() time.Time { return now } + t.Cleanup(func() { pagesSourceScanNow = previous }) +} + +func containsDispatchedSource(dispatched []scannerDispatchedSync, sourceID uint, revision string) bool { + for _, item := range dispatched { + if item.SourceID == sourceID && (revision == "" || item.Revision == revision) { + return true + } + } + return false +} diff --git a/internal/apps/openflare/pages/source_sync.go b/internal/apps/openflare/pages/source_sync.go index 4ef52431..972305dd 100644 --- a/internal/apps/openflare/pages/source_sync.go +++ b/internal/apps/openflare/pages/source_sync.go @@ -25,10 +25,11 @@ import ( ) const ( - pagesSourceTriggerManualSync = "manual_sync" - pagesSourceCreatedBySystem = "system:pages-source-sync" - pagesSourceHeartbeatInterval = pagesSourceSyncLeaseDuration / 3 - pagesSourceCleanupTimeout = 15 * time.Second + pagesSourceTriggerManualSync = "manual_sync" + pagesSourceTriggerScheduledAutoUpdate = "scheduled_auto_update" + pagesSourceCreatedBySystem = "system:pages-source-sync" + pagesSourceHeartbeatInterval = pagesSourceSyncLeaseDuration / 3 + pagesSourceCleanupTimeout = 15 * time.Second ) var ( @@ -154,16 +155,17 @@ func sourceCleanupContext(ctx context.Context) (context.Context, context.CancelF return context.WithTimeout(context.WithoutCancel(ctx), pagesSourceCleanupTimeout) } -func syncRemoteSource( +func syncRemoteSourceWithTrigger( ctx context.Context, snapshot *sourceExecutionSnapshot, actor string, + triggerType string, ) (outcome *sourceSyncOutcome, resultErr error) { if snapshot == nil || snapshot.SourceType != PagesSourceTypeRemoteURL { return nil, errors.New(errPagesSourceTypeUnsupported) } actor = strings.TrimSpace(actor) - if actor == "" { + if actor == "" || !validSourceDeploymentTrigger(triggerType) { return nil, errors.New(errPagesSourceActionInvalid) } defer func() { @@ -222,7 +224,7 @@ func syncRemoteSource( } task.AppendLog(ctx, "[activate] 正在原子切换生产部署") - deployment, reused, referenced, err := commitSourceDeployment( + deployment, reused, referenced, err := commitSourceDeploymentWithTrigger( ctx, snapshot, prepared.Candidate.Checksum, @@ -230,6 +232,7 @@ func syncRemoteSource( prepared.Detail, prepared.DetailJSON, actor, + triggerType, prepared.Manifest, ingestState.Result, ingestState.HasIngest, @@ -382,7 +385,7 @@ func findSourceDeployment( return &deployment, nil } -func commitSourceDeployment( +func commitSourceDeploymentWithTrigger( ctx context.Context, snapshot *sourceExecutionSnapshot, revision string, @@ -390,6 +393,7 @@ func commitSourceDeployment( detail sourceDetail, detailJSON string, actor string, + triggerType string, manifest *deploymentManifest, ingestResult upload.IngestResult, hasIngest bool, @@ -398,6 +402,9 @@ func commitSourceDeployment( if snapshot == nil || manifest == nil { return nil, false, false, errors.New(errPagesSourceSyncFailed) } + if !validSourceDeploymentTrigger(triggerType) { + return nil, false, false, errors.New(errPagesSourceActionInvalid) + } var committed model.PagesDeployment reused := false ingestReferenced := false @@ -407,7 +414,8 @@ func commitSourceDeployment( return err } target, targetReused, err := resolveSourceDeploymentTx( - tx, state, revision, packageChecksum, detail, detailJSON, actor, manifest, ingestResult, hasIngest, + tx, state, revision, packageChecksum, detail, detailJSON, actor, triggerType, + manifest, ingestResult, hasIngest, ) if err != nil { return err @@ -499,6 +507,7 @@ func resolveSourceDeploymentTx( detail sourceDetail, detailJSON string, actor string, + triggerType string, manifest *deploymentManifest, ingestResult upload.IngestResult, hasIngest bool, @@ -520,7 +529,7 @@ func resolveSourceDeploymentTx( return nil, false, errSourceFinalFence } return createSourceDeploymentTx( - tx, state, revision, packageChecksum, detail, detailJSON, actor, manifest, ingestResult, + tx, state, revision, packageChecksum, detail, detailJSON, actor, triggerType, manifest, ingestResult, ) } @@ -532,6 +541,7 @@ func createSourceDeploymentTx( detail sourceDetail, detailJSON string, actor string, + triggerType string, manifest *deploymentManifest, ingestResult upload.IngestResult, ) (*model.PagesDeployment, bool, error) { @@ -558,7 +568,7 @@ func createSourceDeploymentTx( SourceRevision: &revisionValue, SourceLabel: sourceDetailLabel(detail), SourceMeta: detailJSON, - TriggerType: pagesSourceTriggerManualSync, + TriggerType: triggerType, } result := tx.Clauses(clause.OnConflict{DoNothing: true}).Create(target) if result.Error != nil { @@ -573,6 +583,10 @@ func createSourceDeploymentTx( return target, false, nil } +func validSourceDeploymentTrigger(triggerType string) bool { + return triggerType == pagesSourceTriggerManualSync || triggerType == pagesSourceTriggerScheduledAutoUpdate +} + func reloadSourceDeploymentTx( tx *gorm.DB, projectID uint, @@ -657,7 +671,7 @@ func activateSourceDeploymentTx( state.Source.ReleaseSelector == githubReleaseSelectorLatest { next := nextGitHubCheckAt(finishedAt, state.Source.ID, state.Source.CheckIntervalMinutes) if nextCheckNotBefore != nil && nextCheckNotBefore.After(next) { - next = nextCheckNotBefore.UTC() + next = nextCheckNotBefore.In(finishedAt.Location()) } nextCheckAt = &next } diff --git a/internal/apps/openflare/pages/source_tasks.go b/internal/apps/openflare/pages/source_tasks.go index 5e72b383..9d8fdf65 100644 --- a/internal/apps/openflare/pages/source_tasks.go +++ b/internal/apps/openflare/pages/source_tasks.go @@ -51,6 +51,7 @@ type SourceActionPayload struct { ConfigVersion int `json:"config_version"` Action string `json:"action"` Actor string `json:"actor"` + TriggerType string `json:"trigger_type"` TargetRevision string `json:"target_revision"` ConfirmedRevision string `json:"confirmed_revision"` } @@ -71,22 +72,59 @@ func (h *SourceActionHandler) ValidatePayload(payload []byte) ([]byte, error) { } input.Action = strings.TrimSpace(input.Action) input.Actor = strings.TrimSpace(input.Actor) + input.TriggerType = strings.TrimSpace(input.TriggerType) input.TargetRevision = strings.TrimSpace(input.TargetRevision) input.ConfirmedRevision = strings.TrimSpace(input.ConfirmedRevision) - if input.SourceID == 0 || input.ConfigVersion <= 0 || - (input.Action != sourceActionCheck && input.Action != sourceActionSync) || - !validPagesSourceActor(input.Actor) || - !validOptionalSourceRevision(input.TargetRevision) || - !validOptionalSourceRevision(input.ConfirmedRevision) || - (input.Action == sourceActionCheck && (input.TargetRevision != "" || input.ConfirmedRevision != "")) || - (input.TargetRevision != "" && input.ConfirmedRevision != "") || - (input.TargetRevision != "" && input.Actor != pagesSourceCreatedBySystem) || - (input.ConfirmedRevision != "" && !strings.HasPrefix(input.Actor, "user:")) { + if input.Action == sourceActionSync && input.TriggerType == "" { + // Keep already queued Phase 2 payloads valid while making every new + // dispatch carry an explicit deployment trigger. + input.TriggerType = pagesSourceTriggerManualSync + } + if !validSourceActionPayload(input) { return nil, errors.New(errPagesSourceActionInvalid) } return json.Marshal(input) } +func validSourceActionPayload(input SourceActionPayload) bool { + if input.SourceID == 0 || input.ConfigVersion <= 0 { + return false + } + if input.Action != sourceActionCheck && input.Action != sourceActionSync { + return false + } + if !validPagesSourceActor(input.Actor) { + return false + } + if !validOptionalSourceRevision(input.TargetRevision) || + !validOptionalSourceRevision(input.ConfirmedRevision) { + return false + } + if input.Action == sourceActionCheck { + return input.TriggerType == "" && input.TargetRevision == "" && input.ConfirmedRevision == "" + } + return validSourceSyncPayload(input) +} + +func validSourceSyncPayload(input SourceActionPayload) bool { + if !validSourceDeploymentTrigger(input.TriggerType) || + (input.TargetRevision != "" && input.ConfirmedRevision != "") { + return false + } + switch input.TriggerType { + case pagesSourceTriggerScheduledAutoUpdate: + return input.Actor == pagesSourceCreatedBySystem && + input.TargetRevision != "" && input.ConfirmedRevision == "" + case pagesSourceTriggerManualSync: + if input.TargetRevision != "" || !strings.HasPrefix(input.Actor, "user:") { + return false + } + return true + default: + return false + } +} + // Execute validates again inside the worker and performs the source action. func (h *SourceActionHandler) Execute(ctx context.Context, payload []byte) (*task.TaskResult, error) { normalized, err := h.ValidatePayload(payload) @@ -169,9 +207,11 @@ func executeSourceSyncAction( var result *sourceSyncOutcome var err error if source.SourceType == PagesSourceTypeGitHubRelease { - result, err = syncGitHubSource(ctx, snapshot, input.Actor, input.TargetRevision, input.ConfirmedRevision) + result, err = syncGitHubSourceWithTrigger( + ctx, snapshot, input.Actor, input.TargetRevision, input.ConfirmedRevision, input.TriggerType, + ) } else { - result, err = syncRemoteSource(ctx, snapshot, input.Actor) + result, err = syncRemoteSourceWithTrigger(ctx, snapshot, input.Actor, input.TriggerType) } if err != nil { logger.ErrorF(ctx, "[PagesSource] sync failed: project_id=%d source_id=%d error=%v", snapshot.ProjectID, snapshot.SourceID, err) @@ -333,6 +373,25 @@ func dispatchSourceActionSnapshot( targetRevision string, confirmedRevision string, triggeredBy string, +) (*SourceActionReceipt, error) { + triggerType := "" + if action == sourceActionSync { + triggerType = pagesSourceTriggerManualSync + } + return dispatchSourceActionSnapshotWithTrigger( + ctx, source, action, actor, triggerType, targetRevision, confirmedRevision, triggeredBy, + ) +} + +func dispatchSourceActionSnapshotWithTrigger( + ctx context.Context, + source model.PagesProjectSource, + action string, + actor string, + triggerType string, + targetRevision string, + confirmedRevision string, + triggeredBy string, ) (*SourceActionReceipt, error) { if task.AsynqClient == nil { return nil, errors.New(errPagesSourceTaskDispatchFailed) @@ -343,6 +402,7 @@ func dispatchSourceActionSnapshot( ConfigVersion: source.ConfigVersion, Action: action, Actor: actor, + TriggerType: triggerType, TargetRevision: targetRevision, ConfirmedRevision: confirmedRevision, }) diff --git a/internal/apps/openflare/pages/source_tasks_test.go b/internal/apps/openflare/pages/source_tasks_test.go index 92458dc6..d51b9af7 100644 --- a/internal/apps/openflare/pages/source_tasks_test.go +++ b/internal/apps/openflare/pages/source_tasks_test.go @@ -19,6 +19,7 @@ func TestSourceActionPayloadValidationIsStrictAndCredentialFree(t *testing.T) { ConfigVersion: 3, Action: sourceActionSync, Actor: "user:42", + TriggerType: pagesSourceTriggerManualSync, } raw, err := json.Marshal(valid) if err != nil { @@ -35,6 +36,23 @@ func TestSourceActionPayloadValidationIsStrictAndCredentialFree(t *testing.T) { if got != valid { t.Errorf("ValidatePayload(valid) = %+v, want %+v", got, valid) } + legacy := valid + legacy.TriggerType = "" + legacyRaw, err := json.Marshal(legacy) + if err != nil { + t.Fatalf("json.Marshal(legacy payload) error = %v, want nil", err) + } + legacyNormalized, err := handler.ValidatePayload(legacyRaw) + if err != nil { + t.Fatalf("ValidatePayload(legacy payload) error = %v, want nil", err) + } + var legacyGot SourceActionPayload + if err := json.Unmarshal(legacyNormalized, &legacyGot); err != nil { + t.Fatalf("json.Unmarshal(legacy normalized payload) error = %v, want nil", err) + } + if legacyGot.TriggerType != pagesSourceTriggerManualSync { + t.Errorf("legacy payload trigger_type = %q, want %q", legacyGot.TriggerType, pagesSourceTriggerManualSync) + } for _, forbidden := range []string{"remote_url", "content_config_version", "expected_revision", "lease_token", "etag"} { if strings.Contains(string(normalized), forbidden) { t.Errorf("normalized payload = %s, want no forbidden field %q", normalized, forbidden) diff --git a/internal/db/migrator/goose/postgres/202607190002_seed_pages_source_scan.sql b/internal/db/migrator/goose/postgres/202607190002_seed_pages_source_scan.sql new file mode 100644 index 00000000..d50e45b3 --- /dev/null +++ b/internal/db/migrator/goose/postgres/202607190002_seed_pages_source_scan.sql @@ -0,0 +1,38 @@ +-- +goose Up +-- Earlier built-in schedules used explicit IDs, so advance the identity only +-- when it trails either existing rows or an already-higher sequence value. +SELECT setval( + pg_get_serial_sequence('w_schedules', 'id'), + GREATEST( + 1, + COALESCE((SELECT MAX(id) FROM w_schedules), 0), + COALESCE(( + SELECT sequences.last_value + FROM pg_sequences AS sequences + WHERE format('%I.%I', sequences.schemaname, sequences.sequencename)::regclass = + pg_get_serial_sequence('w_schedules', 'id')::regclass + ), 0) + ), + TRUE +); + +INSERT INTO w_schedules (name, task_type, cron, payload, is_active, created_at, updated_at) +SELECT + 'OpenFlare Pages 部署源扫描', + 'of_pages_source_scan', + '*/5 * * * *', + '{}', + TRUE, + CURRENT_TIMESTAMP, + CURRENT_TIMESTAMP +WHERE NOT EXISTS ( + SELECT 1 FROM w_schedules WHERE task_type = 'of_pages_source_scan' +); + +-- +goose Down +DELETE FROM w_schedules +WHERE task_type = 'of_pages_source_scan' + AND name = 'OpenFlare Pages 部署源扫描' + AND cron = '*/5 * * * *' + AND payload = '{}' + AND is_active = TRUE; diff --git a/internal/db/migrator/goose/sqlite/202607190002_seed_pages_source_scan.sql b/internal/db/migrator/goose/sqlite/202607190002_seed_pages_source_scan.sql new file mode 100644 index 00000000..dd0235bf --- /dev/null +++ b/internal/db/migrator/goose/sqlite/202607190002_seed_pages_source_scan.sql @@ -0,0 +1,21 @@ +-- +goose Up +INSERT INTO w_schedules (name, task_type, cron, payload, is_active, created_at, updated_at) +SELECT + 'OpenFlare Pages 部署源扫描', + 'of_pages_source_scan', + '*/5 * * * *', + '{}', + 1, + CURRENT_TIMESTAMP, + CURRENT_TIMESTAMP +WHERE NOT EXISTS ( + SELECT 1 FROM w_schedules WHERE task_type = 'of_pages_source_scan' +); + +-- +goose Down +DELETE FROM w_schedules +WHERE task_type = 'of_pages_source_scan' + AND name = 'OpenFlare Pages 部署源扫描' + AND cron = '*/5 * * * *' + AND payload = '{}' + AND is_active = 1; diff --git a/internal/db/migrator/pages_source_scan_migration_test.go b/internal/db/migrator/pages_source_scan_migration_test.go new file mode 100644 index 00000000..3ee01b33 --- /dev/null +++ b/internal/db/migrator/pages_source_scan_migration_test.go @@ -0,0 +1,160 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package migrator + +import ( + "database/sql" + "fmt" + "os" + "strings" + "testing" + "time" + + "github.com/Rain-kl/Wavelet/internal/model" + "github.com/glebarez/sqlite" + "github.com/pressly/goose/v3" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "gorm.io/driver/postgres" + "gorm.io/gorm" +) + +const ( + pagesSourceScanPreviousMigration = int64(202607190001) + pagesSourceScanMigration = int64(202607190002) + pagesSourceScanTaskType = "of_pages_source_scan" +) + +func TestPagesSourceScanScheduleMigrationSQLite(t *testing.T) { + dbPath := t.TempDir() + "/pages-source-scan-migration.db" + gormDB, err := gorm.Open(sqlite.Open(dbPath), &gorm.Config{ + DisableForeignKeyConstraintWhenMigrating: true, + }) + require.NoError(t, err) + sqlDB, err := gormDB.DB() + require.NoError(t, err) + sqlDB.SetMaxOpenConns(1) + t.Cleanup(func() { require.NoError(t, sqlDB.Close()) }) + + runPagesSourceScanScheduleMigration(t, gormDB, sqlDB, dialectSqlite, "goose/sqlite") +} + +func TestPagesSourceScanScheduleMigrationPostgres(t *testing.T) { + dsn := strings.TrimSpace(os.Getenv("OPENFLARE_TEST_POSTGRES_DSN")) + if dsn == "" { + t.Skip("OPENFLARE_TEST_POSTGRES_DSN is not set") + } + gormDB, err := gorm.Open(postgres.Open(dsn), &gorm.Config{ + DisableForeignKeyConstraintWhenMigrating: true, + }) + require.NoError(t, err) + sqlDB, err := gormDB.DB() + require.NoError(t, err) + sqlDB.SetMaxOpenConns(1) + + schema := fmt.Sprintf("pages_source_scan_migration_%d", time.Now().UnixNano()) + require.Regexp(t, `^[a-z0-9_]+$`, schema) + require.NoError(t, gormDB.Exec(`CREATE SCHEMA "`+schema+`"`).Error) + require.NoError(t, gormDB.Exec(`SET search_path TO "`+schema+`"`).Error) + t.Cleanup(func() { + assert.NoError(t, gormDB.Exec("SET search_path TO public").Error) + assert.NoError(t, gormDB.Exec(`DROP SCHEMA IF EXISTS "`+schema+`" CASCADE`).Error) + assert.NoError(t, sqlDB.Close()) + }) + + runPagesSourceScanScheduleMigration(t, gormDB, sqlDB, dialectPostgres, "goose/postgres") +} + +func runPagesSourceScanScheduleMigration( + t *testing.T, + gormDB *gorm.DB, + sqlDB *sql.DB, + dialect string, + dir string, +) { + t.Helper() + goose.SetBaseFS(migrationFS) + require.NoError(t, goose.SetDialect(dialect)) + require.NoError(t, goose.UpTo(sqlDB, dir, pagesSourceScanPreviousMigration)) + + var previousMaxID uint64 + require.NoError(t, gormDB.Table("w_schedules").Select("COALESCE(MAX(id), 0)").Scan(&previousMaxID).Error) + require.NoError(t, goose.UpTo(sqlDB, dir, pagesSourceScanMigration)) + seeded := assertPagesSourceScanSchedule(t, gormDB) + assert.NotZero(t, seeded.ID) + if dialect == dialectPostgres { + assert.Greater(t, seeded.ID, previousMaxID) + } + + require.NoError(t, goose.DownTo(sqlDB, dir, pagesSourceScanPreviousMigration)) + assertPagesSourceScanScheduleMissing(t, gormDB) + + custom := model.Schedule{ + ID: 900001, + Name: "用户保留的 Pages 扫描任务", + TaskType: pagesSourceScanTaskType, + Cron: "0 * * * *", + Payload: `{"custom":true}`, + IsActive: false, + } + require.NoError(t, gormDB.Create(&custom).Error) + require.NoError(t, goose.UpTo(sqlDB, dir, pagesSourceScanMigration)) + var schedules []model.Schedule + require.NoError(t, gormDB.Where("task_type = ?", pagesSourceScanTaskType).Find(&schedules).Error) + require.Len(t, schedules, 1) + assert.Equal(t, custom.ID, schedules[0].ID) + assert.Equal(t, custom.Name, schedules[0].Name) + + require.NoError(t, goose.DownTo(sqlDB, dir, pagesSourceScanPreviousMigration)) + var retained model.Schedule + require.NoError(t, gormDB.First(&retained, custom.ID).Error) + assert.Equal(t, custom.TaskType, retained.TaskType) +} + +func TestPagesSourceScanScheduleMigrationsUseDatabaseGeneratedIDs(t *testing.T) { + for _, name := range []string{ + "goose/postgres/202607190002_seed_pages_source_scan.sql", + "goose/sqlite/202607190002_seed_pages_source_scan.sql", + } { + t.Run(name, func(t *testing.T) { + content, err := migrationFS.ReadFile(name) + require.NoError(t, err) + normalized := strings.ToLower(string(content)) + assert.NotContains(t, normalized, "insert into w_schedules (id,") + assert.NotContains(t, normalized, "coalesce(max(id)") + }) + } + + postgresContent, err := migrationFS.ReadFile("goose/postgres/202607190002_seed_pages_source_scan.sql") + require.NoError(t, err) + compactPostgres := strings.Join(strings.Fields(strings.ToLower(string(postgresContent))), " ") + assert.Contains( + t, + compactPostgres, + "select setval( pg_get_serial_sequence('w_schedules', 'id'), greatest( 1,", + "sequence synchronization must retain a valid lower bound for an empty table", + ) +} + +func assertPagesSourceScanSchedule(t *testing.T, gormDB *gorm.DB) model.Schedule { + t.Helper() + var schedules []model.Schedule + require.NoError(t, gormDB.Where("task_type = ?", pagesSourceScanTaskType).Find(&schedules).Error) + require.Len(t, schedules, 1) + schedule := schedules[0] + assert.Equal(t, "OpenFlare Pages 部署源扫描", schedule.Name) + assert.Equal(t, "*/5 * * * *", schedule.Cron) + assert.Equal(t, "{}", schedule.Payload) + assert.True(t, schedule.IsActive) + return schedule +} + +func assertPagesSourceScanScheduleMissing(t *testing.T, gormDB *gorm.DB) { + t.Helper() + var count int64 + require.NoError(t, gormDB.Model(&model.Schedule{}). + Where("task_type = ?", pagesSourceScanTaskType). + Count(&count).Error) + assert.Zero(t, count) +} diff --git a/internal/task/handlers/register.go b/internal/task/handlers/register.go index 641d9dca..005bfd9e 100644 --- a/internal/task/handlers/register.go +++ b/internal/task/handlers/register.go @@ -53,6 +53,9 @@ func Register() { task.RegisterTaskMeta(openflare.UptimeKumaSyncMeta) // pages source actions are only dispatched by the Pages domain API/scanner. + task.RegisterHandler(pages.PagesSourceScanTask, &pages.SourceScanHandler{}) + task.RegisterTaskMeta(pages.PagesSourceScanMeta) + task.RegisterHandler(pages.PagesSourceActionTask, &pages.SourceActionHandler{}) task.RegisterTaskMeta(pages.PagesSourceActionMeta)