| 1 | package main |
| 2 | |
| 3 | import ( |
| 4 | "bytes" |
| 5 | "crypto/sha256" |
| 6 | "encoding/hex" |
| 7 | "encoding/json" |
| 8 | "fmt" |
| 9 | "io/fs" |
| 10 | "maps" |
| 11 | "os" |
| 12 | "path/filepath" |
| 13 | "reflect" |
| 14 | "testing" |
| 15 | |
| 16 | "reasonix/internal/config" |
| 17 | "reasonix/internal/provider" |
| 18 | "reasonix/internal/session" |
| 19 | ) |
| 20 | |
| 21 | // These files were emitted by the actual tagged writers, including their event |
| 22 | // logs, digests and external content objects. SOURCE.json pins the provenance; |
| 23 | // no fixture is constructed using the current writer's idea of an old format. |
| 24 | func taggedHistoryFixture(t *testing.T, version string) string { |
| 25 | t.Helper() |
| 26 | root := filepath.Join("testdata", "session-history-upgrade", version) |
| 27 | var source struct { |
| 28 | Tag, Commit string |
| 29 | Files map[string]string |
| 30 | } |
| 31 | readTaggedJSON(t, filepath.Join(root, "SOURCE.json"), &source) |
| 32 | if source.Tag != "desktop-v"+version || len(source.Commit) != 40 || len(source.Files) == 0 { |
| 33 | t.Fatal("missing historical source provenance") |
| 34 | } |
| 35 | for name, digest := range source.Files { |
| 36 | body, err := os.ReadFile(filepath.Join(root, name)) |
| 37 | if err != nil { |
| 38 | t.Fatal(err) |
| 39 | } |
| 40 | sum := sha256.Sum256(body) |
| 41 | if hex.EncodeToString(sum[:]) != digest { |
| 42 | t.Fatalf("historical fixture changed: %s", name) |
| 43 | } |
| 44 | } |
| 45 | return root |
| 46 | } |
| 47 | |
| 48 | func readTaggedJSON(t *testing.T, path string, value any) { |
| 49 | t.Helper() |
| 50 | body, err := os.ReadFile(path) |
| 51 | if err != nil { |
| 52 | t.Fatal(err) |
| 53 | } |
| 54 | if err := json.Unmarshal(body, value); err != nil { |
| 55 | t.Fatal(err) |
| 56 | } |
| 57 | } |
| 58 | |
| 59 | func copyTaggedDirectory(t *testing.T, source, destination string) map[string][]byte { |
| 60 | t.Helper() |
| 61 | before := map[string][]byte{} |
| 62 | err := filepath.WalkDir(source, func(path string, entry fs.DirEntry, err error) error { |
| 63 | if err != nil { |
| 64 | return err |
| 65 | } |
| 66 | relative, err := filepath.Rel(source, path) |
| 67 | if err != nil { |
| 68 | return err |
| 69 | } |
| 70 | target := filepath.Join(destination, relative) |
| 71 | if entry.IsDir() { |
| 72 | return os.MkdirAll(target, 0700) |
| 73 | } |
| 74 | body, err := os.ReadFile(path) |
| 75 | if err != nil { |
| 76 | return err |
| 77 | } |
| 78 | before[target] = body |
| 79 | return os.WriteFile(target, body, 0600) |
| 80 | }) |
| 81 | if err != nil { |
| 82 | t.Fatal(err) |
| 83 | } |
| 84 | return before |
| 85 | } |
| 86 | |
| 87 | func assertTaggedMessages(t *testing.T, got, want []provider.Message) { |
| 88 | t.Helper() |
| 89 | if len(got) != len(want) { |
| 90 | t.Fatalf("historical messages: got %d, want %d", len(got), len(want)) |
| 91 | } |
| 92 | for i := range want { |
| 93 | if !reflect.DeepEqual(got[i], want[i]) { |
| 94 | // Do not print the large tool output on failure. |
| 95 | g, _ := json.Marshal(got[i]) |
| 96 | w, _ := json.Marshal(want[i]) |
| 97 | t.Fatalf("message %d (%s) changed: got sha256=%x want=%x", i, want[i].ID, sha256.Sum256(g), sha256.Sum256(w)) |
| 98 | } |
| 99 | } |
| 100 | } |
| 101 | |
| 102 | func assertTaggedPages(t *testing.T, app *App, ref session.SessionRef, want []provider.Message) { |
| 103 | t.Helper() |
| 104 | query := app.desktopSessionService("").Query() |
| 105 | // Await the public synchronous projection query. This matrix validates |
| 106 | // durable reachability; asynchronous cold-read latency has separate tests. |
| 107 | if _, err := query.HistoryShape(t.Context(), ref); err != nil { |
| 108 | t.Fatal(err) |
| 109 | } |
| 110 | request := session.HistoryWindowRequest{Anchor: "newest", Limit: 17} |
| 111 | seen := map[string]bool{} |
| 112 | for pages := 0; ; pages++ { |
| 113 | if pages > len(want) { |
| 114 | t.Fatal("history cursor failed to advance") |
| 115 | } |
| 116 | page, err := query.ReadHistoryWindow(t.Context(), ref, request) |
| 117 | if err != nil || page.Status != "ready" || len(page.Messages) > 17 { |
| 118 | t.Fatalf("historical page: %s %v", page.Status, err) |
| 119 | } |
| 120 | for _, message := range page.Messages { |
| 121 | if seen[message.MessageID] { |
| 122 | t.Fatalf("duplicate history identity %s", message.MessageID) |
| 123 | } |
| 124 | seen[message.MessageID] = true |
| 125 | } |
| 126 | if !page.HasOlder { |
| 127 | break |
| 128 | } |
| 129 | if page.OlderCursor == "" { |
| 130 | t.Fatal("unreachable older history") |
| 131 | } |
| 132 | request = session.HistoryWindowRequest{Anchor: "cursor", Cursor: page.OlderCursor, Limit: 17} |
| 133 | } |
| 134 | for _, message := range want { |
| 135 | if message.Role != provider.RoleSystem && !seen[message.ID] { |
| 136 | t.Fatalf("historical message unreachable: %s", message.ID) |
| 137 | } |
| 138 | } |
| 139 | } |
| 140 | |
| 141 | func TestTaggedHistory1383To1385(t *testing.T) { taggedHistoryReleaseUpgradeMatrix(t, 3, 5) } |
| 142 | func TestTaggedHistory1386To1387(t *testing.T) { taggedHistoryReleaseUpgradeMatrix(t, 6, 7) } |
| 143 | func TestTaggedHistory1388To1389(t *testing.T) { taggedHistoryReleaseUpgradeMatrix(t, 8, 9) } |
| 144 | func TestTaggedHistory13810To13811(t *testing.T) { taggedHistoryReleaseUpgradeMatrix(t, 10, 11) } |
| 145 | |
| 146 | func taggedHistoryReleaseUpgradeMatrix(t *testing.T, first, last int) { |
| 147 | t.Helper() |
| 148 | for release := first; release <= last; release++ { |
| 149 | version := fmt.Sprintf("1.38.%d", release) |
| 150 | for _, scope := range []string{"global", "project"} { |
| 151 | t.Run(version+"/"+scope, func(t *testing.T) { |
| 152 | isolateDesktopUserDirs(t) |
| 153 | fixture := taggedHistoryFixture(t, version) |
| 154 | legacyRoot, canonicalRoot := config.SessionDir(), config.SessionStoreDir() |
| 155 | if scope == "project" { |
| 156 | workspace := filepath.Join(t.TempDir(), "历史项目 # % 中文") |
| 157 | if err := os.MkdirAll(workspace, 0700); err != nil { |
| 158 | t.Fatal(err) |
| 159 | } |
| 160 | legacyRoot, canonicalRoot = desktopSessionDir(workspace), config.ProjectSessionStoreDir(workspace) |
| 161 | if err := saveProjectsFile(desktopProjectFile{Projects: []desktopProject{{Root: workspace}}}); err != nil { |
| 162 | t.Fatal(err) |
| 163 | } |
| 164 | } |
| 165 | before := copyTaggedDirectory(t, filepath.Join(fixture, "legacy"), legacyRoot) |
| 166 | if release >= 8 { |
| 167 | maps.Copy(before, copyTaggedDirectory(t, filepath.Join(fixture, "canonical"), canonicalRoot)) |
| 168 | } |
| 169 | var identities map[string]string |
| 170 | for restart := range 3 { |
| 171 | app := NewApp() |
| 172 | t.Cleanup(app.closeSessionServices) |
| 173 | t.Cleanup(func() { _ = app.sessionUIStore().Close() }) |
| 174 | if err := app.migrateDesktopSessionsV5(t.Context()); err != nil { |
| 175 | t.Fatal(err) |
| 176 | } |
| 177 | ledger, err := readDesktopMigrationLedger() |
| 178 | if err != nil { |
| 179 | t.Fatal(err) |
| 180 | } |
| 181 | current := map[string]string{} |
| 182 | for _, name := range []string{"complete", "branch", "interrupted", "canonical-complete", "canonical-compacted", "canonical-interrupted"} { |
| 183 | canonical := len(name) > 10 && name[:10] == "canonical-" |
| 184 | if canonical && release < 8 { |
| 185 | continue |
| 186 | } |
| 187 | key := desktopLegacyMigrationKey(filepath.Join(legacyRoot, name+".jsonl")) |
| 188 | expected := name + "-expected.json" |
| 189 | if canonical { |
| 190 | key, expected = desktopCanonicalMigrationKey(canonicalRoot, name), "canonical-expected.json" |
| 191 | } |
| 192 | record := ledger.Records[key] |
| 193 | if record.Status != "completed" || record.TargetSessionID == "" { |
| 194 | t.Fatalf("%s failed adoption: %+v", name, record) |
| 195 | } |
| 196 | current[name] = record.TargetSessionID |
| 197 | ref := session.SessionRef{HostID: localDesktopHostID, SessionID: record.TargetSessionID} |
| 198 | snapshot, err := app.desktopSessionService("").Query().Snapshot(t.Context(), ref) |
| 199 | if err != nil { |
| 200 | t.Fatal(err) |
| 201 | } |
| 202 | var messages []provider.Message |
| 203 | readTaggedJSON(t, filepath.Join(fixture, expected), &messages) |
| 204 | assertTaggedMessages(t, snapshot.Projection.Messages, messages) |
| 205 | if canonical && !bytes.Contains(snapshot.Projection.GoalState, []byte(`"futureGoal":{"keep":true}`)) { |
| 206 | t.Fatal("canonical goal state lost") |
| 207 | } |
| 208 | if name == "canonical-compacted" && (len(snapshot.Projection.ModelMessages) != 1 || snapshot.Projection.ModelMessages[0].Content != "Persisted compaction summary") { |
| 209 | t.Fatal("compaction replaced by full history on upgrade") |
| 210 | } |
| 211 | assertTaggedPages(t, app, ref, messages) |
| 212 | input, err := app.GetSessionComposerState(ref) |
| 213 | if err != nil || input.HistoryChanged { |
| 214 | t.Fatalf("input recovery incorrectly blocked: %+v %v", input, err) |
| 215 | } |
| 216 | if restart == 0 { |
| 217 | _, err = app.SaveSessionComposerState(SessionComposerSaveRequest{Ref: ref, ExpectedRevision: input.Revision, ContentVersion: 1, ContentJSON: `{"text":"new unsent input on historical session"}`}) |
| 218 | if err != nil { |
| 219 | t.Fatal(err) |
| 220 | } |
| 221 | } else if input.ContentJSON != `{"text":"new unsent input on historical session"}` { |
| 222 | t.Fatal("historical session input lost on restart") |
| 223 | } |
| 224 | } |
| 225 | if restart > 0 && !reflect.DeepEqual(identities, current) { |
| 226 | t.Fatal("restart remapped historical sessions") |
| 227 | } |
| 228 | identities = current |
| 229 | assertMigrationSourceSnapshot(t, before) |
| 230 | app.closeSessionServices() |
| 231 | if err := app.sessionUIStore().Close(); err != nil { |
| 232 | t.Fatal(err) |
| 233 | } |
| 234 | } |
| 235 | }) |
| 236 | } |
| 237 | } |
| 238 | } |
| 239 |