🐛 fix(driver-agent): 移除 IPC 单行 8 MiB 上限,避免大批次请求打死连接

- 请求主循环由 bufio.Scanner 改为 bufio.Reader.ReadString:Scanner 的单行上限硬编码 8 MiB,
  超限时 Scan() 返回 false 使循环退出、进程终止,该连接从此永久不可用(主进程后续所有
  查询/元数据操作都返回 EOF,只能手动重连)。而主进程写入端本就无上限:一个 1000 行的导入
  批次或一个大 JSON/CLOB 单元格都能轻易超过 8 MiB
- 循环抽出为 serveAgentRequests(input, writer, runtimeState) 以便测试,并正确处理末行无换行符
- 补 3 项回归测试:12 MiB 超大行后续请求仍被处理、非法 JSON 只回错误不终止循环、末行无换行符
- 已实测确认前提:8 MiB 上限的 Scanner 对 12 MiB 单行读到 0 行并返回 token too long
This commit is contained in:
Syngnat
2026-07-26 20:42:50 +08:00
parent 9074cd9707
commit 4c93e7f175
2 changed files with 168 additions and 12 deletions

View File

@@ -0,0 +1,134 @@
package main
import (
"bufio"
"bytes"
"encoding/json"
"strings"
"testing"
"GoNavi-Wails/internal/db"
)
// decodeAgentResponses 把 JSON-lines 输出解析为响应列表。
func decodeAgentResponses(t *testing.T, payload []byte) []agentResponse {
t.Helper()
var out []agentResponse
for _, line := range strings.Split(strings.TrimSpace(string(payload)), "\n") {
line = strings.TrimSpace(line)
if line == "" {
continue
}
var resp agentResponse
if err := json.Unmarshal([]byte(line), &resp); err != nil {
t.Fatalf("解析响应失败:%v\n原始行前 200 字节:%.200s", err, line)
}
out = append(out, resp)
}
return out
}
// TestServeAgentRequestsHandlesLinesLargerThanLegacyScannerLimit 覆盖超大请求行。
//
// 回归背景agent 原先用 bufio.Scanner 且单行上限硬编码 8 MiB而主进程把整个请求
// (含 1000 行导入批次的 ChangeSet序列化成一行且无任何上限。超限时 Scanner 报
// token too long、Scan() 返回 false主循环退出、进程终止该连接从此永久不可用
// 主进程后续所有查询/元数据操作都返回「读取驱动代理响应失败EOF」只能手动重连。
func TestServeAgentRequestsHandlesLinesLargerThanLegacyScannerLimit(t *testing.T) {
const legacyScannerLimit = 8 << 20
// 构造一个远超旧上限的合法请求:用一个超大字符串参数把单行撑到 ~12 MiB。
huge := strings.Repeat("x", 12<<20)
payload, err := json.Marshal(agentRequest{ID: 1, Method: agentMethodMetadata, SessionID: huge})
if err != nil {
t.Fatalf("构造请求失败:%v", err)
}
if len(payload) <= legacyScannerLimit {
t.Fatalf("测试请求只有 %d 字节,未超过旧的 %d 字节上限,无法覆盖该回归", len(payload), legacyScannerLimit)
}
// 超大行之后再跟一个正常请求:若主循环因超限退出,第二个请求就不会被处理。
followUp, err := json.Marshal(agentRequest{ID: 2, Method: agentMethodMetadata})
if err != nil {
t.Fatalf("构造后续请求失败:%v", err)
}
input := bytes.NewReader(append(append(payload, '\n'), append(followUp, '\n')...))
var out bytes.Buffer
writer := bufio.NewWriter(&out)
runtimeState := &agentRuntime{sessions: make(map[string]db.StatementExecer)}
if err := serveAgentRequests(input, writer, runtimeState); err != nil {
t.Fatalf("serveAgentRequests 返回错误:%v", err)
}
if err := writer.Flush(); err != nil {
t.Fatalf("Flush 失败:%v", err)
}
responses := decodeAgentResponses(t, out.Bytes())
if len(responses) != 2 {
t.Fatalf("响应条数 = %d期望 2超大行后主循环提前退出连接被打死", len(responses))
}
if responses[0].ID != 1 {
t.Errorf("第一条响应 ID = %d期望 1", responses[0].ID)
}
if responses[1].ID != 2 {
t.Errorf("第二条响应 ID = %d期望 2后续请求未被处理说明循环已退出", responses[1].ID)
}
}
// TestServeAgentRequestsReportsInvalidJSONWithoutExiting 非法 JSON 只回错误响应,不得终止循环。
func TestServeAgentRequestsReportsInvalidJSONWithoutExiting(t *testing.T) {
followUp, err := json.Marshal(agentRequest{ID: 9, Method: agentMethodMetadata})
if err != nil {
t.Fatalf("构造请求失败:%v", err)
}
input := strings.NewReader("{not json}\n\n" + string(followUp) + "\n")
var out bytes.Buffer
writer := bufio.NewWriter(&out)
runtimeState := &agentRuntime{sessions: make(map[string]db.StatementExecer)}
if err := serveAgentRequests(input, writer, runtimeState); err != nil {
t.Fatalf("serveAgentRequests 返回错误:%v", err)
}
if err := writer.Flush(); err != nil {
t.Fatalf("Flush 失败:%v", err)
}
responses := decodeAgentResponses(t, out.Bytes())
if len(responses) != 2 {
t.Fatalf("响应条数 = %d期望 2解析失败一条 + 后续正常一条)", len(responses))
}
if responses[0].Success {
t.Error("非法 JSON 的响应不应是 Success")
}
if responses[1].ID != 9 || !responses[1].Success {
t.Errorf("后续正常请求未被处理:%#v", responses[1])
}
}
// TestServeAgentRequestsHandlesFinalLineWithoutNewline 末行无换行符时仍须处理。
func TestServeAgentRequestsHandlesFinalLineWithoutNewline(t *testing.T) {
payload, err := json.Marshal(agentRequest{ID: 5, Method: agentMethodMetadata})
if err != nil {
t.Fatalf("构造请求失败:%v", err)
}
input := bytes.NewReader(payload) // 故意不加末尾换行
var out bytes.Buffer
writer := bufio.NewWriter(&out)
runtimeState := &agentRuntime{sessions: make(map[string]db.StatementExecer)}
if err := serveAgentRequests(input, writer, runtimeState); err != nil {
t.Fatalf("serveAgentRequests 返回错误:%v", err)
}
if err := writer.Flush(); err != nil {
t.Fatalf("Flush 失败:%v", err)
}
responses := decodeAgentResponses(t, out.Bytes())
if len(responses) != 1 || responses[0].ID != 5 {
t.Fatalf("末行无换行符的请求未被处理:%#v", responses)
}
}

