feat: restore Xiaohongshu read-only collection

This commit is contained in:
2026-09-14 20:05:32 +08:00
parent 44a28954cf
commit 4c37d0c9bf
13 changed files with 1539 additions and 39 deletions
+10 -32
View File
@@ -980,7 +980,7 @@ func verifyCreatorAccount(ctx context.Context, store *creator.Store, phaseAStore
if err != nil {
return creator.LoginResult{}, err
}
if account.Platform != creator.PlatformDouyin || profile.Platform != creator.PlatformDouyin || account.AuthorizationStatus != "authorized" || profile.PlatformAccountKey == "" || account.PlatformAccountKey != profile.PlatformAccountKey {
if account.Platform != profile.Platform || (account.Platform != creator.PlatformDouyin && account.Platform != creator.PlatformXiaohongshu) || account.AuthorizationStatus != "authorized" || profile.PlatformAccountKey == "" || account.PlatformAccountKey != profile.PlatformAccountKey {
return creator.LoginResult{}, creator.ErrConflict
}
environment, err := hubStore.GetEnvironmentContextForAccount(ctx, accountID)
@@ -994,8 +994,7 @@ func verifyCreatorAccount(ctx context.Context, store *creator.Store, phaseAStore
if err != nil {
return creator.LoginResult{}, fmt.Errorf("%w: gateway unavailable: %v", creator.ErrUnavailable, err)
}
browser := creatorGatewayBrowser{gateway: gateway, environment: environment}
uid, err := browser.Identity(ctx, profile.PlatformAccountKey)
uid, err := verifyCreatorPlatformIdentity(ctx, account.Platform, gateway, environment, profile.PlatformAccountKey)
if err != nil {
return creator.LoginResult{}, fmt.Errorf("%w: verify the manually logged-in browser identity: %v", creator.ErrConflict, err)
}
@@ -1126,8 +1125,8 @@ func syncCreatorCompetitorWithClaim(ctx context.Context, store *creator.Store, p
markErr := store.MarkCompetitorSync(ctx, competitorID, leaseToken, "blocked", "", blockErr.Error(), nil)
return creator.CollectionReport{}, errors.Join(blockErr, markErr)
}
if competitor.Platform != creator.PlatformDouyin {
return blocked(fmt.Errorf("%w: 小红书采集器尚未完成平台能力验证", creator.ErrUnavailable))
if competitor.Platform != creator.PlatformDouyin && competitor.Platform != creator.PlatformXiaohongshu {
return blocked(fmt.Errorf("%w: unsupported creator platform %s", creator.ErrUnavailable, competitor.Platform))
}
account, err := phaseAStore.GetAccount(ctx, accountID)
if err != nil {
@@ -1154,16 +1153,10 @@ func syncCreatorCompetitorWithClaim(ctx context.Context, store *creator.Store, p
if err != nil {
return blocked(fmt.Errorf("%w: gateway unavailable: %v", creator.ErrUnavailable, err))
}
browser := creatorGatewayBrowser{gateway: gateway, environment: environment}
if _, identityErr := browser.Identity(ctx, account.PlatformAccountKey); identityErr != nil {
return blocked(fmt.Errorf("%w: verify the manually logged-in browser identity: %v", creator.ErrConflict, identityErr))
}
collector := douyin.CreatorCollector{Browser: browser, AccountKey: competitor.PlatformAccountKey, SourceType: creator.SourceCompetitor, SourceID: competitor.ID}
canonicalSecUID, err := collector.CanonicalSecUID(ctx, account.PlatformAccountKey)
collector, _, err := newCreatorCollector(ctx, competitor.Platform, gateway, environment, account.PlatformAccountKey, creator.SourceCompetitor, competitor.ID)
if err != nil {
return blocked(fmt.Errorf("%w: account identity verification failed: %v", creator.ErrConflict, err))
}
collector.AccountKey = canonicalSecUID
collectionNow := now
if competitor.NextSyncAt != nil && !competitor.NextSyncAt.After(now) {
collectionNow = competitor.NextSyncAt.UTC()
@@ -1279,7 +1272,7 @@ func refreshCreatorMetricWork(ctx context.Context, store *creator.Store, phaseAS
if err != nil {
return err
}
if account.Platform != creator.PlatformDouyin || account.AuthorizationStatus != "authorized" || profile.BusinessStatus != "normal" || profile.LoginStatus != "logged_in" {
if (account.Platform != creator.PlatformDouyin && account.Platform != creator.PlatformXiaohongshu) || account.AuthorizationStatus != "authorized" || profile.BusinessStatus != "normal" || profile.LoginStatus != "logged_in" {
return creator.ErrConflict
}
environment, err := hubStore.GetEnvironmentContextForAccount(ctx, accountID)
@@ -1290,19 +1283,13 @@ func refreshCreatorMetricWork(ctx context.Context, store *creator.Store, phaseAS
if err != nil {
return fmt.Errorf("%w: gateway unavailable: %v", creator.ErrUnavailable, err)
}
browser := creatorGatewayBrowser{gateway: gateway, environment: environment}
if _, identityErr := browser.Identity(ctx, account.PlatformAccountKey); identityErr != nil {
return fmt.Errorf("%w: verify the manually logged-in browser identity: %v", creator.ErrConflict, identityErr)
}
collector := douyin.CreatorCollector{Browser: browser, AccountKey: account.PlatformAccountKey, SourceType: work.SourceType, SourceID: work.SourceID}
canonical, err := collector.CanonicalSecUID(ctx, account.PlatformAccountKey)
collector, collectionKey, err := newCreatorCollector(ctx, work.Platform, gateway, environment, account.PlatformAccountKey, work.SourceType, work.SourceID)
if err != nil {
return fmt.Errorf("%w: account identity verification failed: %v", creator.ErrConflict, err)
}
collector.AccountKey = canonical
cursor := ""
for page := 0; page < 100; page++ {
result, pageErr := collector.ListWorks(ctx, canonical, cursor)
result, pageErr := collector.ListWorks(ctx, collectionKey, cursor)
if pageErr != nil {
return pageErr
}
@@ -1335,7 +1322,7 @@ func syncCreatorOwned(ctx context.Context, store *creator.Store, phaseAStore *ph
if err != nil {
return err
}
if account.Platform != creator.PlatformDouyin || account.AuthorizationStatus != "authorized" {
if account.Platform != creator.PlatformDouyin && account.Platform != creator.PlatformXiaohongshu || account.AuthorizationStatus != "authorized" {
return creator.ErrConflict
}
profile, err := store.GetAccountProfile(ctx, accountID)
@@ -1349,9 +1336,6 @@ func syncCreatorOwned(ctx context.Context, store *creator.Store, phaseAStore *ph
if err != nil {
return err
}
blockOwned := func(blockErr error) error {
return errors.Join(blockErr, store.MarkCollectionBlocked(ctx, creator.SourceOwned, accountID, blockErr.Error(), now, settings.LookbackDays))
}
environment, err := hubStore.GetEnvironmentContextForAccount(ctx, accountID)
if err != nil {
return fmt.Errorf("%w: account environment unavailable: %v", creator.ErrUnavailable, err)
@@ -1363,10 +1347,6 @@ func syncCreatorOwned(ctx context.Context, store *creator.Store, phaseAStore *ph
if err != nil {
return fmt.Errorf("%w: gateway unavailable: %v", creator.ErrUnavailable, err)
}
browser := creatorGatewayBrowser{gateway: gateway, environment: environment}
if _, identityErr := browser.Identity(ctx, account.PlatformAccountKey); identityErr != nil {
return blockOwned(fmt.Errorf("%w: verify the manually logged-in browser identity: %v", creator.ErrConflict, identityErr))
}
syncLease, err := store.ClaimSourceSync(ctx, creator.SourceOwned, account.ID)
if err != nil {
return err
@@ -1376,13 +1356,11 @@ func syncCreatorOwned(ctx context.Context, store *creator.Store, phaseAStore *ph
logrus.WithError(releaseErr).WithField("account_id", account.ID).Warn("creator source sync lease release failed")
}
}()
collector := douyin.CreatorCollector{Browser: browser, AccountKey: account.PlatformAccountKey, SourceType: creator.SourceOwned, SourceID: account.ID}
canonicalSecUID, err := collector.CanonicalSecUID(ctx, account.PlatformAccountKey)
collector, _, err := newCreatorCollector(ctx, account.Platform, gateway, environment, account.PlatformAccountKey, creator.SourceOwned, account.ID)
if err != nil {
blockErr := store.MarkCollectionBlocked(ctx, creator.SourceOwned, account.ID, err.Error(), now, settings.LookbackDays)
return errors.Join(fmt.Errorf("%w: account identity verification failed: %v", creator.ErrConflict, err), blockErr)
}
collector.AccountKey = canonicalSecUID
_, collectionNow, windowErr := store.NextCollectionWindow(ctx, creator.SourceOwned, account.ID, now, time.Duration(settings.NewWorkIntervalSeconds)*time.Second, settings.LookbackDays)
if windowErr != nil {
return windowErr
+15 -5
View File
@@ -20,7 +20,7 @@ type creatorMaterialDownloader struct {
}
func (downloader creatorMaterialDownloader) Download(ctx context.Context, work creator.Work, destination string) error {
if work.Platform != creator.PlatformDouyin || downloader.store == nil || downloader.phaseAStore == nil || downloader.hubStore == nil {
if (work.Platform != creator.PlatformDouyin && work.Platform != creator.PlatformXiaohongshu) || downloader.store == nil || downloader.phaseAStore == nil || downloader.hubStore == nil {
return fmt.Errorf("%w: creator media gateway is unavailable", creator.ErrUnavailable)
}
accountID := work.SourceID
@@ -39,7 +39,7 @@ func (downloader creatorMaterialDownloader) Download(ctx context.Context, work c
if err != nil {
return err
}
if account.Platform != creator.PlatformDouyin || profile.Platform != creator.PlatformDouyin || account.AuthorizationStatus != "authorized" || profile.LoginStatus != "logged_in" || account.PlatformAccountKey != profile.PlatformAccountKey {
if account.Platform != work.Platform || profile.Platform != work.Platform || account.AuthorizationStatus != "authorized" || profile.LoginStatus != "logged_in" || account.PlatformAccountKey != profile.PlatformAccountKey {
return fmt.Errorf("%w: media account identity is not verified", creator.ErrConflict)
}
environment, err := downloader.hubStore.GetEnvironmentContextForAccount(ctx, accountID)
@@ -50,11 +50,21 @@ func (downloader creatorMaterialDownloader) Download(ctx context.Context, work c
if err != nil {
return fmt.Errorf("%w: media gateway unavailable: %v", creator.ErrUnavailable, err)
}
browser := creatorGatewayBrowser{gateway: gateway, environment: environment}
if _, err := browser.Identity(ctx, profile.PlatformAccountKey); err != nil {
if _, err := verifyCreatorPlatformIdentity(ctx, work.Platform, gateway, environment, profile.PlatformAccountKey); err != nil {
return fmt.Errorf("%w: media browser identity verification failed: %v", creator.ErrConflict, err)
}
return browser.Media(ctx, work.OriginalURL, destination)
if work.Platform == creator.PlatformDouyin {
return (creatorGatewayBrowser{gateway: gateway, environment: environment}).Media(ctx, work.OriginalURL, destination)
}
data, contentType, err := (xiaohongshuGatewayBrowser{gateway: gateway, environment: environment}).Media(ctx, work.OriginalURL)
if err != nil {
return err
}
contentType = strings.ToLower(strings.TrimSpace(strings.SplitN(contentType, ";", 2)[0]))
if !strings.HasPrefix(contentType, "video/") && contentType != "application/octet-stream" {
return fmt.Errorf("xiaohongshu media response is not a video")
}
return writeCreatorMedia(destination, data)
}
func processCreatorMaterial(ctx context.Context, store *creator.Store, phaseAStore *phasea.Store, hubStore *hub.Store, workID string) (creator.MaterialJob, error) {
+192
View File
@@ -0,0 +1,192 @@
package main
import (
"context"
"encoding/base64"
"encoding/json"
"errors"
"fmt"
"net/http"
"net/url"
"strings"
"time"
"git.ipao.vip/rogee/creator-hub/internal/creator"
"git.ipao.vip/rogee/creator-hub/internal/douyin"
"git.ipao.vip/rogee/creator-hub/internal/hub"
"git.ipao.vip/rogee/creator-hub/internal/xiaohongshu"
)
type xiaohongshuGatewayBrowser struct {
gateway hub.Gateway
environment hub.EnvironmentContext
}
func (browser xiaohongshuGatewayBrowser) generation() (map[string]any, error) {
request, err := (douyinGatewayBrowser{gateway: browser.gateway, environment: browser.environment}).request()
if err != nil {
return nil, err
}
return map[string]any{
"binding_version": request.BindingVersion,
"runtime_id": request.RuntimeID,
"network_id": request.NetworkID,
"network_exit_id": request.NetworkExitID,
}, nil
}
func (browser xiaohongshuGatewayBrowser) Get(ctx context.Context, target string) (xiaohongshu.Response, error) {
request, err := browser.generation()
if err != nil {
return xiaohongshu.Response{}, err
}
request["url"] = target
status, body, err := gatewayCall(ctx, browser.gateway, http.MethodPost,
"/v1/browsers/"+url.PathEscape(browser.environment.Alias)+"/xiaohongshu/get", request, 30*time.Second)
if err != nil || status != http.StatusOK {
return xiaohongshu.Response{}, errors.New("restricted Xiaohongshu browser operation failed")
}
return decodeXiaohongshuResponse(body)
}
func (browser xiaohongshuGatewayBrowser) Post(ctx context.Context, target string, payload []byte) (xiaohongshu.Response, error) {
if len(payload) == 0 || len(payload) > 4<<20 {
return xiaohongshu.Response{}, errors.New("invalid Xiaohongshu browser body")
}
var bodyValue any
if err := json.Unmarshal(payload, &bodyValue); err != nil {
return xiaohongshu.Response{}, fmt.Errorf("invalid Xiaohongshu browser body: %w", err)
}
if _, ok := bodyValue.(map[string]any); !ok {
return xiaohongshu.Response{}, errors.New("Xiaohongshu browser body must be an object")
}
request, err := browser.generation()
if err != nil {
return xiaohongshu.Response{}, err
}
request["url"] = target
request["body"] = bodyValue
status, responseBody, err := gatewayCall(ctx, browser.gateway, http.MethodPost,
"/v1/browsers/"+url.PathEscape(browser.environment.Alias)+"/xiaohongshu/post", request, 30*time.Second)
if err != nil || status != http.StatusOK {
return xiaohongshu.Response{}, errors.New("restricted Xiaohongshu browser POST failed")
}
return decodeXiaohongshuResponse(responseBody)
}
func (browser xiaohongshuGatewayBrowser) Identity(ctx context.Context, expectedKey string) (string, error) {
request, err := browser.generation()
if err != nil {
return "", err
}
request["expected_account_key"] = expectedKey
status, body, err := gatewayCall(ctx, browser.gateway, http.MethodPost,
"/v1/browsers/"+url.PathEscape(browser.environment.Alias)+"/xiaohongshu/identity", request, 30*time.Second)
if err != nil || status != http.StatusOK {
return "", errors.New("Xiaohongshu identity verification failed")
}
var identity struct {
UID string `json:"uid"`
}
if err := json.Unmarshal(body, &identity); err != nil || strings.TrimSpace(identity.UID) == "" {
return "", errors.New("Xiaohongshu identity response omitted uid")
}
return identity.UID, nil
}
func (browser xiaohongshuGatewayBrowser) Media(ctx context.Context, target string) ([]byte, string, error) {
request, err := browser.generation()
if err != nil {
return nil, "", err
}
request["url"] = target
status, body, err := gatewayCall(ctx, browser.gateway, http.MethodPost,
"/v1/browsers/"+url.PathEscape(browser.environment.Alias)+"/xiaohongshu/media", request, 90*time.Second)
if err != nil || status != http.StatusOK {
return nil, "", errors.New("restricted Xiaohongshu media request failed")
}
var response struct {
Status int `json:"status"`
ContentType string `json:"content_type"`
BodyBase64 string `json:"body_base64"`
}
if err := json.Unmarshal(body, &response); err != nil || response.Status < 200 || response.Status >= 300 || response.BodyBase64 == "" {
return nil, "", errors.New("Xiaohongshu media response is invalid")
}
data, err := decodeBase64(response.BodyBase64)
if err != nil {
return nil, "", err
}
return data, response.ContentType, nil
}
func decodeXiaohongshuResponse(body []byte) (xiaohongshu.Response, error) {
var response struct {
Status int `json:"status"`
Body string `json:"body"`
Challenge string `json:"challenge"`
}
if err := json.Unmarshal(body, &response); err != nil || response.Status < 100 || response.Status > 599 {
return xiaohongshu.Response{}, errors.New("restricted Xiaohongshu browser returned an invalid response")
}
return xiaohongshu.Response{Status: response.Status, Body: []byte(response.Body), Challenge: response.Challenge}, nil
}
func newCreatorCollector(ctx context.Context, platform string, gateway hub.Gateway, environment hub.EnvironmentContext, accountKey, sourceType, sourceID string) (creator.PlatformCollector, string, error) {
switch platform {
case creator.PlatformDouyin:
browser := creatorGatewayBrowser{gateway: gateway, environment: environment}
uid, err := browser.Identity(ctx, accountKey)
if err != nil {
return nil, "", err
}
collector := douyinCollector(browser, accountKey, sourceType, sourceID)
canonical, err := collector.CanonicalSecUID(ctx, uid)
if err != nil {
return nil, "", err
}
collector.AccountKey = canonical
return &collector, canonical, nil
case creator.PlatformXiaohongshu:
browser := xiaohongshuGatewayBrowser{gateway: gateway, environment: environment}
uid, err := browser.Identity(ctx, accountKey)
if err != nil {
return nil, "", err
}
return &xiaohongshu.Collector{Browser: browser, AccountKey: uid, SourceType: sourceType, SourceID: sourceID}, uid, nil
default:
return nil, "", fmt.Errorf("%w: unsupported creator platform %s", creator.ErrUnavailable, platform)
}
}
func douyinCollector(browser creatorGatewayBrowser, accountKey, sourceType, sourceID string) douyin.CreatorCollector {
return douyin.CreatorCollector{Browser: browser, AccountKey: accountKey, SourceType: sourceType, SourceID: sourceID}
}
func verifyCreatorPlatformIdentity(ctx context.Context, platform string, gateway hub.Gateway, environment hub.EnvironmentContext, expectedKey string) (string, error) {
switch platform {
case creator.PlatformDouyin:
return (creatorGatewayBrowser{gateway: gateway, environment: environment}).Identity(ctx, expectedKey)
case creator.PlatformXiaohongshu:
return (xiaohongshuGatewayBrowser{gateway: gateway, environment: environment}).Identity(ctx, expectedKey)
default:
return "", fmt.Errorf("%w: unsupported creator platform %s", creator.ErrUnavailable, platform)
}
}
func decodeBase64(value string) ([]byte, error) {
const maxEncoded = 96 << 20
if len(value) > maxEncoded {
return nil, errors.New("media response is too large")
}
data, err := base64.StdEncoding.DecodeString(value)
if err != nil {
return nil, fmt.Errorf("decode media response: %w", err)
}
if len(data) > maxCreatorMediaBytes {
return nil, errors.New("media response is too large")
}
return data, nil
}
var _ xiaohongshu.Browser = xiaohongshuGatewayBrowser{}
+53
View File
@@ -0,0 +1,53 @@
package main
import (
"context"
"encoding/json"
"net/http"
"net/http/httptest"
"testing"
"git.ipao.vip/rogee/creator-hub/internal/hub"
)
const testXiaohongshuIdentityURL = "https://edith.xiaohongshu.com/api/sns/web/v2/user/me"
func TestXiaohongshuGatewayBrowserFencesAccountGeneration(t *testing.T) {
server := httptest.NewServer(http.HandlerFunc(func(response http.ResponseWriter, request *http.Request) {
if request.Header.Get("Authorization") != "Bearer gateway-token-1" {
t.Fatal("missing gateway authorization")
}
var body map[string]any
if json.NewDecoder(request.Body).Decode(&body) != nil || body["binding_version"] != float64(2) || body["runtime_id"] != "runtime-a" || body["network_id"] != "network-a" || body["network_exit_id"] != "exit-a" {
t.Fatalf("generation fence missing: %#v", body)
}
if request.URL.Path != "/v1/browsers/account-a/xiaohongshu/get" || body["url"] != testXiaohongshuIdentityURL {
t.Fatalf("unexpected request: path=%s body=%#v", request.URL.Path, body)
}
_ = json.NewEncoder(response).Encode(map[string]any{"status": 200, "body": `{"success":true}`, "challenge": ""})
}))
defer server.Close()
browser := xiaohongshuGatewayBrowser{gateway: hub.Gateway{Endpoint: server.URL, Token: "gateway-token-1"}, environment: readyDouyinEnvironment()}
result, err := browser.Get(context.Background(), testXiaohongshuIdentityURL)
if err != nil || result.Status != 200 || string(result.Body) != `{"success":true}` {
t.Fatalf("unexpected result: %#v err=%v", result, err)
}
}
func TestXiaohongshuGatewayBrowserPostCarriesJSONBody(t *testing.T) {
server := httptest.NewServer(http.HandlerFunc(func(response http.ResponseWriter, request *http.Request) {
var body map[string]any
if json.NewDecoder(request.Body).Decode(&body) != nil || body["url"] != "https://so.xiaohongshu.com/api/sns/web/v2/search/notes" {
t.Fatalf("unexpected request body: %#v", body)
}
if _, ok := body["body"].(map[string]any); !ok {
t.Fatalf("missing nested request body: %#v", body)
}
_ = json.NewEncoder(response).Encode(map[string]any{"status": 200, "body": `{"success":true}`, "challenge": ""})
}))
defer server.Close()
browser := xiaohongshuGatewayBrowser{gateway: hub.Gateway{Endpoint: server.URL, Token: "gateway-token-1"}, environment: readyDouyinEnvironment()}
if _, err := browser.Post(context.Background(), "https://so.xiaohongshu.com/api/sns/web/v2/search/notes", []byte(`{"keyword":"x"}`)); err != nil {
t.Fatalf("post failed: %v", err)
}
}
+118
View File
@@ -1285,6 +1285,124 @@ def detect_challenge(status: int, body: str) -> str:
return ""
XHS_ORIGIN = "https://www.xiaohongshu.com"
XHS_API_ORIGIN = "https://edith.xiaohongshu.com"
XHS_SEARCH_ORIGIN = "https://so.xiaohongshu.com"
XHS_IDENTITY_URL = XHS_API_ORIGIN + "/api/sns/web/v2/user/me"
XHS_ALLOWED_HOSTS = frozenset({"www.xiaohongshu.com", "edith.xiaohongshu.com", "so.xiaohongshu.com"})
class XiaohongshuBrowser(DouyinBrowser):
def __init__(self, endpoint=None) -> None:
super().__init__(
endpoint,
origin=XHS_ORIGIN,
url_validator=is_xiaohongshu_url,
media_validator=is_xiaohongshu_media_url,
)
def post(self, alias: str, target: str, body: bytes) -> BrowserResponse:
if not is_xiaohongshu_url(target) or len(body) > RESPONSE_LIMIT:
raise DouyinError("restricted Xiaohongshu POST request is invalid")
try:
body_text = body.decode("utf-8")
except UnicodeDecodeError as exc:
raise DouyinError("restricted Xiaohongshu POST body is not UTF-8") from exc
with self.connection(alias) as cdp:
if cdp.evaluate("location.origin") != self.origin:
raise DouyinError("restricted browser origin changed")
expression = f"""(async()=>{{
const r=await fetch({json.dumps(target)},{{method:'POST',headers:{{'content-type':'application/json'}},body:{json.dumps(body_text)},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, responseBody='';
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;responseBody+=decoder.decode(item.value,{{stream:true}});
}}
responseBody+=decoder.decode();return {{status:r.status,body:responseBody,too_large:false}};
}})()"""
result = cdp.evaluate(expression)
if (
not isinstance(result, dict)
or result.get("too_large")
or not isinstance(result.get("status"), int)
):
raise DouyinError("restricted Xiaohongshu POST failed")
status = result["status"]
if 300 <= status < 400:
raise DouyinError("restricted Xiaohongshu POST redirected")
response_body = result.get("body")
if not isinstance(response_body, str):
raise DouyinError("restricted Xiaohongshu POST returned invalid body")
return BrowserResponse(status, response_body, detect_challenge(status, response_body))
def identity(self, alias: str, expected_uid: str | None = None) -> dict:
response = self.get(alias, XHS_IDENTITY_URL)
try:
payload = json.loads(response.body)
except json.JSONDecodeError as exc:
raise DouyinError("Xiaohongshu identity response is invalid") from exc
data = payload.get("data") if isinstance(payload, dict) else None
user_info = data.get("user_info") if isinstance(data, dict) else None
user_id = data.get("user_id", "") if isinstance(data, dict) else ""
nickname = data.get("nickname", "") if isinstance(data, dict) else ""
if isinstance(user_info, dict):
user_id = user_id or user_info.get("user_id", "")
nickname = nickname or user_info.get("nickname", "")
success = payload.get("success") if isinstance(payload, dict) else None
if (
response.status != 200
or not isinstance(payload, dict)
or not isinstance(success, bool)
or not success
or not isinstance(user_id, str)
or not ACCOUNT_KEY_RE.fullmatch(user_id)
or nickname is not None
and not isinstance(nickname, str)
):
raise DouyinError("Xiaohongshu login is not valid")
if expected_uid and user_id != expected_uid:
raise DouyinError("Xiaohongshu identity does not match the expected account")
return {"uid": user_id, "user_id": user_id, "nickname": nickname or ""}
def is_xiaohongshu_media_url(value: object) -> bool:
if not isinstance(value, str):
return False
try:
parsed = urlsplit(value)
port = parsed.port
except (TypeError, ValueError):
return False
return (
parsed.scheme == "https"
and parsed.hostname == "www.xiaohongshu.com"
and port is None
and parsed.username is None
and parsed.password is None
and parsed.fragment == ""
and parsed.path.startswith("/explore/")
)
def is_xiaohongshu_url(value: object) -> bool:
if not isinstance(value, str):
return False
try:
parsed = urlsplit(value)
port = parsed.port
except (TypeError, ValueError):
return False
return (
parsed.scheme == "https"
and parsed.hostname in XHS_ALLOWED_HOSTS
and port is None
and parsed.username is None
and parsed.password is None
and parsed.fragment == ""
)
def notice_ids(event: dict) -> list[str]:
try:
payload = json.loads(event["payload"])
+184
View File
@@ -46,6 +46,7 @@ from .douyin import (
DouyinBrowser,
DouyinError,
SubscriptionManager,
XiaohongshuBrowser,
)
from .proxy import ProxyExit, ProxyRegistry
@@ -64,6 +65,12 @@ DOUYIN_IDENTITY_PATH = "/aweme/v1/web/user/profile/self/"
DOUYIN_IDENTITY_URL = IDENTITY_URL
DOUYIN_WORKS_PATH = WORKS_PATH
DOUYIN_COMMENTS_PATH = COMMENTS_PATH
XHS_ACCOUNT_KEY_RE = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._:@/-]{0,127}$")
XHS_IDENTITY_PATH = "/api/sns/web/v2/user/me"
XHS_USER_POSTED_PATH = "/api/sns/web/v1/user_posted"
XHS_COMMENTS_PATH = "/api/sns/web/v2/comment/page"
XHS_SEARCH_PATH = "/api/sns/web/v2/search/notes"
XHS_FEED_PATH = "/api/sns/web/v1/feed"
def _noop() -> None:
@@ -85,12 +92,14 @@ class Gateway:
token: str,
self_name: str,
browser: DouyinBrowser | None = None,
xiaohongshu_browser: XiaohongshuBrowser | None = None,
) -> None:
self.docker = docker
self.network = network
self.token = token
self.self_name = self_name
self.browser = browser or DouyinBrowser(self._browser_endpoint)
self.xiaohongshu_browser = xiaohongshu_browser or XiaohongshuBrowser(self._browser_endpoint)
self.proxies = ProxyRegistry()
self.reservations = AliasReservationManager(docker, self_name)
self.subscriptions = SubscriptionManager(self.browser)
@@ -666,6 +675,76 @@ class Gateway:
)
return identity
def get_xiaohongshu(self, alias: str, input: dict) -> dict:
target = input.get("url", "")
if not valid_xiaohongshu_generation(input) or not valid_xiaohongshu_url(target):
raise RequestError("invalid restricted Xiaohongshu request", 400)
with self._alias_lock(alias):
self._require_douyin_generation(alias, input)
try:
response = self.xiaohongshu_browser.get(alias, target)
self._require_douyin_generation(alias, input)
except DouyinError as exc:
LOG.warning("Xiaohongshu GET failed alias=%s reason=%s", alias, str(exc))
raise RequestError("restricted Xiaohongshu operation failed") from exc
return {"status": response.status, "body": response.body, "challenge": response.challenge}
def post_xiaohongshu(self, alias: str, input: dict) -> dict:
target = input.get("url", "")
body = input.get("body")
if (
not valid_xiaohongshu_generation(input)
or not valid_xhs_post_url(target)
or not isinstance(body, dict)
):
raise RequestError("invalid restricted Xiaohongshu POST request", 400)
try:
encoded = json.dumps(body, ensure_ascii=False, separators=(",", ":")).encode()
except (TypeError, ValueError) as exc:
raise RequestError("invalid restricted Xiaohongshu POST body", 400) from exc
with self._alias_lock(alias):
self._require_douyin_generation(alias, input)
try:
response = self.xiaohongshu_browser.post(alias, target, encoded)
self._require_douyin_generation(alias, input)
except DouyinError as exc:
LOG.warning("Xiaohongshu POST failed alias=%s reason=%s", alias, str(exc))
raise RequestError("restricted Xiaohongshu operation failed") from exc
return {"status": response.status, "body": response.body, "challenge": response.challenge}
def get_xiaohongshu_media(self, alias: str, input: dict) -> dict:
target = input.get("url", "")
if not valid_xiaohongshu_generation(input) or not valid_xiaohongshu_media_url(target):
raise RequestError("invalid restricted Xiaohongshu media request", 400)
with self._alias_lock(alias):
self._require_douyin_generation(alias, input)
try:
response = self.xiaohongshu_browser.get_media(alias, target)
self._require_douyin_generation(alias, input)
except DouyinError as exc:
LOG.warning("Xiaohongshu media download failed alias=%s reason=%s", alias, str(exc))
raise RequestError("restricted Xiaohongshu media download failed") from exc
return {"status": response.status, "content_type": response.content_type, "body_base64": response.body_base64}
def xiaohongshu_identity(self, alias: str, input: dict) -> dict:
expected_account_key = input.get("expected_account_key", "")
if (
not valid_xiaohongshu_generation(input)
or not isinstance(expected_account_key, str)
or not XHS_ACCOUNT_KEY_RE.fullmatch(expected_account_key)
):
raise RequestError("invalid Xiaohongshu identity request", 400)
with self._alias_lock(alias):
self._require_douyin_generation(alias, input)
try:
identity = self.xiaohongshu_browser.identity(alias)
except DouyinError as exc:
LOG.warning("Xiaohongshu identity verification failed alias=%s reason=%s", alias, str(exc))
raise RequestError("Xiaohongshu login identity could not be verified") from exc
if identity.get("uid") != expected_account_key:
raise RequestError("Xiaohongshu identity does not match the expected account", 409)
return identity
def douyin_action(self, alias: str, input: dict) -> dict:
expected_uid = input.get("expected_uid", "")
action = input.get("action", "")
@@ -1130,6 +1209,20 @@ class GatewayHandler(BaseHTTPRequestHandler):
if method == "POST" and action == "proxy":
gateway.restore_proxy(alias, body)
return None
match = re.fullmatch(
r"/v1/browsers/([a-z0-9][a-z0-9-]{0,31})/xiaohongshu/(get|post|media|identity)",
path,
)
if match:
alias, action = match.groups()
if action == "get" and method == "POST":
return gateway.get_xiaohongshu(alias, body)
if action == "post" and method == "POST":
return gateway.post_xiaohongshu(alias, body)
if action == "media" and method == "POST":
return gateway.get_xiaohongshu_media(alias, body)
if action == "identity" and method == "POST":
return gateway.xiaohongshu_identity(alias, body)
match = re.fullmatch(
r"/v1/browsers/([a-z0-9][a-z0-9-]{0,31})/douyin/(get|media|identity|action|events)",
path,
@@ -1419,6 +1512,97 @@ def valid_douyin_generation(value: dict) -> bool:
)
def valid_xiaohongshu_generation(value: dict) -> bool:
return valid_douyin_generation(value)
def valid_xhs_query(query: object, allowed: set[str], required: set[str] | None = None) -> bool:
if not isinstance(query, dict) or not isinstance(allowed, set):
return False
required = required or set()
if not required.issubset(query) or not set(query).issubset(allowed):
return False
for key, values in query.items():
if not isinstance(key, str) or not isinstance(values, list) or len(values) != 1:
return False
if not isinstance(values[0], str) or len(values[0]) > 2048 or "\r" in values[0] or "\n" in values[0]:
return False
return True
def _valid_xhs_host(parsed: object, host: str) -> bool:
return (
getattr(parsed, "scheme", "") == "https"
and getattr(parsed, "hostname", None) == host
and getattr(parsed, "port", None) is None
and getattr(parsed, "username", None) is None
and getattr(parsed, "password", None) is None
and getattr(parsed, "fragment", "") == ""
)
def valid_xhs_url(raw: object) -> bool:
if not isinstance(raw, str):
return False
try:
parsed = urlsplit(raw)
query = parse_qs(parsed.query, keep_blank_values=True)
except ValueError:
return False
if _valid_xhs_host(parsed, "edith.xiaohongshu.com") and parsed.path == XHS_IDENTITY_PATH:
return not query
if _valid_xhs_host(parsed, "edith.xiaohongshu.com") and parsed.path == XHS_USER_POSTED_PATH:
return valid_xhs_query(
query,
{"user_id", "cursor", "num", "image_formats", "xsec_source", "xsec_token"},
{"user_id", "num"},
) and bool(XHS_ACCOUNT_KEY_RE.fullmatch(query["user_id"][0])) and query["num"] == ["30"]
if _valid_xhs_host(parsed, "edith.xiaohongshu.com") and parsed.path == XHS_COMMENTS_PATH:
return valid_xhs_query(
query,
{"note_id", "cursor", "top_comment_id", "image_formats", "xsec_source", "xsec_token"},
{"note_id", "cursor", "top_comment_id"},
) and bool(XHS_ACCOUNT_KEY_RE.fullmatch(query["note_id"][0]))
return False
def valid_xiaohongshu_url(raw: object) -> bool:
return valid_xhs_url(raw)
def valid_xhs_post_url(raw: object) -> bool:
if not isinstance(raw, str):
return False
try:
parsed = urlsplit(raw)
query = parse_qs(parsed.query, keep_blank_values=True)
except ValueError:
return False
return (
_valid_xhs_host(parsed, "so.xiaohongshu.com")
and parsed.path == XHS_SEARCH_PATH
and not query
) or (
_valid_xhs_host(parsed, "edith.xiaohongshu.com")
and parsed.path == XHS_FEED_PATH
and not query
)
def valid_xiaohongshu_media_url(raw: object) -> bool:
if not isinstance(raw, str):
return False
try:
parsed = urlsplit(raw)
query = parse_qs(parsed.query, keep_blank_values=True)
except ValueError:
return False
if not _valid_xhs_host(parsed, "www.xiaohongshu.com"):
return False
parts = parsed.path.strip("/").split("/")
return len(parts) == 2 and parts[0] == "explore" and bool(XHS_ACCOUNT_KEY_RE.fullmatch(parts[1])) and valid_xhs_query(query, {"xsec_source", "xsec_token"})
def valid_douyin_url(raw: object) -> bool:
if not isinstance(raw, str):
return False
+65
View File
@@ -0,0 +1,65 @@
from __future__ import annotations
import unittest
from typing import Any, cast
from unittest.mock import Mock
from . import gateway as gateway_module
Gateway = gateway_module.Gateway
RequestError = gateway_module.RequestError
valid_xhs_post_url = gateway_module.valid_xhs_post_url
valid_xhs_url = gateway_module.valid_xhs_url
valid_xiaohongshu_url = gateway_module.valid_xiaohongshu_url
valid_xiaohongshu_media_url = gateway_module.valid_xiaohongshu_media_url
valid_xiaohongshu_generation = gateway_module.valid_xiaohongshu_generation
class XiaohongshuValidationTests(unittest.TestCase):
def test_read_urls_use_explicit_host_path_and_query_allowlist(self) -> None:
self.assertTrue(valid_xhs_url("https://edith.xiaohongshu.com/api/sns/web/v2/user/me"))
self.assertTrue(valid_xiaohongshu_url("https://edith.xiaohongshu.com/api/sns/web/v2/user/me"))
self.assertTrue(
valid_xhs_url(
"https://edith.xiaohongshu.com/api/sns/web/v1/user_posted?user_id=u-1&cursor=&num=30&xsec_source=pc_user"
)
)
self.assertFalse(valid_xhs_url("https://edith.xiaohongshu.com/api/sns/web/v1/user_posted?user_id=u-1&num=10"))
self.assertFalse(valid_xhs_url("https://edith.xiaohongshu.com.evil/api/sns/web/v2/user/me"))
self.assertTrue(valid_xhs_post_url("https://so.xiaohongshu.com/api/sns/web/v2/search/notes"))
self.assertTrue(valid_xhs_post_url("https://edith.xiaohongshu.com/api/sns/web/v1/feed"))
self.assertTrue(valid_xiaohongshu_media_url("https://www.xiaohongshu.com/explore/n-1?xsec_source=pc_search"))
self.assertFalse(valid_xiaohongshu_media_url("https://www.xiaohongshu.com/explore/n-1#fragment"))
def test_generation_shape_matches_existing_browser_fence(self) -> None:
self.assertTrue(
valid_xiaohongshu_generation(
{"binding_version": 1, "runtime_id": "a" * 64, "network_id": "network", "network_exit_id": ""}
)
)
self.assertFalse(valid_xiaohongshu_generation({"binding_version": 1, "runtime_id": "runtime", "network_id": "network"}))
class XiaohongshuRouteTests(unittest.TestCase):
def test_read_only_routes_dispatch_without_action_or_event_routes(self) -> None:
handler = gateway_module.GatewayHandler.__new__(gateway_module.GatewayHandler)
gateway = Mock()
gateway.get_xiaohongshu.return_value = {"status": 200}
gateway.post_xiaohongshu.return_value = {"status": 200}
gateway.get_xiaohongshu_media.return_value = {"status": 200}
gateway.xiaohongshu_identity.return_value = {"uid": "u-1"}
server = Mock()
server.gateway = gateway
cast(Any, handler).server = server
cast(Any, handler).server_as_gateway = lambda: server
self.assertEqual(handler._route("POST", "/v1/browsers/account-a/xiaohongshu/get", {}, {}), {"status": 200})
self.assertEqual(handler._route("POST", "/v1/browsers/account-a/xiaohongshu/post", {}, {}), {"status": 200})
self.assertEqual(handler._route("POST", "/v1/browsers/account-a/xiaohongshu/media", {}, {}), {"status": 200})
self.assertEqual(handler._route("POST", "/v1/browsers/account-a/xiaohongshu/identity", {}, {}), {"uid": "u-1"})
with self.assertRaises(RequestError):
handler._route("POST", "/v1/browsers/account-a/xiaohongshu/action", {}, {})
gateway.get_xiaohongshu.assert_called_once_with("account-a", {})
gateway.post_xiaohongshu.assert_called_once_with("account-a", {})
gateway.get_xiaohongshu_media.assert_called_once_with("account-a", {})
gateway.xiaohongshu_identity.assert_called_once_with("account-a", {})
+27
View File
@@ -0,0 +1,27 @@
"""Xiaohongshu gateway facade.
The browser implementation lives next to the existing Douyin browser so both
platforms share the CDP transport and response limits without duplicating it.
"""
from .douyin import (
XHS_ALLOWED_HOSTS,
XHS_API_ORIGIN,
XHS_IDENTITY_URL,
XHS_ORIGIN,
XHS_SEARCH_ORIGIN,
XiaohongshuBrowser,
is_xiaohongshu_media_url,
is_xiaohongshu_url,
)
__all__ = [
"XHS_ALLOWED_HOSTS",
"XHS_API_ORIGIN",
"XHS_IDENTITY_URL",
"XHS_ORIGIN",
"XHS_SEARCH_ORIGIN",
"XiaohongshuBrowser",
"is_xiaohongshu_media_url",
"is_xiaohongshu_url",
]