feat(creator): 竞品作品同步改用匿名浏览器并限制并发
- 竞品同步不再依赖自有采集账号:每次同步按名额(并发上限 2)创建匿名浏览器, 导航至抖音首页等待访客 cookie 稳定后经网关受限 fetch 拉取作品 - 调度器对到期竞品并发同步,名额耗尽时排队等待 - 平台风控语义(实测,与公开爬虫项目一致):作品列表匿名仅能取第一页 (16-20 条),后续页一律空 envelope;空 envelope 视为正常分页终止 - 评论解析适配真实响应:reply_id="0" 为顶层评论、作者取 user 对象、 纯图片评论(text 空但带 image_list)记为「[图片]」 - fetch 签名与白名单:注入 byted_acrawler.frontierSign 的 a_bogus, 网关放宽响应上限至 8MB 并允许签名参数
This commit is contained in:
@@ -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:
|
||||
|
||||
@@ -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))
|
||||
)
|
||||
|
||||
|
||||
|
||||
@@ -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(
|
||||
[
|
||||
|
||||
@@ -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")
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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()) }},
|
||||
|
||||
@@ -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"},
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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"`
|
||||
}
|
||||
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user