onesvm-browser-server/server/internal/dock/cdp_engine.go
chii 9b689b2476 feat: 落地节点指纹与 Cookie 罐,并按现网能力更新消费/接手文档
公开页抓取改为 Chrome 136 自洽身份 + 每节点 SQLite 养罐,Trafilatura 走 curl_cffi;MCP/README/接手说明与 09-02 现网复测对齐,避免消费方继续抄过期的站点三分表。

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-09-02 17:04:25 +08:00

325 lines
11 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

// cdp_engine.go:CDP 引擎调用序列(lightpanda / headless-shell 共用)。
//
// 复用声明:方法序列、超时、EXTRACT_JS 提取脚本与 detectVendor 正则全部移植自
// bench/site-matrix/cdp_fetch.mjs(Go 重写);bodyText → markdown 走「纯文本
// 段落保持」包装(\n\n 保段),trafilatura 型 markdown 引擎不适用浏览器引擎。
package dock
import (
"context"
"encoding/json"
"fmt"
"net/http"
"regexp"
"strings"
"sync"
"time"
"onesvm.com/onesvm/browser-server/internal/contract"
"onesvm.com/onesvm/browser-server/internal/fingerprint"
)
// cdpUA 默认模版 UA(Chrome 136;有 Session 时用模版覆盖)。
var cdpUA = fingerprint.Allocate("cdp-default").UserAgent
// extractJS 页面提取脚本(逐字移植 cdp_fetch.mjs EXTRACT_JS)。
const extractJS = `(() => {
const title = document.title || "";
const bodyText = (document.body && document.body.innerText) ? document.body.innerText : "";
const html = document.documentElement ? document.documentElement.outerHTML : "";
const lc = (title + "\n" + bodyText).toLowerCase();
return {
title,
text: bodyText.slice(0, 8000),
htmlLen: html.length,
textLen: bodyText.length,
htmlHead: html.slice(0, 4000),
finalUrl: location.href,
readyState: document.readyState,
looksBlocked: /(robot|captcha|sorry|click the button|continue shopping|just a moment|attention required|verify you are human|are you a human|access denied)/i.test(lc),
};
})()`
// detectVendor 反爬特征判定(逐条移植 cdp_fetch.mjs detectVendor 正则族)。
// 返回 cloudflare|waf|paywall|captcha|redirect|empty|none。
func detectVendor(title, text, htmlHead, headersBlob string, looksBlocked bool, status, textLen int) string {
blob := strings.ToLower(title + "\n" + text + "\n" + htmlHead + "\n" + headersBlob)
switch {
case regexp.MustCompile(`cloudflare|cf-challenge|just a moment|cf-turnstile|attention required`).MatchString(blob):
return "cloudflare"
case regexp.MustCompile(`datadome|captcha-delivery`).MatchString(blob):
return "waf"
case regexp.MustCompile(`akamai`).MatchString(blob) && regexp.MustCompile(`access denied`).MatchString(blob):
return "waf"
case regexp.MustCompile(`amazon|opfcaptcha|validatecaptcha|enter the characters you see`).MatchString(blob) &&
regexp.MustCompile(`robot|captcha`).MatchString(blob):
return "waf"
case regexp.MustCompile(`paywall|subscribe to continue|become a member`).MatchString(blob):
return "paywall"
case regexp.MustCompile(`recaptcha|hcaptcha|verify you are human|captcha`).MatchString(blob) &&
regexp.MustCompile(`robot|human|verify`).MatchString(blob):
return "captcha"
case looksBlocked:
return "waf"
case status >= 300 && status < 400:
return "redirect"
case textLen == 0 || textLen < 80:
return "empty"
}
return "none"
}
// cdpProbe 探活 CDP 端点(GET /json/version)。返回 (ok, message)。
func cdpProbe(ctx context.Context, host string) (bool, string) {
hc := &http.Client{Timeout: cdpHTTPTimeout}
req, err := http.NewRequestWithContext(ctx, http.MethodGet, "http://"+host+"/json/version", nil)
if err != nil {
return false, err.Error()
}
resp, err := hc.Do(req)
if err != nil {
return false, err.Error()
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
return false, fmt.Sprintf("/json/version HTTP %d", resp.StatusCode)
}
return true, "ok"
}
// extractResult Runtime.evaluate 返回形状(golden 测试锁死)。
type extractResult struct {
Title string `json:"title"`
Text string `json:"text"`
HTMLLen int `json:"htmlLen"`
TextLen int `json:"textLen"`
HTMLHead string `json:"htmlHead"`
FinalURL string `json:"finalUrl"`
ReadyState string `json:"readyState"`
LooksBlocked bool `json:"looksBlocked"`
}
// navHeaderBits Network.responseReceived 采集的反爬相关头(cf-*/server/x-amzn…)。
type navCollector struct {
mu sync.Mutex
status int
finalURL string
headerBits []string
loadFired bool
navError string
}
// navTimeout 单页导航超时(design §4.3:read 15–30s,取 cdp_fetch 默认 15s 档下限;
// bench cdp_fetch.mjs 默认 20s——取 20s 对齐 bench 实测参数)。
const navTimeout = 20 * time.Second
// cdpFetchPage 完整一次 CDP 抓取:createTarget→attach→enable→UA→navigate→
// 等 load→evaluate 提取→closeTarget。返回 RawResult 或 blocked/upstream 错误。
func cdpFetchPage(ctx context.Context, host, pageURL string, sess *contract.SessionAttach) (*contract.RawResult, *contract.ErrBody) {
c, err := newCdpClient(host, 8*time.Second)
if err != nil {
return nil, &contract.ErrBody{Code: contract.CodeUpstream, Message: "CDP 端点不可达: " + err.Error()}
}
defer c.Close()
var sessionID string
send := func(method string, params any, timeout time.Duration) (json.RawMessage, error) {
return c.Send(ctx, method, params, sessionID, timeout)
}
// Target.createTarget + attachToTarget(flatten)(cdp_fetch.mjs 同序列)。
created, err := c.Send(ctx, "Target.createTarget", map[string]any{"url": "about:blank"}, "", 8*time.Second)
if err != nil {
return nil, &contract.ErrBody{Code: contract.CodeUpstream, Message: "Target.createTarget: " + err.Error()}
}
var tgt struct {
TargetID string `json:"targetId"`
}
_ = json.Unmarshal(created, &tgt)
attached, err := c.Send(ctx, "Target.attachToTarget", map[string]any{"targetId": tgt.TargetID, "flatten": true}, "", 8*time.Second)
if err != nil {
return nil, &contract.ErrBody{Code: contract.CodeUpstream, Message: "Target.attachToTarget: " + err.Error()}
}
var att struct {
SessionID string `json:"sessionId"`
}
_ = json.Unmarshal(attached, &att)
sessionID = att.SessionID
// Page/Runtime/Network enable + UA override(失败不致命,cdp_fetch.mjs 同容错)。
for _, m := range []string{"Page.enable", "Runtime.enable", "Network.enable"} {
_, _ = send(m, map[string]any{}, 8*time.Second)
}
cdpApplySession(send, sess)
// 事件收集器:loadEventFired / responseReceived / navigate errorText。
col := &navCollector{}
done := make(chan struct{})
go func() {
defer close(done)
for ev := range c.Events() {
switch ev.Method {
case "Page.loadEventFired":
col.mu.Lock()
col.loadFired = true
col.mu.Unlock()
case "Network.responseReceived":
var p struct {
Type string `json:"type"`
Response struct {
URL string `json:"url"`
Status int `json:"status"`
MimeType string `json:"mimeType"`
Headers map[string]string `json:"headers"`
} `json:"response"`
}
if json.Unmarshal(ev.Params, &p) != nil {
continue
}
col.mu.Lock()
if p.Type == "Document" || !col.fired() || strings.EqualFold(p.Response.URL, pageURL) {
if col.status == 0 || p.Type == "Document" {
col.status = p.Response.Status
if p.Response.URL != "" {
col.finalURL = p.Response.URL
}
}
n := 0
for k, v := range p.Response.Headers {
if headerHit(k) {
col.headerBits = append(col.headerBits, k+":"+v)
}
n++
if n > 64 {
break
}
}
}
col.mu.Unlock()
}
}
}()
// Page.navigate → 等 loadEventFired(cdp_fetch.mjs 轮询节奏:100ms 步进)。
navDeadline := time.Now().Add(navTimeout)
_, navErr := send("Page.navigate", map[string]any{"url": pageURL}, navTimeout+5*time.Second)
if navErr != nil {
return nil, &contract.ErrBody{Code: contract.CodeUpstream, Message: "Page.navigate: " + navErr.Error()}
}
_ = navErr
for time.Now().Before(navDeadline) {
col.mu.Lock()
lf := col.loadFired
col.mu.Unlock()
if lf {
break
}
select {
case <-ctx.Done():
return nil, &contract.ErrBody{Code: contract.CodeTimeout, Message: "导航等待取消"}
case <-time.After(100 * time.Millisecond):
}
}
// 稳定窗口:load 后 800ms / 超时后 1500ms(cdp_fetch.mjs 同值)。
settle := 1500 * time.Millisecond
if col.isLoaded() {
settle = 800 * time.Millisecond
}
select {
case <-ctx.Done():
case <-time.After(settle):
}
// Runtime.evaluate 提取。
var ex extractResult
evRaw, evErr := send("Runtime.evaluate", map[string]any{
"expression": extractJS,
"returnByValue": true,
"awaitPromise": true,
}, 10*time.Second)
if evErr == nil {
var ev struct {
Result struct {
Value *extractResult `json:"value"`
} `json:"result"`
}
if json.Unmarshal(evRaw, &ev) == nil && ev.Result.Value != nil {
ex = *ev.Result.Value
}
}
gotCookies := cdpCollectCookies(send, firstNonEmpty(ex.FinalURL, col.finalURL, pageURL))
// closeTarget 收尾(失败忽略)。
_, _ = c.Send(ctx, "Target.closeTarget", map[string]any{"targetId": tgt.TargetID}, "", 5*time.Second)
// blocked 判定(detectVendor,cdp_fetch.mjs 正则族移植)。
headersBlob := strings.Join(col.snapshotHeaders(), "\n")
vendor := detectVendor(ex.Title, ex.Text, ex.HTMLHead, headersBlob, ex.LooksBlocked, col.status, ex.TextLen)
if vendor != "none" && vendor != "redirect" {
return nil, &contract.ErrBody{Code: contract.CodeBlocked,
Message: fmt.Sprintf("目标站拦截(%s):vendor=%s status=%d", pageURL, vendor, col.status)}
}
// bodyText → markdown:纯文本段落保持(\n\n 保段;浏览器引擎不做 DOM→md 差分)。
md := textToMarkdown(ex.Text)
return &contract.RawResult{
Title: ex.Title,
Text: ex.Text,
Markdown: md,
FinalURL: firstNonEmpty(ex.FinalURL, col.finalURL, pageURL),
StatusCode: col.status,
Engine: "cdp",
Extra: map[string]any{
"vendor": vendor,
"text_len": ex.TextLen,
"load_fired": col.isLoaded(),
},
SetCookies: gotCookies,
}, nil
}
// firstNonEmpty 取第一个非空串。
func firstNonEmpty(vals ...string) string {
for _, v := range vals {
if v != "" {
return v
}
}
return ""
}
// headerHit 反爬相关头名匹配(cf-*/server/x-amzn/x-cache/location/refresh)。
var headerHitRe = regexp.MustCompile(`(?i)cf-|server|x-amzn|x-cache|location|refresh`)
func headerHit(k string) bool { return headerHitRe.MatchString(k) }
// snapshotHeaders 拷贝头片段。
func (n *navCollector) snapshotHeaders() []string {
n.mu.Lock()
defer n.mu.Unlock()
out := make([]string, len(n.headerBits))
copy(out, n.headerBits)
return out
}
func (n *navCollector) isLoaded() bool { n.mu.Lock(); defer n.mu.Unlock(); return n.loadFired }
func (n *navCollector) fired() bool { return n.loadFired }
// textToMarkdown 纯文本 → markdown 包装:非空行按段落以 \n\n 连接(bodyText 保段)。
func textToMarkdown(text string) string {
lines := strings.Split(strings.ReplaceAll(text, "\r\n", "\n"), "\n")
var paras []string
cur := make([]string, 0, 8)
flush := func() {
if len(cur) > 0 {
paras = append(paras, strings.Join(cur, "\n"))
cur = cur[:0]
}
}
for _, ln := range lines {
if strings.TrimSpace(ln) == "" {
flush()
continue
}
cur = append(cur, ln)
}
flush()
return strings.Join(paras, "\n\n")
}