From 01ae13b5db1c46ef9ab5b876357d855c9aacc97c Mon Sep 17 00:00:00 2001 From: Juan Antonio Osorio Date: Tue, 28 Jul 2026 07:38:45 +0000 Subject: [PATCH 1/2] Fix thv list race that poisons starting workloads MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ListWorkloads snapshotted the runtime before the status files. A workload whose file flips to running between the two snapshots — every thv run does this, since the detached child starts the container and then writes running — appears as "file running, runtime missing", and handleRuntimeMissing persists unhealthy into its status file. Nothing recovers a non-running file status against a live container, so one concurrent `thv list` permanently marks a healthy, freshly started workload unhealthy. This is user-facing state corruption; it also caused the E2E Tests Core (vmcp) failure on main (run 30335474873), where a backend that had demonstrably started and stayed alive was reported "not found in runtime" for two minutes. Read the files before the runtime: a file that says running at time T implies its container existed before T, so a later runtime miss is genuine. The old order was introduced without stated rationale in e6e29dc3. Guard the unhealthy write itself: re-check under the file lock and skip when the file is gone or no longer says running/stopped. Reordering alone would let a list that brackets a whole `thv rm` recreate the just-deleted status file as an unhealthy ghost, because SetWorkloadStatus creates missing files. Give waitForIsolatedMCPServer a state dump through its isolated env on timeout; the upgrade-spec failure on the same run was undiagnosable because the generic timeout message carries no status context. Refs #4432 — the same write-on-read class; this fixes the startup trigger but thv list still writes on read, so that issue stays open. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_01GvBj7nNWqDQ3rKYNeY6aam --- pkg/workloads/statuses/file_status.go | 62 +++++++++-- pkg/workloads/statuses/file_status_test.go | 123 +++++++++++++++++++++ test/e2e/upgrade_e2e_test.go | 34 +++++- 3 files changed, 207 insertions(+), 12 deletions(-) diff --git a/pkg/workloads/statuses/file_status.go b/pkg/workloads/statuses/file_status.go index b5f77e4f9c..d118476759 100644 --- a/pkg/workloads/statuses/file_status.go +++ b/pkg/workloads/statuses/file_status.go @@ -279,16 +279,26 @@ func (f *fileStatusManager) ListWorkloads(ctx context.Context, listAll bool, lab return nil, fmt.Errorf("failed to parse label filters: %w", err) } - // Get workloads from runtime - runtimeContainers, err := f.runtime.ListWorkloads(ctx) + // Get workloads from files BEFORE snapshotting the runtime. The order is + // load-bearing: a workload's status file only flips to `running` after its + // container has been created and started, so a file that says `running` in + // a snapshot taken at time T implies the container existed before T. If a + // later runtime snapshot lacks it, the container is genuinely gone and + // handleRuntimeMissing's unhealthy verdict is correct. With the snapshots + // the other way round, a workload that starts between the runtime snapshot + // and the file read looks like "file running, runtime missing" and gets its + // status file permanently overwritten as unhealthy — a healthy, freshly + // started workload poisoned by a concurrent `thv list` (see #4432 for the + // broader write-on-read hazard in this path). + fileWorkloadsWithPID, err := f.getWorkloadsFromFiles() if err != nil { - return nil, fmt.Errorf("failed to list workloads from runtime: %w", err) + return nil, fmt.Errorf("failed to get workloads from files: %w", err) } - // Get workloads from files - fileWorkloadsWithPID, err := f.getWorkloadsFromFiles() + // Get workloads from runtime + runtimeContainers, err := f.runtime.ListWorkloads(ctx) if err != nil { - return nil, fmt.Errorf("failed to get workloads from files: %w", err) + return nil, fmt.Errorf("failed to list workloads from runtime: %w", err) } // TODO: Fetch the runconfig if present to populate additional fields like package, tool type, group etc. @@ -973,9 +983,12 @@ func (f *fileStatusManager) handleRuntimeMissing( if fileWorkload.Status == rt.WorkloadStatusRunning || fileWorkload.Status == rt.WorkloadStatusStopped { // The workload cannot be running or stopped if the runtime container is not found contextMsg := fmt.Sprintf("workload %s not found in runtime, marking as unhealthy", workloadName) - if err := f.SetWorkloadStatus(ctx, workloadName, rt.WorkloadStatusUnhealthy, contextMsg); err != nil { + if err := f.setUnhealthyIfStillPresent(ctx, workloadName, contextMsg); err != nil { return core.Workload{}, err } + // The returned view reflects what this list observed regardless of + // whether the guarded write persisted anything (the file may have been + // removed or repurposed by a concurrent operation in the meantime). fileWorkload.Status = rt.WorkloadStatusUnhealthy } @@ -983,6 +996,41 @@ func (f *fileStatusManager) handleRuntimeMissing( return fileWorkload, nil } +// setUnhealthyIfStillPresent marks a workload unhealthy only if, re-checked +// under the file lock, its status file still exists and still claims running +// or stopped. handleRuntimeMissing decides to mark unhealthy from a snapshot +// that is stale by the time the write happens; an unconditional +// SetWorkloadStatus would recreate the file of a workload that `thv rm` +// deleted in the meantime (SetWorkloadStatus creates missing files), leaving +// a ghost "unhealthy" entry behind, or clobber a fresher status such as +// `removing`. Skipping in those cases is not an error: the persisted state is +// already newer than the observation that triggered the write. +func (f *fileStatusManager) setUnhealthyIfStillPresent(ctx context.Context, workloadName, contextMsg string) error { + return f.withFileLock(ctx, workloadName, func(statusFilePath string) error { + if _, err := os.Stat(statusFilePath); os.IsNotExist(err) { + slog.Debug("skipping unhealthy mark: status file no longer exists", "workload", workloadName) + return nil + } else if err != nil { + return fmt.Errorf("failed to check status file for workload %s: %w", workloadName, err) + } + + statusFile, err := f.readStatusFile(statusFilePath) + if err != nil { + return fmt.Errorf("failed to read existing status for workload %s: %w", workloadName, err) + } + if statusFile.Status != rt.WorkloadStatusRunning && statusFile.Status != rt.WorkloadStatusStopped { + slog.Debug("skipping unhealthy mark: status changed since observation", + "workload", workloadName, "status", statusFile.Status) + return nil + } + + statusFile.Status = rt.WorkloadStatusUnhealthy + statusFile.StatusContext = contextMsg + statusFile.UpdatedAt = time.Now() + return f.writeStatusFile(statusFilePath, *statusFile) + }) +} + // isProxyUnhealthy checks if the proxy process is running for the workload. // Returns (unhealthyWorkload, true) if proxy is not running, (emptyWorkload, false) if proxy is healthy or not applicable. func (f *fileStatusManager) isProxyUnhealthy( diff --git a/pkg/workloads/statuses/file_status_test.go b/pkg/workloads/statuses/file_status_test.go index 883ee0ba5e..d9546af162 100644 --- a/pkg/workloads/statuses/file_status_test.go +++ b/pkg/workloads/statuses/file_status_test.go @@ -1204,6 +1204,129 @@ func TestFileStatusManager_ListWorkloads_WithValidation(t *testing.T) { assert.True(t, isProxyValidated, "Proxy down workload should be detected as unhealthy") } +// TestFileStatusManager_ListWorkloads_DoesNotPoisonWorkloadStartingDuringList +// pins the snapshot ordering inside ListWorkloads: status files must be read +// before the runtime is queried. With the inverted order, a workload whose +// status file flips starting → running between the two snapshots (a `thv run` +// finishing concurrently) is seen as "file running, runtime missing", and +// handleRuntimeMissing permanently overwrites its status file as unhealthy — +// poisoning a healthy, freshly started workload. The mock flips the file to +// running from inside the runtime-list call, which under the correct ordering +// is the later of the two snapshots. +func TestFileStatusManager_ListWorkloads_DoesNotPoisonWorkloadStartingDuringList(t *testing.T) { + t.Parallel() + + ctrl := gomock.NewController(t) + defer ctrl.Finish() + + manager, mockRuntime, mockRunConfigStore := newTestFileStatusManager(t, ctrl) + ctx := context.Background() + + mockRunConfigStore.EXPECT().Exists(gomock.Any(), gomock.Any()).Return(false, nil).AnyTimes() + + const workloadName = "starting-workload" + require.NoError(t, manager.SetWorkloadStatus(ctx, workloadName, rt.WorkloadStatusStarting, "container starting")) + + // Simulate the concurrent startup completing between the file snapshot and + // the runtime snapshot: flip the file to running and return a runtime view + // from before the container existed. + mockRuntime.EXPECT().ListWorkloads(gomock.Any()).DoAndReturn(func(callCtx context.Context) ([]rt.ContainerInfo, error) { + require.NoError(t, manager.SetWorkloadStatus(callCtx, workloadName, rt.WorkloadStatusRunning, "")) + return []rt.ContainerInfo{}, nil + }) + + workloads, err := manager.ListWorkloads(ctx, true, nil) + require.NoError(t, err) + require.Len(t, workloads, 1) + assert.Equal(t, workloadName, workloads[0].Name) + assert.Equal(t, rt.WorkloadStatusStarting, workloads[0].Status, + "a workload still starting in the file snapshot must be reported as starting, not poisoned") + + // The persisted status file must still say running: ListWorkloads must not + // have overwritten the concurrent writer's state with unhealthy. + data, err := os.ReadFile(manager.getStatusFilePath(workloadName)) + require.NoError(t, err) + var persisted workloadStatusFile + require.NoError(t, json.Unmarshal(data, &persisted)) + assert.Equal(t, rt.WorkloadStatusRunning, persisted.Status, + "ListWorkloads must not overwrite the status file of a workload that finished starting mid-list") +} + +// TestFileStatusManager_ListWorkloads_UnhealthyMarkIsGuarded pins the guarded +// write in handleRuntimeMissing: by the time the merge decides to mark a +// workload unhealthy, its snapshot is stale, and a concurrent `thv rm` may +// have deleted the status file (an unconditional SetWorkloadStatus would +// recreate it as a ghost) or moved it to `removing` (which must not be +// clobbered by a stale unhealthy verdict). Both mutations are injected from +// inside the runtime-list call, which is the later snapshot. +func TestFileStatusManager_ListWorkloads_UnhealthyMarkIsGuarded(t *testing.T) { + t.Parallel() + + const workloadName = "guarded-workload" + + tests := []struct { + name string + mutate func(ctx context.Context, m *fileStatusManager) + wantFileExists bool + wantFileStatus rt.WorkloadStatus + wantFileContext string + }{ + { + name: "file deleted mid-list is not resurrected", + mutate: func(ctx context.Context, m *fileStatusManager) { + require.NoError(t, m.DeleteWorkloadStatus(ctx, workloadName)) + }, + wantFileExists: false, + }, + { + name: "status moved to removing mid-list is not clobbered", + mutate: func(ctx context.Context, m *fileStatusManager) { + require.NoError(t, m.SetWorkloadStatus(ctx, workloadName, rt.WorkloadStatusRemoving, "being removed")) + }, + wantFileExists: true, + wantFileStatus: rt.WorkloadStatusRemoving, + wantFileContext: "being removed", + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + t.Parallel() + + ctrl := gomock.NewController(t) + defer ctrl.Finish() + + manager, mockRuntime, mockRunConfigStore := newTestFileStatusManager(t, ctrl) + ctx := context.Background() + + mockRunConfigStore.EXPECT().Exists(gomock.Any(), gomock.Any()).Return(false, nil).AnyTimes() + + require.NoError(t, manager.SetWorkloadStatus(ctx, workloadName, rt.WorkloadStatusRunning, "container started")) + + mockRuntime.EXPECT().ListWorkloads(gomock.Any()).DoAndReturn(func(callCtx context.Context) ([]rt.ContainerInfo, error) { + tt.mutate(callCtx, manager) + return []rt.ContainerInfo{}, nil + }) + + _, err := manager.ListWorkloads(ctx, true, nil) + require.NoError(t, err) + + data, err := os.ReadFile(manager.getStatusFilePath(workloadName)) + if !tt.wantFileExists { + require.True(t, os.IsNotExist(err), + "status file deleted by a concurrent rm must not be recreated by ListWorkloads") + return + } + require.NoError(t, err) + var persisted workloadStatusFile + require.NoError(t, json.Unmarshal(data, &persisted)) + assert.Equal(t, tt.wantFileStatus, persisted.Status, + "a status written after the merge's observation must not be clobbered") + assert.Equal(t, tt.wantFileContext, persisted.StatusContext) + }) + } +} + func TestFileStatusManager_GetWorkload_vs_ListWorkloads_Consistency(t *testing.T) { t.Parallel() diff --git a/test/e2e/upgrade_e2e_test.go b/test/e2e/upgrade_e2e_test.go index a18fbe80f0..b48cedf2b5 100644 --- a/test/e2e/upgrade_e2e_test.go +++ b/test/e2e/upgrade_e2e_test.go @@ -5,6 +5,7 @@ package e2e_test import ( "encoding/json" + "fmt" "os" "path/filepath" "time" @@ -248,12 +249,35 @@ var _ = Describe("Upgrade Command", Label("core", "upgrade", "e2e"), func() { // read from the named workload's own record. func waitForIsolatedMCPServer(thvCmd func(args ...string) *e2e.THVCommand, serverName string, timeout time.Duration) { GinkgoHelper() - Eventually(func() bool { + var lastObserved string + ok := func() bool { workload, err := e2e.FindWorkload(thvCmd, serverName) - if err != nil { + switch { + case err != nil: + lastObserved = fmt.Sprintf("thv list failed: %v", err) + return false + case workload == nil: + lastObserved = "not listed" + return false + case workload.Status == runtime.WorkloadStatusRunning: + return true + default: + lastObserved = fmt.Sprintf("status %q (%s)", workload.Status, workload.StatusContext) return false } - return workload != nil && workload.Status == runtime.WorkloadStatusRunning - }, timeout, 1*time.Second).Should(BeTrue(), - "workload %q should be running within %s", serverName, timeout) + } + // On timeout, dump state through the SAME isolated env before failing — + // e2e.DebugServerState would query the real ToolHive config and show + // nothing. Without this dump a CI failure here is undiagnosable: it reads + // as a bare "not running" with no way to tell a stuck startup from an + // errored workload (which is exactly what happened on main run 30335474873). + Eventually(ok, timeout, 1*time.Second).Should(BeTrue(), func() string { + stdout, stderr, err := thvCmd("list", "--all").Run() + GinkgoWriter.Printf("isolated thv list output:\nStdout: %s\nStderr: %s\nError: %v\n", stdout, stderr, err) + logs, stderr, err := thvCmd("logs", serverName).Run() + GinkgoWriter.Printf("Server logs:\n%s\nStderr: %s\nError: %v\n", logs, stderr, err) + proxyLogs, stderr, err := thvCmd("logs", serverName, "--proxy").Run() + GinkgoWriter.Printf("Proxy logs:\n%s\nStderr: %s\nError: %v\n", proxyLogs, stderr, err) + return fmt.Sprintf("workload %q should be running within %s; last observed: %s", serverName, timeout, lastObserved) + }) } From 89ceb6449aab6f80c1c7398fb1856347a49458d3 Mon Sep 17 00:00:00 2001 From: Juan Antonio Osorio Date: Tue, 28 Jul 2026 07:39:02 +0000 Subject: [PATCH 2/2] Deflake Envoy AllowPort spec's allowed-port leg MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The AllowPort spec fetched https://example.com exactly once, immediately after workload readiness. The allowed leg needs a working DNS -> TLS -> upstream chain through Envoy's dynamic_forward_proxy, and the first egress request after startup can race the isolation stack's own warm-up, failing instantly with a local error before anything reaches the network. Observed twice on 2026-07-28 (PR heads 1588e6a4 and 5c901a9d, within the same minute): IsError ~40ms and ~190ms after readiness — far too fast for genuine external TLS trouble — while the same spec passed on main minutes later. The property this spec pins is steady-state reachability of an allowed port; first-request-after-boot success is not something the DFP path promises. Retry the allowed leg until the deadline, and surface the fetch tool's error text on failure — the diagnostic that distinguishes a warm-up race from a genuinely broken AllowPort, which the old boolean-only assertion discarded. A broken AllowPort still fails, now with evidence. The blocked port-80 leg is unchanged: denial is enforced locally and deterministically once the stack is warm, which the preceding allowed leg has just proven. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_01GvBj7nNWqDQ3rKYNeY6aam --- test/e2e/network_isolation_envoy_test.go | 24 +++++++++++++++++++++--- 1 file changed, 21 insertions(+), 3 deletions(-) diff --git a/test/e2e/network_isolation_envoy_test.go b/test/e2e/network_isolation_envoy_test.go index eddea95b13..2121934fda 100644 --- a/test/e2e/network_isolation_envoy_test.go +++ b/test/e2e/network_isolation_envoy_test.go @@ -243,9 +243,27 @@ var _ = Describe("NetworkIsolationEnvoy", Label("proxy", "network", "isolation", server := startFetchServer("port", profile) By("HTTPS request (port 443) must succeed") - httpsResult := fetchThrough(server, "https://example.com", 30*time.Second) - Expect(httpsResult.IsError).To(BeFalse(), - "https://example.com must be allowed (port 443 is in AllowPort)") + // The allowed leg needs a working DNS → TLS → example.com chain through + // Envoy's dynamic_forward_proxy. The very first egress request after + // startup can race the isolation stack's own warm-up (DFP DNS cache, + // the per-workload DNS container), failing instantly with a local 503 + // long before anything has actually reached the network — observed in + // CI as an IsError result within ~40ms of readiness. The property this + // spec pins is steady-state reachability of an allowed port, not + // first-request-after-boot success (which Envoy does not promise), so + // retry until the deadline and surface the tool's error text — the one + // diagnostic that distinguishes a warm-up race from a broken AllowPort. + var lastHTTPSError string + Eventually(func() bool { + httpsResult := fetchThrough(server, "https://example.com", 30*time.Second) + if httpsResult.IsError { + lastHTTPSError = resultText(httpsResult) + return false + } + return true + }, 60*time.Second, 2*time.Second).Should(BeTrue(), func() string { + return "https://example.com must be allowed (port 443 is in AllowPort); last fetch error: " + lastHTTPSError + }) By("HTTP request (port 80) must be blocked") httpResult := fetchThrough(server, "http://example.com", 30*time.Second)