Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions apps/server/src/persistence/Migrations.ts
Original file line number Diff line number Diff line change
Expand Up @@ -65,6 +65,8 @@ import Migration0050 from "./Migrations/050_ProjectionThreadPullRequests.ts";
import Migration0051 from "./Migrations/051_ProjectionThreadMessageContext.ts";
import Migration0052 from "./Migrations/052_ProjectionThreadTitleState.ts";
import Migration0053 from "./Migrations/053_PullRequestFilesViewed.ts";
import Migration0054 from "./Migrations/054_ProviderSessionHistory.ts";
import Migration0055 from "./Migrations/055_ThreadRouteEvents.ts";

/**
* Migration loader with all migrations defined inline.
Expand Down Expand Up @@ -130,6 +132,8 @@ const migrationEntries = [
[51, "ProjectionThreadMessageContext", Migration0051],
[52, "ProjectionThreadTitleState", Migration0052],
[53, "PullRequestFilesViewed", Migration0053],
[54, "ProviderSessionHistory", Migration0054],
[55, "ThreadRouteEvents", Migration0055],
] as const;

export const migrationManifest = migrationEntries.map(([id, name]) => [id, name] as const);
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,207 @@
import { assert, it } from "@effect/vitest";
import * as Effect from "effect/Effect";
import * as Layer from "effect/Layer";
import * as SqlClient from "effect/unstable/sql/SqlClient";
import * as NodeSqliteClient from "@t3tools/shared/nodeSqliteClient";

import { migrationManifest, runMigrations } from "../Migrations.ts";

const layer = it.layer(Layer.mergeAll(NodeSqliteClient.layer({ filename: ":memory:" })));

interface HistoryRow {
readonly threadId: string;
readonly providerName: string;
readonly providerInstanceKey: string;
readonly nativeSessionId: string;
readonly parentNativeSessionId: string | null;
readonly origin: string;
readonly firstSeenAt: string;
readonly lastSeenAt: string;
}

layer("054_ProviderSessionHistory", (it) => {
it.effect("backfills the current cursor with runtime nativeSessionIdOf semantics", () =>
Effect.gen(function* () {
const sql = yield* SqlClient.SqlClient;

yield* runMigrations({ toMigrationInclusive: 53 });

const insertRuntime = (
threadId: string,
providerName: string,
cursor: string | null,
providerInstanceId: string | null = providerName,
) => {
const instanceJson = providerInstanceId === null ? null : providerInstanceId;
return sql`
INSERT INTO provider_session_runtime (
thread_id,
provider_name,
provider_instance_id,
adapter_key,
runtime_mode,
status,
last_seen_at,
resume_cursor_json,
runtime_payload_json
)
VALUES (
${threadId},
${providerName},
${instanceJson},
${providerName},
'full-access',
'running',
'2026-09-23T09:00:00.000Z',
${cursor},
NULL
)
`;
};

// Valid text resume.
yield* insertRuntime("thread-resume", "claudeAgent", '{"resume":"resume-id"}');
// Empty/whitespace resume falls through to a valid threadId.
yield* insertRuntime(
"thread-empty-resume",
"claudeAgent",
'{"resume":"","threadId":"thread-id-fallback"}',
);
// Non-string resume falls through to a valid threadId.
yield* insertRuntime(
"thread-nonstring-resume",
"codex",
'{"resume":123,"threadId":"thread-id-2"}',
);
// Blank resume and threadId fall through to a valid sessionId.
yield* insertRuntime(
"thread-session-fallback",
"codex",
'{"resume":"","threadId":" ","sessionId":"session-id-3"}',
);
// Whitespace-only text fields are ignored entirely.
yield* insertRuntime("thread-whitespace", "codex", '{"resume":" ","threadId":" "}');
// Non-string candidates never become an id.
yield* insertRuntime("thread-number-only", "codex", '{"resume":123}');
yield* insertRuntime("thread-bool-only", "codex", '{"threadId":true}');
yield* insertRuntime("thread-object-only", "codex", '{"sessionId":{"nested":1}}');
// No recognised id field.
yield* insertRuntime("thread-unknown", "codex", '{"opaque":true}');
// No cursor at all.
yield* insertRuntime("thread-absent", "codex", null);
// A null provider instance uses the deterministic empty-string key.
yield* insertRuntime("thread-null-instance", "codex", '{"resume":"null-instance-id"}', null);
// A whitespace-padded instance id is trimmed in the key.
yield* insertRuntime(
"thread-trimmed-instance",
"codex",
'{"resume":"trimmed-id"}',
" codex-x ",
);

yield* runMigrations({ toMigrationInclusive: 55 });

const rows = yield* sql<HistoryRow>`
SELECT
thread_id AS "threadId",
provider_name AS "providerName",
provider_instance_key AS "providerInstanceKey",
native_session_id AS "nativeSessionId",
parent_native_session_id AS "parentNativeSessionId",
origin,
first_seen_at AS "firstSeenAt",
last_seen_at AS "lastSeenAt"
FROM provider_session_history
ORDER BY native_session_id ASC
`;

assert.deepStrictEqual(rows, [
{
threadId: "thread-null-instance",
providerName: "codex",
providerInstanceKey: "",
nativeSessionId: "null-instance-id",
parentNativeSessionId: null,
origin: "runtimeCursor",
firstSeenAt: "2026-09-23T09:00:00.000Z",
lastSeenAt: "2026-09-23T09:00:00.000Z",
},
{
threadId: "thread-resume",
providerName: "claudeAgent",
providerInstanceKey: "claudeAgent",
nativeSessionId: "resume-id",
parentNativeSessionId: null,
origin: "runtimeCursor",
firstSeenAt: "2026-09-23T09:00:00.000Z",
lastSeenAt: "2026-09-23T09:00:00.000Z",
},
{
threadId: "thread-session-fallback",
providerName: "codex",
providerInstanceKey: "codex",
nativeSessionId: "session-id-3",
parentNativeSessionId: null,
origin: "runtimeCursor",
firstSeenAt: "2026-09-23T09:00:00.000Z",
lastSeenAt: "2026-09-23T09:00:00.000Z",
},
{
threadId: "thread-nonstring-resume",
providerName: "codex",
providerInstanceKey: "codex",
nativeSessionId: "thread-id-2",
parentNativeSessionId: null,
origin: "runtimeCursor",
firstSeenAt: "2026-09-23T09:00:00.000Z",
lastSeenAt: "2026-09-23T09:00:00.000Z",
},
{
threadId: "thread-empty-resume",
providerName: "claudeAgent",
providerInstanceKey: "claudeAgent",
nativeSessionId: "thread-id-fallback",
parentNativeSessionId: null,
origin: "runtimeCursor",
firstSeenAt: "2026-09-23T09:00:00.000Z",
lastSeenAt: "2026-09-23T09:00:00.000Z",
},
{
threadId: "thread-trimmed-instance",
providerName: "codex",
providerInstanceKey: "codex-x",
nativeSessionId: "trimmed-id",
parentNativeSessionId: null,
origin: "runtimeCursor",
firstSeenAt: "2026-09-23T09:00:00.000Z",
lastSeenAt: "2026-09-23T09:00:00.000Z",
},
]);

const routeEventCount = yield* sql<{ readonly count: number }>`
SELECT COUNT(*) AS count FROM thread_route_events
`;
assert.equal(routeEventCount[0]!.count, 0);

// The unique key is on the durable identity (including the provider
// instance key), so a repeat cannot duplicate.
const indexes = yield* sql<{ readonly name: string }>`
PRAGMA index_list(provider_session_history)
`;
assert.ok(indexes.length > 0);

// Migration rerun is idempotent: no new rows, no duplicate backfill.
yield* runMigrations({ toMigrationInclusive: 55 });
const afterRerun = yield* sql<{ readonly count: number }>`
SELECT COUNT(*) AS count FROM provider_session_history
`;
assert.equal(afterRerun[0]!.count, rows.length);
}),
);

it("registers both migrations in the manifest", () => {
const entries = new Set(migrationManifest.map(([id, name]) => `${id}_${name}`));
assert.ok(entries.has("54_ProviderSessionHistory"));
assert.ok(entries.has("55_ThreadRouteEvents"));
});
});
118 changes: 118 additions & 0 deletions apps/server/src/persistence/Migrations/054_ProviderSessionHistory.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,118 @@
import * as SqlClient from "effect/unstable/sql/SqlClient";
import * as Effect from "effect/Effect";

/**
* Append-only native-session identity history per T3 thread.
*
* `provider_session_runtime` keeps a single current `resume_cursor_json`. A
* resume, fork, or model switch overwrites it, so earlier native sessions that
* contributed to a thread become unrecoverable. This table records one row per
* durable identity so a thread can answer which native sessions it used even
* after the cursor moved on. Repeated observations of the same identity only
* advance `last_seen_at`; they never replace a different session's row.
*
* Identity includes the configured provider instance, not just the provider
* name. Two instances of one driver can expose the same native session id on
* one thread (for example native OpenCode Go and a CLIProxyAPI loopback), and
* collapsing them would erase which instance produced the session. SQLite
* treats `NULL` as distinct in `UNIQUE`, so the instance is stored as a
* normalized non-null `provider_instance_key`: a trimmed instance id, or `""`
* for an unknown/null instance. `""` is a deterministic bucket, so two
* unknown-instance observations collapse to one row while a real instance
* stays distinct. The runtime writer computes the same key in JS via
* `normalizeProviderInstanceKey`.
*
* The backfill seeds the table from whatever cursor already exists so an
* upgraded database does not start empty. Its id selection mirrors the runtime
* `nativeSessionIdOf` semantics exactly: only a JSON text value that is
* non-empty after trimming is accepted, and the precedence `resume` →
* `threadId` → `sessionId` falls through to the next candidate on any absent,
* null, non-string, or blank value. A numeric/boolean/object cursor never
* becomes an id.
*/
export default Effect.gen(function* () {
const sql = yield* SqlClient.SqlClient;

yield* sql`
CREATE TABLE IF NOT EXISTS provider_session_history (
history_id INTEGER PRIMARY KEY AUTOINCREMENT,
thread_id TEXT NOT NULL,
provider_name TEXT NOT NULL,
provider_instance_id TEXT,
provider_instance_key TEXT NOT NULL,
adapter_key TEXT NOT NULL,
native_session_id TEXT NOT NULL,
parent_native_session_id TEXT,
origin TEXT NOT NULL,
first_seen_at TEXT NOT NULL,
last_seen_at TEXT NOT NULL,
UNIQUE (thread_id, provider_name, provider_instance_key, native_session_id)
)
`;

yield* sql`
CREATE INDEX IF NOT EXISTS idx_provider_session_history_thread
ON provider_session_history(thread_id, first_seen_at)
`;

yield* sql`
INSERT OR IGNORE INTO provider_session_history (
thread_id,
provider_name,
provider_instance_id,
provider_instance_key,
adapter_key,
native_session_id,
parent_native_session_id,
origin,
first_seen_at,
last_seen_at
)
SELECT
current.thread_id,
current.provider_name,
current.provider_instance_id,
COALESCE(NULLIF(TRIM(current.provider_instance_id), ''), ''),
current.adapter_key,
current.native_session_id,
NULL,
'runtimeCursor',
current.last_seen_at,
current.last_seen_at
FROM (
SELECT
runtime.thread_id,
runtime.provider_name,
runtime.provider_instance_id,
runtime.adapter_key,
runtime.last_seen_at,
CASE
WHEN json_type(runtime.cursor, '$.resume') = 'text'
AND TRIM(json_extract(runtime.cursor, '$.resume')) <> ''
THEN TRIM(json_extract(runtime.cursor, '$.resume'))
WHEN json_type(runtime.cursor, '$.threadId') = 'text'
AND TRIM(json_extract(runtime.cursor, '$.threadId')) <> ''
THEN TRIM(json_extract(runtime.cursor, '$.threadId'))
WHEN json_type(runtime.cursor, '$.sessionId') = 'text'
AND TRIM(json_extract(runtime.cursor, '$.sessionId')) <> ''
THEN TRIM(json_extract(runtime.cursor, '$.sessionId'))
ELSE NULL
END AS native_session_id
FROM (
SELECT
thread_id,
provider_name,
provider_instance_id,
adapter_key,
last_seen_at,
CASE
WHEN resume_cursor_json IS NOT NULL AND json_valid(resume_cursor_json)
THEN resume_cursor_json
ELSE NULL
END AS cursor
FROM provider_session_runtime
) AS runtime
) AS current
WHERE current.native_session_id IS NOT NULL
`;
});
49 changes: 49 additions & 0 deletions apps/server/src/persistence/Migrations/055_ThreadRouteEvents.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,49 @@
import * as SqlClient from "effect/unstable/sql/SqlClient";
import * as Effect from "effect/Effect";

