返回 DeepSeek-Reasonix
persistence_test.go
根目录 / internal / session / persistence_test.go
1 package session
2
3 import (
4 "bytes"
5 "context"
6 "errors"
7 "io"
8 "os"
9 "path/filepath"
10 "sync"
11 "testing"
12 "time"
13 )
14
15 func TestFilesystemCreateIsExclusiveAndOpenNeverCreates(t *testing.T) {
16 root := filepath.Join(t.TempDir(), "sessions-v4")
17 persistence := NewFilesystemPersistence(root)
18
19 if _, err := persistence.Open("missing", ReadWrite); !errors.Is(err, ErrSessionNotFound) {
20 t.Fatalf("Open missing error = %v, want ErrSessionNotFound", err)
21 }
22 if _, err := os.Stat(filepath.Join(root, "missing")); !os.IsNotExist(err) {
23 t.Fatalf("Open created missing session: %v", err)
24 }
25
26 first, err := persistence.Create(CreateOptions{SessionID: "created"})
27 if err != nil {
28 t.Fatal(err)
29 }
30 defer first.Close(context.Background())
31 if _, err := persistence.Create(CreateOptions{SessionID: "created"}); !errors.Is(err, ErrSessionExists) {
32 t.Fatalf("duplicate Create error = %v, want ErrSessionExists", err)
33 }
34 }
35
36 func TestFilesystemReadOnlyQueryWritesIndexOutsideSessionDirectory(t *testing.T) {
37 root := filepath.Join(t.TempDir(), "sessions-v4")
38 persistence := NewFilesystemPersistence(root)
39 writer, err := persistence.Create(CreateOptions{SessionID: "readonly"})
40 if err != nil {
41 t.Fatal(err)
42 }
43 if _, err := writer.Append(t.Context(), Batch{OperationID: "one", Events: []Event{{Kind: "diagnostic", Optional: true}}}); err != nil {
44 t.Fatal(err)
45 }
46 if _, err := writer.Flush(t.Context()); err != nil {
47 t.Fatal(err)
48 }
49 if err := writer.Close(t.Context()); err != nil {
50 t.Fatal(err)
51 }
52
53 sessionDir := filepath.Join(root, "readonly")
54 before, err := os.ReadDir(sessionDir)
55 if err != nil {
56 t.Fatal(err)
57 }
58 reader, err := persistence.Open("readonly", ReadOnly)
59 if err != nil {
60 t.Fatal(err)
61 }
62 if _, err := reader.Read(t.Context(), 0, 10); err != nil {
63 t.Fatal(err)
64 }
65 if err := reader.Close(t.Context()); err != nil {
66 t.Fatal(err)
67 }
68 after, err := os.ReadDir(sessionDir)
69 if err != nil {
70 t.Fatal(err)
71 }
72 if len(after) != len(before) {
73 t.Fatalf("read-only query changed session directory: before=%d after=%d", len(before), len(after))
74 }
75 if _, err := os.Stat(filepath.Join(root, ".query-cache", "readonly", "events.offset-index.json")); err != nil {
76 t.Fatalf("external query index: %v", err)
77 }
78 }
79
80 type manualTimer struct {
81 mu sync.Mutex
82 stopped bool
83 fire func()
84 }
85
86 func (t *manualTimer) Stop() bool {
87 t.mu.Lock()
88 defer t.mu.Unlock()
89 if t.stopped {
90 return false
91 }
92 t.stopped = true
93 return true
94 }
95
96 func (t *manualTimer) trigger() {
97 t.mu.Lock()
98 if t.stopped {
99 t.mu.Unlock()
100 return
101 }
102 t.stopped = true
103 fire := t.fire
104 t.mu.Unlock()
105 fire()
106 }
107
108 type manualScheduler struct {
109 mu sync.Mutex
110 durations []time.Duration
111 timers []*manualTimer
112 }
113
114 func (s *manualScheduler) after(delay time.Duration, fire func()) timerHandle {
115 s.mu.Lock()
116 defer s.mu.Unlock()
117 timer := &manualTimer{fire: fire}
118 s.durations = append(s.durations, delay)
119 s.timers = append(s.timers, timer)
120 return timer
121 }
122
123 func (s *manualScheduler) count() int {
124 s.mu.Lock()
125 defer s.mu.Unlock()
126 return len(s.timers)
127 }
128
129 func (s *manualScheduler) trigger(index int) {
130 s.mu.Lock()
131 timer := s.timers[index]
132 s.mu.Unlock()
133 timer.trigger()
134 }
135
136 func TestAppendCommitsToMemoryBeforeDurability(t *testing.T) {
137 dir := filepath.Join(t.TempDir(), "session")
138 scheduler := &manualScheduler{}
139 s, err := OpenWithOptions(dir, "s", OpenOptions{AfterFunc: scheduler.after})
140 if err != nil {
141 t.Fatal(err)
142 }
143 t.Cleanup(func() { _ = s.Close(context.Background()) })
144
145 commit, err := s.Append(t.Context(), Batch{OperationID: "turn-1", Events: []Event{{Kind: "turn/start"}}})
146 if err != nil {
147 t.Fatal(err)
148 }
149 if commit.FirstSequence != 1 || commit.LastSequence() != 1 {
150 t.Fatalf("commit = %+v", commit)
151 }
152 snapshot := s.Snapshot()
153 if snapshot.EventSequence != 1 || snapshot.DurableSequence != 0 || snapshot.Projection.TurnID == "" {
154 t.Fatalf("live snapshot = %+v", snapshot)
155 }
156 commits, err := Replay(dir, nil)
157 if err != nil {
158 t.Fatal(err)
159 }
160 if len(commits) != 0 {
161 t.Fatalf("cold replay observed %d unflushed commits", len(commits))
162 }
163 if scheduler.count() != 1 || scheduler.durations[0] != 200*time.Millisecond {
164 t.Fatalf("scheduled drains = %d at %v", scheduler.count(), scheduler.durations)
165 }
166 }
167
168 func TestPreparePublishesLargePayloadBeforeAcceptance(t *testing.T) {
169 dir := filepath.Join(t.TempDir(), "session")
170 scheduler := &manualScheduler{}
171 s, err := OpenWithOptions(dir, "s", OpenOptions{AfterFunc: scheduler.after})
172 if err != nil {
173 t.Fatal(err)
174 }
175 t.Cleanup(func() { _ = s.Close(context.Background()) })
176 payload := append([]byte(`{"text":"`), bytes.Repeat([]byte("x"), v4InlinePayloadBytes+1)...)
177 payload = append(payload, []byte(`"}`)...)
178 prepared, err := s.PrepareBatchContext(t.Context(), "large", Batch{Events: []Event{{Kind: "diagnostic", Payload: payload}}})
179 if err != nil {
180 t.Fatal(err)
181 }
182 if s.EventSequence() != 0 {
183 t.Fatal("prepare advanced accepted sequence")
184 }
185 if len(prepared.storedEvents) != 1 || prepared.storedEvents[0].PayloadRef == nil || len(prepared.storedEvents[0].Payload) != 0 {
186 t.Fatalf("prepared storage event = %#v", prepared.storedEvents)
187 }
188 if err := s.contentStore().Verify(t.Context(), *prepared.storedEvents[0].PayloadRef); err != nil {
189 t.Fatalf("content was not durable before acceptance: %v", err)
190 }
191 if _, err := s.CommitPrepared(prepared); err != nil {
192 t.Fatal(err)
193 }
194 s.binding.mu.Lock()
195 queued := cloneCommit(s.binding.queue[0])
196 s.binding.mu.Unlock()
197 if len(queued.Events[0].Payload) != 0 || queued.Events[0].PayloadRef == nil {
198 t.Fatalf("write queue retained large body: %#v", queued.Events[0])
199 }
200 if got := s.Snapshot().Projection.CommittedSequence; got != 1 {
201 t.Fatalf("logical projection sequence = %d", got)
202 }
203 }
204
205 func TestPendingHotBudgetBackpressureIsCancellable(t *testing.T) {
206 binding := newPersistenceBinding(nil, t.TempDir(), 0, OpenOptions{})
207 first, err := binding.reserve(t.Context(), pendingHotBytes*4)
208 if err != nil {
209 t.Fatal(err)
210 }
211 ctx, cancel := context.WithCancel(t.Context())
212 cancel()
213 if _, err := binding.reserve(ctx, 1); !errors.Is(err, context.Canceled) {
214 t.Fatalf("backpressure error = %v, want context cancellation", err)
215 }
216 first.release()
217 second, err := binding.reserve(t.Context(), 1)
218 if err != nil {
219 t.Fatalf("capacity was not returned: %v", err)
220 }
221 second.release()
222 }
223
224 func TestRejectedPersistenceAcceptanceDoesNotMutateSession(t *testing.T) {
225 dir := filepath.Join(t.TempDir(), "session")
226 s, err := Open(dir, "s")
227 if err != nil {
228 t.Fatal(err)
229 }
230 t.Cleanup(func() { _ = s.Close(context.Background()) })
231
232 // Simulate the binding admission boundary closing immediately before a
233 // prepared business batch is accepted. The Session projection and sequence
234 // must remain unchanged; there is no accepted fact without a queued copy.
235 prepared, err := s.PrepareBatch("one", Batch{Events: []Event{{Kind: "turn/start"}}})
236 if err != nil {
237 t.Fatal(err)
238 }
239 s.binding.stopAccepting()
240 if _, err := s.CommitPrepared(prepared); err == nil {
241 t.Fatal("commit succeeded after persistence admission closed")
242 }
243 if snapshot := s.Snapshot(); snapshot.EventSequence != 0 || snapshot.Projection.TurnID != "" {
244 t.Fatalf("rejected commit changed session state: %+v", snapshot)
245 }
246 }
247
248 func TestBatchDeadlineIsFixedAndDoesNotReset(t *testing.T) {
249 dir := filepath.Join(t.TempDir(), "session")
250 scheduler := &manualScheduler{}
251 s, err := OpenWithOptions(dir, "s", OpenOptions{AfterFunc: scheduler.after})
252 if err != nil {
253 t.Fatal(err)
254 }
255 t.Cleanup(func() { _ = s.Close(context.Background()) })
256
257 if _, err := s.Append(t.Context(), Batch{OperationID: "one", Events: []Event{{Kind: "turn/start"}}}); err != nil {
258 t.Fatal(err)
259 }
260 if _, err := s.Append(t.Context(), Batch{OperationID: "two", Events: []Event{{Kind: "turn/end", Payload: []byte(`{"status":"completed"}`)}}}); err != nil {
261 t.Fatal(err)
262 }
263 if scheduler.count() != 1 {
264 t.Fatalf("later append reset the batch deadline: timers=%d", scheduler.count())
265 }
266 scheduler.trigger(0)
267 if got := s.Snapshot().DurableSequence; got != 2 {
268 t.Fatalf("durable sequence after scheduled drain = %d", got)
269 }
270 }
271
272 func TestFlushDrainsWritesThatArriveDuringDrain(t *testing.T) {
273 dir := filepath.Join(t.TempDir(), "session")
274 started := make(chan struct{})
275 release := make(chan struct{})
276 var once sync.Once
277 s, err := OpenWithOptions(dir, "s", OpenOptions{
278 Write: func(ctx context.Context, w io.Writer, data []byte) error {
279 once.Do(func() {
280 close(started)
281 <-release
282 })
283 return writeAllContext(ctx, w, data)
284 },
285 })
286 if err != nil {
287 t.Fatal(err)
288 }
289 t.Cleanup(func() { _ = s.Close(context.Background()) })
290 if _, err := s.Append(t.Context(), Batch{OperationID: "one", Events: []Event{{Kind: "turn/start"}}}); err != nil {
291 t.Fatal(err)
292 }
293 done := make(chan error, 1)
294 go func() {
295 _, flushErr := s.Flush(context.Background())
296 done <- flushErr
297 }()
298 <-started
299 if _, err := s.Append(t.Context(), Batch{OperationID: "two", Events: []Event{{Kind: "turn/end", Payload: []byte(`{"status":"completed"}`)}}}); err != nil {
300 t.Fatal(err)
301 }
302 close(release)
303 if err := <-done; err != nil {
304 t.Fatal(err)
305 }
306 if got := s.Snapshot().DurableSequence; got != 2 {
307 t.Fatalf("flush stopped before concurrent append: durable=%d", got)
308 }
309 }
310
311 func TestCancellingOneFlushWaiterDoesNotCancelTheSharedWrite(t *testing.T) {
312 dir := filepath.Join(t.TempDir(), "session")
313 started := make(chan struct{})
314 release := make(chan struct{})
315 var mu sync.Mutex
316 writes := 0
317 s, err := OpenWithOptions(dir, "s", OpenOptions{
318 Write: func(ctx context.Context, w io.Writer, data []byte) error {
319 mu.Lock()
320 writes++
321 if writes == 1 {
322 close(started)
323 }
324 mu.Unlock()
325 <-release
326 return writeAllContext(ctx, w, data)
327 },
328 })
329 if err != nil {
330 t.Fatal(err)
331 }
332 t.Cleanup(func() { _ = s.Close(context.Background()) })
333 if _, err := s.Append(t.Context(), Batch{OperationID: "one", Events: []Event{{Kind: "turn/start"}}}); err != nil {
334 t.Fatal(err)
335 }
336
337 firstCtx, cancelFirst := context.WithCancel(context.Background())
338 first := make(chan error, 1)
339 go func() {
340 _, flushErr := s.Flush(firstCtx)
341 first <- flushErr
342 }()
343 <-started
344 second := make(chan error, 1)
345 go func() {
346 _, flushErr := s.Flush(context.Background())
347 second <- flushErr
348 }()
349 cancelFirst()
350 if err := <-first; !errors.Is(err, context.Canceled) {
351 t.Fatalf("cancelled waiter error = %v", err)
352 }
353 close(release)
354 if err := <-second; err != nil {
355 t.Fatal(err)
356 }
357 mu.Lock()
358 defer mu.Unlock()
359 if writes != 1 || s.Snapshot().DurableSequence != 1 {
360 t.Fatalf("writes=%d snapshot=%+v", writes, s.Snapshot())
361 }
362 }
363
364 func TestScheduledDrainImmediatelyConsumesWritesThatArriveDuringDrain(t *testing.T) {
365 dir := filepath.Join(t.TempDir(), "session")
366 scheduler := &manualScheduler{}
367 started := make(chan struct{})
368 release := make(chan struct{})
369 var once sync.Once
370 s, err := OpenWithOptions(dir, "s", OpenOptions{
371 AfterFunc: scheduler.after,
372 Write: func(ctx context.Context, w io.Writer, data []byte) error {
373 once.Do(func() {
374 close(started)
375 <-release
376 })
377 return writeAllContext(ctx, w, data)
378 },
379 })
380 if err != nil {
381 t.Fatal(err)
382 }
383 t.Cleanup(func() { _ = s.Close(context.Background()) })
384 if _, err := s.Append(t.Context(), Batch{OperationID: "one", Events: []Event{{Kind: "turn/start"}}}); err != nil {
385 t.Fatal(err)
386 }
387 go scheduler.trigger(0)
388 <-started
389 if _, err := s.Append(t.Context(), Batch{OperationID: "two", Events: []Event{{Kind: "turn/end", Payload: []byte(`{"status":"completed"}`)}}}); err != nil {
390 t.Fatal(err)
391 }
392 close(release)
393 deadline := time.After(2 * time.Second)
394 for s.Snapshot().DurableSequence != 2 {
395 select {
396 case <-deadline:
397 t.Fatalf("scheduled drain stopped before concurrent append: %+v", s.Snapshot())
398 default:
399 time.Sleep(time.Millisecond)
400 }
401 }
402 if scheduler.count() != 1 {
403 t.Fatalf("active drain scheduled a second batching window: %d", scheduler.count())
404 }
405 }
406
407 func TestBackgroundFailurePausesRetryAndExplicitFlushRetries(t *testing.T) {
408 dir := filepath.Join(t.TempDir(), "session")
409 scheduler := &manualScheduler{}
410 injected := errors.New("disk full")
411 var mu sync.Mutex
412 writes := 0
413 s, err := OpenWithOptions(dir, "s", OpenOptions{
414 AfterFunc: scheduler.after,
415 Write: func(ctx context.Context, w io.Writer, data []byte) error {
416 mu.Lock()
417 writes++
418 attempt := writes
419 mu.Unlock()
420 if attempt == 1 {
421 return injected
422 }
423 return writeAllContext(ctx, w, data)
424 },
425 })
426 if err != nil {
427 t.Fatal(err)
428 }
429 t.Cleanup(func() { _ = s.Close(context.Background()) })
430 if _, err := s.Append(t.Context(), Batch{OperationID: "one", Events: []Event{{Kind: "turn/start"}}}); err != nil {
431 t.Fatal(err)
432 }
433 scheduler.trigger(0)
434 if snapshot := s.Snapshot(); snapshot.DurableSequence != 0 || snapshot.PersistenceStatus != PersistenceFailed {
435 t.Fatalf("snapshot after background failure = %+v", snapshot)
436 }
437 if _, err := s.Append(t.Context(), Batch{OperationID: "two", Events: []Event{{Kind: "turn/end", Payload: []byte(`{"status":"completed"}`)}}}); err != nil {
438 t.Fatal(err)
439 }
440 if scheduler.count() != 1 {
441 t.Fatalf("failed writer scheduled an automatic retry: %d timers", scheduler.count())
442 }
443 receipt, err := s.Flush(t.Context())
444 if err != nil {
445 t.Fatal(err)
446 }
447 if receipt.DurableSequence != 2 || s.Snapshot().PersistenceStatus != PersistenceReady {
448 t.Fatalf("flush receipt/snapshot = %+v / %+v", receipt, s.Snapshot())
449 }
450 }
451
452 func TestOperationRetryIsIdempotentAndConflictFails(t *testing.T) {
453 dir := filepath.Join(t.TempDir(), "session")
454 s, err := Open(dir, "s")
455 if err != nil {
456 t.Fatal(err)
457 }
458 t.Cleanup(func() { _ = s.Close(context.Background()) })
459 one := Batch{OperationID: "same", Events: []Event{{Kind: "turn/start"}}}
460 first, err := s.Append(t.Context(), one)
461 if err != nil {
462 t.Fatal(err)
463 }
464 second, err := s.Append(t.Context(), one)
465 if err != nil {
466 t.Fatal(err)
467 }
468 if first.ID != second.ID || s.Snapshot().EventSequence != 1 {
469 t.Fatalf("retry appended twice: first=%+v second=%+v snapshot=%+v", first, second, s.Snapshot())
470 }
471 _, err = s.Append(t.Context(), Batch{OperationID: "same", TurnID: "different-turn", Events: []Event{{Kind: "turn/start"}}})
472 if !errors.Is(err, ErrOperationConflict) {
473 t.Fatalf("same payload for a different turn error = %v", err)
474 }
475 _, err = s.Append(t.Context(), Batch{OperationID: "same", Events: []Event{{Kind: "turn/end", Payload: []byte(`{"status":"completed"}`)}}})
476 if !errors.Is(err, ErrOperationConflict) {
477 t.Fatalf("conflicting retry error = %v", err)
478 }
479 if _, err := s.Flush(t.Context()); err != nil {
480 t.Fatal(err)
481 }
482 if len(s.commits) != 0 {
483 t.Fatalf("durable commits remained in the runtime: %d", len(s.commits))
484 }
485 if err := s.Close(t.Context()); err != nil {
486 t.Fatal(err)
487 }
488 reopened, err := Open(dir, "s")
489 if err != nil {
490 t.Fatal(err)
491 }
492 t.Cleanup(func() { _ = reopened.Close(context.Background()) })
493 if len(reopened.commits) != 0 || len(reopened.operations["same"].commit.Events) != 0 {
494 t.Fatalf("reopen retained durable bodies: commits=%d operation=%+v", len(reopened.commits), reopened.operations["same"])
495 }
496 retried, err := reopened.Append(t.Context(), one)
497 if err != nil || retried.ID != first.ID || retried.FirstSequence != first.FirstSequence || len(retried.Events) != 1 {
498 t.Fatalf("reopened retry = %+v, %v", retried, err)
499 }
500 page, err := reopened.AcceptedPage(t.Context(), 0, 10)
501 if err != nil || len(page.Commits) != 1 || page.Commits[0].ID != first.ID {
502 t.Fatalf("accepted durable page = %+v, %v", page, err)
503 }
504 }
505
506 func TestExplicitFlushRepairsPreservedPartialAppendWithoutDuplicateCommit(t *testing.T) {
507 dir := filepath.Join(t.TempDir(), "session")
508 injected := errors.New("connection to filesystem interrupted")
509 writes := 0
510 s, err := OpenWithOptions(dir, "s", OpenOptions{Write: func(ctx context.Context, w io.Writer, data []byte) error {
511 writes++
512 if writes == 1 {
513 if _, err := w.Write(data[:len(data)/2]); err != nil {
514 return err
515 }
516 return injected
517 }
518 return writeAllContext(ctx, w, data)
519 }})
520 if err != nil {
521 t.Fatal(err)
522 }
523 t.Cleanup(func() { _ = s.Close(context.Background()) })
524 if _, err := s.Append(t.Context(), Batch{OperationID: "one", Events: []Event{{Kind: "turn/start"}}}); err != nil {
525 t.Fatal(err)
526 }
527 if _, err := s.Flush(t.Context()); !errors.Is(err, ErrPersistenceUncertain) {
528 t.Fatalf("partial append error = %v", err)
529 }
530 if snapshot := s.Snapshot(); snapshot.DurableSequence != 0 || snapshot.PersistenceStatus != PersistenceUncertain {
531 t.Fatalf("partial append snapshot = %+v", snapshot)
532 }
533 receipt, err := s.Flush(t.Context())
534 if err != nil {
535 t.Fatal(err)
536 }
537 if receipt.DurableSequence != 1 || writes != 2 {
538 t.Fatalf("repair receipt=%+v writes=%d", receipt, writes)
539 }
540 commits, err := Replay(dir, nil)
541 if err != nil {
542 t.Fatal(err)
543 }
544 if len(commits) != 1 || commits[0].OperationID != "one" {
545 t.Fatalf("repaired commits = %+v", commits)
546 }
547 backups, err := filepath.Glob(filepath.Join(dir, "events.uncertain-*.tail"))
548 if err != nil || len(backups) != 1 {
549 t.Fatalf("partial tail backups=%v err=%v", backups, err)
550 }
551 if info, err := os.Stat(backups[0]); err != nil || info.Size() == 0 {
552 t.Fatalf("partial tail backup info=%v err=%v", info, err)
553 }
554 }
555
556 func TestExplicitFlushRecognizesCompleteAppendAfterWriteError(t *testing.T) {
557 dir := filepath.Join(t.TempDir(), "session")
558 writes := 0
559 s, err := OpenWithOptions(dir, "s", OpenOptions{Write: func(ctx context.Context, w io.Writer, data []byte) error {
560 writes++
561 if err := writeAllContext(ctx, w, data); err != nil {
562 return err
563 }
564 if writes == 1 {
565 return errors.New("late write acknowledgement lost")
566 }
567 return nil
568 }})
569 if err != nil {
570 t.Fatal(err)
571 }
572 t.Cleanup(func() { _ = s.Close(context.Background()) })
573 if _, err := s.Append(t.Context(), Batch{OperationID: "one", Events: []Event{{Kind: "turn/start"}}}); err != nil {
574 t.Fatal(err)
575 }
576 // A full exact record is proved and synced during the first call, so the
577 // write implementation's late error is treated as success immediately.
578 if receipt, err := s.Flush(t.Context()); err != nil || receipt.DurableSequence != 1 {
579 t.Fatalf("flush receipt=%+v err=%v", receipt, err)
580 }
581 if writes != 1 {
582 t.Fatalf("complete uncertain append was written %d times", writes)
583 }
584 commits, err := Replay(dir, nil)
585 if err != nil || len(commits) != 1 {
586 t.Fatalf("commits=%+v err=%v", commits, err)
587 }
588 }
589
590 func TestExplicitFlushRetriesOnlySyncAfterUncertainFsync(t *testing.T) {
591 dir := filepath.Join(t.TempDir(), "session")
592 writes, syncs := 0, 0
593 s, err := OpenWithOptions(dir, "s", OpenOptions{
594 Write: func(ctx context.Context, w io.Writer, data []byte) error {
595 writes++
596 return writeAllContext(ctx, w, data)
597 },
598 Sync: func(file *os.File) error {
599 syncs++
600 if syncs == 1 {
601 return errors.New("fsync interrupted")
602 }
603 return file.Sync()
604 },
605 })
606 if err != nil {
607 t.Fatal(err)
608 }
609 t.Cleanup(func() { _ = s.Close(context.Background()) })
610 if _, err := s.Append(t.Context(), Batch{OperationID: "one", Events: []Event{{Kind: "turn/start"}}}); err != nil {
611 t.Fatal(err)
612 }
613 if _, err := s.Flush(t.Context()); !errors.Is(err, ErrPersistenceUncertain) {
614 t.Fatalf("first fsync error = %v", err)
615 }
616 if receipt, err := s.Flush(t.Context()); err != nil || receipt.DurableSequence != 1 {
617 t.Fatalf("retry receipt=%+v err=%v", receipt, err)
618 }
619 if writes != 1 || syncs != 2 {
620 t.Fatalf("writes=%d syncs=%d, want 1/2", writes, syncs)
621 }
622 }
623
624 func TestConfirmedUncertainPrefixDoesNotMarkLaterBatchDurable(t *testing.T) {
625 dir := filepath.Join(t.TempDir(), "session")
626 writes, syncs := 0, 0
627 s, err := OpenWithOptions(dir, "s", OpenOptions{
628 Write: func(ctx context.Context, w io.Writer, data []byte) error {
629 writes++
630 if err := writeAllContext(ctx, w, data); err != nil {
631 return err
632 }
633 if writes == 1 {
634 return errors.New("append acknowledgement lost")
635 }
636 return nil
637 },
638 Sync: func(file *os.File) error {
639 syncs++
640 if syncs == 1 {
641 return errors.New("first sync unavailable")
642 }
643 return file.Sync()
644 },
645 })
646 if err != nil {
647 t.Fatal(err)
648 }
649 t.Cleanup(func() { _ = s.Close(context.Background()) })
650 if _, err := s.Append(t.Context(), Batch{OperationID: "a", Events: []Event{{Kind: "turn/start"}}}); err != nil {
651 t.Fatal(err)
652 }
653 if _, err := s.Flush(t.Context()); !errors.Is(err, ErrPersistenceUncertain) {
654 t.Fatalf("first flush error = %v", err)
655 }
656 if _, err := s.Append(t.Context(), Batch{OperationID: "b", Events: []Event{{Kind: "turn/end", Payload: []byte(`{"status":"completed"}`)}}}); err != nil {
657 t.Fatal(err)
658 }
659 receipt, err := s.Flush(t.Context())
660 if err != nil {
661 t.Fatal(err)
662 }
663 if receipt.DurableSequence != 2 || writes != 2 {
664 t.Fatalf("receipt=%+v writes=%d syncs=%d", receipt, writes, syncs)
665 }
666 commits, err := Replay(dir, nil)
667 if err != nil || len(commits) != 2 || commits[0].OperationID != "a" || commits[1].OperationID != "b" {
668 t.Fatalf("commits=%+v err=%v", commits, err)
669 }
670 }
671
672 func TestCloseReturnsStableFailureAndReleasesOwnership(t *testing.T) {
673 dir := filepath.Join(t.TempDir(), "session")
674 injected := errors.New("disk unavailable")
675 s, err := OpenWithOptions(dir, "s", OpenOptions{Write: func(context.Context, io.Writer, []byte) error {
676 return injected
677 }})
678 if err != nil {
679 t.Fatal(err)
680 }
681 if _, err := s.Append(t.Context(), Batch{OperationID: "one", Events: []Event{{Kind: "turn/start"}}}); err != nil {
682 t.Fatal(err)
683 }
684 first := s.Close(t.Context())
685 second := s.Close(context.Background())
686 if !errors.Is(first, injected) || first.Error() != second.Error() {
687 t.Fatalf("close errors first=%v second=%v", first, second)
688 }
689 if reopened, err := Open(dir, "s"); err != nil {
690 t.Fatalf("close did not release writer ownership: %v", err)
691 } else {
692 _ = reopened.Close(context.Background())
693 }
694 }
695
695 lines GO