Skip to content

Commit 809489a

Browse files
GautamBytesGautam Manchandani
andauthored
Batch session event persistence (#526)
Co-authored-by: Gautam Manchandani <gautammanch@Gautams-MacBook-Air.local>
1 parent 56147a7 commit 809489a

4 files changed

Lines changed: 391 additions & 48 deletions

File tree

Lines changed: 233 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,233 @@
1+
package sessions
2+
3+
import (
4+
"encoding/json"
5+
"fmt"
6+
"reflect"
7+
"strings"
8+
"sync"
9+
"testing"
10+
"time"
11+
)
12+
13+
func TestStoreAppendEventsBatchesSequencesAndMetadata(t *testing.T) {
14+
store := NewStore(StoreOptions{RootDir: t.TempDir(), Now: sequenceClock([]time.Time{
15+
time.Date(2026, 6, 4, 15, 0, 0, 0, time.UTC),
16+
time.Date(2026, 6, 4, 15, 0, 1, 0, time.UTC),
17+
time.Date(2026, 6, 4, 15, 0, 2, 0, time.UTC),
18+
time.Date(2026, 6, 4, 15, 0, 3, 0, time.UTC),
19+
})})
20+
session, err := store.Create(CreateInput{SessionID: "batch"})
21+
if err != nil {
22+
t.Fatalf("Create returned error: %v", err)
23+
}
24+
25+
appended, err := store.AppendEvents(session.SessionID, []AppendEventInput{
26+
{Type: EventMessage, Payload: map[string]any{"role": "user", "content": "one"}},
27+
{Type: EventToolCall, Payload: map[string]any{"id": "call_1", "name": "read_file"}},
28+
{Type: EventToolResult, Payload: map[string]any{"toolCallId": "call_1", "status": "ok"}},
29+
})
30+
if err != nil {
31+
t.Fatalf("AppendEvents returned error: %v", err)
32+
}
33+
if len(appended) != 3 {
34+
t.Fatalf("expected 3 appended events, got %#v", appended)
35+
}
36+
for index, event := range appended {
37+
wantSequence := index + 1
38+
if event.Sequence != wantSequence || event.ID != fmt.Sprintf("batch:%d", wantSequence) {
39+
t.Fatalf("event identity mismatch at %d: %#v", index, event)
40+
}
41+
}
42+
if appended[0].CreatedAt != "2026-06-04T15:00:01Z" || appended[2].CreatedAt != "2026-06-04T15:00:03Z" {
43+
t.Fatalf("unexpected event timestamps: %#v", appended)
44+
}
45+
46+
loaded, err := store.Get(session.SessionID)
47+
if err != nil {
48+
t.Fatalf("Get returned error: %v", err)
49+
}
50+
if loaded == nil || loaded.EventCount != 3 || loaded.LastEventType != EventToolResult || loaded.UpdatedAt != appended[2].CreatedAt {
51+
t.Fatalf("metadata not updated from final batch event: %#v", loaded)
52+
}
53+
events, err := store.ReadEvents(session.SessionID)
54+
if err != nil {
55+
t.Fatalf("ReadEvents returned error: %v", err)
56+
}
57+
if !reflect.DeepEqual(eventTypesForTest(events), []EventType{EventMessage, EventToolCall, EventToolResult}) {
58+
t.Fatalf("unexpected event types: %#v", events)
59+
}
60+
}
61+
62+
func TestStoreAppendEventsEmptyBatchDoesNotRewriteMetadata(t *testing.T) {
63+
store := NewStore(StoreOptions{RootDir: t.TempDir(), Now: fixedClock("2026-06-04T15:10:00Z")})
64+
session, err := store.Create(CreateInput{SessionID: "empty_batch", Title: "unchanged"})
65+
if err != nil {
66+
t.Fatalf("Create returned error: %v", err)
67+
}
68+
69+
appended, err := store.AppendEvents(session.SessionID, nil)
70+
if err != nil {
71+
t.Fatalf("AppendEvents empty returned error: %v", err)
72+
}
73+
if len(appended) != 0 {
74+
t.Fatalf("empty AppendEvents returned %#v", appended)
75+
}
76+
loaded, err := store.Get(session.SessionID)
77+
if err != nil {
78+
t.Fatalf("Get returned error: %v", err)
79+
}
80+
if loaded == nil || !reflect.DeepEqual(*loaded, session) {
81+
t.Fatalf("empty AppendEvents rewrote metadata: before=%#v after=%#v", session, loaded)
82+
}
83+
}
84+
85+
func TestStoreAppendEventsInvalidPayloadDoesNotPartiallyAppend(t *testing.T) {
86+
store := NewStore(StoreOptions{RootDir: t.TempDir(), Now: fixedClock("2026-06-04T15:20:00Z")})
87+
session, err := store.Create(CreateInput{SessionID: "invalid_batch"})
88+
if err != nil {
89+
t.Fatalf("Create returned error: %v", err)
90+
}
91+
92+
_, err = store.AppendEvents(session.SessionID, []AppendEventInput{
93+
{Type: EventMessage, Payload: map[string]any{"content": "valid first"}},
94+
{Type: EventMessage, Payload: json.RawMessage(`{"broken"`)},
95+
})
96+
if err == nil || !strings.Contains(err.Error(), "invalid raw JSON payload") {
97+
t.Fatalf("expected invalid raw payload error, got %v", err)
98+
}
99+
events, err := store.ReadEvents(session.SessionID)
100+
if err != nil {
101+
t.Fatalf("ReadEvents returned error: %v", err)
102+
}
103+
if len(events) != 0 {
104+
t.Fatalf("invalid batch should not append partial events: %#v", events)
105+
}
106+
loaded, err := store.Get(session.SessionID)
107+
if err != nil {
108+
t.Fatalf("Get returned error: %v", err)
109+
}
110+
if loaded == nil || loaded.EventCount != 0 || loaded.LastEventType != "" {
111+
t.Fatalf("invalid batch should not update metadata: %#v", loaded)
112+
}
113+
}
114+
115+
func TestStoreAppendEventsSequencesAfterMetadataLag(t *testing.T) {
116+
store := NewStore(StoreOptions{RootDir: t.TempDir(), Now: fixedClock("2026-06-04T15:30:00Z")})
117+
session, err := store.Create(CreateInput{SessionID: "lagging"})
118+
if err != nil {
119+
t.Fatalf("Create returned error: %v", err)
120+
}
121+
if _, err := store.AppendEvent(session.SessionID, AppendEventInput{Type: EventMessage, Payload: map[string]any{"content": "already durable"}}); err != nil {
122+
t.Fatalf("AppendEvent returned error: %v", err)
123+
}
124+
lagging, err := store.Get(session.SessionID)
125+
if err != nil {
126+
t.Fatalf("Get returned error: %v", err)
127+
}
128+
lagging.EventCount = 0
129+
lagging.LastEventType = ""
130+
if err := store.writeMetadata(*lagging); err != nil {
131+
t.Fatalf("writeMetadata returned error: %v", err)
132+
}
133+
134+
appended, err := store.AppendEvents(session.SessionID, []AppendEventInput{
135+
{Type: EventToolCall, Payload: map[string]any{"id": "call"}},
136+
{Type: EventToolResult, Payload: map[string]any{"toolCallId": "call"}},
137+
})
138+
if err != nil {
139+
t.Fatalf("AppendEvents returned error: %v", err)
140+
}
141+
if len(appended) != 2 || appended[0].Sequence != 2 || appended[1].Sequence != 3 {
142+
t.Fatalf("batch did not sequence after durable log tail: %#v", appended)
143+
}
144+
loaded, err := store.Get(session.SessionID)
145+
if err != nil {
146+
t.Fatalf("Get returned error: %v", err)
147+
}
148+
if loaded == nil || loaded.EventCount != 3 || loaded.LastEventType != EventToolResult {
149+
t.Fatalf("metadata not repaired by batch append: %#v", loaded)
150+
}
151+
}
152+
153+
func TestStoreAppendEventsSerializesConcurrentBatches(t *testing.T) {
154+
store := NewStore(StoreOptions{RootDir: t.TempDir(), Now: fixedClock("2026-06-04T15:40:00Z")})
155+
session, err := store.Create(CreateInput{SessionID: "concurrent_batches"})
156+
if err != nil {
157+
t.Fatalf("Create returned error: %v", err)
158+
}
159+
160+
const writers = 8
161+
const perBatch = 3
162+
var wg sync.WaitGroup
163+
errs := make(chan error, writers)
164+
for writer := 0; writer < writers; writer++ {
165+
wg.Add(1)
166+
go func(writer int) {
167+
defer wg.Done()
168+
appended, err := store.AppendEvents(session.SessionID, []AppendEventInput{
169+
{Type: EventMessage, Payload: map[string]int{"writer": writer, "index": 0}},
170+
{Type: EventMessage, Payload: map[string]int{"writer": writer, "index": 1}},
171+
{Type: EventMessage, Payload: map[string]int{"writer": writer, "index": 2}},
172+
})
173+
if err != nil {
174+
errs <- err
175+
return
176+
}
177+
if len(appended) != perBatch {
178+
errs <- fmt.Errorf("writer %d appended %d events", writer, len(appended))
179+
return
180+
}
181+
for index := 1; index < len(appended); index++ {
182+
if appended[index].Sequence != appended[index-1].Sequence+1 {
183+
errs <- fmt.Errorf("writer %d batch was interleaved: %#v", writer, appended)
184+
return
185+
}
186+
}
187+
errs <- nil
188+
}(writer)
189+
}
190+
wg.Wait()
191+
close(errs)
192+
for err := range errs {
193+
if err != nil {
194+
t.Fatalf("AppendEvents returned error: %v", err)
195+
}
196+
}
197+
198+
events, err := store.ReadEvents(session.SessionID)
199+
if err != nil {
200+
t.Fatalf("ReadEvents returned error: %v", err)
201+
}
202+
total := writers * perBatch
203+
if len(events) != total {
204+
t.Fatalf("expected %d events, got %d", total, len(events))
205+
}
206+
seen := map[int]bool{}
207+
for _, event := range events {
208+
if seen[event.Sequence] {
209+
t.Fatalf("duplicate sequence %d in %#v", event.Sequence, events)
210+
}
211+
seen[event.Sequence] = true
212+
}
213+
for sequence := 1; sequence <= total; sequence++ {
214+
if !seen[sequence] {
215+
t.Fatalf("missing sequence %d in %#v", sequence, events)
216+
}
217+
}
218+
loaded, err := store.Get(session.SessionID)
219+
if err != nil {
220+
t.Fatalf("Get returned error: %v", err)
221+
}
222+
if loaded == nil || loaded.EventCount != total || loaded.LastEventType != EventMessage {
223+
t.Fatalf("metadata not updated after concurrent batches: %#v", loaded)
224+
}
225+
}
226+
227+
func eventTypesForTest(events []Event) []EventType {
228+
types := make([]EventType, 0, len(events))
229+
for _, event := range events {
230+
types = append(types, event.Type)
231+
}
232+
return types
233+
}

0 commit comments

Comments
 (0)