diff --git a/internal/controlplane/api/creator.go b/internal/controlplane/api/creator.go index 84f2536..2cbd8c9 100644 --- a/internal/controlplane/api/creator.go +++ b/internal/controlplane/api/creator.go @@ -1285,6 +1285,24 @@ func previewDouyinCompetitor(ctx context.Context, store *creator.Store, phaseASt }, nil } +// refreshCompetitorProfile 拉取抖音公开主页画像并回填监控账号画像字段。 +// profile/other 响应自带粉丝/关注/获赞/作品数与签名,此前未回填导致详情页画像恒为空。 +func refreshCompetitorProfile(ctx context.Context, store *creator.Store, browser creatorGatewayBrowser, competitor creator.Competitor) error { + profile, err := (douyin.CreatorCollector{Browser: browser}).ResolveTarget(ctx, competitor.PlatformAccountKey) + if err != nil { + return err + } + _, err = store.UpdateCompetitorProfile(ctx, competitor.ID, creator.CompetitorProfileUpdate{ + AvatarURL: profile.AvatarURL, + Signature: profile.Signature, + FollowerCount: profile.FollowerCount, + FollowingCount: profile.FollowingCount, + TotalFavorited: profile.TotalFavorited, + AwemeCount: profile.AwemeCount, + }) + return err +} + func syncCreatorCompetitor(ctx context.Context, store *creator.Store, hubStore *hub.Store, competitorID string) (creator.CollectionReport, error) { return syncCreatorCompetitorWithClaim(ctx, store, hubStore, competitorID, true) } @@ -1349,6 +1367,11 @@ func syncCreatorCompetitorWithClaim(ctx context.Context, store *creator.Store, h } report, collectErr := store.CollectSource(ctx, competitor.Platform, creator.SourceCompetitor, competitor.ID, collector, collectionNow) if collectErr == nil { + // 采集成功后回填监控账号画像(粉丝/关注/获赞/作品数/作者介绍); + // profile 拉取失败只记日志,不影响本次作品同步结果。 + if profileErr := refreshCompetitorProfile(ctx, store, creatorGatewayBrowser{gateway: lease.gateway, environment: lease.environment}, competitor); profileErr != nil { + logrus.WithError(profileErr).WithField("competitor_id", competitorID).Warn("creator competitor profile refresh failed") + } // 采集成功后把缺封面的作品封面缓存到本地(失败只记日志,不影响本次同步结果)。 if coverErr := cacheCompetitorWorkCovers(ctx, store, creatorGatewayBrowser{gateway: lease.gateway, environment: lease.environment}, competitor.ID); coverErr != nil { logrus.WithError(coverErr).WithField("competitor_id", competitorID).Warn("creator competitor cover cache failed") diff --git a/internal/creator/content.go b/internal/creator/content.go index e7995a3..2e712b6 100644 --- a/internal/creator/content.go +++ b/internal/creator/content.go @@ -72,20 +72,21 @@ func (s *Store) CreateCompetitor(ctx context.Context, input CompetitorInput) (Co // competitorColumns 是 creator_competitor 的完整列清单(含画像属性),各查询共用。 const competitorColumns = `competitor_id, platform, platform_account_key, unique_id, nickname, avatar_url, homepage_url, tags, - enabled, follower_count, following_count, aweme_count, sync_status, sync_cursor, sync_error, sync_lease_until, + signature, enabled, follower_count, following_count, total_favorited, aweme_count, sync_status, sync_cursor, sync_error, sync_lease_until, last_sync_at, next_sync_at, created_at, updated_at` type competitorScan struct { - competitor Competitor - tags pgtype.FlatArray[string] - leaseUntil, lastSync, nextSync sql.NullTime - follower, following, aweme sql.NullInt64 + competitor Competitor + tags pgtype.FlatArray[string] + leaseUntil, lastSync, nextSync sql.NullTime + follower, following, favorited sql.NullInt64 + aweme sql.NullInt64 } func competitorScanDestinations(scan *competitorScan) []any { return []any{&scan.competitor.ID, &scan.competitor.Platform, &scan.competitor.PlatformAccountKey, &scan.competitor.UniqueID, &scan.competitor.Nickname, - &scan.competitor.AvatarURL, &scan.competitor.HomepageURL, pgtype.NewMap().SQLScanner(&scan.tags), &scan.competitor.Enabled, - &scan.follower, &scan.following, &scan.aweme, &scan.competitor.SyncStatus, &scan.competitor.SyncCursor, + &scan.competitor.AvatarURL, &scan.competitor.HomepageURL, pgtype.NewMap().SQLScanner(&scan.tags), &scan.competitor.Signature, &scan.competitor.Enabled, + &scan.follower, &scan.following, &scan.favorited, &scan.aweme, &scan.competitor.SyncStatus, &scan.competitor.SyncCursor, &scan.competitor.SyncError, &scan.leaseUntil, &scan.lastSync, &scan.nextSync, &scan.competitor.CreatedAt, &scan.competitor.UpdatedAt} } @@ -93,6 +94,7 @@ func (scan *competitorScan) materialize() Competitor { scan.competitor.Tags = []string(scan.tags) scan.competitor.FollowerCount = nullableInt64(scan.follower) scan.competitor.FollowingCount = nullableInt64(scan.following) + scan.competitor.TotalFavorited = nullableInt64(scan.favorited) scan.competitor.AwemeCount = nullableInt64(scan.aweme) scan.competitor.SyncLeaseUntil = nullableTime(scan.leaseUntil) scan.competitor.LastSyncAt = nullableTime(scan.lastSync) @@ -166,18 +168,37 @@ func scanCompetitorView(scanner interface{ Scan(...any) error }) (CompetitorView return result, nil } -// UpdateCompetitorProfile 更新监控账号画像属性(头像与粉丝/关注/作品计数)。 -// 当前仅提供数据通道(字段定义与接口返回),采集侧回填由后续接入。 -func (s *Store) UpdateCompetitorProfile(ctx context.Context, id, avatarURL string, followerCount, followingCount, awemeCount int64) (Competitor, error) { +// CompetitorProfileUpdate 监控账号画像回填值;nil 计数表示本次未采集到,保留原值。 +type CompetitorProfileUpdate struct { + AvatarURL string + Signature string + FollowerCount *int64 + FollowingCount *int64 + TotalFavorited *int64 + AwemeCount *int64 +} + +// UpdateCompetitorProfile 回填监控账号画像(头像/签名/粉丝/关注/获赞/作品计数)。 +func (s *Store) UpdateCompetitorProfile(ctx context.Context, id string, update CompetitorProfileUpdate) (Competitor, error) { id = strings.TrimSpace(id) - avatarURL = strings.TrimSpace(avatarURL) - if id == "" || utf8.RuneCountInString(avatarURL) > 2000 || followerCount < 0 || followingCount < 0 || awemeCount < 0 { + avatarURL := strings.TrimSpace(update.AvatarURL) + signature := strings.TrimSpace(update.Signature) + if id == "" || utf8.RuneCountInString(avatarURL) > 2000 || utf8.RuneCountInString(signature) > 2000 { return Competitor{}, ErrInvalid } + for _, value := range []*int64{update.FollowerCount, update.FollowingCount, update.TotalFavorited, update.AwemeCount} { + if value != nil && *value < 0 { + return Competitor{}, ErrInvalid + } + } result, err := s.db.ExecContext(ctx, ` UPDATE creator_competitor - SET avatar_url = $2, follower_count = $3, following_count = $4, aweme_count = $5, updated_at = now() - WHERE competitor_id = $1`, id, avatarURL, followerCount, followingCount, awemeCount) + SET avatar_url = CASE WHEN $2 = '' THEN avatar_url ELSE $2 END, + signature = CASE WHEN $3 = '' THEN signature ELSE $3 END, + follower_count = COALESCE($4, follower_count), following_count = COALESCE($5, following_count), + total_favorited = COALESCE($6, total_favorited), aweme_count = COALESCE($7, aweme_count), + updated_at = now() + WHERE competitor_id = $1`, id, avatarURL, signature, update.FollowerCount, update.FollowingCount, update.TotalFavorited, update.AwemeCount) if err != nil { return Competitor{}, databaseError(err) } diff --git a/internal/creator/integration_test.go b/internal/creator/integration_test.go index f37a847..a690925 100644 --- a/internal/creator/integration_test.go +++ b/internal/creator/integration_test.go @@ -171,6 +171,49 @@ func TestCompetitorShareJobLifecycle(t *testing.T) { } // 强制同步抢锁失败的分类:租约被占 → ErrSyncInProgress;竞品停用 → ErrSyncDisabled。 +func TestCompetitorProfileUpdateRoundTrip(t *testing.T) { + store, _, ctx := openCreatorIntegrationStore(t) + stamp := fmt.Sprintf("%d", time.Now().UnixNano()) + competitor, err := store.UpsertCompetitor(ctx, CompetitorInput{ + Platform: PlatformDouyin, + PlatformAccountKey: "sec_uid_profile_" + stamp, + Nickname: "Profile Target", + HomepageURL: "https://www.douyin.com/user/sec_uid_profile_" + stamp, + }) + if err != nil { + t.Fatal(err) + } + follower, following, favorited, aweme := int64(123), int64(45), int64(98765), int64(7) + updated, err := store.UpdateCompetitorProfile(ctx, competitor.ID, CompetitorProfileUpdate{ + AvatarURL: "https://example.invalid/avatar", + Signature: "第一行介绍\n第二行介绍", + FollowerCount: &follower, + FollowingCount: &following, + TotalFavorited: &favorited, + AwemeCount: &aweme, + }) + if err != nil { + t.Fatal(err) + } + if updated.Signature != "第一行介绍\n第二行介绍" || updated.FollowerCount == nil || *updated.FollowerCount != 123 || + updated.FollowingCount == nil || *updated.FollowingCount != 45 || updated.TotalFavorited == nil || *updated.TotalFavorited != 98765 || + updated.AwemeCount == nil || *updated.AwemeCount != 7 { + t.Fatalf("profile update: %+v", updated) + } + // nil 计数保留原值,不置 NULL。 + reread, err := store.UpdateCompetitorProfile(ctx, competitor.ID, CompetitorProfileUpdate{Signature: "新签名"}) + if err != nil || reread.Signature != "新签名" || reread.FollowerCount == nil || *reread.FollowerCount != 123 || reread.TotalFavorited == nil || *reread.TotalFavorited != 98765 { + t.Fatalf("profile partial update: %+v err=%v", reread, err) + } + negative := int64(-1) + if _, err := store.UpdateCompetitorProfile(ctx, competitor.ID, CompetitorProfileUpdate{FollowerCount: &negative}); err == nil { + t.Fatal("negative follower count accepted") + } + if _, err := store.UpdateCompetitorProfile(ctx, "missing_"+stamp, CompetitorProfileUpdate{}); err == nil { + t.Fatal("missing competitor accepted") + } +} + func TestCompetitorSyncClaimFailureClassification(t *testing.T) { store, _, ctx := openCreatorIntegrationStore(t) stamp := fmt.Sprintf("%d", time.Now().UnixNano()) @@ -378,7 +421,8 @@ func TestCreatorPostgresContentAndWorkflow(t *testing.T) { if views[0].WorkCount != 1 || views[0].LatestPublishedAt == nil { t.Fatalf("competitor view aggregates: work_count=%d latest=%v", views[0].WorkCount, views[0].LatestPublishedAt) } - if _, err := store.UpdateCompetitorProfile(ctx, competitor.ID, "https://example.invalid/avatar.jpg", 1200, 300, 88); err != nil { + follower, following, aweme := int64(1200), int64(300), int64(88) + if _, err := store.UpdateCompetitorProfile(ctx, competitor.ID, CompetitorProfileUpdate{AvatarURL: "https://example.invalid/avatar.jpg", FollowerCount: &follower, FollowingCount: &following, AwemeCount: &aweme}); err != nil { t.Fatal(err) } views, err = store.ListCompetitorsWithProfile(ctx, "") diff --git a/internal/creator/models.go b/internal/creator/models.go index 5c5bc1a..f322579 100644 --- a/internal/creator/models.go +++ b/internal/creator/models.go @@ -103,9 +103,11 @@ type Competitor struct { AvatarURL string `json:"avatar_url"` HomepageURL string `json:"homepage_url"` Tags []string `json:"tags"` + Signature string `json:"signature"` Enabled bool `json:"enabled"` FollowerCount *int64 `json:"follower_count,omitempty"` FollowingCount *int64 `json:"following_count,omitempty"` + TotalFavorited *int64 `json:"total_favorited,omitempty"` AwemeCount *int64 `json:"aweme_count,omitempty"` SyncStatus string `json:"sync_status"` SyncCursor string `json:"sync_cursor,omitempty"` diff --git a/internal/environment/migration043_probe_test.go b/internal/environment/migration043_probe_test.go index 77aebc6..febb324 100644 --- a/internal/environment/migration043_probe_test.go +++ b/internal/environment/migration043_probe_test.go @@ -71,5 +71,5 @@ func TestMigration043ConsolidatedSchemaShape(t *testing.T) { assertDatabaseCount(t, db, `SELECT count(*) FROM information_schema.columns WHERE table_schema = current_schema() AND table_name = 'audit_event' AND column_name IN ('confirmation_id','confirmation_version','attempt_id','task_id','runtime_instance_id')`, 0) // 统一登记表:1-38(除 36)、1017-1045、43 全部登记 - assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 52) + assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 53) } diff --git a/internal/environment/migration_test.go b/internal/environment/migration_test.go index a44e30a..19e079e 100644 --- a/internal/environment/migration_test.go +++ b/internal/environment/migration_test.go @@ -30,7 +30,7 @@ func TestUnifiedAccountMigration(t *testing.T) { t.Fatal(err) } defer db.Close() - assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 52) + assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 53) assertDatabaseCount(t, db, `SELECT count(*) FROM information_schema.tables WHERE table_schema = current_schema() AND table_name IN ('social_account', 'browser_env', 'network_exit', 'environment_binding')`, 3) assertDatabaseCount(t, db, `SELECT count(*) FROM information_schema.tables WHERE table_schema = current_schema() AND table_name = 'browser_image'`, 0) @@ -42,7 +42,7 @@ func TestUnifiedAccountMigration(t *testing.T) { store = openFullyMigratedHub(t, ctx, testURL) store.Close() - assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 52) + assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 53) }) t.Run("legacy migration 013 without account secrets is repaired forward", func(t *testing.T) { diff --git a/internal/environment/migrations/1046_competitor_profile_signature.sql b/internal/environment/migrations/1046_competitor_profile_signature.sql new file mode 100644 index 0000000..0757d62 --- /dev/null +++ b/internal/environment/migrations/1046_competitor_profile_signature.sql @@ -0,0 +1,13 @@ +-- 监控账号画像补全:作者介绍与获赞总数(同步时经抖音 profile/other 接口回填)。 + +ALTER TABLE creator_competitor + ADD COLUMN IF NOT EXISTS signature text NOT NULL DEFAULT '', + ADD COLUMN IF NOT EXISTS total_favorited bigint; + +ALTER TABLE creator_competitor + ADD CONSTRAINT creator_competitor_counts_non_negative CHECK ( + (follower_count IS NULL OR follower_count >= 0) + AND (following_count IS NULL OR following_count >= 0) + AND (total_favorited IS NULL OR total_favorited >= 0) + AND (aweme_count IS NULL OR aweme_count >= 0) + ); diff --git a/internal/environment/store.go b/internal/environment/store.go index a727e69..b537ad4 100644 --- a/internal/environment/store.go +++ b/internal/environment/store.go @@ -176,6 +176,9 @@ var migration1044 string //go:embed migrations/1045_drop_authorization_status.sql var migration1045 string +//go:embed migrations/1046_competitor_profile_signature.sql +var migration1046 string + var ( ErrConflict = errors.New("resource conflicts with existing state") ErrInvalid = errors.New("invalid environment input") @@ -315,7 +318,7 @@ func (s *Store) migrate(ctx context.Context) error { {1023, migration1023}, {1024, migration1024}, {1025, migration1025}, {1026, migration1026}, {1027, migration1027}, {1028, migration1028}, {1029, migration1029}, {1030, migration1030}, {1031, migration1031}, {1032, migration1032}, {1033, migration1033}, {1034, migration1034}, {1035, migration1035}, {1036, migration1036}, {1037, migration1037}, {1038, migration1038}, {1039, migration1039}, {1040, migration1040}, - {1041, migration1041}, {1042, migration1042}, {1043, migration1043}, {1044, migration1044}, {1045, migration1045}, + {1041, migration1041}, {1042, migration1042}, {1043, migration1043}, {1044, migration1044}, {1045, migration1045}, {1046, migration1046}, {43, migration043}, {44, migration044}} { var applied bool if err := tx.QueryRowContext(ctx, `SELECT EXISTS (SELECT 1 FROM schema_migration WHERE version = $1)`, migration.version).Scan(&applied); err != nil { diff --git a/internal/platform/douyin/connector.go b/internal/platform/douyin/connector.go index 22cedf6..25122ab 100644 --- a/internal/platform/douyin/connector.go +++ b/internal/platform/douyin/connector.go @@ -316,6 +316,7 @@ type identityEnvelope struct { AvatarThumb *struct { URLList []string `json:"url_list"` } `json:"avatar_thumb"` + Signature string `json:"signature"` FollowerCount *int64 `json:"follower_count"` FollowingCount *int64 `json:"following_count"` AwemeCount *int64 `json:"aweme_count"` @@ -327,6 +328,14 @@ type identityEnvelope struct { // flexibleInt 兼容抖音返回的数字与字符串两种计数形态(total_favorited 实测为大数字符串)。 type flexibleInt int64 +func flexibleIntPtr(value *flexibleInt) *int64 { + if value == nil { + return nil + } + converted := int64(*value) + return &converted +} + func (value *flexibleInt) UnmarshalJSON(data []byte) error { text := strings.TrimSpace(string(data)) if text == "null" { diff --git a/internal/platform/douyin/creator_collector.go b/internal/platform/douyin/creator_collector.go index f263db2..210e5e0 100644 --- a/internal/platform/douyin/creator_collector.go +++ b/internal/platform/douyin/creator_collector.go @@ -27,11 +27,16 @@ type CreatorCollector struct { } type TargetProfile struct { - UID string - SecUID string - UniqueID string - Nickname string - AvatarURL string + UID string + SecUID string + UniqueID string + Nickname string + AvatarURL string + Signature string + FollowerCount *int64 + FollowingCount *int64 + TotalFavorited *int64 + AwemeCount *int64 } // WorkStatistics 是作品详情接口返回的指标快照,用于增长曲线时序采集。 @@ -174,7 +179,11 @@ func (c CreatorCollector) ResolveTarget(ctx context.Context, expectedKey string) if identity.User.AvatarThumb != nil && len(identity.User.AvatarThumb.URLList) > 0 { avatarURL = identity.User.AvatarThumb.URLList[0] } - return TargetProfile{UID: identity.User.UID, SecUID: identity.User.SecUID, UniqueID: identity.User.UniqueID, Nickname: identity.User.Nickname, AvatarURL: avatarURL}, nil + return TargetProfile{ + UID: identity.User.UID, SecUID: identity.User.SecUID, UniqueID: identity.User.UniqueID, Nickname: identity.User.Nickname, AvatarURL: avatarURL, + Signature: identity.User.Signature, FollowerCount: identity.User.FollowerCount, FollowingCount: identity.User.FollowingCount, + TotalFavorited: flexibleIntPtr(identity.User.TotalFavorited), AwemeCount: identity.User.AwemeCount, + }, nil } func (c CreatorCollector) ResolveWork(ctx context.Context, workKey string) (TargetProfile, error) { diff --git a/internal/platform/douyin/creator_collector_test.go b/internal/platform/douyin/creator_collector_test.go index bd3a627..45a2ee0 100644 --- a/internal/platform/douyin/creator_collector_test.go +++ b/internal/platform/douyin/creator_collector_test.go @@ -37,6 +37,33 @@ func TestResolveTargetReturnsVerifiedProfileFields(t *testing.T) { } } +func TestResolveTargetExtractsPublicProfileMetrics(t *testing.T) { + browser := &collectorBrowser{response: Response{Status: 200, Body: []byte(`{"status_code":0,"user":{"uid":"2328120603967913","sec_uid":"MS4wLjABAAAA9f_a7k0bzVizLYXlpC7R61EIaqJ8Ordug7yp7AB8fGKuuF8Fzqk5_DM-eutXnPIK","unique_id":"96332518739","nickname":"目标","signature":"第一行介绍\n第二行介绍","follower_count":123,"following_count":45,"aweme_count":7,"total_favorited":"98765"}}`)}} + profile, err := (CreatorCollector{Browser: browser}).ResolveTarget(context.Background(), "MS4wLjABAAAA9f_a7k0bzVizLYXlpC7R61EIaqJ8Ordug7yp7AB8fGKuuF8Fzqk5_DM-eutXnPIK") + if err != nil { + t.Fatal(err) + } + if profile.Signature != "第一行介绍\n第二行介绍" { + t.Fatalf("signature: %q", profile.Signature) + } + if profile.FollowerCount == nil || *profile.FollowerCount != 123 || profile.FollowingCount == nil || *profile.FollowingCount != 45 || + profile.AwemeCount == nil || *profile.AwemeCount != 7 { + t.Fatalf("counts: %+v", profile) + } + // total_favorited 实测为大数字符串,必须经 flexibleInt 兼容解析。 + if profile.TotalFavorited == nil || *profile.TotalFavorited != 98765 { + t.Fatalf("total_favorited: %+v", profile.TotalFavorited) + } +} + +func TestResolveTargetToleratesMissingPublicMetrics(t *testing.T) { + browser := &collectorBrowser{response: Response{Status: 200, Body: []byte(`{"status_code":0,"user":{"uid":"2328120603967913","sec_uid":"MS4wLjABAAAA9f_a7k0bzVizLYXlpC7R61EIaqJ8Ordug7yp7AB8fGKuuF8Fzqk5_DM-eutXnPIK","nickname":"目标"}}`)}} + profile, err := (CreatorCollector{Browser: browser}).ResolveTarget(context.Background(), "MS4wLjABAAAA9f_a7k0bzVizLYXlpC7R61EIaqJ8Ordug7yp7AB8fGKuuF8Fzqk5_DM-eutXnPIK") + if err != nil || profile.Signature != "" || profile.FollowerCount != nil || profile.TotalFavorited != nil { + t.Fatalf("target profile without metrics: %+v err=%v", profile, err) + } +} + func TestCanonicalTargetSecUIDRejectsUniqueIDLookup(t *testing.T) { browser := &collectorBrowser{} _, err := (CreatorCollector{Browser: browser}).CanonicalTargetSecUID(context.Background(), "creator_handle")