| 1 | package main |
| 2 | |
| 3 | import ( |
| 4 | "context" |
| 5 | "encoding/json" |
| 6 | "fmt" |
| 7 | "testing" |
| 8 | |
| 9 | "reasonix/internal/provider" |
| 10 | "reasonix/internal/session" |
| 11 | ) |
| 12 | |
| 13 | type workspaceInfoReaderFunc func(context.Context, session.SessionRef) (session.SessionInfo, error) |
| 14 | |
| 15 | func (f workspaceInfoReaderFunc) Stat(ctx context.Context, ref session.SessionRef) (session.SessionInfo, error) { |
| 16 | return f(ctx, ref) |
| 17 | } |
| 18 | |
| 19 | func TestProjectTopicSnapshotAllowsContinuousSessionWrites(t *testing.T) { |
| 20 | for _, pinnedOnly := range []bool{false, true} { |
| 21 | t.Run(fmt.Sprintf("pinned=%v", pinnedOnly), func(t *testing.T) { |
| 22 | app, root, refs := canonicalOrganizationFixture(t, "a", "b") |
| 23 | service := app.desktopSessionService("") |
| 24 | if pinnedOnly { |
| 25 | if err := app.workspaceRegistry().UpdatePresentation(t.Context(), []string{"a", "b"}, nil, &pinnedOnly); err != nil { |
| 26 | t.Fatal(err) |
| 27 | } |
| 28 | } |
| 29 | reads := 0 |
| 30 | observed := map[string]string{} |
| 31 | reader := workspaceInfoReaderFunc(func(ctx context.Context, ref session.SessionRef) (session.SessionInfo, error) { |
| 32 | info, err := service.Query().Stat(ctx, ref) |
| 33 | if err != nil { |
| 34 | return info, err |
| 35 | } |
| 36 | reads++ |
| 37 | observed[ref.SessionID] = info.Title |
| 38 | // Force a real owner write after every metadata observation. A |
| 39 | // double collect can never settle, regardless of retry count. |
| 40 | if err := service.SetTitle(ctx, ref, fmt.Sprintf("Live write %d", reads)); err != nil { |
| 41 | t.Fatal(err) |
| 42 | } |
| 43 | return info, nil |
| 44 | }) |
| 45 | req := ProjectTopicPageRequest{Scope: "project", WorkspaceRoot: root, Limit: 1, pinnedOnly: pinnedOnly} |
| 46 | for range 3 { |
| 47 | page, err := app.readProjectTopicPage(req, reader) |
| 48 | if err != nil || len(page.Items) != 1 || page.NextCursor == "" { |
| 49 | t.Fatalf("live first page = %+v, %v", page, err) |
| 50 | } |
| 51 | row := page.Items[0] |
| 52 | if title := observed[row.Session.SessionID]; title != "" && row.Label != title { |
| 53 | t.Fatalf("page mixed captured and later metadata: %+v, observed %q", row, title) |
| 54 | } |
| 55 | continuation := req |
| 56 | continuation.Cursor = page.NextCursor |
| 57 | second, err := app.ListProjectTopics(continuation) |
| 58 | if err != nil || len(second.Items) != 1 || second.Items[0].Key == row.Key || second.SnapshotID != page.SnapshotID { |
| 59 | t.Fatalf("write disturbed frozen continuation: %+v %v", second, err) |
| 60 | } |
| 61 | } |
| 62 | if reads != 3*len(refs) { |
| 63 | t.Fatalf("metadata reads = %d, want one observation per session per page", reads) |
| 64 | } |
| 65 | fresh, err := app.ListProjectTopics(req) |
| 66 | if err != nil || len(fresh.Items) != 1 { |
| 67 | t.Fatalf("refresh = %+v, %v", fresh, err) |
| 68 | } |
| 69 | info, err := service.Query().Stat(t.Context(), *fresh.Items[0].Session) |
| 70 | if err != nil || fresh.Items[0].Label != info.Title { |
| 71 | t.Fatalf("refresh did not observe latest title: %+v, %+v, %v", fresh, info, err) |
| 72 | } |
| 73 | }) |
| 74 | } |
| 75 | } |
| 76 | |
| 77 | func TestProjectTopicCursorIgnoresUnrelatedWorkspaceMutation(t *testing.T) { |
| 78 | app, root, _ := canonicalOrganizationFixture(t, "a", "b") |
| 79 | req := ProjectTopicPageRequest{Scope: "project", WorkspaceRoot: root, Limit: 1} |
| 80 | first, err := app.ListProjectTopics(req) |
| 81 | if err != nil || first.NextCursor == "" { |
| 82 | t.Fatalf("first page = %+v, %v", first, err) |
| 83 | } |
| 84 | if _, err := app.ensureDesktopWorkspace(t.Context(), "project", t.TempDir()); err != nil { |
| 85 | t.Fatal(err) |
| 86 | } |
| 87 | req.Cursor = first.NextCursor |
| 88 | second, err := app.ListProjectTopics(req) |
| 89 | if err != nil || len(second.Items) != 1 || second.Items[0].Key == first.Items[0].Key { |
| 90 | t.Fatalf("unrelated workspace invalidated pagination: %+v, %v", second, err) |
| 91 | } |
| 92 | } |
| 93 | |
| 94 | func TestProjectTopicCursorFreezesNewResult(t *testing.T) { |
| 95 | app, root, refs := canonicalOrganizationFixture(t, "a", "b") |
| 96 | req := ProjectTopicPageRequest{Scope: "project", WorkspaceRoot: root, Limit: 1} |
| 97 | first, err := app.ListProjectTopics(req) |
| 98 | if err != nil || first.NextCursor == "" { |
| 99 | t.Fatalf("first page = %+v, %v", first, err) |
| 100 | } |
| 101 | runtime, ok := app.desktopSessionService("").Runtime(refs["a"]) |
| 102 | if !ok { |
| 103 | t.Fatal("fixture runtime missing") |
| 104 | } |
| 105 | payload, err := json.Marshal(map[string]any{"message": provider.Message{ID: "answer", Role: provider.RoleAssistant, Content: "Build finished"}}) |
| 106 | if err != nil { |
| 107 | t.Fatal(err) |
| 108 | } |
| 109 | commit, err := runtime.Session().Append(t.Context(), session.Batch{OperationID: "build", TurnID: "build", Events: []session.Event{ |
| 110 | {Kind: "turn/start"}, |
| 111 | {Kind: "message/complete", Payload: payload}, |
| 112 | {Kind: "turn/end", Payload: json.RawMessage(`{"status":"completed"}`)}, |
| 113 | }}) |
| 114 | if err != nil { |
| 115 | t.Fatal(err) |
| 116 | } |
| 117 | req.Cursor = first.NextCursor |
| 118 | second, err := app.ListProjectTopics(req) |
| 119 | if err != nil || second.SnapshotID != first.SnapshotID || len(second.Items) != 1 || second.Items[0].Key == first.Items[0].Key { |
| 120 | t.Fatalf("new result disturbed snapshot: %+v %v", second, err) |
| 121 | } |
| 122 | req.Cursor, req.Limit = "", 10 |
| 123 | fresh, err := app.ListProjectTopics(req) |
| 124 | if err != nil { |
| 125 | t.Fatal(err) |
| 126 | } |
| 127 | for _, row := range fresh.Items { |
| 128 | if row.Session.SessionID == "a" && row.ResultSequence == commit.FirstSequence+uint64(commit.EventCount)-1 { |
| 129 | return |
| 130 | } |
| 131 | } |
| 132 | t.Fatalf("refresh omitted the new result: %+v", fresh) |
| 133 | } |
| 134 | |
| 135 | func TestProjectTopicCursorIgnoresUnrelatedCatalogScan(t *testing.T) { |
| 136 | app, root, _ := canonicalOrganizationFixture(t, "a", "b") |
| 137 | otherRoot := t.TempDir() |
| 138 | installSessionCatalogForTest(t, app, otherRoot, "project", otherRoot) |
| 139 | req := ProjectTopicPageRequest{Scope: "project", WorkspaceRoot: root, Limit: 1} |
| 140 | first, err := app.ListProjectTopics(req) |
| 141 | if err != nil || first.NextCursor == "" { |
| 142 | t.Fatalf("first page = %+v, %v", first, err) |
| 143 | } |
| 144 | before := app.currentSessionCatalogStatus().Revision |
| 145 | writeTopicSession(t, otherRoot, "other.jsonl", "other-topic", "Other project", otherRoot) |
| 146 | reconcileSessionCatalogForTest(t, app, otherRoot, "project", otherRoot) |
| 147 | if app.currentSessionCatalogStatus().Revision <= before { |
| 148 | t.Fatal("fixture did not advance catalog revision") |
| 149 | } |
| 150 | req.Cursor = first.NextCursor |
| 151 | second, err := app.ListProjectTopics(req) |
| 152 | if err != nil || len(second.Items) != 1 || second.Items[0].Key == first.Items[0].Key { |
| 153 | t.Fatalf("unrelated catalog scan invalidated pagination: %+v, %v", second, err) |
| 154 | } |
| 155 | } |
| 156 | |
| 157 | func TestProjectTopicSnapshotPreservesCapturedOrganization(t *testing.T) { |
| 158 | app, root, refs := canonicalOrganizationFixture(t, "a", "b", "c") |
| 159 | keys := []string{} |
| 160 | for _, id := range []string{"a", "b", "c"} { |
| 161 | ref := refs[id] |
| 162 | keys = append(keys, projectNodeSessionKey(ProjectNode{Session: &ref})) |
| 163 | } |
| 164 | if err := app.ReorderSessions("project", root, keys); err != nil { |
| 165 | t.Fatal(err) |
| 166 | } |
| 167 | changed := false |
| 168 | reader := workspaceInfoReaderFunc(func(ctx context.Context, ref session.SessionRef) (session.SessionInfo, error) { |
| 169 | if !changed { |
| 170 | changed = true |
| 171 | if err := app.ReorderSessions("project", root, []string{keys[2], keys[0], keys[1]}); err != nil { |
| 172 | t.Fatal(err) |
| 173 | } |
| 174 | } |
| 175 | return app.desktopSessionService("").Query().Stat(ctx, ref) |
| 176 | }) |
| 177 | req := ProjectTopicPageRequest{Scope: "project", WorkspaceRoot: root, Limit: 1} |
| 178 | first, err := app.readProjectTopicPage(req, reader) |
| 179 | if err != nil || len(first.Items) != 1 || first.Items[0].Session.SessionID != "a" { |
| 180 | t.Fatalf("captured order = %+v, %v", first, err) |
| 181 | } |
| 182 | req.Cursor = first.NextCursor |
| 183 | second, err := app.ListProjectTopics(req) |
| 184 | if err != nil || len(second.Items) != 1 || second.Items[0].Session.SessionID != "b" { |
| 185 | t.Fatalf("reorder disturbed captured order: %+v %v", second, err) |
| 186 | } |
| 187 | req.Cursor = "" |
| 188 | fresh, err := app.ListProjectTopics(req) |
| 189 | if err != nil || len(fresh.Items) != 1 || fresh.Items[0].Session.SessionID != "c" { |
| 190 | t.Fatalf("refreshed order = %+v, %v", fresh, err) |
| 191 | } |
| 192 | } |
| 193 | |
| 194 | func TestProjectTopicSnapshotRejectsArchivedMember(t *testing.T) { |
| 195 | app, root, refs := canonicalOrganizationFixture(t, "a", "b", "c") |
| 196 | req := ProjectTopicPageRequest{Scope: "project", WorkspaceRoot: root, Limit: 1} |
| 197 | page, err := app.ListProjectTopics(req) |
| 198 | if err != nil { |
| 199 | t.Fatal(err) |
| 200 | } |
| 201 | if err := app.workspaceRegistry().ArchiveSession(t.Context(), refs["b"].SessionID); err != nil { |
| 202 | t.Fatal(err) |
| 203 | } |
| 204 | req.Cursor = page.NextCursor |
| 205 | if _, err := app.ListProjectTopics(req); err == nil { |
| 206 | t.Fatal("archived member accepted by frozen read") |
| 207 | } |
| 208 | req.Cursor = "" |
| 209 | req.Limit = 10 |
| 210 | fresh, err := app.ListProjectTopics(req) |
| 211 | if err != nil { |
| 212 | t.Fatal(err) |
| 213 | } |
| 214 | for _, row := range fresh.Items { |
| 215 | if row.Session != nil && row.Session.SessionID == "b" { |
| 216 | t.Fatal("archived row resurrected") |
| 217 | } |
| 218 | } |
| 219 | } |
| 220 |