test(creator): 账号监控视图与采集状态集成断言(聚合/下次窗口/failed 错误可见)

This commit is contained in:
2026-09-27 22:20:45 +08:00
parent fb1db68f13
commit 8c1fce3819
+79
View File
@@ -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()