| 1 | package main |
| 2 | |
| 3 | import ( |
| 4 | "context" |
| 5 | "os" |
| 6 | "path/filepath" |
| 7 | "strings" |
| 8 | "testing" |
| 9 | "time" |
| 10 | |
| 11 | "reasonix/internal/agent" |
| 12 | "reasonix/internal/config" |
| 13 | "reasonix/internal/provider" |
| 14 | "reasonix/internal/sessioncatalog" |
| 15 | ) |
| 16 | |
| 17 | const recoveryChainLength = 9 |
| 18 | |
| 19 | func countDirs(t *testing.T, dir string) int { |
| 20 | t.Helper() |
| 21 | entries, err := os.ReadDir(dir) |
| 22 | if err != nil && !os.IsNotExist(err) { |
| 23 | t.Fatal(err) |
| 24 | } |
| 25 | n := 0 |
| 26 | for _, e := range entries { |
| 27 | if e.IsDir() && !strings.HasPrefix(e.Name(), ".") { |
| 28 | n++ |
| 29 | } |
| 30 | } |
| 31 | return n |
| 32 | } |
| 33 | |
| 34 | // Recovery snapshots written by older builds chain each ParentID to the |
| 35 | // previous snapshot, so one action must follow the chain, not only the |
| 36 | // snapshots that share a parent. |
| 37 | func seedChainedRecoveryApp(t *testing.T) (*App, func(string) []ProjectNode, []ProjectNode) { |
| 38 | t.Helper() |
| 39 | isolateDesktopUserDirs(t) |
| 40 | dir := config.SessionDir() |
| 41 | parent := filepath.Join(dir, "20260803-140947.000000000-fake-model.jsonl") |
| 42 | p := agent.NewSession("system") |
| 43 | p.Add(provider.Message{ID: "q0", Role: provider.RoleUser, Content: "long running conversation"}) |
| 44 | if err := p.Save(parent); err != nil { |
| 45 | t.Fatal(err) |
| 46 | } |
| 47 | prev := parent |
| 48 | for i := range recoveryChainLength { |
| 49 | s := agent.NewSession("system") |
| 50 | s.Add(provider.Message{ID: "q0", Role: provider.RoleUser, Content: "long running conversation"}) |
| 51 | s.Add(provider.Message{ID: "a", Role: provider.RoleAssistant, Content: strings.Repeat("x", i+1)}) |
| 52 | meta := agent.BranchMeta{Scope: "global", TopicID: "topic_chain", TopicTitle: "long running conversation", ParentID: agent.BranchID(prev)} |
| 53 | info, err := s.SaveConflictRecoveryBranch(agent.RecoveryBranchOptions{OriginalPath: prev, Reason: "conflict", BranchMeta: meta}) |
| 54 | if err != nil { |
| 55 | t.Fatal(err) |
| 56 | } |
| 57 | prev = info.Path |
| 58 | } |
| 59 | _ = os.Remove(parent) |
| 60 | other := filepath.Join(dir, "20260826-120000.000000000-fake-model.jsonl") |
| 61 | o := agent.NewSession("system") |
| 62 | o.Add(provider.Message{ID: "q0", Role: provider.RoleUser, Content: "another conversation"}) |
| 63 | if err := o.Save(other); err != nil { |
| 64 | t.Fatal(err) |
| 65 | } |
| 66 | if _, err := o.SaveConflictRecoveryBranch(agent.RecoveryBranchOptions{OriginalPath: other, Reason: "conflict", |
| 67 | BranchMeta: agent.BranchMeta{Scope: "global", TopicID: "topic_other", TopicTitle: "another conversation"}}); err != nil { |
| 68 | t.Fatal(err) |
| 69 | } |
| 70 | |
| 71 | app := NewApp() |
| 72 | app.ctx = t.Context() |
| 73 | pinDesktopSessionRoot(t, app) |
| 74 | installNoopRuntimeEvents(app) |
| 75 | t.Cleanup(app.closeSessionServices) |
| 76 | catalog, err := sessioncatalog.Open(context.Background(), sessioncatalog.Options{InMemory: true, DisableRepair: true, MetadataOnly: true}) |
| 77 | if err != nil { |
| 78 | t.Fatal(err) |
| 79 | } |
| 80 | app.sessionCatalog.Store(catalog) |
| 81 | t.Cleanup(func() { app.stopSessionCatalog(time.Second) }) |
| 82 | rows := func(topic string) []ProjectNode { |
| 83 | t.Helper() |
| 84 | if err := catalog.ReconcileDirectory(context.Background(), sessioncatalog.DirectoryTarget{Path: dir, Scope: "global"}); err != nil { |
| 85 | t.Fatal(err) |
| 86 | } |
| 87 | page, err := app.ListProjectTopics(ProjectTopicPageRequest{Scope: "global", Limit: 50}) |
| 88 | if err != nil { |
| 89 | t.Fatal(err) |
| 90 | } |
| 91 | app.ReleaseReadSnapshot(page.SnapshotID) |
| 92 | var out []ProjectNode |
| 93 | for _, n := range page.Items { |
| 94 | if n.Recovered && n.TopicID == topic { |
| 95 | out = append(out, n) |
| 96 | } |
| 97 | } |
| 98 | return out |
| 99 | } |
| 100 | before := rows("topic_chain") |
| 101 | if len(before) != recoveryChainLength { |
| 102 | t.Fatalf("recovered rows before archive = %d, want %d", len(before), recoveryChainLength) |
| 103 | } |
| 104 | return app, rows, before |
| 105 | } |
| 106 | |
| 107 | func archiveRow(t *testing.T, app *App, n ProjectNode) SessionMutationResult { |
| 108 | t.Helper() |
| 109 | res, err := app.ArchiveSessionTarget(SessionSelector{Ref: n.Session, Source: n.Source, SessionPath: n.SessionPath}) |
| 110 | if err != nil { |
| 111 | t.Fatal(err) |
| 112 | } |
| 113 | return res |
| 114 | } |
| 115 | |
| 116 | func TestArchivingChainedRecoveryRowArchivesWholeLineageOnce(t *testing.T) { |
| 117 | app, rows, before := seedChainedRecoveryApp(t) |
| 118 | archive := func(n ProjectNode) { |
| 119 | t.Helper() |
| 120 | if _, err := app.ArchiveSessionTarget(SessionSelector{Ref: n.Session, Source: n.Source, SessionPath: n.SessionPath}); err != nil { |
| 121 | t.Fatal(err) |
| 122 | } |
| 123 | } |
| 124 | archive(before[0]) |
| 125 | if left := rows("topic_chain"); len(left) != 0 { |
| 126 | t.Fatalf("one archive action left %d of %d chained recovery snapshots listed", len(left), recoveryChainLength) |
| 127 | } |
| 128 | store := app.desktopSessions.root |
| 129 | imported := countDirs(t, store) |
| 130 | if imported != recoveryChainLength { |
| 131 | t.Fatalf("archived copies = %d, want one per snapshot (%d)", imported, recoveryChainLength) |
| 132 | } |
| 133 | archive(before[len(before)-1]) |
| 134 | if got := countDirs(t, store); got != imported { |
| 135 | t.Fatalf("repeating the action wrote %d new session directories", got-imported) |
| 136 | } |
| 137 | if kept := rows("topic_other"); len(kept) != 1 { |
| 138 | t.Fatalf("an unrelated recovered conversation was archived: %d rows left, want 1", len(kept)) |
| 139 | } |
| 140 | } |
| 141 | |
| 142 | func TestArchivingRecoveryLineageReleasesRuntimeLockBetweenSnapshots(t *testing.T) { |
| 143 | app, rows, before := seedChainedRecoveryApp(t) |
| 144 | attempts, held := 0, 0 |
| 145 | app.runtimeMutationBeforeLockHook = func(string) { |
| 146 | attempts++ |
| 147 | if !app.runtimeAdmissionMu.TryLock() { |
| 148 | held++ |
| 149 | return |
| 150 | } |
| 151 | app.runtimeAdmissionMu.Unlock() |
| 152 | } |
| 153 | archiveRow(t, app, before[0]) |
| 154 | if attempts != recoveryChainLength { |
| 155 | t.Fatalf("lock taken %d times, want once per snapshot (%d)", attempts, recoveryChainLength) |
| 156 | } |
| 157 | if held != 0 { |
| 158 | t.Fatalf("runtime lock stayed held between %d snapshots", held) |
| 159 | } |
| 160 | if left := rows("topic_chain"); len(left) != 0 { |
| 161 | t.Fatalf("%d snapshots left listed", len(left)) |
| 162 | } |
| 163 | } |
| 164 | |
| 165 | func TestArchivingRecoveryLineageReportsPartialAndResumes(t *testing.T) { |
| 166 | app, rows, before := seedChainedRecoveryApp(t) |
| 167 | calls := 0 |
| 168 | app.runtimeMutationBeforeLockHook = func(string) { |
| 169 | calls++ |
| 170 | switch calls { |
| 171 | case 4: |
| 172 | app.runtimeAdmissionMu.Lock() |
| 173 | case 5: |
| 174 | app.runtimeAdmissionMu.Unlock() |
| 175 | } |
| 176 | } |
| 177 | res := archiveRow(t, app, before[0]) |
| 178 | if !res.Committed || res.Outcome != sessionOutcomeArchivedPartial || res.PendingSiblings != 1 { |
| 179 | t.Fatalf("receipt = %+v, want committed archived_partial with 1 pending", res) |
| 180 | } |
| 181 | if left := rows("topic_chain"); len(left) != 1 { |
| 182 | t.Fatalf("%d snapshots left after partial run, want 1", len(left)) |
| 183 | } |
| 184 | app.runtimeMutationBeforeLockHook = nil |
| 185 | again := archiveRow(t, app, rows("topic_chain")[0]) |
| 186 | if again.Outcome == sessionOutcomeArchivedPartial { |
| 187 | t.Fatalf("retry still partial: %+v", again) |
| 188 | } |
| 189 | if left := rows("topic_chain"); len(left) != 0 { |
| 190 | t.Fatalf("retry left %d snapshots", len(left)) |
| 191 | } |
| 192 | if got := countDirs(t, app.desktopSessions.root); got != recoveryChainLength { |
| 193 | t.Fatalf("archived copies = %d, want %d", got, recoveryChainLength) |
| 194 | } |
| 195 | } |
| 196 |