feat(daemon-beta t5/t6): rag_ingest_card 收卡 workflow(零LLM零credential)+direct 迴圈接萃取器——本地萃卡只上卡、原文不出機;e2e 測試驗「原文不洩上雲」+失敗重試語意

This commit is contained in:
2026-07-24 12:23:52 +08:00
parent be47afe26e
commit 058d19b743
2 changed files with 174 additions and 3 deletions
+64 -3
View File
@@ -42,10 +42,16 @@ type DirectConfig struct {
Namespace string `json:"namespace"` // 租戶 namespace(必填),如 demo
APIKey string `json:"api_key"` // X-Arcrun-API-Key(空=沿用 namespacedemo 慣例)
Library string `json:"library"` // 藏書地圖歸庫鍵(空=kb
IngestWF string `json:"ingest_workflow"` // 直送萃取 workflow 名(空=rag_ingest_direct
IngestWF string `json:"ingest_workflow"` // 直送萃取 workflow 名(空=rag_ingest_directextractor 模式不用
RemovedWF string `json:"removed_workflow"` // 下架 workflow 名(空=rag_takedown_direct;吃 {page_name,path}
PollSec int `json:"poll_interval_sec"` // 輪詢間隔秒(空/05
MaxRemoved float64 `json:"max_removed_ratio"` // 大量刪除防呆門檻(空/00.4
// —— 四步定稿(daemon-beta t3/t4/t6):本地萃卡模式 ——
Extractor string `json:"extractor,omitempty"` // "claude""gemma";空=舊制(內容直送雲端萃
ClaudeBin string `json:"claude_bin,omitempty"` // claude 執行檔(空=PATH 找 claude
GeminiAPIKey string `json:"gemini_api_key,omitempty"` // gemma 路的用戶 key
LLMModel string `json:"llm_model,omitempty"` // gemma 路模型(空=gemma-4-31b-it
CardIngestWF string `json:"card_ingest_workflow,omitempty"` // 收卡 workflow(空=rag_ingest_card
PollSec int `json:"poll_interval_sec"` // 輪詢間隔秒(空/05
MaxRemoved float64 `json:"max_removed_ratio"` // 大量刪除防呆門檻(空/0=0.4)
}
// LoadDirectConfig 讀設定檔並補預設值 + 基本驗證。
@@ -86,6 +92,15 @@ func LoadDirectConfig(path string) (*DirectConfig, error) {
if c.RemovedWF == "" {
c.RemovedWF = "rag_takedown_direct"
}
if c.CardIngestWF == "" {
c.CardIngestWF = "rag_ingest_card"
}
if c.LLMModel == "" {
c.LLMModel = defaultLLMModel
}
if c.Extractor != "" && c.Extractor != "claude" && c.Extractor != "gemma" {
return nil, fmt.Errorf("extractor 只能是 claude / gemma(或留空走舊制),got %q", c.Extractor)
}
if c.PollSec <= 0 {
c.PollSec = 5
}
@@ -253,6 +268,52 @@ func runDirectOnceRoot(cfg *DirectConfig, root string, dryRun bool) ([]DirectRes
results = append(results, res)
continue
}
if cfg.Extractor != "" {
// 四步定稿:本地萃卡 → 每張卡 POST rag_ingest_card(原文不出機)
var cards []string
var xerr error
switch cfg.Extractor {
case "claude":
cards, xerr = ExtractWithClaude(cfg.ClaudeBin, absRoot, ev.Path)
case "gemma":
cards, xerr = ExtractWithGemma(cfg.GeminiAPIKey, cfg.LLMModel, absRoot, ev.Path)
}
if xerr != nil {
res.Status, res.Error = "failed", "本地萃取失敗:"+xerr.Error()
results = append(results, res)
exit = 1
continue
}
ok := true
for _, cardRel := range cards {
cardData, cerr := os.ReadFile(filepath.Join(absRoot, filepath.FromSlash(cardRel)))
if cerr != nil {
res.Status, res.Error = "failed", "讀卡片失敗:"+cerr.Error()
ok = false
break
}
status, _, perr := cfg.postJSON(cfg.triggerURL(cfg.CardIngestWF), map[string]any{
"page_name": pageNameOf(cardRel),
"path": cardRel,
"card_content": string(cardData),
"library": cfg.Library,
})
res.HTTPStatus = status
if perr != nil {
res.Status, res.Error = "failed", perr.Error()
ok = false
break
}
}
if ok {
res.Status = "ingested"
m.MarkIngested(ev.Path, ev.SourceHash, now)
} else {
exit = 1
}
results = append(results, res)
continue
}
status, _, perr := cfg.postJSON(cfg.triggerURL(cfg.IngestWF), map[string]any{
"page_name": pageNameOf(ev.Path),
"path": ev.Path,
+110
View File
@@ -0,0 +1,110 @@
// direct_extract_test.go — task 6extractor 模式端到端(本地萃卡→POST rag_ingest_card)。
package main
import (
"encoding/json"
"io"
"net/http"
"net/http/httptest"
"os"
"path/filepath"
"strings"
"testing"
)
// 完整鏈(claude stub 版):丟原稿 → 萃卡落地本地 → 只有「卡片」被 POST 到 rag_ingest_card
// → 原文從未離開本機 → manifest 標 ingested(下一輪不重送)。
func TestDirectExtractorModeE2E(t *testing.T) {
root := t.TempDir()
if err := os.WriteFile(filepath.Join(root, "報銷規則.md"), []byte("# 原稿機密內容 XYZZY"), 0o644); err != nil {
t.Fatal(err)
}
// 假 cypher:收 rag_ingest_card、驗 payload、記帳
var posted []map[string]any
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if !strings.HasSuffix(r.URL.Path, "/webhooks/named/demo/rag_ingest_card/trigger") {
t.Errorf("打錯端點:%s", r.URL.Path)
}
body, _ := io.ReadAll(r.Body)
var m map[string]any
_ = json.Unmarshal(body, &m)
posted = append(posted, m)
_ = json.NewEncoder(w).Encode(map[string]any{"success": true})
}))
defer srv.Close()
// stub claude:把原稿萃成卡(模擬 /wiki-capture 行為)
stubDir := t.TempDir()
stub := filepath.Join(stubDir, "claude")
script := "#!/bin/sh\nmkdir -p system-dev/wiki/cards\nprintf '# 報銷規則\\n## 一句話定義\\n測試卡\\n## 關聯\\n- 報銷規則 >> 屬於 >> 財務\\n' > 'system-dev/wiki/cards/報銷規則.md'\n"
if err := os.WriteFile(stub, []byte(script), 0o755); err != nil {
t.Fatal(err)
}
cfg := &DirectConfig{
WatchFolders: []string{root},
Manifest: filepath.Join(t.TempDir(), "m.json"),
CypherURL: srv.URL, Namespace: "demo", APIKey: "demo",
Library: "kb", Extractor: "claude", ClaudeBin: stub,
CardIngestWF: "rag_ingest_card", MaxRemoved: DefaultMaxRemovedRatio,
}
results, exit, _ := RunDirectOnce(cfg, false)
if exit != 0 {
t.Fatalf("exit=%d results=%+v", exit, results)
}
if len(results) != 1 || results[0].Status != "ingested" {
t.Fatalf("results=%+v", results)
}
// 卡片落地本地(用戶看得到自己的 wiki)
if _, err := os.Stat(filepath.Join(root, "system-dev", "wiki", "cards", "報銷規則.md")); err != nil {
t.Fatalf("卡片未落地:%v", err)
}
// 上雲的是卡片、不是原文
if len(posted) != 1 {
t.Fatalf("應恰好 POST 一張卡,got %d", len(posted))
}
cc, _ := posted[0]["card_content"].(string)
if !strings.Contains(cc, "## 一句話定義") {
t.Fatalf("card_content 不是卡片:%.80s", cc)
}
if strings.Contains(cc, "XYZZY") {
t.Fatal("原文內容洩上雲=違反四步定稿邊界")
}
if p, _ := posted[0]["path"].(string); p != "system-dev/wiki/cards/報銷規則.md" {
t.Fatalf("path=%q", p)
}
// 第二輪:原稿沒變 → 不重萃不重送
results2, exit2, _ := RunDirectOnce(cfg, false)
if exit2 != 0 || len(results2) != 0 || len(posted) != 1 {
t.Fatalf("第二輪應零事件:results=%+v posted=%d", results2, len(posted))
}
}
// 萃取失敗=該檔標 failed、exit=1、manifest 不標(下輪重試),其他檔不受影響。
func TestDirectExtractorFailKeepsRetry(t *testing.T) {
root := t.TempDir()
if err := os.WriteFile(filepath.Join(root, "a.md"), []byte("x"), 0o644); err != nil {
t.Fatal(err)
}
stubDir := t.TempDir()
stub := filepath.Join(stubDir, "claude")
if err := os.WriteFile(stub, []byte("#!/bin/sh\nexit 3\n"), 0o755); err != nil {
t.Fatal(err)
}
cfg := &DirectConfig{
WatchFolders: []string{root},
Manifest: filepath.Join(t.TempDir(), "m.json"),
CypherURL: "https://x.example", Namespace: "demo", APIKey: "demo",
Extractor: "claude", ClaudeBin: stub, MaxRemoved: DefaultMaxRemovedRatio,
}
results, exit, _ := RunDirectOnce(cfg, false)
if exit != 1 || len(results) != 1 || results[0].Status != "failed" {
t.Fatalf("exit=%d results=%+v", exit, results)
}
// 再跑一輪:仍是同一個事件(manifest 沒標 ingested=會重試)
results2, _, _ := RunDirectOnce(cfg, false)
if len(results2) != 1 {
t.Fatalf("失敗檔應重試:%+v", results2)
}
}