Files
arcrun-collector/direct.go
T
Leo 1a3e974941 t178:托盤標籤與錯誤訊息對齊實際行為(leo 封測者實撞)
【leo 08-04】封測者是台大資工碩士,仍搞不清楚 ⇒「一般人就完蛋了」。
他的托盤顯示「oscar · Claude」但他根本沒有 Claude,
同時錯誤說「Gemini 金鑰是空的」——**兩個訊息互相矛盾**。

① 標籤與行為脫鉤(accountEngineLabel)
   direct.go 早就把 claude 正規化成 gemma(t176:地端先只支援 Gemini),
   但標籤照 config 舊字串念 ⇒ 顯示 Claude、實際走 gemma。
   舊值來源=雲端舊版下發後留在 config.json 的殘留(t176 擋了新寫入、沒洗舊值)。
   修:claude 也顯示「· Gemini」,與萃取實際走的路一致。

② 錯誤訊息沒說設定在哪
   舊:「請在設定裡輸入 Gemini API Key」——用戶找不到入口。
   新:「點托盤選單的『AI 設定…』貼上金鑰(免費申請:aistudio.google.com/apikey)」
   ——錯誤訊息本身就要能當 onboarding。

測試:TestAccountEngineLabel 兩則斷言翻轉成新規格(claude→仍顯示 Gemini),
並註明「若變回 · Claude 代表標籤又和萃取實際走的路脫鉤」;兩模組全綠。

⚠️ 未送達:daemon 執行檔要打包 v0.15.5 出貨才會到用戶手上。

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-08-04 11:33:13 +08:00

859 lines
33 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 頂層。
// t126Extractor/GeminiAPIKey/LLMModel 支援帳號層覆蓋——帳號有值時優先,空值繼承機器層。
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
// t126:每帳號獨立的萃取設定(空值繼承 DirectConfig 頂層)
Extractor string `json:"extractor,omitempty"`
GeminiAPIKey string `json:"gemini_api_key,omitempty"`
LLMModel string `json:"llm_model,omitempty"`
}
// 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(), // 正規化後的監看清單
}}
}
// t126 遷移:把頂層金鑰複製到每個沒有金鑰的帳號(複製非搬移,頂層保留當預設;冪等)。
// 帳號已有自己的值(非空)→ 不覆蓋,讓帳號層設定永遠優先。
for i := range c.Accounts {
if strings.TrimSpace(c.Accounts[i].Extractor) == "" && strings.TrimSpace(c.Extractor) != "" {
c.Accounts[i].Extractor = c.Extractor
}
if strings.TrimSpace(c.Accounts[i].GeminiAPIKey) == "" && strings.TrimSpace(c.GeminiAPIKey) != "" {
c.Accounts[i].GeminiAPIKey = c.GeminiAPIKey
}
if strings.TrimSpace(c.Accounts[i].LLMModel) == "" && strings.TrimSpace(c.LLMModel) != "" {
c.Accounts[i].LLMModel = c.LLMModel
}
}
// 驗證
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)
}
// t1492026-07-29 leo 實測揪出):**多帳號的看守資料夾在 Accounts[].WatchFolders**
// 先前只讀頂層 ⇒ 只在 accounts[] 設定的用戶,Folders() 回 nil ⇒ 掃描清單空 ⇒
// 按「立刻同步」毫無反應、log 顯示 "folders": null、results: []leo 就是這樣卡住的)。
// 頂層是舊制(單帳號);新制多帳號一律走 Accounts ⇒ 兩邊都要收。
for _, a := range c.Accounts {
for _, p := range a.WatchFolders {
add(expandHome(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 等。
// t126:帳號層 Extractor/GeminiAPIKey/LLMModel 有值時優先覆蓋機器層(空字串不算「有值」)。
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"
}
// t126:帳號層有值時優先(空字串繼承機器層,已由 sub := *c 複製)
if strings.TrimSpace(acc.Extractor) != "" {
sub.Extractor = acc.Extractor
}
if strings.TrimSpace(acc.GeminiAPIKey) != "" {
sub.GeminiAPIKey = acc.GeminiAPIKey
}
if strings.TrimSpace(acc.LLMModel) != "" {
sub.LLMModel = acc.LLMModel
}
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 := ""
// t176leo 08-03 拍板):**地端先只支援 Gemini**(「地端先限制 Gemini API Key 配合客戶要求」)。
// 殘留的 extractor:"claude"(雲端舊版下發、或舊 config 殘留)一律當 gemma 處理。
// 這不是「自動偵測有無 claude」(那是 leo 07-27 已否決的 B 案,見 daemon-beta/tasks.md:641),
// 而是「claude 路整條先不支援」——之後要裝回來,把這段拿掉即可。
if cfg.Extractor == "claude" {
cfg.Extractor = "gemma"
}
if cfg.Extractor == "gemma" {
if strings.TrimSpace(cfg.GeminiAPIKey) == "" {
extractorOK = false
// t178leo 08-04:封測者是台大資工碩士都卡住 ⇒「一般人就完蛋了」):
// 舊訊息「請在設定裡輸入 Gemini API Key」**沒說設定在哪** ⇒ 用戶找不到入口。
// 錯誤訊息本身就要能當 onboarding:指名選單項、指名去哪申請。
extractorError = "還沒設定 Gemini API Key ⇒ 點托盤選單的「AI 設定…」貼上金鑰" +
"(免費申請:aistudio.google.com/apikey"
}
}
// 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":
// t176claude 先不支援,等同 gemma(與 RunDirectOnce 的正規化保持一致,避免兩處漂移)。
fallthrough
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 "gemma":
cards, xerr = ExtractWithGemma(cfg.GeminiAPIKey, cfg.LLMModel, absRoot, ev.Path)
default:
// t176claude 路先不支援(RunDirectOnce 開頭已把 "claude" 正規化為 "gemma")。
// 走到這裡代表 config 有沒見過的值——誠實報錯,不要靜默跳過(禁假綠)。
xerr = fmt.Errorf("不支援的萃取方式 %q(目前只支援 Gemini", cfg.Extractor)
}
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
}
// t176t92 的 claude_bin 回寫已隨 claude 萃取路一併退役——RunDirectOnce 不再解析
// claude 執行檔,cfg.ClaudeBin 不會被改寫,故沒有東西需要回寫。
// (留著空轉的回寫邏輯會讓未來的人以為 claude 路還活著=錯誤的環境信號。)
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,跨平台。首輪立即跑。
// 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()
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
}