公开页抓取改为 Chrome 136 自洽身份 + 每节点 SQLite 养罐,Trafilatura 走 curl_cffi;MCP/README/接手说明与 09-02 现网复测对齐,避免消费方继续抄过期的站点三分表。 Co-authored-by: Cursor <cursoragent@cursor.com>
325 lines
11 KiB
Go
325 lines
11 KiB
Go
// 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")
|
||
}
|