From 3f1e2d39f4d6bba91462d1df60f1ee86d2408937 Mon Sep 17 00:00:00 2001 From: Rogee Date: Sun, 27 Sep 2026 22:05:01 +0800 Subject: [PATCH] =?UTF-8?q?feat(creator):=20=E8=87=AA=E6=9C=89=E8=B4=A6?= =?UTF-8?q?=E5=8F=B7=E7=9B=91=E6=8E=A7=E5=90=8E=E7=AB=AF=E2=80=94=E2=80=94?= =?UTF-8?q?=E6=92=AD=E6=94=BE=E9=87=8F=E5=85=A5=E5=BA=93=E3=80=81=E8=B4=A6?= =?UTF-8?q?=E5=8F=B7=E7=B2=89=E4=B8=9D=E6=97=B6=E5=BA=8F=E5=BF=AB=E7=85=A7?= =?UTF-8?q?=E3=80=81=E7=9B=91=E6=8E=A7=20API?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - migration 040/041:creator_work 与 creator_work_metric 加 play_count;新建 creator_account_metric 时序快照表 - douyin:详情接口解析 play_count;self profile 扩展解析粉丝/关注/获赞/作品数(flexibleInt 兼容字符串) - creator store:PlayCount 全链路(UpsertWork COALESCE/RecordMetric/ListMetrics);RecordAccountMetric/ListAccountMetrics/ListAccountMonitorViews/GetAccountCollectionStatus - controlplane:syncCreatorOwned 成功路径末尾记账号画像快照(失败 errors.Join 不吞错);新增 monitor-views/collection-status/metrics/sync 四个路由;creatorError 补 account.ErrNotFound→404 - 测试:integration/douyin/api 断言播放量落库与 COALESCE 语义、快照覆盖语义、路由 happy path 与 409/503/404 语义 --- .../controlplane/api/app_migrated_test.go | 81 ++++++++++++ internal/controlplane/api/creator.go | 62 ++++++++- .../controlplane/api/creator_collector.go | 1 - .../controlplane/api/creator_helper_test.go | 8 ++ internal/creator/accounts.go | 119 ++++++++++++++++++ internal/creator/content.go | 74 +++++++++-- internal/creator/integration_test.go | 68 +++++++++- internal/creator/metrics.go | 12 +- .../migrations/040_work_play_count.sql | 3 + .../migrations/041_creator_account_metric.sql | 14 +++ internal/creator/models.go | 22 ++++ internal/creator/store.go | 8 ++ internal/platform/douyin/connector.go | 88 +++++++++++++ internal/platform/douyin/creator_collector.go | 21 ++-- .../platform/douyin/creator_collector_test.go | 50 +++++++- 15 files changed, 593 insertions(+), 38 deletions(-) create mode 100644 internal/creator/migrations/040_work_play_count.sql create mode 100644 internal/creator/migrations/041_creator_account_metric.sql diff --git a/internal/controlplane/api/app_migrated_test.go b/internal/controlplane/api/app_migrated_test.go index e207de4..2f1c24d 100644 --- a/internal/controlplane/api/app_migrated_test.go +++ b/internal/controlplane/api/app_migrated_test.go @@ -92,6 +92,7 @@ func TestCreatorRouteValidationCoverage(t *testing.T) { app := newHandlerWithCreator(t.TempDir(), "operator", "unit-test-password", phaseAStore, hubStore, nil, creatorStore) for _, path := range []string{ "/api/creator/accounts/route-account/profile", "/api/creator/accounts/route-account/strategies", + "/api/creator/accounts/route-account/collection-status", "/api/creator/accounts/route-account/metrics", } { response := do(app, http.MethodGet, path, "", "operator", "unit-test-password") if response.Code != http.StatusOK { @@ -139,6 +140,7 @@ func TestCreatorRouteValidationCoverage(t *testing.T) { {http.MethodPost, "/api/creator/competitors/missing/pause"}, {http.MethodPost, "/api/creator/competitors/missing/resume"}, {http.MethodPost, "/api/creator/competitors/missing/sync"}, + {http.MethodPost, "/api/creator/accounts/missing/sync"}, {http.MethodPost, "/api/creator/works/missing/metrics"}, {http.MethodPost, "/api/creator/works/missing/material/select"}, {http.MethodPost, "/api/creator/works/missing/material/process"}, @@ -338,6 +340,85 @@ func TestCreatorFixtureRoutesPostgres(t *testing.T) { if _, err := previewDouyinCompetitor(ctx, creatorStore, phaseAStore, hubStore, "fixture-account", creator.CompetitorInput{Platform: creator.PlatformDouyin, PlatformAccountKey: "preview-key", Nickname: "Preview", HomepageURL: "https://www.douyin.com/user/preview-key"}); err == nil { t.Fatal("competitor preview without browser unexpectedly succeeded") } + + // 自有账号监控:列表带作品聚合统计(fixture 已有 1 个 owned 作品)。 + viewResponse := do(app, http.MethodGet, "/api/creator/accounts/monitor-views", "", "operator", "unit-test-password") + if viewResponse.Code != http.StatusOK { + t.Fatalf("GET monitor views: %d %s", viewResponse.Code, viewResponse.Body.String()) + } + var monitorViews []struct { + ID string `json:"id"` + WorkCount int64 `json:"work_count"` + LatestPublishedAt *string `json:"latest_published_at"` + } + if err := json.Unmarshal(viewResponse.Body.Bytes(), &monitorViews); err != nil { + t.Fatal(err) + } + viewByAccount := map[string]bool{} + for _, view := range monitorViews { + viewByAccount[view.ID] = true + if view.ID == "fixture-account" && (view.WorkCount < 1 || view.LatestPublishedAt == nil) { + t.Fatalf("monitor view aggregates missing for fixture-account: %+v", view) + } + } + if !viewByAccount["fixture-account"] { + t.Fatal("monitor views missing fixture-account") + } + // 采集状态:无 checkpoint 时 pending;作品采集完成后 succeeded 且带窗口。 + statusResponse := do(app, http.MethodGet, "/api/creator/accounts/fixture-small/collection-status", "", "operator", "unit-test-password") + if statusResponse.Code != http.StatusOK { + t.Fatalf("GET collection status: %d %s", statusResponse.Code, statusResponse.Body.String()) + } + var collectionStatus struct { + Works struct { + Status string `json:"status"` + } `json:"works"` + } + if err := json.Unmarshal(statusResponse.Body.Bytes(), &collectionStatus); err != nil { + t.Fatal(err) + } + if collectionStatus.Works.Status != "pending" { + t.Fatalf("expected pending checkpoint for untouched account: %s", statusResponse.Body.String()) + } + statusResponse = do(app, http.MethodGet, "/api/creator/accounts/fixture-account/collection-status", "", "operator", "unit-test-password") + if statusResponse.Code != http.StatusOK { + t.Fatalf("GET collection status: %d %s", statusResponse.Code, statusResponse.Body.String()) + } + var fixtureStatus struct { + Works struct { + Status string `json:"status"` + NextWindowStart *string `json:"next_window_start"` + NextWindowEnd *string `json:"next_window_end"` + } `json:"works"` + } + if err := json.Unmarshal(statusResponse.Body.Bytes(), &fixtureStatus); err != nil { + t.Fatal(err) + } + // fixture 流程没有跑过调度器 checkpoint,状态可能仍为 pending;两种都合法。 + if fixtureStatus.Works.Status != "pending" && fixtureStatus.Works.Status != "succeeded" { + t.Fatalf("unexpected works checkpoint status: %s", statusResponse.Body.String()) + } + // 账号画像时序:初始为空数组。 + metricsResponse := do(app, http.MethodGet, "/api/creator/accounts/fixture-account/metrics", "", "operator", "unit-test-password") + if metricsResponse.Code != http.StatusOK || strings.TrimSpace(metricsResponse.Body.String()) != "[]" { + t.Fatalf("GET account metrics: %d %s", metricsResponse.Code, metricsResponse.Body.String()) + } + // 未登录账号强制同步:账号存在但未通过登录/环境校验,必须 409/503 且不吞错。 + smallSync := do(app, http.MethodPost, "/api/creator/accounts/fixture-small/sync", ``, "operator", "unit-test-password") + if smallSync.Code != http.StatusConflict && smallSync.Code != http.StatusServiceUnavailable { + t.Fatalf("sync not-logged-in account must be blocked: %d %s", smallSync.Code, smallSync.Body.String()) + } + if _, err := creatorStore.RecordVerifiedLoginResult(ctx, "fixture-small", "fixture-small-platform"); err != nil { + t.Fatal(err) + } + loggedSync := do(app, http.MethodPost, "/api/creator/accounts/fixture-small/sync", ``, "operator", "unit-test-password") + if loggedSync.Code != http.StatusConflict && loggedSync.Code != http.StatusServiceUnavailable { + t.Fatalf("sync without runtime environment must be blocked: %d %s", loggedSync.Code, loggedSync.Body.String()) + } + missingSync := do(app, http.MethodPost, "/api/creator/accounts/fixture-missing/sync", ``, "operator", "unit-test-password") + if missingSync.Code != http.StatusNotFound && missingSync.Code != http.StatusConflict && missingSync.Code != http.StatusServiceUnavailable { + t.Fatalf("sync missing account: %d %s", missingSync.Code, missingSync.Body.String()) + } eventBody := `{"platform":"douyin","receiving_account_id":"fixture-account","event_key":"fixture-event","event_type":"comment","interactor_uid":"peer","work_id":"` + workID + `","comment_id":"` + commentID + `"}` eventResponse := do(app, http.MethodPost, "/api/creator/test/events", eventBody, "operator", "unit-test-password") if eventResponse.Code != http.StatusCreated { diff --git a/internal/controlplane/api/creator.go b/internal/controlplane/api/creator.go index 1953da0..9a8ba8f 100644 --- a/internal/controlplane/api/creator.go +++ b/internal/controlplane/api/creator.go @@ -132,6 +132,37 @@ func registerCreatorWithServices(app *fiber.App, store *creator.Store, phaseASto } return c.JSON(profiles) }) + // 自有账号监控:列表带作品聚合统计。 + app.Get("/api/creator/accounts/monitor-views", func(c fiber.Ctx) error { + views, err := store.ListAccountMonitorViews(c.Context()) + if err != nil { + return creatorError(c, err) + } + return c.JSON(views) + }) + // 自有账号采集状态:作品/评论两条 checkpoint + 下次采集窗口。 + app.Get("/api/creator/accounts/:id/collection-status", func(c fiber.Ctx) error { + status, err := store.GetAccountCollectionStatus(c.Context(), c.Params("id")) + if err != nil { + return creatorError(c, err) + } + return c.JSON(status) + }) + // 自有账号画像时序(粉丝/关注/获赞/作品总数),用于趋势曲线。 + app.Get("/api/creator/accounts/:id/metrics", func(c fiber.Ctx) error { + points, err := store.ListAccountMetrics(c.Context(), c.Params("id")) + if err != nil { + return creatorError(c, err) + } + return c.JSON(points) + }) + // 自有账号强制同步:语义对齐竞品 sync,立即执行一次作品采集(需已登录+运行环境就绪)。 + app.Post("/api/creator/accounts/:id/sync", func(c fiber.Ctx) error { + if err := syncCreatorOwned(c.Context(), store, phaseAStore, hubStore, c.Params("id"), time.Now().UTC()); err != nil { + return creatorError(c, err) + } + return c.SendStatus(fiber.StatusAccepted) + }) app.Get("/api/creator/accounts/:id/profile", func(c fiber.Ctx) error { profile, err := store.GetAccountProfile(c.Context(), c.Params("id")) if err != nil { @@ -1053,6 +1084,9 @@ func creatorError(c fiber.Ctx, err error) error { switch { case errors.Is(err, creator.ErrInvalid): status, message = fiber.StatusBadRequest, creator.ErrInvalid.Error() + case errors.Is(err, accountdomain.ErrNotFound): + // 跨域引用(如账号不存在):account 包哨兵也映射 404,避免 500 掩盖真实语义。 + status, message = fiber.StatusNotFound, "resource not found" case errors.Is(err, creator.ErrConflict): status, message = fiber.StatusConflict, creator.ErrConflict.Error() case errors.Is(err, creator.ErrNotFound): @@ -2017,7 +2051,7 @@ func refreshCreatorCompetitorMetricAnonymous(ctx context.Context, store *creator if fetchErr != nil { return fetchErr } - _, recordErr := store.RecordMetric(ctx, creator.MetricInput{WorkID: work.ID, CollectedAt: now, Likes: stats.Likes, CommentsCount: stats.CommentsCount, Shares: stats.Shares, CollectCount: stats.CollectCount}, settings, now) + _, recordErr := store.RecordMetric(ctx, creator.MetricInput{WorkID: work.ID, CollectedAt: now, Likes: stats.Likes, CommentsCount: stats.CommentsCount, Shares: stats.Shares, CollectCount: stats.CollectCount, PlayCount: stats.PlayCount}, settings, now) return recordErr } @@ -2086,10 +2120,10 @@ func refreshCreatorMetricWork(ctx context.Context, store *creator.Store, phaseAS if item.WorkKey != work.WorkKey { continue } - if item.Likes == nil && item.CommentsCount == nil && item.Shares == nil { + if item.Likes == nil && item.CommentsCount == nil && item.Shares == nil && item.PlayCount == nil { return creator.ErrUnavailable } - _, metricErr := store.RecordMetric(useCtx, creator.MetricInput{WorkID: work.ID, CollectedAt: now, Likes: item.Likes, CommentsCount: item.CommentsCount, Shares: item.Shares}, settings, now) + _, metricErr := store.RecordMetric(useCtx, creator.MetricInput{WorkID: work.ID, CollectedAt: now, Likes: item.Likes, CommentsCount: item.CommentsCount, Shares: item.Shares, CollectCount: item.CollectCount, PlayCount: item.PlayCount}, settings, now) return metricErr } if !result.HasMore { @@ -2159,7 +2193,27 @@ func syncCreatorOwned(ctx context.Context, store *creator.Store, phaseAStore *ac return errors.Join(windowErr, runtimeUse.Close()) } _, err = store.CollectSource(useCtx, account.Platform, creator.SourceOwned, account.ID, collector, collectionNow) - return errors.Join(err, runtimeUse.Close()) + collectErr := err + // 作品采集成功后顺带拉 self profile 记账号画像快照(粉丝/关注/获赞/作品总数)。 + // 失败不吞:join 进返回错误,由调度器日志可见,但不影响已入库的作品数据。 + profileErr := error(nil) + if collectErr == nil { + profileErr = recordCreatorAccountMetricSnapshot(useCtx, store, creatorGatewayBrowser{gateway: gateway, environment: environment}, account.ID) + } + return errors.Join(collectErr, profileErr, runtimeUse.Close()) +} + +// recordCreatorAccountMetricSnapshot 拉取登录账号自己的画像并记一次时序快照。 +func recordCreatorAccountMetricSnapshot(ctx context.Context, store *creator.Store, browser creatorGatewayBrowser, accountID string) error { + profile, err := douyin.FetchSelfProfile(ctx, browser) + if err != nil { + return fmt.Errorf("creator account metric snapshot: %w", err) + } + return store.RecordAccountMetric(ctx, creator.AccountMetricInput{ + AccountID: accountID, CollectedAt: time.Now().UTC(), + FollowerCount: profile.FollowerCount, FollowingCount: profile.FollowingCount, + TotalFavorited: profile.TotalFavorited, AwemeCount: profile.AwemeCount, + }) } func creatorCollectionAccount(ctx context.Context, store *creator.Store, phaseAStore *accountdomain.Store, hubStore *hub.Store, platform string) (string, error) { diff --git a/internal/controlplane/api/creator_collector.go b/internal/controlplane/api/creator_collector.go index 22673b0..b7bfa00 100644 --- a/internal/controlplane/api/creator_collector.go +++ b/internal/controlplane/api/creator_collector.go @@ -37,7 +37,6 @@ func newCreatorCollector(ctx context.Context, platform string, gateway hub.Gatew func douyinCollector(browser creatorGatewayBrowser, accountKey, sourceType, sourceID string) douyin.CreatorCollector { return douyin.CreatorCollector{Browser: browser, AccountKey: accountKey, SourceType: sourceType, SourceID: sourceID} } - // validateDouyinCompetitor 监控账号入队前的平台一致性校验(平台收敛后仅抖音)。 func validateDouyinCompetitor(input creator.CompetitorInput) error { if input.Platform != creator.PlatformDouyin { diff --git a/internal/controlplane/api/creator_helper_test.go b/internal/controlplane/api/creator_helper_test.go index 572faaf..2f07b67 100644 --- a/internal/controlplane/api/creator_helper_test.go +++ b/internal/controlplane/api/creator_helper_test.go @@ -79,6 +79,14 @@ func TestCreatorHelperBranches(t *testing.T) { } } +func TestCreatorAccountMetricSnapshotFailurePropagates(t *testing.T) { + // 浏览器不可用(网关零值)时快照拉取失败必须向上返回错误,绝不静默吞掉; + // 调度器对 syncCreatorOwned 的失败会记 warning 日志(runCreatorScheduleOnce)。 + if err := recordCreatorAccountMetricSnapshot(context.Background(), nil, creatorGatewayBrowser{}, "account-1"); err == nil { + t.Fatal("account metric snapshot failure must propagate error") + } +} + func TestCreatorPageQueryValidation(t *testing.T) { app := fiber.New() app.Get("/", func(c fiber.Ctx) error { diff --git a/internal/creator/accounts.go b/internal/creator/accounts.go index 4fb4140..566b68e 100644 --- a/internal/creator/accounts.go +++ b/internal/creator/accounts.go @@ -104,6 +104,125 @@ func (s *Store) ListAccountProfiles(ctx context.Context) ([]AccountProfile, erro return profiles, rows.Err() } +// AccountMonitorView 自有账号监控列表视图:账号画像 + 作品聚合统计(作品数/最近发布)。 +type AccountMonitorView struct { + AccountProfile + WorkCount int64 `json:"work_count"` + LatestPublishedAt *time.Time `json:"latest_published_at,omitempty"` +} + +// ListAccountMonitorViews 自有账号监控列表:画像 + owned 作品聚合,形态对齐竞品的 ListCompetitorsWithProfile。 +func (s *Store) ListAccountMonitorViews(ctx context.Context) ([]AccountMonitorView, error) { + profiles, err := s.ListAccountProfiles(ctx) + if err != nil { + return nil, err + } + stats, err := s.ownedWorkStats(ctx) + if err != nil { + return nil, err + } + views := make([]AccountMonitorView, 0, len(profiles)) + for _, profile := range profiles { + view := AccountMonitorView{AccountProfile: profile} + if stat, exists := stats[profile.ID]; exists { + view.WorkCount, view.LatestPublishedAt = stat.WorkCount, stat.LatestPublishedAt + } + views = append(views, view) + } + return views, nil +} + +type ownedWorkStat struct { + WorkCount int64 + LatestPublishedAt *time.Time +} + +// ownedWorkStats 按 source_id 聚合 owned 作品数与最近发布时间(creator_work_source 与 creator_work 已保持一致)。 +func (s *Store) ownedWorkStats(ctx context.Context) (map[string]ownedWorkStat, error) { + rows, err := s.db.QueryContext(ctx, ` + SELECT ws.source_id, COUNT(*) AS work_count, MAX(w.published_at) AS latest_published_at + FROM creator_work_source ws JOIN creator_work w ON w.id = ws.work_id + WHERE ws.source_type = 'owned' + GROUP BY ws.source_id`) + if err != nil { + return nil, databaseError(err) + } + defer rows.Close() + stats := make(map[string]ownedWorkStat) + for rows.Next() { + var stat ownedWorkStat + var sourceID string + var latest sql.NullTime + if err := rows.Scan(&sourceID, &stat.WorkCount, &latest); err != nil { + return nil, err + } + stat.LatestPublishedAt = nullableTime(latest) + stats[sourceID] = stat + } + return stats, rows.Err() +} + +// AccountCollectionStatus 自有账号采集状态(checkpoint 形态,对齐竞品的 sync_status 展示语义)。 +type AccountCollectionStatus struct { + Works AccountCheckpointStatus `json:"works"` + Comments AccountCheckpointStatus `json:"comments"` +} + +type AccountCheckpointStatus struct { + Status string `json:"status"` + LastCompletedAt *time.Time `json:"last_completed_at,omitempty"` + LastError string `json:"last_error,omitempty"` + WindowEnd *time.Time `json:"window_end,omitempty"` + NextWindowStart *time.Time `json:"next_window_start,omitempty"` + NextWindowEnd *time.Time `json:"next_window_end,omitempty"` +} + +// GetAccountCollectionStatus 返回自有账号作品/评论两条采集 checkpoint 的状态与下次采集窗口。 +func (s *Store) GetAccountCollectionStatus(ctx context.Context, accountID string) (AccountCollectionStatus, error) { + accountID = strings.TrimSpace(accountID) + if accountID == "" { + return AccountCollectionStatus{}, ErrInvalid + } + status := AccountCollectionStatus{} + for _, kind := range []struct { + name string + pointer *AccountCheckpointStatus + }{ + {name: "works", pointer: &status.Works}, + {name: "comments", pointer: &status.Comments}, + } { + var state AccountCheckpointStatus + var completed, windowEnd sql.NullTime + err := s.db.QueryRowContext(ctx, ` + SELECT status, last_completed_at, last_error, window_end + FROM creator_collection_checkpoint + WHERE source_type = $1 AND source_id = $2 AND collection_kind = $3`, + SourceOwned, accountID, kind.name).Scan(&state.Status, &completed, &state.LastError, &windowEnd) + if errors.Is(err, sql.ErrNoRows) { + // 尚未建立 checkpoint:账号还没被调度器采集过。 + *kind.pointer = AccountCheckpointStatus{Status: "pending"} + continue + } + if err != nil { + return AccountCollectionStatus{}, databaseError(err) + } + state.LastCompletedAt = nullableTime(completed) + state.WindowEnd = nullableTime(windowEnd) + if state.Status == "succeeded" && windowEnd.Valid { + // 下次采集窗口按固定间隔网格推进(与调度器 NextCollectionWindow 同口径)。 + settings, settingsErr := s.GetSettings(ctx) + if settingsErr != nil { + return AccountCollectionStatus{}, settingsErr + } + nextEnd := NextFixedRun(windowEnd.Time.UTC(), time.Now().UTC(), time.Duration(settings.NewWorkIntervalSeconds)*time.Second) + start := nextEnd.Add(-time.Duration(settings.LookbackDays) * 24 * time.Hour) + state.NextWindowStart, state.NextWindowEnd = &start, &nextEnd + } + *kind.pointer = state + } + return status, nil +} + func validateProfileUpdate(input AccountProfileUpdate) error { if input.RealNameStatus != "unknown" && input.RealNameStatus != "not_real_name" && input.RealNameStatus != "recorded" { return ErrInvalid diff --git a/internal/creator/content.go b/internal/creator/content.go index a898a13..325b02c 100644 --- a/internal/creator/content.go +++ b/internal/creator/content.go @@ -368,6 +368,9 @@ func (s *Store) UpsertWork(ctx context.Context, input WorkInput, now time.Time) if input.Likes != nil && *input.Likes < 0 || input.CommentsCount != nil && *input.CommentsCount < 0 || input.Shares != nil && *input.Shares < 0 { return Work{}, false, ErrInvalid } + if input.CollectCount != nil && *input.CollectCount < 0 || input.PlayCount != nil && *input.PlayCount < 0 { + return Work{}, false, ErrInvalid + } tx, err := s.db.BeginTx(ctx, nil) if err != nil { @@ -379,8 +382,8 @@ func (s *Store) UpsertWork(ctx context.Context, input WorkInput, now time.Time) var inserted bool err = tx.QueryRowContext(ctx, ` INSERT INTO creator_work (id, platform, work_key, source_type, source_id, author_name, title, body, - published_at, published_at_status, original_url, cover_url, raw_payload, likes, comments_count, shares, collect_count) - VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17) + published_at, published_at_status, original_url, cover_url, raw_payload, likes, comments_count, shares, collect_count, play_count) + VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17, $18) ON CONFLICT (platform, work_key) DO UPDATE SET author_name = CASE WHEN EXCLUDED.author_name = '' THEN creator_work.author_name ELSE EXCLUDED.author_name END, title = CASE WHEN EXCLUDED.title = '' THEN creator_work.title ELSE EXCLUDED.title END, @@ -394,10 +397,11 @@ func (s *Store) UpsertWork(ctx context.Context, input WorkInput, now time.Time) comments_count = COALESCE(EXCLUDED.comments_count, creator_work.comments_count), shares = COALESCE(EXCLUDED.shares, creator_work.shares), collect_count = COALESCE(EXCLUDED.collect_count, creator_work.collect_count), + play_count = COALESCE(EXCLUDED.play_count, creator_work.play_count), updated_at = now() RETURNING id, (xmax = 0)`, id, input.Platform, input.WorkKey, input.SourceType, input.SourceID, input.AuthorName, input.Title, input.Body, input.PublishedAt, status, input.OriginalURL, input.CoverURL, - nullableRawPayload(input.RawPayload), input.Likes, input.CommentsCount, input.Shares, input.CollectCount).Scan(&returnedID, &inserted) + nullableRawPayload(input.RawPayload), input.Likes, input.CommentsCount, input.Shares, input.CollectCount, input.PlayCount).Scan(&returnedID, &inserted) if err != nil { return Work{}, false, databaseError(err) } @@ -432,11 +436,11 @@ func nullableRawPayload(value string) any { func scanWork(scanner interface{ Scan(...any) error }) (Work, error) { var result Work var publishedAt, latestAt, nextAt sql.NullTime - var likes, commentsCount, shares, collectCount sql.NullInt64 + var likes, commentsCount, shares, collectCount, playCount sql.NullInt64 var rawPayload sql.NullString if err := scanner.Scan(&result.ID, &result.Platform, &result.WorkKey, &result.SourceType, &result.SourceID, &result.AuthorName, &result.Title, &result.Body, &publishedAt, &result.PublishedAtStatus, - &result.OriginalURL, &result.CoverURL, &rawPayload, &likes, &commentsCount, &shares, &collectCount, &latestAt, &nextAt, + &result.OriginalURL, &result.CoverURL, &rawPayload, &likes, &commentsCount, &shares, &collectCount, &playCount, &latestAt, &nextAt, &result.MetricStopReason, &result.CreatedAt, &result.UpdatedAt); err != nil { return Work{}, err } @@ -445,13 +449,13 @@ func scanWork(scanner interface{ Scan(...any) error }) (Work, error) { result.RawPayload = rawPayload.String } result.Likes, result.CommentsCount, result.Shares = nullableInt64(likes), nullableInt64(commentsCount), nullableInt64(shares) - result.CollectCount = nullableInt64(collectCount) + result.CollectCount, result.PlayCount = nullableInt64(collectCount), nullableInt64(playCount) result.LatestMetricsAt, result.NextMetricAt = nullableTime(latestAt), nullableTime(nextAt) return result, nil } const workSelect = `SELECT id, platform, work_key, source_type, source_id, author_name, title, body, - published_at, published_at_status, original_url, cover_url, raw_payload, likes, comments_count, shares, collect_count, + published_at, published_at_status, original_url, cover_url, raw_payload, likes, comments_count, shares, collect_count, play_count, latest_metrics_at, next_metric_at, metric_stop_reason, created_at, updated_at FROM creator_work` func (s *Store) loadWorkSources(ctx context.Context, work *Work) error { @@ -590,7 +594,7 @@ func (s *Store) ListWorks(ctx context.Context, filter WorkFilter) ([]Work, error } func (s *Store) RecordMetric(ctx context.Context, input MetricInput, settings Settings, now time.Time) (MetricPoint, error) { - if input.WorkID == "" || input.CollectedAt.IsZero() || input.Likes != nil && *input.Likes < 0 || input.CommentsCount != nil && *input.CommentsCount < 0 || input.Shares != nil && *input.Shares < 0 || input.CollectCount != nil && *input.CollectCount < 0 { + if input.WorkID == "" || input.CollectedAt.IsZero() || input.Likes != nil && *input.Likes < 0 || input.CommentsCount != nil && *input.CommentsCount < 0 || input.Shares != nil && *input.Shares < 0 || input.CollectCount != nil && *input.CollectCount < 0 || input.PlayCount != nil && *input.PlayCount < 0 { return MetricPoint{}, ErrInvalid } if err := ValidateSettings(SettingsUpdate{LookbackDays: settings.LookbackDays, NewWorkIntervalSeconds: settings.NewWorkIntervalSeconds, MetricInitialIntervalSeconds: settings.MetricInitialIntervalSeconds, MetricMultiplier: settings.MetricMultiplier, MetricMaxIntervalSeconds: settings.MetricMaxIntervalSeconds, MetricAgeSeconds: settings.MetricAgeSeconds}); err != nil { @@ -609,8 +613,54 @@ func nullableArg(value time.Time) any { return value.UTC() } +// RecordAccountMetric 落库账号画像快照。同一 (account_id, collected_at) 重复采集覆盖更新。 +func (s *Store) RecordAccountMetric(ctx context.Context, input AccountMetricInput) error { + if input.AccountID == "" || input.CollectedAt.IsZero() || input.FollowerCount != nil && *input.FollowerCount < 0 || + input.FollowingCount != nil && *input.FollowingCount < 0 || input.TotalFavorited != nil && *input.TotalFavorited < 0 || input.AwemeCount != nil && *input.AwemeCount < 0 { + return ErrInvalid + } + if _, err := s.db.ExecContext(ctx, ` + INSERT INTO creator_account_metric (account_id, collected_at, follower_count, following_count, total_favorited, aweme_count) + VALUES ($1, $2, $3, $4, $5, $6) + ON CONFLICT (account_id, collected_at) DO UPDATE SET + follower_count = EXCLUDED.follower_count, following_count = EXCLUDED.following_count, + total_favorited = EXCLUDED.total_favorited, aweme_count = EXCLUDED.aweme_count`, + input.AccountID, input.CollectedAt.UTC(), input.FollowerCount, input.FollowingCount, input.TotalFavorited, input.AwemeCount); err != nil { + return databaseError(err) + } + return nil +} + +// ListAccountMetrics 返回账号画像时序(时间升序),用于粉丝/获赞趋势曲线。 +func (s *Store) ListAccountMetrics(ctx context.Context, accountID string) ([]AccountMetricPoint, error) { + accountID = strings.TrimSpace(accountID) + if accountID == "" { + return nil, ErrInvalid + } + rows, err := s.db.QueryContext(ctx, ` + SELECT collected_at, follower_count, following_count, total_favorited, aweme_count + FROM creator_account_metric WHERE account_id = $1 ORDER BY collected_at`, accountID) + if err != nil { + return nil, databaseError(err) + } + defer rows.Close() + result := make([]AccountMetricPoint, 0) + for rows.Next() { + var point AccountMetricPoint + var follower, following, favorited, aweme sql.NullInt64 + if err := rows.Scan(&point.CollectedAt, &follower, &following, &favorited, &aweme); err != nil { + return nil, err + } + point.CollectedAt = point.CollectedAt.UTC() + point.FollowerCount, point.FollowingCount = nullableInt64(follower), nullableInt64(following) + point.TotalFavorited, point.AwemeCount = nullableInt64(favorited), nullableInt64(aweme) + result = append(result, point) + } + return result, rows.Err() +} + func (s *Store) ListMetrics(ctx context.Context, workID string) ([]MetricPoint, error) { - rows, err := s.db.QueryContext(ctx, `SELECT collected_at, likes, comments_count, shares, collect_count FROM creator_work_metric WHERE work_id = $1 ORDER BY collected_at`, workID) + rows, err := s.db.QueryContext(ctx, `SELECT collected_at, likes, comments_count, shares, collect_count, play_count FROM creator_work_metric WHERE work_id = $1 ORDER BY collected_at`, workID) if err != nil { return nil, databaseError(err) } @@ -618,13 +668,13 @@ func (s *Store) ListMetrics(ctx context.Context, workID string) ([]MetricPoint, result := make([]MetricPoint, 0) for rows.Next() { var point MetricPoint - var likes, commentsCount, shares, collectCount sql.NullInt64 - if err := rows.Scan(&point.CollectedAt, &likes, &commentsCount, &shares, &collectCount); err != nil { + var likes, commentsCount, shares, collectCount, playCount sql.NullInt64 + if err := rows.Scan(&point.CollectedAt, &likes, &commentsCount, &shares, &collectCount, &playCount); err != nil { return nil, err } point.CollectedAt = point.CollectedAt.UTC() point.Likes, point.CommentsCount, point.Shares = nullableInt64(likes), nullableInt64(commentsCount), nullableInt64(shares) - point.CollectCount = nullableInt64(collectCount) + point.CollectCount, point.PlayCount = nullableInt64(collectCount), nullableInt64(playCount) result = append(result, point) } return result, rows.Err() diff --git a/internal/creator/integration_test.go b/internal/creator/integration_test.go index 2a00bdc..426ae44 100644 --- a/internal/creator/integration_test.go +++ b/internal/creator/integration_test.go @@ -252,11 +252,27 @@ func TestCreatorPostgresContentAndWorkflow(t *testing.T) { } published := now.Add(-2 * time.Hour) likes, comments, shares := int64(10), int64(2), int64(1) + plays := int64(777) workKey := "creator-it-work-" + stamp - work, inserted, err := store.UpsertWork(ctx, WorkInput{Platform: PlatformDouyin, WorkKey: workKey, SourceType: SourceOwned, SourceID: bigID, Title: "Topic", Body: "Need consultation", PublishedAt: &published, PublishedAtStatus: "verified", OriginalURL: "https://www.douyin.com/video/" + workKey, Likes: &likes, CommentsCount: &comments, Shares: &shares}, now) + work, inserted, err := store.UpsertWork(ctx, WorkInput{Platform: PlatformDouyin, WorkKey: workKey, SourceType: SourceOwned, SourceID: bigID, Title: "Topic", Body: "Need consultation", PublishedAt: &published, PublishedAtStatus: "verified", OriginalURL: "https://www.douyin.com/video/" + workKey, Likes: &likes, CommentsCount: &comments, Shares: &shares, PlayCount: &plays}, now) if err != nil || !inserted { t.Fatalf("insert work: work=%+v inserted=%v err=%v", work, inserted, err) } + if work.PlayCount == nil || *work.PlayCount != 777 { + t.Fatalf("play_count not persisted on insert: %+v", work) + } + // COALESCE 语义:二次 upsert 未带 play_count 时保留旧值。 + if _, _, err := store.UpsertWork(ctx, WorkInput{Platform: PlatformDouyin, WorkKey: workKey, SourceType: SourceOwned, SourceID: bigID, Title: "Topic-renamed", PublishedAt: &published, PublishedAtStatus: "verified"}, now); err != nil { + t.Fatal(err) + } + retained, err := store.GetWorkByKey(ctx, PlatformDouyin, workKey) + if err != nil || retained.PlayCount == nil || *retained.PlayCount != 777 { + t.Fatalf("play_count must be retained via COALESCE: work=%+v err=%v", retained, err) + } + // 负播放量拒绝。 + if _, _, err := store.UpsertWork(ctx, WorkInput{Platform: PlatformDouyin, WorkKey: workKey + "-neg", SourceType: SourceOwned, SourceID: bigID, Title: "neg", PlayCount: &[]int64{-1}[0]}, now); !errors.Is(err, ErrInvalid) { + t.Fatalf("negative play_count must be rejected: err=%v", err) + } competitor, err := store.CreateCompetitor(ctx, CompetitorInput{Platform: PlatformDouyin, PlatformAccountKey: "sec_uid_competitor_" + stamp, UniqueID: "competitor_" + stamp, Nickname: "Competitor", HomepageURL: "https://www.douyin.com/user/sec_uid_competitor_" + stamp, Tags: []string{"重点监测"}}) if err != nil { t.Fatal(err) @@ -355,7 +371,8 @@ func TestCreatorPostgresContentAndWorkflow(t *testing.T) { t.Fatalf("comment deduplication failed: duplicate=%v err=%v", duplicate, err) } collects := int64(4) - if _, err := store.RecordMetric(ctx, MetricInput{WorkID: work.ID, CollectedAt: now, Likes: &likes, CommentsCount: &comments, Shares: &shares, CollectCount: &collects}, Settings{LookbackDays: 30, NewWorkIntervalSeconds: 1800, MetricInitialIntervalSeconds: 3600, MetricMultiplier: 2, MetricMaxIntervalSeconds: 86400, MetricAgeSeconds: 30 * 24 * 60 * 60}, now); err != nil { + playSnap := int64(5566) + if _, err := store.RecordMetric(ctx, MetricInput{WorkID: work.ID, CollectedAt: now, Likes: &likes, CommentsCount: &comments, Shares: &shares, CollectCount: &collects, PlayCount: &playSnap}, Settings{LookbackDays: 30, NewWorkIntervalSeconds: 1800, MetricInitialIntervalSeconds: 3600, MetricMultiplier: 2, MetricMaxIntervalSeconds: 86400, MetricAgeSeconds: 30 * 24 * 60 * 60}, now); err != nil { t.Fatal(err) } metrics, err := store.ListMetrics(ctx, work.ID) @@ -365,6 +382,9 @@ func TestCreatorPostgresContentAndWorkflow(t *testing.T) { if metrics[0].CollectCount == nil || *metrics[0].CollectCount != 4 { t.Fatalf("collect_count not persisted: %+v", metrics[0]) } + if metrics[0].PlayCount == nil || *metrics[0].PlayCount != 5566 { + t.Fatalf("play_count not persisted in metric point: %+v", metrics[0]) + } // 增速计算:独立作品(避免干扰上方计划推进断言),两个采样点。 growthPublished := now.Add(-2 * time.Hour) growthLikes0, growthLikes1 := int64(100), int64(400) @@ -451,6 +471,50 @@ func TestCreatorPostgresContentAndWorkflow(t *testing.T) { } } +func TestCreatorPostgresAccountMetricSnapshot(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) + follower, following, favorited, aweme := int64(1200), int64(56), int64(56789), int64(34) + if err := store.RecordAccountMetric(ctx, AccountMetricInput{AccountID: accountID, CollectedAt: now, FollowerCount: &follower, FollowingCount: &following, TotalFavorited: &favorited, AwemeCount: &aweme}); err != nil { + t.Fatal(err) + } + points, err := store.ListAccountMetrics(ctx, accountID) + if err != nil || len(points) != 1 { + t.Fatalf("list account metrics: points=%+v err=%v", points, err) + } + if points[0].FollowerCount == nil || *points[0].FollowerCount != 1200 || points[0].TotalFavorited == nil || *points[0].TotalFavorited != 56789 { + t.Fatalf("account metric snapshot: %+v", points[0]) + } + // UNIQUE(account_id, collected_at) 覆盖语义:同一时刻重采更新而非新增。 + newFollower := int64(1300) + if err := store.RecordAccountMetric(ctx, AccountMetricInput{AccountID: accountID, CollectedAt: now, FollowerCount: &newFollower}); err != nil { + t.Fatal(err) + } + points, err = store.ListAccountMetrics(ctx, accountID) + if err != nil || len(points) != 1 || points[0].FollowerCount == nil || *points[0].FollowerCount != 1300 { + t.Fatalf("account metric upsert must overwrite same timestamp: points=%+v err=%v", points, err) + } + // 不同时刻追加为独立点。 + if err := store.RecordAccountMetric(ctx, AccountMetricInput{AccountID: accountID, CollectedAt: now.Add(30 * time.Minute), FollowerCount: &newFollower}); err != nil { + t.Fatal(err) + } + points, err = store.ListAccountMetrics(ctx, accountID) + if err != nil || len(points) != 2 || !points[0].CollectedAt.Before(points[1].CollectedAt) { + t.Fatalf("account metric timeline must append: points=%+v err=%v", points, err) + } + // 负值拒绝。 + bad := int64(-1) + if err := store.RecordAccountMetric(ctx, AccountMetricInput{AccountID: accountID, CollectedAt: now, FollowerCount: &bad}); !errors.Is(err, ErrInvalid) { + t.Fatalf("negative follower_count must be rejected: err=%v", err) + } + // 账号不存在时 FK 拒绝。 + if err := store.RecordAccountMetric(ctx, AccountMetricInput{AccountID: "account-missing-" + fmt.Sprint(stamp), CollectedAt: now}); err == nil { + t.Fatal("missing account must violate FK") + } +} + func TestCreatorPostgresCollectionAllowsMutedReadAccount(t *testing.T) { store, phaseAStore, ctx := openCreatorIntegrationStore(t) stamp := time.Now().UnixNano() diff --git a/internal/creator/metrics.go b/internal/creator/metrics.go index e23949a..3e8a7ef 100644 --- a/internal/creator/metrics.go +++ b/internal/creator/metrics.go @@ -141,9 +141,9 @@ func (s *Store) recordMetricWithPlan(ctx context.Context, input MetricInput, set if err := tx.QueryRowContext(ctx, `SELECT published_at, published_at_status, latest_metrics_at FROM creator_work WHERE id=$1 FOR UPDATE`, input.WorkID).Scan(&publishedAt, &publishedStatus, &latestAt); err != nil { return MetricPoint{}, rowError(err) } - point := MetricPoint{CollectedAt: collectedAt, Likes: input.Likes, CommentsCount: input.CommentsCount, Shares: input.Shares, CollectCount: input.CollectCount} + point := MetricPoint{CollectedAt: collectedAt, Likes: input.Likes, CommentsCount: input.CommentsCount, Shares: input.Shares, CollectCount: input.CollectCount, PlayCount: input.PlayCount} if !publishedAt.Valid || publishedStatus != "verified" { - if _, err := tx.ExecContext(ctx, `UPDATE creator_work SET likes=$2, comments_count=$3, shares=$4, collect_count=$5, latest_metrics_at=$6, next_metric_at=NULL, metric_stop_reason='published_at_pending_verification', updated_at=now() WHERE id=$1`, input.WorkID, input.Likes, input.CommentsCount, input.Shares, input.CollectCount, collectedAt); err != nil { + if _, err := tx.ExecContext(ctx, `UPDATE creator_work SET likes=$2, comments_count=$3, shares=$4, collect_count=$5, play_count=$6, latest_metrics_at=$7, next_metric_at=NULL, metric_stop_reason='published_at_pending_verification', updated_at=now() WHERE id=$1`, input.WorkID, input.Likes, input.CommentsCount, input.Shares, input.CollectCount, input.PlayCount, collectedAt); err != nil { return MetricPoint{}, databaseError(err) } if err := tx.Commit(); err != nil { @@ -189,7 +189,7 @@ func (s *Store) recordMetricWithPlan(ctx context.Context, input MetricInput, set } if !latestAt.Valid { nextAt, nextReason := NextMetricAt(published, now, time.Duration(settings.MetricInitialIntervalSeconds)*time.Second, time.Duration(settings.MetricMaxIntervalSeconds)*time.Second, settings.MetricMultiplier, time.Duration(settings.MetricAgeSeconds)*time.Second) - if _, err := tx.ExecContext(ctx, `INSERT INTO creator_work_metric (work_id,collected_at,likes,comments_count,shares,collect_count) VALUES ($1,$2,$3,$4,$5,$6) ON CONFLICT (work_id,collected_at) DO UPDATE SET likes=EXCLUDED.likes,comments_count=EXCLUDED.comments_count,shares=EXCLUDED.shares,collect_count=EXCLUDED.collect_count`, input.WorkID, collectedAt, input.Likes, input.CommentsCount, input.Shares, input.CollectCount); err != nil { + if _, err := tx.ExecContext(ctx, `INSERT INTO creator_work_metric (work_id,collected_at,likes,comments_count,shares,collect_count,play_count) VALUES ($1,$2,$3,$4,$5,$6,$7) ON CONFLICT (work_id,collected_at) DO UPDATE SET likes=EXCLUDED.likes,comments_count=EXCLUDED.comments_count,shares=EXCLUDED.shares,collect_count=EXCLUDED.collect_count,play_count=EXCLUDED.play_count`, input.WorkID, collectedAt, input.Likes, input.CommentsCount, input.Shares, input.CollectCount, input.PlayCount); err != nil { return MetricPoint{}, databaseError(err) } nextInterval := time.Duration(settings.MetricMaxIntervalSeconds) * time.Second @@ -210,7 +210,7 @@ func (s *Store) recordMetricWithPlan(ctx context.Context, input MetricInput, set } else { nextReason = "" } - if _, err := tx.ExecContext(ctx, `UPDATE creator_work SET likes=$2,comments_count=$3,shares=$4,collect_count=$5,latest_metrics_at=$6,next_metric_at=$7,metric_stop_reason=$8,updated_at=now() WHERE id=$1`, input.WorkID, input.Likes, input.CommentsCount, input.Shares, input.CollectCount, collectedAt, nullableArg(nextAt), nextReason); err != nil { + if _, err := tx.ExecContext(ctx, `UPDATE creator_work SET likes=$2,comments_count=$3,shares=$4,collect_count=$5,play_count=$6,latest_metrics_at=$7,next_metric_at=$8,metric_stop_reason=$9,updated_at=now() WHERE id=$1`, input.WorkID, input.Likes, input.CommentsCount, input.Shares, input.CollectCount, input.PlayCount, collectedAt, nullableArg(nextAt), nextReason); err != nil { return MetricPoint{}, databaseError(err) } if err := tx.Commit(); err != nil { @@ -243,13 +243,13 @@ func (s *Store) recordMetricWithPlan(ctx context.Context, input MetricInput, set } else { stopReason = "" } - if _, err := tx.ExecContext(ctx, `INSERT INTO creator_work_metric (work_id,collected_at,likes,comments_count,shares,collect_count) VALUES ($1,$2,$3,$4,$5,$6) ON CONFLICT (work_id,collected_at) DO UPDATE SET likes=EXCLUDED.likes,comments_count=EXCLUDED.comments_count,shares=EXCLUDED.shares,collect_count=EXCLUDED.collect_count`, input.WorkID, collectedAt, input.Likes, input.CommentsCount, input.Shares, input.CollectCount); err != nil { + if _, err := tx.ExecContext(ctx, `INSERT INTO creator_work_metric (work_id,collected_at,likes,comments_count,shares,collect_count,play_count) VALUES ($1,$2,$3,$4,$5,$6,$7) ON CONFLICT (work_id,collected_at) DO UPDATE SET likes=EXCLUDED.likes,comments_count=EXCLUDED.comments_count,shares=EXCLUDED.shares,collect_count=EXCLUDED.collect_count,play_count=EXCLUDED.play_count`, input.WorkID, collectedAt, input.Likes, input.CommentsCount, input.Shares, input.CollectCount, input.PlayCount); err != nil { return MetricPoint{}, databaseError(err) } if _, err := tx.ExecContext(ctx, `UPDATE creator_metric_plan SET next_plan_at=$2, interval_seconds=$3, point_index=point_index+$4, stopped=$5, stop_reason=$6, updated_at=now() WHERE work_id=$1`, input.WorkID, nullableArg(nextAt), nextInterval, steps, stopped, stopReason); err != nil { return MetricPoint{}, databaseError(err) } - if _, err := tx.ExecContext(ctx, `UPDATE creator_work SET likes=$2,comments_count=$3,shares=$4,collect_count=$5,latest_metrics_at=$6,next_metric_at=$7,metric_stop_reason=$8,updated_at=now() WHERE id=$1`, input.WorkID, input.Likes, input.CommentsCount, input.Shares, input.CollectCount, collectedAt, nullableArg(nextAt), stopReason); err != nil { + if _, err := tx.ExecContext(ctx, `UPDATE creator_work SET likes=$2,comments_count=$3,shares=$4,collect_count=$5,play_count=$6,latest_metrics_at=$7,next_metric_at=$8,metric_stop_reason=$9,updated_at=now() WHERE id=$1`, input.WorkID, input.Likes, input.CommentsCount, input.Shares, input.CollectCount, input.PlayCount, collectedAt, nullableArg(nextAt), stopReason); err != nil { return MetricPoint{}, databaseError(err) } if err := tx.Commit(); err != nil { diff --git a/internal/creator/migrations/040_work_play_count.sql b/internal/creator/migrations/040_work_play_count.sql new file mode 100644 index 0000000..2f4cb99 --- /dev/null +++ b/internal/creator/migrations/040_work_play_count.sql @@ -0,0 +1,3 @@ +-- 播放量指标入库:作品列表快照与指标时序各加 play_count(2026-09-27 自有账号监控设计定稿)。 +ALTER TABLE creator_work ADD COLUMN IF NOT EXISTS play_count bigint; +ALTER TABLE creator_work_metric ADD COLUMN IF NOT EXISTS play_count bigint; diff --git a/internal/creator/migrations/041_creator_account_metric.sql b/internal/creator/migrations/041_creator_account_metric.sql new file mode 100644 index 0000000..9709f5a --- /dev/null +++ b/internal/creator/migrations/041_creator_account_metric.sql @@ -0,0 +1,14 @@ +-- 自有账号监控:账号级指标时序快照(粉丝/关注/获赞总数/作品总数)。 +-- 采集时机:syncCreatorOwned 成功路径末尾顺带拉 self profile,与作品采集同节奏。 +CREATE TABLE IF NOT EXISTS creator_account_metric ( + id bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY, + account_id text NOT NULL REFERENCES social_account(id) ON DELETE CASCADE, + collected_at timestamptz NOT NULL DEFAULT now(), + follower_count bigint, + following_count bigint, + total_favorited bigint, + aweme_count bigint, + UNIQUE (account_id, collected_at) +); +CREATE INDEX IF NOT EXISTS creator_account_metric_account_idx + ON creator_account_metric (account_id, collected_at); diff --git a/internal/creator/models.go b/internal/creator/models.go index 164053f..311278f 100644 --- a/internal/creator/models.go +++ b/internal/creator/models.go @@ -193,6 +193,7 @@ type Work struct { CommentsCount *int64 `json:"comments_count"` Shares *int64 `json:"shares"` CollectCount *int64 `json:"collect_count"` + PlayCount *int64 `json:"play_count"` LikesGrowth *int64 `json:"likes_growth,omitempty"` LikesGrowthCoverage string `json:"likes_growth_coverage,omitempty"` // full=窗口基线完整 partial=以首点近似 LatestMetricsAt *time.Time `json:"latest_metrics_at,omitempty"` @@ -220,6 +221,7 @@ type WorkInput struct { CommentsCount *int64 `json:"comments_count"` Shares *int64 `json:"shares"` CollectCount *int64 `json:"collect_count"` + PlayCount *int64 `json:"play_count"` RawPayload string `json:"-"` } @@ -243,6 +245,7 @@ type MetricInput struct { CommentsCount *int64 Shares *int64 CollectCount *int64 + PlayCount *int64 } type MetricPoint struct { @@ -251,6 +254,25 @@ type MetricPoint struct { CommentsCount *int64 `json:"comments_count"` Shares *int64 `json:"shares"` CollectCount *int64 `json:"collect_count"` + PlayCount *int64 `json:"play_count"` +} + +// AccountMetricInput 自有账号画像快照(粉丝/关注/获赞/作品总数),随作品采集同节奏落库。 +type AccountMetricInput struct { + AccountID string + CollectedAt time.Time + FollowerCount *int64 + FollowingCount *int64 + TotalFavorited *int64 + AwemeCount *int64 +} + +type AccountMetricPoint struct { + CollectedAt time.Time `json:"collected_at"` + FollowerCount *int64 `json:"follower_count"` + FollowingCount *int64 `json:"following_count"` + TotalFavorited *int64 `json:"total_favorited"` + AwemeCount *int64 `json:"aweme_count"` } type MaterialJob struct { diff --git a/internal/creator/store.go b/internal/creator/store.go index 2875e0a..09db35a 100644 --- a/internal/creator/store.go +++ b/internal/creator/store.go @@ -89,6 +89,12 @@ var migration038 string //go:embed migrations/039_work_collect_count.sql var migration039 string +//go:embed migrations/040_work_play_count.sql +var migration040 string + +//go:embed migrations/041_creator_account_metric.sql +var migration041 string + type SecretReference struct { ID string Provider string @@ -189,6 +195,8 @@ func (s *Store) migrate(ctx context.Context) error { {version: 37, sql: migration037}, {version: 38, sql: migration038}, {version: 39, sql: migration039}, + {version: 40, sql: migration040}, + {version: 41, sql: migration041}, } for _, migration := range migrations { var applied bool diff --git a/internal/platform/douyin/connector.go b/internal/platform/douyin/connector.go index b0e2899..2a472c2 100644 --- a/internal/platform/douyin/connector.go +++ b/internal/platform/douyin/connector.go @@ -5,6 +5,7 @@ import ( "context" "encoding/json" "errors" + "fmt" "io" "net/http" "net/url" @@ -315,9 +316,39 @@ type identityEnvelope struct { AvatarThumb *struct { URLList []string `json:"url_list"` } `json:"avatar_thumb"` + FollowerCount *int64 `json:"follower_count"` + FollowingCount *int64 `json:"following_count"` + AwemeCount *int64 `json:"aweme_count"` + TotalFavorited *flexibleInt `json:"total_favorited"` } `json:"user"` } +// flexibleInt 兼容抖音返回的数字与字符串两种计数形态(total_favorited 实测为大数字符串)。 +type flexibleInt int64 + +func (value *flexibleInt) UnmarshalJSON(data []byte) error { + text := strings.TrimSpace(string(data)) + if text == "null" { + return nil + } + if len(text) >= 2 && text[0] == '"' && text[len(text)-1] == '"' { + unquoted, err := strconv.Unquote(text) + if err != nil { + return errors.New("douyin flexible int string invalid") + } + text = strings.TrimSpace(unquoted) + } + if text == "" { + return nil + } + parsed, err := strconv.ParseInt(text, 10, 64) + if err != nil || parsed < 0 { + return errors.New("douyin flexible int invalid") + } + *value = flexibleInt(parsed) + return nil +} + func parseIdentity(body []byte) (identityEnvelope, bool) { var identity identityEnvelope if len(body) > 1<<20 || json.Unmarshal(body, &identity) != nil || identity.StatusCode == nil || *identity.StatusCode != 0 || identity.User == nil || @@ -325,9 +356,66 @@ func parseIdentity(body []byte) (identityEnvelope, bool) { (identity.User.UniqueID != "" && !keyPattern.MatchString(identity.User.UniqueID)) { return identityEnvelope{}, false } + for _, value := range []*int64{identity.User.FollowerCount, identity.User.FollowingCount, identity.User.AwemeCount} { + if value != nil && *value < 0 { + return identityEnvelope{}, false + } + } return identity, true } +// SelfProfile 是自有账号自己的画像快照(self profile 接口),用于账号级指标时序采集。 +type SelfProfile struct { + UID string + SecUID string + UniqueID string + Nickname string + AvatarURL string + FollowerCount *int64 + FollowingCount *int64 + TotalFavorited *int64 + AwemeCount *int64 +} + +// FetchSelfProfile 经网关受限 fetch 拉取登录账号自己的画像(复用 self profile 身份接口,返回体自带粉丝/关注/获赞/作品数)。 +func FetchSelfProfile(ctx context.Context, browser Browser) (SelfProfile, error) { + if browser == nil { + return SelfProfile{}, fmt.Errorf("%w: invalid self profile request", ErrInvalid) + } + response, err := browser.Get(ctx, identityEndpoint) + if err != nil { + return SelfProfile{}, err + } + if err := creatorResponseError(response, "self profile"); err != nil { + return SelfProfile{}, err + } + profile, ok := parseSelfProfile(response.Body) + if !ok { + return SelfProfile{}, fmt.Errorf("%w: invalid douyin self profile response (status=%d body_len=%d head=%.160s)", ErrInvalid, response.Status, len(response.Body), response.Body) + } + return profile, nil +} + +func parseSelfProfile(body []byte) (SelfProfile, bool) { + identity, ok := parseIdentity(body) + if !ok { + return SelfProfile{}, false + } + avatarURL := "" + if identity.User.AvatarThumb != nil && len(identity.User.AvatarThumb.URLList) > 0 { + avatarURL = identity.User.AvatarThumb.URLList[0] + } + var totalFavorited *int64 + if identity.User.TotalFavorited != nil { + value := int64(*identity.User.TotalFavorited) + totalFavorited = &value + } + return SelfProfile{ + UID: identity.User.UID, SecUID: identity.User.SecUID, UniqueID: identity.User.UniqueID, Nickname: identity.User.Nickname, AvatarURL: avatarURL, + FollowerCount: identity.User.FollowerCount, FollowingCount: identity.User.FollowingCount, TotalFavorited: totalFavorited, AwemeCount: identity.User.AwemeCount, + }, true +} + type douyinBool bool func (value *douyinBool) UnmarshalJSON(data []byte) error { diff --git a/internal/platform/douyin/creator_collector.go b/internal/platform/douyin/creator_collector.go index 79d38f5..6db3c21 100644 --- a/internal/platform/douyin/creator_collector.go +++ b/internal/platform/douyin/creator_collector.go @@ -40,6 +40,7 @@ type WorkStatistics struct { CommentsCount *int64 Shares *int64 CollectCount *int64 + PlayCount *int64 } // FetchWorkStatistics 经网关受限 fetch 拉取作品详情并解析指标快照。 @@ -71,15 +72,15 @@ func parseWorkStatistics(body []byte, expectedWorkKey string) (WorkStatistics, b return WorkStatistics{}, false } stats := envelope.AwemeDetail.Statistics - for _, value := range []*int64{stats.DiggCount, stats.CommentCount, stats.ShareCount, stats.CollectCount} { + for _, value := range []*int64{stats.DiggCount, stats.CommentCount, stats.ShareCount, stats.CollectCount, stats.PlayCount} { if value != nil && *value < 0 { return WorkStatistics{}, false } } - if stats.DiggCount == nil && stats.CommentCount == nil && stats.ShareCount == nil && stats.CollectCount == nil { + if stats.DiggCount == nil && stats.CommentCount == nil && stats.ShareCount == nil && stats.CollectCount == nil && stats.PlayCount == nil { return WorkStatistics{}, false } - return WorkStatistics{Likes: stats.DiggCount, CommentsCount: stats.CommentCount, Shares: stats.ShareCount, CollectCount: stats.CollectCount}, true + return WorkStatistics{Likes: stats.DiggCount, CommentsCount: stats.CommentCount, Shares: stats.ShareCount, CollectCount: stats.CollectCount, PlayCount: stats.PlayCount}, true } type workDetailEnvelope struct { @@ -100,6 +101,7 @@ type workDetailEnvelope struct { CommentCount *int64 `json:"comment_count"` ShareCount *int64 `json:"share_count"` CollectCount *int64 `json:"collect_count"` + PlayCount *int64 `json:"play_count"` } `json:"statistics"` } `json:"aweme_detail"` } @@ -247,7 +249,7 @@ func (c CreatorCollector) ListWorks(ctx context.Context, accountKey, cursor stri value := time.Unix(*work.CreatedAt, 0).UTC() published = &value } - likes, comments, shares, collects := work.DiggCount, work.CommentCount, work.ShareCount, work.CollectCount + likes, comments, shares, collects, plays := work.DiggCount, work.CommentCount, work.ShareCount, work.CollectCount, work.PlayCount sourceType, sourceID := c.SourceType, c.SourceID if sourceType == "" { sourceType = creator.SourceCompetitor @@ -261,7 +263,7 @@ func (c CreatorCollector) ListWorks(ctx context.Context, accountKey, cursor stri } else if published != nil { status = "verified" } - items = append(items, creator.WorkInput{Platform: creator.PlatformDouyin, WorkKey: work.ID, SourceType: sourceType, SourceID: sourceID, Body: work.Description, PublishedAt: published, PublishedAtStatus: status, OriginalURL: "https://www.douyin.com/video/" + work.ID, Likes: likes, CommentsCount: comments, Shares: shares, CollectCount: collects}) + items = append(items, creator.WorkInput{Platform: creator.PlatformDouyin, WorkKey: work.ID, SourceType: sourceType, SourceID: sourceID, Body: work.Description, PublishedAt: published, PublishedAtStatus: status, OriginalURL: "https://www.douyin.com/video/" + work.ID, Likes: likes, CommentsCount: comments, Shares: shares, CollectCount: collects, PlayCount: plays}) } page := creator.WorkPage{Items: items, HasMore: hasMore} if nextCursor != nil { @@ -343,6 +345,7 @@ type creatorWorkPageItem struct { CommentCount *int64 ShareCount *int64 CollectCount *int64 + PlayCount *int64 } func parseCreatorWorksPage(body []byte) ([]creatorWorkPageItem, bool, *int64, bool) { @@ -369,11 +372,11 @@ func parseCreatorWorksPage(body []byte) ([]creatorWorkPageItem, bool, *int64, bo return nil, false, nil, false } seen[item.ID] = struct{}{} - var likes, comments, shares, collects *int64 + var likes, comments, shares, collects, plays *int64 if item.Statistics != nil { - likes, comments, shares, collects = item.Statistics.DiggCount, item.Statistics.CommentCount, item.Statistics.ShareCount, item.Statistics.CollectCount + likes, comments, shares, collects, plays = item.Statistics.DiggCount, item.Statistics.CommentCount, item.Statistics.ShareCount, item.Statistics.CollectCount, item.Statistics.PlayCount } - for _, value := range []*int64{likes, comments, shares, collects} { + for _, value := range []*int64{likes, comments, shares, collects, plays} { if value != nil && *value < 0 { return nil, false, nil, false } @@ -383,7 +386,7 @@ func parseCreatorWorksPage(body []byte) ([]creatorWorkPageItem, bool, *int64, bo if createdAtInvalid { createdAt = nil } - items = append(items, creatorWorkPageItem{ID: item.ID, Description: item.Description, CreatedAt: createdAt, CreatedAtInvalid: createdAtInvalid, DiggCount: likes, CommentCount: comments, ShareCount: shares, CollectCount: collects}) + items = append(items, creatorWorkPageItem{ID: item.ID, Description: item.Description, CreatedAt: createdAt, CreatedAtInvalid: createdAtInvalid, DiggCount: likes, CommentCount: comments, ShareCount: shares, CollectCount: collects, PlayCount: plays}) } // 仅当响应带有效 max_cursor 时才声明翻页;无 cursor 的 has_more 不产生翻页游标。 hasMore := envelope.HasMore != nil && bool(*envelope.HasMore) && envelope.MaxCursor != nil diff --git a/internal/platform/douyin/creator_collector_test.go b/internal/platform/douyin/creator_collector_test.go index 215685d..72b4bba 100644 --- a/internal/platform/douyin/creator_collector_test.go +++ b/internal/platform/douyin/creator_collector_test.go @@ -137,7 +137,7 @@ func TestParseCreatorWorksPageAcceptsNumericHasMore(t *testing.T) { func TestParseCreatorWorksPageToleratesHasMoreWithoutCursor(t *testing.T) { // 匿名 works 响应实测:has_more=1 且不带 max_cursor 字段;第一页必须可用。 - body := []byte(`{"status_code":0,"has_more":1,"aweme_list":[{"aweme_id":"123","desc":"first page","create_time":1700000000,"statistics":{"digg_count":1,"comment_count":2,"share_count":3,"collect_count":4}}]}`) + body := []byte(`{"status_code":0,"has_more":1,"aweme_list":[{"aweme_id":"123","desc":"first page","create_time":1700000000,"statistics":{"digg_count":1,"comment_count":2,"share_count":3,"collect_count":4,"play_count":55}}]}`) works, hasMore, cursor, ok := parseCreatorWorksPage(body) if !ok || hasMore || cursor != nil || len(works) != 1 || works[0].ID != "123" { t.Fatalf("has_more without cursor must still yield first page: ok=%v hasMore=%v cursor=%v works=%+v", ok, hasMore, cursor, works) @@ -145,13 +145,16 @@ func TestParseCreatorWorksPageToleratesHasMoreWithoutCursor(t *testing.T) { if works[0].CollectCount == nil || *works[0].CollectCount != 4 { t.Fatalf("collect_count not parsed: %+v", works[0]) } + if works[0].PlayCount == nil || *works[0].PlayCount != 55 { + t.Fatalf("play_count not parsed: %+v", works[0]) + } } func TestParseWorkStatisticsExtractsMetrics(t *testing.T) { - // 匿名 detail 响应实测:statistics 携带完整指标快照。 - body := []byte(`{"status_code":0,"aweme_detail":{"aweme_id":"123","statistics":{"digg_count":6155,"comment_count":235,"share_count":578,"collect_count":2292}}}`) + // 匿名 detail 响应实测:statistics 携带完整指标快照(含播放量)。 + body := []byte(`{"status_code":0,"aweme_detail":{"aweme_id":"123","statistics":{"digg_count":6155,"comment_count":235,"share_count":578,"collect_count":2292,"play_count":991234}}}`) stats, ok := parseWorkStatistics(body, "123") - if !ok || stats.Likes == nil || *stats.Likes != 6155 || stats.CommentsCount == nil || *stats.CommentsCount != 235 || stats.Shares == nil || *stats.Shares != 578 || stats.CollectCount == nil || *stats.CollectCount != 2292 { + if !ok || stats.Likes == nil || *stats.Likes != 6155 || stats.CommentsCount == nil || *stats.CommentsCount != 235 || stats.Shares == nil || *stats.Shares != 578 || stats.CollectCount == nil || *stats.CollectCount != 2292 || stats.PlayCount == nil || *stats.PlayCount != 991234 { t.Fatalf("unexpected statistics: ok=%v %+v", ok, stats) } if _, ok := parseWorkStatistics(body, "999"); ok { @@ -160,6 +163,15 @@ func TestParseWorkStatisticsExtractsMetrics(t *testing.T) { if _, ok := parseWorkStatistics([]byte(`{"status_code":0,"aweme_detail":{"aweme_id":"123"}}`), "123"); ok { t.Fatal("accepted missing statistics") } + // 仅播放量可用也视为有效快照(其它指标可能缺失)。 + playOnly := []byte(`{"status_code":0,"aweme_detail":{"aweme_id":"123","statistics":{"play_count":88}}}`) + playStats, ok := parseWorkStatistics(playOnly, "123") + if !ok || playStats.PlayCount == nil || *playStats.PlayCount != 88 { + t.Fatalf("play-only statistics rejected: ok=%v %+v", ok, playStats) + } + if _, ok := parseWorkStatistics([]byte(`{"status_code":0,"aweme_detail":{"aweme_id":"123","statistics":{"play_count":-1}}}`), "123"); ok { + t.Fatal("accepted negative play_count") + } } func TestFetchWorkStatisticsUsesDetailEndpoint(t *testing.T) { @@ -174,6 +186,36 @@ func TestFetchWorkStatisticsUsesDetailEndpoint(t *testing.T) { } } +func TestFetchSelfProfileExtractsAccountMetrics(t *testing.T) { + // self profile 返回体自带账号级计数;total_favorited 实测为大数字符串。 + browser := &collectorBrowser{response: Response{Status: 200, Body: []byte(`{"status_code":0,"user":{"uid":"2328120603967913","sec_uid":"MS4wLjABAAAA9f_a7k0bzVizLYXlpC7R61EIaqJ8Ordug7yp7AB8fGKuuF8Fzqk5_DM-eutXnPIK","unique_id":"96332518739","nickname":"self","follower_count":1200,"following_count":56,"aweme_count":34,"total_favorited":"56789"}}`)}} + profile, err := FetchSelfProfile(context.Background(), browser) + if err != nil { + t.Fatalf("fetch self profile: %v", err) + } + if profile.FollowerCount == nil || *profile.FollowerCount != 1200 || profile.FollowingCount == nil || *profile.FollowingCount != 56 || + profile.AwemeCount == nil || *profile.AwemeCount != 34 || profile.TotalFavorited == nil || *profile.TotalFavorited != 56789 { + t.Fatalf("unexpected self profile: %+v", profile) + } + parsed, err := url.Parse(browser.url) + if err != nil || parsed.Path != "/aweme/v1/web/user/profile/self/" { + t.Fatalf("unexpected request url: %s", browser.url) + } +} + +func TestParseSelfProfileRejectsInvalidCounts(t *testing.T) { + secUID := "MS4wLjABAAAA9f_a7k0bzVizLYXlpC7R61EIaqJ8Ordug7yp7AB8fGKuuF8Fzqk5_DM-eutXnPIK" + if _, ok := parseSelfProfile([]byte(`{"status_code":0,"user":{"uid":"2328120603967913","sec_uid":"` + secUID + `","follower_count":-5}}`)); ok { + t.Fatal("accepted negative follower_count") + } + if _, ok := parseSelfProfile([]byte(`{"status_code":0,"user":{"uid":"2328120603967913","sec_uid":"` + secUID + `","total_favorited":"abc"}}`)); ok { + t.Fatal("accepted invalid total_favorited") + } + if _, ok := parseSelfProfile([]byte(`{"status_code":1,"user":{}}`)); ok { + t.Fatal("accepted error status") + } +} + func TestParseCreatorWorksPageTreatsBareEnvelopeAsEmptyPage(t *testing.T) { // 平台风控语义(TikTokDownloader #706 同样现象):翻页被拦截时返回空 // envelope {"status_code":0},无 aweme_list/has_more/max_cursor。视为正常空页。