Files
arcrun-collector/direct.go
T
Leo 1738acaca7 fix(t108 🔴🔴): 逐帳號同步遺失萃取設定→原文直送——三層修+契約保險絲
事故(leo 機 07-28 17:3x 總管現場抓到):t104 config 重寫吃掉機器層 extractor/key、
per-account DirectConfig 不繼承 ⇒ extractor 空 ⇒ 6 筆走 rag_ingest_direct 原文出機。
修:①config 讀改寫全程保留既有欄位(帶測試)②帳號層無值垂直繼承機器層
extractor/claude_bin/gemini_api_key/llm_model/CardIngestWF/RemovedWF(帶 fake server 測試:
必須收到 rag_ingest_card 非 direct)③契約保險絲:extractor 空時非 .md/.txt 禁直送、
標 failed「萃取器未設定,已跳過(不直送原文)」——同類 bug 永不再成外洩。
三模組 go test 全綠(總管親跑)。(實作=子 CC;驗證+commit=總管)
2026-07-28 17:59:42 +08:00

829 lines
31 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"
"net/url"
"os"
"path/filepath"
"strings"
"time"
)
// AccountConfig 單一帳號的連線設定(t104 多帳號同時看守)。
// 每個帳號代表一個 Arcrun 知識庫實例;Extractor/Manifest 等機器層級設定住在 DirectConfig 頂層。
type AccountConfig struct {
InstanceName string `json:"instance_name,omitempty"`
Email string `json:"email,omitempty"`
CypherURL string `json:"cypher_url"`
Namespace string `json:"namespace"`
APIKey string `json:"api_key,omitempty"`
WatchFolders []string `json:"watch_folders,omitempty"` // 此帳號看守的資料夾(多根)
Libraries map[string]string `json:"libraries,omitempty"` // 資料夾→庫對映(t52
}
// DirectConfig 是 direct 模式的設定檔(JSON)。設定只走檔案/環境,不落 code。
type DirectConfig struct {
// t104:多帳號清單(新制)。有值時頂層連線欄位僅保留讀取相容——
// 啟動時若無 Accounts 但有舊的頂層 CypherURLLoadDirectConfig 自動包成 Accounts[0]。
Accounts []AccountConfig `json:"accounts,omitempty"`
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,omitempty"` // 實例 cypher base(舊制;新制走 Accounts
Namespace string `json:"namespace,omitempty"` // 租戶 namespace(舊制;新制走 Accounts
APIKey string `json:"api_key,omitempty"` // X-Arcrun-API-Key(舊制)
Email string `json:"email,omitempty"` // 實例主身分(舊制)
InstanceName string `json:"instance_name,omitempty"` // 暱稱(舊制)
Library string `json:"library"` // 藏書地圖歸庫鍵(空=kb;per-folder 未指定時的後備)
// t52:每個看守資料夾對應自己的庫(key=絕對路徑,value=庫名);新制走 AccountConfig.Libraries。
Libraries map[string]string `json:"libraries,omitempty"`
IngestWF string `json:"ingest_workflow"` // 直送萃取 workflow 名(空=rag_ingest_direct
RemovedWF string `json:"removed_workflow"` // 下架 workflow 名(空=rag_takedown_direct
// —— 四步定稿(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 改生路徑穩定雜湊鍵(t89)。
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 > 路徑穩定雜湊鍵(t89)。
//
// t89leo 2026-07-28 實測):純中文資料夾名(如「官方總圖」)slug 後為空字串,
// 若退到 c.Library/"kb",兩個不同純中文資料夾會塌縮進同一個庫(靜默混庫)。
// 後端庫名鍵限 A-Za-z0-9_-Arcrun KBDB 約束),故改生 lib_+sha256(絕對路徑) 前 6 hex——
// 路徑穩定所以鍵穩定;兩個不同路徑必然不同鍵。
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
}
sum := sha256.Sum256([]byte(absRoot))
return "lib_" + hex.EncodeToString(sum[:3])
}
// 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 讀設定檔並補預設值 + 基本驗證。
// t104 向後相容遷移:若無 Accounts 但有舊的頂層 CypherURL,自動包成 Accounts[0](記憶體遷移;
// 磁碟回寫由呼叫端在合適時機(如 saveDirectConfig)完成)。
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)
}
// t39:路徑欄位一律展開 `~/`(先展開才能正確計算 Folders())
c.Manifest = expandHome(c.Manifest)
c.WatchFolder = expandHome(c.WatchFolder)
for i, p := range c.WatchFolders {
c.WatchFolders[i] = expandHome(p)
}
for i := range c.Accounts {
for j, p := range c.Accounts[i].WatchFolders {
c.Accounts[i].WatchFolders[j] = expandHome(p)
}
}
// t104:舊格式遷移——頂層 CypherURL → Accounts[0](冪等:有 Accounts 就跳過)
if len(c.Accounts) == 0 && c.CypherURL != "" {
c.Accounts = []AccountConfig{{
InstanceName: c.InstanceName,
Email: c.Email,
CypherURL: c.CypherURL,
Namespace: c.Namespace,
APIKey: c.APIKey,
Libraries: c.Libraries,
WatchFolders: c.Folders(), // 正規化後的監看清單
}}
}
// 驗證
var missing []string
if c.Manifest == "" {
missing = append(missing, "manifest")
}
if len(c.Accounts) == 0 {
missing = append(missing, "accounts(或 cypher_url 連線設定)")
}
for i, acc := range c.Accounts {
if acc.CypherURL == "" {
missing = append(missing, fmt.Sprintf("accounts[%d] 缺 cypher_url", i))
}
if acc.Namespace == "" {
missing = append(missing, fmt.Sprintf("accounts[%d] 缺 namespace", i))
}
}
if len(missing) > 0 {
return nil, fmt.Errorf("config 缺必填欄位:%s", strings.Join(missing, ", "))
}
// 補各帳號預設值
for i := range c.Accounts {
if c.Accounts[i].APIKey == "" {
c.Accounts[i].APIKey = c.Accounts[i].Namespace
}
c.Accounts[i].CypherURL = strings.TrimSuffix(c.Accounts[i].CypherURL, "/")
}
// 頂層後備值(舊制相容或被 makeAccountSubConfig 繼承)
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
}
if c.CypherURL != "" {
c.CypherURL = strings.TrimSuffix(c.CypherURL, "/")
}
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
}
// instanceHostOf extracts the host from a CypherURL to use as a per-instance
// distinguisher in manifest paths (t86b). Falls back to the full URL if parsing fails.
func instanceHostOf(cypherURL string) string {
u, err := url.Parse(cypherURL)
if err != nil || u.Host == "" {
return cypherURL
}
return u.Host
}
// manifestPathFor 回傳某根的 manifest 路徑。
// t86b:雜湊改為 sha256(instanceHost + "\n" + absRoot),讓每個(實例, 資料夾)組合對應
// 獨立帳本——換知識庫實例後同資料夾不再重用舊帳本,避免「全部視為已同步」靜默跳過。
// 舊格式遷移見 migrateManifestIfNeeded:啟動/掃描時若新名不存在但舊名存在則 rename 過來。
func (c *DirectConfig) manifestPathFor(absRoot string) string {
host := instanceHostOf(c.CypherURL)
sum := sha256.Sum256([]byte(host + "\n" + absRoot))
ext := filepath.Ext(c.Manifest)
return strings.TrimSuffix(c.Manifest, ext) + "-" + hex.EncodeToString(sum[:4]) + ext
}
// oldManifestPaths 回傳 t86b 之前版本對同一個 absRoot 會產生的 manifest 路徑,
// 供遷移時確認是否有舊帳本需要搬移。
// 優先序:①舊多根(只含路徑雜湊)→ ②舊單根(直用 cfg.Manifest)。
func (c *DirectConfig) oldManifestPaths(absRoot string) []string {
ext := filepath.Ext(c.Manifest)
base := strings.TrimSuffix(c.Manifest, ext)
// 舊多根公式:sha256(absRoot only)
sumPathOnly := sha256.Sum256([]byte(absRoot))
pathOnlyPath := base + "-" + hex.EncodeToString(sumPathOnly[:4]) + ext
// 舊單根公式:直用 cfg.Manifest
return []string{pathOnlyPath, c.Manifest}
}
// migrateManifestIfNeeded 在 newPath 不存在時,把最先找到的舊格式 manifest rename 過來。
// 一次性、冪等:newPath 已存在時直接 return;rename 失敗靜默忽略(最多這一輪重傳,不影響正確性)。
func (c *DirectConfig) migrateManifestIfNeeded(absRoot, newPath string) {
if _, err := os.Stat(newPath); err == nil {
return // 新路徑已存在,無需遷移
}
for _, oldPath := range c.oldManifestPaths(absRoot) {
if oldPath == newPath {
continue
}
if _, err := os.Stat(oldPath); err == nil {
_ = os.Rename(oldPath, newPath)
return
}
}
}
// 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 {
Account string `json:"account,omitempty"` // t104: cypher_url host(多帳號標的)
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"`
}
// makeAccountSubConfig 從帳號設定建出單帳號用的 DirectConfig,繼承機器層級欄位(t104)。
// 用於 RunDirectOnce 逐帳號掃描,每帳號得到獨立的 CypherURL/Namespace/WatchFolders 等。
func (c *DirectConfig) makeAccountSubConfig(acc AccountConfig) *DirectConfig {
sub := *c // 複製機器層級欄位
sub.Accounts = nil
sub.CypherURL = strings.TrimSuffix(acc.CypherURL, "/")
sub.Namespace = acc.Namespace
sub.APIKey = acc.APIKey
if sub.APIKey == "" {
sub.APIKey = sub.Namespace
}
sub.Email = acc.Email
sub.InstanceName = acc.InstanceName
sub.WatchFolder = ""
sub.WatchFolders = acc.WatchFolders
sub.Libraries = acc.Libraries
if sub.Library == "" {
sub.Library = "kb"
}
return &sub
}
// RunDirectOnce 對每個帳號的每個監看根掃一輪並彙總結果(t104 多帳號同時看守)。
// 單帳號行為與舊制完全相同(含 manifest 路徑)。回傳彙總結果與退出碼建議(任一根失敗=1)。
// 額外:
// - 預檢 extractor 可用性(t92-②),有 fallback 時更新 cfg.ClaudeBinin-memory,呼叫端存檔)。
// - 每輪結束寫 ~/.arcrun-rag/status.jsont91 狀態可見性,含 per-account 雲端版本)。
// - 一個帳號失敗不擋其他帳號繼續同步(t104 隔離)。
func RunDirectOnce(cfg *DirectConfig, dryRun bool) ([]DirectResult, int, *TriggerPayload) {
results := []DirectResult{}
exit := 0
var lastPayload *TriggerPayload
// t92-②:預檢 extractor(機器層級),有 fallback 路徑時就地更新 cfg.ClaudeBin。
extractorOK := true
extractorError := ""
if cfg.Extractor == "claude" {
resolved, ferr := FindClaudeBin(cfg.ClaudeBin)
if ferr != nil {
extractorOK = false
extractorError = "找不到 Claude 指令——請確認 Claude Code 已安裝,或改用 Gemma 萃取路"
} else if resolved != cfg.ClaudeBin {
cfg.ClaudeBin = resolved // in-memory 回寫;runDirect 偵到變化才存磁碟
}
} else if cfg.Extractor == "gemma" {
if strings.TrimSpace(cfg.GeminiAPIKey) == "" {
extractorOK = false
extractorError = "金鑰是空的——請在設定裡輸入 Gemini API Key"
}
}
// t104:解出有效帳號清單(向後相容:無 Accounts 但有頂層 CypherURL 時視為單帳號)
accounts := cfg.Accounts
if len(accounts) == 0 && cfg.CypherURL != "" {
accounts = []AccountConfig{{
InstanceName: cfg.InstanceName,
Email: cfg.Email,
CypherURL: cfg.CypherURL,
Namespace: cfg.Namespace,
APIKey: cfg.APIKey,
Libraries: cfg.Libraries,
WatchFolders: cfg.Folders(),
}}
}
accountDetails := map[string]AccountSyncStatus{}
for _, acc := range accounts {
if acc.CypherURL == "" || acc.Namespace == "" {
continue // 跳過設定不完整的帳號
}
accCfg := cfg.makeAccountSubConfig(acc)
accHost := instanceHostOf(acc.CypherURL)
// t103per-account 雲端版本偵測
cloudVer, cloudOK := fetchCloudVersion(accCfg.CypherURL)
accSt := AccountSyncStatus{
LastSync: time.Now().Format(time.RFC3339),
CloudVersion: cloudVer,
CloudCheckOK: cloudOK,
}
multi := len(accCfg.Folders()) > 1
for _, root := range accCfg.Folders() {
r, e, p := runDirectOnceRoot(accCfg, root, dryRun)
if multi {
for i := range r {
r[i].Root = root
}
}
for i := range r {
r[i].Account = accHost // t104: 標明所屬帳號
if cfg.Extractor != "" {
switch r[i].Status {
case "ingested":
accSt.ExtractedOK++
case "failed":
accSt.ExtractFailed++
}
}
}
results = append(results, r...)
if e != 0 {
exit = e // 任一帳號任一根失敗=整體 exit 1,但不停其他帳號
}
if p != nil {
lastPayload = p
}
}
accountDetails[accHost] = accSt
}
// t91:每輪寫狀態檔(含 per-account 雲端版本與萃取計數)。
if !dryRun && cfg.Manifest != "" {
st := SyncStatus{
LastSync: time.Now().Format(time.RFC3339),
ExtractorOK: extractorOK,
ExtractorError: extractorError,
AccountDetails: accountDetails,
}
// 頂層彙總(向後相容:單帳號時填頂層欄位讓舊版 tray 仍能讀)
if cfg.Extractor != "" {
for _, r := range results {
switch r.Status {
case "ingested":
st.ExtractedOK++
case "failed":
st.ExtractFailed++
st.Failures = append(st.Failures, ExtractFail{
Path: r.Path,
Error: shortError(r.Error),
})
}
}
}
// 單帳號時把 cloud version 也填頂層(向後相容)
if len(accountDetails) == 1 {
for _, v := range accountDetails {
st.CloudVersion = v.CloudVersion
st.CloudCheckOK = v.CloudCheckOK
break
}
}
if serr := SaveSyncStatus(StatusFilePath(cfg.Manifest), st); serr != nil {
fmt.Fprintf(os.Stderr, "status 寫入失敗(不擋看守):%v\n", serr)
}
}
return results, exit, lastPayload
}
// shortError 把錯誤字串截為一句話(供 UI 顯示,不要超過 120 字)。
func shortError(msg string) string {
if len([]rune(msg)) <= 120 {
return msg
}
runes := []rune(msg)
return string(runes[:120]) + "…"
}
// CheckExtractor 預檢萃取器是否可用(不執行萃取、不打 API)。
// 只檢「可執行檔存在且可執行」或「金鑰非空」。
// extractor 空(舊制直送)一律回 (true, "")。
func CheckExtractor(cfg *DirectConfig) (ok bool, errMsg string) {
switch cfg.Extractor {
case "claude":
if _, err := FindClaudeBin(cfg.ClaudeBin); err != nil {
return false, "找不到 Claude 指令——請確認 Claude Code 已安裝,或改用 Gemma 萃取路"
}
return true, ""
case "gemma":
if strings.TrimSpace(cfg.GeminiAPIKey) == "" {
return false, "金鑰是空的——請在設定裡輸入 Gemini API Key"
}
return true, ""
default:
return true, ""
}
}
// saveDirectConfig 把 DirectConfig 回寫到 configPatht92:找到 claude fallback 路徑後持久化)。
// 只寫 claude_bin 等萃取相關欄位不會影響用戶的其他設定(JSON 完整覆蓋整個 config)。
func saveDirectConfig(configPath string, cfg *DirectConfig) error {
data, err := json.MarshalIndent(cfg, "", " ")
if err != nil {
return err
}
return os.WriteFile(configPath, data, 0o600)
}
// 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
}
cfg.migrateManifestIfNeeded(absRoot, absManifest) // t86b:一次性遷移舊格式帳本
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"
// 記下是誰萃的(t73/leo 07-27):換萃取器時才分辨得出哪些卡是舊的。
m.MarkIngestedBy(ev.Path, ev.SourceHash, now, cfg.Extractor)
} else {
exit = 1
}
results = append(results, res)
continue
}
// t108 防禦閘:extractor 未設定時,非 .md/.txt 檔禁止直送原文(原文外洩保險絲)。
// 讓同類 bug 永遠不再變成內容外洩,而是明確的 failed 狀態。
if ext := strings.ToLower(filepath.Ext(ev.Path)); ext != ".md" && ext != ".txt" {
res.Status = "failed"
res.Error = "萃取器未設定,已跳過(不直送原文)"
results = append(results, res)
exit = 1
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
}
// t92:第一輪若 claude 找到了 fallback 路徑,把更新後的 claude_bin 存回 config 檔(下次直達)。
origClaudeBin := cfg.ClaudeBin
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
}
// claude_bin 回寫:只在找到 fallback 路徑時才存(避免頻繁寫磁碟)
persistClaudeBinIfChanged := func() {
if cfg.ClaudeBin != origClaudeBin && cfg.ClaudeBin != "" && *configPath != "" {
if serr := saveDirectConfig(*configPath, cfg); serr != nil {
fmt.Fprintf(os.Stderr, "claude_bin 回寫 config 失敗(不擋看守):%v\n", serr)
} else {
origClaudeBin = cfg.ClaudeBin // 只存一次
}
}
}
// 四步定稿第 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 {
code := runOne()
persistClaudeBinIfChanged()
return code
}
// 常駐輪詢:純 stdlib,跨平台。首輪立即跑。
// t98:每 1 秒檢查一次訊號檔(SyncNowSignalPath),命中即立刻跑一輪並刪檔;
// 否則依 PollSec 間隔照舊定時跑。這讓 tray「立刻同步」按鈕能即時觸發,
// 不需引入 IPC/socket 等平台依賴。
fmt.Fprintf(os.Stderr, "collector direct daemon 啟動:監看 %s → %s(每 %ds 掃一輪)\n",
strings.Join(cfg.Folders(), "、"), cfg.triggerURL(cfg.IngestWF), cfg.PollSec)
runOne()
persistClaudeBinIfChanged() // 第一輪後立即回寫(下次重啟直達)
signalPath := SyncNowSignalPath(cfg.Manifest)
pollInterval := time.Duration(cfg.PollSec) * time.Second
lastRun := time.Now()
const checkInterval = time.Second
for {
time.Sleep(checkInterval)
if consumeSyncNowSignal(signalPath) {
runOne()
lastRun = time.Now()
continue
}
if time.Since(lastRun) >= pollInterval {
runOne()
lastRun = time.Now()
}
}
}
// consumeSyncNowSignal 檢查訊號檔是否存在:存在則刪除並回 true(呼叫端立刻跑一輪同步);
// 不存在回 false。刪除失敗也回 true——確保本輪至少跑一次,下次若殘留再刪。
func consumeSyncNowSignal(path string) bool {
if _, err := os.Stat(path); err != nil {
return false
}
_ = os.Remove(path)
return true
}