feat: unify account environments and work collection

This commit is contained in:
2026-10-06 13:00:29 +08:00
parent b72727451d
commit 950b23da32
22 changed files with 935 additions and 93 deletions
+18 -2
View File
@@ -234,8 +234,10 @@ func (s *Store) ownedWorkStats(ctx context.Context) (map[string]ownedWorkStat, e
// AccountCollectionStatus 自有账号采集状态(checkpoint 形态,对齐竞品的 sync_status 展示语义)。
type AccountCollectionStatus struct {
Works AccountCheckpointStatus `json:"works"`
Comments AccountCheckpointStatus `json:"comments"`
WorkCount int64 `json:"work_count"`
AwemeCount *int64 `json:"aweme_count"`
Works AccountCheckpointStatus `json:"works"`
Comments AccountCheckpointStatus `json:"comments"`
}
type AccountCheckpointStatus struct {
@@ -254,6 +256,20 @@ func (s *Store) GetAccountCollectionStatus(ctx context.Context, accountID string
return AccountCollectionStatus{}, ErrInvalid
}
status := AccountCollectionStatus{}
var total sql.NullInt64
err := s.db.QueryRowContext(ctx, `
SELECT
(SELECT COUNT(*) FROM creator_work WHERE source_type = $1 AND source_id = $2),
(SELECT metric.aweme_count FROM creator_account_metric metric
JOIN social_account account ON account.id = metric.account_id
WHERE account.account_id = $2 AND metric.aweme_count IS NOT NULL
ORDER BY metric.collected_at DESC LIMIT 1)`, SourceOwned, accountID).Scan(&status.WorkCount, &total)
if err != nil {
return AccountCollectionStatus{}, databaseError(err)
}
if total.Valid {
status.AwemeCount = &total.Int64
}
for _, kind := range []struct {
name string
pointer *AccountCheckpointStatus
+3 -1
View File
@@ -403,7 +403,9 @@ func (s *Store) CollectSource(ctx context.Context, platform, sourceType, sourceI
return err
}
report.WorksSeen++
if work.PublishedAt != nil && work.PublishedAt.Before(report.WindowStart) {
// Owned accounts keep their entire work library; the lookback window
// still bounds competitor works and the later comments phase.
if sourceType == SourceCompetitor && work.PublishedAt != nil && work.PublishedAt.Before(report.WindowStart) {
continue
}
work.Platform, work.SourceType, work.SourceID = platform, sourceType, sourceID
@@ -0,0 +1,145 @@
package creator
import (
"context"
"encoding/json"
"fmt"
"testing"
"time"
)
type ownedHistoryCollector struct {
pages map[string]WorkPage
workCursors []string
commentWorks []string
}
func (c *ownedHistoryCollector) ListWorks(_ context.Context, _, cursor string) (WorkPage, error) {
c.workCursors = append(c.workCursors, cursor)
page, ok := c.pages[cursor]
if !ok {
return WorkPage{}, fmt.Errorf("unexpected cursor %q", cursor)
}
return page, nil
}
func (c *ownedHistoryCollector) ListTopLevelComments(_ context.Context, workKey, _ string) (CommentPage, error) {
c.commentWorks = append(c.commentWorks, workKey)
return CommentPage{}, nil
}
func TestCreatorPostgresOwnedCollectionIncludesHistoryAndProgress(t *testing.T) {
store, accounts, ctx := openCreatorIntegrationStore(t)
accountID := createIntegrationAccount(t, ctx, accounts, fmt.Sprint(time.Now().UnixNano()))
now := time.Now().UTC().Truncate(time.Second)
settings, err := store.GetSettings(ctx)
if err != nil {
t.Fatal(err)
}
old := now.Add(-time.Duration(settings.LookbackDays+365) * 24 * time.Hour)
recent := now.Add(-time.Hour)
pages := map[string]WorkPage{"": {NextCursor: "history", HasMore: true}}
for i := 0; i < 28; i++ {
published := old
if i == 0 {
published = recent
}
work := WorkInput{WorkKey: fmt.Sprintf("owned-history-%02d", i), Title: fmt.Sprintf("作品 %d", i), PublishedAt: &published, PublishedAtStatus: "verified"}
cursor := "history"
if i < 14 {
cursor = ""
}
page := pages[cursor]
page.Items = append(page.Items, work)
pages[cursor] = page
}
collector := &ownedHistoryCollector{pages: pages}
total := int64(30)
if err := store.RecordAccountMetric(ctx, AccountMetricInput{AccountID: accountID, CollectedAt: now.Add(-time.Minute), AwemeCount: &total}); err != nil {
t.Fatal(err)
}
// A newer snapshot lacking aweme_count must not erase the known platform total.
followers := int64(100)
if err := store.RecordAccountMetric(ctx, AccountMetricInput{AccountID: accountID, CollectedAt: now, FollowerCount: &followers}); err != nil {
t.Fatal(err)
}
report, err := store.CollectSource(ctx, PlatformDouyin, SourceOwned, accountID, collector, now)
if err != nil || !report.PaginationComplete || report.WorksSeen != 28 || report.WorksSaved != 28 {
t.Fatalf("full history collection: report=%+v err=%v", report, err)
}
if len(collector.workCursors) != 2 || collector.workCursors[1] != "history" {
t.Fatalf("history pagination: %v", collector.workCursors)
}
// History is listed, but comments remain limited to the configured window.
if len(collector.commentWorks) != 1 || collector.commentWorks[0] != "owned-history-00" {
t.Fatalf("historical comments should not be scheduled: %v", collector.commentWorks)
}
first, err := store.ListWorksPage(ctx, WorkFilter{SourceType: SourceOwned, SourceID: accountID}, 1, 25)
if err != nil || len(first.Data) != 25 || first.Total != 28 || !first.HasNext {
t.Fatalf("first history page: page=%+v err=%v", first, err)
}
second, err := store.ListWorksPage(ctx, WorkFilter{SourceType: SourceOwned, SourceID: accountID}, 2, 25)
if err != nil || len(second.Data) != 3 || second.Total != 28 || second.HasNext {
t.Fatalf("last history page: page=%+v err=%v", second, err)
}
status, err := store.GetAccountCollectionStatus(ctx, accountID)
if err != nil {
t.Fatal(err)
}
// Read JSON fields to keep the regression test compilable before implementation.
assertCollectionCounts(t, status, 28, &total)
views, err := store.ListAccountMonitorViews(ctx)
if err != nil || len(views) != 1 || views[0].WorkCount != 28 || views[0].AwemeCount == nil || *views[0].AwemeCount != total {
t.Fatalf("list/detail count consistency: views=%+v err=%v", views, err)
}
// A full subsequent run must not count duplicate works twice.
if _, err := store.CollectSource(ctx, PlatformDouyin, SourceOwned, accountID, collector, now.Add(time.Second)); err != nil {
t.Fatal(err)
}
status, err = store.GetAccountCollectionStatus(ctx, accountID)
if err != nil {
t.Fatal(err)
}
assertCollectionCounts(t, status, 28, &total)
}
func assertCollectionCounts(t *testing.T, status AccountCollectionStatus, collected int64, total *int64) {
t.Helper()
data, err := json.Marshal(status)
if err != nil {
t.Fatal(err)
}
var fields map[string]json.RawMessage
if err := json.Unmarshal(data, &fields); err != nil {
t.Fatal(err)
}
if string(fields["work_count"]) != fmt.Sprint(collected) {
t.Fatalf("collected count: status=%s want=%d", data, collected)
}
want := "null"
if total != nil {
want = fmt.Sprint(*total)
}
if string(fields["aweme_count"]) != want {
t.Fatalf("platform total: status=%s want=%s", data, want)
}
}
func TestCreatorPostgresCollectionProgressDistinguishesUnknownAndZero(t *testing.T) {
store, accounts, ctx := openCreatorIntegrationStore(t)
accountID := createIntegrationAccount(t, ctx, accounts, fmt.Sprint(time.Now().UnixNano()))
status, err := store.GetAccountCollectionStatus(ctx, accountID)
if err != nil || status.Works.Status != "pending" {
t.Fatalf("pending status: %+v err=%v", status, err)
}
assertCollectionCounts(t, status, 0, nil)
zero := int64(0)
if err := store.RecordAccountMetric(ctx, AccountMetricInput{AccountID: accountID, CollectedAt: time.Now().UTC(), AwemeCount: &zero}); err != nil {
t.Fatal(err)
}
status, err = store.GetAccountCollectionStatus(ctx, accountID)
if err != nil {
t.Fatal(err)
}
assertCollectionCounts(t, status, 0, &zero)
}