| 1 | package workspacestate |
| 2 | |
| 3 | import ( |
| 4 | "context" |
| 5 | "encoding/json" |
| 6 | "errors" |
| 7 | "fmt" |
| 8 | "os" |
| 9 | "path/filepath" |
| 10 | "runtime" |
| 11 | "testing" |
| 12 | ) |
| 13 | |
| 14 | func TestReadSnapshotOverlappingValidationSharesFlightAndCancellation(t *testing.T) { |
| 15 | store := NewStore(filepath.Join(t.TempDir(), "registry.json")) |
| 16 | writeReadSnapshotFixture(t, store.Path(), newState()) |
| 17 | type result struct { |
| 18 | snapshot *ReadSnapshot |
| 19 | err error |
| 20 | } |
| 21 | results := make(chan result, 8) |
| 22 | canceled, cancel := context.WithCancel(t.Context()) |
| 23 | defer cancel() |
| 24 | func() { |
| 25 | // Hold the authoritative read boundary so all callers deterministically |
| 26 | // join the same verification before any file can be read. |
| 27 | store.mu.Lock() |
| 28 | defer store.mu.Unlock() |
| 29 | for i := range 8 { |
| 30 | ctx := t.Context() |
| 31 | if i == 0 { |
| 32 | ctx = canceled |
| 33 | } |
| 34 | go func() { |
| 35 | snapshot, err := store.VerifySnapshot(ctx) |
| 36 | results <- result{snapshot, err} |
| 37 | }() |
| 38 | } |
| 39 | for { |
| 40 | store.verificationMu.Lock() |
| 41 | joined := store.verification != nil && store.verification.readers == 8 |
| 42 | store.verificationMu.Unlock() |
| 43 | if joined { |
| 44 | break |
| 45 | } |
| 46 | if err := t.Context().Err(); err != nil { |
| 47 | t.Fatal(err) |
| 48 | } |
| 49 | runtime.Gosched() |
| 50 | } |
| 51 | cancel() |
| 52 | }() |
| 53 | var shared *ReadSnapshot |
| 54 | cancellations := 0 |
| 55 | for range 8 { |
| 56 | r := <-results |
| 57 | if errors.Is(r.err, context.Canceled) { |
| 58 | cancellations++ |
| 59 | continue |
| 60 | } |
| 61 | if r.err != nil { |
| 62 | t.Fatal(r.err) |
| 63 | } |
| 64 | if shared != nil && shared != r.snapshot { |
| 65 | t.Fatal("overlapping validation did not share a snapshot") |
| 66 | } |
| 67 | shared = r.snapshot |
| 68 | } |
| 69 | if cancellations != 1 || shared == nil { |
| 70 | t.Fatalf("cancellations=%d snapshot=%v", cancellations, shared) |
| 71 | } |
| 72 | } |
| 73 | |
| 74 | func writeReadSnapshotFixture(t testing.TB, path string, state State) { |
| 75 | t.Helper() |
| 76 | normalize(&state) |
| 77 | body, err := json.Marshal(state) |
| 78 | if err != nil { |
| 79 | t.Fatal(err) |
| 80 | } |
| 81 | if err := os.WriteFile(path, body, 0600); err != nil { |
| 82 | t.Fatal(err) |
| 83 | } |
| 84 | } |
| 85 | |
| 86 | func TestReadSnapshotKeepsIdentityAndDetectsLegacyReplacement(t *testing.T) { |
| 87 | path := filepath.Join(t.TempDir(), "registry.json") |
| 88 | state := newState() |
| 89 | state.WorkspaceIDs = []string{"a"} |
| 90 | state.Workspaces["a"] = Workspace{ID: "a", Root: "/a", SessionIDs: []string{"s"}} |
| 91 | state.Presentation["s"] = Presentation{TopicID: "topic", Title: "before"} |
| 92 | writeReadSnapshotFixture(t, path, state) |
| 93 | store := NewStore(path) |
| 94 | first, err := store.VerifySnapshot(t.Context()) |
| 95 | if err != nil { |
| 96 | t.Fatal(err) |
| 97 | } |
| 98 | second, err := store.VerifySnapshot(t.Context()) |
| 99 | if err != nil || first != second { |
| 100 | t.Fatalf("unchanged bytes rebuilt snapshot: %v", err) |
| 101 | } |
| 102 | // Rendering a previously published view must not wait for the mutex held |
| 103 | // by disk verification or a durable mutation. |
| 104 | store.mu.Lock() |
| 105 | display := store.PublishedSnapshot() |
| 106 | store.mu.Unlock() |
| 107 | if display != first { |
| 108 | t.Fatal("published display snapshot changed without a publication") |
| 109 | } |
| 110 | info, err := os.Stat(path) |
| 111 | if err != nil { |
| 112 | t.Fatal(err) |
| 113 | } |
| 114 | state.Workspaces["a"] = Workspace{ID: "a", Root: "/b", SessionIDs: []string{"s"}} |
| 115 | writeReadSnapshotFixture(t, path, state) // Same generation and file length. |
| 116 | if err := os.Chtimes(path, info.ModTime(), info.ModTime()); err != nil { |
| 117 | t.Fatal(err) |
| 118 | } |
| 119 | if store.PublishedSnapshot() != first { |
| 120 | t.Fatal("display read unexpectedly performs external validation") |
| 121 | } |
| 122 | current, err := store.VerifySnapshot(t.Context()) |
| 123 | if err != nil { |
| 124 | t.Fatal(err) |
| 125 | } |
| 126 | if current == first || current.Session("s").Workspace.Root != "/b" || first.Session("s").Workspace.Root != "/a" { |
| 127 | t.Fatal("legacy replacement was missed or changed a retained snapshot") |
| 128 | } |
| 129 | if got := current.Session("s"); len(got.Workspace.SessionIDs) != 0 || got.Workspace.Organization != nil || !got.Registered { |
| 130 | t.Fatalf("lookup must be bounded metadata: %+v", got) |
| 131 | } |
| 132 | if err := os.WriteFile(path, []byte("{broken"), 0600); err != nil { |
| 133 | t.Fatal(err) |
| 134 | } |
| 135 | if _, err := store.VerifySnapshot(t.Context()); err == nil { |
| 136 | t.Fatal("corrupt authority accepted") |
| 137 | } |
| 138 | if store.PublishedSnapshot() != current { |
| 139 | t.Fatal("failed verification replaced display snapshot") |
| 140 | } |
| 141 | if err := store.RenameWorkspace(t.Context(), "a", "unsafe"); err == nil { |
| 142 | t.Fatal("mutation overwrote corrupt authority") |
| 143 | } |
| 144 | } |
| 145 | |
| 146 | func TestReadSnapshotCommitPublicationAndConflict(t *testing.T) { |
| 147 | path := filepath.Join(t.TempDir(), "registry.json") |
| 148 | state := newState() |
| 149 | state.WorkspaceIDs = []string{"a"} |
| 150 | state.Workspaces["a"] = Workspace{ID: "a", Title: "old", SessionIDs: []string{"s", "other"}} |
| 151 | state.Presentation["s"], state.Presentation["other"] = Presentation{TopicID: "topic"}, Presentation{TopicID: "topic"} |
| 152 | writeReadSnapshotFixture(t, path, state) |
| 153 | store := NewStore(path) |
| 154 | first, err := store.VerifySnapshot(t.Context()) |
| 155 | if err != nil { |
| 156 | t.Fatal(err) |
| 157 | } |
| 158 | if !first.Session("s").SharedTopic { |
| 159 | t.Fatal("shared active topic missing") |
| 160 | } |
| 161 | if err := store.RenameWorkspace(t.Context(), "a", "new"); err != nil { |
| 162 | t.Fatal(err) |
| 163 | } |
| 164 | current := store.PublishedSnapshot() |
| 165 | if current == first || current.Session("s").Workspace.Title != "new" || first.Session("s").Workspace.Title != "old" { |
| 166 | t.Fatal("durable commit did not publish independently owned snapshot") |
| 167 | } |
| 168 | verified, err := store.VerifySnapshot(t.Context()) |
| 169 | if err != nil || verified != current { |
| 170 | t.Fatalf("committed bytes rebuilt snapshot: %v", err) |
| 171 | } |
| 172 | state.SessionStates["other"] = SessionState{Lifecycle: Archived} |
| 173 | writeReadSnapshotFixture(t, path, state) |
| 174 | archived, err := store.VerifySnapshot(t.Context()) |
| 175 | if err != nil { |
| 176 | t.Fatal(err) |
| 177 | } |
| 178 | if archived.Session("s").SharedTopic { |
| 179 | t.Fatal("archived member counted as a shared active topic") |
| 180 | } |
| 181 | // Reject an externally persisted ownership conflict before publication. |
| 182 | state.WorkspaceIDs = append(state.WorkspaceIDs, "b") |
| 183 | state.Workspaces["b"] = Workspace{ID: "b", SessionIDs: []string{"s"}} |
| 184 | writeReadSnapshotFixture(t, path, state) |
| 185 | if _, err := store.VerifySnapshot(t.Context()); err == nil || store.PublishedSnapshot() != archived { |
| 186 | t.Fatal("conflicting membership replaced the verified snapshot") |
| 187 | } |
| 188 | if err := os.Remove(path); err != nil { |
| 189 | t.Fatal(err) |
| 190 | } |
| 191 | empty, err := store.VerifySnapshot(t.Context()) |
| 192 | if err != nil || empty.Session("s").Registered { |
| 193 | t.Fatalf("removed registry retained authority: %v", err) |
| 194 | } |
| 195 | } |
| 196 | |
| 197 | func BenchmarkPublishedSessionLookup(b *testing.B) { |
| 198 | for _, count := range []int{100, 10_000, 100_000} { |
| 199 | for _, projects := range []int{1, 100, 1000} { |
| 200 | b.Run(fmt.Sprintf("sessions=%d/projects=%d", count, projects), func(b *testing.B) { |
| 201 | state := newState() |
| 202 | for i := range projects { |
| 203 | id := fmt.Sprintf("p%d", i) |
| 204 | state.WorkspaceIDs = append(state.WorkspaceIDs, id) |
| 205 | state.Workspaces[id] = Workspace{ID: id} |
| 206 | } |
| 207 | for i := range count { |
| 208 | id, project := fmt.Sprintf("s%d", i), fmt.Sprintf("p%d", i%projects) |
| 209 | w := state.Workspaces[project] |
| 210 | w.SessionIDs = append(w.SessionIDs, id) |
| 211 | state.Workspaces[project] = w |
| 212 | state.SessionStates[id] = SessionState{Lifecycle: Active} |
| 213 | } |
| 214 | store := NewStore(filepath.Join(b.TempDir(), "registry.json")) |
| 215 | store.publishSnapshotLocked(nil, state) |
| 216 | b.ReportAllocs() |
| 217 | b.ResetTimer() |
| 218 | for b.Loop() { |
| 219 | if !store.PublishedSnapshot().Session("s0").Registered { |
| 220 | b.Fatal("missing session") |
| 221 | } |
| 222 | } |
| 223 | }) |
| 224 | } |
| 225 | } |
| 226 | } |
| 227 | |
| 228 | func TestReadSnapshotDoesNotRetainMutableTransactionInputs(t *testing.T) { |
| 229 | store := NewStore(filepath.Join(t.TempDir(), "registry.json")) |
| 230 | members := []string{"s"} |
| 231 | if err := store.mutate(t.Context(), func(state *State) error { |
| 232 | state.WorkspaceIDs = []string{"a"} |
| 233 | state.Workspaces["a"] = Workspace{ID: "a", SessionIDs: members} |
| 234 | return nil |
| 235 | }); err != nil { |
| 236 | t.Fatal(err) |
| 237 | } |
| 238 | snapshot := store.PublishedSnapshot() |
| 239 | members[0] = "replaced" |
| 240 | if snapshot.Session("s").Workspace.ID != "a" || snapshot.Session("replaced").Workspace.ID != "" { |
| 241 | t.Fatal("transaction input mutated the published snapshot") |
| 242 | } |
| 243 | loaded, err := store.Load(t.Context()) |
| 244 | if err != nil { |
| 245 | t.Fatal(err) |
| 246 | } |
| 247 | w := loaded.Workspaces["a"] |
| 248 | w.SessionIDs[0] = "another" |
| 249 | loaded.Presentation["s"] = Presentation{Title: "changed"} |
| 250 | if snapshot.Session("s").Presentation.Title != "" { |
| 251 | t.Fatal("compatibility Load leaked the snapshot's maps") |
| 252 | } |
| 253 | } |
| 254 |