- 指标采集不再依赖已登录自有账号:竞品作品到期后经匿名浏览器拉取 抖音作品详情接口(detail 返回完整点赞/收藏/转发/评论快照) - 补采收藏数:creator_work 与 creator_work_metric 新增 collect_count (迁移 039),works 列表与 detail 解析均映射该字段 - 新增点赞增速:LoadWorksGrowth 按时间窗计算基线增量(full/partial 覆盖标注),works API 支持 growth_hours/min_likes_growth 筛选排序 - 前端监控详情页:作品卡片展示 24h 点赞增量、支持按增量筛选爆品; 作品 Drawer 新增增长曲线折线图(@ant-design/charts)与收藏列
This commit is contained in:
@@ -1017,6 +1017,21 @@ func workFilter(c fiber.Ctx) (creator.WorkFilter, error) {
|
||||
}
|
||||
*field.target = &parsed
|
||||
}
|
||||
// 爆品筛选:growth_hours 开启点赞增量计算(默认 24h),min_likes_growth 设增量下限。
|
||||
if value := c.Query("growth_hours"); value != "" {
|
||||
parsed, err := strconv.ParseInt(value, 10, 64)
|
||||
if err != nil {
|
||||
return creator.WorkFilter{}, creator.ErrInvalid
|
||||
}
|
||||
filter.Growth.Hours = parsed
|
||||
}
|
||||
if value := c.Query("min_likes_growth"); value != "" {
|
||||
parsed, err := strconv.ParseInt(value, 10, 64)
|
||||
if err != nil {
|
||||
return creator.WorkFilter{}, creator.ErrInvalid
|
||||
}
|
||||
filter.Growth.MinLikes = &parsed
|
||||
}
|
||||
return filter, nil
|
||||
}
|
||||
|
||||
@@ -1962,21 +1977,50 @@ func runCreatorMetricScheduleOnce(ctx context.Context, store *creator.Store, pha
|
||||
return err
|
||||
}
|
||||
for _, work := range works {
|
||||
accountID := work.SourceID
|
||||
var refreshErr error
|
||||
if work.SourceType == creator.SourceCompetitor {
|
||||
accountID, err = creatorCollectionAccount(ctx, store, phaseAStore, hubStore, work.Platform)
|
||||
if err != nil {
|
||||
logrus.WithError(err).WithField("work_id", work.ID).Warn("creator metric refresh account unavailable")
|
||||
continue
|
||||
}
|
||||
// 竞品作品指标是公开数据:匿名浏览器拉详情接口快照,不依赖已登录自有账号。
|
||||
refreshErr = refreshCreatorCompetitorMetricAnonymous(ctx, store, hubStore, work, settings, now)
|
||||
} else {
|
||||
accountID := work.SourceID
|
||||
refreshErr = refreshCreatorMetricWork(ctx, store, phaseAStore, hubStore, work, accountID, settings, now)
|
||||
}
|
||||
if err := refreshCreatorMetricWork(ctx, store, phaseAStore, hubStore, work, accountID, settings, now); err != nil {
|
||||
logrus.WithError(err).WithField("work_id", work.ID).Warn("creator metric refresh failed")
|
||||
if refreshErr != nil {
|
||||
logrus.WithError(refreshErr).WithField("work_id", work.ID).Warn("creator metric refresh failed")
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func refreshCreatorCompetitorMetricAnonymous(ctx context.Context, store *creator.Store, hubStore *hub.Store, work creator.Work, settings creator.Settings, now time.Time) error {
|
||||
competitor, err := store.GetCompetitor(ctx, work.SourceID)
|
||||
if err != nil {
|
||||
if errors.Is(err, creator.ErrNotFound) {
|
||||
return store.StopMetricPlan(ctx, work.ID, "source_deleted")
|
||||
}
|
||||
return err
|
||||
}
|
||||
if !competitor.Enabled {
|
||||
return store.StopMetricPlan(ctx, work.ID, "source_disabled")
|
||||
}
|
||||
if competitor.Platform != work.Platform || competitor.Platform != creator.PlatformDouyin {
|
||||
return creator.ErrConflict
|
||||
}
|
||||
lease, err := newAnonymousBrowser(ctx, hubStore)
|
||||
if err != nil {
|
||||
return fmt.Errorf("%w: anonymous browser unavailable: %v", creator.ErrUnavailable, err)
|
||||
}
|
||||
stats, fetchErr := douyin.FetchWorkStatistics(ctx, creatorGatewayBrowser{gateway: lease.gateway, environment: lease.environment}, work.WorkKey)
|
||||
if closeErr := lease.close(); closeErr != nil {
|
||||
fetchErr = errors.Join(fetchErr, closeErr)
|
||||
}
|
||||
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)
|
||||
return recordErr
|
||||
}
|
||||
|
||||
func refreshCreatorMetricWork(ctx context.Context, store *creator.Store, phaseAStore *accountdomain.Store, hubStore *hub.Store, work creator.Work, accountID string, settings creator.Settings, now time.Time) (resultErr error) {
|
||||
if work.SourceType == creator.SourceCompetitor {
|
||||
competitor, err := store.GetCompetitor(ctx, work.SourceID)
|
||||
|
||||
+30
-10
@@ -379,8 +379,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)
|
||||
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16)
|
||||
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)
|
||||
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,
|
||||
@@ -390,10 +390,14 @@ func (s *Store) UpsertWork(ctx context.Context, input WorkInput, now time.Time)
|
||||
original_url = CASE WHEN EXCLUDED.original_url = '' THEN creator_work.original_url ELSE EXCLUDED.original_url END,
|
||||
cover_url = CASE WHEN EXCLUDED.cover_url = '' THEN creator_work.cover_url ELSE EXCLUDED.cover_url END,
|
||||
raw_payload = COALESCE(EXCLUDED.raw_payload, creator_work.raw_payload),
|
||||
likes = COALESCE(EXCLUDED.likes, creator_work.likes),
|
||||
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),
|
||||
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).Scan(&returnedID, &inserted)
|
||||
nullableRawPayload(input.RawPayload), input.Likes, input.CommentsCount, input.Shares, input.CollectCount).Scan(&returnedID, &inserted)
|
||||
if err != nil {
|
||||
return Work{}, false, databaseError(err)
|
||||
}
|
||||
@@ -428,11 +432,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 sql.NullInt64
|
||||
var likes, commentsCount, shares, collectCount 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, &latestAt, &nextAt,
|
||||
&result.OriginalURL, &result.CoverURL, &rawPayload, &likes, &commentsCount, &shares, &collectCount, &latestAt, &nextAt,
|
||||
&result.MetricStopReason, &result.CreatedAt, &result.UpdatedAt); err != nil {
|
||||
return Work{}, err
|
||||
}
|
||||
@@ -441,12 +445,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.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,
|
||||
published_at, published_at_status, original_url, cover_url, raw_payload, likes, comments_count, shares, collect_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 {
|
||||
@@ -585,7 +590,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 {
|
||||
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 {
|
||||
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 {
|
||||
@@ -605,7 +610,7 @@ func nullableArg(value time.Time) any {
|
||||
}
|
||||
|
||||
func (s *Store) ListMetrics(ctx context.Context, workID string) ([]MetricPoint, error) {
|
||||
rows, err := s.db.QueryContext(ctx, `SELECT collected_at, likes, comments_count, shares 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 FROM creator_work_metric WHERE work_id = $1 ORDER BY collected_at`, workID)
|
||||
if err != nil {
|
||||
return nil, databaseError(err)
|
||||
}
|
||||
@@ -613,12 +618,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 sql.NullInt64
|
||||
if err := rows.Scan(&point.CollectedAt, &likes, &commentsCount, &shares); err != nil {
|
||||
var likes, commentsCount, shares, collectCount sql.NullInt64
|
||||
if err := rows.Scan(&point.CollectedAt, &likes, &commentsCount, &shares, &collectCount); 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)
|
||||
result = append(result, point)
|
||||
}
|
||||
return result, rows.Err()
|
||||
@@ -977,6 +983,20 @@ func (s *Store) ListWorksPage(ctx context.Context, filter WorkFilter, page, page
|
||||
if err != nil {
|
||||
return Page[Work]{}, err
|
||||
}
|
||||
if filter.Growth.Hours != 0 || filter.Growth.MinLikes != nil {
|
||||
ids := make([]string, 0, len(items))
|
||||
for _, item := range items {
|
||||
ids = append(ids, item.ID)
|
||||
}
|
||||
growth, growthErr := s.LoadWorksGrowth(ctx, ids, filter.Growth)
|
||||
if growthErr != nil {
|
||||
return Page[Work]{}, growthErr
|
||||
}
|
||||
items, err = ApplyWorkGrowthFilter(items, growth, filter.Growth)
|
||||
if err != nil {
|
||||
return Page[Work]{}, err
|
||||
}
|
||||
}
|
||||
return slicePage(items, page, pageSize)
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,162 @@
|
||||
package creator
|
||||
|
||||
import (
|
||||
"context"
|
||||
"sort"
|
||||
"time"
|
||||
)
|
||||
|
||||
// WorkGrowth 是作品在观察窗口内的点赞增量,用于爆品筛选。
|
||||
// Baseline 为窗口起点前最后一个采样点;无基线时以最早采样点近似。
|
||||
type WorkGrowth struct {
|
||||
Likes *int64 `json:"likes_growth"`
|
||||
Coverage string `json:"likes_growth_coverage"` // full=有窗口起点基线;partial=以首点近似
|
||||
}
|
||||
|
||||
// WorkGrowthFilter 增速筛选条件。
|
||||
type WorkGrowthFilter struct {
|
||||
Hours int64 // 观察窗口长度,默认 24
|
||||
MinLikes *int64
|
||||
}
|
||||
|
||||
// normalizeGrowthFilter 校验并补默认值。
|
||||
func normalizeGrowthFilter(filter WorkGrowthFilter) (WorkGrowthFilter, error) {
|
||||
if filter.Hours == 0 {
|
||||
filter.Hours = 24
|
||||
}
|
||||
if filter.Hours < 0 || filter.Hours > 24*30 {
|
||||
return WorkGrowthFilter{}, ErrInvalid
|
||||
}
|
||||
if filter.MinLikes != nil && *filter.MinLikes < 0 {
|
||||
return WorkGrowthFilter{}, ErrInvalid
|
||||
}
|
||||
return filter, nil
|
||||
}
|
||||
|
||||
// computeWorkGrowth 从按时间升序的采样点计算观察窗口内的点赞增量。
|
||||
// points 为空返回零值(无数据);不足两个点或无基线时 coverage=partial。
|
||||
func computeWorkGrowth(points []MetricPoint, now time.Time, hours int64) WorkGrowth {
|
||||
// hours 为小时数(接口参数同名同义),窗口起点 = now - hours 小时。
|
||||
cutoff := now.Add(-time.Duration(hours) * time.Hour)
|
||||
var baseline, latest *MetricPoint
|
||||
for index := range points {
|
||||
point := &points[index]
|
||||
if point.Likes == nil {
|
||||
continue
|
||||
}
|
||||
if !point.CollectedAt.Before(cutoff) {
|
||||
// 窗口起点及之后的点属于窗口内;升序遍历取最晚点作为 latest。
|
||||
latest = point
|
||||
continue
|
||||
}
|
||||
// 窗口起点前最近点作为基线(升序遍历,最后赋值即最近)。
|
||||
baseline = point
|
||||
}
|
||||
if latest == nil {
|
||||
// 窗口内无采样点:以最后一个窗口外点近似(数据稀疏时的降级)。
|
||||
for index := range points {
|
||||
point := &points[index]
|
||||
if point.Likes == nil {
|
||||
continue
|
||||
}
|
||||
latest = point
|
||||
}
|
||||
if latest == nil {
|
||||
return WorkGrowth{}
|
||||
}
|
||||
}
|
||||
if baseline == nil {
|
||||
// 无窗口起点前基线:以窗口内最早采样点近似基线(partial)。
|
||||
for index := range points {
|
||||
point := &points[index]
|
||||
if point.Likes == nil {
|
||||
continue
|
||||
}
|
||||
if !point.CollectedAt.Before(cutoff) {
|
||||
baseline = point
|
||||
break
|
||||
}
|
||||
}
|
||||
if baseline == nil || baseline == latest {
|
||||
// 仅有单点:无法计算增量。
|
||||
return WorkGrowth{Likes: int64Ptr(0), Coverage: "partial"}
|
||||
}
|
||||
return WorkGrowth{Likes: int64Ptr(*latest.Likes - *baseline.Likes), Coverage: "partial"}
|
||||
}
|
||||
if baseline == latest {
|
||||
return WorkGrowth{Likes: int64Ptr(0), Coverage: "partial"}
|
||||
}
|
||||
return WorkGrowth{Likes: int64Ptr(*latest.Likes - *baseline.Likes), Coverage: "full"}
|
||||
}
|
||||
|
||||
// LoadWorksGrowth 批量加载作品的点赞增量。返回 map[workID]WorkGrowth;无采样点的作品不在 map 中。
|
||||
func (s *Store) LoadWorksGrowth(ctx context.Context, workIDs []string, filter WorkGrowthFilter) (map[string]WorkGrowth, error) {
|
||||
filter, err := normalizeGrowthFilter(filter)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if len(workIDs) == 0 {
|
||||
return map[string]WorkGrowth{}, nil
|
||||
}
|
||||
rows, err := s.db.QueryContext(ctx, `SELECT work_id, collected_at, likes FROM creator_work_metric WHERE work_id = ANY($1) AND likes IS NOT NULL ORDER BY work_id, collected_at`, workIDs)
|
||||
if err != nil {
|
||||
return nil, databaseError(err)
|
||||
}
|
||||
defer rows.Close()
|
||||
byWork := make(map[string][]MetricPoint)
|
||||
for rows.Next() {
|
||||
var workID string
|
||||
var point MetricPoint
|
||||
var likes int64
|
||||
if err := rows.Scan(&workID, &point.CollectedAt, &likes); err != nil {
|
||||
return nil, databaseError(err)
|
||||
}
|
||||
point.CollectedAt = point.CollectedAt.UTC()
|
||||
value := likes
|
||||
point.Likes = &value
|
||||
byWork[workID] = append(byWork[workID], point)
|
||||
}
|
||||
if err := rows.Err(); err != nil {
|
||||
return nil, databaseError(err)
|
||||
}
|
||||
now := time.Now().UTC()
|
||||
result := make(map[string]WorkGrowth, len(byWork))
|
||||
for workID, points := range byWork {
|
||||
result[workID] = computeWorkGrowth(points, now, filter.Hours)
|
||||
}
|
||||
return result, nil
|
||||
}
|
||||
|
||||
// ApplyWorkGrowthFilter 为作品列表附加增速并按条件过滤、排序(增速降序)。
|
||||
// 只有设置了 MinLikes 才改变排序;否则保持原顺序仅附加增速。
|
||||
func ApplyWorkGrowthFilter(works []Work, growth map[string]WorkGrowth, filter WorkGrowthFilter) ([]Work, error) {
|
||||
filter, err := normalizeGrowthFilter(filter)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
result := make([]Work, 0, len(works))
|
||||
for _, work := range works {
|
||||
if stat, ok := growth[work.ID]; ok {
|
||||
work.LikesGrowth = stat.Likes
|
||||
work.LikesGrowthCoverage = stat.Coverage
|
||||
}
|
||||
if filter.MinLikes != nil {
|
||||
if work.LikesGrowth == nil || *work.LikesGrowth < *filter.MinLikes {
|
||||
continue
|
||||
}
|
||||
}
|
||||
result = append(result, work)
|
||||
}
|
||||
if filter.MinLikes != nil {
|
||||
sort.SliceStable(result, func(i, j int) bool {
|
||||
gi, gj := result[i].LikesGrowth, result[j].LikesGrowth
|
||||
if gi == nil || gj == nil {
|
||||
return gj == nil && gi != nil
|
||||
}
|
||||
return *gi > *gj
|
||||
})
|
||||
}
|
||||
return result, nil
|
||||
}
|
||||
|
||||
func int64Ptr(value int64) *int64 { return &value }
|
||||
@@ -0,0 +1,82 @@
|
||||
package creator
|
||||
|
||||
import (
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
func TestComputeWorkGrowth(t *testing.T) {
|
||||
now := time.Date(2026, 9, 27, 12, 0, 0, 0, time.UTC)
|
||||
likes := func(at time.Time, value int64) MetricPoint {
|
||||
return MetricPoint{CollectedAt: at, Likes: &value}
|
||||
}
|
||||
tooOld := now.Add(-72 * time.Hour)
|
||||
|
||||
t.Run("full baseline inside window", func(t *testing.T) {
|
||||
points := []MetricPoint{
|
||||
likes(tooOld, 100),
|
||||
likes(now.Add(-25*time.Hour), 200), // 窗口起点前最近点 → 基线
|
||||
likes(now.Add(-2*time.Hour), 500),
|
||||
}
|
||||
growth := computeWorkGrowth(points, now, 24)
|
||||
if growth.Likes == nil || *growth.Likes != 300 || growth.Coverage != "full" {
|
||||
t.Fatalf("unexpected growth: %+v", growth)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("partial when only one point", func(t *testing.T) {
|
||||
points := []MetricPoint{likes(now.Add(-2*time.Hour), 500)}
|
||||
growth := computeWorkGrowth(points, now, 24)
|
||||
if growth.Likes == nil || *growth.Likes != 0 || growth.Coverage != "partial" {
|
||||
t.Fatalf("unexpected growth: %+v", growth)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("partial falls back to first point", func(t *testing.T) {
|
||||
points := []MetricPoint{likes(now.Add(-20*time.Hour), 100), likes(now.Add(-2*time.Hour), 500)}
|
||||
growth := computeWorkGrowth(points, now, 24)
|
||||
// 两个点都在窗口内、窗口前无基线 → 以首点近似,partial。
|
||||
if growth.Likes == nil || *growth.Likes != 400 || growth.Coverage != "partial" {
|
||||
t.Fatalf("unexpected growth: %+v likes=%v", growth, *growth.Likes)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("old pre-window baseline counts as full", func(t *testing.T) {
|
||||
points := []MetricPoint{likes(tooOld, 100), likes(now.Add(-2*time.Hour), 500)}
|
||||
growth := computeWorkGrowth(points, now, 24)
|
||||
// 窗口前存在真实基线(72h 前)→ full;增量跨 72h 但可比。
|
||||
if growth.Likes == nil || *growth.Likes != 400 || growth.Coverage != "full" {
|
||||
t.Fatalf("unexpected growth: %+v", growth)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("no points", func(t *testing.T) {
|
||||
if growth := (WorkGrowth{}); computeWorkGrowth(nil, now, 24) != growth {
|
||||
t.Fatal("empty points must yield zero growth")
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
func TestApplyWorkGrowthFilterSortsAndFilters(t *testing.T) {
|
||||
mk := func(id string, growth int64) (Work, WorkGrowth) {
|
||||
value := growth
|
||||
return Work{ID: id}, WorkGrowth{Likes: &value, Coverage: "full"}
|
||||
}
|
||||
w1, g1 := mk("a", 100)
|
||||
w2, g2 := mk("b", 900)
|
||||
w3, _ := mk("c", 0) // 无采样点:不附加增速
|
||||
works := []Work{w1, w2, w3}
|
||||
growth := map[string]WorkGrowth{"a": g1, "b": g2}
|
||||
|
||||
min := int64(50)
|
||||
filtered, err := ApplyWorkGrowthFilter(works, growth, WorkGrowthFilter{Hours: 24, MinLikes: &min})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(filtered) != 2 || filtered[0].ID != "b" || filtered[1].ID != "a" {
|
||||
t.Fatalf("unexpected filter/sort result: %+v", filtered)
|
||||
}
|
||||
if filtered[0].LikesGrowth == nil || *filtered[0].LikesGrowth != 900 {
|
||||
t.Fatalf("growth not attached: %+v", filtered[0])
|
||||
}
|
||||
}
|
||||
@@ -354,13 +354,44 @@ func TestCreatorPostgresContentAndWorkflow(t *testing.T) {
|
||||
if _, duplicate, err := store.SaveComment(ctx, CommentInput{Platform: PlatformDouyin, CommentKey: comment.CommentKey, WorkID: work.ID, Content: "咨询价格", CommentType: "top_level"}); err != nil || duplicate {
|
||||
t.Fatalf("comment deduplication failed: duplicate=%v err=%v", duplicate, err)
|
||||
}
|
||||
if _, err := store.RecordMetric(ctx, MetricInput{WorkID: work.ID, CollectedAt: now, Likes: &likes, CommentsCount: &comments, Shares: &shares}, Settings{LookbackDays: 30, NewWorkIntervalSeconds: 1800, MetricInitialIntervalSeconds: 3600, MetricMultiplier: 2, MetricMaxIntervalSeconds: 86400, MetricAgeSeconds: 30 * 24 * 60 * 60}, now); err != nil {
|
||||
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 {
|
||||
t.Fatal(err)
|
||||
}
|
||||
metrics, err := store.ListMetrics(ctx, work.ID)
|
||||
if err != nil || len(metrics) != 1 {
|
||||
t.Fatalf("initial metric snapshot failed: metrics=%+v err=%v", metrics, err)
|
||||
}
|
||||
if metrics[0].CollectCount == nil || *metrics[0].CollectCount != 4 {
|
||||
t.Fatalf("collect_count not persisted: %+v", metrics[0])
|
||||
}
|
||||
// 增速计算:独立作品(避免干扰上方计划推进断言),两个采样点。
|
||||
growthPublished := now.Add(-2 * time.Hour)
|
||||
growthLikes0, growthLikes1 := int64(100), int64(400)
|
||||
growthWork, _, err := store.UpsertWork(ctx, WorkInput{Platform: PlatformDouyin, WorkKey: "creator-it-growth-work-" + stamp, SourceType: SourceOwned, SourceID: smallID, Title: "Growth", Body: "body", PublishedAt: &growthPublished, PublishedAtStatus: "verified", Likes: &growthLikes0}, growthPublished.Add(time.Hour))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
growthSettings := Settings{LookbackDays: 30, NewWorkIntervalSeconds: 1800, MetricInitialIntervalSeconds: 3600, MetricMultiplier: 2, MetricMaxIntervalSeconds: 86400, MetricAgeSeconds: 30 * 24 * 60 * 60}
|
||||
growthPoint0 := now.Add(-90 * time.Minute) // 首点(晚于发布时间)
|
||||
if _, err := store.recordMetricWithPlan(ctx, MetricInput{WorkID: growthWork.ID, CollectedAt: growthPoint0, Likes: &growthLikes0}, growthSettings, growthPoint0); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
growthPoint1 := growthPoint0.Add(2 * time.Hour) // 第二个计划点:next = 首点+2h
|
||||
if _, err := store.recordMetricWithPlan(ctx, MetricInput{WorkID: growthWork.ID, CollectedAt: growthPoint1, Likes: &growthLikes1}, growthSettings, growthPoint1); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
growth, err := store.LoadWorksGrowth(ctx, []string{growthWork.ID, work.ID}, WorkGrowthFilter{Hours: 24})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
stat, ok := growth[growthWork.ID]
|
||||
if !ok || stat.Likes == nil || *stat.Likes != 300 || stat.Coverage != "partial" {
|
||||
t.Fatalf("growth not computed: stat=%+v ok=%v", stat, ok)
|
||||
}
|
||||
if stat, ok := growth[work.ID]; !ok || stat.Likes == nil || *stat.Likes != 0 || stat.Coverage != "partial" {
|
||||
t.Fatalf("single-point work growth: %+v ok=%v", stat, ok)
|
||||
}
|
||||
if _, err := store.RecordMetric(ctx, MetricInput{WorkID: work.ID, CollectedAt: now.Add(-time.Minute), Likes: &likes, CommentsCount: &comments, Shares: &shares}, Settings{LookbackDays: 30, NewWorkIntervalSeconds: 1800, MetricInitialIntervalSeconds: 3600, MetricMultiplier: 2, MetricMaxIntervalSeconds: 86400, MetricAgeSeconds: 30 * 24 * 60 * 60}, now); !errors.Is(err, ErrConflict) {
|
||||
t.Fatalf("expected early metric point to be rejected, got %v", err)
|
||||
}
|
||||
|
||||
@@ -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}
|
||||
point := MetricPoint{CollectedAt: collectedAt, Likes: input.Likes, CommentsCount: input.CommentsCount, Shares: input.Shares, CollectCount: input.CollectCount}
|
||||
if !publishedAt.Valid || publishedStatus != "verified" {
|
||||
if _, err := tx.ExecContext(ctx, `UPDATE creator_work SET likes=$2, comments_count=$3, shares=$4, latest_metrics_at=$5, 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, collectedAt); err != nil {
|
||||
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 {
|
||||
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) VALUES ($1,$2,$3,$4,$5) ON CONFLICT (work_id,collected_at) DO UPDATE SET likes=EXCLUDED.likes,comments_count=EXCLUDED.comments_count,shares=EXCLUDED.shares`, input.WorkID, collectedAt, input.Likes, input.CommentsCount, input.Shares); err != nil {
|
||||
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 {
|
||||
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,latest_metrics_at=$5,next_metric_at=$6,metric_stop_reason=$7,updated_at=now() WHERE id=$1`, input.WorkID, input.Likes, input.CommentsCount, input.Shares, 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,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 {
|
||||
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) VALUES ($1,$2,$3,$4,$5) ON CONFLICT (work_id,collected_at) DO UPDATE SET likes=EXCLUDED.likes,comments_count=EXCLUDED.comments_count,shares=EXCLUDED.shares`, input.WorkID, collectedAt, input.Likes, input.CommentsCount, input.Shares); err != nil {
|
||||
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 {
|
||||
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,latest_metrics_at=$5,next_metric_at=$6,metric_stop_reason=$7,updated_at=now() WHERE id=$1`, input.WorkID, input.Likes, input.CommentsCount, input.Shares, 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,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 {
|
||||
return MetricPoint{}, databaseError(err)
|
||||
}
|
||||
if err := tx.Commit(); err != nil {
|
||||
|
||||
@@ -0,0 +1,4 @@
|
||||
-- 作品指标补采收藏数(collect_count):抖音 works/detail 接口 statistics 均携带,
|
||||
-- 此前解析与存储均未覆盖,导致增长曲线缺收藏维度。
|
||||
ALTER TABLE creator_work ADD COLUMN collect_count bigint;
|
||||
ALTER TABLE creator_work_metric ADD COLUMN collect_count bigint;
|
||||
@@ -192,6 +192,9 @@ type Work struct {
|
||||
Likes *int64 `json:"likes"`
|
||||
CommentsCount *int64 `json:"comments_count"`
|
||||
Shares *int64 `json:"shares"`
|
||||
CollectCount *int64 `json:"collect_count"`
|
||||
LikesGrowth *int64 `json:"likes_growth,omitempty"`
|
||||
LikesGrowthCoverage string `json:"likes_growth_coverage,omitempty"` // full=窗口基线完整 partial=以首点近似
|
||||
LatestMetricsAt *time.Time `json:"latest_metrics_at,omitempty"`
|
||||
NextMetricAt *time.Time `json:"next_metric_at,omitempty"`
|
||||
MetricStopReason string `json:"metric_stop_reason,omitempty"`
|
||||
@@ -216,6 +219,7 @@ type WorkInput struct {
|
||||
Likes *int64 `json:"likes"`
|
||||
CommentsCount *int64 `json:"comments_count"`
|
||||
Shares *int64 `json:"shares"`
|
||||
CollectCount *int64 `json:"collect_count"`
|
||||
RawPayload string `json:"-"`
|
||||
}
|
||||
|
||||
@@ -229,6 +233,7 @@ type WorkFilter struct {
|
||||
MinLikes *int64
|
||||
MinComments *int64
|
||||
MinShares *int64
|
||||
Growth WorkGrowthFilter // 点赞增量筛选(爆品),空窗口不启用
|
||||
}
|
||||
|
||||
type MetricInput struct {
|
||||
@@ -237,6 +242,7 @@ type MetricInput struct {
|
||||
Likes *int64
|
||||
CommentsCount *int64
|
||||
Shares *int64
|
||||
CollectCount *int64
|
||||
}
|
||||
|
||||
type MetricPoint struct {
|
||||
@@ -244,6 +250,7 @@ type MetricPoint struct {
|
||||
Likes *int64 `json:"likes"`
|
||||
CommentsCount *int64 `json:"comments_count"`
|
||||
Shares *int64 `json:"shares"`
|
||||
CollectCount *int64 `json:"collect_count"`
|
||||
}
|
||||
|
||||
type MaterialJob struct {
|
||||
|
||||
@@ -86,6 +86,9 @@ var migration037 string
|
||||
//go:embed migrations/038_douyin_only_competitor_profile.sql
|
||||
var migration038 string
|
||||
|
||||
//go:embed migrations/039_work_collect_count.sql
|
||||
var migration039 string
|
||||
|
||||
type SecretReference struct {
|
||||
ID string
|
||||
Provider string
|
||||
@@ -185,6 +188,7 @@ func (s *Store) migrate(ctx context.Context) error {
|
||||
{version: 36, sql: migration036},
|
||||
{version: 37, sql: migration037},
|
||||
{version: 38, sql: migration038},
|
||||
{version: 39, sql: migration039},
|
||||
}
|
||||
for _, migration := range migrations {
|
||||
var applied bool
|
||||
|
||||
@@ -89,6 +89,7 @@ type Work struct {
|
||||
DiggCount int64 `json:"digg_count"`
|
||||
CommentCount int64 `json:"comment_count"`
|
||||
ShareCount int64 `json:"share_count"`
|
||||
CollectCount int64 `json:"collect_count"`
|
||||
PlayCount int64 `json:"play_count"`
|
||||
}
|
||||
|
||||
@@ -355,6 +356,7 @@ type worksEnvelope struct {
|
||||
DiggCount *int64 `json:"digg_count"`
|
||||
CommentCount *int64 `json:"comment_count"`
|
||||
ShareCount *int64 `json:"share_count"`
|
||||
CollectCount *int64 `json:"collect_count"`
|
||||
PlayCount *int64 `json:"play_count"`
|
||||
} `json:"statistics"`
|
||||
} `json:"aweme_list"`
|
||||
|
||||
@@ -34,6 +34,54 @@ type TargetProfile struct {
|
||||
AvatarURL string
|
||||
}
|
||||
|
||||
// WorkStatistics 是作品详情接口返回的指标快照,用于增长曲线时序采集。
|
||||
type WorkStatistics struct {
|
||||
Likes *int64
|
||||
CommentsCount *int64
|
||||
Shares *int64
|
||||
CollectCount *int64
|
||||
}
|
||||
|
||||
// FetchWorkStatistics 经网关受限 fetch 拉取作品详情并解析指标快照。
|
||||
// 作品指标是公开数据,匿名浏览器即可访问(实测 detail 返回完整 statistics)。
|
||||
func FetchWorkStatistics(ctx context.Context, browser Browser, workKey string) (WorkStatistics, error) {
|
||||
if browser == nil || !workKeyPattern.MatchString(workKey) {
|
||||
return WorkStatistics{}, fmt.Errorf("%w: invalid work statistics request", ErrInvalid)
|
||||
}
|
||||
query := douyinAPIQuery()
|
||||
query.Set("aweme_id", workKey)
|
||||
response, err := browser.Get(ctx, workDetailEndpoint+"?"+query.Encode())
|
||||
if err != nil {
|
||||
return WorkStatistics{}, err
|
||||
}
|
||||
if err := creatorResponseError(response, "work detail"); err != nil {
|
||||
return WorkStatistics{}, err
|
||||
}
|
||||
stats, ok := parseWorkStatistics(response.Body, workKey)
|
||||
if !ok {
|
||||
return WorkStatistics{}, fmt.Errorf("%w: invalid douyin work detail response (status=%d body_len=%d head=%.160s)", ErrInvalid, response.Status, len(response.Body), response.Body)
|
||||
}
|
||||
return stats, nil
|
||||
}
|
||||
|
||||
func parseWorkStatistics(body []byte, expectedWorkKey string) (WorkStatistics, bool) {
|
||||
var envelope workDetailEnvelope
|
||||
if len(body) > 2<<20 || json.Unmarshal(body, &envelope) != nil || envelope.StatusCode == nil || *envelope.StatusCode != 0 ||
|
||||
envelope.AwemeDetail == nil || envelope.AwemeDetail.AwemeID != expectedWorkKey || envelope.AwemeDetail.Statistics == nil {
|
||||
return WorkStatistics{}, false
|
||||
}
|
||||
stats := envelope.AwemeDetail.Statistics
|
||||
for _, value := range []*int64{stats.DiggCount, stats.CommentCount, stats.ShareCount, stats.CollectCount} {
|
||||
if value != nil && *value < 0 {
|
||||
return WorkStatistics{}, false
|
||||
}
|
||||
}
|
||||
if stats.DiggCount == nil && stats.CommentCount == nil && stats.ShareCount == nil && stats.CollectCount == nil {
|
||||
return WorkStatistics{}, false
|
||||
}
|
||||
return WorkStatistics{Likes: stats.DiggCount, CommentsCount: stats.CommentCount, Shares: stats.ShareCount, CollectCount: stats.CollectCount}, true
|
||||
}
|
||||
|
||||
type workDetailEnvelope struct {
|
||||
StatusCode *int `json:"status_code"`
|
||||
AwemeDetail *struct {
|
||||
@@ -47,6 +95,12 @@ type workDetailEnvelope struct {
|
||||
URLList []string `json:"url_list"`
|
||||
} `json:"avatar_thumb"`
|
||||
} `json:"author"`
|
||||
Statistics *struct {
|
||||
DiggCount *int64 `json:"digg_count"`
|
||||
CommentCount *int64 `json:"comment_count"`
|
||||
ShareCount *int64 `json:"share_count"`
|
||||
CollectCount *int64 `json:"collect_count"`
|
||||
} `json:"statistics"`
|
||||
} `json:"aweme_detail"`
|
||||
}
|
||||
|
||||
@@ -193,7 +247,7 @@ func (c CreatorCollector) ListWorks(ctx context.Context, accountKey, cursor stri
|
||||
value := time.Unix(*work.CreatedAt, 0).UTC()
|
||||
published = &value
|
||||
}
|
||||
likes, comments, shares := work.DiggCount, work.CommentCount, work.ShareCount
|
||||
likes, comments, shares, collects := work.DiggCount, work.CommentCount, work.ShareCount, work.CollectCount
|
||||
sourceType, sourceID := c.SourceType, c.SourceID
|
||||
if sourceType == "" {
|
||||
sourceType = creator.SourceCompetitor
|
||||
@@ -207,7 +261,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})
|
||||
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})
|
||||
}
|
||||
page := creator.WorkPage{Items: items, HasMore: hasMore}
|
||||
if nextCursor != nil {
|
||||
@@ -288,6 +342,7 @@ type creatorWorkPageItem struct {
|
||||
DiggCount *int64
|
||||
CommentCount *int64
|
||||
ShareCount *int64
|
||||
CollectCount *int64
|
||||
}
|
||||
|
||||
func parseCreatorWorksPage(body []byte) ([]creatorWorkPageItem, bool, *int64, bool) {
|
||||
@@ -314,11 +369,11 @@ func parseCreatorWorksPage(body []byte) ([]creatorWorkPageItem, bool, *int64, bo
|
||||
return nil, false, nil, false
|
||||
}
|
||||
seen[item.ID] = struct{}{}
|
||||
var likes, comments, shares *int64
|
||||
var likes, comments, shares, collects *int64
|
||||
if item.Statistics != nil {
|
||||
likes, comments, shares = item.Statistics.DiggCount, item.Statistics.CommentCount, item.Statistics.ShareCount
|
||||
likes, comments, shares, collects = item.Statistics.DiggCount, item.Statistics.CommentCount, item.Statistics.ShareCount, item.Statistics.CollectCount
|
||||
}
|
||||
for _, value := range []*int64{likes, comments, shares} {
|
||||
for _, value := range []*int64{likes, comments, shares, collects} {
|
||||
if value != nil && *value < 0 {
|
||||
return nil, false, nil, false
|
||||
}
|
||||
@@ -328,7 +383,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})
|
||||
items = append(items, creatorWorkPageItem{ID: item.ID, Description: item.Description, CreatedAt: createdAt, CreatedAtInvalid: createdAtInvalid, DiggCount: likes, CommentCount: comments, ShareCount: shares, CollectCount: collects})
|
||||
}
|
||||
// 仅当响应带有效 max_cursor 时才声明翻页;无 cursor 的 has_more 不产生翻页游标。
|
||||
hasMore := envelope.HasMore != nil && bool(*envelope.HasMore) && envelope.MaxCursor != nil
|
||||
|
||||
@@ -137,11 +137,41 @@ 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}}]}`)
|
||||
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}}]}`)
|
||||
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)
|
||||
}
|
||||
if works[0].CollectCount == nil || *works[0].CollectCount != 4 {
|
||||
t.Fatalf("collect_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}}}`)
|
||||
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 {
|
||||
t.Fatalf("unexpected statistics: ok=%v %+v", ok, stats)
|
||||
}
|
||||
if _, ok := parseWorkStatistics(body, "999"); ok {
|
||||
t.Fatal("accepted statistics for another work")
|
||||
}
|
||||
if _, ok := parseWorkStatistics([]byte(`{"status_code":0,"aweme_detail":{"aweme_id":"123"}}`), "123"); ok {
|
||||
t.Fatal("accepted missing statistics")
|
||||
}
|
||||
}
|
||||
|
||||
func TestFetchWorkStatisticsUsesDetailEndpoint(t *testing.T) {
|
||||
browser := &collectorBrowser{response: Response{Status: 200, Body: []byte(`{"status_code":0,"aweme_detail":{"aweme_id":"123","statistics":{"digg_count":10}}}`)}}
|
||||
stats, err := FetchWorkStatistics(context.Background(), browser, "123")
|
||||
if err != nil || stats.Likes == nil || *stats.Likes != 10 {
|
||||
t.Fatalf("fetch statistics: %+v err=%v", stats, err)
|
||||
}
|
||||
parsed, err := url.Parse(browser.url)
|
||||
if err != nil || parsed.Path != "/aweme/v1/web/aweme/detail/" || parsed.Query().Get("aweme_id") != "123" {
|
||||
t.Fatalf("unexpected request url: %s", browser.url)
|
||||
}
|
||||
}
|
||||
|
||||
func TestParseCreatorWorksPageTreatsBareEnvelopeAsEmptyPage(t *testing.T) {
|
||||
|
||||
Generated
+1137
-451
File diff suppressed because it is too large
Load Diff
@@ -7,6 +7,7 @@
|
||||
"start": "max preview"
|
||||
},
|
||||
"dependencies": {
|
||||
"@ant-design/charts": "^2.6.7",
|
||||
"@ant-design/icons": "^6.3.4",
|
||||
"@ant-design/pro-components": "^3.1.14-7",
|
||||
"@umijs/max": "^4.7.19",
|
||||
|
||||
@@ -1,10 +1,11 @@
|
||||
// 监控账号详情:形态对齐账号详情页(页头槽 + 卡片流)——
|
||||
// 页头 info 为账号标识(头像 + 昵称 + 同步状态 + 平台),actions 为返回/同步/暂停恢复;
|
||||
// 主体两张 Card:账号画像(Descriptions)与作品列表(Row/Col card 网格,分页)。
|
||||
// 点击作品弹出 Drawer 展示统计信息:作品详情 + 指标采集记录表。
|
||||
// 主体两张 Card:账号画像(Descriptions)与作品列表(Row/Col card 网格,分页,支持按点赞增速筛选爆品)。
|
||||
// 点击作品弹出 Drawer 展示统计信息:作品详情 + 增长曲线折线图 + 指标采集记录表。
|
||||
// 作品列表走 GET /creator/works?source_id=<competitor_id>,定时采集由后端调度器维护。
|
||||
import { useCallback, useEffect, useState } from 'react';
|
||||
import { useCallback, useEffect, useMemo, useState } from 'react';
|
||||
import { history, useParams } from '@umijs/max';
|
||||
import { Line } from '@ant-design/charts';
|
||||
import {
|
||||
Alert,
|
||||
Avatar,
|
||||
@@ -15,6 +16,7 @@ import {
|
||||
Drawer,
|
||||
Empty,
|
||||
Flex,
|
||||
InputNumber,
|
||||
Row,
|
||||
Space,
|
||||
Spin,
|
||||
@@ -23,7 +25,7 @@ import {
|
||||
Typography,
|
||||
message,
|
||||
} from 'antd';
|
||||
import { CommentOutlined, LikeOutlined, ReloadOutlined, ShareAltOutlined } from '@ant-design/icons';
|
||||
import { CommentOutlined, FireOutlined, LikeOutlined, ReloadOutlined, ShareAltOutlined } from '@ant-design/icons';
|
||||
import { creatorAction, creatorGet, getOne, getList } from '@/services/api';
|
||||
import { conflictMessage, dateTime, platformLabel } from '@/utils/helpers';
|
||||
import { usePageActions, usePageInfo } from '@/components/PageActions';
|
||||
@@ -60,6 +62,42 @@ function formatCount(value?: number | null): string {
|
||||
return value >= 10000 ? `${(value / 10000).toFixed(1)}w` : `${value}`;
|
||||
}
|
||||
|
||||
const metricSeriesMeta = [
|
||||
{ key: 'likes', label: '点赞' },
|
||||
{ key: 'collect_count', label: '收藏' },
|
||||
{ key: 'comments_count', label: '评论' },
|
||||
{ key: 'shares', label: '转发' },
|
||||
] as const;
|
||||
|
||||
// 增长曲线:把指标时序点展开成长表(时间 × 指标 × 数值)供折线图渲染。
|
||||
function GrowthChart({ metrics }: { metrics: any[] }) {
|
||||
const data = useMemo(
|
||||
() =>
|
||||
metrics.flatMap((point) =>
|
||||
metricSeriesMeta
|
||||
.filter((series) => point[series.key] != null)
|
||||
.map((series) => ({
|
||||
time: dateTime(point.collected_at),
|
||||
metric: series.label,
|
||||
value: point[series.key],
|
||||
})),
|
||||
),
|
||||
[metrics],
|
||||
);
|
||||
if (!data.length) return null;
|
||||
return (
|
||||
<Line
|
||||
data={data}
|
||||
xField="time"
|
||||
yField="value"
|
||||
seriesField="metric"
|
||||
height={240}
|
||||
yAxis={{ labelFormatter: (v: number) => (v >= 10000 ? `${(v / 10000).toFixed(1)}w` : `${v}`) }}
|
||||
legend={{ position: 'top' }}
|
||||
/>
|
||||
);
|
||||
}
|
||||
|
||||
export default function Page() {
|
||||
const { id = '' } = useParams<{ id: string }>();
|
||||
const [competitor, setCompetitor] = useState<CompetitorDetail | null>(null);
|
||||
@@ -74,6 +112,7 @@ export default function Page() {
|
||||
const [detailWork, setDetailWork] = useState<any>(null);
|
||||
const [workMetrics, setWorkMetrics] = useState<any[]>([]);
|
||||
const [metricsPending, setMetricsPending] = useState(false);
|
||||
const [minGrowth, setMinGrowth] = useState<number | null>(null);
|
||||
const [messageApi, contextHolder] = message.useMessage();
|
||||
|
||||
const loadCompetitor = useCallback(async () => {
|
||||
@@ -92,7 +131,13 @@ export default function Page() {
|
||||
setWorksPending(true);
|
||||
setWorksError(null);
|
||||
try {
|
||||
const result = await getList({ resource: 'creator-works', page: workPage, pageSize: 25, filters: { source_id: id } });
|
||||
const filters: Record<string, string | number> = { source_id: id };
|
||||
if (minGrowth !== null) {
|
||||
// 爆品筛选:growth_hours 开启后端 24h 点赞增量计算,min_likes_growth 设下限。
|
||||
filters.growth_hours = 24;
|
||||
filters.min_likes_growth = minGrowth;
|
||||
}
|
||||
const result = await getList({ resource: 'creator-works', page: workPage, pageSize: 25, filters });
|
||||
setWorks(result.data);
|
||||
setWorkPageInfo({ total: result.total, hasNext: Boolean(result.hasNext) });
|
||||
} catch (loadError) {
|
||||
@@ -100,7 +145,7 @@ export default function Page() {
|
||||
} finally {
|
||||
setWorksPending(false);
|
||||
}
|
||||
}, [id, workPage]);
|
||||
}, [id, workPage, minGrowth]);
|
||||
|
||||
useEffect(() => {
|
||||
loadCompetitor();
|
||||
@@ -229,6 +274,21 @@ export default function Page() {
|
||||
>
|
||||
<Flex justify="space-between" align="center" style={{ marginBottom: 16 }}>
|
||||
<Typography.Text type="secondary">共 {workPageInfo.total} 个作品;作品由系统按采集间隔自动同步。</Typography.Text>
|
||||
<Flex align="center" gap={8}>
|
||||
<FireOutlined style={{ color: 'rgba(0,0,0,0.45)' }} />
|
||||
<Typography.Text type="secondary">24h 点赞增量 ≥</Typography.Text>
|
||||
<InputNumber
|
||||
size="small"
|
||||
min={0}
|
||||
placeholder="不限"
|
||||
style={{ width: 96 }}
|
||||
value={minGrowth}
|
||||
onChange={(value) => {
|
||||
setMinGrowth(value ?? null);
|
||||
setWorkPage(1);
|
||||
}}
|
||||
/>
|
||||
</Flex>
|
||||
</Flex>
|
||||
{worksError ? (
|
||||
<Alert
|
||||
@@ -270,6 +330,11 @@ export default function Page() {
|
||||
<Typography.Text style={{ fontSize: 12 }}><CommentOutlined /> {work.comments_count ?? '—'}</Typography.Text>
|
||||
<Typography.Text style={{ fontSize: 12 }}><ShareAltOutlined /> {work.shares ?? '—'}</Typography.Text>
|
||||
</Flex>
|
||||
{work.likes_growth != null && work.likes_growth > 0 ? (
|
||||
<Typography.Text type="warning" style={{ fontSize: 12 }}>
|
||||
<FireOutlined /> 24h +{work.likes_growth}
|
||||
</Typography.Text>
|
||||
) : null}
|
||||
</Flex>
|
||||
</Card>
|
||||
</Col>
|
||||
@@ -309,6 +374,7 @@ export default function Page() {
|
||||
<Typography.Link href={detailWork.original_url} target="_blank">打开原文</Typography.Link>
|
||||
) : null}
|
||||
</Space>
|
||||
{workMetrics.length >= 2 ? <GrowthChart metrics={workMetrics} /> : null}
|
||||
<Table
|
||||
rowKey="collected_at"
|
||||
size="small"
|
||||
@@ -318,6 +384,7 @@ export default function Page() {
|
||||
columns={[
|
||||
{ title: '采集时间', dataIndex: 'collected_at', render: (value: string) => dateTime(value) },
|
||||
{ title: '点赞', dataIndex: 'likes', render: (v) => v ?? '—' },
|
||||
{ title: '收藏', dataIndex: 'collect_count', render: (v) => v ?? '—' },
|
||||
{ title: '评论', dataIndex: 'comments_count', render: (v) => v ?? '—' },
|
||||
{ title: '转发', dataIndex: 'shares', render: (v) => v ?? '—' },
|
||||
]}
|
||||
|
||||
@@ -33,7 +33,7 @@ const filterKeys: Partial<Record<Resource, string[]>> = {
|
||||
tasks: ['account_id', 'draft_id', 'state'],
|
||||
'creator-competitors': ['platform'],
|
||||
'creator-competitor-share-jobs': ['platform', 'status'],
|
||||
'creator-works': ['platform', 'source_id', 'source_type', 'published_at_status', 'published_after', 'published_before', 'min_likes', 'min_comments', 'min_shares'],
|
||||
'creator-works': ['platform', 'source_id', 'source_type', 'published_at_status', 'published_after', 'published_before', 'min_likes', 'min_comments', 'min_shares', 'growth_hours', 'min_likes_growth'],
|
||||
'creator-comments': ['platform', 'work_id'],
|
||||
'creator-leads': ['platform'],
|
||||
'creator-operations': ['account_id'],
|
||||
|
||||
Reference in New Issue
Block a user