f8450815d3
每輪 GET /health 取 bundle_version(5s timeout 失敗靜默);空或日期<minCloudBuilt → 托盤「⚠ 知識庫需要更新(點我)」開 install.arcrun.dev;結果進 status.json。 兩模組 go test 全綠(總管親跑)。leo:「daemon 和雲端是連動的」——自此用戶只看托盤。 (實作=子 CC;驗證+commit=總管)
699 lines
27 KiB
Go
699 lines
27 KiB
Go
// 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 workflow(LLM 萃卡 → 機械切塊 → 寫 kbdb,全在 Arcrun workflow 裡完成)。
|
||
// 刪檔 → 把 removed 事件(collector-trigger.v1)POST 進實例的 rag_ingest workflow removed 分支
|
||
// (只按 page_name 讀 kbdb blocks 並標 deprecated,不碰 R2)。
|
||
//
|
||
// dogfooding(D29 daemon 薄殼豁免):本檔只做「監看/讀檔/算 hash/HTTP POST」——原生 Go。
|
||
// 萃取/切塊/RAG 一律在實例 workflow,daemon 內零 LLM/切塊邏輯。
|
||
//
|
||
// 用法:
|
||
//
|
||
// collector direct --config <config.json> [--once] [--dry-run]
|
||
//
|
||
// --once:掃一輪就退出(測試/cron 用);預設常駐輪詢(poll_interval_sec)。
|
||
// --dry-run:只列出會送出的動作,不 POST、不寫 manifest。
|
||
//
|
||
// 跨平台:純 stdlib、輪詢式偵測(不依賴 fsnotify)=零 CGo,darwin/arm64、windows/amd64 直接交叉編譯。
|
||
package main
|
||
|
||
import (
|
||
"bytes"
|
||
"crypto/sha256"
|
||
"encoding/hex"
|
||
"encoding/json"
|
||
"fmt"
|
||
"io"
|
||
"net/http"
|
||
"net/url"
|
||
"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(空=沿用 namespace,demo 慣例)
|
||
Email string `json:"email,omitempty"` // 實例主身分(t26;人人記得自己的 email,CF 全程隱形)
|
||
InstanceName string `json:"instance_name,omitempty"` // 暱稱(t26 選配;不取就顯示 email)
|
||
Library string `json:"library"` // 藏書地圖歸庫鍵(空=kb;per-folder 未指定時的後備)
|
||
// t52(leo 2026-07-25 裁決:「資料夾=庫」——企業有 10 個資料夾,財務/人事不准任何人看,
|
||
// 全塞一個庫=找死):每個看守資料夾對應自己的庫。key=資料夾絕對路徑,value=庫名。
|
||
// 未列出的資料夾=用資料夾名 slug 當庫名(librarySlug);再不行才退回 Library。
|
||
Libraries map[string]string `json:"libraries,omitempty"`
|
||
IngestWF string `json:"ingest_workflow"` // 直送萃取 workflow 名(空=rag_ingest_direct;extractor 模式不用)
|
||
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"` // 輪詢間隔秒(空/0=5)
|
||
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)。
|
||
//
|
||
// t89(leo 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 讀設定檔並補預設值 + 基本驗證。
|
||
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
|
||
}
|
||
|
||
// 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 {
|
||
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)。
|
||
// 額外:
|
||
// - 預檢 extractor 可用性(t92-②),有 fallback 時更新 cfg.ClaudeBin(in-memory,呼叫端存檔)。
|
||
// - 每輪結束寫 ~/.arcrun-rag/status.json(t91 狀態可見性)。
|
||
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"
|
||
}
|
||
}
|
||
|
||
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
|
||
}
|
||
}
|
||
|
||
// t91:每輪寫狀態檔(只有 extractor 模式才有意義的計數;direct 雲端萃模式 extracted_ok=0)。
|
||
if !dryRun && cfg.Manifest != "" {
|
||
// t103:每輪順手 GET /health 取雲端版本(5s timeout,失敗靜默)。
|
||
cloudVer, cloudOK := fetchCloudVersion(cfg.CypherURL)
|
||
st := SyncStatus{
|
||
LastSync: time.Now().Format(time.RFC3339),
|
||
ExtractorOK: extractorOK,
|
||
ExtractorError: extractorError,
|
||
CloudVersion: cloudVer,
|
||
CloudCheckOK: cloudOK,
|
||
}
|
||
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),
|
||
})
|
||
}
|
||
}
|
||
}
|
||
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 回寫到 configPath(t92:找到 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、下架 removed,2xx 後回寫該根 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
|
||
}
|
||
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"
|
||
// t15:extractor 模式雲端下架成功後,同步清掉本地萃出的卡
|
||
//(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
|
||
}
|