View File

@@ -4,7 +4,9 @@ import (
"bufio"
"context"
"encoding/json"
"errors"
"fmt"
"io"
"os"
"reflect"
"runtime"
@@ -121,16 +123,42 @@ func main() {
debug.SetGCPercent(50)
db.InitMemorySoftLimit(db.MemorySoftLimitInitialBytes)
scanner := bufio.NewScanner(os.Stdin)
scanner.Buffer(make([]byte, 0, 16<<10), 8<<20)
writer := bufio.NewWriter(os.Stdout)
defer writer.Flush()
runtimeState := &agentRuntime{
sessions: make(map[string]db.StatementExecer),
}
for scanner.Scan() {
line := strings.TrimSpace(scanner.Text())
readErr := serveAgentRequests(os.Stdin, writer, runtimeState)
runtimeState.close()
if readErr != nil {
fmt.Fprintf(os.Stderr, "读取请求失败:%v\n", readErr)
}
}
// serveAgentRequests 是 driver-agent 的 JSON-lines 请求主循环。
//
// 用 bufio.Reader 而非 bufio.ScannerScanner 有单行长度上限(原为 8 MiB超限时
// Scan() 直接返回 false 使主循环退出、进程终止,该连接从此永久不可用(主进程侧随后所有
// 请求都拿到 EOF只能手动重连。而主进程写入端没有任何上限一个 1000 行的导入批次或
// 一个大 JSON/CLOB 单元格都能轻易超过 8 MiB。bufio.Reader.ReadString 会按需增长,无单行上限。
//
// 返回非 nil 表示读取过程出现了非 EOF 的真实错误。
func serveAgentRequests(input io.Reader, writer *bufio.Writer, runtimeState *agentRuntime) error {
reader := bufio.NewReaderSize(input, 16<<10)
for {
raw, err := reader.ReadString('\n')
if err != nil && len(raw) == 0 {
if errors.Is(err, io.EOF) {
return nil
}
return err
}
// err != nil 但 raw 非空:最后一行没有换行符,仍需处理;
// 下一轮 ReadString 会立即以 len(raw)==0 返回并结束循环。
line := strings.TrimSpace(raw)
if line == "" {
continue
}
@@ -148,7 +176,7 @@ func main() {
if strings.TrimSpace(req.Method) == agentMethodStreamQuery {
if err := handleStreamRequest(runtimeState, req, writer); err != nil {
fmt.Fprintf(os.Stderr, "写入流式响应失败:%v\n", err)
break
return nil
}
continue
}
@@ -156,18 +184,12 @@ func main() {
resp := handleRequest(runtimeState, req)
if err := writeResponse(writer, resp); err != nil {
fmt.Fprintf(os.Stderr, "写入响应失败:%v\n", err)
break
return nil
}
if strings.TrimSpace(req.Method) == agentMethodQuery {
maybeReleaseAgentMemory("query-response", countAgentResponseRows(resp.Data))
}
}
runtimeState.close()
if err := scanner.Err(); err != nil {
fmt.Fprintf(os.Stderr, "读取请求失败:%v\n", err)
}
}
func handleRequest(runtimeState *agentRuntime, req agentRequest) agentResponse {