diff --git a/docs/docs.go b/docs/docs.go index 2f5ca21a..d37cf169 100644 --- a/docs/docs.go +++ b/docs/docs.go @@ -2595,7 +2595,7 @@ const docTemplate = `{ "SessionCookie": [] } ], - "description": "返回系统所有的定时任务配置列表,包括名称、关联的异步任务类型、Cron 表达式和启用状态,需要管理员权限", + "description": "返回管理员可管理的定时任务配置列表,包括名称、关联的异步任务类型、Cron 表达式和启用状态;系统内部排程不会暴露,需要管理员权限", "produces": [ "application/json" ], @@ -2860,6 +2860,12 @@ const docTemplate = `{ "$ref": "#/definitions/response.Any" } }, + "404": { + "description": "定时任务不存在", + "schema": { + "$ref": "#/definitions/response.Any" + } + }, "500": { "description": "删除定时任务失败", "schema": { diff --git a/docs/swagger.json b/docs/swagger.json index cc0f02c5..e02e9b98 100644 --- a/docs/swagger.json +++ b/docs/swagger.json @@ -2588,7 +2588,7 @@ "SessionCookie": [] } ], - "description": "返回系统所有的定时任务配置列表,包括名称、关联的异步任务类型、Cron 表达式和启用状态,需要管理员权限", + "description": "返回管理员可管理的定时任务配置列表,包括名称、关联的异步任务类型、Cron 表达式和启用状态;系统内部排程不会暴露,需要管理员权限", "produces": [ "application/json" ], @@ -2853,6 +2853,12 @@ "$ref": "#/definitions/response.Any" } }, + "404": { + "description": "定时任务不存在", + "schema": { + "$ref": "#/definitions/response.Any" + } + }, "500": { "description": "删除定时任务失败", "schema": { diff --git a/docs/swagger.yaml b/docs/swagger.yaml index edce34e0..3c9e9cdf 100644 --- a/docs/swagger.yaml +++ b/docs/swagger.yaml @@ -5717,7 +5717,7 @@ paths: - admin /api/v1/admin/tasks/schedules: get: - description: 返回系统所有的定时任务配置列表,包括名称、关联的异步任务类型、Cron 表达式和启用状态,需要管理员权限 + description: 返回管理员可管理的定时任务配置列表,包括名称、关联的异步任务类型、Cron 表达式和启用状态;系统内部排程不会暴露,需要管理员权限 produces: - application/json responses: @@ -5822,6 +5822,10 @@ paths: description: 无管理员权限 schema: $ref: '#/definitions/response.Any' + "404": + description: 定时任务不存在 + schema: + $ref: '#/definitions/response.Any' "500": description: 删除定时任务失败 schema: diff --git a/internal/apps/admin/task/routers.go b/internal/apps/admin/task/routers.go index 594f1eb2..78d391a5 100644 --- a/internal/apps/admin/task/routers.go +++ b/internal/apps/admin/task/routers.go @@ -202,7 +202,7 @@ func RetryTask(c *gin.Context) { // ListSchedules 获取定时任务列表 // @Summary 获取定时任务列表 -// @Description 返回系统所有的定时任务配置列表,包括名称、关联的异步任务类型、Cron 表达式和启用状态,需要管理员权限 +// @Description 返回管理员可管理的定时任务配置列表,包括名称、关联的异步任务类型、Cron 表达式和启用状态;系统内部排程不会暴露,需要管理员权限 // @Tags admin // @Produce json // @Security SessionCookie @@ -216,7 +216,15 @@ func ListSchedules(c *gin.Context) { response.AbortInternal(c, err.Error()) return } - c.JSON(http.StatusOK, response.OK(schedules)) + visible := make([]model.Schedule, 0, len(schedules)) + for _, schedule := range schedules { + meta := task.GetTaskMeta(schedule.TaskType) + if meta != nil && meta.InternalOnly { + continue + } + visible = append(visible, schedule) + } + c.JSON(http.StatusOK, response.OK(visible)) } // CreateScheduleRequest 创建定时任务请求 @@ -405,6 +413,7 @@ func getAdminTaskMeta(taskType string) *task.TaskMeta { // @Failure 400 {object} response.Any "参数错误" // @Failure 401 {object} response.Any "未登录" // @Failure 403 {object} response.Any "无管理员权限" +// @Failure 404 {object} response.Any "定时任务不存在" // @Failure 500 {object} response.Any "删除定时任务失败" // @Router /api/v1/admin/tasks/schedules/{id} [delete] func DeleteSchedule(c *gin.Context) { @@ -413,6 +422,15 @@ func DeleteSchedule(c *gin.Context) { response.AbortBadRequest(c, "无效的定时任务ID") return } + schedule, err := model.GetScheduleByID(c.Request.Context(), id) + if err != nil { + response.AbortNotFound(c, ScheduleNotFound) + return + } + if meta := task.GetTaskMeta(schedule.TaskType); meta != nil && meta.InternalOnly { + response.AbortBadRequest(c, InvalidTaskType) + return + } if err := model.DeleteSchedule(c.Request.Context(), id); err != nil { response.AbortInternal(c, fmt.Sprintf("%s: %v", ScheduleDeleteFailed, err)) diff --git a/internal/apps/admin/task/routers_test.go b/internal/apps/admin/task/routers_test.go index 694861e7..3cd63fb9 100644 --- a/internal/apps/admin/task/routers_test.go +++ b/internal/apps/admin/task/routers_test.go @@ -75,8 +75,10 @@ func setupTestRouter(authUser *model.User) *gin.Engine { adminGroup.GET("/tasks/executions", ListTaskExecutions) adminGroup.GET("/tasks/executions/:id", GetTaskExecution) adminGroup.POST("/tasks/executions/:id/retry", RetryTask) + adminGroup.GET("/tasks/schedules", ListSchedules) adminGroup.POST("/tasks/schedules", CreateSchedule) adminGroup.PUT("/tasks/schedules/:id", UpdateSchedule) + adminGroup.DELETE("/tasks/schedules/:id", DeleteSchedule) return r } @@ -137,6 +139,39 @@ func TestInternalOnlyTaskAdminBoundaries(t *testing.T) { router := setupTestRouter(adminUser) ctx := context.Background() + t.Run("list hides internal-only schedule", func(t *testing.T) { + internalSchedule := &model.Schedule{ + Name: "隐藏的系统内部排程", + TaskType: testInternalOnlyTaskType, + Cron: "*/5 * * * *", + Payload: "{}", + IsActive: true, + } + publicSchedule := &model.Schedule{ + Name: "可见的公开排程", + TaskType: uploadtask.TaskTypeSystemCleanup, + Cron: "0 * * * *", + Payload: "{}", + IsActive: true, + } + require.NoError(t, model.CreateSchedule(ctx, internalSchedule)) + require.NoError(t, model.CreateSchedule(ctx, publicSchedule)) + + req := httptest.NewRequest(http.MethodGet, "/api/v1/admin/tasks/schedules", nil) + w := httptest.NewRecorder() + router.ServeHTTP(w, req) + + assert.Equal(t, http.StatusOK, w.Code) + var resp response.Any + require.NoError(t, json.Unmarshal(w.Body.Bytes(), &resp)) + data, err := json.Marshal(resp.Data) + require.NoError(t, err) + var schedules []model.Schedule + require.NoError(t, json.Unmarshal(data, &schedules)) + assert.NotContains(t, scheduleIDs(schedules), internalSchedule.ID) + assert.Contains(t, scheduleIDs(schedules), publicSchedule.ID) + }) + t.Run("dispatch rejects internal-only task", func(t *testing.T) { body, err := json.Marshal(DispatchTaskRequest{TaskType: testInternalOnlyTaskType}) require.NoError(t, err) @@ -239,6 +274,74 @@ func TestInternalOnlyTaskAdminBoundaries(t *testing.T) { assert.Equal(t, "公开排程", unchanged.Name) assert.Equal(t, uploadtask.TaskTypeSystemCleanup, unchanged.TaskType) }) + + t.Run("delete rejects internal-only schedule", func(t *testing.T) { + schedule := &model.Schedule{ + Name: "不可删除的系统内部排程", + TaskType: testInternalOnlyTaskType, + Cron: "*/5 * * * *", + IsActive: true, + } + require.NoError(t, model.CreateSchedule(ctx, schedule)) + req := httptest.NewRequest( + http.MethodDelete, + fmt.Sprintf("/api/v1/admin/tasks/schedules/%d", schedule.ID), + nil, + ) + w := httptest.NewRecorder() + + router.ServeHTTP(w, req) + + assert.Equal(t, http.StatusBadRequest, w.Code) + var resp response.Any + require.NoError(t, json.Unmarshal(w.Body.Bytes(), &resp)) + assert.Equal(t, InvalidTaskType, resp.ErrorMsg) + preserved, err := model.GetScheduleByID(ctx, schedule.ID) + require.NoError(t, err) + assert.Equal(t, testInternalOnlyTaskType, preserved.TaskType) + }) + + t.Run("delete missing schedule returns not found", func(t *testing.T) { + req := httptest.NewRequest(http.MethodDelete, "/api/v1/admin/tasks/schedules/999999", nil) + w := httptest.NewRecorder() + + router.ServeHTTP(w, req) + + assert.Equal(t, http.StatusNotFound, w.Code) + var resp response.Any + require.NoError(t, json.Unmarshal(w.Body.Bytes(), &resp)) + assert.Equal(t, ScheduleNotFound, resp.ErrorMsg) + }) + + t.Run("delete public schedule remains allowed", func(t *testing.T) { + schedule := &model.Schedule{ + Name: "可删除的公开排程", + TaskType: uploadtask.TaskTypeSystemCleanup, + Cron: "0 * * * *", + IsActive: false, + } + require.NoError(t, model.CreateSchedule(ctx, schedule)) + req := httptest.NewRequest( + http.MethodDelete, + fmt.Sprintf("/api/v1/admin/tasks/schedules/%d", schedule.ID), + nil, + ) + w := httptest.NewRecorder() + + router.ServeHTTP(w, req) + + assert.Equal(t, http.StatusOK, w.Code) + _, err := model.GetScheduleByID(ctx, schedule.ID) + assert.Error(t, err) + }) +} + +func scheduleIDs(schedules []model.Schedule) []uint64 { + ids := make([]uint64, 0, len(schedules)) + for _, schedule := range schedules { + ids = append(ids, schedule.ID) + } + return ids } func TestDispatchTask(t *testing.T) { diff --git a/internal/apps/openflare/pages/source_orphan_cleanup.go b/internal/apps/openflare/pages/source_orphan_cleanup.go new file mode 100644 index 00000000..6f80380f --- /dev/null +++ b/internal/apps/openflare/pages/source_orphan_cleanup.go @@ -0,0 +1,328 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package pages + +import ( + "context" + "errors" + "fmt" + "strconv" + "time" + + "github.com/Rain-kl/Wavelet/internal/apps/upload" + "github.com/Rain-kl/Wavelet/internal/db" + "github.com/Rain-kl/Wavelet/internal/model" + "github.com/Rain-kl/Wavelet/internal/repository" + "github.com/Rain-kl/Wavelet/pkg/logger" + "gorm.io/gorm" + "gorm.io/gorm/clause" +) + +const pagesOrphanUploadIsolation = 2 * time.Hour + +// PagesOrphanCleanupSummary describes one bounded delayed compensation pass. +// Every candidate is counted exactly once in one outcome field. +// +//nolint:revive // Keep the domain-qualified exported name for scanner/task result clarity. +type PagesOrphanCleanupSummary struct { + Candidates int `json:"candidates"` + Reconciled int `json:"reconciled"` + Referenced int `json:"referenced"` + LeaseBusy int `json:"lease_busy"` + InvalidMarker int `json:"invalid_marker"` + Skipped int `json:"skipped"` + Failed int `json:"failed"` +} + +type pagesOrphanMarker struct { + ProjectID uint + SourceID *uint +} + +type pagesOrphanCleanupOutcome uint8 + +const ( + pagesOrphanCleanupSkipped pagesOrphanCleanupOutcome = iota + pagesOrphanCleanupReconciled + pagesOrphanCleanupReferenced + pagesOrphanCleanupLeaseBusy + pagesOrphanCleanupInvalidMarker +) + +// ReconcilePagesOrphanUploads performs one bounded delayed compensation pass. +// Individual candidate failures are counted and logged so they do not prevent +// the scanner from continuing with source checks. +func ReconcilePagesOrphanUploads( + ctx context.Context, + now time.Time, +) (PagesOrphanCleanupSummary, error) { + if now.IsZero() { + now = time.Now() + } + cutoff := now.UTC().Add(-pagesOrphanUploadIsolation) + systemUser := repository.GetSystemUser(ctx) + candidates, err := model.ListPagesOrphanUploadCandidates(ctx, model.PagesOrphanUploadCandidateQuery{ + SystemUserID: systemUser.ID, + UploadType: upload.ReservedPagesDeploymentType, + Marker: pagesIngestMarkerV2, + CreatedBefore: cutoff, + }) + if err != nil { + return PagesOrphanCleanupSummary{}, err + } + + summary := PagesOrphanCleanupSummary{Candidates: len(candidates)} + for index := range candidates { + if err := ctx.Err(); err != nil { + return summary, err + } + candidate := &candidates[index] + marker, err := parsePagesOrphanMarker(candidate.Metadata) + if err != nil { + summary.InvalidMarker++ + logger.WarnF(ctx, "[PagesSource] orphan upload marker invalid: upload_id=%d error=%v", candidate.ID, err) + continue + } + + outcome, err := reconcilePagesOrphanUploadCandidate( + ctx, + candidate, + marker, + systemUser.ID, + cutoff, + ) + if err != nil { + summary.Failed++ + logger.WarnF(ctx, "[PagesSource] orphan upload reconciliation failed: upload_id=%d error=%v", candidate.ID, err) + continue + } + summary.add(outcome) + } + return summary, nil +} + +func (summary *PagesOrphanCleanupSummary) add(outcome pagesOrphanCleanupOutcome) { + switch outcome { + case pagesOrphanCleanupReconciled: + summary.Reconciled++ + case pagesOrphanCleanupReferenced: + summary.Referenced++ + case pagesOrphanCleanupLeaseBusy: + summary.LeaseBusy++ + case pagesOrphanCleanupInvalidMarker: + summary.InvalidMarker++ + default: + summary.Skipped++ + } +} + +func reconcilePagesOrphanUploadCandidate( + ctx context.Context, + candidate *model.Upload, + marker pagesOrphanMarker, + systemUserID uint64, + cutoff time.Time, +) (pagesOrphanCleanupOutcome, error) { + if candidate == nil || candidate.ID == 0 { + return pagesOrphanCleanupSkipped, nil + } + + outcome := pagesOrphanCleanupSkipped + uploadLocked := false + err := db.DB(ctx).Transaction(func(tx *gorm.DB) error { + scopeOutcome, proceed, err := lockPagesOrphanCleanupScope(ctx, tx, candidate.ID, marker) + if err != nil { + return err + } + if !proceed { + outcome = scopeOutcome + return nil + } + lockedOutcome, locked, err := reconcileLockedPagesOrphanUpload( + ctx, + tx, + candidate.ID, + marker, + systemUserID, + cutoff, + ) + if err != nil { + return err + } + outcome = lockedOutcome + uploadLocked = locked + return nil + }) + if err != nil { + return pagesOrphanCleanupSkipped, err + } + if uploadLocked { + // Also heal a prior post-commit cache invalidation interruption when the + // status transition was an idempotent no-op. + upload.InvalidateUploadMetaCache(ctx, candidate.ID) + } + return outcome, nil +} + +func lockPagesOrphanCleanupScope( + ctx context.Context, + tx *gorm.DB, + uploadID uint64, + marker pagesOrphanMarker, +) (pagesOrphanCleanupOutcome, bool, error) { + var project model.PagesProject + if _, err := lockOptionalPagesCleanupRecord(tx, &project, "id = ?", marker.ProjectID); err != nil { + return pagesOrphanCleanupSkipped, false, err + } + if marker.SourceID == nil { + return pagesOrphanCleanupSkipped, true, nil + } + + var source model.PagesProjectSource + sourceExists, err := lockOptionalPagesCleanupRecord(tx, &source, "id = ?", *marker.SourceID) + if err != nil { + return pagesOrphanCleanupSkipped, false, err + } + if !sourceExists { + return pagesOrphanCleanupSkipped, true, nil + } + if source.ProjectID != marker.ProjectID { + logger.WarnF(ctx, + "[PagesSource] orphan upload source ownership mismatch: upload_id=%d project_id=%d source_id=%d source_project_id=%d", + uploadID, + marker.ProjectID, + *marker.SourceID, + source.ProjectID, + ) + return pagesOrphanCleanupInvalidMarker, false, nil + } + + var runtime model.PagesProjectSourceRuntime + runtimeExists, err := lockOptionalPagesCleanupRecord(tx, &runtime, "source_id = ?", source.ID) + if err != nil { + return pagesOrphanCleanupSkipped, false, err + } + // Read the real clock only after obtaining the runtime row lock. The scanner + // snapshot time is only an isolation cutoff and may be stale after lock wait. + leaseCheckedAt := time.Now() + if runtimeExists && runtime.LeaseExpiresAt != nil && runtime.LeaseExpiresAt.After(leaseCheckedAt) { + return pagesOrphanCleanupLeaseBusy, false, nil + } + return pagesOrphanCleanupSkipped, true, nil +} + +func reconcileLockedPagesOrphanUpload( + ctx context.Context, + tx *gorm.DB, + uploadID uint64, + marker pagesOrphanMarker, + systemUserID uint64, + cutoff time.Time, +) (pagesOrphanCleanupOutcome, bool, error) { + var lockedUpload model.Upload + found, err := lockOptionalPagesCleanupRecord(tx, &lockedUpload, "id = ?", uploadID) + if err != nil || !found { + return pagesOrphanCleanupSkipped, false, err + } + + lockedMarker, err := parsePagesOrphanMarker(lockedUpload.Metadata) + if err != nil { + logger.WarnF(ctx, "[PagesSource] orphan upload marker changed or invalid: upload_id=%d error=%v", uploadID, err) + return pagesOrphanCleanupInvalidMarker, true, nil + } + if lockedUpload.Status != model.UploadStatusUsed || + lockedUpload.UserID != systemUserID || + lockedUpload.Type != upload.ReservedPagesDeploymentType || + !lockedUpload.CreatedAt.Before(cutoff) { + return pagesOrphanCleanupSkipped, true, nil + } + if lockedMarker.ProjectID != marker.ProjectID || !sameOptionalPagesSourceID(lockedMarker.SourceID, marker.SourceID) { + logger.WarnF(ctx, "[PagesSource] orphan upload marker changed during reconciliation: upload_id=%d", uploadID) + return pagesOrphanCleanupInvalidMarker, true, nil + } + + var references int64 + if err := tx.Model(&model.PagesDeployment{}). + Where("upload_id = ?", lockedUpload.ID). + Count(&references).Error; err != nil { + return pagesOrphanCleanupSkipped, true, err + } + if references > 0 { + return pagesOrphanCleanupReferenced, true, nil + } + + transitioned, err := upload.RemoveLockedTx(tx, &lockedUpload) + if err != nil { + return pagesOrphanCleanupSkipped, true, err + } + if transitioned { + return pagesOrphanCleanupReconciled, true, nil + } + return pagesOrphanCleanupSkipped, true, nil +} + +func lockOptionalPagesCleanupRecord( + tx *gorm.DB, + value any, + query string, + args ...any, +) (bool, error) { + err := tx.Clauses(clause.Locking{Strength: pagesRowLockStrength}). + Where(query, args...). + First(value).Error + if err == nil { + return true, nil + } + if errors.Is(err, gorm.ErrRecordNotFound) { + return false, nil + } + return false, err +} + +func parsePagesOrphanMarker(metadata model.UploadMetadata) (pagesOrphanMarker, error) { + if metadata.Extra == nil { + return pagesOrphanMarker{}, errors.New("pages marker metadata missing") + } + marker, ok := metadata.Extra[pagesIngestMarkerKey].(string) + if !ok || marker != pagesIngestMarkerV2 { + return pagesOrphanMarker{}, errors.New("pages marker version invalid") + } + projectID, err := parsePagesOrphanMetadataID(metadata.Extra, pagesProjectIDMetadataKey) + if err != nil { + return pagesOrphanMarker{}, err + } + result := pagesOrphanMarker{ProjectID: projectID} + if _, exists := metadata.Extra[pagesSourceIDMetadataKey]; exists { + sourceID, err := parsePagesOrphanMetadataID(metadata.Extra, pagesSourceIDMetadataKey) + if err != nil { + return pagesOrphanMarker{}, err + } + result.SourceID = &sourceID + } + return result, nil +} + +func parsePagesOrphanMetadataID(extra map[string]any, key string) (uint, error) { + raw, exists := extra[key] + if !exists { + return 0, fmt.Errorf("pages marker %s missing", key) + } + value, ok := raw.(string) + if !ok || value == "" { + return 0, fmt.Errorf("pages marker %s must be a decimal string", key) + } + parsed, err := strconv.ParseUint(value, 10, 64) + maxModelID := uint64(^uint(0) >> 1) + if err != nil || parsed == 0 || parsed > maxModelID || strconv.FormatUint(parsed, 10) != value { + return 0, fmt.Errorf("pages marker %s is not a canonical non-zero decimal ID", key) + } + return uint(parsed), nil +} + +func sameOptionalPagesSourceID(left, right *uint) bool { + if left == nil || right == nil { + return left == nil && right == nil + } + return *left == *right +} diff --git a/internal/apps/openflare/pages/source_orphan_cleanup_test.go b/internal/apps/openflare/pages/source_orphan_cleanup_test.go new file mode 100644 index 00000000..2919c299 --- /dev/null +++ b/internal/apps/openflare/pages/source_orphan_cleanup_test.go @@ -0,0 +1,381 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package pages + +import ( + "context" + "errors" + "strconv" + "testing" + "time" + + "github.com/Rain-kl/Wavelet/internal/apps/upload" + "github.com/Rain-kl/Wavelet/internal/db" + "github.com/Rain-kl/Wavelet/internal/model" + "gorm.io/gorm" + "gorm.io/gorm/clause" +) + +func TestParsePagesOrphanMarker(t *testing.T) { + tests := []struct { + name string + extra map[string]any + wantProject uint + wantSource uint + wantSourceOK bool + wantErr bool + }{ + { + name: "manual source marker", + extra: map[string]any{ + pagesIngestMarkerKey: pagesIngestMarkerV2, + pagesProjectIDMetadataKey: "12", + }, + wantProject: 12, + }, + { + name: "persistent source marker", + extra: map[string]any{ + pagesIngestMarkerKey: pagesIngestMarkerV2, + pagesProjectIDMetadataKey: "12", + pagesSourceIDMetadataKey: "34", + }, + wantProject: 12, + wantSource: 34, + wantSourceOK: true, + }, + { + name: "project ID must be canonical decimal", + extra: map[string]any{ + pagesIngestMarkerKey: pagesIngestMarkerV2, + pagesProjectIDMetadataKey: "012", + }, + wantErr: true, + }, + { + name: "source ID must be a string", + extra: map[string]any{ + pagesIngestMarkerKey: pagesIngestMarkerV2, + pagesProjectIDMetadataKey: "12", + pagesSourceIDMetadataKey: float64(34), + }, + wantErr: true, + }, + { + name: "zero ID rejected", + extra: map[string]any{ + pagesIngestMarkerKey: pagesIngestMarkerV2, + pagesProjectIDMetadataKey: "0", + }, + wantErr: true, + }, + { + name: "wrong marker version rejected", + extra: map[string]any{ + pagesIngestMarkerKey: "pages_deployment_v1", + pagesProjectIDMetadataKey: "12", + }, + wantErr: true, + }, + } + + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + got, err := parsePagesOrphanMarker(model.UploadMetadata{Extra: test.extra}) + if gotErr := err != nil; gotErr != test.wantErr { + t.Fatalf("parsePagesOrphanMarker(%v) error = %v, want error presence = %t", test.extra, err, test.wantErr) + } + if test.wantErr { + return + } + if got.ProjectID != test.wantProject { + t.Errorf("parsePagesOrphanMarker(%v).ProjectID = %d, want %d", test.extra, got.ProjectID, test.wantProject) + } + if gotSourceOK := got.SourceID != nil; gotSourceOK != test.wantSourceOK { + t.Fatalf("parsePagesOrphanMarker(%v).SourceID presence = %t, want %t", test.extra, gotSourceOK, test.wantSourceOK) + } + if got.SourceID != nil && *got.SourceID != test.wantSource { + t.Errorf("parsePagesOrphanMarker(%v).SourceID = %d, want %d", test.extra, *got.SourceID, test.wantSource) + } + }) + } +} + +func TestReconcilePagesOrphanUploadsDeletesEligibleUploadOnce(t *testing.T) { + cleanup := setupPagesTestDB(t) + defer cleanup() + ctx := context.Background() + now := time.Now().UTC() + project := createPagesOrphanProject(t, ctx, "eligible-orphan") + candidate := createPagesOrphanUpload(t, ctx, now.Add(-3*time.Hour), project.ID, nil) + if err := upload.RebuildUploadStats(ctx); err != nil { + t.Fatalf("RebuildUploadStats() error = %v, want nil", err) + } + + summary, err := ReconcilePagesOrphanUploads(ctx, now) + if err != nil { + t.Fatalf("ReconcilePagesOrphanUploads() error = %v, want nil", err) + } + if summary.Candidates != 1 || summary.Reconciled != 1 || cleanupOutcomeTotal(summary) != 1 { + t.Errorf("ReconcilePagesOrphanUploads() summary = %+v, want one reconciled candidate", summary) + } + assertPagesCleanupUploadStatus(t, ctx, candidate.ID, model.UploadStatusDeleted) + assertPagesCleanupTotalStat(t, ctx, 0) + + second, err := ReconcilePagesOrphanUploads(ctx, now.Add(time.Minute)) + if err != nil { + t.Fatalf("second ReconcilePagesOrphanUploads() error = %v, want nil", err) + } + if second.Candidates != 0 || cleanupOutcomeTotal(second) != 0 { + t.Errorf("second ReconcilePagesOrphanUploads() summary = %+v, want empty", second) + } + assertPagesCleanupTotalStat(t, ctx, 0) +} + +func TestReconcilePagesOrphanUploadsAllowsDeletedProjectAndSource(t *testing.T) { + cleanup := setupPagesTestDB(t) + defer cleanup() + ctx := context.Background() + now := time.Now().UTC() + missingSourceID := uint(9876) + candidate := createPagesOrphanUpload(t, ctx, now.Add(-3*time.Hour), 8765, &missingSourceID) + + summary, err := ReconcilePagesOrphanUploads(ctx, now) + if err != nil { + t.Fatalf("ReconcilePagesOrphanUploads() error = %v, want nil", err) + } + if summary.Candidates != 1 || summary.Reconciled != 1 || cleanupOutcomeTotal(summary) != 1 { + t.Errorf("ReconcilePagesOrphanUploads() summary = %+v, want deleted project/source treated as one orphan", summary) + } + assertPagesCleanupUploadStatus(t, ctx, candidate.ID, model.UploadStatusDeleted) +} + +func TestReconcilePagesOrphanUploadsSkipsBusyLeaseAndSourceMismatch(t *testing.T) { + t.Run("unexpired source lease", func(t *testing.T) { + cleanup := setupPagesTestDB(t) + defer cleanup() + ctx := context.Background() + realNow := time.Now().UTC() + // A deliberately future scanner snapshot proves lease freshness uses the + // real clock after the runtime lock, not this isolation-cutoff input. + scannerNow := realNow.Add(24 * time.Hour) + project := createPagesOrphanProject(t, ctx, "busy-orphan") + source := createPagesOrphanSource(t, ctx, project.ID) + future := realNow.Add(time.Hour) + if err := db.DB(ctx).Create(&model.PagesProjectSourceRuntime{ + SourceID: source.ID, + LeaseToken: "busy-worker", + LeaseExpiresAt: &future, + }).Error; err != nil { + t.Fatalf("create busy source runtime error = %v, want nil", err) + } + candidate := createPagesOrphanUpload(t, ctx, realNow.Add(-3*time.Hour), project.ID, &source.ID) + + summary, err := ReconcilePagesOrphanUploads(ctx, scannerNow) + if err != nil { + t.Fatalf("ReconcilePagesOrphanUploads() error = %v, want nil", err) + } + if summary.LeaseBusy != 1 || cleanupOutcomeTotal(summary) != 1 { + t.Errorf("ReconcilePagesOrphanUploads() summary = %+v, want one lease-busy candidate", summary) + } + assertPagesCleanupUploadStatus(t, ctx, candidate.ID, model.UploadStatusUsed) + }) + + t.Run("source belongs to another project", func(t *testing.T) { + cleanup := setupPagesTestDB(t) + defer cleanup() + ctx := context.Background() + now := time.Now().UTC() + markerProject := createPagesOrphanProject(t, ctx, "marker-project") + actualProject := createPagesOrphanProject(t, ctx, "actual-project") + source := createPagesOrphanSource(t, ctx, actualProject.ID) + candidate := createPagesOrphanUpload(t, ctx, now.Add(-3*time.Hour), markerProject.ID, &source.ID) + + summary, err := ReconcilePagesOrphanUploads(ctx, now) + if err != nil { + t.Fatalf("ReconcilePagesOrphanUploads() error = %v, want nil", err) + } + if summary.InvalidMarker != 1 || cleanupOutcomeTotal(summary) != 1 { + t.Errorf("ReconcilePagesOrphanUploads() summary = %+v, want one ownership mismatch", summary) + } + assertPagesCleanupUploadStatus(t, ctx, candidate.ID, model.UploadStatusUsed) + }) +} + +func TestReconcilePagesOrphanUploadsRejectsMalformedMarker(t *testing.T) { + cleanup := setupPagesTestDB(t) + defer cleanup() + ctx := context.Background() + now := time.Now().UTC() + candidate := createPagesOrphanUpload(t, ctx, now.Add(-3*time.Hour), 1, nil) + metadata := candidate.Metadata + metadata.Extra[pagesProjectIDMetadataKey] = "01" + candidate.Metadata = metadata + if err := db.DB(ctx).Save(candidate).Error; err != nil { + t.Fatalf("seed malformed candidate marker error = %v, want nil", err) + } + + summary, err := ReconcilePagesOrphanUploads(ctx, now) + if err != nil { + t.Fatalf("ReconcilePagesOrphanUploads() error = %v, want nil", err) + } + if summary.InvalidMarker != 1 || cleanupOutcomeTotal(summary) != 1 { + t.Errorf("ReconcilePagesOrphanUploads() summary = %+v, want one invalid marker", summary) + } + assertPagesCleanupUploadStatus(t, ctx, candidate.ID, model.UploadStatusUsed) +} + +func TestPagesOrphanCleanupAndDeploymentCommitInterleavings(t *testing.T) { + t.Run("deployment reference commits first", func(t *testing.T) { + cleanup := setupPagesTestDB(t) + defer cleanup() + ctx := context.Background() + now := time.Now().UTC() + project := createPagesOrphanProject(t, ctx, "deployment-first") + candidate := createPagesOrphanUpload(t, ctx, now.Add(-3*time.Hour), project.ID, nil) + marker, err := parsePagesOrphanMarker(candidate.Metadata) + if err != nil { + t.Fatalf("parsePagesOrphanMarker() error = %v, want nil", err) + } + if err := db.DB(ctx).Create(&model.PagesDeployment{ + ProjectID: project.ID, + DeploymentNumber: 1, + Checksum: "deployment-first", + Status: model.PagesDeploymentStatusUploaded, + UploadID: candidate.ID, + }).Error; err != nil { + t.Fatalf("create deployment reference error = %v, want nil", err) + } + + outcome, err := reconcilePagesOrphanUploadCandidate(ctx, candidate, marker, 999, now.Add(-2*time.Hour)) + if err != nil { + t.Fatalf("reconcilePagesOrphanUploadCandidate() error = %v, want nil", err) + } + if outcome != pagesOrphanCleanupReferenced { + t.Errorf("reconcilePagesOrphanUploadCandidate() outcome = %d, want %d", outcome, pagesOrphanCleanupReferenced) + } + assertPagesCleanupUploadStatus(t, ctx, candidate.ID, model.UploadStatusUsed) + }) + + t.Run("cleanup commits first", func(t *testing.T) { + cleanup := setupPagesTestDB(t) + defer cleanup() + ctx := context.Background() + now := time.Now().UTC() + project := createPagesOrphanProject(t, ctx, "cleanup-first") + candidate := createPagesOrphanUpload(t, ctx, now.Add(-3*time.Hour), project.ID, nil) + + summary, err := ReconcilePagesOrphanUploads(ctx, now) + if err != nil { + t.Fatalf("ReconcilePagesOrphanUploads() error = %v, want nil", err) + } + if summary.Reconciled != 1 { + t.Fatalf("ReconcilePagesOrphanUploads() summary = %+v, want one reconciled candidate", summary) + } + + target := &model.PagesDeployment{ProjectID: project.ID, UploadID: candidate.ID} + err = db.DB(ctx).Transaction(func(tx *gorm.DB) error { + var lockedProject model.PagesProject + if err := tx.Clauses(clause.Locking{Strength: pagesRowLockStrength}).First(&lockedProject, project.ID).Error; err != nil { + return err + } + return lockSourceDeploymentUploadsTx(tx, target, upload.IngestResult{}, false) + }) + if !errors.Is(err, errSourceFinalFence) { + t.Errorf("final deployment upload lock after cleanup error = %v, want %v", err, errSourceFinalFence) + } + var references int64 + if err := db.DB(ctx).Model(&model.PagesDeployment{}).Where("upload_id = ?", candidate.ID).Count(&references).Error; err != nil { + t.Fatalf("count deployment references error = %v, want nil", err) + } + if references != 0 { + t.Errorf("deployment references after cleanup-first interleaving = %d, want 0", references) + } + }) +} + +func cleanupOutcomeTotal(summary PagesOrphanCleanupSummary) int { + return summary.Reconciled + summary.Referenced + summary.LeaseBusy + summary.InvalidMarker + summary.Skipped + summary.Failed +} + +func createPagesOrphanProject(t *testing.T, ctx context.Context, slug string) *model.PagesProject { + t.Helper() + project := &model.PagesProject{Name: slug, Slug: slug, Enabled: true} + if err := db.DB(ctx).Create(project).Error; err != nil { + t.Fatalf("create Pages orphan project %q error = %v, want nil", slug, err) + } + return project +} + +func createPagesOrphanSource(t *testing.T, ctx context.Context, projectID uint) *model.PagesProjectSource { + t.Helper() + source := &model.PagesProjectSource{ + ProjectID: projectID, + SourceType: PagesSourceTypeRemoteURL, + ConfigVersion: 1, + SourceIdentity: "orphan-source-identity", + } + if err := db.DB(ctx).Create(source).Error; err != nil { + t.Fatalf("create Pages orphan source for project %d error = %v, want nil", projectID, err) + } + return source +} + +func createPagesOrphanUpload( + t *testing.T, + ctx context.Context, + createdAt time.Time, + projectID uint, + sourceID *uint, +) *model.Upload { + t.Helper() + extra := map[string]any{ + pagesIngestMarkerKey: pagesIngestMarkerV2, + pagesProjectIDMetadataKey: strconv.FormatUint(uint64(projectID), 10), + } + if sourceID != nil { + extra[pagesSourceIDMetadataKey] = strconv.FormatUint(uint64(*sourceID), 10) + } + candidate := &model.Upload{ + UserID: 999, + FileName: "site.zip", + FilePath: "pages/orphan-site.zip", + FileSize: 64, + MimeType: "application/zip", + Extension: "zip", + Hash: "orphan-checksum", + Type: upload.ReservedPagesDeploymentType, + Status: model.UploadStatusUsed, + AccessMode: 0, + Metadata: model.UploadMetadata{Extra: extra}, + CreatedAt: createdAt, + UpdatedAt: createdAt, + } + if err := db.DB(ctx).Create(candidate).Error; err != nil { + t.Fatalf("create Pages orphan upload error = %v, want nil", err) + } + return candidate +} + +func assertPagesCleanupUploadStatus(t *testing.T, ctx context.Context, uploadID uint64, want model.UploadStatus) { + t.Helper() + var got model.Upload + if err := db.DB(ctx).First(&got, uploadID).Error; err != nil { + t.Fatalf("load upload %d error = %v, want nil", uploadID, err) + } + if got.Status != want { + t.Errorf("upload %d status = %q, want %q", uploadID, got.Status, want) + } +} + +func assertPagesCleanupTotalStat(t *testing.T, ctx context.Context, want int64) { + t.Helper() + var stat model.UploadStat + if err := db.DB(ctx).Where("dimension = ? AND stat_key = ?", model.UploadStatDimensionTotal, "").First(&stat).Error; err != nil { + t.Fatalf("load total upload stat error = %v, want nil", err) + } + if stat.FileCount != want { + t.Errorf("total upload stat FileCount = %d, want %d", stat.FileCount, want) + } +} diff --git a/internal/model/openflare_pages_cleanup.go b/internal/model/openflare_pages_cleanup.go new file mode 100644 index 00000000..a740a54e --- /dev/null +++ b/internal/model/openflare_pages_cleanup.go @@ -0,0 +1,75 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package model + +import ( + "context" + "errors" + "time" + + "github.com/Rain-kl/Wavelet/internal/db" +) + +const ( + // PagesOrphanUploadCandidateLimit bounds one delayed Pages upload cleanup pass. + PagesOrphanUploadCandidateLimit = 100 + + pagesOrphanMarkerPredicatePostgres = "w_uploads.metadata #>> '{extra,pages_ingest_marker}' = ?" + pagesOrphanMarkerPredicateSQLite = "CASE WHEN json_valid(w_uploads.metadata) THEN json_extract(w_uploads.metadata, '$.extra.pages_ingest_marker') ELSE NULL END = ?" +) + +// PagesOrphanUploadCandidateQuery describes the fail-closed SQL candidate set +// for delayed Pages upload compensation. +type PagesOrphanUploadCandidateQuery struct { + SystemUserID uint64 + UploadType string + Marker string + CreatedBefore time.Time +} + +// ListPagesOrphanUploadCandidates returns at most 100 unreferenced, isolated +// Pages V2 upload records. Callers must still lock and recheck every condition +// before deleting a candidate. +func ListPagesOrphanUploadCandidates( + ctx context.Context, + input PagesOrphanUploadCandidateQuery, +) ([]Upload, error) { + if input.SystemUserID == 0 || input.UploadType == "" || input.Marker == "" || input.CreatedBefore.IsZero() { + return nil, errors.New("invalid pages orphan upload candidate query") + } + markerPredicate, err := pagesOrphanMarkerPredicate(db.DB(ctx).Name()) + if err != nil { + return nil, err + } + + deploymentTable := (PagesDeployment{}).TableName() + uploadTable := (Upload{}).TableName() + var candidates []Upload + err = db.DB(ctx). + Model(&Upload{}). + Where(uploadTable+".status = ?", UploadStatusUsed). + Where(uploadTable+".user_id = ?", input.SystemUserID). + Where(uploadTable+".type = ?", input.UploadType). + Where(uploadTable+".created_at < ?", input.CreatedBefore). + Where(markerPredicate, input.Marker). + Where("NOT EXISTS (SELECT 1 FROM " + deploymentTable + " WHERE " + deploymentTable + ".upload_id = " + uploadTable + ".id)"). + Order(uploadTable + ".id ASC"). + Limit(PagesOrphanUploadCandidateLimit). + Find(&candidates).Error + if err != nil { + return nil, err + } + return candidates, nil +} + +func pagesOrphanMarkerPredicate(dialect string) (string, error) { + switch dialect { + case "postgres": + return pagesOrphanMarkerPredicatePostgres, nil + case "sqlite": + return pagesOrphanMarkerPredicateSQLite, nil + default: + return "", errors.New("unsupported database dialect for Pages orphan cleanup") + } +} diff --git a/internal/model/openflare_pages_cleanup_test.go b/internal/model/openflare_pages_cleanup_test.go new file mode 100644 index 00000000..cc249d15 --- /dev/null +++ b/internal/model/openflare_pages_cleanup_test.go @@ -0,0 +1,193 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package model + +import ( + "context" + "strings" + "testing" + "time" + + "github.com/Rain-kl/Wavelet/internal/db" + "github.com/glebarez/sqlite" + "gorm.io/gorm" +) + +func TestPagesOrphanMarkerPredicate(t *testing.T) { + tests := []struct { + name string + dialect string + want string + wantErr bool + }{ + { + name: "postgres jsonb path", + dialect: "postgres", + want: "metadata #>> '{extra,pages_ingest_marker}'", + }, + { + name: "sqlite guarded json extract", + dialect: "sqlite", + want: "CASE WHEN json_valid(w_uploads.metadata) THEN json_extract", + }, + { + name: "unknown dialect rejected", + dialect: "mysql", + wantErr: true, + }, + } + + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + got, err := pagesOrphanMarkerPredicate(test.dialect) + if gotErr := err != nil; gotErr != test.wantErr { + t.Fatalf("pagesOrphanMarkerPredicate(%q) error = %v, want error presence = %t", test.dialect, err, test.wantErr) + } + if test.want != "" && !strings.Contains(got, test.want) { + t.Errorf("pagesOrphanMarkerPredicate(%q) = %q, want substring %q", test.dialect, got, test.want) + } + }) + } +} + +func TestListPagesOrphanUploadCandidatesFiltersAndLimits(t *testing.T) { + ctx := context.Background() + gormDB := setupPagesCleanupModelTestDB(t) + cutoff := time.Now().UTC().Add(-2 * time.Hour) + old := cutoff.Add(-time.Minute) + marker := UploadMetadata{Extra: map[string]any{ + "pages_ingest_marker": "pages_deployment_v2", + "pages_project_id": "1", + }} + + valid := make([]Upload, 0, PagesOrphanUploadCandidateLimit+1) + for index := 0; index < PagesOrphanUploadCandidateLimit+1; index++ { + valid = append(valid, pagesCleanupModelUpload(uint64(index+100), 999, "openflare_pages_deployment", UploadStatusUsed, old, marker)) + } + if err := gormDB.Create(&valid).Error; err != nil { + t.Fatalf("create valid candidates error = %v, want nil", err) + } + + referenced := pagesCleanupModelUpload(1, 999, "openflare_pages_deployment", UploadStatusUsed, old, marker) + wrongOwner := pagesCleanupModelUpload(2, 1000, "openflare_pages_deployment", UploadStatusUsed, old, marker) + wrongType := pagesCleanupModelUpload(3, 999, "generic", UploadStatusUsed, old, marker) + wrongStatus := pagesCleanupModelUpload(4, 999, "openflare_pages_deployment", UploadStatusPending, old, marker) + fresh := pagesCleanupModelUpload(5, 999, "openflare_pages_deployment", UploadStatusUsed, cutoff, marker) + wrongMarker := pagesCleanupModelUpload(6, 999, "openflare_pages_deployment", UploadStatusUsed, old, UploadMetadata{Extra: map[string]any{ + "pages_ingest_marker": "pages_deployment_v1", + "pages_project_id": "1", + }}) + for _, upload := range []Upload{referenced, wrongOwner, wrongType, wrongStatus, fresh, wrongMarker} { + if err := gormDB.Create(&upload).Error; err != nil { + t.Fatalf("create filtered upload %d error = %v, want nil", upload.ID, err) + } + } + if err := gormDB.Create(&PagesDeployment{ + ProjectID: 1, + DeploymentNumber: 1, + Checksum: "referenced", + Status: PagesDeploymentStatusUploaded, + UploadID: referenced.ID, + }).Error; err != nil { + t.Fatalf("create referenced deployment error = %v, want nil", err) + } + + invalidJSON := pagesCleanupModelUpload(7, 999, "openflare_pages_deployment", UploadStatusUsed, old, marker) + if err := gormDB.Create(&invalidJSON).Error; err != nil { + t.Fatalf("create invalid JSON upload error = %v, want nil", err) + } + if err := gormDB.Table((Upload{}).TableName()).Where("id = ?", invalidJSON.ID). + UpdateColumn("metadata", "{invalid").Error; err != nil { + t.Fatalf("corrupt upload metadata error = %v, want nil", err) + } + + got, err := ListPagesOrphanUploadCandidates(ctx, PagesOrphanUploadCandidateQuery{ + SystemUserID: 999, + UploadType: "openflare_pages_deployment", + Marker: "pages_deployment_v2", + CreatedBefore: cutoff, + }) + if err != nil { + t.Fatalf("ListPagesOrphanUploadCandidates() error = %v, want nil", err) + } + if len(got) != PagesOrphanUploadCandidateLimit { + t.Fatalf("ListPagesOrphanUploadCandidates() count = %d, want %d", len(got), PagesOrphanUploadCandidateLimit) + } + for index, candidate := range got { + wantID := uint64(index + 100) + if candidate.ID != wantID { + t.Errorf("ListPagesOrphanUploadCandidates()[%d].ID = %d, want %d", index, candidate.ID, wantID) + } + } +} + +func TestListPagesOrphanUploadCandidatesSkipsInvalidSQLiteJSON(t *testing.T) { + ctx := context.Background() + gormDB := setupPagesCleanupModelTestDB(t) + cutoff := time.Now().UTC().Add(-2 * time.Hour) + upload := pagesCleanupModelUpload(1, 999, "openflare_pages_deployment", UploadStatusUsed, cutoff.Add(-time.Minute), UploadMetadata{}) + if err := gormDB.Create(&upload).Error; err != nil { + t.Fatalf("create invalid JSON candidate error = %v, want nil", err) + } + if err := gormDB.Table((Upload{}).TableName()).Where("id = ?", upload.ID). + UpdateColumn("metadata", "{invalid").Error; err != nil { + t.Fatalf("corrupt upload metadata error = %v, want nil", err) + } + + got, err := ListPagesOrphanUploadCandidates(ctx, PagesOrphanUploadCandidateQuery{ + SystemUserID: 999, + UploadType: "openflare_pages_deployment", + Marker: "pages_deployment_v2", + CreatedBefore: cutoff, + }) + if err != nil { + t.Fatalf("ListPagesOrphanUploadCandidates(invalid JSON) error = %v, want nil", err) + } + if len(got) != 0 { + t.Errorf("ListPagesOrphanUploadCandidates(invalid JSON) count = %d, want 0", len(got)) + } +} + +func setupPagesCleanupModelTestDB(t *testing.T) *gorm.DB { + t.Helper() + + gormDB, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{ + DisableForeignKeyConstraintWhenMigrating: true, + }) + if err != nil { + t.Fatalf("open Pages cleanup model test database error = %v, want nil", err) + } + if err := gormDB.AutoMigrate(&Upload{}, &PagesDeployment{}); err != nil { + t.Fatalf("migrate Pages cleanup model test database error = %v, want nil", err) + } + db.SetDB(gormDB) + t.Cleanup(func() { db.SetDB(nil) }) + return gormDB +} + +func pagesCleanupModelUpload( + id uint64, + userID uint64, + uploadType string, + status UploadStatus, + createdAt time.Time, + metadata UploadMetadata, +) Upload { + return Upload{ + ID: id, + UserID: userID, + FileName: "site.zip", + FilePath: "pages/site.zip", + FileSize: 10, + MimeType: "application/zip", + Extension: "zip", + Hash: "checksum", + Type: uploadType, + Status: status, + AccessMode: 0, + Metadata: metadata, + CreatedAt: createdAt, + UpdatedAt: createdAt, + } +}