Files
arcrun-collector/manifest.go
T
Leo 47580f4ba5 t195:失敗重試加指數退避+上限——一個壞檔不再拖住整個資料夾
leo 2026-08-05:「有封測者在等,一直出錯」「已經等了幾天了」

## 病(leo 實撞,log 實證)
`小果被AFTEE詐貸.pdf` 雲端 401 失敗後,**manifest 完全不記失敗**
⇒ 下輪掃描又當「新檔」⇒ **1387 輪、跨 11 小時**,每輪 3.2~3.8 秒全在撞同一面牆;
且它排在佇列前面 ⇒ **整個資料夾的同步被一個壞檔拖住**。
leo:「原先萃檔案速度也快,現在花了十幾分才萃完」——**萃取沒變慢,慢的是重試**
(log 的 `FOREACH 所有 5 項目均失敗` 證明卡早就萃好了)。

**這是獨立於 401 的架構缺陷**:401 修好了,下次換別的錯照樣卡死。

## 修
· ManifestEntry 加 FailCount / LastFailAt / NextRetry
· MarkFailed():退避階梯 1m→5m→15m→1h→6h(之後維持 6h)
· ShouldRetry():退避窗口內跳過;連續失敗 8 次暫停自動重試
  force(使用者按「立刻同步」)忽略退避與上限——**人明確要求不該被機器擋住**
· MarkIngestedBy() 成功時清空失敗狀態(下次再壞從第一階重算)
· direct.go 三個失敗出口都記退避(讀檔失敗/萃取失敗/上傳失敗)
· retrySkipReason():跳過時說人話,不靜默(同 t195 燈號誠實原則)

## 🔴 真兇其實有兩層——第二層才是關鍵
只加退避欄位**沒有用**:`scan.go` 每輪都**重建** ManifestEntry,
原本只 carry IngestedHash/IngestedAt ⇒ 我寫進去的 fail_count 下一輪就被抹掉
⇒ 退避永遠停在「第 1 次失敗」=等同沒有退避。
(順帶發現 ExtractedBy(t73「誰萃的」)原本也一直悄悄丟失。)
⇒ scan.go carry 補齊四個欄位,並留註解:**日後新增跨輪欄位必須加在這裡**。

## 驗(真實跑,非只有單元測試)
單元測試 6 項全過(退避窗口/指數遞增 60/300/900/3600/21600/21600/
上限停止/force 忽略/成功清空/壞檔不連累同輪其他檔)
實跑(cypher_url 指向不存在主機製造必失敗):
  第 1 輪 failed
  第 2 輪 skipped「上次失敗(第 1 次),1m0s 後重試」
  第 3 輪 skipped(同上)
  第 4 輪 skipped(同上)
  模擬退避到期 → 確實重送、fail_count=2、退避升為 300s 
go test ./... 全綠;go vet 通過。

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-08-05 14:46:03 +08:00

