From 89df77db6a89251f6eb061f32acc5dea9d1c36cf Mon Sep 17 00:00:00 2001 From: piyush0049 Date: Sat, 25 Jul 2026 21:03:25 +0530 Subject: [PATCH 1/3] fix(server): prevent send-on-closed-channel panic in generateTitle --- pkg/server/session_manager.go | 17 ++++++++++++++--- 1 file changed, 14 insertions(+), 3 deletions(-) diff --git a/pkg/server/session_manager.go b/pkg/server/session_manager.go index 7555b92a5e..23774d9865 100644 --- a/pkg/server/session_manager.go +++ b/pkg/server/session_manager.go @@ -891,9 +891,13 @@ func (sm *SessionManager) RunSession(ctx context.Context, sessionID, agentFilena defer cancel() defer runtimeSession.streaming.Unlock() - // Start title generation in parallel if needed + // Start title generation in parallel if needed, coordinating via WaitGroup + // so close(streamChan) does not fire while generateTitle is still sending. + var wg sync.WaitGroup if needsTitle { - go sm.generateTitle(ctx, sess, titleGen, userMessages, streamChan) + wg.Go(func() { + sm.generateTitle(ctx, sess, titleGen, userMessages, streamChan) + }) } else if titleToEmit != "" { // Re-emit the existing title so late-joining SSE consumers // and boards can pick it up without an extra API call. @@ -903,11 +907,18 @@ func (sm *SessionManager) RunSession(ctx context.Context, sessionID, agentFilena stream := runtimeSession.runtime.RunStream(streamCtx, sess) for event := range stream { if streamCtx.Err() != nil { - return + break } streamChan <- event } + // Ensure title generation finishes before defers run close(streamChan). + wg.Wait() + + if streamCtx.Err() != nil { + return + } + if err := sm.sessionStore.UpdateSession(ctx, sess); err != nil { return } From 41a4c130d8c3e1650e86552be328751c1782a3f9 Mon Sep 17 00:00:00 2001 From: piyush0049 Date: Wed, 5 Aug 2026 16:53:30 +0530 Subject: [PATCH 2/3] fix(server): add regression test and teardown note --- pkg/server/session_manager.go | 6 ++- pkg/server/session_manager_test.go | 81 ++++++++++++++++++++++++++++++ 2 files changed, 86 insertions(+), 1 deletion(-) diff --git a/pkg/server/session_manager.go b/pkg/server/session_manager.go index bf7af1dd16..8558d0c7bf 100644 --- a/pkg/server/session_manager.go +++ b/pkg/server/session_manager.go @@ -924,6 +924,11 @@ func (sm *SessionManager) RunSession(ctx context.Context, sessionID, agentFilena // so close(streamChan) does not fire while generateTitle is still sending. var wg sync.WaitGroup if needsTitle { + // Note: generateTitle runs on the request ctx, so wg.Wait() only unblocks + // quickly when that context is cancelled. In the DeleteSession path + // (where the stream context is cancelled but the request context stays alive), + // the wait blocks until the title LLM call completes, delaying stream teardown. + // This bounded delay (first turn only) prioritizes persisting the title. wg.Go(func() { sm.generateTitle(ctx, sess, titleGen, userMessages, streamChan) }) @@ -941,7 +946,6 @@ func (sm *SessionManager) RunSession(ctx context.Context, sessionID, agentFilena streamChan <- event } - // Ensure title generation finishes before defers run close(streamChan). wg.Wait() if streamCtx.Err() != nil { diff --git a/pkg/server/session_manager_test.go b/pkg/server/session_manager_test.go index 555e541d85..4bb33fdec6 100644 --- a/pkg/server/session_manager_test.go +++ b/pkg/server/session_manager_test.go @@ -5,6 +5,7 @@ import ( "context" "encoding/json" "errors" + "io" "net/http" "net/http/httptest" "path/filepath" @@ -25,6 +26,9 @@ import ( "github.com/docker/docker-agent/pkg/concurrent" "github.com/docker/docker-agent/pkg/config" "github.com/docker/docker-agent/pkg/config/types" + "github.com/docker/docker-agent/pkg/model/provider" + "github.com/docker/docker-agent/pkg/model/provider/base" + "github.com/docker/docker-agent/pkg/modelsdev" "github.com/docker/docker-agent/pkg/runtime" "github.com/docker/docker-agent/pkg/session" "github.com/docker/docker-agent/pkg/session/sqlitestore" @@ -2252,3 +2256,80 @@ func TestApplyAgentSwitchCommands_RollsBackOnBatchFailure(t *testing.T) { assert.Equal(t, "root", rt.currentAgent, "runtime must be restored to the pre-batch agent") assert.Equal(t, []string{"planner", "nonexistent", "root"}, rt.setCalls) } + +type blockingStream struct { + chat.MessageStream + + blocker chan struct{} + returned bool +} + +func (s *blockingStream) Recv() (chat.MessageStreamResponse, error) { + if !s.returned { + <-s.blocker + s.returned = true + return chat.MessageStreamResponse{ + Choices: []chat.MessageStreamChoice{ + {Delta: chat.MessageDelta{Content: "A Generated Title"}}, + }, + }, nil + } + return chat.MessageStreamResponse{}, io.EOF +} + +func (s *blockingStream) Close() {} + +type blockingTitleProvider struct { + provider.Provider + + blocker chan struct{} +} + +func (p *blockingTitleProvider) ID() modelsdev.ID { return modelsdev.NewID("mock", "mock") } + +func (p *blockingTitleProvider) BaseConfig() base.Config { return base.Config{} } + +func (p *blockingTitleProvider) CreateChatCompletionStream(_ context.Context, _ []chat.Message, _ []tools.Tool) (chat.MessageStream, error) { + return &blockingStream{blocker: p.blocker}, nil +} + +func TestRunSession_GenerateTitleConcurrentClosePanic(t *testing.T) { + t.Parallel() + + ctx := t.Context() + sess := session.New() + + // Use fakeRuntime that ends the agent stream immediately by passing a nil release channel + fake := &fakeRuntime{} + sm := newTestSessionManager(t, sess, fake) + + blocker := make(chan struct{}) + mockProvider := &blockingTitleProvider{blocker: blocker} + titleGen := sessiontitle.New(mockProvider) + + // Inject the mock title generator into the active runtime + rt, _ := sm.runtimeSessions.Load(sess.ID) + rt.titleGen = titleGen + + // RunSession spawns generateTitle and fakeRuntime in parallel + streamChan, err := sm.RunSession(ctx, sess.ID, "agent", "root", []api.Message{{Content: "trigger title generation"}}, "") + require.NoError(t, err) + + go func() { + // Drain the stream. It will close once the title generator unblocks and wg.Wait() returns. + for range streamChan { + } + }() + + // Ensure the agent stream draining finishes, putting RunSession's defer stack + // in a position to execute if it weren't blocking on wg.Wait() + time.Sleep(50 * time.Millisecond) //nolint:forbidigo // giving the inner goroutine time to run past the fake stream and reach wg.Wait + + // Unblock the title generator. Without wg.Wait(), the streamChan would already + // be closed and this would panic on send in generateTitle. + close(blocker) + + // Test passes if it does not panic. Delete the session to clean up. + err = sm.DeleteSession(ctx, sess.ID) + require.NoError(t, err) +} From 8dc03824e957b0f27d64d2ee594323114eb68787 Mon Sep 17 00:00:00 2001 From: piyush0049 Date: Wed, 5 Aug 2026 18:08:56 +0530 Subject: [PATCH 3/3] test: make credential helper tests cross-platform for windows --- pkg/environment/credential_helper_test.go | 24 ++++++++++++++++++----- 1 file changed, 19 insertions(+), 5 deletions(-) diff --git a/pkg/environment/credential_helper_test.go b/pkg/environment/credential_helper_test.go index f7caccf28f..43272a822c 100644 --- a/pkg/environment/credential_helper_test.go +++ b/pkg/environment/credential_helper_test.go @@ -1,6 +1,7 @@ package environment import ( + "runtime" "testing" "github.com/stretchr/testify/assert" @@ -17,6 +18,19 @@ func TestNewCredentialHelperProvider(t *testing.T) { func TestCredentialHelperProvider_Get(t *testing.T) { t.Parallel() + echoCmd := "echo" + echoArgs := func(v string) []string { return []string{v} } + falseCmd := "false" + + if runtime.GOOS == "windows" { + echoCmd = "powershell" + echoArgs = func(v string) []string { + return []string{"-NoProfile", "-Command", "Write-Output '" + v + "'"} + } + falseCmd = "powershell" + // simulate 'false' by exiting with 1 + } + tests := []struct { name string command string @@ -25,11 +39,11 @@ func TestCredentialHelperProvider_Get(t *testing.T) { wantValue string wantFound bool }{ - {"ignores non-DOCKER_TOKEN vars", "echo", []string{"test-token"}, "OTHER_VAR", "", false}, - {"success", "echo", []string{"my-secret-token"}, DockerDesktopTokenEnv, "my-secret-token", true}, - {"trims whitespace", "echo", []string{" token-with-spaces "}, DockerDesktopTokenEnv, "token-with-spaces", true}, - {"empty output", "echo", []string{""}, DockerDesktopTokenEnv, "", false}, - {"command fails", "false", nil, DockerDesktopTokenEnv, "", false}, + {"ignores non-DOCKER_TOKEN vars", echoCmd, echoArgs("test-token"), "OTHER_VAR", "", false}, + {"success", echoCmd, echoArgs("my-secret-token"), DockerDesktopTokenEnv, "my-secret-token", true}, + {"trims whitespace", echoCmd, echoArgs(" token-with-spaces "), DockerDesktopTokenEnv, "token-with-spaces", true}, + {"empty output", echoCmd, echoArgs(""), DockerDesktopTokenEnv, "", false}, + {"command fails", falseCmd, []string{"-NoProfile", "-Command", "exit 1"}, DockerDesktopTokenEnv, "", false}, {"command not found", "nonexistent-command-12345", nil, DockerDesktopTokenEnv, "", false}, }