From 4180ec104c00231e07dd4922bdbfe383191718a5 Mon Sep 17 00:00:00 2001 From: Rogee Date: Tue, 6 Oct 2026 11:55:36 +0800 Subject: [PATCH] fix: cache owned work covers as UID-scoped static image files --- .env.example | 3 + .gitignore | 1 + AGENTS.md | 2 + compose.yaml | 3 + internal/controlplane/api/creator.go | 41 ++--- .../api/creator_internal_coverage_test.go | 39 +++-- .../api/owned_cover_cache_test.go | 97 +++++++++++ internal/creator/content.go | 66 +------- internal/creator/coverage_integration_test.go | 53 +++--- internal/creator/covers.go | 157 ++++++++++++++++++ internal/creator/covers_integration_test.go | 113 +++++++++++++ internal/creator/integration_test.go | 1 + internal/creator/store.go | 18 +- .../environment/migration043_probe_test.go | 12 +- internal/environment/migration_test.go | 4 +- .../1049_work_covers_static_files.sql | 3 + internal/environment/store.go | 5 +- scripts/dev-backend.mjs | 1 + 18 files changed, 485 insertions(+), 134 deletions(-) create mode 100644 internal/controlplane/api/owned_cover_cache_test.go create mode 100644 internal/creator/covers.go create mode 100644 internal/creator/covers_integration_test.go create mode 100644 internal/environment/migrations/1049_work_covers_static_files.sql diff --git a/.env.example b/.env.example index 687a081..3dffd6c 100644 --- a/.env.example +++ b/.env.example @@ -5,6 +5,9 @@ CONTROL_PLANE_USERNAME=admin CONTROL_PLANE_PASSWORD=change-me CREATORHUB_CREDENTIAL_MASTER_KEY= # openssl rand -base64 32 +# 作品封面目录:保存为 <账号 UID>/<作品 ID>.<图片扩展名> +CREATOR_COVER_DIR=.data/covers + # 数据库 CREATORHUB_POSTGRES_PORT=15433 diff --git a/.gitignore b/.gitignore index d6c8aa4..cf94366 100644 --- a/.gitignore +++ b/.gitignore @@ -34,6 +34,7 @@ __pycache__/ .env.* !.env.example .dev-credentials/ +.data/ # 工具缓存与本机数据 .codegraph/ diff --git a/AGENTS.md b/AGENTS.md index 909175f..7208dd0 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -41,6 +41,8 @@ 环境先行:创建时只创建抖音浏览器环境,不填写昵称、UID 或 Cookie,不创建占位账号。未登录环境在账号列表独立显示。首次浏览器身份核验成功后,同一事务创建/关联账号并同步真实 UID、昵称、头像、抖音号与 secUID;同一 UID 只能绑定一个账号,已有绑定禁止换绑或自动覆盖。后续同 UID 核验更新平台资料,不覆盖本地备注和业务设置。浏览器 profile_id 与指纹 seed 在环境创建时确定,账号绑定不得改变它们;旧环境保持原 profile_id 和 seed,不重建浏览器或清空 Cookie。指纹 seed 由环境独立序列分配,不再依赖账号 ID。 +作品封面:自有账号与监测账号采集后都下载封面,图片保存为静态文件 `<账号 UID>/<作品 ID>.<实际图片扩展名>`,不把图片字节存入数据库。目录由 `CREATOR_COVER_DIR` 配置,开发默认 `.data/covers`,容器使用持久卷;作品封面接口只读取本地文件,不使用远程图片兜底。下载与文件错误必须可追踪并反映到同步结果;同一来源每次同步最多回填 40 张历史缺失封面。 + 账号创建交互:在“我的账号”列表通过 Modal 创建,网关必选,浏览器指纹为默认折叠的可选配置;创建成功关闭弹窗并刷新列表,不保留独立新增页面。未登录环境及空昵称仅在前端显示“待登录”,不将此文案写入数据库。 前端框架:Umi Max 4.7 + React 19; diff --git a/compose.yaml b/compose.yaml index ea28d1a..8fa4d86 100644 --- a/compose.yaml +++ b/compose.yaml @@ -10,6 +10,7 @@ services: BAILIAN_API_KEY: ${BAILIAN_API_KEY:-} BAILIAN_BASE_URL: ${BAILIAN_BASE_URL:-} CREATOR_MEDIA_DIR: /var/lib/creatorhub/materials + CREATOR_COVER_DIR: /var/lib/creatorhub/covers CREATOR_TRANSCRIPTION_BIN: ${CREATOR_TRANSCRIPTION_BIN:-} ports: - "${CREATORHUB_PORT:-8080}:8080" @@ -17,6 +18,7 @@ services: volumes: - creatorhub_credentials:/var/lib/creatorhub/credentials - creatorhub_materials:/var/lib/creatorhub/materials + - creatorhub_covers:/var/lib/creatorhub/covers tmpfs: - /tmp:size=16m,noexec,nosuid,nodev cap_drop: [ALL] @@ -72,3 +74,4 @@ volumes: creatorhub_postgres: creatorhub_credentials: creatorhub_materials: + creatorhub_covers: diff --git a/internal/controlplane/api/creator.go b/internal/controlplane/api/creator.go index 909c7d1..daa11d2 100644 --- a/internal/controlplane/api/creator.go +++ b/internal/controlplane/api/creator.go @@ -433,13 +433,12 @@ func registerCreatorWithServices(app *fiber.App, store *creator.Store, phaseASto return c.JSON(point) }) app.Get("/api/creator/works/:id/cover", func(c fiber.Ctx) error { - contentType, data, err := store.GetWorkCover(c.Context(), c.Params("id"), "cover") + path, err := store.GetWorkCover(c.Context(), c.Params("id")) if err != nil { return creatorError(c, err) } - c.Set(fiber.HeaderContentType, contentType) c.Set(fiber.HeaderCacheControl, "private, max-age=86400") - return c.Send(data) + return c.SendFile(path) }) app.Get("/api/creator/comments", func(c fiber.Ctx) error { page, pageSize, paged, err := creatorPageQuery(c) @@ -967,26 +966,31 @@ func (browser creatorGatewayBrowser) GetImage(ctx context.Context, target string const maxCreatorCoverBytes = 4 << 20 -// cacheCompetitorWorkCovers 把竞品作品封面拉取到本地缓存(带签名的 douyinpic URL 会过期,前端读本地缓存)。 -// 单次同步最多回填 40 张,超出部分留待后续同步自愈,避免首次采集时长时间占用匿名浏览器。 -func cacheCompetitorWorkCovers(ctx context.Context, store *creator.Store, browser creatorGatewayBrowser, competitorID string) error { - works, err := store.ListWorksMissingCover(ctx, competitorID) +// cacheWorkCovers persists owned and monitored covers as UID/work-ID image files. +func cacheWorkCovers(ctx context.Context, store *creator.Store, browser creatorGatewayBrowser, sourceType, sourceID string) error { + works, err := store.ListWorksMissingCover(ctx, sourceType, sourceID) if err != nil { return err } + cached := 0 + var failures []error for _, work := range works { - contentType, data, fetchErr := browser.GetImage(ctx, work.CoverURL) - if fetchErr != nil { - return fmt.Errorf("fetch cover for work %s: %w", work.ID, fetchErr) + contentType, data, workErr := browser.GetImage(ctx, work.CoverURL) + if workErr == nil { + workErr = store.SaveWorkCover(ctx, work.ID, contentType, data) } - if saveErr := store.SaveWorkCover(ctx, work.ID, "cover", contentType, data); saveErr != nil { - return fmt.Errorf("save cover for work %s: %w", work.ID, saveErr) + if workErr != nil { + failure := fmt.Errorf("cache cover for work %s (%s): %w", work.ID, work.WorkKey, workErr) + logrus.WithError(failure).WithFields(logrus.Fields{"source_type": sourceType, "source_id": sourceID, "work_id": work.ID}).Error("creator cover cache failed") + failures = append(failures, failure) + continue } + cached++ } if len(works) > 0 { - logrus.WithFields(logrus.Fields{"competitor_id": competitorID, "covers_cached": len(works)}).Info("creator competitor covers cached") + logrus.WithFields(logrus.Fields{"source_type": sourceType, "source_id": sourceID, "covers_requested": len(works), "covers_cached": cached, "covers_failed": len(failures)}).Info("creator covers cached") } - return nil + return errors.Join(failures...) } func (browser creatorGatewayBrowser) Resolve(ctx context.Context, target string) (string, error) { @@ -1457,11 +1461,9 @@ func syncCreatorCompetitorWithClaim(ctx context.Context, store *creator.Store, h 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") - } } + coverErr := cacheWorkCovers(ctx, store, creatorGatewayBrowser{gateway: lease.gateway, environment: lease.environment}, creator.SourceCompetitor, competitor.ID) + collectErr = errors.Join(collectErr, coverErr) if closeErr := lease.close(); closeErr != nil { collectErr = errors.Join(collectErr, closeErr) } @@ -1775,7 +1777,8 @@ func syncCreatorOwned(ctx context.Context, store *creator.Store, phaseAStore *ac if collectErr == nil { profileErr = recordCreatorAccountMetricSnapshot(ctx, store, creatorGatewayBrowser{gateway: gateway, environment: environment}, account.ID) } - return errors.Join(collectErr, profileErr) + coverErr := cacheWorkCovers(ctx, store, creatorGatewayBrowser{gateway: gateway, environment: environment}, creator.SourceOwned, account.ID) + return errors.Join(collectErr, profileErr, coverErr) } // recordCreatorAccountMetricSnapshot 拉取登录账号自己的画像并记一次时序快照。 diff --git a/internal/controlplane/api/creator_internal_coverage_test.go b/internal/controlplane/api/creator_internal_coverage_test.go index c2668a3..9fe299c 100644 --- a/internal/controlplane/api/creator_internal_coverage_test.go +++ b/internal/controlplane/api/creator_internal_coverage_test.go @@ -8,6 +8,7 @@ import ( "net/http" "net/http/httptest" "os" + "path/filepath" "strings" "testing" "time" @@ -135,7 +136,7 @@ func TestReconcileRuntimeLeasesExported(t *testing.T) { } // TestCacheCompetitorWorkCovers 在隔离 PG 上驱动封面回填链路: -// 缺封面作品经浏览器 GetImage 拉取并落 creator_work_cover。 +// 缺封面作品经浏览器 GetImage 拉取并保存为静态文件。 func TestCacheCompetitorWorkCovers(t *testing.T) { databaseURL := os.Getenv("CREATORHUB_POSTGRES_TEST_URL") if databaseURL == "" { @@ -168,12 +169,16 @@ func TestCacheCompetitorWorkCovers(t *testing.T) { gateway: hub.Gateway{Endpoint: server.URL, Token: "token"}, environment: hub.EnvironmentContext{Env: hub.Env{Alias: "browser"}, RuntimeID: "dddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddd", BindingVersion: 1}, } - if err := cacheCompetitorWorkCovers(ctx, store, browser, competitor.ID); err != nil { + if err := cacheWorkCovers(ctx, store, browser, creator.SourceCompetitor, competitor.ID); err != nil { t.Fatalf("cache covers: %v", err) } - contentType, data, err := store.GetWorkCover(ctx, work.ID, "cover") - if err != nil || contentType != "image/png" || string(data) != "cover-bytes" { - t.Fatalf("cover: type=%q data=%q err=%v", contentType, data, err) + path, err := store.GetWorkCover(ctx, work.ID) + if err != nil { + t.Fatal(err) + } + data, err := os.ReadFile(path) + if err != nil || filepath.Ext(path) != ".png" || string(data) != "cover-bytes" { + t.Fatalf("cover: path=%q data=%q err=%v", path, data, err) } // 网关故障 → 报错且不落库(新增一件缺封面作品,否则无可拉取对象直接成功) @@ -202,10 +207,10 @@ func TestCacheCompetitorWorkCovers(t *testing.T) { })) defer failing.Close() browser.gateway.Endpoint = failing.URL - if err := cacheCompetitorWorkCovers(ctx, store, browser, other.ID); err == nil { + if err := cacheWorkCovers(ctx, store, browser, creator.SourceCompetitor, other.ID); err == nil { t.Fatal("failing gateway must surface an error") } - missing, err := store.ListWorksMissingCover(ctx, other.ID) + missing, err := store.ListWorksMissingCover(ctx, creator.SourceCompetitor, other.ID) if err != nil || len(missing) != 1 { t.Fatalf("failed fetch must not save covers: missing=%d err=%v", len(missing), err) } @@ -213,7 +218,19 @@ func TestCacheCompetitorWorkCovers(t *testing.T) { func openCreatorIntegrationStoreForAPITest(t *testing.T, databaseURL string) (*creator.Store, *account.Store, context.Context) { t.Helper() + t.Setenv("CREATOR_COVER_DIR", t.TempDir()) + databaseURL = isolatedControlPlaneDatabaseURL(t, databaseURL) ctx := context.Background() + accountStore, err := account.Open(ctx, databaseURL) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = accountStore.Close() }) + hubStore, err := hub.Open(ctx, databaseURL) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = hubStore.Close() }) creatorStore, err := creator.Open(ctx, databaseURL) if err != nil { t.Fatal(err) @@ -223,13 +240,7 @@ func openCreatorIntegrationStoreForAPITest(t *testing.T, databaseURL string) (*c t.Fatal(err) } creatorStore.SetSecretBridge(nil) - // hub.Open 铺迁移链(social_account 等 phase A 表)。 - hubStore, err := hub.Open(ctx, databaseURL) - if err != nil { - t.Fatal(err) - } - t.Cleanup(func() { _ = hubStore.Close() }) - return creatorStore, nil, ctx + return creatorStore, accountStore, ctx } // TestNewAnonymousBrowserCreateFailureCleanup 驱动匿名浏览器创建失败路径: diff --git a/internal/controlplane/api/owned_cover_cache_test.go b/internal/controlplane/api/owned_cover_cache_test.go new file mode 100644 index 0000000..4fb8229 --- /dev/null +++ b/internal/controlplane/api/owned_cover_cache_test.go @@ -0,0 +1,97 @@ +package api + +import ( + "bytes" + "encoding/base64" + "encoding/json" + "fmt" + "io" + "net/http" + "net/http/httptest" + "os" + "path/filepath" + "strings" + "testing" + "time" + + "git.ipao.vip/rogee/creator-hub/internal/account" + "git.ipao.vip/rogee/creator-hub/internal/creator" + hub "git.ipao.vip/rogee/creator-hub/internal/environment" + "github.com/gofiber/fiber/v3" +) + +func TestOwnedCoversCacheFilesAndReportIndividualFailures(t *testing.T) { + databaseURL := os.Getenv("CREATORHUB_POSTGRES_TEST_URL") + if databaseURL == "" { + t.Skip("set CREATORHUB_POSTGRES_TEST_URL for isolated PostgreSQL coverage") + } + store, accounts, ctx := openCreatorIntegrationStoreForAPITest(t, databaseURL) + owner := account.Account{ + ID: "owned-cover-account", Name: "封面账号", Platform: creator.PlatformDouyin, PlatformAccountKey: "12345678901", + CredentialReference: account.CredentialReference{ID: "cover-test-credential", Provider: "os_keyring"}, + CredentialKey: "creatorhub/owned-cover-account/cookies", + } + if err := accounts.CreateAccount(ctx, owner, &testCredentialBridge{values: make(map[string]string)}); err != nil { + t.Fatal(err) + } + works := make(map[string]creator.Work) + for _, key := range []string{"good", "broken"} { + work, _, err := store.UpsertWork(ctx, creator.WorkInput{SourceType: creator.SourceOwned, SourceID: owner.ID, Platform: creator.PlatformDouyin, WorkKey: "owned-cover-" + key, CoverURL: "https://p3-pc-sign.douyinpic.com/" + key + ".png"}, time.Now().UTC()) + if err != nil { + t.Fatal(err) + } + works[key] = work + } + requests := 0 + gateway := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/v1/browsers/owned-cover-alias/douyin/image" { + http.NotFound(w, r) + return + } + requests++ + var payload struct { + URL string `json:"url"` + } + if err := json.NewDecoder(r.Body).Decode(&payload); err != nil { + t.Error(err) + w.WriteHeader(400) + return + } + if strings.Contains(payload.URL, "broken") { + w.WriteHeader(502) + _, _ = io.WriteString(w, "image upstream unavailable") + return + } + _ = json.NewEncoder(w).Encode(map[string]any{"status": 200, "content_type": "image/png", "body_base64": base64.StdEncoding.EncodeToString([]byte("owned-cover-bytes"))}) + })) + t.Cleanup(gateway.Close) + browser := creatorGatewayBrowser{gateway: hub.Gateway{Endpoint: gateway.URL, Token: "test-token"}, environment: hub.EnvironmentContext{Env: hub.Env{Alias: "owned-cover-alias"}, RuntimeID: "runtime-cover", RuntimeNetworkID: "network-cover", BindingVersion: 1}} + if err := cacheWorkCovers(ctx, store, browser, creator.SourceOwned, owner.ID); err == nil || !strings.Contains(err.Error(), works["broken"].ID) { + t.Fatalf("individual failure was hidden: %v", err) + } + if requests != 2 { + t.Fatalf("failure stopped other covers from downloading: %d", requests) + } + path, err := store.GetWorkCover(ctx, works["good"].ID) + if err != nil { + t.Fatal(err) + } + if filepath.Base(filepath.Dir(path)) != owner.PlatformAccountKey || filepath.Base(path) != "owned-cover-good.png" { + t.Fatalf("wrong UID/work-ID layout: %s", path) + } + missing, err := store.ListWorksMissingCover(ctx, creator.SourceOwned, owner.ID) + if err != nil || len(missing) != 1 || missing[0].ID != works["broken"].ID { + t.Fatalf("cache retry targets: %v %v", missing, err) + } + app := fiber.New() + RegisterCreator(app, store, accounts, nil) + response, err := app.Test(httptest.NewRequest(http.MethodGet, fmt.Sprintf("/api/creator/works/%s/cover", works["good"].ID), nil)) + if err != nil { + t.Fatal(err) + } + defer response.Body.Close() + body, err := io.ReadAll(response.Body) + if err != nil || response.StatusCode != 200 || !strings.HasPrefix(response.Header.Get("Content-Type"), "image/png") || !bytes.Equal(body, []byte("owned-cover-bytes")) { + t.Fatalf("static cover response: %d %s %q %v", response.StatusCode, response.Header.Get("Content-Type"), body, err) + } +} diff --git a/internal/creator/content.go b/internal/creator/content.go index 2e712b6..9bcf68a 100644 --- a/internal/creator/content.go +++ b/internal/creator/content.go @@ -76,11 +76,11 @@ const competitorColumns = `competitor_id, platform, platform_account_key, unique 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, favorited sql.NullInt64 - 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 { @@ -597,30 +597,6 @@ func (s *Store) ListWorks(ctx context.Context, filter WorkFilter) ([]Work, error return result, nil } -// ListWorksMissingCover 返回已采集到远程封面但本地尚无缓存的作品(用于封面回填)。 -// 单次返回上限 40 条,超出部分由后续同步继续回填。 -func (s *Store) ListWorksMissingCover(ctx context.Context, sourceID string) ([]Work, error) { - if sourceID == "" { - return nil, ErrInvalid - } - rows, err := s.db.QueryContext(ctx, workSelect+` WHERE source_type = $1 AND source_id = $2 AND cover_url <> '' - AND NOT EXISTS (SELECT 1 FROM creator_work_cover c WHERE c.work_id = creator_work.id AND c.variant = 'cover') - ORDER BY created_at DESC LIMIT 40`, SourceCompetitor, sourceID) - if err != nil { - return nil, databaseError(err) - } - defer rows.Close() - result := make([]Work, 0) - for rows.Next() { - item, err := scanWork(rows) - if err != nil { - return nil, err - } - result = append(result, item) - } - return result, rows.Err() -} - 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 || input.PlayCount != nil && *input.PlayCount < 0 { return MetricPoint{}, ErrInvalid @@ -773,38 +749,6 @@ func (s *Store) GetCommentByKey(ctx context.Context, platform, commentKey string return result, rowError(err) } -// SaveWorkCover 缓存作品封面图(带签名的 douyinpic URL 会过期,存字节后经本地接口展示)。 -func (s *Store) SaveWorkCover(ctx context.Context, workID, variant, contentType string, data []byte) error { - if workID == "" || variant == "" || contentType == "" || len(data) == 0 || len(data) > maxWorkCoverBytes { - return ErrInvalid - } - _, err := s.db.ExecContext(ctx, ` - INSERT INTO creator_work_cover (work_id, variant, content_type, data, fetched_at) - SELECT w.id, $2, $3, $4, now() FROM creator_work w WHERE w.work_id = $1 - ON CONFLICT (work_id, variant) DO UPDATE SET content_type = EXCLUDED.content_type, data = EXCLUDED.data, fetched_at = now()`, - workID, variant, contentType, data) - if err != nil { - return databaseError(err) - } - return nil -} - -func (s *Store) GetWorkCover(ctx context.Context, workID, variant string) (contentType string, data []byte, err error) { - if workID == "" || variant == "" { - return "", nil, ErrInvalid - } - err = s.db.QueryRowContext(ctx, `SELECT cover.content_type, cover.data FROM creator_work_cover cover JOIN creator_work w ON w.id = cover.work_id WHERE w.work_id = $1 AND cover.variant = $2`, workID, variant).Scan(&contentType, &data) - if errors.Is(err, sql.ErrNoRows) { - return "", nil, ErrNotFound - } - if err != nil { - return "", nil, databaseError(err) - } - return contentType, data, nil -} - -const maxWorkCoverBytes = 4 << 20 - func pageBounds(page, pageSize int) (int, int, error) { if page < 1 || pageSize < 1 || pageSize > 100 { return 0, 0, ErrInvalid diff --git a/internal/creator/coverage_integration_test.go b/internal/creator/coverage_integration_test.go index 9cb1e73..9556660 100644 --- a/internal/creator/coverage_integration_test.go +++ b/internal/creator/coverage_integration_test.go @@ -3,6 +3,8 @@ package creator import ( "errors" "fmt" + "os" + "path/filepath" "testing" "time" ) @@ -15,7 +17,7 @@ func TestCreatorPostgresWorkCoverCache(t *testing.T) { HomepageURL: "https://www.douyin.com/user/cover" + stamp, }) if err != nil { - t.Fatalf("upsert competitor: %v", err) + t.Fatal(err) } now := time.Now().UTC().Add(-time.Hour) work, inserted, err := store.UpsertWork(ctx, WorkInput{ @@ -26,37 +28,32 @@ func TestCreatorPostgresWorkCoverCache(t *testing.T) { if err != nil || !inserted { t.Fatalf("upsert work: inserted=%v err=%v", inserted, err) } - // 采集到封面 URL 但尚未缓存 → 出现在回填列表。 - missing, err := store.ListWorksMissingCover(ctx, competitor.ID) + missing, err := store.ListWorksMissingCover(ctx, SourceCompetitor, competitor.ID) if err != nil || len(missing) != 1 || missing[0].ID != work.ID { - t.Fatalf("list missing covers: works=%+v err=%v", missing, err) + t.Fatalf("missing covers: %v %v", missing, err) } - if err := store.SaveWorkCover(ctx, work.ID, "cover", "image/jpeg", []byte("jpeg-bytes")); err != nil { - t.Fatalf("save cover: %v", err) + for _, item := range []struct{ mediaType, extension, bytes string }{{"image/jpeg", ".jpg", "jpeg-bytes"}, {"image/webp", ".webp", "webp-bytes"}} { + if err := store.SaveWorkCover(ctx, work.ID, item.mediaType, []byte(item.bytes)); err != nil { + t.Fatal(err) + } + path, err := store.GetWorkCover(ctx, work.ID) + if err != nil { + t.Fatal(err) + } + data, err := os.ReadFile(path) + if err != nil || filepath.Ext(path) != item.extension || string(data) != item.bytes { + t.Fatalf("cover: %s %q %v", path, data, err) + } + missing, err = store.ListWorksMissingCover(ctx, SourceCompetitor, competitor.ID) + if err != nil || len(missing) != 0 { + t.Fatalf("cached works still missing: %v %v", missing, err) + } } - contentType, data, err := store.GetWorkCover(ctx, work.ID, "cover") - if err != nil || contentType != "image/jpeg" || string(data) != "jpeg-bytes" { - t.Fatalf("get cover: contentType=%q data=%q err=%v", contentType, data, err) + if _, err := store.GetWorkCover(ctx, "missing"); !errors.Is(err, ErrNotFound) { + t.Fatalf("missing cover: %v", err) } - // 缓存后不再出现在回填列表。 - missing, err = store.ListWorksMissingCover(ctx, competitor.ID) - if err != nil || len(missing) != 0 { - t.Fatalf("missing covers after cache: works=%+v err=%v", missing, err) - } - // 重复写入覆盖旧内容(新采集周期封面可能更新)。 - if err := store.SaveWorkCover(ctx, work.ID, "cover", "image/webp", []byte("webp-bytes")); err != nil { - t.Fatalf("overwrite cover: %v", err) - } - contentType, data, err = store.GetWorkCover(ctx, work.ID, "cover") - if err != nil || contentType != "image/webp" || string(data) != "webp-bytes" { - t.Fatalf("overwritten cover: contentType=%q data=%q err=%v", contentType, data, err) - } - if _, _, err := store.GetWorkCover(ctx, "missing", "cover"); !errors.Is(err, ErrNotFound) { - t.Fatalf("missing cover must be ErrNotFound: err=%v", err) - } - // 非法输入拒绝。 - if err := store.SaveWorkCover(ctx, work.ID, "cover", "image/jpeg", nil); !errors.Is(err, ErrInvalid) { - t.Fatalf("empty cover data must be rejected: err=%v", err) + if err := store.SaveWorkCover(ctx, work.ID, "image/jpeg", nil); !errors.Is(err, ErrInvalid) { + t.Fatalf("empty cover accepted: %v", err) } } diff --git a/internal/creator/covers.go b/internal/creator/covers.go new file mode 100644 index 0000000..b218ed7 --- /dev/null +++ b/internal/creator/covers.go @@ -0,0 +1,157 @@ +package creator + +import ( + "context" + "errors" + "fmt" + "mime" + "os" + "path/filepath" + "regexp" + + "github.com/sirupsen/logrus" +) + +const maxWorkCoverBytes = 4 << 20 + +var coverComponent = regexp.MustCompile(`^[A-Za-z0-9_-]+$`) +var coverExtensions = map[string]string{ + "image/jpeg": ".jpg", "image/png": ".png", "image/webp": ".webp", "image/gif": ".gif", "image/avif": ".avif", +} + +func (s *Store) workCoverBase(ctx context.Context, workID string) (string, error) { + if workID == "" { + return "", ErrInvalid + } + var uid, workKey string + err := s.db.QueryRowContext(ctx, ` + SELECT COALESCE(a.platform_account_key, c.platform_account_key, ''), w.work_key + FROM creator_work w + LEFT JOIN social_account a ON w.source_type = 'owned' AND a.account_id = w.source_id + LEFT JOIN creator_competitor c ON w.source_type = 'competitor' AND c.competitor_id = w.source_id + WHERE w.work_id = $1`, workID).Scan(&uid, &workKey) + if err != nil { + return "", rowError(err) + } + if !coverComponent.MatchString(uid) || !coverComponent.MatchString(workKey) { + return "", fmt.Errorf("%w: cover requires a valid account UID and platform work ID", ErrInvalid) + } + return filepath.Join(s.coverDirectory, uid, workKey), nil +} + +// SaveWorkCover writes /. atomically. +func (s *Store) SaveWorkCover(ctx context.Context, workID, contentType string, data []byte) error { + mediaType, _, err := mime.ParseMediaType(contentType) + extension, supported := coverExtensions[mediaType] + if err != nil || !supported || len(data) == 0 || len(data) > maxWorkCoverBytes { + return ErrInvalid + } + base, err := s.workCoverBase(ctx, workID) + if err != nil { + return err + } + if err := os.MkdirAll(filepath.Dir(base), 0755); err != nil { + return fmt.Errorf("create cover directory: %w", err) + } + file, err := os.CreateTemp(filepath.Dir(base), ".cover-*") + if err != nil { + return fmt.Errorf("create cover file: %w", err) + } + defer func() { + if err := os.Remove(file.Name()); err != nil && !errors.Is(err, os.ErrNotExist) { + logrus.WithError(err).WithField("file", file.Name()).Warn("temporary work cover cleanup failed") + } + }() + if _, err := file.Write(data); err != nil { + return errors.Join(fmt.Errorf("write cover: %w", err), file.Close()) + } + if err := file.Sync(); err != nil { + return errors.Join(fmt.Errorf("flush cover: %w", err), file.Close()) + } + if err := file.Close(); err != nil { + return fmt.Errorf("close cover: %w", err) + } + if err := os.Rename(file.Name(), base+extension); err != nil { + return fmt.Errorf("publish cover: %w", err) + } + for _, other := range coverExtensions { + if other == extension { + continue + } + if err := os.Remove(base + other); err != nil && !errors.Is(err, os.ErrNotExist) { + return fmt.Errorf("remove replaced cover: %w", err) + } + } + return nil +} + +// GetWorkCover returns the local file; image bytes are never stored in PostgreSQL. +func (s *Store) GetWorkCover(ctx context.Context, workID string) (string, error) { + base, err := s.workCoverBase(ctx, workID) + if err != nil { + return "", err + } + var result string + for _, extension := range coverExtensions { + path := base + extension + info, err := os.Stat(path) + if errors.Is(err, os.ErrNotExist) { + continue + } + if err != nil { + return "", fmt.Errorf("stat cover: %w", err) + } + if !info.Mode().IsRegular() || info.Size() == 0 { + return "", fmt.Errorf("%w: cover file is not a nonempty regular file", ErrInvalid) + } + if result != "" { + return "", fmt.Errorf("%w: multiple cover formats exist for work %s", ErrConflict, workID) + } + result = path + } + if result == "" { + return "", ErrNotFound + } + return result, nil +} + +// ListWorksMissingCover backfills both owned and monitored sources, at most 40 per sync. +func (s *Store) ListWorksMissingCover(ctx context.Context, sourceType, sourceID string) ([]Work, error) { + if (sourceType != SourceOwned && sourceType != SourceCompetitor) || sourceID == "" { + return nil, ErrInvalid + } + rows, err := s.db.QueryContext(ctx, workSelect+` WHERE source_type = $1 AND source_id = $2 AND cover_url <> '' ORDER BY created_at DESC`, sourceType, sourceID) + if err != nil { + return nil, databaseError(err) + } + defer rows.Close() + works := make([]Work, 0) + for rows.Next() { + work, err := scanWork(rows) + if err != nil { + return nil, err + } + works = append(works, work) + } + if err := rows.Err(); err != nil { + return nil, err + } + if err := rows.Close(); err != nil { + return nil, err + } + missing := make([]Work, 0) + for _, work := range works { + _, err := s.GetWorkCover(ctx, work.ID) + if err == nil { + continue + } + if !errors.Is(err, ErrNotFound) { + return nil, err + } + missing = append(missing, work) + if len(missing) == 40 { + break + } + } + return missing, nil +} diff --git a/internal/creator/covers_integration_test.go b/internal/creator/covers_integration_test.go new file mode 100644 index 0000000..fca9f2a --- /dev/null +++ b/internal/creator/covers_integration_test.go @@ -0,0 +1,113 @@ +package creator + +import ( + "errors" + "os" + "path/filepath" + "testing" + "time" +) + +func TestWorkCoversUseUIDWorkIDFilesForOwnedAndCompetitor(t *testing.T) { + store, accounts, ctx := openCreatorIntegrationStore(t) + accountID := createIntegrationAccount(t, ctx, accounts, "cover-files") + owned, err := accounts.GetAccount(ctx, accountID) + if err != nil { + t.Fatal(err) + } + competitor, err := store.CreateCompetitor(ctx, CompetitorInput{Platform: PlatformDouyin, PlatformAccountKey: "cover-target-uid", Nickname: "封面账号", HomepageURL: "https://www.douyin.com/user/cover-target"}) + if err != nil { + t.Fatal(err) + } + for _, source := range []struct{ kind, id, uid string }{{SourceOwned, owned.ID, owned.PlatformAccountKey}, {SourceCompetitor, competitor.ID, competitor.PlatformAccountKey}} { + work, _, err := store.UpsertWork(ctx, WorkInput{SourceType: source.kind, SourceID: source.id, Platform: PlatformDouyin, WorkKey: "7312345678901234567-" + source.kind, Title: "封面作品", CoverURL: "https://p3.douyinpic.com/cover"}, time.Now().UTC()) + if err != nil { + t.Fatal(err) + } + missing, err := store.ListWorksMissingCover(ctx, source.kind, source.id) + if err != nil || len(missing) != 1 { + t.Fatalf("missing: %v %v", missing, err) + } + if _, err := store.GetWorkCover(ctx, work.ID); !errors.Is(err, ErrNotFound) { + t.Fatalf("missing cover: %v", err) + } + if err := store.SaveWorkCover(ctx, work.ID, "image/png", []byte("first-cover")); err != nil { + t.Fatal(err) + } + path, err := store.GetWorkCover(ctx, work.ID) + if err != nil { + t.Fatal(err) + } + want := filepath.Join(store.coverDirectory, source.uid, work.WorkKey+".png") + if path != want { + t.Fatalf("cover path=%q want=%q", path, want) + } + data, err := os.ReadFile(path) + if err != nil || string(data) != "first-cover" { + t.Fatalf("file: %q %v", data, err) + } + missing, err = store.ListWorksMissingCover(ctx, source.kind, source.id) + if err != nil || len(missing) != 0 { + t.Fatalf("cached work still missing: %v %v", missing, err) + } + if err := store.SaveWorkCover(ctx, work.ID, "image/webp", []byte("replacement")); err != nil { + t.Fatal(err) + } + if _, err := os.Stat(want); !errors.Is(err, os.ErrNotExist) { + t.Fatalf("old extension still exists: %v", err) + } + path, err = store.GetWorkCover(ctx, work.ID) + if err != nil || filepath.Ext(path) != ".webp" { + t.Fatalf("replacement path: %q %v", path, err) + } + data, err = os.ReadFile(path) + if err != nil || string(data) != "replacement" { + t.Fatalf("replacement: %q %v", data, err) + } + } + var obsoleteTableExists bool + if err := store.db.QueryRowContext(ctx, `SELECT to_regclass('creator_work_cover') IS NOT NULL`).Scan(&obsoleteTableExists); err != nil || obsoleteTableExists { + t.Fatalf("image-byte table remains: %v %v", obsoleteTableExists, err) + } +} + +func TestWorkCoverFilesExposeFailuresAndRejectInvalidNames(t *testing.T) { + store, accounts, ctx := openCreatorIntegrationStore(t) + accountID := createIntegrationAccount(t, ctx, accounts, "cover-errors") + work, _, err := store.UpsertWork(ctx, WorkInput{SourceType: SourceOwned, SourceID: accountID, Platform: PlatformDouyin, WorkKey: "file-errors", CoverURL: "https://p3.douyinpic.com/cover"}, time.Now().UTC()) + if err != nil { + t.Fatal(err) + } + for _, item := range []struct { + contentType string + data []byte + }{{"text/html", []byte("wrong")}, {"image/png", nil}, {"image/png", make([]byte, maxWorkCoverBytes+1)}} { + if err := store.SaveWorkCover(ctx, work.ID, item.contentType, item.data); !errors.Is(err, ErrInvalid) { + t.Fatalf("invalid image accepted: %v", err) + } + } + if _, err := store.GetWorkCover(ctx, "missing-work"); !errors.Is(err, ErrNotFound) { + t.Fatalf("unknown work: %v", err) + } + blocked := filepath.Join(t.TempDir(), "not-a-directory") + if err := os.WriteFile(blocked, []byte("blocked"), 0600); err != nil { + t.Fatal(err) + } + store.coverDirectory = blocked + if err := store.SaveWorkCover(ctx, work.ID, "image/png", []byte("cover")); err == nil { + t.Fatal("filesystem error was swallowed") + } + if _, err := store.GetWorkCover(ctx, work.ID); err == nil || errors.Is(err, ErrNotFound) { + t.Fatalf("filesystem error treated as normal cache miss: %v", err) + } + bad, _, err := store.UpsertWork(ctx, WorkInput{SourceType: SourceOwned, SourceID: accountID, Platform: PlatformDouyin, WorkKey: "../outside"}, time.Now().UTC()) + if err != nil { + t.Fatal(err) + } + if err := store.SaveWorkCover(ctx, bad.ID, "image/png", []byte("cover")); !errors.Is(err, ErrInvalid) { + t.Fatalf("invalid work filename accepted: %v", err) + } + if _, err := store.ListWorksMissingCover(ctx, "wrong-source", accountID); !errors.Is(err, ErrInvalid) { + t.Fatalf("invalid source: %v", err) + } +} diff --git a/internal/creator/integration_test.go b/internal/creator/integration_test.go index a690925..623d61e 100644 --- a/internal/creator/integration_test.go +++ b/internal/creator/integration_test.go @@ -64,6 +64,7 @@ func (c integrationCollector) ListTopLevelComments(_ context.Context, _, cursor func openCreatorIntegrationStore(t *testing.T) (*Store, *account.Store, context.Context) { t.Helper() + t.Setenv("CREATOR_COVER_DIR", t.TempDir()) databaseURL := os.Getenv("CREATORHUB_POSTGRES_TEST_URL") if databaseURL == "" { t.Skip("set CREATORHUB_POSTGRES_TEST_URL to run CreatorHub PostgreSQL integration coverage") diff --git a/internal/creator/store.go b/internal/creator/store.go index f513664..320a50f 100644 --- a/internal/creator/store.go +++ b/internal/creator/store.go @@ -8,6 +8,8 @@ import ( "encoding/json" "errors" "fmt" + "os" + "path/filepath" "sync" "time" @@ -29,8 +31,9 @@ type SecretBridge interface { } type Store struct { - db *sql.DB - secrets SecretBridge + db *sql.DB + coverDirectory string + secrets SecretBridge // Automatic writes for one execution account stay serialized while event // ingestion and unrelated accounts remain concurrent. automaticLocks sync.Map // map[string]*sync.Mutex @@ -48,7 +51,16 @@ func Open(ctx context.Context, databaseURL string) (*Store, error) { _ = db.Close() return nil, errors.New("connect to creator database") } - store := &Store{db: db} + directory := os.Getenv("CREATOR_COVER_DIR") + if directory == "" { + directory = filepath.Join(".data", "covers") + } + absolute, err := filepath.Abs(directory) + if err != nil { + _ = db.Close() + return nil, fmt.Errorf("resolve creator cover directory: %w", err) + } + store := &Store{db: db, coverDirectory: absolute} return store, nil } diff --git a/internal/environment/migration043_probe_test.go b/internal/environment/migration043_probe_test.go index 34eeb96..ae0173d 100644 --- a/internal/environment/migration043_probe_test.go +++ b/internal/environment/migration043_probe_test.go @@ -29,19 +29,19 @@ func TestMigration043ConsolidatedSchemaShape(t *testing.T) { } defer db.Close() - // 最终 19 张表 + // 封面改为静态文件,最终 18 张表 assertDatabaseCount(t, db, `SELECT count(*) FROM information_schema.tables WHERE table_schema = current_schema() AND table_name IN ('social_account','gateway','network_exit','browser_env','audit_event', - 'creator_competitor','creator_competitor_share_job','creator_work','creator_work_metric','creator_work_cover', + 'creator_competitor','creator_competitor_share_job','creator_work','creator_work_metric', 'creator_comment','creator_comment_rule_result','creator_lead_rule','creator_account_profile','creator_account_password', - 'creator_account_metric','creator_collection_checkpoint','creator_settings','schema_migration')`, 19) + 'creator_account_metric','creator_collection_checkpoint','creator_settings','schema_migration')`, 18) // 已删表不再存在 assertDatabaseCount(t, db, `SELECT count(*) FROM information_schema.tables WHERE table_schema = current_schema() AND table_name IN ('runtime_instance','environment_binding','credential_reference','content_draft','confirmation', 'operation_task','execution_attempt','runtime_use_lease','creator_event','creator_strategy','creator_operation', 'creator_conversation','creator_message','creator_cooldown','creator_listener_state','creator_listener_boundary', 'creator_relation','creator_material_job','creator_event_strategy_trace','creator_work_source','creator_metric_plan', - 'creator_source_sync_lease','creator_schema_migration')`, 0) + 'creator_source_sync_lease','creator_schema_migration','creator_work_cover')`, 0) // browser_env 合并列 assertDatabaseCount(t, db, `SELECT count(*) FROM information_schema.columns WHERE table_schema = current_schema() AND table_name = 'browser_env' AND column_name IN ('account_id','exit_id','gateway_id','runtime_cleanup_pending', @@ -70,8 +70,8 @@ func TestMigration043ConsolidatedSchemaShape(t *testing.T) { // audit_event 任务列已删 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) - // 统一登记表包含后续的首次登录 UID 绑定迁移。 - assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 55) + // 统一登记表包含首次登录 UID 绑定及封面静态文件迁移。 + assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 56) assertDatabaseCount(t, db, `SELECT count(*) FROM information_schema.columns WHERE table_schema = current_schema() AND table_name = 'social_account' AND column_name = 'platform_account_key' AND is_nullable = 'YES'`, 1) } diff --git a/internal/environment/migration_test.go b/internal/environment/migration_test.go index c932302..aa4273a 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`, 55) + assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 56) 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`, 55) + assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 56) }) t.Run("legacy migration 013 without account secrets is repaired forward", func(t *testing.T) { diff --git a/internal/environment/migrations/1049_work_covers_static_files.sql b/internal/environment/migrations/1049_work_covers_static_files.sql new file mode 100644 index 0000000..5210fe7 --- /dev/null +++ b/internal/environment/migrations/1049_work_covers_static_files.sql @@ -0,0 +1,3 @@ +-- Work cover bytes are now stored as /. files. +-- Export existing covers before applying this migration; no database fallback is kept. +DROP TABLE creator_work_cover; diff --git a/internal/environment/store.go b/internal/environment/store.go index 4d5e4ee..dba7e2b 100644 --- a/internal/environment/store.go +++ b/internal/environment/store.go @@ -185,6 +185,9 @@ var migration1047 string //go:embed migrations/1048_environment_first_login.sql var migration1048 string +//go:embed migrations/1049_work_covers_static_files.sql +var migration1049 string + var ( ErrConflict = errors.New("resource conflicts with existing state") ErrInvalid = errors.New("invalid environment input") @@ -325,7 +328,7 @@ func (s *Store) migrate(ctx context.Context) error { {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}, {1046, migration1046}, - {43, migration043}, {44, migration044}, {1047, migration1047}, {1048, migration1048}} { + {43, migration043}, {44, migration044}, {1047, migration1047}, {1048, migration1048}, {1049, migration1049}} { var applied bool if err := tx.QueryRowContext(ctx, `SELECT EXISTS (SELECT 1 FROM schema_migration WHERE version = $1)`, migration.version).Scan(&applied); err != nil { return errors.New("read environment schema migration state") diff --git a/scripts/dev-backend.mjs b/scripts/dev-backend.mjs index 571c331..d93dc6a 100755 --- a/scripts/dev-backend.mjs +++ b/scripts/dev-backend.mjs @@ -34,6 +34,7 @@ if (!overrides.DATABASE_URL) { overrides.DATABASE_URL = `postgres://creatorhub@127.0.0.1:${postgresPort}/creatorhub?sslmode=disable`; } if (!overrides.WEB_DIR) overrides.WEB_DIR = path.join(root, "web", "dist"); +if (!overrides.CREATOR_COVER_DIR) overrides.CREATOR_COVER_DIR = path.join(root, ".data", "covers"); if (!overrides.CREATORHUB_CREDENTIAL_STORE_DIR) overrides.CREATORHUB_CREDENTIAL_STORE_DIR = path.join(root, ".dev-credentials"); if (!overrides.LOG_LEVEL) overrides.LOG_LEVEL = "debug"; if (!overrides.NATIVE_GATEWAY_ENDPOINT) overrides.NATIVE_GATEWAY_ENDPOINT = "http://127.0.0.1:28187";