diff --git a/cmd/odek/audit_serve_test.go b/cmd/odek/audit_serve_test.go index ec80784..17dc612 100644 --- a/cmd/odek/audit_serve_test.go +++ b/cmd/odek/audit_serve_test.go @@ -130,8 +130,7 @@ func TestAudit_WebSocketInvalidModelRejectedEarly(t *testing.T) { t.Fatalf("session.NewStore: %v", err) } ln, mux := buildServeMux(t, store) - defer ln.Close() - go func() { _ = serveOnListener(ln, mux) }() + defer startServeTest(t, ln, mux)() waitForHTTP(t, ln.Addr().String()) wsUpgradeLimiter.reset() diff --git a/cmd/odek/next_security_vulnerabilities_test.go b/cmd/odek/next_security_vulnerabilities_test.go index 9c35cd1..941d0bc 100644 --- a/cmd/odek/next_security_vulnerabilities_test.go +++ b/cmd/odek/next_security_vulnerabilities_test.go @@ -242,9 +242,8 @@ func TestServe_CSRF_AllowsLocalhostOrigin(t *testing.T) { func TestServe_API_RequiresServeToken(t *testing.T) { store := newTestSessionStore(t) ln, mux := buildServeMux(t, store) - defer ln.Close() - go serveOnListener(ln, mux) + defer startServeTest(t, ln, mux)() waitForHTTP(t, ln.Addr().String()) resp, err := http.Get("http://" + ln.Addr().String() + "/api/sessions") @@ -260,13 +259,12 @@ func TestServe_API_RequiresServeToken(t *testing.T) { func TestServe_API_RequiresLocalHost(t *testing.T) { store := newTestSessionStore(t) ln, mux := buildServeMux(t, store) - defer ln.Close() testTokenMu.Lock() token := testLastToken testTokenMu.Unlock() - go serveOnListener(ln, mux) + defer startServeTest(t, ln, mux)() waitForHTTP(t, ln.Addr().String()) req, _ := http.NewRequest(http.MethodGet, "http://"+ln.Addr().String()+"/api/sessions", nil) @@ -285,13 +283,12 @@ func TestServe_API_RequiresLocalHost(t *testing.T) { func TestServe_API_AcceptsServeTokenHeader(t *testing.T) { store := newTestSessionStore(t) ln, mux := buildServeMux(t, store) - defer ln.Close() testTokenMu.Lock() token := testLastToken testTokenMu.Unlock() - go serveOnListener(ln, mux) + defer startServeTest(t, ln, mux)() waitForHTTP(t, ln.Addr().String()) req, _ := http.NewRequest(http.MethodGet, "http://"+ln.Addr().String()+"/api/sessions", nil) @@ -309,13 +306,12 @@ func TestServe_API_AcceptsServeTokenHeader(t *testing.T) { func TestServe_API_AcceptsServeTokenCookie(t *testing.T) { store := newTestSessionStore(t) ln, mux := buildServeMux(t, store) - defer ln.Close() testTokenMu.Lock() token := testLastToken testTokenMu.Unlock() - go serveOnListener(ln, mux) + defer startServeTest(t, ln, mux)() waitForHTTP(t, ln.Addr().String()) req, _ := http.NewRequest(http.MethodGet, "http://"+ln.Addr().String()+"/api/sessions", nil) diff --git a/cmd/odek/serve_api_test.go b/cmd/odek/serve_api_test.go index 4f20a4b..68a5a65 100644 --- a/cmd/odek/serve_api_test.go +++ b/cmd/odek/serve_api_test.go @@ -885,9 +885,8 @@ func TestServe_E2E_NoRepetitiveResponses(t *testing.T) { store := newTestSessionStore(t) ln, mux := buildServeMux(t, store) - defer ln.Close() - go serveOnListener(ln, mux) + defer startServeTest(t, ln, mux)() waitForHTTP(t, ln.Addr().String()) conn := dialTestWS(t, ln.Addr().String()) @@ -964,9 +963,8 @@ func TestServe_E2E_SessionMessagesStoredWithoutSystemInjections(t *testing.T) { // Build a mux that also includes the session-by-ID route so we can inspect it. ln, mux := buildServeMuxWithSessionByID(t, store) - defer ln.Close() - go serveOnListener(ln, mux) + defer startServeTest(t, ln, mux)() waitForHTTP(t, ln.Addr().String()) conn := dialTestWS(t, ln.Addr().String()) diff --git a/cmd/odek/serve_api_v2_test.go b/cmd/odek/serve_api_v2_test.go index f9940b7..68565a0 100644 --- a/cmd/odek/serve_api_v2_test.go +++ b/cmd/odek/serve_api_v2_test.go @@ -573,8 +573,7 @@ func TestServe_E2E_ServerInfoHelloAndPingPong(t *testing.T) { store := newTestSessionStore(t) ln, mux := buildServeMuxV2(t, store, nil) - defer ln.Close() - go func() { _ = serveOnListener(ln, mux) }() + defer startServeTest(t, ln, mux)() waitForHTTP(t, ln.Addr().String()) wsUpgradeLimiter.reset() @@ -613,8 +612,7 @@ func TestServe_E2E_ServerKeepalive(t *testing.T) { store := newTestSessionStore(t) ln, mux := buildServeMuxV2(t, store, nil) - defer ln.Close() - go func() { _ = serveOnListener(ln, mux) }() + defer startServeTest(t, ln, mux)() waitForHTTP(t, ln.Addr().String()) wsUpgradeLimiter.reset() @@ -644,8 +642,7 @@ func TestServe_E2E_SessionSwitchMessage(t *testing.T) { } ln, mux := buildServeMuxV2(t, store, nil) - defer ln.Close() - go func() { _ = serveOnListener(ln, mux) }() + defer startServeTest(t, ln, mux)() waitForHTTP(t, ln.Addr().String()) wsUpgradeLimiter.reset() @@ -684,8 +681,7 @@ func TestServe_E2E_WSCancelMessage(t *testing.T) { } ln, mux := buildServeMuxV2(t, store, nil) - defer ln.Close() - go func() { _ = serveOnListener(ln, mux) }() + defer startServeTest(t, ln, mux)() waitForHTTP(t, ln.Addr().String()) wsUpgradeLimiter.reset() @@ -729,8 +725,7 @@ func TestServe_E2E_StreamDeltas(t *testing.T) { ln, mux := buildServeMuxV2(t, store, func(rc *config.ResolvedConfig) { rc.Stream = true }) - defer ln.Close() - go func() { _ = serveOnListener(ln, mux) }() + defer startServeTest(t, ln, mux)() waitForHTTP(t, ln.Addr().String()) wsUpgradeLimiter.reset() @@ -796,8 +791,7 @@ func TestServe_E2E_StreamFallbackKeepsBulkPath(t *testing.T) { ln, mux := buildServeMuxV2(t, store, func(rc *config.ResolvedConfig) { rc.Stream = true }) - defer ln.Close() - go func() { _ = serveOnListener(ln, mux) }() + defer startServeTest(t, ln, mux)() waitForHTTP(t, ln.Addr().String()) wsUpgradeLimiter.reset() diff --git a/cmd/odek/serve_approval_test.go b/cmd/odek/serve_approval_test.go index ad1749b..b5142a3 100644 --- a/cmd/odek/serve_approval_test.go +++ b/cmd/odek/serve_approval_test.go @@ -89,9 +89,8 @@ func TestServe_E2E_ApprovalRoundTrip(t *testing.T) { store := newTestSessionStore(t) ln, mux := buildServeMuxPromptAll(t, store) - defer ln.Close() - go func() { _ = serveOnListener(ln, mux) }() + defer startServeTest(t, ln, mux)() waitForHTTP(t, ln.Addr().String()) conn := dialTestWS(t, ln.Addr().String()) diff --git a/cmd/odek/serve_buffer_bleed_test.go b/cmd/odek/serve_buffer_bleed_test.go index e06e404..652a964 100644 --- a/cmd/odek/serve_buffer_bleed_test.go +++ b/cmd/odek/serve_buffer_bleed_test.go @@ -81,9 +81,7 @@ func TestServe_E2E_PromptPathSessionSwitch_ClearsStaleBuffer(t *testing.T) { } ln, mux := buildServeMux(t, store) - defer ln.Close() - errCh := make(chan error, 1) - go func() { errCh <- serveOnListener(ln, mux) }() + defer startServeTest(t, ln, mux)() waitForHTTP(t, ln.Addr().String()) wsUpgradeLimiter.reset() diff --git a/cmd/odek/serve_cancel_test.go b/cmd/odek/serve_cancel_test.go index 06c418e..1c2cc7c 100644 --- a/cmd/odek/serve_cancel_test.go +++ b/cmd/odek/serve_cancel_test.go @@ -161,8 +161,7 @@ func TestServe_E2E_WSCancelInterruptsApprovalWait(t *testing.T) { store := newTestSessionStore(t) ln, mux := buildServeMuxPromptAll(t, store) - defer ln.Close() - go func() { _ = serveOnListener(ln, mux) }() + defer startServeTest(t, ln, mux)() waitForHTTP(t, ln.Addr().String()) wsUpgradeLimiter.reset() @@ -199,8 +198,7 @@ func TestServe_E2E_RESTCancelInterruptsApprovalWait(t *testing.T) { store := newTestSessionStore(t) ln, mux := buildServeMuxPromptAll(t, store) - defer ln.Close() - go func() { _ = serveOnListener(ln, mux) }() + defer startServeTest(t, ln, mux)() waitForHTTP(t, ln.Addr().String()) wsUpgradeLimiter.reset() @@ -335,8 +333,7 @@ func TestServe_E2E_WSCancelRunningPromptTerminalEvent(t *testing.T) { store := newTestSessionStore(t) ln, mux := buildServeMuxV2(t, store, nil) - defer ln.Close() - go func() { _ = serveOnListener(ln, mux) }() + defer startServeTest(t, ln, mux)() waitForHTTP(t, ln.Addr().String()) wsUpgradeLimiter.reset() @@ -384,8 +381,7 @@ func TestServe_E2E_ApprovalWorksAfterCancel(t *testing.T) { store := newTestSessionStore(t) ln, mux := buildServeMuxPromptAll(t, store) - defer ln.Close() - go func() { _ = serveOnListener(ln, mux) }() + defer startServeTest(t, ln, mux)() waitForHTTP(t, ln.Addr().String()) wsUpgradeLimiter.reset() @@ -491,9 +487,8 @@ func TestServe_E2E_CancelDuringSetupWindowHonored(t *testing.T) { gate := &gatingResolver{entered: make(chan struct{}, 1), release: make(chan struct{})} resReg := resource.NewRegistry(gate) ln, mux := buildServeMuxWithResolvers(t, store, resReg) - defer ln.Close() + defer startServeTest(t, ln, mux)() defer close(gate.release) - go func() { _ = serveOnListener(ln, mux) }() waitForHTTP(t, ln.Addr().String()) wsUpgradeLimiter.reset() diff --git a/cmd/odek/serve_lifecycle_test.go b/cmd/odek/serve_lifecycle_test.go new file mode 100644 index 0000000..26206aa --- /dev/null +++ b/cmd/odek/serve_lifecycle_test.go @@ -0,0 +1,31 @@ +package main + +import ( + "errors" + "net" + "net/http" + "testing" + "time" +) + +// startServeTest starts the real server and returns a stop function that waits +// for its handlers and session writes to finish. Defer the returned function +// before opening clients, so shutdown precedes environment and temp-dir cleanup. +func startServeTest(t *testing.T, ln net.Listener, mux *http.ServeMux) func() { + t.Helper() + done := make(chan error, 1) + go func() { done <- serveOnListener(ln, mux) }() + return func() { + t.Helper() + http.DefaultClient.CloseIdleConnections() + _ = ln.Close() + select { + case err := <-done: + if err != nil && !errors.Is(err, net.ErrClosed) { + t.Errorf("serve shutdown: %v", err) + } + case <-time.After(20 * time.Second): + t.Fatal("serve shutdown did not finish before test cleanup") + } + } +} diff --git a/cmd/odek/serve_surface_fixes_test.go b/cmd/odek/serve_surface_fixes_test.go index 7208af8..cc54d32 100644 --- a/cmd/odek/serve_surface_fixes_test.go +++ b/cmd/odek/serve_surface_fixes_test.go @@ -272,8 +272,7 @@ func TestServe_E2E_PingPongUsesConfigSnapshotNotLiveModel(t *testing.T) { store := newTestSessionStore(t) ln, mux := buildServeMuxV2(t, store, func(rc *config.ResolvedConfig) { rc.Model = "initial-model" }) - defer ln.Close() - go func() { _ = serveOnListener(ln, mux) }() + defer startServeTest(t, ln, mux)() waitForHTTP(t, ln.Addr().String()) wsUpgradeLimiter.reset() diff --git a/cmd/odek/serve_test.go b/cmd/odek/serve_test.go index d2a658f..0a5e8c5 100644 --- a/cmd/odek/serve_test.go +++ b/cmd/odek/serve_test.go @@ -6,7 +6,6 @@ import ( "encoding/json" "fmt" "io" - "io/fs" "net" "net/http" "net/http/httptest" @@ -583,34 +582,8 @@ func TestServe_E2E_WebSocketPipeline(t *testing.T) { mux.HandleFunc("/api/resources", handleResourceSearch(resourceReg)) mux.HandleFunc("/api/sessions", handleSessionList(store)) - // Start serving on the pre-created listener in a goroutine - errCh := make(chan error, 1) - go func() { - errCh <- serveOnListener(ln, mux) - }() - defer ln.Close() - - // Wait for server to be ready - var httpReady bool - for i := 0; i < 20; i++ { - time.Sleep(250 * time.Millisecond) - resp, err := http.Get("http://" + addr + "/") - if err == nil && resp.StatusCode == 200 { - resp.Body.Close() - httpReady = true - break - } - } - if !httpReady { - // Check if serveCmd returned an error immediately - select { - case err := <-errCh: - t.Fatalf("server exited before ready: %v", err) - default: - t.Fatal("server not ready after 5s") - } - } - defer ln.Close() + defer startServeTest(t, ln, mux)() + waitForHTTP(t, addr) // 1. Connect via WebSocket // Reset the per-IP upgrade limiter so this E2E test is not throttled by @@ -703,13 +676,9 @@ func TestServe_E2E_FullWebUIFlow(t *testing.T) { // 4. Build the real odek serve mux with mock config ln, mux := buildServeMux(t, store) - defer ln.Close() // 5. Start serving - errCh := make(chan error, 1) - go func() { - errCh <- serveOnListener(ln, mux) - }() + defer startServeTest(t, ln, mux)() // 6. Wait for HTTP ready waitForHTTP(t, ln.Addr().String()) @@ -851,37 +820,6 @@ func newTestSessionStore(t *testing.T) *session.Store { return store } -// waitForOdekTreeQuiet blocks until a write newer than `since` has been seen -// under dir and no file has been modified for quietWindow (bounded by -// timeout). The serve handler issues its final store.Save — session file plus -// vector-index files — and runs the learn loop AFTER sending the "done" -// event (serve.go), so a test that returns on "done" can race those writes -// and fail t.TempDir cleanup with "directory not empty". Best-effort: a -// timeout is not a test failure. -func waitForOdekTreeQuiet(dir string, since time.Time, quietWindow, timeout time.Duration) { - deadline := time.Now().Add(timeout) - sawWrite := false - for { - var latest time.Time - _ = filepath.WalkDir(dir, func(_ string, d fs.DirEntry, err error) error { - if err != nil || d.IsDir() { - return nil - } - if fi, err := d.Info(); err == nil && fi.ModTime().After(latest) { - latest = fi.ModTime() - } - return nil - }) - if latest.After(since) { - sawWrite = true - } - if (sawWrite && time.Since(latest) >= quietWindow) || time.Now().After(deadline) { - return - } - time.Sleep(50 * time.Millisecond) - } -} - // mockLLM creates an httptest.Server that handles both the model discovery call // (GET /models — called once by odek.New during agent initialization) and chat // completion requests. The chat handler receives a 1-based callCount that @@ -1248,10 +1186,8 @@ func TestServe_E2E_MultiToolCall(t *testing.T) { store := newTestSessionStore(t) ln, mux := buildServeMux(t, store) - defer ln.Close() - errCh := make(chan error, 1) - go func() { errCh <- serveOnListener(ln, mux) }() + defer startServeTest(t, ln, mux)() waitForHTTP(t, ln.Addr().String()) conn := dialTestWS(t, ln.Addr().String()) @@ -1369,10 +1305,8 @@ func TestServe_E2E_TokenStats(t *testing.T) { store := newTestSessionStore(t) ln, mux := buildServeMux(t, store) - defer ln.Close() - errCh := make(chan error, 1) - go func() { errCh <- serveOnListener(ln, mux) }() + defer startServeTest(t, ln, mux)() waitForHTTP(t, ln.Addr().String()) conn := dialTestWS(t, ln.Addr().String()) @@ -1482,7 +1416,6 @@ statsCheck: conn.SetReadDeadline(time.Now().Add(15 * time.Second)) var done2 map[string]any - var done2At time.Time for i := 0; i < 15; i++ { var raw []byte if err := golangws.Message.Receive(conn, &raw); err != nil { @@ -1496,7 +1429,6 @@ statsCheck: } if evt["type"] == "done" { done2 = evt - done2At = time.Now() goto sessionCheck } if evt["type"] == "error" { @@ -1506,10 +1438,6 @@ statsCheck: t.Fatal("did not receive second done event") sessionCheck: - // The final store.Save runs after "done" is sent — wait for the writes - // to settle so t.TempDir cleanup doesn't race them. - waitForOdekTreeQuiet(filepath.Join(os.Getenv("HOME"), ".odek"), done2At, 300*time.Millisecond, 10*time.Second) - if done2 == nil { t.Fatal("second done event not found") } @@ -1563,10 +1491,8 @@ func TestServe_E2E_UsageEvents(t *testing.T) { store := newTestSessionStore(t) ln, mux := buildServeMux(t, store) - defer ln.Close() - errCh := make(chan error, 1) - go func() { errCh <- serveOnListener(ln, mux) }() + defer startServeTest(t, ln, mux)() waitForHTTP(t, ln.Addr().String()) conn := dialTestWS(t, ln.Addr().String()) @@ -1580,7 +1506,6 @@ func TestServe_E2E_UsageEvents(t *testing.T) { conn.SetReadDeadline(time.Now().Add(15 * time.Second)) var usages []map[string]any - var doneAt time.Time for i := 0; i < 20; i++ { var raw []byte if err := golangws.Message.Receive(conn, &raw); err != nil { @@ -1596,7 +1521,6 @@ func TestServe_E2E_UsageEvents(t *testing.T) { case "usage": usages = append(usages, evt) case "done": - doneAt = time.Now() goto usageCheck case "error": t.Fatalf("unexpected error: %v", evt["message"]) @@ -1605,11 +1529,6 @@ func TestServe_E2E_UsageEvents(t *testing.T) { t.Fatal("did not receive done event") usageCheck: - // The final store.Save (session file + vector index) and the learn loop - // run after "done" is sent — wait for those writes to settle so - // t.TempDir cleanup doesn't race them ("directory not empty"). - waitForOdekTreeQuiet(filepath.Join(os.Getenv("HOME"), ".odek"), doneAt, 300*time.Millisecond, 10*time.Second) - if len(usages) < 2 { t.Fatalf("got %d usage events before done, want at least 2 (one per LLM turn)", len(usages)) } @@ -1652,10 +1571,8 @@ func TestServe_E2E_LiveToolEvents(t *testing.T) { store := newTestSessionStore(t) ln, mux := buildServeMux(t, store) - defer ln.Close() - errCh := make(chan error, 1) - go func() { errCh <- serveOnListener(ln, mux) }() + defer startServeTest(t, ln, mux)() waitForHTTP(t, ln.Addr().String()) conn := dialTestWS(t, ln.Addr().String()) @@ -1671,7 +1588,6 @@ func TestServe_E2E_LiveToolEvents(t *testing.T) { var toolCallTime, toolResultTime, doneTime int var eventOrder []string - var doneAt time.Time // Additive protocol frames may precede done; the read deadline bounds the wait. for i := 0; i < 100; i++ { @@ -1699,7 +1615,6 @@ func TestServe_E2E_LiveToolEvents(t *testing.T) { } case "done": doneTime = i - doneAt = time.Now() goto doneCheck case "error": t.Fatalf("unexpected error: %v", evt["message"]) @@ -1707,10 +1622,6 @@ func TestServe_E2E_LiveToolEvents(t *testing.T) { } doneCheck: - // The final store.Save runs after "done" is sent — wait for the writes - // to settle so t.TempDir cleanup doesn't race them. - waitForOdekTreeQuiet(filepath.Join(os.Getenv("HOME"), ".odek"), doneAt, 300*time.Millisecond, 10*time.Second) - t.Logf("Event order: %v", eventOrder) if toolCallTime == 0 { @@ -2024,11 +1935,10 @@ func TestServe_E2E_CancelWithMockLLM(t *testing.T) { store := newTestSessionStore(t) ln, mux := buildServeMux(t, store) - defer ln.Close() - errCh := make(chan error, 1) - go func() { errCh <- serveOnListener(ln, mux) }() + defer startServeTest(t, ln, mux)() waitForHTTP(t, ln.Addr().String()) + wsUpgradeLimiter.reset() conn := dialTestWS(t, ln.Addr().String()) defer conn.Close() @@ -2153,11 +2063,10 @@ func TestServe_Cancel_CannotCrossSessions(t *testing.T) { store := newTestSessionStore(t) ln, mux := buildServeMux(t, store) - defer ln.Close() - errCh := make(chan error, 1) - go func() { errCh <- serveOnListener(ln, mux) }() + defer startServeTest(t, ln, mux)() waitForHTTP(t, ln.Addr().String()) + wsUpgradeLimiter.reset() // Victim connection: starts first, hits the slow LLM path. victimConn := dialTestWS(t, ln.Addr().String()) @@ -2375,12 +2284,8 @@ func TestServe_E2E_AttachmentsWrappedAsUntrusted(t *testing.T) { store := newTestSessionStore(t) ln, mux := buildServeMux(t, store) - defer ln.Close() - errCh := make(chan error, 1) - go func() { - errCh <- serveOnListener(ln, mux) - }() + defer startServeTest(t, ln, mux)() waitForHTTP(t, ln.Addr().String()) conn := dialTestWS(t, ln.Addr().String()) @@ -2467,10 +2372,8 @@ func TestServe_E2E_AttachmentsWrappedAsUntrusted(t *testing.T) { func TestServe_CSRF_TokenRequired(t *testing.T) { store := newTestSessionStore(t) ln, mux := buildServeMux(t, store) - defer ln.Close() - errCh := make(chan error, 1) - go func() { errCh <- serveOnListener(ln, mux) }() + defer startServeTest(t, ln, mux)() waitForHTTP(t, ln.Addr().String()) addr := ln.Addr().String() @@ -2565,10 +2468,8 @@ func TestServe_E2E_PromptSizeCap(t *testing.T) { store := newTestSessionStore(t) ln, mux := buildServeMux(t, store) - defer ln.Close() - errCh := make(chan error, 1) - go func() { errCh <- serveOnListener(ln, mux) }() + defer startServeTest(t, ln, mux)() waitForHTTP(t, ln.Addr().String()) conn := dialTestWS(t, ln.Addr().String()) @@ -2616,10 +2517,8 @@ func TestServe_E2E_InvalidModelIDRejected(t *testing.T) { store := newTestSessionStore(t) ln, mux := buildServeMux(t, store) - defer ln.Close() - errCh := make(chan error, 1) - go func() { errCh <- serveOnListener(ln, mux) }() + defer startServeTest(t, ln, mux)() waitForHTTP(t, ln.Addr().String()) conn := dialTestWS(t, ln.Addr().String()) diff --git a/cmd/odek/turn_started_test.go b/cmd/odek/turn_started_test.go index e3efe1b..e48727b 100644 --- a/cmd/odek/turn_started_test.go +++ b/cmd/odek/turn_started_test.go @@ -88,8 +88,7 @@ func startTurnServer(t *testing.T) *golangws.Conn { store := newTestSessionStore(t) ln, mux := buildServeMuxV2(t, store, nil) - t.Cleanup(func() { ln.Close() }) - go func() { _ = serveOnListener(ln, mux) }() + t.Cleanup(startServeTest(t, ln, mux)) waitForHTTP(t, ln.Addr().String()) wsUpgradeLimiter.reset()