diff --git a/internal/creator/integration_test.go b/internal/creator/integration_test.go index 426ae44..5f46dbe 100644 --- a/internal/creator/integration_test.go +++ b/internal/creator/integration_test.go @@ -515,6 +515,85 @@ func TestCreatorPostgresAccountMetricSnapshot(t *testing.T) { } } +func TestCreatorPostgresAccountMonitorViewAndCollectionStatus(t *testing.T) { + store, phaseAStore, ctx := openCreatorIntegrationStore(t) + stamp := fmt.Sprintf("%d", time.Now().UnixNano()) + accountID := createIntegrationAccount(t, ctx, phaseAStore, "metric"+stamp) + now := time.Now().UTC().Truncate(time.Second) + + // 无 checkpoint:采集状态为 pending。 + status, err := store.GetAccountCollectionStatus(ctx, accountID) + if err != nil || status.Works.Status != "pending" || status.Comments.Status != "pending" { + t.Fatalf("pending collection status: %+v err=%v", status, err) + } + // 无作品:监控视图作品数为 0。 + views, err := store.ListAccountMonitorViews(ctx) + if err != nil { + t.Fatal(err) + } + viewFor := func(id string) *AccountMonitorView { + for index := range views { + if views[index].ID == id { + return &views[index] + } + } + return nil + } + view := viewFor(accountID) + if view == nil || view.WorkCount != 0 || view.LatestPublishedAt != nil { + t.Fatalf("monitor view without works: %+v", view) + } + + // 写入 checkpoint(succeeded + 窗口)与 owned 作品后:状态带下次窗口,视图聚合作品数。 + lease, err := store.beginCheckpoint(ctx, SourceOwned, accountID, "works", now.Add(-24*time.Hour), now) + if err != nil { + t.Fatal(err) + } + if err := store.finishCheckpoint(ctx, SourceOwned, accountID, "works", lease, "succeeded", ""); err != nil { + t.Fatal(err) + } + published := now.Add(-time.Hour) + if _, _, err := store.UpsertWork(ctx, WorkInput{Platform: PlatformDouyin, WorkKey: "monitor-view-" + stamp, SourceType: SourceOwned, SourceID: accountID, Title: "MonitorView", PublishedAt: &published, PublishedAtStatus: "verified"}, now); err != nil { + t.Fatal(err) + } + status, err = store.GetAccountCollectionStatus(ctx, accountID) + if err != nil || status.Works.Status != "succeeded" || status.Works.LastCompletedAt == nil || + status.Works.NextWindowStart == nil || status.Works.NextWindowEnd == nil { + t.Fatalf("succeeded collection status must carry next window: %+v err=%v", status, err) + } + if !status.Works.NextWindowStart.Before(*status.Works.NextWindowEnd) { + t.Fatalf("next window start must precede end: %+v", status.Works) + } + view = viewFor(accountID) + if view == nil || view.WorkCount != 1 || view.LatestPublishedAt == nil { + // 插入作品后重拉视图:聚合来自 creator_work_source 实时 JOIN。 + views, err = store.ListAccountMonitorViews(ctx) + if err != nil { + t.Fatal(err) + } + view = viewFor(accountID) + if view == nil || view.WorkCount != 1 || view.LatestPublishedAt == nil { + t.Fatalf("monitor view aggregates: %+v", view) + } + } + // failed checkpoint 带错误信息(普通错误落在 failed,不是 blocked)。 + lease, err = store.beginCheckpoint(ctx, SourceOwned, accountID, "comments", now.Add(-24*time.Hour), now) + if err != nil { + t.Fatal(err) + } + if err := store.failCheckpoint(ctx, SourceOwned, accountID, "comments", lease, errors.New("boom")); err == nil { + t.Fatal("failCheckpoint must return the primary error") + } + status, err = store.GetAccountCollectionStatus(ctx, accountID) + if err != nil || status.Comments.Status != "failed" || status.Comments.LastError == "" { + t.Fatalf("failed comments checkpoint surfaces error: %+v err=%v", status.Comments, err) + } + // 空账号 ID 拒绝。 + if _, err := store.GetAccountCollectionStatus(ctx, " "); !errors.Is(err, ErrInvalid) { + t.Fatalf("blank account id must be rejected: err=%v", err) + } +} + func TestCreatorPostgresCollectionAllowsMutedReadAccount(t *testing.T) { store, phaseAStore, ctx := openCreatorIntegrationStore(t) stamp := time.Now().UnixNano()