Files
arcrun-collector/direct.go
T

521 lines
20 KiB
Go
Raw 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.
// direct.go — daemon「直送萃取、無 R2」同步模式(SDD ingest-hash-trigger task 11)。
//
// 既有 sync 走 collector → R2 → rag_ingest(需 R2 bucket=綁卡)。direct 模式繞開 R2/Gitea
// 監看資料夾 → 偵測新增/改動檔(沿用 Scan 的 hash 差異偵測)→ 讀檔內容 inline POST 進實例的
// rag_ingest_direct workflowLLM 萃卡 → 機械切塊 → 寫 kbdb,全在 Arcrun workflow 裡完成)。
// 刪檔 → 把 removed 事件(collector-trigger.v1POST 進實例的 rag_ingest workflow removed 分支
// (只按 page_name 讀 kbdb blocks 並標 deprecated,不碰 R2)。
//
// dogfoodingD29 daemon 薄殼豁免):本檔只做「監看/讀檔/算 hash/HTTP POST」——原生 Go。
// 萃取/切塊/RAG 一律在實例 workflowdaemon 內零 LLM/切塊邏輯。
//
// 用法:
//
// collector direct --config <config.json> [--once] [--dry-run]
//
// --once:掃一輪就退出(測試/cron 用);預設常駐輪詢(poll_interval_sec)。
// --dry-run:只列出會送出的動作,不 POST、不寫 manifest。
//
// 跨平台:純 stdlib、輪詢式偵測(不依賴 fsnotify)=零 CGodarwin/arm64、windows/amd64 直接交叉編譯。
package main
import (
"bytes"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"fmt"
"io"
"net/http"
"os"
"path/filepath"
"strings"
"time"
)
// DirectConfig 是 direct 模式的設定檔(JSON)。設定只走檔案/環境,不落 code。
type DirectConfig struct {
WatchFolder string `json:"watch_folder,omitempty"` // 監看的知識資料夾(單數舊制;與 watch_folders 至少填一)
WatchFolders []string `json:"watch_folders,omitempty"` // 監看的知識資料夾清單(daemon-beta task 1 多資料夾)
Manifest string `json:"manifest"` // manifest JSON 路徑(必填;多資料夾時為基底名,每根一份帶尾碼)
CypherURL string `json:"cypher_url"` // 實例 cypher base(必填),如 https://arcrun-cypher-executor.<acct>.workers.dev
Namespace string `json:"namespace"` // 租戶 namespace(必填),如 demo
APIKey string `json:"api_key"` // X-Arcrun-API-Key(空=沿用 namespacedemo 慣例)
Email string `json:"email,omitempty"` // 實例主身分(t26;人人記得自己的 email,CF 全程隱形)
InstanceName string `json:"instance_name,omitempty"` // 暱稱(t26 選配;不取就顯示 email)
Library string `json:"library"` // 藏書地圖歸庫鍵(空=kb;per-folder 未指定時的後備)
// t52leo 2026-07-25 裁決:「資料夾=庫」——企業有 10 個資料夾,財務/人事不准任何人看,
// 全塞一個庫=找死):每個看守資料夾對應自己的庫。key=資料夾絕對路徑,value=庫名。
// 未列出的資料夾=用資料夾名 slug 當庫名(librarySlug);再不行才退回 Library。
Libraries map[string]string `json:"libraries,omitempty"`
IngestWF string `json:"ingest_workflow"` // 直送萃取 workflow 名(空=rag_ingest_directextractor 模式不用)
RemovedWF string `json:"removed_workflow"` // 下架 workflow 名(空=rag_takedown_direct;吃 {page_name,path}
// —— 四步定稿(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)
}
// librarySlug 把資料夾名轉成合法庫名(A-Za-z0-9_-;中文等非 ASCII 轉為底線分段)。
// 空結果(如純中文名)回 "",交由 libraryFor 退回後備值——不硬造出無意義的庫名。
func librarySlug(folder string) string {
base := filepath.Base(strings.TrimRight(folder, string(filepath.Separator)))
var b strings.Builder
lastUnderscore := false
for _, r := range base {
switch {
case (r >= 'a' && r <= 'z') || (r >= '0' && r <= '9') || r == '-' || r == '_':
b.WriteRune(r)
lastUnderscore = false
case r >= 'A' && r <= 'Z':
b.WriteRune(r + 32) // 統一小寫(庫名大小寫不敏感比對較不易出錯)
lastUnderscore = false
default:
if !lastUnderscore && b.Len() > 0 {
b.WriteRune('_')
lastUnderscore = true
}
}
}
return strings.Trim(b.String(), "_")
}
// libraryFor 決定某個看守資料夾的資料該蓋哪個庫章(t52)。
// 優先序:config.Libraries 明列 > 資料夾名 slug > config.Library(後備,預設 kb)。
func (c *DirectConfig) libraryFor(absRoot string) string {
if c.Libraries != nil {
if lib, ok := c.Libraries[absRoot]; ok && strings.TrimSpace(lib) != "" {
return strings.TrimSpace(lib)
}
}
if slug := librarySlug(absRoot); slug != "" {
return slug
}
if c.Library != "" {
return c.Library
}
return "kb"
}
// expandHome 把開頭的 `~/`(或單獨的 `~`)展開成使用者家目錄的絕對路徑。
//
// t39(07-24 安裝器實案):成功頁下載的 config.json 把 manifest 寫成 `~/.arcrun-rag/manifest.json`——
// 那是給人看的寫法,**Go 不會展開波浪號**,os.ReadFile 會去找一個字面上叫 "~" 的資料夾。
// 修在 daemon 端而不是前端:用戶自己手打 config、或把設定搬到別台機器時,`~/` 都該會動;
// 前端也拿不到用戶的家目錄。展開失敗(拿不到家目錄)就原樣返回,讓後續錯誤照常誠實浮現。
func expandHome(p string) string {
if p != "~" && !strings.HasPrefix(p, "~/") {
return p
}
home, err := os.UserHomeDir()
if err != nil || home == "" {
return p
}
if p == "~" {
return home
}
return filepath.Join(home, strings.TrimPrefix(p, "~/"))
}
// LoadDirectConfig 讀設定檔並補預設值 + 基本驗證。
func LoadDirectConfig(path string) (*DirectConfig, error) {
data, err := os.ReadFile(path)
if err != nil {
return nil, fmt.Errorf("讀 config 失敗:%w", err)
}
var c DirectConfig
if err := json.Unmarshal(data, &c); err != nil {
return nil, fmt.Errorf("config JSON 解析失敗:%w", err)
}
var missing []string
if c.WatchFolder == "" && len(c.WatchFolders) == 0 {
missing = append(missing, "watch_folder(或 watch_folders")
}
if c.Manifest == "" {
missing = append(missing, "manifest")
}
if c.CypherURL == "" {
missing = append(missing, "cypher_url")
}
if c.Namespace == "" {
missing = append(missing, "namespace")
}
if len(missing) > 0 {
return nil, fmt.Errorf("config 缺必填欄位:%s", strings.Join(missing, ", "))
}
if c.APIKey == "" {
c.APIKey = c.Namespace
}
if c.Library == "" {
c.Library = "kb"
}
if c.IngestWF == "" {
c.IngestWF = "rag_ingest_direct"
}
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
}
if c.MaxRemoved <= 0 {
c.MaxRemoved = DefaultMaxRemovedRatio
}
c.CypherURL = strings.TrimSuffix(c.CypherURL, "/")
// t39:路徑欄位一律展開 `~/`(config 是給人填/人讀的,波浪號是人的寫法)
c.Manifest = expandHome(c.Manifest)
c.WatchFolder = expandHome(c.WatchFolder)
for i, p := range c.WatchFolders {
c.WatchFolders[i] = expandHome(p)
}
return &c, nil
}
// Folders 回傳監看根清單(正規化:單數舊制併入、去重、保序)。
func (c *DirectConfig) Folders() []string {
seen := map[string]bool{}
var out []string
add := func(p string) {
if p == "" || seen[p] {
return
}
seen[p] = true
out = append(out, p)
}
add(c.WatchFolder)
for _, p := range c.WatchFolders {
add(p)
}
return out
}
// manifestPathFor 回傳某根的 manifest 路徑。單根=沿用 cfg.Manifest(升級不丟既有狀態);
// 多根=每根一份,基底名加 root 絕對路徑的 sha256 前 8 碼尾碼(路徑穩定=尾碼穩定)。
func (c *DirectConfig) manifestPathFor(absRoot string) string {
folders := c.Folders()
if len(folders) <= 1 {
return c.Manifest
}
sum := sha256.Sum256([]byte(absRoot))
ext := filepath.Ext(c.Manifest)
return strings.TrimSuffix(c.Manifest, ext) + "-" + hex.EncodeToString(sum[:4]) + ext
}
// directHTTP 是 direct 模式共用的 HTTP client(萃取 workflow 可能同步跑 LLM,放寬 timeout)。
var directHTTP = &http.Client{Timeout: 300 * time.Second}
// triggerURL 組出 named-webhook 觸發完整 URL。
func (c *DirectConfig) triggerURL(workflow string) string {
return fmt.Sprintf("%s/webhooks/named/%s/%s/trigger", c.CypherURL, c.Namespace, workflow)
}
// postJSON POST 一個 JSON body 到 url,回傳 HTTP 狀態碼與回應片段。非 2xx 視為錯誤。
func (c *DirectConfig) postJSON(url string, body any) (int, string, error) {
data, err := json.Marshal(body)
if err != nil {
return 0, "", err
}
req, err := http.NewRequest(http.MethodPost, url, bytes.NewReader(data))
if err != nil {
return 0, "", err
}
req.Header.Set("Content-Type", "application/json")
req.Header.Set("X-Arcrun-API-Key", c.APIKey)
resp, err := directHTTP.Do(req)
if err != nil {
return 0, "", err
}
defer resp.Body.Close()
snippet, _ := io.ReadAll(io.LimitReader(resp.Body, 1024))
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
return resp.StatusCode, string(snippet), fmt.Errorf("HTTP %d%s", resp.StatusCode, strings.TrimSpace(string(snippet)))
}
return resp.StatusCode, string(snippet), nil
}
// pageNameOf 從相對路徑導出頁名(basename 去副檔名),與 rag_ingest collect_changed 的 pageOf 同語意。
func pageNameOf(relPath string) string {
base := relPath
if i := strings.LastIndex(base, "/"); i >= 0 {
base = base[i+1:]
}
return strings.TrimSuffix(base, filepath.Ext(base))
}
// DirectResult 是單一事件的直送結果(隨每輪 log 輸出)。
type DirectResult struct {
Root string `json:"root,omitempty"` // 多資料夾時標明事件屬於哪個根
Type string `json:"type"`
Path string `json:"path"`
Status string `json:"status"` // ingested | removed | planned | failed | skipped
HTTPStatus int `json:"http_status,omitempty"`
Error string `json:"error,omitempty"`
}
// RunDirectOnce 對每個監看根掃一輪並彙總結果(daemon-beta task 1 多資料夾)。
// 單根行為與舊制完全相同(含 manifest 路徑)。回傳彙總結果與退出碼建議(任一根失敗=1)。
func RunDirectOnce(cfg *DirectConfig, dryRun bool) ([]DirectResult, int, *TriggerPayload) {
results := []DirectResult{}
exit := 0
var lastPayload *TriggerPayload
multi := len(cfg.Folders()) > 1
for _, root := range cfg.Folders() {
r, e, p := runDirectOnceRoot(cfg, root, dryRun)
if multi {
for i := range r {
r[i].Root = root
}
}
results = append(results, r...)
if e != 0 {
exit = e
}
if p != nil {
lastPayload = p
}
}
return results, exit, lastPayload
}
// runDirectOnceRoot 對單一根掃一輪、直送 added/modified/renamed、下架 removed2xx 後回寫該根 manifest。
func runDirectOnceRoot(cfg *DirectConfig, root string, dryRun bool) ([]DirectResult, int, *TriggerPayload) {
results := []DirectResult{}
exit := 0
absRoot, err := filepath.Abs(root)
if err != nil {
return append(results, DirectResult{Status: "failed", Error: err.Error()}), 1, nil
}
absManifest, err := filepath.Abs(cfg.manifestPathFor(absRoot))
if err != nil {
return append(results, DirectResult{Status: "failed", Error: err.Error()}), 1, nil
}
m, err := LoadManifest(absManifest, absRoot)
if err != nil {
return append(results, DirectResult{Status: "failed", Error: err.Error()}), 1, nil
}
payload, err := Scan(absRoot, m, ScanOptions{
MaxRemovedRatio: cfg.MaxRemoved,
SkipPaths: map[string]bool{
absManifest: true,
// template 代裝的根層 CLAUDE.md 是 CC 設定檔,永遠不是用戶知識(task 2)
filepath.Join(absRoot, "CLAUDE.md"): true,
},
// template 代裝後 system-dev/wiki 產物區)不得被當原稿掃進 ingest(task 2
SkipDirNames: map[string]bool{"system-dev": true},
})
if err != nil {
return append(results, DirectResult{Status: "failed", Error: err.Error()}), 1, nil
}
now := time.Now().Unix()
for _, ev := range payload.Events {
switch ev.Type {
case "added", "modified", "renamed":
// renamed 在 direct 模式視同 added:內容未變但為求 kbdb 有這頁名的卡,重送一次萃取
//(頁名可能改變=要新頁名的卡)。冪等由 kbdb 端承擔(同頁名覆蓋語意)。
res := DirectResult{Type: ev.Type, Path: ev.Path}
full := filepath.Join(absRoot, filepath.FromSlash(ev.Path))
content, rerr := os.ReadFile(full)
if rerr != nil {
res.Status, res.Error = "failed", "讀檔失敗:"+rerr.Error()
results = append(results, res)
exit = 1
continue
}
if dryRun {
res.Status = "planned"
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
}
// path 帶「原檔路徑」不是卡片路徑(07-24 真機第五枚坑):
// source_uri=kb://<path> 是 takedown 的比對鍵,也是 B4 溯源該指的原文——
// 帶卡片路徑會讓「刪原檔→下架」永遠 0 命中。
status, _, perr := cfg.postJSON(cfg.triggerURL(cfg.CardIngestWF), map[string]any{
"page_name": pageNameOf(cardRel),
"path": ev.Path,
"card_content": string(cardData),
"library": cfg.libraryFor(absRoot),
})
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,
"content": string(content),
"library": cfg.libraryFor(absRoot),
})
res.HTTPStatus = status
if perr != nil {
res.Status, res.Error = "failed", perr.Error()
exit = 1
} else {
res.Status = "ingested"
m.MarkIngested(ev.Path, ev.SourceHash, now) // 2xx 才回寫(下輪不重送)
}
results = append(results, res)
case "removed":
res := DirectResult{Type: ev.Type, Path: ev.Path}
if dryRun {
res.Status = "planned"
results = append(results, res)
continue
}
// 下架=POST {page_name, path} 進 rag_takedown_direct(按 page_name 讀 kbdb blocks
// 標 deprecated,不碰 R2;獨立於 rag_ingest 的 __CARDS_PREFIX__ 閘——direct 模式檔在
// 資料夾根,會被 rag_ingest 的前綴閘擋掉,故自帶不含前綴閘的下架 workflow)。
status, _, perr := cfg.postJSON(cfg.triggerURL(cfg.RemovedWF), map[string]any{
"page_name": pageNameOf(ev.Path),
"path": ev.Path,
})
res.HTTPStatus = status
if perr != nil {
res.Status, res.Error = "failed", perr.Error()
exit = 1
} else {
res.Status = "removed"
// t15extractor 模式雲端下架成功後,同步清掉本地萃出的卡
//system-dev/wiki/cards/<頁名>.md),保持本地 wiki 與雲端一致。
// 存在才刪;刪失敗只記 warning 不擋(下架本體已成功)。
if cfg.Extractor != "" {
cardAbs := filepath.Join(absRoot, "system-dev", "wiki", "cards", pageNameOf(ev.Path)+".md")
if _, serr := os.Stat(cardAbs); serr == nil {
if rerr := os.Remove(cardAbs); rerr != nil {
results = append(results, DirectResult{
Type: "warning", Path: cardAbs, Status: "skipped",
Error: "本地卡刪除失敗(不擋下架):" + rerr.Error(),
})
}
}
}
}
results = append(results, res)
}
}
// 防呆警告輪:Scan 已壓下 removed 事件,這裡只回報警告不下架。
for _, w := range payload.Warnings {
results = append(results, DirectResult{Type: "warning", Status: "skipped", Error: w.Code + ": " + w.Message})
}
if !dryRun {
if err := m.Save(absManifest); err != nil {
results = append(results, DirectResult{Status: "failed", Error: "manifest 存檔失敗:" + err.Error()})
exit = 1
}
}
return results, exit, payload
}
// runDirect 是 `collector direct` 子命令主體。
func runDirect(args []string) int {
fs := newFlagSet()
configPath := fs.String("config", "", "direct 模式設定檔(JSON)路徑(必填)")
once := fs.Bool("once", false, "掃一輪即退出(測試/cron;預設常駐輪詢)")
dryRun := fs.Bool("dry-run", false, "只列出會送出的動作,不 POST、不寫 manifest")
if err := fs.Parse(args); err != nil {
return 2
}
if *configPath == "" {
fmt.Fprintln(os.Stderr, "錯誤:--config 為必填")
return 2
}
cfg, err := LoadDirectConfig(*configPath)
if err != nil {
fmt.Fprintln(os.Stderr, "collector direct:", err)
return 2
}
runOne := func() int {
results, exit, _ := RunDirectOnce(cfg, *dryRun)
out, _ := json.MarshalIndent(struct {
At string `json:"at"`
Folders []string `json:"folders"`
Results []DirectResult `json:"results"`
}{time.Now().Format(time.RFC3339), cfg.Folders(), results}, "", " ")
fmt.Println(string(out))
return exit
}
// 四步定稿第 1 步:daemon 代裝 template——常駐看守前確保每根都鋪好(冪等,不覆寫既有檔)。
// dry-run/--once 測試情境不代裝(不留副作用),由 template-install 子命令顯式做。
if !*once && !*dryRun {
for _, root := range cfg.Folders() {
if TemplateInstalled(root) {
continue
}
if res, ierr := InstallTemplate(root); ierr != nil {
fmt.Fprintf(os.Stderr, "template 代裝失敗(%s):%v\n", root, ierr)
} else {
fmt.Fprintf(os.Stderr, "template v%s 已鋪進 %s(新 %d 檔)\n", res.Version, root, len(res.Installed))
}
}
}
if *once {
return runOne()
}
// 常駐輪詢:純 stdlib ticker,跨平台。首輪立即跑。
fmt.Fprintf(os.Stderr, "collector direct daemon 啟動:監看 %s → %s(每 %ds 掃一輪)\n",
strings.Join(cfg.Folders(), "、"), cfg.triggerURL(cfg.IngestWF), cfg.PollSec)
runOne()
ticker := time.NewTicker(time.Duration(cfg.PollSec) * time.Second)
defer ticker.Stop()
for range ticker.C {
runOne()
}
return 0
}