feat(pages): 增加来源扫描与自动更新

为 GitHub latest 来源增加五分钟 scanner、按来源间隔检查、精确 revision 自动发布与租约恢复。\n记录退避和投递统计,并为 PostgreSQL 与 SQLite 幂等创建内部排程。
This commit is contained in:
deqiying
2026-07-19 19:09:21 +08:00
parent 848884d8cd
commit 999428cf9a
16 changed files with 1493 additions and 65 deletions
+2 -2
View File
@@ -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
+1
View File
@@ -61,6 +61,7 @@ const (
errPagesSourceActionInvalid = "pages 部署源任务参数无效"
errPagesSourceActionStale = "pages 部署源配置已变化,本次任务已跳过"
errPagesSourceLeaseLost = "pages 部署源任务执行权已失效"
errPagesSourceLeaseExpired = "上次 pages 部署源任务租约已过期"
errPagesSourceSyncFailed = "pages 部署源同步失败"
errPagesSourceTaskDispatchFailed = "pages 部署源任务入队失败"
errPagesSourceInternal = "pages 部署源操作失败,请稍后重试"
+11 -7
View File
@@ -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
}
@@ -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,
}
}
@@ -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)
@@ -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,
)
}
@@ -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 {
@@ -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
}
@@ -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
}
+26 -12
View File
@@ -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
}
+71 -11
View File
@@ -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,
})
@@ -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)
@@ -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;
@@ -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;
@@ -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)
}
+3
View File
@@ -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)