Устойчивость раскладки и переходов linking (MAJOR-4, MINOR-7)
Закрывает две связанные дыры «claim-then-side-effect» в раскладке хардлинками. MINOR-7: transition глотал ошибку записи состояния — на путях Apply и авто-раскладки выполнение продолжалось к хардлинкам при незакоммиченном claim перехода в linking, а финальный linking→done отклонялся графом (задача застревала со stale-планом). Выделен transitionErr, возвращающий ошибку; Apply и finishRecognition прерываются ДО linkPlan при провале claim. Обёртка transition (void) сохранена для fire-and-forget переходов — соседние функции воркера не тронуты. MAJOR-4: (A) провал CreateFileLinks после создания хардлинков больше не оставляет задачу в linking голым return — уводим в review (код persist), повтор Apply идемпотентен. (B) новый шаг pollOnce sweepLinking возвращает осиротевшие после краха linking-задачи в review (код interrupted) на тике и старте; любая linking под w.mu устарела по построению. Восстановлен инвариант «у каждого нетерминального состояния есть владелец». Граф переходов не тронут (ребро linking→review уже объявлено). Тесты: провал claim не создаёт хардлинков; провал учёта уводит в review; sweep осиротевшего linking. OpenSpec-change linking-transition-robustness (дельты file-layout, state-reconciliation). Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -194,7 +194,13 @@ func (w *Worker) finishRecognition(ctx context.Context, id string, res recognize
|
||||
if res.Decision.Auto && !forceReview && w.layouter != nil {
|
||||
plan := applyOverrides(res.Plan, overrides)
|
||||
lctx := w.scoped(ctx, capFileLayout, id, d.PrimaryInfohash())
|
||||
w.transition(lctx, *d, store.StateLinking, "", "")
|
||||
// Claim перехода в linking должен закоммититься до хардлинков (MINOR-7):
|
||||
// при провале записи не линкуем — задача остаётся в recognizing, и
|
||||
// поллинг-цикл (recognizePending) повторит распознавание/авто-раскладку.
|
||||
if err := w.transitionErr(lctx, *d, store.StateLinking, "", ""); err != nil {
|
||||
logctx.From(lctx).Warn("auto-apply claim failed, left for recognizing", "error", err)
|
||||
return
|
||||
}
|
||||
if err := w.linkPlan(lctx, d, plan, tag, savePath); err != nil {
|
||||
logctx.From(lctx).Warn("auto-apply failed, left for review", "error", err)
|
||||
}
|
||||
@@ -253,7 +259,13 @@ func (w *Worker) Apply(ctx context.Context, id string) error {
|
||||
return fmt.Errorf("apply: торрент ещё качается: %w", ErrNotReady)
|
||||
}
|
||||
|
||||
w.transition(ctx, *d, store.StateLinking, "", "")
|
||||
// Claim перехода в linking ОБЯЗАН закоммититься до создания хардлинков: при
|
||||
// провале записи не линкуем (иначе ссылки лягут при задаче в review, а
|
||||
// финальный linking→done граф отклонит — MINOR-7). Задача остаётся в
|
||||
// review/deferred, повтор безопасен.
|
||||
if err := w.transitionErr(ctx, *d, store.StateLinking, "", ""); err != nil {
|
||||
return fmt.Errorf("apply: %w", err)
|
||||
}
|
||||
if err := w.linkPlan(ctx, d, plan, tag, translatePath(t.SavePath, w.cfg.PathMap)); err != nil {
|
||||
return fmt.Errorf("apply: %w", err)
|
||||
}
|
||||
@@ -288,6 +300,12 @@ func (w *Worker) linkPlan(ctx context.Context, d *store.Download, plan recognize
|
||||
}
|
||||
if len(fl) > 0 {
|
||||
if err := w.store.CreateFileLinks(ctx, fl); err != nil {
|
||||
// Хардлинки уже на диске, но их учёт не записан (транзиентная ошибка
|
||||
// SQLite). НЕ оставляем задачу в linking (осиротела бы до sweep, а
|
||||
// файлы висели бы без file_link — MAJOR-4): уводим в review с
|
||||
// причиной. Повторный Apply идемпотентен — Apply вернёт StatusExists
|
||||
// на уже созданных ссылках и допишет учёт.
|
||||
w.transition(ctx, *d, store.StateReview, "persist", err.Error())
|
||||
return fmt.Errorf("persist links: %w", err)
|
||||
}
|
||||
// Инвариант «один целевой путь — один владелец»: забираем владение
|
||||
|
||||
@@ -310,6 +310,10 @@ type memStore struct {
|
||||
links []store.FileLink
|
||||
candidates []store.MetadataCandidate
|
||||
torrents map[string][]byte
|
||||
|
||||
// Инъекция сбоев (для тестов устойчивости раскладки).
|
||||
failCreateLinks error // CreateFileLinks вернёт эту ошибку
|
||||
failSetState func(store.State) error // SetDownloadState вернёт ошибку для перехода
|
||||
}
|
||||
|
||||
func newMemStore() *memStore {
|
||||
@@ -436,6 +440,11 @@ func (m *memStore) GetDownload(_ context.Context, id string) (*store.Download, e
|
||||
}
|
||||
|
||||
func (m *memStore) SetDownloadState(_ context.Context, id string, st store.State, code, msg string) error {
|
||||
if m.failSetState != nil {
|
||||
if err := m.failSetState(st); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
d := m.downloads[id]
|
||||
d.State = st
|
||||
d.ErrorCode = store.NullString(code)
|
||||
@@ -523,6 +532,9 @@ func (m *memStore) ListOverrides(_ context.Context, id string) (map[string]strin
|
||||
}
|
||||
|
||||
func (m *memStore) CreateFileLinks(_ context.Context, links []store.FileLink) error {
|
||||
if m.failCreateLinks != nil {
|
||||
return m.failCreateLinks
|
||||
}
|
||||
m.links = append(m.links, links...)
|
||||
return nil
|
||||
}
|
||||
@@ -1547,3 +1559,103 @@ func TestToLayoutPlan_SrcPrefixIsSavePath(t *testing.T) {
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// --- Устойчивость раскладки и переходов (MAJOR-4, MINOR-7) ---
|
||||
|
||||
// TestApply_ClaimFailureAbortsBeforeLinks (MINOR-7): если запись claim перехода
|
||||
// в linking падает, Apply ОБЯЗАН прерваться ДО создания хардлинков — иначе
|
||||
// ссылки лягут при задаче в review, а финальный linking→done граф отклонит.
|
||||
func TestApply_ClaimFailureAbortsBeforeLinks(t *testing.T) {
|
||||
f := newApplyFixture(t, seriesResult().Plan)
|
||||
f.st.failSetState = func(st store.State) error {
|
||||
if st == store.StateLinking {
|
||||
return errors.New("boom: claim persist failed")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
err := f.w.Apply(context.Background(), "1")
|
||||
if err == nil {
|
||||
t.Fatal("Apply must fail when linking claim persist fails")
|
||||
}
|
||||
if f.st.downloads["1"].State != store.StateReview {
|
||||
t.Errorf("state = %q, want review (claim not committed)", f.st.downloads["1"].State)
|
||||
}
|
||||
if len(f.st.links) != 0 {
|
||||
t.Errorf("file_links = %d, want 0 (no linking before committed claim)", len(f.st.links))
|
||||
}
|
||||
// Хардлинки на диск НЕ созданы — раскладка не запускалась.
|
||||
dst := filepath.Join(f.series, "Show (2006)", "Season 02", "Show (2006) S02E01.mkv")
|
||||
if _, statErr := os.Stat(dst); statErr == nil {
|
||||
t.Errorf("hardlink %q created despite failed claim", dst)
|
||||
}
|
||||
}
|
||||
|
||||
// TestApply_PersistFailureLeavesReview (MAJOR-4 A): хардлинки созданы, но
|
||||
// CreateFileLinks упал транзиентно → задача не должна застрять в linking; уходит
|
||||
// в review с причиной, повтор идемпотентен.
|
||||
func TestApply_PersistFailureLeavesReview(t *testing.T) {
|
||||
f := newApplyFixture(t, seriesResult().Plan)
|
||||
f.st.failCreateLinks = errors.New("boom: sqlite busy")
|
||||
|
||||
err := f.w.Apply(context.Background(), "1")
|
||||
if err == nil {
|
||||
t.Fatal("Apply must fail when persisting links fails")
|
||||
}
|
||||
if f.st.downloads["1"].State != store.StateReview {
|
||||
t.Fatalf("state = %q, want review (not stranded in linking)", f.st.downloads["1"].State)
|
||||
}
|
||||
if f.st.downloads["1"].ErrorCode.String != "persist" {
|
||||
t.Errorf("error_code = %q, want persist", f.st.downloads["1"].ErrorCode.String)
|
||||
}
|
||||
// Хардлинки уже на диске (учёт лишь не записан) — повторный Apply их допишет.
|
||||
dst := filepath.Join(f.series, "Show (2006)", "Season 02", "Show (2006) S02E01.mkv")
|
||||
if _, statErr := os.Stat(dst); statErr != nil {
|
||||
t.Errorf("expected hardlink on disk despite persist failure: %v", statErr)
|
||||
}
|
||||
}
|
||||
|
||||
// TestSweepLinking_OrphanedToReview (MAJOR-4 B): задача, застрявшая в linking
|
||||
// после краха, на тике/старте возвращается в review с причиной.
|
||||
func TestSweepLinking_OrphanedToReview(t *testing.T) {
|
||||
st := newMemStore()
|
||||
d := completedDownload("1")
|
||||
d.State = store.StateLinking
|
||||
st.put(d)
|
||||
w := testWorkerWith(st, &fakeQbt{}, &fakeRecognizer{}, nil)
|
||||
|
||||
w.sweepLinking(context.Background())
|
||||
|
||||
got := st.downloads["1"]
|
||||
if got.State != store.StateReview {
|
||||
t.Fatalf("state = %q, want review", got.State)
|
||||
}
|
||||
if got.ErrorCode.String != "interrupted" {
|
||||
t.Errorf("error_code = %q, want interrupted", got.ErrorCode.String)
|
||||
}
|
||||
if got.ErrorMsg.String == "" {
|
||||
t.Error("expected error_msg with reason")
|
||||
}
|
||||
}
|
||||
|
||||
// TestSweepLinking_LeavesOtherStates: sweep трогает только linking, прочие
|
||||
// состояния (в т.ч. done) не задевает.
|
||||
func TestSweepLinking_LeavesOtherStates(t *testing.T) {
|
||||
st := newMemStore()
|
||||
done := completedDownload("1")
|
||||
done.State = store.StateDone
|
||||
st.put(done)
|
||||
review := completedDownload("2")
|
||||
review.State = store.StateReview
|
||||
st.put(review)
|
||||
w := testWorkerWith(st, &fakeQbt{}, &fakeRecognizer{}, nil)
|
||||
|
||||
w.sweepLinking(context.Background())
|
||||
|
||||
if st.downloads["1"].State != store.StateDone {
|
||||
t.Errorf("done task moved to %q", st.downloads["1"].State)
|
||||
}
|
||||
if st.downloads["2"].State != store.StateReview {
|
||||
t.Errorf("review task moved to %q", st.downloads["2"].State)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -322,12 +322,39 @@ func (w *Worker) pollOnce(ctx context.Context) {
|
||||
// Быстрый приём отложил добавление в qBittorrent: подхватываем пойманные
|
||||
// (catched) загрузки и добавляем их (сеть — вне блокировки переходов).
|
||||
w.processCatched(ctx)
|
||||
// Восстанавливаем задачи, застрявшие в linking после краха между claim и
|
||||
// финальным переходом (иначе их не листит никто — вечный лимбо).
|
||||
w.sweepLinking(ctx)
|
||||
// Ф3: распознаём завершённые загрузки (и перезапускаем по подсказке).
|
||||
if w.recognizer != nil {
|
||||
w.recognizePending(ctx)
|
||||
}
|
||||
}
|
||||
|
||||
// sweepLinking восстанавливает задачи, застрявшие в состоянии linking. Любая
|
||||
// linking-задача, видимая под w.mu, устарела по построению: активная раскладка
|
||||
// (linkPlan) держит w.mu на всё время и завершает переход из linking ДО отпускания
|
||||
// замка — значит эта задача осталась в linking после краха процесса между claim
|
||||
// (переходом в linking) и финальным переходом. Возвращаем её в review с причиной;
|
||||
// человек повторит Apply (linkPlan идемпотентен), и незаписанный учёт хардлинков
|
||||
// допишется. Так у linking появляется владелец на рестарте/тике — инвариант «у
|
||||
// каждого нетерминального состояния есть владелец» (как recognizePending для
|
||||
// recognizing). Выполняется на каждом тике и на старте (первый pollOnce до цикла).
|
||||
func (w *Worker) sweepLinking(ctx context.Context) {
|
||||
w.mu.Lock()
|
||||
defer w.mu.Unlock()
|
||||
stuck, err := w.store.ListDownloadsByState(ctx, store.StateLinking)
|
||||
if err != nil {
|
||||
w.log.Warn("sweep linking list failed", "capability", capFileLayout, "error", err)
|
||||
return
|
||||
}
|
||||
for _, d := range stuck {
|
||||
lctx := w.scoped(ctx, capFileLayout, d.ID, d.PrimaryInfohash())
|
||||
w.transition(lctx, d, store.StateReview, "interrupted",
|
||||
"прерванная раскладка, повтори применение")
|
||||
}
|
||||
}
|
||||
|
||||
// processCatched — асинхронный шаг добавления пойманных загрузок в qBittorrent.
|
||||
// Для каждой catched: (предохранитель) если висит дольше catch_timeout — уводим
|
||||
// в failed; иначе выводим имя и добавляем в qBit. Медленные вызовы (LLM-namer,
|
||||
@@ -641,14 +668,28 @@ func (w *Worker) retriedFloor(d store.Download, basis time.Time) time.Time {
|
||||
return basis
|
||||
}
|
||||
|
||||
// transition пишет новое состояние и логирует переход.
|
||||
// transition пишет новое состояние и логирует переход. Fire-and-forget обёртка
|
||||
// над transitionErr: применяется там, где переход терминален для шага — за ним
|
||||
// нет побочного эффекта, зависящего от факта записи claim (reconcile, таймауты,
|
||||
// команды ревью, финальные переходы linkPlan, sweep). Ошибку записи гасит (её
|
||||
// уже залогировал transitionErr).
|
||||
func (w *Worker) transition(ctx context.Context, d store.Download, state store.State, code, msg string) {
|
||||
_ = w.transitionErr(ctx, d, state, code, msg)
|
||||
}
|
||||
|
||||
// transitionErr пишет новое состояние, шлёт пинги/скан, логирует переход и
|
||||
// ВОЗВРАЩАЕТ ошибку записи. На claim-then-side-effect путях (ручное Apply,
|
||||
// авто-раскладка в finishRecognition) провал claim перехода в `linking` ОБЯЗАН
|
||||
// прервать выполнение ДО побочных эффектов (хардлинков): иначе ссылки лягут при
|
||||
// незакоммиченном claim, а финальный переход из фактического (не `linking`)
|
||||
// состояния граф отклонит — задача застрянет со stale-планом (MINOR-7).
|
||||
func (w *Worker) transitionErr(ctx context.Context, d store.Download, state store.State, code, msg string) error {
|
||||
// FromOr, а не From: если вызывающий не завёл scoped-логгер, падаем на
|
||||
// w.log (настроенный), а не на slog.Default().
|
||||
log := logctx.FromOr(ctx, w.log)
|
||||
if err := w.store.SetDownloadState(ctx, d.ID, state, code, msg); err != nil {
|
||||
log.Error("state transition failed", "from", d.State, "to", state, "error", err)
|
||||
return
|
||||
return fmt.Errorf("transition %s → %s: %w", d.State, state, err)
|
||||
}
|
||||
log.Info("state transition", "from", d.State, "to", state, "code", code)
|
||||
|
||||
@@ -681,6 +722,7 @@ func (w *Worker) transition(ctx context.Context, d store.Download, state store.S
|
||||
gctx := w.scoped(context.Background(), capFileLayout, d.ID, d.PrimaryInfohash())
|
||||
go func() { _ = w.scanner.RefreshLibraries(gctx) }()
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// shouldNotifyFail дебаунсит повторные уведомления о падении одной задачи
|
||||
|
||||
Reference in New Issue
Block a user