From a5078be6b9f66725b5d5114b92ece341f876898b Mon Sep 17 00:00:00 2001 From: Theo Browne Date: Mon, 13 Jul 2026 17:09:59 -0700 Subject: [PATCH 1/3] Skip undecodable provider runtime rows when listing sessions MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A single stale row in provider_session_runtime (e.g. written by an older build with a status value that is no longer in the schema) made ProviderSessionRuntimeRepository.list() fail wholesale, which silently disabled every consumer that enumerates sessions — most notably the session reaper, whose sweep failed on every run. list() now logs a warning with the offending threadId and skips the row instead. Point reads (getByThreadId) still surface decode errors so a corrupt row for a specific thread remains visible. Co-Authored-By: Claude Fable 5 --- .../src/persistence/ProviderSessionRuntime.ts | 19 +++++++++---- .../RepositoryErrorCorrelation.test.ts | 28 ++++++++++++------- 2 files changed, 31 insertions(+), 16 deletions(-) diff --git a/apps/server/src/persistence/ProviderSessionRuntime.ts b/apps/server/src/persistence/ProviderSessionRuntime.ts index a3475d2f190..fe81c4ed3be 100644 --- a/apps/server/src/persistence/ProviderSessionRuntime.ts +++ b/apps/server/src/persistence/ProviderSessionRuntime.ts @@ -280,18 +280,25 @@ export const make = Effect.gen(function* () { ), ), Effect.flatMap((rows) => + // Skip rows that no longer decode (e.g. written by an older build) + // instead of failing the whole list — one stale row must not disable + // every consumer that enumerates sessions, such as the reaper. Effect.forEach(rows, (row) => decodeRuntimeRow(row).pipe( - Effect.mapError((cause) => - PersistenceDecodeError.fromSchemaError( - "ProviderSessionRuntimeRepository.list:decodeRows", - cause, - { threadId: row.threadId }, - ), + Effect.map(Option.some), + Effect.catch((cause) => + Effect.logWarning("provider.session.runtime.row-skipped", { + error: PersistenceDecodeError.fromSchemaError( + "ProviderSessionRuntimeRepository.list:decodeRows", + cause, + { threadId: row.threadId }, + ), + }).pipe(Effect.as(Option.none())), ), ), ), ), + Effect.map((decoded) => decoded.filter(Option.isSome).map((row) => row.value)), ); const deleteByThreadId: ProviderSessionRuntimeRepository["Service"]["deleteByThreadId"] = ( diff --git a/apps/server/src/persistence/RepositoryErrorCorrelation.test.ts b/apps/server/src/persistence/RepositoryErrorCorrelation.test.ts index f7425200fd1..379b06e2a22 100644 --- a/apps/server/src/persistence/RepositoryErrorCorrelation.test.ts +++ b/apps/server/src/persistence/RepositoryErrorCorrelation.test.ts @@ -182,7 +182,7 @@ describe("persistence error correlation", () => { }).pipe(Effect.provide(authPairingLinkLayer)), ); - it.effect("correlates provider runtime SQL and per-row decode failures by thread", () => + it.effect("skips undecodable provider runtime rows and correlates SQL failures by thread", () => Effect.gen(function* () { const runtimes = yield* ProviderSessionRuntime.ProviderSessionRuntimeRepository; const sql = yield* SqlClient.SqlClient; @@ -215,16 +215,24 @@ describe("persistence error correlation", () => { ) `; - const decodeError = yield* Effect.flip(runtimes.list()); - assert.instanceOf(decodeError, PersistenceErrors.PersistenceDecodeError); - assert.deepStrictEqual(decodeError.correlation, { threadId }); - assert.equal( - decodeError.message, - `Decode error in ProviderSessionRuntimeRepository.list:decodeRows: ${decodeError.issue}`, + const validThreadId = ThreadId.make("thread-valid"); + yield* runtimes.upsert({ + threadId: validThreadId, + providerName: "codex", + providerInstanceId: null, + adapterKey: "codex", + runtimeMode: "full-access", + status: "running", + lastSeenAt, + resumeCursor: null, + runtimePayload: null, + }); + + const listed = yield* runtimes.list(); + assert.deepStrictEqual( + listed.map((runtime) => runtime.threadId), + [validThreadId], ); - assert.notInclude(decodeError.issue, runtimePayload); - assert.notInclude(decodeError.message, runtimePayload); - assert.notInclude(decodeError.message, lastSeenAt); yield* sql`DROP TABLE provider_session_runtime`; const sqlFailure = yield* Effect.flip( From 2f3904a2c6e6415a8d8a290dc96d93de5af9d0d7 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Fri, 17 Jul 2026 22:03:58 +0200 Subject: [PATCH 2/3] fix(server): sanitize skipped runtime row logging Co-authored-by: codex --- apps/server/src/persistence/ProviderSessionRuntime.ts | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/apps/server/src/persistence/ProviderSessionRuntime.ts b/apps/server/src/persistence/ProviderSessionRuntime.ts index fe81c4ed3be..661e0800668 100644 --- a/apps/server/src/persistence/ProviderSessionRuntime.ts +++ b/apps/server/src/persistence/ProviderSessionRuntime.ts @@ -288,11 +288,12 @@ export const make = Effect.gen(function* () { Effect.map(Option.some), Effect.catch((cause) => Effect.logWarning("provider.session.runtime.row-skipped", { + threadId: row.threadId, error: PersistenceDecodeError.fromSchemaError( "ProviderSessionRuntimeRepository.list:decodeRows", cause, { threadId: row.threadId }, - ), + ).message, }).pipe(Effect.as(Option.none())), ), ), From 3f61538a82dd00d53ac4aea89cc8b57420a3ab53 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Fri, 17 Jul 2026 23:22:41 +0200 Subject: [PATCH 3/3] refactor(server): collect decoded runtime rows once Use Effect Array.filterMap to discard skipped rows and unwrap successful decodes in a single traversal. Co-authored-by: codex --- apps/server/src/persistence/ProviderSessionRuntime.ts | 8 +++++++- 1 file changed, 7 insertions(+), 1 deletion(-) diff --git a/apps/server/src/persistence/ProviderSessionRuntime.ts b/apps/server/src/persistence/ProviderSessionRuntime.ts index 661e0800668..2ccdd862522 100644 --- a/apps/server/src/persistence/ProviderSessionRuntime.ts +++ b/apps/server/src/persistence/ProviderSessionRuntime.ts @@ -1,7 +1,9 @@ +import * as Arr from "effect/Array"; import * as Context from "effect/Context"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; +import * as Result from "effect/Result"; import * as Schema from "effect/Schema"; import * as Struct from "effect/Struct"; import * as SqlClient from "effect/unstable/sql/SqlClient"; @@ -299,7 +301,11 @@ export const make = Effect.gen(function* () { ), ), ), - Effect.map((decoded) => decoded.filter(Option.isSome).map((row) => row.value)), + Effect.map((decoded) => + Arr.filterMap(decoded, (row) => + Option.isSome(row) ? Result.succeed(row.value) : Result.failVoid, + ), + ), ); const deleteByThreadId: ProviderSessionRuntimeRepository["Service"]["deleteByThreadId"] = (