fix(pages): 增加部署包孤儿补偿

按项目、来源、运行时与上传记录锁序补偿异常中断遗留的部署包。\n同时隐藏并保护系统内部排程,避免通用任务管理入口修改 scanner。
This commit is contained in:
deqiying
2026-07-19 19:08:43 +08:00
parent c39a3edcc3
commit 848884d8cd
9 changed files with 1119 additions and 5 deletions
+7 -1
View File
@@ -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": {
+7 -1
View File
@@ -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": {
+5 -1
View File
@@ -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:
+20 -2
View File
@@ -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))
+103
View File
@@ -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) {
@@ -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
}
@@ -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)
}
}
+75
View File
@@ -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")
}
}
@@ -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,
}
}