feat(creator): 自有账号监控后端——播放量入库、账号粉丝时序快照、监控 API

- 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 语义
This commit is contained in:
2026-09-27 22:05:01 +08:00
parent 01d1d44a15
commit 3f1e2d39f4
15 changed files with 593 additions and 38 deletions
@@ -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 {
+58 -4
View File
@@ -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) {
@@ -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 {
@@ -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 {
+119
View File
@@ -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
+62 -12
View File
@@ -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()
+66 -2
View File
@@ -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()
+6 -6
View File
@@ -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 {
@@ -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;
@@ -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);
+22
View File
@@ -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 {
+8
View File
@@ -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
+88
View File
@@ -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 {
+12 -9
View File
@@ -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
@@ -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。视为正常空页。