From e40f6b64a5ba70037e7a1dd6a33b15443dfb586e Mon Sep 17 00:00:00 2001 From: Mawen Salignat-Moandal Date: Sun, 30 Aug 2026 22:55:40 +0200 Subject: [PATCH 1/3] test(runtime): add persistence observer streaming benchmarks Document the per-chunk AddMessage/UpdateMessage contract and establish in-memory vs SQLite baselines for streaming assistant persistence. --- .../persistence_observer_bench_test.go | 161 ++++++++++++++++++ 1 file changed, 161 insertions(+) create mode 100644 pkg/runtime/persistence_observer_bench_test.go diff --git a/pkg/runtime/persistence_observer_bench_test.go b/pkg/runtime/persistence_observer_bench_test.go new file mode 100644 index 000000000..40f20eb00 --- /dev/null +++ b/pkg/runtime/persistence_observer_bench_test.go @@ -0,0 +1,161 @@ +package runtime + +import ( + "context" + "database/sql" + "sync/atomic" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + _ "modernc.org/sqlite" + + "github.com/docker/docker-agent/pkg/chat" + "github.com/docker/docker-agent/pkg/session" +) + +// streamingBenchChunks is the number of AgentChoice deltas exercised per +// benchmark iteration — roughly the order of magnitude of a long assistant turn. +const streamingBenchChunks = 500 + +// countingStore wraps an in-memory store and records AddMessage / UpdateMessage +// calls so tests can pin the per-chunk persistence contract. +type countingStore struct { + *session.InMemorySessionStore + addCalls atomic.Int64 + updateCalls atomic.Int64 +} + +func newCountingStore() *countingStore { + return &countingStore{ + InMemorySessionStore: session.NewInMemorySessionStore().(*session.InMemorySessionStore), + } +} + +func (s *countingStore) AddMessage(ctx context.Context, sessionID string, msg *session.Message) (int64, error) { + s.addCalls.Add(1) + return s.InMemorySessionStore.AddMessage(ctx, sessionID, msg) +} + +func (s *countingStore) UpdateMessage(ctx context.Context, messageID int64, msg *session.Message) error { + s.updateCalls.Add(1) + return s.InMemorySessionStore.UpdateMessage(ctx, messageID, msg) +} + +func setupPersistenceObserverBench(tb testing.TB) (*PersistenceObserver, *session.Session, *session.InMemorySessionStore) { + tb.Helper() + store := session.NewInMemorySessionStore().(*session.InMemorySessionStore) + obs := newPersistenceObserver(store) + require.NotNil(tb, obs) + + sess := session.New(session.WithID("bench-session"), session.WithUserMessage("hi")) + require.NoError(tb, store.AddSession(context.Background(), sess)) + return obs, sess, store +} + +func emitStreamingChunks(ctx context.Context, obs *PersistenceObserver, sess *session.Session, chunks int) { + for range chunks { + obs.OnEvent(ctx, sess, AgentChoice("root", sess.ID, "tok")) + } +} + +func finalizeStreamingMessage(ctx context.Context, obs *PersistenceObserver, sess *session.Session) { + obs.OnEvent(ctx, sess, MessageAdded(sess.ID, &session.Message{ + AgentName: "root", + Message: chat.Message{ + Role: chat.MessageRoleAssistant, + Content: "done", + }, + }, "root")) +} + +// TestPersistenceObserver_UpdateCountPerChunk documents that each streaming +// delta triggers a store write: one AddMessage on the first chunk, then one +// UpdateMessage per subsequent chunk until MessageAddedEvent finalises the row. +func TestPersistenceObserver_UpdateCountPerChunk(t *testing.T) { + t.Parallel() + + const chunks = 100 + ctx := t.Context() + + store := newCountingStore() + obs := newPersistenceObserver(store) + require.NotNil(t, obs) + + sess := session.New(session.WithID("s1"), session.WithUserMessage("hi")) + require.NoError(t, store.AddSession(ctx, sess)) + + emitStreamingChunks(ctx, obs, sess, chunks) + + assert.Equal(t, int64(1), store.addCalls.Load(), "first chunk should INSERT") + assert.Equal(t, int64(chunks-1), store.updateCalls.Load(), "each later chunk should UPDATE") + + finalizeStreamingMessage(ctx, obs, sess) + // MessageAdded with an existing streaming row issues one more UpdateMessage. + assert.Equal(t, int64(chunks), store.updateCalls.Load()) +} + +// TestPersistenceObserver_StreamingContentAccumulates verifies mid-stream +// persistence keeps the growing assistant text in the store. +func TestPersistenceObserver_StreamingContentAccumulates(t *testing.T) { + t.Parallel() + + ctx := t.Context() + store := session.NewInMemorySessionStore() + obs := newPersistenceObserver(store) + require.NotNil(t, obs) + + sess := session.New(session.WithID("s1"), session.WithUserMessage("hi")) + require.NoError(t, store.AddSession(ctx, sess)) + + obs.OnEvent(ctx, sess, AgentChoice("root", sess.ID, "hel")) + obs.OnEvent(ctx, sess, AgentChoice("root", sess.ID, "lo")) + + reloaded, err := store.GetSession(ctx, sess.ID) + require.NoError(t, err) + require.Len(t, reloaded.Messages, 2) // user + streaming assistant + require.NotNil(t, reloaded.Messages[1].Message) + assert.Equal(t, "hello", reloaded.Messages[1].Message.Message.Content) +} + +func BenchmarkPersistenceObserver_StreamingChunks(b *testing.B) { + ctx := b.Context() + obs, sess, _ := setupPersistenceObserverBench(b) + + b.ReportAllocs() + b.ResetTimer() + for range b.N { + emitStreamingChunks(ctx, obs, sess, streamingBenchChunks) + finalizeStreamingMessage(ctx, obs, sess) + } +} + +func BenchmarkPersistenceObserver_StreamingChunks_SQLite(b *testing.B) { + ctx := b.Context() + store := openBenchSQLiteStore(b) + obs := newPersistenceObserver(store) + require.NotNil(b, obs) + + sess := session.New(session.WithID("bench-sqlite"), session.WithUserMessage("hi")) + require.NoError(b, store.AddSession(ctx, sess)) + + b.ReportAllocs() + b.ResetTimer() + for range b.N { + emitStreamingChunks(ctx, obs, sess, streamingBenchChunks) + finalizeStreamingMessage(ctx, obs, sess) + } +} + +func openBenchSQLiteStore(b *testing.B) *session.SQLiteSessionStore { + b.Helper() + db, err := sql.Open("sqlite", ":memory:") + require.NoError(b, err) + b.Cleanup(func() { _ = db.Close() }) + db.SetMaxOpenConns(1) + + store, err := session.NewSQLiteSessionStoreFromDB(b.Context(), db) + require.NoError(b, err) + b.Cleanup(func() { _ = store.Close() }) + return store +} From a0ddcd4a47cced4fbf7986ebf92e4ac0c90e06b2 Mon Sep 17 00:00:00 2001 From: Mawen Salignat-Moandal Date: Mon, 31 Aug 2026 18:18:26 +0200 Subject: [PATCH 2/3] test(runtime): bound in-memory persistence observer bench store MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Reset the session each iteration so UpdateMessage's O(sessions×messages) scan cannot make ns/op a function of b.N. --- .../persistence_observer_bench_test.go | 22 +++++++++++++------ 1 file changed, 15 insertions(+), 7 deletions(-) diff --git a/pkg/runtime/persistence_observer_bench_test.go b/pkg/runtime/persistence_observer_bench_test.go index 40f20eb00..c5759517f 100644 --- a/pkg/runtime/persistence_observer_bench_test.go +++ b/pkg/runtime/persistence_observer_bench_test.go @@ -3,6 +3,7 @@ package runtime import ( "context" "database/sql" + "strconv" "sync/atomic" "testing" @@ -42,15 +43,12 @@ func (s *countingStore) UpdateMessage(ctx context.Context, messageID int64, msg return s.InMemorySessionStore.UpdateMessage(ctx, messageID, msg) } -func setupPersistenceObserverBench(tb testing.TB) (*PersistenceObserver, *session.Session, *session.InMemorySessionStore) { +func setupPersistenceObserverBench(tb testing.TB) (*PersistenceObserver, *session.InMemorySessionStore) { tb.Helper() store := session.NewInMemorySessionStore().(*session.InMemorySessionStore) obs := newPersistenceObserver(store) require.NotNil(tb, obs) - - sess := session.New(session.WithID("bench-session"), session.WithUserMessage("hi")) - require.NoError(tb, store.AddSession(context.Background(), sess)) - return obs, sess, store + return obs, store } func emitStreamingChunks(ctx context.Context, obs *PersistenceObserver, sess *session.Session, chunks int) { @@ -120,13 +118,23 @@ func TestPersistenceObserver_StreamingContentAccumulates(t *testing.T) { func BenchmarkPersistenceObserver_StreamingChunks(b *testing.B) { ctx := b.Context() - obs, sess, _ := setupPersistenceObserverBench(b) + obs, store := setupPersistenceObserverBench(b) b.ReportAllocs() b.ResetTimer() - for range b.N { + // Bound the store each iteration: UpdateMessage scans every session's + // messages, so a shared session (or accumulating sessions) makes ns/op + // a function of b.N rather than a per-turn baseline. + for i := range b.N { + sess := session.New(session.WithID(strconv.Itoa(i)), session.WithUserMessage("hi")) + if err := store.AddSession(ctx, sess); err != nil { + b.Fatal(err) + } emitStreamingChunks(ctx, obs, sess, streamingBenchChunks) finalizeStreamingMessage(ctx, obs, sess) + if err := store.DeleteSession(ctx, sess.ID); err != nil { + b.Fatal(err) + } } } From 29b4b7f42ad3c298c4be85275dbde67fd4ac129b Mon Sep 17 00:00:00 2001 From: Mawen Salignat-Moandal Date: Mon, 31 Aug 2026 18:42:35 +0200 Subject: [PATCH 3/3] test(runtime): silence sqlite migration logs in persistence benches Discard the default slog handler during SQLite store setup so harness re-invocations do not interleave migration INFO with -bench output. --- pkg/runtime/persistence_observer_bench_test.go | 7 +++++++ 1 file changed, 7 insertions(+) diff --git a/pkg/runtime/persistence_observer_bench_test.go b/pkg/runtime/persistence_observer_bench_test.go index c5759517f..665753121 100644 --- a/pkg/runtime/persistence_observer_bench_test.go +++ b/pkg/runtime/persistence_observer_bench_test.go @@ -3,6 +3,7 @@ package runtime import ( "context" "database/sql" + "log/slog" "strconv" "sync/atomic" "testing" @@ -157,6 +158,12 @@ func BenchmarkPersistenceObserver_StreamingChunks_SQLite(b *testing.B) { func openBenchSQLiteStore(b *testing.B) *session.SQLiteSessionStore { b.Helper() + // NewSQLiteSessionStoreFromDB logs migration INFO on every harness + // re-invocation; discard it so -bench output stays machine-parseable. + prev := slog.Default() + slog.SetDefault(slog.New(slog.DiscardHandler)) + b.Cleanup(func() { slog.SetDefault(prev) }) + db, err := sql.Open("sqlite", ":memory:") require.NoError(b, err) b.Cleanup(func() { _ = db.Close() })