/**
* Append-only pre-execution route and experiment metadata.
*
* This is low-cardinality metadata about a routing decision, not conversation
* content: which provider/model/effort was requested, the pre-execution task
* stratum, the experiment/cohort, the readable manager/agent identifiers, the
* route event kind, and any escalation reason. T3 only carries the metadata;
* the canonical policy that chooses a route lives in agent-config.
*
* `route_event_kind` is nullable so an automatic "this is what was requested"
* event is distinguishable from a declared canary/fallback/escalation event.
* Observed (actual) values are deliberately NOT stored here: they must come
* from measured usage, never be copied from the request.
*
* `escalation_reason` is a bounded code/slug, not free text; the writer rejects
* anything that is not a low-cardinality token. `selection_conflict` marks an
* automatic request row whose stored selection disagreed with a later
* observation, so a conflict is surfaced rather than silently overwritten.
*/
export default Effect.gen(function* () {
const sql = yield* SqlClient.SqlClient;

yield* sql`
CREATE TABLE IF NOT EXISTS thread_route_events (
event_id TEXT PRIMARY KEY,
thread_id TEXT NOT NULL,
native_session_id TEXT,
route_event_kind TEXT,
task_stratum TEXT NOT NULL,
experiment_id TEXT,
manager_id TEXT,
agent_id TEXT,
requested_provider TEXT,
requested_model TEXT,
requested_effort TEXT,
escalation_reason TEXT,
selection_conflict INTEGER NOT NULL DEFAULT 0,
recorded_at TEXT NOT NULL
)
`;

yield* sql`
CREATE INDEX IF NOT EXISTS idx_thread_route_events_thread
ON thread_route_events(thread_id, recorded_at)
`;
});
Loading
Loading