208 lines
8.3 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.
// manifest.go — collector 本機 manifest 讀寫(SDD ingest-hash-trigger design §2)。
// manifest 只由 collector 寫入,雲端不得回寫(design §6-4)。
package main
import (
"crypto/rand"
"encoding/json"
"errors"
"fmt"
"os"
"path/filepath"
"time"
)
// ManifestEntry 是單一檔案在 manifest 裡的狀態。
// - mtime 只決定「要不要重算 hash」(fast-path),絕不作為內容變更判準。
// - ingested_hash 與 content_hash 分開存=可表達「已變更但尚未成功 ingest」;
// 上傳失敗不更新 ingested_hash,天然可重試。本階段(不接網路)永不寫 ingested_hash
// 留給之後的上傳/webhook task 在成功後回寫。
type ManifestEntry struct {
ContentHash string `json:"content_hash"`
Size int64 `json:"size"`
Mtime int64 `json:"mtime"`
IngestedHash string `json:"ingested_hash,omitempty"`
IngestedAt int64 `json:"ingested_at,omitempty"`
// ExtractedBy=這張卡是誰萃的("claude""gemma"""=無萃取器的直送路)。
//
// 為什麼要記(leo 2026-07-27 問「如果已經萃過了它知道嗎?」時發現的缺口):
// 原本 manifest 只記「萃過了」不記「誰萃的」。防重複本來就成立(hash 相同就跳過),
// 但**換萃取器時無從分辨哪些卡是舊萃取器產的**——claude 萃的卡和 gemma 萃的卡
// 品質不同卻混在同一個知識庫裡,想重萃也不知道該重萃哪些。
//
// 本欄位**不改變任何現有行為**(純記錄),但它是「換萃取器要不要重萃」這個決定的前提:
// 沒有它,之後想分辨就永遠分辨不了(舊資料補不回來)。
ExtractedBy string `json:"extracted_by,omitempty"`
// ── 失敗退避(t195,2026-08-05)────────────────────────────────────────────
// 病(leo 實撞):`小果被AFTEE詐貸.pdf` 因雲端 401 失敗後,**manifest 完全不記失敗**
// ⇒ 下一輪掃描又把它當「新檔」⇒ 無退避、無上限地重試。
// 實測:**1387 輪、跨 11 小時**,每輪 3.2~3.8 秒全花在撞同一面牆上;
// 而且它排在佇列前面 ⇒ **一個壞檔就把整個資料夾的同步拖住**。
// leo:「不只推上去慢,原先萃檔案速度也快,現在也花了十幾分才萃完」——
// **萃取根本沒變慢,慢的是重試**log 的 `FOREACH 所有 5 項目均失敗` 證明卡早就萃好了)。
//
// 解:記下失敗次數與下次可重試時間,指數退避+次數上限。
// 這是**獨立於當次錯誤**的架構缺陷——401 修好了,下次換別的錯照樣卡死,所以要修這裡。
FailCount int `json:"fail_count,omitempty"` // 連續失敗次數(成功即歸零)
LastFailAt int64 `json:"last_fail_at,omitempty"` // 最後一次失敗的 unix 秒
NextRetry int64 `json:"next_retry,omitempty"` // 早於這個時間不重試(0=可立即重試)
}
// retryBackoff 退避階梯:1m → 5m → 15m → 1h → 6h,之後每次 6h。
var retryBackoff = []time.Duration{
1 * time.Minute, 5 * time.Minute, 15 * time.Minute, 1 * time.Hour, 6 * time.Hour,
}
// MaxFailBeforeSkip 連續失敗達此次數 → 該檔暫停自動重試,不再拖住佇列。
// 使用者改檔(content hash 變)或按「立刻同步」時仍會重試(見 MarkFailedShouldRetry)。
const MaxFailBeforeSkip = 8
// Manifest 對應一個被勾選的資料夾。
type Manifest struct {
FolderID string `json:"folder_id"`
Root string `json:"root"`
Entries map[string]*ManifestEntry `json:"entries"`
}
// newUUID 產生 RFC 4122 v4 UUID(純 stdlib)。
func newUUID() (string, error) {
var b [16]byte
if _, err := rand.Read(b[:]); err != nil {
return "", err
}
b[6] = (b[6] & 0x0f) | 0x40 // version 4
b[8] = (b[8] & 0x3f) | 0x80 // variant 10
return fmt.Sprintf("%x-%x-%x-%x-%x", b[0:4], b[4:6], b[6:8], b[8:10], b[10:16]), nil
}
// LoadManifest 讀 manifest 檔;不存在=新資料夾,生成 folder_id。
func LoadManifest(path, root string) (*Manifest, error) {
data, err := os.ReadFile(path)
if err != nil {
if errors.Is(err, os.ErrNotExist) {
id, uerr := newUUID()
if uerr != nil {
return nil, uerr
}
return &Manifest{FolderID: id, Root: root, Entries: map[string]*ManifestEntry{}}, nil
}
return nil, err
}
var m Manifest
if err := json.Unmarshal(data, &m); err != nil {
return nil, fmt.Errorf("manifest %s 解析失敗(不覆寫、直接報錯): %w", path, err)
}
if m.Entries == nil {
m.Entries = map[string]*ManifestEntry{}
}
if m.FolderID == "" {
id, uerr := newUUID()
if uerr != nil {
return nil, uerr
}
m.FolderID = id
}
if root != "" {
m.Root = root
}
return &m, nil
}
// MarkIngested 是「整條 ingest 鏈成功」後的回寫鉤子(design §2):把該路徑的
// ingested_hash 接上當時送出的 source_hash。**R2 上傳成功不呼叫它**——上傳只是鏈的
// 第一環,要等 task 4named-webhook → ingest workflow)確認成功才回寫;在那之前
// 同檔每輪重發 added/modified=設計內重試,R2 端靠 HEAD no-op 天然冪等、零浪費。
// 回傳 false=路徑已不在 manifest(例如回報前檔案又被改名/刪除),呼叫端自行決定忽略或告警。
func (m *Manifest) MarkIngested(path, sourceHash string, at int64) bool {
return m.MarkIngestedBy(path, sourceHash, at, "")
}
// MarkIngestedBy 同 MarkIngested,另記「這輪是誰萃的」(extractor"claude""gemma""")。
//
// 為什麼另開一支而不是改 MarkIngested 的簽名:MarkIngested 有多個呼叫點
// direct 兩處+trigger.go 的 MarkIngestedEvents+測試),改簽名會擴散破壞。
// 走無萃取器路徑(sync/直送)的呼叫端不必知道這個欄位,維持原簽名最小侵入。
func (m *Manifest) MarkIngestedBy(path, sourceHash string, at int64, extractor string) bool {
e, ok := m.Entries[path]
if !ok {
return false
}
e.IngestedHash = sourceHash
e.IngestedAt = at
e.ExtractedBy = extractor
// 成功即清掉失敗狀態(t195):下次再壞會從第一階退避重新算起。
e.FailCount, e.LastFailAt, e.NextRetry = 0, 0, 0
return true
}
// MarkFailed 記錄一次失敗並排定下次可重試時間(t195 指數退避)。
//
// 為什麼要記在 manifest 而不是記憶體:collector 每輪是獨立 process
// `direct --once` 由看守器反覆拉起),記憶體狀態一輪就沒了——
// 這正是原本「1387 輪重試同一個檔」的原因:每輪都以為自己是第一次。
func (m *Manifest) MarkFailed(path string, at int64) bool {
e, ok := m.Entries[path]
if !ok {
return false
}
e.FailCount++
e.LastFailAt = at
idx := e.FailCount - 1
if idx >= len(retryBackoff) {
idx = len(retryBackoff) - 1
}
e.NextRetry = at + int64(retryBackoff[idx].Seconds())
return true
}
// ShouldRetry 回報「這個檔現在該不該送」。
//
// - 從沒失敗過 → true(行為與 t195 前一字不變)
// - 還在退避窗口內 → false(**這是止血點:壞檔不再每 4 秒撞一次牆拖住整個佇列**)
// - 連續失敗 ≥ MaxFailBeforeSkip → false(暫停自動重試)
//
// force=使用者按「立刻同步」:忽略退避與上限一律重送
// (人明確要求時不該被機器的退避擋住)。
// 另:使用者改檔會讓 content hash 變 → 走的是「內容變更」路徑,本函式不介入。
func (m *Manifest) ShouldRetry(path string, now int64, force bool) bool {
e, ok := m.Entries[path]
if !ok || e.FailCount == 0 {
return true
}
if force {
return true
}
if e.FailCount >= MaxFailBeforeSkip {
return false
}
return now >= e.NextRetry
}
// Save 原子寫入(temp + rename),避免掃描中斷留半個 JSON。
func (m *Manifest) Save(path string) error {
data, err := json.MarshalIndent(m, "", " ")
if err != nil {
return err
}
dir := filepath.Dir(path)
if err := os.MkdirAll(dir, 0o755); err != nil {
return err
}
tmp, err := os.CreateTemp(dir, ".manifest-*.tmp")
if err != nil {
return err
}
tmpName := tmp.Name()
if _, err := tmp.Write(append(data, '\n')); err != nil {
tmp.Close()
os.Remove(tmpName)
return err
}
if err := tmp.Close(); err != nil {
os.Remove(tmpName)
return err
}
return os.Rename(tmpName, path)
}