diff --git a/browser_gateway/platform/douyin.py b/browser_gateway/platform/douyin.py index a511a11..ecbd1a6 100644 --- a/browser_gateway/platform/douyin.py +++ b/browser_gateway/platform/douyin.py @@ -44,12 +44,14 @@ LOGIN_ORIGINS = frozenset( LOGIN_SCREENSHOT_LIMIT = 8 << 20 LOGIN_RENDER_TIMEOUT = 12.0 LOGIN_RENDER_MIN_BYTES = 20 << 10 +ORIGIN_NAVIGATE_TIMEOUT = 20.0 IDENTITY_URL = ( ORIGIN + "/aweme/v1/web/user/profile/self/?aid=6383&device_platform=webapp" ) WORKS_PATH = "/aweme/v1/web/aweme/post/" COMMENTS_PATH = "/aweme/v1/web/comment/list/" -RESPONSE_LIMIT = 1 << 20 +# 作品/评论响应带播放信息等大字段,单页常超 1MB;media 走独立管线不在此限。 +RESPONSE_LIMIT = 8 << 20 MEDIA_RESPONSE_LIMIT = 32 << 20 MEDIA_SOURCE_WAIT_MS = 15000 UID_RE = re.compile(r"^[1-9][0-9]{0,19}$") @@ -205,19 +207,80 @@ class DouyinBrowser: raise DouyinError("browser CDP connection failed") from exc return CDPConnection(socket) + def _navigate_to_origin(self, cdp: CDPConnection) -> None: + """把浏览器带到抖音首页。 + + 匿名浏览器从 about:blank 启动(origin 为 null),必须先导航到抖音才能在 + 该 origin 下 fetch API。首页 JS 会在落地后异步写入访客 cookie(ttwid 等) + 并完成风控环境初始化:过早发起 API fetch 只会拿到空/截断响应 + (实测 <3s 必空,~3s 后才有全量数据)。等 ttwid 出现且后续 cookie + 不再变化后才返回;超时或连接异常不抛错,放行 fetch 让它给出真实错误。 + """ + cdp.command("Page.enable") + cdp.notify("Page.navigate", {"url": ORIGIN_URL}) + deadline = time.monotonic() + ORIGIN_NAVIGATE_TIMEOUT + # 阶段一:等页面到达抖音 origin。 + while time.monotonic() < deadline: + try: + if cdp.evaluate("location.origin") == self.origin: + break + except DouyinError: + return + time.sleep(0.5) + else: + return + # 阶段二:等访客 cookie 稳定(连续 1s 不再变化)。 + settled_since = None + last_cookies = None + while time.monotonic() < deadline: + try: + cookies = cdp.evaluate("document.cookie") or "" + except DouyinError: + return + if "ttwid=" not in cookies: + settled_since, last_cookies = None, None + time.sleep(0.5) + continue + now = time.monotonic() + if cookies != last_cookies: + last_cookies, settled_since = cookies, now + elif now - settled_since >= 1.0: + return + time.sleep(0.5) + def get(self, alias: str, target: str) -> BrowserResponse: with self.connection(alias) as cdp: if cdp.evaluate("location.origin") != self.origin: - raise DouyinError("restricted browser origin changed") + # about:blank(匿名浏览器)或其它页面:先导航到抖音首页再 fetch。 + self._navigate_to_origin(cdp) expression = f"""(async()=>{{ - const r=await fetch({json.dumps(target)},{{credentials:'include',redirect:'error'}}); - if(!r.body)return {{status:r.status,body:'',too_large:false}}; - const reader=r.body.getReader(), decoder=new TextDecoder(); let size=0, body=''; - for(;;){{const item=await reader.read();if(item.done)break; - if(size+item.value.byteLength>={RESPONSE_LIMIT}){{await reader.cancel();return {{too_large:true}};}} - size+=item.value.byteLength;body+=decoder.decode(item.value,{{stream:true}}); + try {{ + // 抖音部分 API(如 aweme/post)要求 a_bogus 签名;首页加载的 byted_acrawler 提供签名器, + // 不存在或失败时退化为原请求(由服务端返回真实错误)。 + let requestUrl = {json.dumps(target)}; + try {{ + const signer = window.byted_acrawler && window.byted_acrawler.frontierSign; + if (typeof signer === 'function') {{ + const query = requestUrl.split('?')[1] || ''; + const sign = signer(query); + const extra = new URLSearchParams(); + if (sign && typeof sign === 'object') for (const key in sign) extra.set(key, String(sign[key])); + const signed = extra.toString(); + if (signed) requestUrl = requestUrl + (requestUrl.includes('?') ? '&' : '?') + signed; + }} + }} catch (signError) {{}} + const r=await fetch(requestUrl,{{credentials:'include',redirect:'error'}}); + if(!r.body)return {{status:r.status,body:'__empty_response__ signer='+(typeof window.byted_acrawler)+' hasFrontier='+(String(typeof (window.byted_acrawler&&window.byted_acrawler.frontierSign))),too_large:false}}; + const reader=r.body.getReader(), decoder=new TextDecoder(); let size=0, body=''; + for(;;){{const item=await reader.read();if(item.done)break; + if(size+item.value.byteLength>={RESPONSE_LIMIT}){{await reader.cancel();return {{too_large:true}};}} + size+=item.value.byteLength;body+=decoder.decode(item.value,{{stream:true}}); + }} + body+=decoder.decode();return {{status:r.status,body,too_large:false}}; + }} catch(e) {{ + // fetch 网络异常/被页面拦截时必须透出原因,否则 CDP 层只剩 undefined,无法定位。 + return {{status:0,body:'__fetch_error__:'+String((e&&e.message)||e),too_large:false}}; }} - body+=decoder.decode();return {{status:r.status,body,too_large:false}}; }})()""" result = None for attempt in range(2): @@ -235,6 +298,7 @@ class DouyinBrowser: or result.get("too_large") or not isinstance(result.get("status"), int) ): + LOG.warning("restricted browser fetch returned unexpected value: %r", result) raise DouyinError("restricted browser fetch failed") status = result["status"] if 300 <= status < 400: diff --git a/browser_gateway/server/http.py b/browser_gateway/server/http.py index f5cd900..ac513c8 100644 --- a/browser_gateway/server/http.py +++ b/browser_gateway/server/http.py @@ -62,6 +62,7 @@ NETWORK_ID_RE = _NETWORK_ID_RE RUNTIME_CLEANUP_SENTINEL = _RUNTIME_CLEANUP_SENTINEL EXIT_ID_RE = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._/-]{0,127}$") DOUYIN_ACCOUNT_KEY_RE = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._:@-]{0,127}$") +DOUYIN_SIGNATURE_RE = re.compile(r"^[A-Za-z0-9_-]{16,512}$") DOUYIN_ORIGIN = "https://www.douyin.com" DOUYIN_IDENTITY_PATH = "/aweme/v1/web/user/profile/self/" DOUYIN_PROFILE_OTHER_PATH = "/aweme/v1/web/user/profile/other/" @@ -749,13 +750,20 @@ def valid_douyin_profile_query(query: dict[str, list[str]]) -> bool: def valid_douyin_api_query( query: dict[str, list[str]], account_field: str, cursor_field: str ) -> bool: + # 浏览器内 byted_acrawler.frontierSign 会给 API 请求追匫a_bogus 签名参数(可选)。 + rest = query + if len(query) == 6 and len(query.get("a_bogus", [])) == 1: + if not DOUYIN_SIGNATURE_RE.fullmatch(query["a_bogus"][0]): + return False + rest = {key: value for key, value in query.items() if key != "a_bogus"} + if len(rest) != 5: + return False return ( - len(query) == 5 - and query.get("aid") == ["6383"] - and query.get("device_platform") == ["webapp"] - and valid_account_key_query(query, account_field) - and query.get("count") == ["20"] - and numeric_cursor(query.get(cursor_field)) + rest.get("aid") == ["6383"] + and rest.get("device_platform") == ["webapp"] + and valid_account_key_query(rest, account_field) + and rest.get("count") == ["20"] + and numeric_cursor(rest.get(cursor_field)) ) diff --git a/browser_gateway/test_gateway.py b/browser_gateway/test_gateway.py index afe253d..10fa03d 100644 --- a/browser_gateway/test_gateway.py +++ b/browser_gateway/test_gateway.py @@ -21,6 +21,7 @@ from .platform.douyin import ( DouyinBrowser, DouyinError, DouyinSubscription, + ORIGIN_URL, SubscriptionManager, ack_expression, action_expression, @@ -653,6 +654,9 @@ class BrowserCDP: return {"frameId": "frame-1"} return {} + def notify(self, method: str, params: dict | None = None) -> None: + self.commands.append((method, params)) + def wait_event(self, method: str, predicate: object, timeout: float = 15.0) -> dict: del predicate, timeout self.events.append(method) @@ -744,6 +748,29 @@ class BrowserTests(unittest.TestCase): self.assertEqual(response.status, 200) self.assertNotIn("Network.setCookies", [method for method, _ in cdp.commands]) + def test_browser_fetch_navigates_blank_runtime_to_origin(self) -> None: + # 匿名浏览器从 about:blank 启动(origin 为 null):get 前必须先导航到抖音首页, + # 并等待访客 cookie(ttwid)就绪,否则 API fetch 返回空响应。 + cdp = BrowserCDP( + [ + "null", + "https://www.douyin.com", + "ttwid=1%7Cabc; s_v_web_id=x", + "ttwid=1%7Cabc; s_v_web_id=x", + "ttwid=1%7Cabc; s_v_web_id=x", + {"status": 200, "body": "{}", "too_large": False}, + ] + ) + browser = DouyinBrowser() + self._with_connection(browser, cdp) + + response = browser.get( + "safe", "https://www.douyin.com/aweme/v1/web/user/profile/self/?aid=6383" + ) + self.assertEqual(response.status, 200) + self.assertIn(("Page.navigate", {"url": ORIGIN_URL}), cdp.commands) + self.assertEqual(cdp.events, []) + def test_browser_fetch_retries_transient_auth_response(self) -> None: cdp = BrowserCDP( [ diff --git a/internal/controlplane/api/app_migrated_test.go b/internal/controlplane/api/app_migrated_test.go index db50811..e207de4 100644 --- a/internal/controlplane/api/app_migrated_test.go +++ b/internal/controlplane/api/app_migrated_test.go @@ -332,8 +332,8 @@ func TestCreatorFixtureRoutesPostgres(t *testing.T) { t.Fatalf("competitor %s: %d %s", action, response.Code, response.Body.String()) } } - if _, err := syncCreatorCompetitorWithClaim(ctx, creatorStore, phaseAStore, hubStore, competitorID, "fixture-account", true); err == nil { - t.Fatal("competitor sync without browser unexpectedly succeeded") + if _, err := syncCreatorCompetitorWithClaim(ctx, creatorStore, hubStore, competitorID, true); err == nil { + t.Fatal("competitor sync without anonymous browser unexpectedly succeeded") } if _, err := previewDouyinCompetitor(ctx, creatorStore, phaseAStore, hubStore, "fixture-account", creator.CompetitorInput{Platform: creator.PlatformDouyin, PlatformAccountKey: "preview-key", Nickname: "Preview", HomepageURL: "https://www.douyin.com/user/preview-key"}); err == nil { t.Fatal("competitor preview without browser unexpectedly succeeded") diff --git a/internal/controlplane/api/creator.go b/internal/controlplane/api/creator.go index db03a6f..8581b73 100644 --- a/internal/controlplane/api/creator.go +++ b/internal/controlplane/api/creator.go @@ -16,6 +16,7 @@ import ( "path/filepath" "strconv" "strings" + "sync" "time" accountdomain "git.ipao.vip/rogee/creator-hub/internal/account" @@ -347,13 +348,8 @@ func registerCreatorWithServices(app *fiber.App, store *creator.Store, phaseASto return c.JSON(item) }) app.Post("/api/creator/competitors/:id/sync", func(c fiber.Ctx) error { - var input struct { - AccountID string `json:"account_id"` - } - if err := decodeCreator(c, &input); err != nil { - return creatorError(c, err) - } - report, err := syncCreatorCompetitor(c.Context(), store, phaseAStore, hubStore, c.Params("id"), input.AccountID) + // 强制同步走临时匿名浏览器,不需要自有账号身份。 + report, err := syncCreatorCompetitor(c.Context(), store, hubStore, c.Params("id")) if err != nil { return creatorError(c, err) } @@ -1371,7 +1367,8 @@ func verifyCreatorAccount(ctx context.Context, store *creator.Store, phaseAStore func (browser creatorGatewayBrowser) Get(ctx context.Context, target string) (douyin.Response, error) { payload := gatewayGenerationPayload(browser.environment) payload["url"] = target - status, body, err := gatewayCall(ctx, browser.gateway, http.MethodPost, "/v1/browsers/"+url.PathEscape(browser.environment.Alias)+"/douyin/get", payload, 30*time.Second) + // 网关側 socket 读超时 5s,1MB+ 响应体的全量传输需要分多次读:客户端超时须远大于服务端单次读超时。 + status, body, err := gatewayCallWithLimit(ctx, browser.gateway, http.MethodPost, "/v1/browsers/"+url.PathEscape(browser.environment.Alias)+"/douyin/get", payload, 120*time.Second, largeGatewayResponseLimit) if err != nil { return douyin.Response{}, err } @@ -1528,9 +1525,16 @@ func previewCompetitorShare(ctx context.Context, store *creator.Store, phaseASto type anonymousBrowserLease struct { gateway hub.Gateway environment hub.EnvironmentContext + // slot 占用全局匿名浏览器并发名额;非 nil 时 close 必须释放,保证任何路径都不会泄漏名额。 + slot chan struct{} } func (lease anonymousBrowserLease) close() error { + defer func() { + if lease.slot != nil { + <-lease.slot + } + }() cleanupContext, cancel := context.WithTimeout(context.Background(), 90*time.Second) defer cancel() payload := gatewayCleanupGenerationPayload(lease.environment) @@ -1575,6 +1579,11 @@ func randomGateway(gateways []hub.Gateway) hub.Gateway { return gateways[rand.Intn(len(gateways))] } +// 匿名浏览器全局并发上限:每个实例都是一个完整的指纹浏览器,同时拉起过多会耗尽机器资源。 +// 所有路径(竞品同步、share 链接解析、调度器并发)共享同一份名额。 +// ponytail: 固定值 2 够本地单机使用;需要按机器规格调整时再提升为配置项。 +var anonymousBrowserSlots = make(chan struct{}, 2) + func newAnonymousBrowser(ctx context.Context, store *hub.Store) (anonymousBrowserLease, error) { if store == nil { return anonymousBrowserLease{}, creator.ErrUnavailable @@ -1591,6 +1600,13 @@ func newAnonymousBrowser(ctx context.Context, store *hub.Store) (anonymousBrowse return anonymousBrowserLease{}, fmt.Errorf("%w: anonymous browser gateway is not configured", creator.ErrUnavailable) } + // 拿不到并发名额就在此等待(ctx 取消可退出),保证浏览器实例数永不起过上限。 + select { + case anonymousBrowserSlots <- struct{}{}: + case <-ctx.Done(): + return anonymousBrowserLease{}, ctx.Err() + } + gateway := randomGateway(gateways) var template hub.Env for _, candidate := range envs { @@ -1619,7 +1635,7 @@ func newAnonymousBrowser(ctx context.Context, store *hub.Store) (anonymousBrowse var created runtimeStatus if json.Unmarshal(body, &created) == nil && created.Alias == environment.Alias && created.BindingVersion == environment.BindingVersion && created.State == "running" && validCreatedRuntime(created, environment, true) { environment.RuntimeID, environment.RuntimeNetworkID = created.ID, created.NetworkID - return anonymousBrowserLease{gateway: gateway, environment: environment}, nil + return anonymousBrowserLease{gateway: gateway, environment: environment, slot: anonymousBrowserSlots}, nil } callErr = errors.New("gateway returned an invalid anonymous runtime generation") } else if callErr == nil { @@ -1628,6 +1644,7 @@ func newAnonymousBrowser(ctx context.Context, store *hub.Store) (anonymousBrowse callErr = gatewayUnreachable(callErr) } cleanupErr := cleanupAnonymousRuntimeAfterCreateFailure(ctx, gateway, environment) + <-anonymousBrowserSlots return anonymousBrowserLease{}, errors.Join(callErr, cleanupErr) } @@ -1750,16 +1767,16 @@ func previewDouyinCompetitor(ctx context.Context, store *creator.Store, phaseASt }, nil } -func syncCreatorCompetitor(ctx context.Context, store *creator.Store, phaseAStore *accountdomain.Store, hubStore *hub.Store, competitorID, accountID string) (creator.CollectionReport, error) { - return syncCreatorCompetitorWithClaim(ctx, store, phaseAStore, hubStore, competitorID, accountID, true) +func syncCreatorCompetitor(ctx context.Context, store *creator.Store, hubStore *hub.Store, competitorID string) (creator.CollectionReport, error) { + return syncCreatorCompetitorWithClaim(ctx, store, hubStore, competitorID, true) } -func syncCreatorCompetitorDue(ctx context.Context, store *creator.Store, phaseAStore *accountdomain.Store, hubStore *hub.Store, competitorID, accountID string) (creator.CollectionReport, error) { - return syncCreatorCompetitorWithClaim(ctx, store, phaseAStore, hubStore, competitorID, accountID, false) +func syncCreatorCompetitorDue(ctx context.Context, store *creator.Store, hubStore *hub.Store, competitorID string) (creator.CollectionReport, error) { + return syncCreatorCompetitorWithClaim(ctx, store, hubStore, competitorID, false) } -func syncCreatorCompetitorWithClaim(ctx context.Context, store *creator.Store, phaseAStore *accountdomain.Store, hubStore *hub.Store, competitorID, accountID string, force bool) (creator.CollectionReport, error) { - if store == nil || phaseAStore == nil || hubStore == nil || accountID == "" { +func syncCreatorCompetitorWithClaim(ctx context.Context, store *creator.Store, hubStore *hub.Store, competitorID string, force bool) (creator.CollectionReport, error) { + if store == nil || hubStore == nil || competitorID == "" { return creator.CollectionReport{}, creator.ErrUnavailable } competitor, err := store.GetCompetitor(ctx, competitorID) @@ -1778,7 +1795,7 @@ func syncCreatorCompetitorWithClaim(ctx context.Context, store *creator.Store, p if !claimed { return creator.CollectionReport{}, creator.ErrConflict } - // blocked(环境暂不可用)也必须排下一次重试:否则 next_sync_at 置空后定时获取对该账号永久失效。 + // blocked(网关不可用等环境暂不可用)也必须排下一次重试:否则 next_sync_at 置空后定时获取对该账号永久失效。 nextRetry := func() *time.Time { nextBase := now if competitor.NextSyncAt != nil { @@ -1797,46 +1814,23 @@ func syncCreatorCompetitorWithClaim(ctx context.Context, store *creator.Store, p if competitor.Platform != creator.PlatformDouyin { return blocked(fmt.Errorf("%w: unsupported creator platform %s", creator.ErrUnavailable, competitor.Platform)) } - account, err := phaseAStore.GetAccount(ctx, accountID) + // 竞品作品是公开数据:用临时匿名浏览器采集,不依赖任何已登录自有账号; + // 并发名额由 newAnonymousBrowser 全局信号量控制(含调度器并发与手动同步)。 + lease, err := newAnonymousBrowser(ctx, hubStore) if err != nil { - return blocked(err) - } - if account.Platform != competitor.Platform || account.AuthorizationStatus != "authorized" { - return blocked(creator.ErrConflict) - } - profile, err := store.GetAccountProfile(ctx, accountID) - if err != nil { - return blocked(err) - } - if (profile.BusinessStatus != "normal" && profile.BusinessStatus != "muted") || profile.LoginStatus != "logged_in" { - return blocked(creator.ErrConflict) - } - environment, err := hubStore.GetEnvironmentContextForAccount(ctx, accountID) - if err != nil { - return blocked(fmt.Errorf("%w: account environment unavailable: %v", creator.ErrUnavailable, err)) - } - if environment.RuntimeID == "" || environment.RuntimeNetworkID == "" || environment.BindingVersion <= 0 { - return blocked(fmt.Errorf("%w: account runtime is not running", creator.ErrUnavailable)) - } - gateway, err := hubStore.GetGateway(ctx, environment.Gateway) - if err != nil { - return blocked(fmt.Errorf("%w: gateway unavailable: %v", creator.ErrUnavailable, err)) - } - useCtx, runtimeUse, err := beginRuntimeUseForEnvironment(ctx, hubStore, environment, "task", "creator-sync-"+competitor.ID) - if err != nil { - return blocked(fmt.Errorf("%w: runtime use unavailable: %v", creator.ErrUnavailable, err)) - } - collector, _, err := newCreatorCollector(useCtx, competitor.Platform, gateway, environment, account.PlatformAccountKey, competitor.PlatformAccountKey, competitor.HomepageURL, creator.SourceCompetitor, competitor.ID) - if err != nil { - return blocked(errors.Join(fmt.Errorf("%w: account identity verification failed: %v", creator.ErrConflict, err), runtimeUse.Close())) + return blocked(fmt.Errorf("%w: anonymous browser unavailable: %v", creator.ErrUnavailable, err)) } + collector := douyinCollector( + creatorGatewayBrowser{gateway: lease.gateway, environment: lease.environment}, + competitor.PlatformAccountKey, creator.SourceCompetitor, competitor.ID, + ) collectionNow := now if competitor.NextSyncAt != nil && !competitor.NextSyncAt.After(now) { collectionNow = competitor.NextSyncAt.UTC() } - report, collectErr := store.CollectSource(useCtx, competitor.Platform, creator.SourceCompetitor, competitor.ID, collector, collectionNow) - if releaseErr := runtimeUse.Close(); releaseErr != nil { - collectErr = errors.Join(collectErr, releaseErr) + report, collectErr := store.CollectSource(ctx, competitor.Platform, creator.SourceCompetitor, competitor.ID, collector, collectionNow) + if closeErr := lease.close(); closeErr != nil { + collectErr = errors.Join(collectErr, closeErr) } if collectErr != nil { nextBase := now @@ -1931,36 +1925,18 @@ func runCreatorScheduleOnce(ctx context.Context, store *creator.Store, phaseASto if err != nil { return err } + var wg sync.WaitGroup for _, competitor := range competitors { - accountID, err := creatorCollectionAccount(ctx, store, phaseAStore, hubStore, competitor.Platform) - if err != nil { - leaseToken, claimed, claimErr := store.ClaimCompetitorSync(ctx, competitor.ID, false, now) - if claimErr != nil { - logrus.WithError(claimErr).WithField("competitor_id", competitor.ID).Warn("creator competitor sync claim failed") - continue + // 多个到期账号并发同步;匿名浏览器名额由全局信号量限制(上限 2),不会同时拉起过多实例。 + wg.Add(1) + go func(competitor creator.Competitor) { + defer wg.Done() + if _, err := syncCreatorCompetitorDue(ctx, store, hubStore, competitor.ID); err != nil { + logrus.WithError(err).WithField("competitor_id", competitor.ID).Warn("creator competitor scheduled sync failed") } - if claimed { - // blocked 也排下一次重试(按新作品间隔),避免无采集账号时永久卡死。 - nextBase := now - if competitor.NextSyncAt != nil { - nextBase = competitor.NextSyncAt.UTC() - } - next := creator.NextFixedRun(nextBase, now, time.Duration(settings.NewWorkIntervalSeconds)*time.Second) - var nextAt *time.Time - if !next.IsZero() { - nextAt = &next - } - if markErr := store.MarkCompetitorSync(ctx, competitor.ID, leaseToken, "blocked", "", err.Error(), nextAt); markErr != nil { - logrus.WithError(markErr).WithField("competitor_id", competitor.ID).Warn("creator competitor sync block update failed") - } - } - logrus.WithError(err).WithField("competitor_id", competitor.ID).Warn("creator competitor sync blocked") - continue - } - if _, err := syncCreatorCompetitorDue(ctx, store, phaseAStore, hubStore, competitor.ID, accountID); err != nil { - logrus.WithError(err).WithField("competitor_id", competitor.ID).Warn("creator competitor scheduled sync failed") - } + }(competitor) } + wg.Wait() ownedAccounts, err := store.ListDueOwnedAccounts(ctx, now, settings.NewWorkIntervalSeconds) if err != nil { return err diff --git a/internal/controlplane/api/creator_helper_test.go b/internal/controlplane/api/creator_helper_test.go index a10d178..572faaf 100644 --- a/internal/controlplane/api/creator_helper_test.go +++ b/internal/controlplane/api/creator_helper_test.go @@ -144,9 +144,9 @@ func TestCreatorSchedulerAndPreviewGuards(t *testing.T) { _, err := previewDouyinCompetitor(ctx, nil, nil, nil, "account", creator.CompetitorInput{}) return err }}, - {"sync competitor", func() error { _, err := syncCreatorCompetitor(ctx, nil, nil, nil, "competitor", "account"); return err }}, + {"sync competitor", func() error { _, err := syncCreatorCompetitor(ctx, nil, nil, "competitor"); return err }}, {"sync due competitor", func() error { - _, err := syncCreatorCompetitorDue(ctx, nil, nil, nil, "competitor", "account") + _, err := syncCreatorCompetitorDue(ctx, nil, nil, "competitor") return err }}, {"sync owned", func() error { return syncCreatorOwned(ctx, nil, nil, nil, "account", time.Now()) }}, diff --git a/internal/controlplane/api/creator_route_validation_test.go b/internal/controlplane/api/creator_route_validation_test.go index 1048983..a24a7e8 100644 --- a/internal/controlplane/api/creator_route_validation_test.go +++ b/internal/controlplane/api/creator_route_validation_test.go @@ -54,7 +54,7 @@ func TestCreatorWriteRoutesRejectMalformedInputBeforeStoreAccess(t *testing.T) { {method: http.MethodPut, path: "/api/creator/strategies/strategy-1"}, {method: http.MethodPost, path: "/api/creator/competitor-share-jobs"}, {method: http.MethodPut, path: "/api/creator/competitors/competitor-1"}, - {method: http.MethodPost, path: "/api/creator/competitors/competitor-1/sync"}, + // sync 路由不解析请求体(强制同步走匿名浏览器,无输入字段),不适用畸形 body 校验。 {method: http.MethodPost, path: "/api/creator/test/works"}, {method: http.MethodPost, path: "/api/creator/works/work-1/metrics"}, {method: http.MethodPost, path: "/api/creator/works/work-1/material/rewrite/confirm"}, diff --git a/internal/controlplane/api/creator_share_test.go b/internal/controlplane/api/creator_share_test.go index ba803b0..4d2fb08 100644 --- a/internal/controlplane/api/creator_share_test.go +++ b/internal/controlplane/api/creator_share_test.go @@ -92,6 +92,40 @@ func TestAnonymousBrowserLeasePurgesRuntimeAndProfile(t *testing.T) { } } +func TestAnonymousBrowserLeaseReleasesConcurrencySlot(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodDelete || r.URL.Path != "/v1/browsers/anon-test" { + t.Fatalf("unexpected cleanup request: %s %s", r.Method, r.URL.Path) + } + w.WriteHeader(http.StatusNoContent) + })) + defer server.Close() + + slots := make(chan struct{}, 1) + slots <- struct{}{} + lease := anonymousBrowserLease{gateway: hub.Gateway{Endpoint: server.URL, Token: "test-token"}, slot: slots, environment: hub.EnvironmentContext{Env: hub.Env{Alias: "anon-test"}, BindingVersion: 1}} + if err := lease.close(); err != nil { + t.Fatal(err) + } + // 释放 = 名额回到池里(能再次非阻塞占用)。 + select { + case slots <- struct{}{}: + default: + t.Fatal("close must release the anonymous browser concurrency slot") + } + + // 创建失败路径也必须释放名额,否则失败几次后永久占满。 + failed := anonymousBrowserLease{gateway: hub.Gateway{Endpoint: "http://127.0.0.1:1", Token: "test-token"}, slot: slots} + if err := failed.close(); err == nil { + t.Fatal("cleanup against unreachable gateway must fail") + } + select { + case slots <- struct{}{}: + default: + t.Fatal("failed close must still release the anonymous browser concurrency slot") + } +} + func TestDouyinWorkKeyFromURL(t *testing.T) { for _, test := range []struct { url string diff --git a/internal/platform/douyin/creator_collector.go b/internal/platform/douyin/creator_collector.go index 54b6d6e..ea20842 100644 --- a/internal/platform/douyin/creator_collector.go +++ b/internal/platform/douyin/creator_collector.go @@ -184,7 +184,7 @@ func (c CreatorCollector) ListWorks(ctx context.Context, accountKey, cursor stri } works, hasMore, nextCursor, ok := parseCreatorWorksPage(response.Body) if !ok { - return creator.WorkPage{}, fmt.Errorf("%w: invalid douyin works response", ErrInvalid) + return creator.WorkPage{}, fmt.Errorf("%w: invalid douyin works response (status=%d body_len=%d head=%.160s)", ErrInvalid, response.Status, len(response.Body), response.Body) } items := make([]creator.WorkInput, 0, len(works)) for _, work := range works { @@ -248,15 +248,27 @@ func parseCreatorCommentsPage(body []byte) (creator.CommentPage, error) { } items := make([]creator.CommentInput, 0, len(envelope.Comments)) for _, item := range envelope.Comments { - if item.ID == "" || strings.TrimSpace(item.Text) == "" || item.ReplyID != "" { + hasImage := len(item.Images) > 0 + if item.ID == "" || !hasImage && strings.TrimSpace(item.Text) == "" || item.ReplyID != "" && item.ReplyID != "0" { return creator.CommentPage{}, fmt.Errorf("%w: non-top-level or incomplete douyin comment", ErrInvalid) } + text := item.Text + if strings.TrimSpace(text) == "" && hasImage { + text = "[图片]" + } + authorUID, authorName := item.UserUID, item.UserName + if authorUID == "" { + authorUID = item.User.UID + } + if authorName == "" { + authorName = item.User.Nickname + } var published *time.Time if item.CreateTime > 0 { value := time.Unix(item.CreateTime, 0).UTC() published = &value } - items = append(items, creator.CommentInput{Platform: creator.PlatformDouyin, CommentKey: item.ID, AuthorUID: item.UserUID, AuthorName: item.UserName, Content: item.Text, PublishedAt: published, CommentType: "top_level"}) + items = append(items, creator.CommentInput{Platform: creator.PlatformDouyin, CommentKey: item.ID, AuthorUID: authorUID, AuthorName: authorName, Content: text, PublishedAt: published, CommentType: "top_level"}) } page := creator.CommentPage{Items: items, HasMore: bool(*envelope.HasMore)} if envelope.Cursor != nil { @@ -279,11 +291,17 @@ type creatorWorkPageItem struct { } func parseCreatorWorksPage(body []byte) ([]creatorWorkPageItem, bool, *int64, bool) { - if len(body) > 4<<20 { + // 作品项携带播放信息等大字段,单页真实响应可达 2MB+(见网关 RESPONSE_LIMIT 同步放宽)。 + if len(body) > 8<<20 { return nil, false, nil, false } var envelope worksEnvelope - if json.Unmarshal(body, &envelope) != nil || envelope.StatusCode == nil || *envelope.StatusCode != 0 || envelope.HasMore == nil || envelope.Works == nil || len(envelope.Works) > 20 || envelope.MaxCursor != nil && *envelope.MaxCursor < 0 || bool(*envelope.HasMore) && envelope.MaxCursor == nil { + // 平台风控语义(2026 年实测,与 TikTokDownloader #706 等公开项目一致): + // - 匿名会话第一页正常(16-20 条,has_more=1 且带 max_cursor); + // - 第二页起一律返回空 envelope {"status_code":0}(无 aweme_list/has_more/max_cursor), + // 登录 Cookie 也无法翻页。无 aweme_list 的空 envelope 视为正常空页(分页终止), + // 而非无效响应;只有响应同时带 max_cursor 时才继续翻页。 + if json.Unmarshal(body, &envelope) != nil || envelope.StatusCode == nil || *envelope.StatusCode != 0 || envelope.HasMore != nil && envelope.Works == nil || len(envelope.Works) > 20 || envelope.MaxCursor != nil && *envelope.MaxCursor < 0 { return nil, false, nil, false } items := make([]creatorWorkPageItem, 0, len(envelope.Works)) @@ -312,7 +330,9 @@ func parseCreatorWorksPage(body []byte) ([]creatorWorkPageItem, bool, *int64, bo } items = append(items, creatorWorkPageItem{ID: item.ID, Description: item.Description, CreatedAt: createdAt, CreatedAtInvalid: createdAtInvalid, DiggCount: likes, CommentCount: comments, ShareCount: shares}) } - return items, bool(*envelope.HasMore), envelope.MaxCursor, true + // 仅当响应带有效 max_cursor 时才声明翻页;无 cursor 的 has_more 不产生翻页游标。 + hasMore := envelope.HasMore != nil && bool(*envelope.HasMore) && envelope.MaxCursor != nil + return items, hasMore, envelope.MaxCursor, true } type commentEnvelope struct { @@ -323,9 +343,17 @@ type commentEnvelope struct { ID string `json:"cid"` Text string `json:"text"` CreateTime int64 `json:"create_time"` - ReplyID string `json:"reply_id"` - UserUID string `json:"user_uid"` - UserName string `json:"user_name"` + // 抖音用 reply_id="0" 表示顶层评论;只有非零 reply_id 才是回复评论。 + ReplyID string `json:"reply_id"` + // 图片评论 text 为空,仅以 image_list 承载内容。 + Images []json.RawMessage `json:"image_list"` + // 登录态响应曾提供顶层 user_uid/user_name;匿名响应作者信息在 user 对象里。 + UserUID string `json:"user_uid"` + UserName string `json:"user_name"` + User struct { + UID string `json:"uid"` + Nickname string `json:"nickname"` + } `json:"user"` } `json:"comments"` } diff --git a/internal/platform/douyin/creator_collector_test.go b/internal/platform/douyin/creator_collector_test.go index 6b77033..5b86afd 100644 --- a/internal/platform/douyin/creator_collector_test.go +++ b/internal/platform/douyin/creator_collector_test.go @@ -84,6 +84,41 @@ func TestParseCreatorCommentsPageAllowsEmptyComments(t *testing.T) { } } +func TestParseCreatorCommentsPageAcceptsZeroReplyIDAndUserObject(t *testing.T) { + // 匿名响应实测:reply_id="0" 表示顶层评论;作者信息在 user 对象里。 + body := []byte(`{"status_code":0,"has_more":true,"cursor":20,"comments":[{"cid":"c1","text":"hi","create_time":1700000000,"reply_id":"0","user":{"uid":"u9","nickname":"昵称"}}]}`) + page, err := parseCreatorCommentsPage(body) + if err != nil { + t.Fatal(err) + } + if !page.HasMore || page.NextCursor != "20" || len(page.Items) != 1 { + t.Fatalf("unexpected comments page: %+v", page) + } + item := page.Items[0] + if item.AuthorUID != "u9" || item.AuthorName != "昵称" || item.CommentType != "top_level" { + t.Fatalf("author fields missing: %+v", item) + } +} + +func TestParseCreatorCommentsPageRejectsReplyComment(t *testing.T) { + body := []byte(`{"status_code":0,"has_more":false,"comments":[{"cid":"c1","text":"reply","reply_id":"123"}]}`) + if _, err := parseCreatorCommentsPage(body); err == nil { + t.Fatal("expected reply comment to be rejected") + } +} + +func TestParseCreatorCommentsPageAcceptsImageComment(t *testing.T) { + // 匿名响应实测:纯图片评论 text 为空、content_type=2,仅以 image_list 承载内容。 + body := []byte(`{"status_code":0,"has_more":false,"comments":[{"cid":"c2","text":"","reply_id":"0","image_list":[{"url_list":["https://example.com/x"]}],"user":{"uid":"u1","nickname":"n"}}]}`) + page, err := parseCreatorCommentsPage(body) + if err != nil { + t.Fatal(err) + } + if len(page.Items) != 1 || page.Items[0].Content != "[图片]" { + t.Fatalf("image comment not accepted: %+v", page.Items) + } +} + func TestParseCreatorWorksPageMarksInvalidTimestamp(t *testing.T) { body := []byte(`{"status_code":0,"has_more":false,"aweme_list":[{"aweme_id":"123","desc":"invalid","create_time":-1}]}`) works, _, _, ok := parseCreatorWorksPage(body) @@ -100,6 +135,24 @@ func TestParseCreatorWorksPageAcceptsNumericHasMore(t *testing.T) { } } +func TestParseCreatorWorksPageToleratesHasMoreWithoutCursor(t *testing.T) { + // 匿名 works 响应实测:has_more=1 且不带 max_cursor 字段;第一页必须可用。 + body := []byte(`{"status_code":0,"has_more":1,"aweme_list":[{"aweme_id":"123","desc":"first page","create_time":1700000000,"statistics":{"digg_count":1,"comment_count":2,"share_count":3}}]}`) + works, hasMore, cursor, ok := parseCreatorWorksPage(body) + if !ok || hasMore || cursor != nil || len(works) != 1 || works[0].ID != "123" { + t.Fatalf("has_more without cursor must still yield first page: ok=%v hasMore=%v cursor=%v works=%+v", ok, hasMore, cursor, works) + } +} + +func TestParseCreatorWorksPageTreatsBareEnvelopeAsEmptyPage(t *testing.T) { + // 平台风控语义(TikTokDownloader #706 同样现象):翻页被拦截时返回空 + // envelope {"status_code":0},无 aweme_list/has_more/max_cursor。视为正常空页。 + works, hasMore, cursor, ok := parseCreatorWorksPage([]byte(`{"status_code":0}`)) + if !ok || hasMore || cursor != nil || len(works) != 0 { + t.Fatalf("bare envelope must be a valid empty page: ok=%v hasMore=%v cursor=%v works=%+v", ok, hasMore, cursor, works) + } +} + func TestParseCreatorWorksPageKeepsPartialMetadata(t *testing.T) { body := []byte(`{"status_code":0,"has_more":false,"aweme_list":[{"aweme_id":"123","desc":"partial","statistics":{"digg_count":7}}]}`) works, hasMore, cursor, ok := parseCreatorWorksPage(body)