fix: cache owned work covers as UID-scoped static image files

This commit is contained in:
2026-10-06 11:56:01 +08:00
parent c213013a27
commit 4180ec104c
18 changed files with 485 additions and 134 deletions
+3
View File
@@ -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
+1
View File
@@ -34,6 +34,7 @@ __pycache__/
.env.*
!.env.example
.dev-credentials/
.data/
# 工具缓存与本机数据
.codegraph/
+2
View File
@@ -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;
+3
View File
@@ -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:
+22 -19
View File
@@ -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 拉取登录账号自己的画像并记一次时序快照。
@@ -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 驱动匿名浏览器创建失败路径:
@@ -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)
}
}
+5 -61
View File
@@ -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
+25 -28
View File
@@ -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)
}
}
+157
View File
@@ -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 <account UID>/<platform work ID>.<image extension> 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
}
+113
View File
@@ -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)
}
}
+1
View File
@@ -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")
+15 -3
View File
@@ -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
}
@@ -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)
}
+2 -2
View File
@@ -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) {
@@ -0,0 +1,3 @@
-- Work cover bytes are now stored as <UID>/<platform work ID>.<extension> files.
-- Export existing covers before applying this migration; no database fallback is kept.
DROP TABLE creator_work_cover;
+4 -1
View File
@@ -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")
+1
View File
@@ -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";