feat(ingest-hash-trigger): 凍結 collector/ingest schema+Go collector 骨架(issue #5, SDD task 1+2)
- schemas/collector-trigger.v1.schema.json:collector→named-webhook 觸發 payload(source_hash/r2_key/renamed/R6 warnings) - schemas/kbdb-ingest-request.v1.schema.json:source_hash 必填缺=400;儲存映射走 metadata_json.$.source_hash(守 kbdb 表不變鐵律);同 hash 重送=200 already_ingested - schemas/MIGRATION-webhook-payload.md:舊 Gitea push payload 逐欄對照(給 task 4 改寫 workflow) - collector/:Go module(純 stdlib)——manifest 讀寫+mtime fast-path 掃描+renamed hash 配對+40% 大量刪除防呆;CLI collector scan - go test ./... 全綠 7/7(added/modified/removed/renamed/防呆五情境+fast-path+manifest 往返) - Node 版 collector 標 legacy 並存(SDD task 4 拆除) Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1,289 @@
|
||||
// scan.go — 掃描迴圈與差異分類(SDD ingest-hash-trigger design §3)。
|
||||
// 事件順序:先把本輪 removed×added 以 content_hash 配對成 renamed(只更新路徑映射),
|
||||
// 再分類其餘 added/modified/removed;removed 數 > manifest 條目 × 門檻(預設 40%)
|
||||
// → removed 全部不執行、改發警告(R6)。本階段不接網路,事件輸出到 stdout。
|
||||
package main
|
||||
|
||||
import (
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"fmt"
|
||||
"io"
|
||||
"io/fs"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sort"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
// 先只認 .md 與常見文件檔(transform/其他格式是之後的 task)。
|
||||
var allowedExt = map[string]bool{
|
||||
".md": true,
|
||||
".markdown": true,
|
||||
".txt": true,
|
||||
".docx": true,
|
||||
".pptx": true,
|
||||
".pdf": true,
|
||||
}
|
||||
|
||||
// ---- 輸出 payload(對應 schemas/collector-trigger.v1.schema.json)----
|
||||
|
||||
type Event struct {
|
||||
Type string `json:"type"`
|
||||
Path string `json:"path"`
|
||||
OldPath string `json:"old_path,omitempty"`
|
||||
SourceHash string `json:"source_hash"`
|
||||
Size *int64 `json:"size,omitempty"`
|
||||
R2Key string `json:"r2_key,omitempty"`
|
||||
}
|
||||
|
||||
type Warning struct {
|
||||
Code string `json:"code"`
|
||||
Message string `json:"message"`
|
||||
RemovedCount int `json:"removed_count,omitempty"`
|
||||
ManifestCount int `json:"manifest_count,omitempty"`
|
||||
ThresholdRatio float64 `json:"threshold_ratio,omitempty"`
|
||||
}
|
||||
|
||||
type TriggerPayload struct {
|
||||
SchemaVersion int `json:"schema_version"`
|
||||
FolderID string `json:"folder_id"`
|
||||
Root string `json:"root,omitempty"`
|
||||
GeneratedAt int64 `json:"generated_at,omitempty"`
|
||||
Events []Event `json:"events"`
|
||||
Warnings []Warning `json:"warnings,omitempty"`
|
||||
}
|
||||
|
||||
type ScanOptions struct {
|
||||
// MaxRemovedRatio:單輪 removed 數 > manifest 條目數 × 本值 → 觸發大量刪除防呆(R6)。
|
||||
MaxRemovedRatio float64
|
||||
// SkipPaths:絕對路徑黑名單(如 manifest 檔自己住在 root 底下時)。
|
||||
SkipPaths map[string]bool
|
||||
}
|
||||
|
||||
const DefaultMaxRemovedRatio = 0.4
|
||||
|
||||
type fileState struct {
|
||||
hash string // sha256:<hex>
|
||||
size int64
|
||||
mtime int64
|
||||
}
|
||||
|
||||
func hashFile(path string) (string, error) {
|
||||
f, err := os.Open(path)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
defer f.Close()
|
||||
h := sha256.New()
|
||||
if _, err := io.Copy(h, f); err != nil {
|
||||
return "", err
|
||||
}
|
||||
return "sha256:" + hex.EncodeToString(h.Sum(nil)), nil
|
||||
}
|
||||
|
||||
func r2KeyOf(sourceHash string) string {
|
||||
return "raw/" + strings.TrimPrefix(sourceHash, "sha256:")
|
||||
}
|
||||
|
||||
// Scan 走訪 root、對照並更新 manifest、產出一輪事件。
|
||||
// manifest 更新原則:content_hash/size/mtime 反映現況;ingested_hash/ingested_at
|
||||
// 只搬運(renamed)與保留(modified),本函式永不設值——那是上傳成功後的事。
|
||||
func Scan(root string, m *Manifest, opts ScanOptions) (*TriggerPayload, error) {
|
||||
if opts.MaxRemovedRatio <= 0 {
|
||||
opts.MaxRemovedRatio = DefaultMaxRemovedRatio
|
||||
}
|
||||
orig := m.Entries
|
||||
manifestCountBefore := len(orig)
|
||||
|
||||
// 1) 走訪檔案系統,建立現況(mtime+size fast-path:沒變→沿用 manifest hash,變了才算 sha256)。
|
||||
current := map[string]fileState{}
|
||||
err := filepath.WalkDir(root, func(p string, d fs.DirEntry, werr error) error {
|
||||
if werr != nil {
|
||||
return werr
|
||||
}
|
||||
name := d.Name()
|
||||
if d.IsDir() {
|
||||
if p != root && strings.HasPrefix(name, ".") {
|
||||
return filepath.SkipDir // 隱藏目錄(.git、.obsidian…)整棵跳過
|
||||
}
|
||||
return nil
|
||||
}
|
||||
if strings.HasPrefix(name, ".") {
|
||||
return nil
|
||||
}
|
||||
if abs, aerr := filepath.Abs(p); aerr == nil && opts.SkipPaths[abs] {
|
||||
return nil
|
||||
}
|
||||
if !allowedExt[strings.ToLower(filepath.Ext(name))] {
|
||||
return nil
|
||||
}
|
||||
info, ierr := d.Info()
|
||||
if ierr != nil {
|
||||
return ierr
|
||||
}
|
||||
rel, rerr := filepath.Rel(root, p)
|
||||
if rerr != nil {
|
||||
return rerr
|
||||
}
|
||||
rel = filepath.ToSlash(rel)
|
||||
st := fileState{size: info.Size(), mtime: info.ModTime().Unix()}
|
||||
if e, ok := orig[rel]; ok && e.ContentHash != "" && e.Mtime == st.mtime && e.Size == st.size {
|
||||
st.hash = e.ContentHash // fast-path:mtime+size 沒變,跳過重算
|
||||
} else {
|
||||
h, herr := hashFile(p)
|
||||
if herr != nil {
|
||||
return herr
|
||||
}
|
||||
st.hash = h
|
||||
}
|
||||
current[rel] = st
|
||||
return nil
|
||||
})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// 2) 初分:added 候選(現況有、manifest 無)與 removed 候選(manifest 有、現況無)。
|
||||
var addedPaths, removedPaths []string
|
||||
for p := range current {
|
||||
if _, ok := orig[p]; !ok {
|
||||
addedPaths = append(addedPaths, p)
|
||||
}
|
||||
}
|
||||
for p := range orig {
|
||||
if _, ok := current[p]; !ok {
|
||||
removedPaths = append(removedPaths, p)
|
||||
}
|
||||
}
|
||||
sort.Strings(addedPaths)
|
||||
sort.Strings(removedPaths)
|
||||
|
||||
// 3) 先配對 renamed(design §3 順序 1):removed×added 以 content_hash 配對,
|
||||
// 配上=只更新路徑映射,不 retire、不重萃、不重傳。同 hash 多候選→排序後貪婪配對(確定性)。
|
||||
removedByHash := map[string][]string{}
|
||||
for _, p := range removedPaths {
|
||||
h := orig[p].ContentHash
|
||||
removedByHash[h] = append(removedByHash[h], p)
|
||||
}
|
||||
renamedOldOf := map[string]string{} // newPath -> oldPath
|
||||
pairedOld := map[string]bool{}
|
||||
var events []Event
|
||||
for _, np := range addedPaths {
|
||||
h := current[np].hash
|
||||
cands := removedByHash[h]
|
||||
if len(cands) == 0 {
|
||||
continue
|
||||
}
|
||||
op := cands[0]
|
||||
removedByHash[h] = cands[1:]
|
||||
pairedOld[op] = true
|
||||
renamedOldOf[np] = op
|
||||
events = append(events, Event{Type: "renamed", Path: np, OldPath: op, SourceHash: h})
|
||||
}
|
||||
|
||||
// 4) added:真新檔+「曾偵測但從未成功 ingest」的檔(重試語意,design §2)。
|
||||
sortedCurrent := make([]string, 0, len(current))
|
||||
for p := range current {
|
||||
sortedCurrent = append(sortedCurrent, p)
|
||||
}
|
||||
sort.Strings(sortedCurrent)
|
||||
addedEvent := func(p string) Event {
|
||||
st := current[p]
|
||||
size := st.size
|
||||
return Event{Type: "added", Path: p, SourceHash: st.hash, Size: &size, R2Key: r2KeyOf(st.hash)}
|
||||
}
|
||||
for _, p := range sortedCurrent {
|
||||
if op, isRenamed := renamedOldOf[p]; isRenamed {
|
||||
if orig[op].IngestedHash == "" { // 改名的檔其實從未 ingest 成功 → 補一發 added
|
||||
events = append(events, addedEvent(p))
|
||||
}
|
||||
continue
|
||||
}
|
||||
if _, existed := orig[p]; !existed {
|
||||
events = append(events, addedEvent(p)) // 真新檔
|
||||
} else if orig[p].IngestedHash == "" {
|
||||
events = append(events, addedEvent(p)) // 上輪偵測過但 ingest 未成功 → 重試
|
||||
}
|
||||
}
|
||||
|
||||
// 5) modified:manifest 有、現況有、content_hash != ingested_hash(design §3 順序 3)。
|
||||
for _, p := range sortedCurrent {
|
||||
e, existed := orig[p]
|
||||
if !existed {
|
||||
continue
|
||||
}
|
||||
if _, isRenamed := renamedOldOf[p]; isRenamed {
|
||||
continue
|
||||
}
|
||||
if e.IngestedHash != "" && current[p].hash != e.IngestedHash {
|
||||
st := current[p]
|
||||
size := st.size
|
||||
events = append(events, Event{Type: "modified", Path: p, SourceHash: st.hash, Size: &size, R2Key: r2KeyOf(st.hash)})
|
||||
}
|
||||
}
|
||||
|
||||
// 6) removed(扣掉已配對走的)+大量刪除防呆(R6)。
|
||||
var finalRemoved []string
|
||||
for _, p := range removedPaths {
|
||||
if !pairedOld[p] {
|
||||
finalRemoved = append(finalRemoved, p)
|
||||
}
|
||||
}
|
||||
var warnings []Warning
|
||||
guardTripped := manifestCountBefore > 0 &&
|
||||
float64(len(finalRemoved)) > opts.MaxRemovedRatio*float64(manifestCountBefore)
|
||||
if guardTripped {
|
||||
warnings = append(warnings, Warning{
|
||||
Code: "mass_delete_guard",
|
||||
Message: fmt.Sprintf(
|
||||
"本輪偵測到 %d/%d 個檔案消失(超過 %.0f%% 門檻)——可能是資料夾未掛載或同步半途。本輪全部「不」下架,請確認資料夾完好後再放行。",
|
||||
len(finalRemoved), manifestCountBefore, opts.MaxRemovedRatio*100),
|
||||
RemovedCount: len(finalRemoved),
|
||||
ManifestCount: manifestCountBefore,
|
||||
ThresholdRatio: opts.MaxRemovedRatio,
|
||||
})
|
||||
} else {
|
||||
for _, p := range finalRemoved {
|
||||
events = append(events, Event{Type: "removed", Path: p, SourceHash: orig[p].ContentHash})
|
||||
}
|
||||
}
|
||||
|
||||
// 7) 更新 manifest(rebuild):現況檔全數收錄;ingested_* 由舊 entry(或 renamed 的舊路徑)搬運。
|
||||
// 防呆觸發時 removed 條目保留(下輪重評、警告會再響,直到人確認或檔案回來)。
|
||||
newEntries := make(map[string]*ManifestEntry, len(current))
|
||||
for p, st := range current {
|
||||
ne := &ManifestEntry{ContentHash: st.hash, Size: st.size, Mtime: st.mtime}
|
||||
var carry *ManifestEntry
|
||||
if op, isRenamed := renamedOldOf[p]; isRenamed {
|
||||
carry = orig[op]
|
||||
} else if e, ok := orig[p]; ok {
|
||||
carry = e
|
||||
}
|
||||
if carry != nil {
|
||||
ne.IngestedHash = carry.IngestedHash
|
||||
ne.IngestedAt = carry.IngestedAt
|
||||
}
|
||||
newEntries[p] = ne
|
||||
}
|
||||
if guardTripped {
|
||||
for _, p := range finalRemoved {
|
||||
newEntries[p] = orig[p]
|
||||
}
|
||||
}
|
||||
m.Entries = newEntries
|
||||
m.Root = root
|
||||
|
||||
if events == nil {
|
||||
events = []Event{}
|
||||
}
|
||||
return &TriggerPayload{
|
||||
SchemaVersion: 1,
|
||||
FolderID: m.FolderID,
|
||||
Root: root,
|
||||
GeneratedAt: time.Now().Unix(),
|
||||
Events: events,
|
||||
Warnings: warnings,
|
||||
}, nil
|
||||
}
|
||||
Reference in New Issue
Block a user