From a6ced32d4559bb3d3aaf24d5b14826e672383f19 Mon Sep 17 00:00:00 2001 From: richblack Date: Fri, 7 Aug 2026 16:58:07 +0800 Subject: [PATCH] =?UTF-8?q?collector=EF=BC=9A=E7=A9=8D=E5=A3=93=E5=88=86?= =?UTF-8?q?=E6=89=B9=EF=BC=8B=E6=96=B0=E6=AA=94=E5=84=AA=E5=85=88=EF=BC=8F?= =?UTF-8?q?=E9=A1=8D=E5=BA=A6=E7=94=A8=E5=AE=8C=E8=AC=9B=E4=BA=BA=E8=A9=B1?= =?UTF-8?q?=E9=99=8D=E9=80=9F=EF=BC=8F=E6=96=B7=E9=BB=9E=E7=BA=8C=E5=82=B3?= =?UTF-8?q?=EF=BC=8F=E5=90=8C=E5=85=A7=E5=AE=B9=E5=A4=9A=E6=A0=BC=E5=BC=8F?= =?UTF-8?q?=E5=8E=BB=E9=87=8D?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 封測事故(Evan):daemon 逐檔萃取上傳、每個檔在雲端產生一次工作流執行, 690 個檔在自己的免費 CF 帳號上短時間內衝出 1,070 次寫入,撞上免費上限 1,000, 額度爆掉、只有 8 個檔成功。雲端那一半(紀錄改走資料層 API)已修好, daemon 這一半原本完全沒有節奏——本次補齊四件事: 1. 上傳節奏(direct_pacing.go):單輪最多處理 MaxEventsPerRun 個事件 (預設 25)、每次觸發雲端前節流 700ms;一輪掃到的事件依檔案 mtime 由新到舊排序,今天寫的永遠優先,積壓慢慢消化不擋日常使用。 2. 額度用完講人話(quota.go):偵測到 Workers AI「10,000 neurons」/ 「4006」等已知上游訊號後,換成三句話(今天已整理幾份/可換模型或 升級 Cloudflare/不花錢也沒關係、今天或明天早上 8:00 會自動恢復), 不出現裸露的錯誤碼;同帳號同輪與下一輪都不再繼續撞牆 (quotaState 全域冷卻,跨資料夾/跨程序重啟持續,直到台灣時間 早上 8:00 額度重置)。 3. 斷點續傳:每個事件處理完立刻寫回 manifest(不再等整輪跑完才存一次), process 被殺掉重開只會接著做真正還沒完成的部分;removed 事件另外 用 preScanEntries 快照保護,下架失敗時不會被其他事件的存檔動作 誤標成「已完成」而永遠不再重試。 4. 同內容多格式去重(scan.go):同一批來源轉出的多種格式(如 leo 給的 資料集 27,164 檔=9,045 md+9,044 json+9,043 html,md/json 同檔名 主幹)依檔名主幹分組,只留優先序最高的一份進事件管線,其餘標記在 DuplicateFormats(不吃三倍額度)。 Co-Authored-By: Claude Opus 5 --- direct.go | 163 +++++++++++++++++- direct_pacing.go | 109 ++++++++++++ direct_pacing_test.go | 384 ++++++++++++++++++++++++++++++++++++++++++ direct_quota_test.go | 154 +++++++++++++++++ quota.go | 120 +++++++++++++ quota_test.go | 110 ++++++++++++ scan.go | 107 ++++++++++++ scan_dedup_test.go | 141 ++++++++++++++++ sync_status.go | 10 ++ 9 files changed, 1291 insertions(+), 7 deletions(-) create mode 100644 direct_pacing.go create mode 100644 direct_pacing_test.go create mode 100644 direct_quota_test.go create mode 100644 quota.go create mode 100644 quota_test.go create mode 100644 scan_dedup_test.go diff --git a/direct.go b/direct.go index a4b9540..b76789e 100644 --- a/direct.go +++ b/direct.go @@ -83,6 +83,10 @@ type DirectConfig struct { 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) + // MaxEventsPerRun(2026-08-07 pacing task):單輪最多處理幾個 added/modified/renamed/ + // removed 事件,空/0=DefaultMaxEventsPerRun。存在理由:巨量積壓(實據 27,164 檔) + // 不該一輪湧完——搭配節流間隔+新檔優先排序,讓積壓慢慢消化,不擋今天剛寫的新檔。 + MaxEventsPerRun int `json:"max_events_per_run,omitempty"` // ForceSync=這一輪是使用者按「立刻同步」觸發的(t195)。 // 為真時忽略失敗退避與次數上限,一律重送——**人明確要求時不該被機器的退避擋住**。 @@ -482,6 +486,15 @@ func RunDirectOnce(cfg *DirectConfig, dryRun bool) ([]DirectResult, int, *Trigge results := []DirectResult{} exit := 0 var lastPayload *TriggerPayload + now := time.Now() // 2026-08-07 pacing task:整輪共用同一個時間點(排序/冷卻判斷一致、好測試) + + // 2026-08-07:提早載入上一輪 status(原本只在函式尾端載入做 CarryForwardActivity)。 + // 額度冷卻與「今天已萃幾份」是**跨輪持續的狀態**(quotaState 見 quota.go), + // 要在處理帳號之前就知道上一輪冷卻到什麼時候、今天已經算到幾份。 + var prevStatus SyncStatus + if cfg.Manifest != "" { + prevStatus, _ = LoadSyncStatus(StatusFilePath(cfg.Manifest)) // 讀不到=零值,等同「沒有上一輪」 + } // t92-②:預檢 extractor(機器層級),有 fallback 路徑時就地更新 cfg.ClaudeBin。 extractorOK := true @@ -597,9 +610,24 @@ func RunDirectOnce(cfg *DirectConfig, dryRun bool) ([]DirectResult, int, *Trigge } } + // 2026-08-07:從上一輪 status 復原這個帳號的額度冷卻/今天已萃份數—— + // 這兩件事是**帳號層級、跨輪持續**的狀態,不是單一資料夾的(額度是雲端 + // 實例/Cloudflare 帳號共用的,一個根撞到,同帳號其他根不該還繼續撞牆)。 + qs := "aState{} + if prevAcc, ok := prevStatus.AccountDetails[accHost]; ok { + if prevAcc.DailyIngestedDate == todayUTC(now) { + qs.DailyCount = prevAcc.DailyIngestedCount + } + if prevAcc.QuotaCooldownUntil != "" { + if t, perr := time.Parse(time.RFC3339, prevAcc.QuotaCooldownUntil); perr == nil { + qs.CooldownUntil = t + } + } + } + multi := len(accCfg.Folders()) > 1 for _, root := range accCfg.Folders() { - r, e, p := runDirectOnceRoot(accCfg, root, dryRun) + r, e, p := runDirectOnceRoot(accCfg, root, dryRun, qs, now) if multi { for i := range r { r[i].Root = root @@ -633,6 +661,23 @@ func RunDirectOnce(cfg *DirectConfig, dryRun bool) ([]DirectResult, int, *Trigge skippedOtherNames = append(skippedOtherNames, p.SkippedOtherNames...) } } + + // 2026-08-07:把這輪(可能剛更新過的)額度冷卻/今天已萃份數寫回, + // 供下一輪 RunDirectOnce(甚至下一次程序啟動——status.json 落地磁碟)復原。 + accSt.DailyIngestedDate = todayUTC(now) + accSt.DailyIngestedCount = qs.DailyCount + if qs.inCooldown(now) { + accSt.QuotaCooldownUntil = qs.CooldownUntil.Format(time.RFC3339) + notice := qs.noticeNow(now) + accSt.QuotaMessage = ¬ice + // 冷卻中一定代表萃取沒就緒——浮到頂層讓托盤「狀態:」直接看得到, + // 不必展開帳號才發現「為什麼今天都沒有動靜」。 + extractorOK = false + extractorError = notice.Combined() + } + // 冷卻已過且本輪沒有新命中 ⇒ QuotaCooldownUntil/QuotaMessage 維持零值, + // 自然清除舊訊息(accSt 每輪重建,不會殘留上一輪的冷卻通知)。 + accountDetails[accHost] = accSt } @@ -690,7 +735,10 @@ func RunDirectOnce(cfg *DirectConfig, dryRun bool) ([]DirectResult, int, *Trigge // 於是 status.failures 是空的 ⇒ 畫面只剩「⚠ N 份失敗」沒有原因。 // 退避訊息裡已經帶了上游真因(retrySkipReason),一併收進來, // 使用者才看得到「為什麼」而不只是「幾份」。 - if strings.Contains(r.Error, "後重試") || strings.Contains(r.Error, "已暫停自動重試") { + // 2026-08-07:額度冷卻的 skip 訊息(quotaState.noticeNow().Combined()) + // 用「會自動恢復」當識別字——同樣要讓使用者看得到原因,不是只看到「幾份」。 + if strings.Contains(r.Error, "後重試") || strings.Contains(r.Error, "已暫停自動重試") || + strings.Contains(r.Error, "會自動恢復") { st.Failures = append(st.Failures, ExtractFail{ Path: r.Path, Error: shortError(r.Error), @@ -699,6 +747,12 @@ func RunDirectOnce(cfg *DirectConfig, dryRun bool) ([]DirectResult, int, *Trigge } } } + // 2026-08-07:巨量積壓場景(實據 27,164 檔)下 Failures 可能暴增到數千筆, + // status.json 不該被撐成幾 MB 的清單——裁到跟 SkippedDocs 一樣的上限, + // 總數仍在 ExtractFailed,UI 可以說「…等 N 個」(同 MaxSkippedListed 的做法)。 + if len(st.Failures) > MaxSkippedListed { + st.Failures = st.Failures[:MaxSkippedListed] + } // 單帳號時把 cloud version 也填頂層(向後相容) if len(accountDetails) == 1 { for _, v := range accountDetails { @@ -711,8 +765,7 @@ func RunDirectOnce(cfg *DirectConfig, dryRun bool) ([]DirectResult, int, *Trigge // 本輪有產出就記成「最近一次有做事」;本輪沒事做則把上一輪的原樣帶下來, // 不要讓「剛整理完 N 份」這個唯一的完成證據在 15 秒後被歸零抹掉。 statusPath := StatusFilePath(cfg.Manifest) - prev, _ := LoadSyncStatus(statusPath) // 讀不到=零值,等同「沒有上一輪」 - CarryForwardActivity(prev, &st) + CarryForwardActivity(prevStatus, &st) // 2026-08-07:沿用函式開頭已載入的 prevStatus,不重讀一次 if serr := SaveSyncStatus(statusPath, st); serr != nil { fmt.Fprintf(os.Stderr, "status 寫入失敗(不擋看守):%v\n", serr) } @@ -810,7 +863,9 @@ func accountsConnected(cfg *DirectConfig) bool { return false } -func runDirectOnceRoot(cfg *DirectConfig, root string, dryRun bool) ([]DirectResult, int, *TriggerPayload) { +// qs:這個帳號本輪共用的額度冷卻狀態(跨同帳號的多個監看根,見 quota.go)。 +// runNow:整輪 RunDirectOnce 共用的時間點(排序/冷卻判斷一致、好測試)。 +func runDirectOnceRoot(cfg *DirectConfig, root string, dryRun bool, qs *quotaState, runNow time.Time) ([]DirectResult, int, *TriggerPayload) { results := []DirectResult{} exit := 0 @@ -827,6 +882,16 @@ func runDirectOnceRoot(cfg *DirectConfig, root string, dryRun bool) ([]DirectRes if err != nil { return append(results, DirectResult{Status: "failed", Error: err.Error()}), 1, nil } + // 2026-08-07 task 3:Scan() 會把 removed 的路徑從 m.Entries 整批拿掉(rebuild 語意, + // 見 scan.go 步驟 7)——但那只是「偵測到不見了」,不代表下架 POST 已經成功。 + // 沒有這份快照的話,本輪只要有任何一個 added/modified 事件先觸發了下面的 + // incremental saveManifest(),就會把「還沒確認下架成功」的路徑一併存進磁碟, + // 下一輪 Scan() 兩邊都找不到它 ⇒ 永遠不會再補發 removed 事件、下架永遠不會重試。 + // 先存一份,Scan() 後把這些路徑「暫時放回去」,直到對應的 removed 事件真的成功。 + preScanEntries := make(map[string]*ManifestEntry, len(m.Entries)) + for k, v := range m.Entries { + preScanEntries[k] = v + } payload, err := Scan(absRoot, m, ScanOptions{ MaxRemovedRatio: cfg.MaxRemoved, SkipPaths: map[string]bool{ @@ -841,13 +906,65 @@ func runDirectOnceRoot(cfg *DirectConfig, root string, dryRun bool) ([]DirectRes return append(results, DirectResult{Status: "failed", Error: err.Error()}), 1, nil } - now := time.Now().Unix() + now := runNow.Unix() + + // 2026-08-07 task 3(斷點續傳):每個事件處理完就立刻存檔,不要等整輪跑完。 + // 舊行為=整個 for 迴圈跑完才 Save 一次——process 在跑到一半被殺掉(重開機、 + // 換版、當機)時,**已經成功的那些也會遺失**,下次重開等於從頭來過, + // 且已經花掉的額度/請求全部白費(正是 leo 要求「不從頭來」要防的事)。 + // 改成每個事件收工就存一次:kill 在任何一刻,磁碟上的 manifest 都反映 + // 「這一刻之前已確定成功的事」,下一輪只會處理真正還沒做完的。 + saveManifest := func() { + if dryRun { + return + } + if serr := m.Save(absManifest); serr != nil { + results = append(results, DirectResult{Status: "failed", Error: "manifest 存檔失敗(斷點續傳可能失效):" + serr.Error()}) + exit = 1 + } + } + + // 承上:把本輪偵測到的 removed 路徑暫時放回 m.Entries,直到迴圈裡真的處理到它、 + // POST 成功才由「removed」分支明確刪除。失敗或本輪還沒輪到(單輪上限)都維持放回的狀態, + // 下一輪自然重新偵測、重新嘗試下架——不會因為別的事件先存檔而被誤永久跳過。 for _, ev := range payload.Events { + if ev.Type == "removed" { + if e, ok := preScanEntries[ev.Path]; ok { + m.Entries[ev.Path] = e + } + } + } + + // 2026-08-07 pacing task 1:新改的檔優先+單輪上限。 + // - 排序:把 payload.Events 依檔案 mtime 由新到舊重排(不動送雲端的 payload 本身, + // 只重排這裡的本機處理佇列——見 sortEventsNewestFirst 註解)。 + // - 上限:巨量積壓(實據 27,164 檔)不該一輪湧完;超過上限的事件本輪不碰, + // manifest 未標 ingested ⇒ 下一輪 Scan() 自然重新出現(且若使用者這期間 + // 寫了新檔,新檔的 mtime 更新,下一輪排序會插到最前面,不會被積壓卡住)。 + orderedEvents := sortEventsNewestFirst(absRoot, payload.Events) + perRunCap := cfg.effectiveMaxEventsPerRun() + deferredCount := 0 + if len(orderedEvents) > perRunCap { + deferredCount = len(orderedEvents) - perRunCap + orderedEvents = orderedEvents[:perRunCap] + } + + for _, ev := range orderedEvents { switch ev.Type { case "added", "modified", "renamed": // renamed 在 direct 模式視同 added:內容未變但為求 kbdb 有這頁名的卡,重送一次萃取 //(頁名可能改變=要新頁名的卡)。冪等由 kbdb 端承擔(同頁名覆蓋語意)。 res := DirectResult{Type: ev.Type, Path: ev.Path} + // 2026-08-07 pacing task 2:帳號還在額度冷卻中 → 這輪連試都不試。 + // 這不是這個檔的問題(不記 FailCount/退避——那是「這個檔」的病歷, + // 額度用完是「整個帳號」的狀態,混在一起會讓退避階梯失真)。 + // 放在 ShouldRetry 之前:冷卻是更高層級的條件,沒必要先算退避訊息又蓋掉。 + if qs.inCooldown(runNow) { + res.Status = "skipped" + res.Error = qs.noticeNow(runNow).Combined() + results = append(results, res) + continue + } // 🔴 t195 止血點:這個檔剛失敗過且還在退避窗口內 → 這輪跳過。 // 沒有這道閘時的實測災情:`小果被AFTEE詐貸.pdf` 因雲端 401 失敗, // 每輪重掃又被當成新檔 ⇒ **1387 輪、跨 11 小時**,且它排在佇列前面, @@ -866,6 +983,7 @@ func runDirectOnceRoot(cfg *DirectConfig, root string, dryRun bool) ([]DirectRes if rerr != nil { res.Status, res.Error = "failed", "讀檔失敗:"+rerr.Error() m.MarkFailed(ev.Path, now, res.Error) // t195:讀不到的檔也退避(權限/被鎖/壞掉的外接碟) + saveManifest() results = append(results, res) exit = 1 continue @@ -875,6 +993,7 @@ func runDirectOnceRoot(cfg *DirectConfig, root string, dryRun bool) ([]DirectRes results = append(results, res) continue } + pace() // 2026-08-07:每次要觸發雲端(萃取/POST)之前先節流一下 if cfg.Extractor != "" { // 四步定稿:本地萃卡 → 每張卡 POST rag_ingest_card(原文不出機) var cards []string @@ -898,10 +1017,21 @@ func runDirectOnceRoot(cfg *DirectConfig, root string, dryRun bool) ([]DirectRes xerr = fmt.Errorf("不支援的萃取方式 %q(支援:workers-ai/gemma)", cfg.Extractor) } if xerr != nil { - res.Status, res.Error = "failed", "本地萃取失敗:"+xerr.Error() + // 2026-08-07 task 2:Workers AI 每日免費額度用完是**已知的上游狀況** + // (wiki mistakes.md 2026-08-06),不是 bug——不能讓使用者看到裸露的 + // 「4006」「HTTP 502」,要換成三句話(成就/出口/保證),且不能再 + // 每輪繼續撞同一面牆(qs.markHit 設定帳號層級的冷卻,下一個事件、 + // 下一輪都會被上面的 qs.inCooldown 擋下,不再嘗試萃取)。 + if isQuotaExhausted(xerr.Error()) { + qs.markHit(runNow, xerr.Error()) + res.Status, res.Error = "failed", qs.noticeNow(runNow).Combined() + } else { + res.Status, res.Error = "failed", "本地萃取失敗:"+xerr.Error() + } // t195:萃取階段失敗同樣要記退避。**這條路徑比上傳更早**, // 漏記的話(連不上知識庫、金鑰壞、模型錯)照樣每輪重撞。 m.MarkFailed(ev.Path, now, res.Error) + saveManifest() results = append(results, res) exit = 1 continue @@ -934,12 +1064,14 @@ func runDirectOnceRoot(cfg *DirectConfig, root string, dryRun bool) ([]DirectRes res.Status = "ingested" // 記下是誰萃的(t73/leo 07-27):換萃取器時才分辨得出哪些卡是舊的。 m.MarkIngestedBy(ev.Path, ev.SourceHash, now, cfg.Extractor) + qs.DailyCount++ // 2026-08-07:今天的成就數(額度訊息「今天已經幫你整理了 N 份」用) } else { // t195:記下失敗並排定退避,否則下輪又把它當新檔重試 //(實撞:1387 輪 × 11 小時全在撞同一面 401 的牆,還拖住整個佇列)。 m.MarkFailed(ev.Path, now, res.Error) exit = 1 } + saveManifest() results = append(results, res) continue } @@ -965,7 +1097,9 @@ func runDirectOnceRoot(cfg *DirectConfig, root string, dryRun bool) ([]DirectRes } else { res.Status = "ingested" m.MarkIngested(ev.Path, ev.SourceHash, now) // 2xx 才回寫(下輪不重送) + qs.DailyCount++ } + saveManifest() results = append(results, res) case "removed": @@ -975,6 +1109,7 @@ func runDirectOnceRoot(cfg *DirectConfig, root string, dryRun bool) ([]DirectRes results = append(results, res) continue } + pace() // 2026-08-07:下架一樣是觸發雲端 workflow,同樣節流 // 下架=POST {page_name, path} 進 rag_takedown_direct(按 page_name 讀 kbdb blocks // 標 deprecated,不碰 R2;獨立於 rag_ingest 的 __CARDS_PREFIX__ 閘——direct 模式檔在 // 資料夾根,會被 rag_ingest 的前綴閘擋掉,故自帶不含前綴閘的下架 workflow)。 @@ -986,8 +1121,14 @@ func runDirectOnceRoot(cfg *DirectConfig, root string, dryRun bool) ([]DirectRes if perr != nil { res.Status, res.Error = "failed", perr.Error() exit = 1 + // 2026-08-07:下架失敗——保持上面「暫時放回」的狀態,不刪、不存檔。 + // 下一輪 Scan() 會重新偵測到這個檔仍然不見了,自然重新補發 removed 事件。 } else { res.Status = "removed" + // 2026-08-07 task 3:下架真的成功了,這時才正式從 manifest 拿掉並存檔 + // (不是 Scan() rebuild 時就拿掉——那時只是「偵測到不見了」,不是「已下架」)。 + delete(m.Entries, ev.Path) + saveManifest() // t15:extractor 模式雲端下架成功後,同步清掉本地萃出的卡 //(system-dev/wiki/cards/<頁名>.md),保持本地 wiki 與雲端一致。 // 存在才刪;刪失敗只記 warning 不擋(下架本體已成功)。 @@ -1012,7 +1153,15 @@ func runDirectOnceRoot(cfg *DirectConfig, root string, dryRun bool) ([]DirectRes results = append(results, DirectResult{Type: "warning", Status: "skipped", Error: w.Code + ": " + w.Message}) } + // 2026-08-07 pacing task 1:本輪因單輪上限被延後的事件——不是失敗,是刻意排隊, + // 不佔用 exit(不該讓「積壓還很多」看起來像出錯了)。 + if deferredCount > 0 { + results = append(results, DirectResult{Type: "info", Status: "skipped", Error: resumeAfterCapMessage(deferredCount)}) + } + if !dryRun { + // 2026-08-07:每個事件處理完都已個別呼叫 m.Save(見下方迴圈內),這裡是最後的 + // 安全網——涵蓋防呆警告輪等不對應單一事件的狀態變化,多存一次無害(原子寫入)。 if err := m.Save(absManifest); err != nil { results = append(results, DirectResult{Status: "failed", Error: "manifest 存檔失敗:" + err.Error()}) exit = 1 diff --git a/direct_pacing.go b/direct_pacing.go new file mode 100644 index 0000000..2e0ece3 --- /dev/null +++ b/direct_pacing.go @@ -0,0 +1,109 @@ +// direct_pacing.go — 上傳節奏(2026-08-07):積壓不要一次湧上去,新改的檔優先。 +// +// 背景(封測事故):Evan 在自己的免費 Cloudflare 帳號上跑,監看資料夾裡約 690 個檔, +// daemon 逐檔萃取上傳、每個檔在雲端產生一次工作流執行 ⇒ 1,070 次寫入撞上免費上限 +// 1,000,額度爆掉,只有 8 個檔成功。雲端那一半(紀錄改走資料層 API)已修好, +// 但只要積壓一次全湧上去,額度一恢復就會立刻再爆一次——daemon 這一半必須自己有節奏。 +// +// 兩個機制,缺一不可: +// 1. 排序:一輪掃到的多個 added/modified/renamed 事件,依檔案 mtime 由新到舊處理 +// (leo 2026-08-07 定的排序判準:「日期最近的最優先,今天寫的一定最近」—— +// 一條佇列就夠,不必維護「新檔」「積壓」兩條)。 +// 2. 節流:每次觸發雲端(萃取或 POST ingest/removed)之間留最小間隔, +// 且單輪最多處理有限筆——**存量慢慢消化不影響日常使用**,但下一輪 Scan() +// 會立刻看到使用者剛寫的新檔並插到最前面,不會被積壓卡住幾小時才輪到。 +package collector + +import ( + "fmt" + "os" + "path/filepath" + "sort" + "time" +) + +// directPaceInterval:兩次觸發雲端 workflow 之間的最小間隔。 +// 正式執行維持節流;測試會在 init 時歸零(見 direct_pacing_test.go), +// 需要驗證節流本身的測試再自行 save/restore 成一個小的非零值。 +var directPaceInterval = 700 * time.Millisecond + +// pace 節流:在每次要觸發雲端(萃取 API/postJSON)之前呼叫。 +// directPaceInterval<=0 時(測試環境預設)為 no-op。 +func pace() { + if directPaceInterval > 0 { + time.Sleep(directPaceInterval) + } +} + +// DefaultMaxEventsPerRun:單輪最多處理幾個 added/modified/renamed/removed 事件。 +// 選這個數字的理由:搭配預設 700ms 節流,一輪約 17.5 秒完成,遠短於使用者能感知的 +// 「卡住」門檻,也讓下一輪 Scan()(PollSec 預設 5s 後)很快就能重新排序、 +// 把使用者剛寫的新檔插到隊伍最前面——不會被巨量積壓鎖住好幾小時才輪得到。 +const DefaultMaxEventsPerRun = 25 + +// effectiveMaxEventsPerRun 回傳這份設定實際生效的單輪事件上限(<=0 時用預設)。 +func (c *DirectConfig) effectiveMaxEventsPerRun() int { + if c.MaxEventsPerRun > 0 { + return c.MaxEventsPerRun + } + return DefaultMaxEventsPerRun +} + +// mtimeOf 讀檔案的修改時間;讀不到(已被刪除等)回零值,排序時會被當最舊、排到最後。 +func mtimeOf(absPath string) time.Time { + info, err := os.Stat(absPath) + if err != nil { + return time.Time{} + } + return info.ModTime() +} + +// sortEventsNewestFirst 把一輪的事件依「新改的檔優先」重排。 +// +// ⚠️ 刻意不改 Event struct(不加 Mtime 欄位):collector-trigger.v1.schema.json +// 是凍結版本、additionalProperties:false,欄位變更=開 v2 新檔(見 schema 檔頭註解)。 +// 這裡只重排 direct.go 本機處理佇列的順序,不動送上雲的 payload 形狀。 +// +// - added/modified/renamed:依檔案現在的 mtime 由新到舊;mtime 相同或讀不到時 +// 以路徑排序,保持每次呼叫結果一致(好測試、好除錯)。 +// - removed:檔案已經不存在,沒有 mtime 可排,一律排在內容事件之後 +// (下架不急,不該搶在使用者剛寫的新檔前面)。 +func sortEventsNewestFirst(absRoot string, events []Event) []Event { + out := make([]Event, len(events)) + copy(out, events) + + kind := func(ev Event) int { + if ev.Type == "removed" { + return 1 // 排在後段 + } + return 0 // added/modified/renamed + } + mt := make(map[string]time.Time, len(out)) + for _, ev := range out { + if ev.Type != "removed" { + mt[ev.Path] = mtimeOf(filepath.Join(absRoot, filepath.FromSlash(ev.Path))) + } + } + + sort.SliceStable(out, func(i, j int) bool { + ki, kj := kind(out[i]), kind(out[j]) + if ki != kj { + return ki < kj + } + if ki == 1 { // 兩者皆 removed:路徑序,穩定即可 + return out[i].Path < out[j].Path + } + ti, tj := mt[out[i].Path], mt[out[j].Path] + if !ti.Equal(tj) { + return ti.After(tj) // 新到舊 + } + return out[i].Path < out[j].Path + }) + return out +} + +// resumeAfterCapMessage 給「這輪因為單輪上限被延後」的事件用的說明—— +// 不是失敗,是刻意排隊;不該讓使用者以為程式壞了或漏了這個檔。 +func resumeAfterCapMessage(remaining int) string { + return fmt.Sprintf("已排入佇列,下一輪會繼續處理(本輪上限已到,還有 %d 筆等待)", remaining) +} diff --git a/direct_pacing_test.go b/direct_pacing_test.go new file mode 100644 index 0000000..94b90a3 --- /dev/null +++ b/direct_pacing_test.go @@ -0,0 +1,384 @@ +// direct_pacing_test.go — 上傳節奏/新檔優先/單輪上限/斷點續傳(2026-08-07)。 +// +// 背景見 direct_pacing.go 檔頭:封測事故實測 1,070 次寫入撞上免費上限 1,000, +// 690 個檔逐檔萃取上傳、daemon 這一半完全沒有節奏。這裡驗四件事: +// 1. 一輪掃到的多個事件依 mtime 新到舊處理(今天寫的最優先)。 +// 2. 兩次觸發雲端之間有節流間隔。 +// 3. 單輪上限存在——巨量積壓不會一次湧完。 +// 4. 已處理的事件會立刻落地 manifest,process 被殺掉重開也不會重做。 +package collector + +import ( + "encoding/json" + "io" + "net/http" + "net/http/httptest" + "os" + "path/filepath" + "strings" + "sync" + "testing" + "time" +) + +// init:測試環境預設把節流歸零,不然每個既有測試都要多等 700ms×N,拖慢整個套件。 +// 需要驗證節流本身的測試(見下)自行 save/restore 成一個很小的非零值。 +func init() { + directPaceInterval = 0 +} + +// ── 1) sortEventsNewestFirst ───────────────────────────────────────────── + +func TestSortEventsNewestFirst_NewestGoesFirst(t *testing.T) { + root := t.TempDir() + old := baseTime + mid := baseTime.Add(time.Hour) + newest := baseTime.Add(2 * time.Hour) + writeFile(t, root, "old.md", "old", old) + writeFile(t, root, "mid.md", "mid", mid) + writeFile(t, root, "newest.md", "newest", newest) + + events := []Event{ + {Type: "added", Path: "old.md"}, + {Type: "added", Path: "newest.md"}, + {Type: "added", Path: "mid.md"}, + } + got := sortEventsNewestFirst(root, events) + want := []string{"newest.md", "mid.md", "old.md"} + for i, w := range want { + if got[i].Path != w { + t.Fatalf("順序=%v,want %v", pathsOf(got), want) + } + } +} + +func pathsOf(evs []Event) []string { + out := make([]string, len(evs)) + for i, e := range evs { + out[i] = e.Path + } + return out +} + +// removed 事件沒有 mtime 可排,一律排在 added/modified/renamed 之後(下架不急)。 +func TestSortEventsNewestFirst_RemovedGoesLast(t *testing.T) { + root := t.TempDir() + writeFile(t, root, "new.md", "x", baseTime.Add(time.Hour)) + events := []Event{ + {Type: "removed", Path: "gone.md"}, + {Type: "added", Path: "new.md"}, + } + got := sortEventsNewestFirst(root, events) + if got[0].Path != "new.md" || got[1].Path != "gone.md" { + t.Fatalf("順序=%v", pathsOf(got)) + } +} + +// ── 2) 節流間隔 ──────────────────────────────────────────────────────────── + +func TestPace_RespectsInterval(t *testing.T) { + old := directPaceInterval + directPaceInterval = 30 * time.Millisecond + defer func() { directPaceInterval = old }() + + start := time.Now() + pace() + pace() + elapsed := time.Since(start) + if elapsed < 55*time.Millisecond { // 兩次 pace,留一點餘裕 + t.Fatalf("兩次 pace() 只花了 %v,節流間隔沒有生效", elapsed) + } +} + +// ── 端到端:積壓一次湧上去 vs 有節奏 ───────────────────────────────────────── +// +// 模擬「大量積壓」場景:8 個檔、單輪上限設 3——驗證①一輪只處理 3 個 +// ②被延後的檔案有交代(不是安靜消失)③下一輪接著處理下一批、依然新到舊。 + +func TestDirect_LargeBacklog_ProcessedInNewestFirstBatches(t *testing.T) { + root := t.TempDir() + // 8 個檔,mtime 依檔名反向遞增(a 最舊,h 最新) + names := []string{"a", "b", "c", "d", "e", "f", "g", "h"} + for i, n := range names { + writeFile(t, root, n+".md", "內容 "+n, baseTime.Add(time.Duration(i)*time.Minute)) + } + + var mu sync.Mutex + var order []string + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + body, _ := io.ReadAll(r.Body) + var m map[string]any + _ = json.Unmarshal(body, &m) + mu.Lock() + order = append(order, m["page_name"].(string)) + mu.Unlock() + _ = json.NewEncoder(w).Encode(map[string]any{"success": true}) + })) + defer srv.Close() + defer gemmaStub(t, func(w http.ResponseWriter, r *http.Request) { + // 這條測試只在意「觸發了幾次、依什麼順序送到假 cypher」,卡片內容固定即可。 + _ = json.NewEncoder(w).Encode(map[string]any{ + "candidates": []map[string]any{{ + "content": map[string]any{"parts": []map[string]any{{"text": "# 卡\n## 一句話定義\n測試\n"}}}, + }}, + }) + })() + + cfg := &DirectConfig{ + WatchFolders: []string{root}, + Manifest: filepath.Join(t.TempDir(), "m.json"), + CypherURL: srv.URL, Namespace: "demo", APIKey: "demo", + Library: "kb", Extractor: "gemma", ExtractorExplicit: true, GeminiAPIKey: "k-test", + CardIngestWF: "rag_ingest_card", MaxRemoved: DefaultMaxRemovedRatio, + MaxEventsPerRun: 3, + } + + // 第一輪:只處理 3 個(上限),且依 mtime 新到舊=h, g, f + results1, exit1, _ := RunDirectOnce(cfg, false) + if exit1 != 0 { + t.Fatalf("exit=%d results=%+v", exit1, results1) + } + ingested1 := ingestedPaths(results1) + if len(ingested1) != 3 { + t.Fatalf("第一輪應處理 3 個,got %d: %+v", len(ingested1), results1) + } + wantFirst := []string{"h.md", "g.md", "f.md"} + for i, w := range wantFirst { + if ingested1[i] != w { + t.Fatalf("第一輪順序=%v,want %v(新改的檔要優先)", ingested1, wantFirst) + } + } + // 有交代被延後了幾筆,不是安靜消失 + if !hasDeferredNotice(results1) { + t.Errorf("應該要交代還有事件被延後,results=%+v", results1) + } + + // 第二輪:接著處理下一批 3 個(e, d, c) + results2, exit2, _ := RunDirectOnce(cfg, false) + if exit2 != 0 { + t.Fatalf("exit=%d results=%+v", exit2, results2) + } + ingested2 := ingestedPaths(results2) + wantSecond := []string{"e.md", "d.md", "c.md"} + if len(ingested2) != 3 { + t.Fatalf("第二輪應處理 3 個,got %d: %+v", len(ingested2), results2) + } + for i, w := range wantSecond { + if ingested2[i] != w { + t.Fatalf("第二輪順序=%v,want %v", ingested2, wantSecond) + } + } + + // 第三輪:剩下 2 個(b, a)全部處理完,沒有再延後 + results3, exit3, _ := RunDirectOnce(cfg, false) + if exit3 != 0 { + t.Fatalf("exit=%d", exit3) + } + ingested3 := ingestedPaths(results3) + if len(ingested3) != 2 || ingested3[0] != "b.md" || ingested3[1] != "a.md" { + t.Fatalf("第三輪=%v,want [b.md a.md]", ingested3) + } + if hasDeferredNotice(results3) { + t.Error("全部處理完不該再有延後交代") + } + + // 全部 8 個都真的送到雲端了(沒有一個被永久漏掉) + mu.Lock() + defer mu.Unlock() + if len(order) != 8 { + t.Fatalf("雲端總共應收到 8 次卡片,got %d: %v", len(order), order) + } +} + +func ingestedPaths(results []DirectResult) []string { + var out []string + for _, r := range results { + if r.Status == "ingested" { + out = append(out, r.Path) + } + } + return out +} + +func hasDeferredNotice(results []DirectResult) bool { + for _, r := range results { + if r.Type == "info" && strings.Contains(r.Error, "已排入佇列") { + return true + } + } + return false +} + +// ── 3) 斷點續傳:process 被殺掉重開,不從頭來 ──────────────────────────────── +// +// 用「單輪上限=1」模擬中斷:一次只做一件事就等同「這一刻的程序被殺掉」, +// 用全新的 RunDirectOnce 呼叫(不共用任何記憶體狀態,manifest 完全從磁碟重讀) +// 模擬「程序重開」。驗證:已經成功的不會被重送,且每一步都真的落地磁碟 +// (不是等到全部做完才存檔)。 +func TestDirect_ResumeAfterInterruption(t *testing.T) { + root := t.TempDir() + writeFile(t, root, "one.md", "內容一", baseTime) + writeFile(t, root, "two.md", "內容二", baseTime.Add(time.Minute)) + writeFile(t, root, "three.md", "內容三", baseTime.Add(2*time.Minute)) + + var mu sync.Mutex + var posted []string + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + body, _ := io.ReadAll(r.Body) + var m map[string]any + _ = json.Unmarshal(body, &m) + mu.Lock() + posted = append(posted, m["path"].(string)) + mu.Unlock() + _ = json.NewEncoder(w).Encode(map[string]any{"success": true}) + })) + defer srv.Close() + defer gemmaStub(t, func(w http.ResponseWriter, r *http.Request) { + _ = json.NewEncoder(w).Encode(map[string]any{ + "candidates": []map[string]any{{ + "content": map[string]any{"parts": []map[string]any{{"text": "# 卡\n## 一句話定義\n測試\n"}}}, + }}, + }) + })() + + manifestPath := filepath.Join(t.TempDir(), "m.json") + newCfg := func() *DirectConfig { + return &DirectConfig{ + WatchFolders: []string{root}, + Manifest: manifestPath, + CypherURL: srv.URL, Namespace: "demo", APIKey: "demo", + Library: "kb", Extractor: "gemma", ExtractorExplicit: true, GeminiAPIKey: "k-test", + CardIngestWF: "rag_ingest_card", MaxRemoved: DefaultMaxRemovedRatio, + MaxEventsPerRun: 1, // 模擬「只做一件事就被打斷」 + } + } + + // 「跑到一半」:只完成第一批(three.md,最新) + if _, exit, _ := RunDirectOnce(newCfg(), false); exit != 0 { + t.Fatal("第一批失敗") + } + + // 直接讀磁碟上的 manifest(不透過任何記憶體物件)——驗證真的已經落地, + // 不是要等三批都跑完才存檔。 + absRoot, _ := filepath.Abs(root) + mp := newCfg().manifestPathFor(absRoot) + m1, err := LoadManifest(mp, absRoot) + if err != nil { + t.Fatalf("讀 manifest 失敗:%v", err) + } + if m1.Entries["three.md"].IngestedHash == "" { + t.Fatal("「跑到一半」之後,已完成的那份應該已經落地 manifest(斷點續傳的前提)") + } + if m1.Entries["two.md"].IngestedHash != "" || m1.Entries["one.md"].IngestedHash != "" { + t.Fatal("還沒輪到的不該被誤標成已完成") + } + + // 「重開程序」:全新呼叫,模擬 process 重啟——應該接著做 two.md,不是重做 three.md + if _, exit, _ := RunDirectOnce(newCfg(), false); exit != 0 { + t.Fatal("第二批失敗") + } + if _, exit, _ := RunDirectOnce(newCfg(), false); exit != 0 { + t.Fatal("第三批失敗") + } + + mu.Lock() + defer mu.Unlock() + if len(posted) != 3 { + t.Fatalf("三個檔應該總共只被送出 3 次(不重送已完成的),got %d: %v", len(posted), posted) + } + seen := map[string]bool{} + for _, p := range posted { + if seen[p] { + t.Fatalf("%q 被重複送出——斷點續傳失效,重開後從頭來了", p) + } + seen[p] = true + } + + // 第四輪:三個都做完了,不該再有任何事件 + results4, _, _ := RunDirectOnce(newCfg(), false) + if len(results4) != 0 { + t.Fatalf("全部完成後應該零事件:%+v", results4) + } +} + +// removed 事件在「已放回但下架失敗」時不能被提早存檔清掉——否則下一輪兩邊都找不到 +// 這個路徑,永遠不會再重試下架。驗證:POST 失敗時,manifest 磁碟版本仍保留該筆, +// 下一輪會重新產生 removed 事件。 +func TestDirect_RemovedRetriesOnFailureEvenWithIncrementalSave(t *testing.T) { + root := t.TempDir() + writeFile(t, root, "keep.md", "留著", baseTime) + writeFile(t, root, "gone.md", "要刪的", baseTime.Add(time.Minute)) + + var failRemoved = true + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if strings.Contains(r.URL.Path, "rag_takedown_direct") && failRemoved { + w.WriteHeader(http.StatusInternalServerError) + _, _ = w.Write([]byte(`{"error":"boom"}`)) + return + } + _ = json.NewEncoder(w).Encode(map[string]any{"success": true}) + })) + defer srv.Close() + defer gemmaStub(t, func(w http.ResponseWriter, r *http.Request) { + _ = json.NewEncoder(w).Encode(map[string]any{ + "candidates": []map[string]any{{ + "content": map[string]any{"parts": []map[string]any{{"text": "# 卡\n## 一句話定義\n測試\n"}}}, + }}, + }) + })() + + cfg := &DirectConfig{ + WatchFolders: []string{root}, + Manifest: filepath.Join(t.TempDir(), "m.json"), + CypherURL: srv.URL, Namespace: "demo", APIKey: "demo", + Library: "kb", Extractor: "gemma", ExtractorExplicit: true, GeminiAPIKey: "k-test", + CardIngestWF: "rag_ingest_card", RemovedWF: "rag_takedown_direct", + MaxRemoved: 1.0, + } + if _, exit, _ := RunDirectOnce(cfg, false); exit != 0 { + t.Fatal("第一輪(兩份都新增)失敗") + } + if err := os.Remove(filepath.Join(root, "gone.md")); err != nil { + t.Fatal(err) + } + + // 第二輪:下架 POST 會失敗(且這輪也會處理 keep.md 的事件?不會,keep.md 已 ingest 過、 + // 沒有變化,不會再產生事件——但仍會呼叫 saveManifest() 若有其他觸發)。 + results, exit, _ := RunDirectOnce(cfg, false) + if exit != 1 { + t.Fatalf("下架失敗應該 exit=1,got %d results=%+v", exit, results) + } + + absRoot, _ := filepath.Abs(root) + mp := cfg.manifestPathFor(absRoot) + m, err := LoadManifest(mp, absRoot) + if err != nil { + t.Fatal(err) + } + if _, ok := m.Entries["gone.md"]; !ok { + t.Fatal("下架失敗時,manifest 不該把這個路徑永久丟掉——否則永遠不會再重試下架") + } + + // 第三輪:換伺服器成功 → 應該重新補發並成功下架 + failRemoved = false + results3, exit3, _ := RunDirectOnce(cfg, false) + if exit3 != 0 { + t.Fatalf("第三輪應成功下架:exit=%d results=%+v", exit3, results3) + } + var removedOK bool + for _, r := range results3 { + if r.Type == "removed" && r.Status == "removed" && r.Path == "gone.md" { + removedOK = true + } + } + if !removedOK { + t.Fatalf("第三輪應該重試下架成功:%+v", results3) + } + m2, err := LoadManifest(mp, absRoot) + if err != nil { + t.Fatal(err) + } + if _, ok := m2.Entries["gone.md"]; ok { + t.Fatal("下架成功後這個路徑應該真的從 manifest 消失") + } +} diff --git a/direct_quota_test.go b/direct_quota_test.go new file mode 100644 index 0000000..54c50f1 --- /dev/null +++ b/direct_quota_test.go @@ -0,0 +1,154 @@ +// direct_quota_test.go — 額度用完要說人話並降速(2026-08-07 task 2)端到端驗證。 +// +// 模擬 wiki mistakes.md 2026-08-06 記錄的真實上游狀況:Workers AI 每日免費額度 +// 10,000 neurons 用完,HTTP 502,訊息含「4006: you have used up your daily free +// allocation of 10,000 neurons」。驗收:畫面(status.json)出現三句話, +// 沒有裸露的錯誤碼;且同一輪/下一輪不再繼續撞同一面牆。 +package collector + +import ( + "encoding/json" + "net/http" + "net/http/httptest" + "path/filepath" + "strings" + "testing" + "time" +) + +// quotaExhaustedServer 對任何 /portal/daemon/extract 請求都回上游額度用完的原始錯誤。 +// calls 只計「真的萃取」請求(body.text 非空)——ProbeWorkersAI 每輪固定會探測一次 +// (送空 text,不燒 LLM 額度,見 probe_workersai.go),不是本測試要驗的「有沒有繼續撞牆」。 +func quotaExhaustedServer(t *testing.T, calls *int) *httptest.Server { + t.Helper() + return httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + var body struct { + Text string `json:"text"` + } + _ = json.NewDecoder(r.Body).Decode(&body) + if strings.TrimSpace(body.Text) != "" { + *calls++ + } + w.WriteHeader(http.StatusBadGateway) // 實據:HTTP 502 + _ = json.NewEncoder(w).Encode(map[string]string{ + "error": "4006: you have used up your daily free allocation of 10,000 neurons", + }) + })) +} + +func TestDirect_QuotaExhausted_ShowsHumanMessageNoRawCode(t *testing.T) { + root := t.TempDir() + writeFile(t, root, "a.md", "內容 A", baseTime) + writeFile(t, root, "b.md", "內容 B", baseTime.Add(time.Minute)) + + var calls int + srv := quotaExhaustedServer(t, &calls) + defer srv.Close() + + manifest := filepath.Join(t.TempDir(), "m.json") + cfg := &DirectConfig{ + WatchFolders: []string{root}, + Manifest: manifest, + CypherURL: srv.URL, Namespace: "demo", APIKey: "demo", + Library: "kb", Extractor: "workers-ai", ExtractorExplicit: true, + MaxRemoved: DefaultMaxRemovedRatio, + } + + results, exit, _ := RunDirectOnce(cfg, false) + if exit != 1 { + t.Fatalf("兩份都該萃取失敗,exit 應為 1,got %d", exit) + } + if len(results) != 2 { + t.Fatalf("results=%+v", results) + } + + banned := []string{"4006", "502", "HTTP", "neurons"} + for _, r := range results { + for _, b := range banned { + if strings.Contains(r.Error, b) { + t.Errorf("結果不該出現裸露的錯誤碼 %q:%q", b, r.Error) + } + } + for _, want := range []string{"今天已經幫你整理了", "升級 Cloudflare", "自動恢復"} { + if !strings.Contains(r.Error, want) { + t.Errorf("缺三句話之一 %q:%q", want, r.Error) + } + } + } + + // status.json 也要是人話,同樣不准有裸碼 + st, err := LoadSyncStatus(StatusFilePath(manifest)) + if err != nil { + t.Fatalf("讀 status.json 失敗:%v", err) + } + if st.ExtractorOK { + t.Error("額度冷卻中,ExtractorOK 應為 false") + } + for _, b := range banned { + if strings.Contains(st.ExtractorError, b) { + t.Errorf("status.ExtractorError 不該有裸碼 %q:%q", b, st.ExtractorError) + } + } + acc, ok := st.AccountDetails[instanceHostOf(srv.URL)] + if !ok { + t.Fatal("找不到帳號的 status") + } + if acc.QuotaMessage == nil { + t.Fatal("應該有 QuotaMessage") + } + if acc.QuotaCooldownUntil == "" { + t.Fatal("應該記下冷卻到什麼時候") + } + for _, b := range banned { + combined := acc.QuotaMessage.Combined() + if strings.Contains(combined, b) { + t.Errorf("QuotaMessage 不該有裸碼 %q:%q", b, combined) + } + } + + // 只有第一次真的打了 API(第二份的 xerr 判斷不需要——一偵測到就進冷卻, + // 第二個事件被上面的 qs.inCooldown 擋下,不再觸網)。 + if calls != 1 { + t.Fatalf("偵測到額度用完後不該繼續撞牆,got %d 次呼叫", calls) + } +} + +// 同一帳號、下一輪(新的 RunDirectOnce 呼叫,模擬下一次輪詢):冷卻還沒過, +// 一次網路呼叫都不該再打——這是「不要一次湧上去、也不要一直撞牆」的核心。 +func TestDirect_QuotaExhausted_NextRunSkipsWithoutNetworkCall(t *testing.T) { + root := t.TempDir() + writeFile(t, root, "a.md", "內容 A", baseTime) + + var calls int + srv := quotaExhaustedServer(t, &calls) + defer srv.Close() + + cfg := &DirectConfig{ + WatchFolders: []string{root}, + Manifest: filepath.Join(t.TempDir(), "m.json"), + CypherURL: srv.URL, Namespace: "demo", APIKey: "demo", + Library: "kb", Extractor: "workers-ai", ExtractorExplicit: true, + MaxRemoved: DefaultMaxRemovedRatio, + } + + if _, exit, _ := RunDirectOnce(cfg, false); exit != 1 { + t.Fatal("第一輪應該偵測到額度用完") + } + if calls != 1 { + t.Fatalf("第一輪應該只打 1 次,got %d", calls) + } + + results2, exit2, _ := RunDirectOnce(cfg, false) + if exit2 != 0 { + t.Fatalf("冷卻中的 skip 不是失敗,exit 應為 0,got %d results=%+v", exit2, results2) + } + if calls != 1 { + t.Fatalf("冷卻中不該再打任何一次 API,got %d 次呼叫", calls) + } + if len(results2) != 1 || results2[0].Status != "skipped" { + t.Fatalf("results2=%+v", results2) + } + if !strings.Contains(results2[0].Error, "自動恢復") { + t.Errorf("跳過訊息也該是三句話:%q", results2[0].Error) + } +} diff --git a/quota.go b/quota.go new file mode 100644 index 0000000..ffe3b9a --- /dev/null +++ b/quota.go @@ -0,0 +1,120 @@ +// quota.go — Workers AI 每日免費額度用完時的人話訊息與全域冷卻(2026-08-07 pacing task)。 +// +// 背景(封測事故):Evan 在自己的免費 Cloudflare 帳號上跑,690 個檔逐檔萃取上傳, +// 每個檔在雲端產生一次工作流執行 ⇒ 1,070 次寫入撞上免費上限 1,000,額度爆掉、整批 429/502, +// 只有 8 個檔成功。雲端那一半(紀錄改走資料層 API)已修好;本檔補的是 daemon 這一半: +// 額度用完時要「說人話+降速」,不是安靜地重試到死或吐一串技術錯誤碼。 +// +// 已知的上游錯誤模式(system-dev/wiki/mistakes.md 2026-08-06): +// +// collector.log 出現 `4006: you have used up your daily free allocation of 10,000 neurons`, +// HTTP 502,額度**每日 UTC 午夜重置**(=台灣時間早上 8:00)。 +// +// leo 定的三句話骨架(缺一不可,且不准出現裸露的錯誤碼): +// +// 成就:「今天已經幫你整理了 N 份」← 先講做到什麼,不是先講失敗 +// 出口:「可以換一個模型,或升級 Cloudflare(每月 5 美元)」← 給選擇不是死路 +// 保證:「不花錢也沒關係,明天/今天早上 8 點會自動恢復、會接著跑」 +package collector + +import ( + "fmt" + "strings" + "time" +) + +// taiwanTZ:固定 UTC+8 偏移,不吃系統 tzdata(跨平台/Windows 常缺時區資料庫; +// 台灣本身沒有日光節約時間,固定偏移在任何時候都正確——不必依賴 time.LoadLocation)。 +var taiwanTZ = time.FixedZone("Asia/Taipei", 8*60*60) + +// isQuotaExhausted 辨識「Workers AI 每日免費額度用完」這個已知上游狀況(非 bug)。 +// 用子字串比對而非單看 HTTP 狀態碼——502 本身太泛用(也可能是別的暫時性錯誤), +// 這段特定文案才是可靠訊號(見 wiki mistakes.md 2026-08-06 記錄的原文)。 +func isQuotaExhausted(errMsg string) bool { + return strings.Contains(errMsg, "10,000 neurons") || + strings.Contains(errMsg, "daily free allocation") || + strings.Contains(errMsg, "4006:") +} + +// nextQuotaResetTaiwan 回傳下一次 Workers AI 每日額度重置的時間點。 +// Cloudflare 在 UTC 午夜重置 ⇒ 下一次 UTC 00:00 就是答案(換算成台灣時間固定是早上 8:00)。 +func nextQuotaResetTaiwan(now time.Time) time.Time { + u := now.UTC() + nextMidnightUTC := time.Date(u.Year(), u.Month(), u.Day(), 0, 0, 0, 0, time.UTC).AddDate(0, 0, 1) + return nextMidnightUTC.In(taiwanTZ) +} + +// QuotaNotice 是額度用完時要給使用者看的三句話(leo 定的骨架,缺一不可)。 +// 三句合起來就是完整訊息,任何管道要顯示都直接串接,不再另外組字串(避免各處措辭漂移)。 +type QuotaNotice struct { + Achievement string `json:"achievement"` // 今天已經幫你整理了 N 份 + ExitOptions string `json:"exit_options"` // 可以換一個模型,或升級 Cloudflare(每月 5 美元) + Guarantee string `json:"guarantee"` // 不花錢也沒關係,今天/明天早上 8:00 會自動恢復 + ResumeAt string `json:"resume_at"` // RFC3339,預期恢復時間(供機器判斷冷卻是否結束) +} + +// Combined 把三句話接成一句完整訊息(給只有單一 error 欄位可用的地方,如 DirectResult.Error)。 +// 用句號分隔——三句話都要出現,缺一不可,且全程不含任何原始錯誤碼。 +func (n QuotaNotice) Combined() string { + return n.Achievement + "。" + n.ExitOptions + "。" + n.Guarantee +} + +// buildQuotaNotice 組出三句話。dailyCount=今天(UTC 日界,與額度重置同一條線)已成功 +// 萃取的份數;resetAt=nextQuotaResetTaiwan 算出的下一次重置時間。 +func buildQuotaNotice(now time.Time, dailyCount int, resetAt time.Time) QuotaNotice { + return QuotaNotice{ + Achievement: fmt.Sprintf("今天已經幫你整理了 %d 份", dailyCount), + ExitOptions: "可以換一個模型,或升級 Cloudflare(每月 5 美元)", + Guarantee: quotaGuaranteeText(now, resetAt), + ResumeAt: resetAt.Format(time.RFC3339), + } +} + +// quotaGuaranteeText 把重置時間換成人話:「今天」或「明天」早上 8:00 +// (不能寫死「明天」——若這一刻台灣時間已經過了午夜、還沒到 8 點,重置其實是「今天」)。 +func quotaGuaranteeText(now, resetAt time.Time) string { + nowTW := now.In(taiwanTZ) + resetTW := resetAt.In(taiwanTZ) + dayWord := "明天" + if nowTW.Year() == resetTW.Year() && nowTW.YearDay() == resetTW.YearDay() { + dayWord = "今天" + } + return fmt.Sprintf("不花錢也沒關係,%s早上 8:00 會自動恢復、會接著跑", dayWord) +} + +// quotaState 是「這個帳號本輪的額度冷卻」共享狀態,跨同一帳號的多個監看根 +// (額度是雲端實例/Cloudflare 帳號層級的,不是單一資料夾的——一個根撞到, +// 同帳號其他根不該還傻傻地繼續撞同一面牆)。由 RunDirectOnce 建立、以指標 +// 傳進每個 runDirectOnceRoot 呼叫,狀態在呼叫之間累積。 +type quotaState struct { + CooldownUntil time.Time // 非零值且晚於 now ⇒ 本帳號本輪不再嘗試萃取 + Hit bool // 本輪是否新偵測到額度用完(避免同一輪反覆覆寫 CooldownUntil) + RawReason string // 上游原文,只供 log/除錯,不進任何使用者看得到的欄位 + DailyCount int // 今天(UTC 日界)已成功萃取的份數,seed 自 status.json 並持續累加 +} + +// inCooldown 回報現在是否仍在額度冷卻窗口內。 +func (qs *quotaState) inCooldown(now time.Time) bool { + return !qs.CooldownUntil.IsZero() && now.Before(qs.CooldownUntil) +} + +// noticeNow 用目前狀態組一份 QuotaNotice。 +func (qs *quotaState) noticeNow(now time.Time) QuotaNotice { + return buildQuotaNotice(now, qs.DailyCount, qs.CooldownUntil) +} + +// markHit 記錄「這一刻偵測到額度用完」,只在本輪第一次命中時真正設定冷卻時間 +// (之後同一輪的其他失敗不再往後推遲 CooldownUntil,避免因為連續撞牆而不斷延後恢復承諾)。 +func (qs *quotaState) markHit(now time.Time, rawReason string) { + if qs.Hit { + return + } + qs.Hit = true + qs.RawReason = rawReason + qs.CooldownUntil = nextQuotaResetTaiwan(now) +} + +// todayUTC 回傳 YYYY-MM-DD(UTC),與 Workers AI 額度重置同一條日界線。 +func todayUTC(now time.Time) string { + return now.UTC().Format("2006-01-02") +} diff --git a/quota_test.go b/quota_test.go new file mode 100644 index 0000000..7cd2c28 --- /dev/null +++ b/quota_test.go @@ -0,0 +1,110 @@ +// quota_test.go — 額度用完人話訊息(2026-08-07 pacing task 2)。 +package collector + +import ( + "strings" + "testing" + "time" +) + +func TestIsQuotaExhausted(t *testing.T) { + cases := []struct { + msg string + want bool + }{ + {"本地萃取失敗:雲端萃取失敗(HTTP 502):4006: you have used up your daily free allocation of 10,000 neurons", true}, + {"雲端萃取失敗:exceeded daily free allocation", true}, + {"連不上你的知識庫:dial tcp: connection refused", false}, + {"檔案裡沒有可抽取的文字", false}, + {"雲端萃取失敗(HTTP 500):internal error", false}, + } + for _, c := range cases { + if got := isQuotaExhausted(c.msg); got != c.want { + t.Errorf("isQuotaExhausted(%q)=%v want %v", c.msg, got, c.want) + } + } +} + +// UTC 午夜重置=台灣時間固定早上 8:00。 +func TestNextQuotaResetTaiwan(t *testing.T) { + // 2026-08-07 15:00 UTC = 2026-08-07 23:00 台灣 → 下次重置 2026-08-08 00:00 UTC = 08-08 08:00 台灣 + now := time.Date(2026, 8, 7, 15, 0, 0, 0, time.UTC) + got := nextQuotaResetTaiwan(now) + want := time.Date(2026, 8, 8, 8, 0, 0, 0, taiwanTZ) + if !got.Equal(want) { + t.Fatalf("got %v want %v", got, want) + } +} + +// 三句話缺一不可,且絕不含裸露的錯誤碼(4006/502/HTTP)。 +func TestBuildQuotaNotice_NoRawErrorCode(t *testing.T) { + now := time.Date(2026, 8, 7, 10, 0, 0, 0, time.UTC) + resetAt := nextQuotaResetTaiwan(now) + notice := buildQuotaNotice(now, 42, resetAt) + + if !strings.Contains(notice.Achievement, "42") { + t.Errorf("成就句沒帶到份數:%q", notice.Achievement) + } + if !strings.Contains(notice.ExitOptions, "升級") || !strings.Contains(notice.ExitOptions, "換") { + t.Errorf("出口句缺選項:%q", notice.ExitOptions) + } + if !strings.Contains(notice.Guarantee, "自動恢復") { + t.Errorf("保證句沒講自動恢復:%q", notice.Guarantee) + } + combined := notice.Combined() + for _, banned := range []string{"4006", "502", "HTTP", "neurons"} { + if strings.Contains(combined, banned) { + t.Errorf("三句話不該出現裸露的錯誤碼 %q:%q", banned, combined) + } + } +} + +// 「今天」vs「明天」:現在台灣時間若已過午夜還沒到 8 點,重置其實是「今天」。 +func TestQuotaGuaranteeText_TodayVsTomorrow(t *testing.T) { + // 台灣時間 08-08 02:00(= UTC 08-07 18:00)→ 下次重置 08-08 08:00 台灣 → 同一天 → 「今天」 + now := time.Date(2026, 8, 7, 18, 0, 0, 0, time.UTC) + resetAt := nextQuotaResetTaiwan(now) + got := quotaGuaranteeText(now, resetAt) + if !strings.Contains(got, "今天") { + t.Fatalf("應該是「今天」:%q(now台灣=%v resetAt台灣=%v)", got, now.In(taiwanTZ), resetAt.In(taiwanTZ)) + } + + // 台灣時間 08-07 23:00(= UTC 08-07 15:00)→ 下次重置 08-08 08:00 台灣 → 隔天 → 「明天」 + now2 := time.Date(2026, 8, 7, 15, 0, 0, 0, time.UTC) + resetAt2 := nextQuotaResetTaiwan(now2) + got2 := quotaGuaranteeText(now2, resetAt2) + if !strings.Contains(got2, "明天") { + t.Fatalf("應該是「明天」:%q", got2) + } +} + +func TestQuotaState_MarkHitOnlyFirstTimeSetsCooldown(t *testing.T) { + qs := "aState{} + now := time.Date(2026, 8, 7, 10, 0, 0, 0, time.UTC) + qs.markHit(now, "第一次原因") + first := qs.CooldownUntil + + later := now.Add(90 * time.Minute) + qs.markHit(later, "第二次原因(同一輪另一個檔也撞到)") + if !qs.CooldownUntil.Equal(first) { + t.Fatalf("第二次命中不該再往後推遲冷卻時間:first=%v got=%v", first, qs.CooldownUntil) + } + if qs.RawReason != "第一次原因" { + t.Fatalf("RawReason 應保留第一次的:%q", qs.RawReason) + } +} + +func TestQuotaState_InCooldown(t *testing.T) { + qs := "aState{} + now := time.Date(2026, 8, 7, 10, 0, 0, 0, time.UTC) + if qs.inCooldown(now) { + t.Fatal("從沒命中過不該在冷卻中") + } + qs.markHit(now, "x") + if !qs.inCooldown(now) { + t.Fatal("剛命中應立刻在冷卻中") + } + if qs.inCooldown(qs.CooldownUntil.Add(time.Second)) { + t.Fatal("過了冷卻時間應該解除") + } +} diff --git a/scan.go b/scan.go index c059f3a..65904c9 100644 --- a/scan.go +++ b/scan.go @@ -114,6 +114,99 @@ type TriggerPayload struct { // 可能存成了別的格式)。只報總數在「幾百張圖」時是對的,在「1 個」時等於沒說。 // ⇒ 少量時就點名,讓使用者自己一眼看出「喔,我存錯格式了」。 SkippedOtherNames []string `json:"-"` + + // DuplicateFormats=本輪偵測到、同檔名主幹的多格式重複(2026-08-07,見 FormatDuplicate)。 + // 同 Skipped:只給本機使用者看,不隨 payload 送雲端(schema additionalProperties:false 會擋)。 + DuplicateFormats []FormatDuplicate `json:"-"` +} + +// FormatDuplicate=同一份內容被偵測到有多種格式並存(同檔名主幹、不同副檔名)。 +// +// 🔴 為什麼要有這個(2026-08-07,leo 實據):封測者 Evan 給的資料集 +// 27,164 檔=9,045 個 md+9,044 個 json+9,043 個 html,是同一批來源轉出的三種格式 +// (`markdown/160-00F3_001.md` 與 `json/160-00F3_001.json` 逐一對應,同錯誤碼、同變體號)。 +// 現行設計不知道這件事,會把三種格式各當一份新檔,萃三次、吃三倍額度—— +// 使用者不會知道要挑一種,這是我們該擋的,不是他該懂的。 +// +// 判準:檔名主幹(去副檔名、轉小寫)在整個看守根內相同 → 視為同一份內容的不同格式匯出, +// 只留優先序最高的一份進事件管線,其餘標記在這裡(不產生 added/modified/renamed 事件)。 +// +// ⚠️ 已知的取捨:純靠檔名主幹比對,不比對內容——兩個不相干的檔恰好同名不同副檔名 +// (如兩個專案各自的 `README.md`/`README.pdf`)會被誤判成同一份。刻意接受這個風險, +// 因為:①這不是靜默丟棄——DuplicateFormats 會被 direct.go 收進 status.json 讓使用者看到 +// 「跳過:與 X 視為同一格式」,看起來不對可以改檔名破解誤判;②不比對,就是 leo 實據的 +// md/json 案例(不同目錄、不同副檔名,唯一共同點正是檔名主幹)根本擋不掉。 +// +// 🔴 刻意不寫進 manifest(同 Skipped 的理由,見上方 SkippedFile 註解):每輪由檔案系統 +// 重算,勝出者若之後消失,下一輪換另一份自然遞補,不留跨輪狀態要維護。 +type FormatDuplicate struct { + Path string `json:"path"` // 被跳過的那份 + KeptPath string `json:"kept_path"` // 真正進事件管線的那份 + Stem string `json:"stem"` // 判定依據:去副檔名、轉小寫後的檔名主幹 +} + +// dedupFormatPriority:格式去重時「留誰」的優先序(越前面越優先留下)。 +// 原則:越接近「使用者原始編輯」的格式排越前面(docx/pptx 是可編輯原稿), +// md/markdown 常是**從別的格式轉出的產物**(leo 實據的資料集正是 html→chm 轉出 md/json), +// 排在轉檔產物之前但在原生辦公格式之後。不在表內的副檔名排最後(理論上不會發生, +// current 只收 allowedExt)。 +var dedupFormatPriority = []string{".docx", ".pptx", ".xlsx", ".csv", ".pdf", ".md", ".markdown", ".txt"} + +func formatPriority(relPath string) int { + ext := strings.ToLower(filepath.Ext(relPath)) + for i, e := range dedupFormatPriority { + if e == ext { + return i + } + } + return len(dedupFormatPriority) +} + +// dedupStemOf 回傳去重判準:檔名主幹(basename 去副檔名),轉小寫(跨平台大小寫不敏感)。 +// 刻意只看 basename、不看目錄——leo 實據的 md/json 恰好活在不同的兄弟目錄 +// (markdown/160-00F3_001.md vs json/160-00F3_001.json),只有 basename 主幹相同。 +func dedupStemOf(relPath string) string { + base := filepath.Base(relPath) + stem := strings.TrimSuffix(base, filepath.Ext(base)) + return strings.ToLower(stem) +} + +// detectFormatDuplicates 在 current(本輪掃到、通過 allowedExt 的檔)裡找出同檔名主幹的分組, +// 每組留優先序最高的一份,其餘回報為 loser(path -> winner path)。 +// 走訪用 stem 字母序=確定性輸出(map 迭代順序不穩定)。 +func detectFormatDuplicates(current map[string]fileState) (map[string]string, []FormatDuplicate) { + byStem := map[string][]string{} + for p := range current { + stem := dedupStemOf(p) + byStem[stem] = append(byStem[stem], p) + } + stems := make([]string, 0, len(byStem)) + for s := range byStem { + stems = append(stems, s) + } + sort.Strings(stems) + + loserOf := map[string]string{} + var dups []FormatDuplicate + for _, stem := range stems { + group := byStem[stem] + if len(group) < 2 { + continue + } + sort.Slice(group, func(i, j int) bool { + pi, pj := formatPriority(group[i]), formatPriority(group[j]) + if pi != pj { + return pi < pj + } + return group[i] < group[j] // 同優先序時字母序,確定性 + }) + winner := group[0] + for _, loser := range group[1:] { + loserOf[loser] = winner + dups = append(dups, FormatDuplicate{Path: loser, KeptPath: winner, Stem: stem}) + } + } + return loserOf, dups } type ScanOptions struct { @@ -239,9 +332,16 @@ func Scan(root string, m *Manifest, opts ScanOptions) (*TriggerPayload, error) { return nil, err } + // 1.5) 同內容多格式去重(2026-08-07):先決定誰是 loser,事件管線全程跳過它們。 + dupLoser, duplicateFormats := detectFormatDuplicates(current) + // 2) 初分:added 候選(現況有、manifest 無)與 removed 候選(manifest 有、現況無)。 + // loser 不進候選——它不該被當成新檔,也不該被當成 rename 的另一端。 var addedPaths, removedPaths []string for p := range current { + if dupLoser[p] != "" { + continue + } if _, ok := orig[p]; !ok { addedPaths = append(addedPaths, p) } @@ -289,6 +389,9 @@ func Scan(root string, m *Manifest, opts ScanOptions) (*TriggerPayload, error) { return Event{Type: "added", Path: p, SourceHash: st.hash, Size: &size, R2Key: r2KeyOf(st.hash)} } for _, p := range sortedCurrent { + if dupLoser[p] != "" { // 同內容的另一格式已在事件管線,這份跳過(1.5) + continue + } if op, isRenamed := renamedOldOf[p]; isRenamed { if orig[op].IngestedHash == "" { // 改名的檔其實從未 ingest 成功 → 補一發 added events = append(events, addedEvent(p)) @@ -304,6 +407,9 @@ func Scan(root string, m *Manifest, opts ScanOptions) (*TriggerPayload, error) { // 5) modified:manifest 有、現況有、content_hash != ingested_hash(design §3 順序 3)。 for _, p := range sortedCurrent { + if dupLoser[p] != "" { // 同內容的另一格式已在事件管線,這份跳過(1.5) + continue + } e, existed := orig[p] if !existed { continue @@ -392,5 +498,6 @@ func Scan(root string, m *Manifest, opts ScanOptions) (*TriggerPayload, error) { Skipped: skipped, SkippedOther: skippedOther, SkippedOtherNames: skippedOtherNames, + DuplicateFormats: duplicateFormats, }, nil } diff --git a/scan_dedup_test.go b/scan_dedup_test.go new file mode 100644 index 0000000..125eb86 --- /dev/null +++ b/scan_dedup_test.go @@ -0,0 +1,141 @@ +// scan_dedup_test.go — 同內容多格式去重(2026-08-07,見 scan.go FormatDuplicate 註解)。 +// +// 實據:leo 給的封測資料集 27,164 檔=9,045 md+9,044 json+9,043 html, +// 是同一批來源轉出的三種格式,md/json 逐一同檔名主幹(不同目錄)。 +// 驗收:「同內容三種格式並存 → 只產一張卡」。 +package collector + +import ( + "os" + "path/filepath" + "testing" +) + +// 三種格式、同檔名主幹、活在兄弟目錄(照實據的 markdown/ vs json/ 結構) +// → 只有一份(依優先序 docx > pptx > xlsx > csv > pdf > md > markdown > txt)進事件管線。 +func TestScan_FormatDuplicate_OnlyOneEventPerStem(t *testing.T) { + root := t.TempDir() + writeFile(t, root, "markdown/160-00F3_001.md", "# 錯誤說明\nM118/M128 不可同時使用", baseTime) + writeFile(t, root, "csv/160-00F3_001.csv", "code,msg\n160-00F3,M118/M128 不可同時使用", baseTime) + writeFile(t, root, "doc/160-00F3_001.docx", "docx 二進位占位", baseTime) + + m := newTestManifest() + payload, err := Scan(root, m, ScanOptions{}) + if err != nil { + t.Fatal(err) + } + + var added []Event + for _, ev := range payload.Events { + if ev.Type == "added" { + added = append(added, ev) + } + } + if len(added) != 1 { + t.Fatalf("三種格式應只產生 1 個 added 事件,got %d: %+v", len(added), added) + } + if added[0].Path != "doc/160-00F3_001.docx" { + t.Fatalf("優先序應留 docx,got %q", added[0].Path) + } + + if len(payload.DuplicateFormats) != 2 { + t.Fatalf("應回報 2 份被跳過的重複格式,got %d: %+v", len(payload.DuplicateFormats), payload.DuplicateFormats) + } + for _, d := range payload.DuplicateFormats { + if d.KeptPath != "doc/160-00F3_001.docx" { + t.Errorf("KeptPath=%q,應指向留下的那份", d.KeptPath) + } + if d.Stem != "160-00f3_001" { + t.Errorf("Stem=%q", d.Stem) + } + } +} + +// 不同檔名主幹的檔案(真正不相干的內容)不受影響——去重不能誤傷正常檔案。 +func TestScan_FormatDuplicate_DifferentStemsUnaffected(t *testing.T) { + root := t.TempDir() + writeFile(t, root, "a.md", "內容 A", baseTime) + writeFile(t, root, "b.pdf", "內容 B", baseTime) + + m := newTestManifest() + payload, err := Scan(root, m, ScanOptions{}) + if err != nil { + t.Fatal(err) + } + added := 0 + for _, ev := range payload.Events { + if ev.Type == "added" { + added++ + } + } + if added != 2 { + t.Fatalf("不同主幹的檔案應各自產生事件,got %d", added) + } + if len(payload.DuplicateFormats) != 0 { + t.Fatalf("不該誤判為重複:%+v", payload.DuplicateFormats) + } +} + +// loser 仍留在 manifest(讓下一輪 mtime/size fast-path 照常運作),只是沒有事件、 +// 且 ingested_hash 永遠不會被設定——不會在後續掃描裡被誤判成「新檔」而不斷回報。 +func TestScan_FormatDuplicate_LoserStaysInManifestNoRepeatedEvents(t *testing.T) { + root := t.TempDir() + writeFile(t, root, "x.docx", "docx 內容", baseTime) + writeFile(t, root, "x.md", "md 內容(同主幹)", baseTime) + + m := newTestManifest() + if _, err := Scan(root, m, ScanOptions{}); err != nil { + t.Fatal(err) + } + if _, ok := m.Entries["x.md"]; !ok { + t.Fatal("loser 應仍被記錄在 manifest(供 fast-path 用)") + } + if m.Entries["x.md"].IngestedHash != "" { + t.Fatal("loser 不該有 ingested_hash(它從未真的被送出)") + } + + // 模擬 winner 已成功 ingest + m.Entries["x.docx"].IngestedHash = m.Entries["x.docx"].ContentHash + + // 第二輪:兩份檔案都沒變 → loser 不該又跑出一個 added 事件 + payload2, err := Scan(root, m, ScanOptions{}) + if err != nil { + t.Fatal(err) + } + for _, ev := range payload2.Events { + if ev.Path == "x.md" { + t.Fatalf("loser 不該重複產生事件:%+v", ev) + } + } +} + +// winner 消失後,loser 自然遞補(下一輪重算分組時 loser 變成該 stem 唯一成員)。 +func TestScan_FormatDuplicate_WinnerRemovedLoserPromoted(t *testing.T) { + root := t.TempDir() + writeFile(t, root, "y.docx", "docx 內容", baseTime) + writeFile(t, root, "y.pdf", "pdf 內容", baseTime) + + m := newTestManifest() + if _, err := Scan(root, m, ScanOptions{}); err != nil { + t.Fatal(err) + } + m.Entries["y.docx"].IngestedHash = m.Entries["y.docx"].ContentHash + + // docx 被刪除 + if err := os.Remove(filepath.Join(root, "y.docx")); err != nil { + t.Fatal(err) + } + payload2, err := Scan(root, m, ScanOptions{MaxRemovedRatio: 1.0}) + if err != nil { + t.Fatal(err) + } + var addedPdf bool + for _, ev := range payload2.Events { + if ev.Type == "added" && ev.Path == "y.pdf" { + addedPdf = true + } + } + if !addedPdf { + t.Fatalf("winner 消失後 loser 應遞補產生 added 事件:%+v", payload2.Events) + } +} diff --git a/sync_status.go b/sync_status.go index f0bf0b1..4f954af 100644 --- a/sync_status.go +++ b/sync_status.go @@ -21,6 +21,16 @@ type AccountSyncStatus struct { // 用戶可能有多個實例、更新進度不同步。只在走 workers-ai 這條路時探測。 CloudAIReady bool `json:"cloud_ai_ready"` CloudAINote string `json:"cloud_ai_note,omitempty"` // 還沒通時的白話說明(含該做什麼) + + // ── 額度冷卻(2026-08-07 pacing task 2)─────────────────────────────────── + // Workers AI 每日免費額度用完時,不能每輪繼續撞同一面牆——這裡記「冷卻到什麼時候」 + // 與「今天已經做了幾份」,跨輪讀回(見 direct.go RunDirectOnce 開頭載入 prevStatus)。 + DailyIngestedDate string `json:"daily_ingested_date,omitempty"` // YYYY-MM-DD(UTC,與額度重置同一條日界線) + DailyIngestedCount int `json:"daily_ingested_count"` // 今天已成功萃取的份數 + QuotaCooldownUntil string `json:"quota_cooldown_until,omitempty"` // RFC3339;非空且未到=本帳號本輪不再嘗試萃取 + // QuotaMessage=額度冷卻中要給使用者看的三句話(見 quota.go QuotaNotice)。 + // 冷卻結束且本輪沒有新命中 ⇒ 每輪重建的 AccountSyncStatus 不會再設它,自然清除。 + QuotaMessage *QuotaNotice `json:"quota_message,omitempty"` } // SyncStatus 彙總每輪同步的萃取結果,持久化至 ~/.arcrun-rag/status.json。