mirror of
https://github.com/VibedByKaKi/t3-code-android-nightly.git
synced 2026-10-09 20:01:15 +02:00
293 lines
14 KiB
TypeScript
293 lines
14 KiB
TypeScript
import * as NodeServices from "@effect/platform-node/NodeServices";
|
|
import { assert, it } from "@effect/vitest";
|
|
import * as Effect from "effect/Effect";
|
|
import * as FileSystem from "effect/FileSystem";
|
|
import * as Path from "effect/Path";
|
|
import * as SqlClient from "effect/sql/SqlClient";
|
|
|
|
import { runMigrations } from "../src/persistence/Migrations.ts";
|
|
import * as NodeSqliteClient from "@t3tools/shared/nodeSqliteClient";
|
|
import { runMigrateDevDb } from "./migrate-dev-db.ts";
|
|
|
|
const withDatabase = <A, E>(
|
|
databasePath: string,
|
|
effect: Effect.Effect<A, E, SqlClient.SqlClient>,
|
|
) => effect.pipe(Effect.provide(NodeSqliteClient.layer({ filename: databasePath })));
|
|
|
|
/** A migrated source db with one V2 thread per lifecycle state. Only
|
|
* `stopped-thread` and its fork qualify for the clone. */
|
|
const createFixtureSource = Effect.fn("createMigrateDevDbFixtureSource")(function* (
|
|
baseDir: string,
|
|
) {
|
|
const fs = yield* FileSystem.FileSystem;
|
|
const path = yield* Path.Path;
|
|
const stateDir = path.join(baseDir, "userdata");
|
|
const databasePath = path.join(stateDir, "statev2.sqlite");
|
|
yield* fs.makeDirectory(stateDir, { recursive: true });
|
|
yield* withDatabase(
|
|
databasePath,
|
|
Effect.gen(function* () {
|
|
const sql = yield* SqlClient.SqlClient;
|
|
yield* runMigrations();
|
|
|
|
yield* sql`INSERT INTO projection_projects
|
|
(project_id, title, workspace_root, scripts_json, created_at, updated_at, deleted_at)
|
|
VALUES
|
|
('project-kept', 'Kept', '/tmp/kept', '[]', '2026-08-01', '2026-08-01', NULL),
|
|
('project-deleted', 'Deleted', '/tmp/deleted', '[]', '2026-08-01', '2026-08-02', '2026-08-02')`;
|
|
|
|
const forkPayload =
|
|
'{"lineage":{"parentThreadId":"stopped-thread","relationshipToParent":"fork","rootThreadId":"stopped-thread"}}';
|
|
const subagentPayload =
|
|
'{"lineage":{"parentThreadId":"subagent-parent","relationshipToParent":"subagent","rootThreadId":"subagent-parent"},"forkedFrom":{"type":"node","nodeId":"node-1"}}';
|
|
// Excluded threads are newer than the kept family, so only the filters
|
|
// can keep them out of a one-family-per-project clone.
|
|
const threads = [
|
|
["stopped-thread", "project-kept", "completed", "{}", "2026-08-01"],
|
|
["fork-thread", "project-kept", "completed", forkPayload, "2026-08-02"],
|
|
["running-thread", "project-kept", "running", "{}", "2026-08-05"],
|
|
["settled-thread", "project-kept", "completed", '{"settledAt":"2026-08-01"}', "2026-08-05"],
|
|
[
|
|
"limit-thread",
|
|
"project-kept",
|
|
"completed",
|
|
'{"limitRecovery":{"autoResume":true}}',
|
|
"2026-08-05",
|
|
],
|
|
// Its result never reached the parent, so startup would deliver it.
|
|
["subagent-parent", "project-kept", "completed", "{}", "2026-08-05"],
|
|
["subagent-child", "project-kept", "completed", subagentPayload, "2026-08-05"],
|
|
["deleted-project-thread", "project-deleted", "completed", "{}", "2026-08-05"],
|
|
] as const;
|
|
for (const [threadId, projectId, runStatus, payload, updatedAt] of threads) {
|
|
yield* sql`INSERT INTO orchestration_v2_projection_threads
|
|
(thread_id, project_id, title, default_provider, runtime_mode, interaction_mode, created_at, updated_at, payload_json)
|
|
VALUES (${threadId}, ${projectId}, ${threadId}, 'codex', 'full-access', 'default', '2026-08-01', ${updatedAt}, ${payload})`;
|
|
yield* sql`INSERT INTO orchestration_v2_projection_runs
|
|
(run_id, thread_id, ordinal, provider, status, requested_at, payload_json)
|
|
VALUES (${`run-${threadId}`}, ${threadId}, 1, 'codex', ${runStatus}, '2026-08-01', '{}')`;
|
|
yield* sql`INSERT INTO orchestration_events
|
|
(event_id, aggregate_kind, stream_id, stream_version, event_type, occurred_at, actor_kind, payload_json, metadata_json)
|
|
VALUES (${`event-${threadId}`}, 'thread', ${threadId}, 0, 'thread.created', '2026-08-01', 'user', '{}', '{}')`;
|
|
}
|
|
// A provider session shared by two threads names its latest writer.
|
|
yield* sql`INSERT INTO orchestration_v2_projection_provider_sessions
|
|
(provider_session_id, thread_id, provider, status, updated_at, payload_json)
|
|
VALUES ('session-shared', 'running-thread', 'codex', 'stopped', '2026-08-01', '{}')`;
|
|
yield* sql`INSERT INTO orchestration_v2_projection_provider_session_bindings
|
|
(provider_session_id, thread_id)
|
|
VALUES ('session-shared', 'running-thread'), ('session-shared', 'stopped-thread')`;
|
|
yield* sql`INSERT INTO orchestration_v2_projection_context_transfers
|
|
(context_transfer_id, source_thread_id, target_thread_id, type, status, updated_at, payload_json)
|
|
VALUES ('transfer-1', 'settled-thread', 'stopped-thread', 'provider_handoff', 'completed', '2026-08-01', '{}')`;
|
|
yield* sql`INSERT INTO scheduled_tasks
|
|
(task_id, title, prompt, enabled, schedule_json, project_id, workspace_strategy_json,
|
|
model_selection_json, runtime_mode, interaction_mode, created_by, creation_source,
|
|
created_at, updated_at, last_run_status, run_count)
|
|
VALUES ('task-1', 'Nightly', 'Run it', 1, '{}', 'project-kept', '{}', '{}',
|
|
'full-access', 'default', 'user', 'user', '2026-08-01', '2026-08-01', 'never', 0)`;
|
|
yield* sql`INSERT INTO auth_sessions (session_id, subject, scopes, method, issued_at, expires_at)
|
|
VALUES ('session-1', 'user', '[]', 'pairing', '2026-08-01', '2027-08-01')`;
|
|
}),
|
|
);
|
|
return databasePath;
|
|
});
|
|
|
|
it.layer(NodeServices.layer)("migrate-dev-db", (it) => {
|
|
it.effect("keeps stopped thread families from live projects and clears pending work", () =>
|
|
Effect.gen(function* () {
|
|
const fs = yield* FileSystem.FileSystem;
|
|
const path = yield* Path.Path;
|
|
const sourceDir = yield* fs.makeTempDirectoryScoped({ prefix: "migrate-dev-db-src-" });
|
|
const destDir = yield* fs.makeTempDirectoryScoped({ prefix: "migrate-dev-db-dest-" });
|
|
const source = yield* createFixtureSource(sourceDir);
|
|
|
|
const result = yield* runMigrateDevDb(
|
|
{ baseDir: destDir, source, projects: 5, threadsPerProject: 1 },
|
|
{ sharedHome: sourceDir },
|
|
);
|
|
|
|
assert.equal(result.databasePath, path.join(destDir, "userdata", "statev2.sqlite"));
|
|
const kept = yield* withDatabase(
|
|
result.databasePath,
|
|
Effect.gen(function* () {
|
|
const sql = yield* SqlClient.SqlClient;
|
|
const threads = yield* sql<{ thread_id: string }>`
|
|
SELECT thread_id FROM orchestration_v2_projection_threads ORDER BY thread_id`;
|
|
const events = yield* sql<{ stream_id: string }>`
|
|
SELECT stream_id FROM orchestration_events ORDER BY stream_id`;
|
|
const sessions = yield* sql<{ provider_session_id: string }>`
|
|
SELECT provider_session_id FROM orchestration_v2_projection_provider_sessions`;
|
|
const [leftovers] = yield* sql<{ auth: number; tasks: number; transfers: number }>`
|
|
SELECT
|
|
(SELECT COUNT(*) FROM auth_sessions) AS auth,
|
|
(SELECT COUNT(*) FROM scheduled_tasks) AS tasks,
|
|
(SELECT COUNT(*) FROM orchestration_v2_projection_context_transfers) AS transfers`;
|
|
return { threads, events, sessions, leftovers };
|
|
}),
|
|
);
|
|
assert.deepStrictEqual(
|
|
kept.threads.map((row) => row.thread_id),
|
|
["fork-thread", "stopped-thread"],
|
|
);
|
|
assert.deepStrictEqual(
|
|
kept.events.map((row) => row.stream_id),
|
|
["fork-thread", "stopped-thread"],
|
|
);
|
|
assert.deepStrictEqual(
|
|
kept.sessions.map((row) => row.provider_session_id),
|
|
["session-shared"],
|
|
);
|
|
assert.deepStrictEqual(kept.leftovers, { auth: 0, tasks: 0, transfers: 0 });
|
|
}),
|
|
);
|
|
|
|
it.effect("never copies a settled thread's rows, and new events append after the source's", () =>
|
|
Effect.gen(function* () {
|
|
const fs = yield* FileSystem.FileSystem;
|
|
const sourceDir = yield* fs.makeTempDirectoryScoped({ prefix: "migrate-dev-db-slice-" });
|
|
const destDir = yield* fs.makeTempDirectoryScoped({ prefix: "migrate-dev-db-slice-dest-" });
|
|
const source = yield* createFixtureSource(sourceDir);
|
|
const [sourceSequence] = yield* withDatabase(
|
|
source,
|
|
Effect.gen(function* () {
|
|
const sql = yield* SqlClient.SqlClient;
|
|
return yield* sql<{ seq: number }>`
|
|
SELECT seq FROM sqlite_sequence WHERE name = 'orchestration_events'`;
|
|
}),
|
|
);
|
|
|
|
// Keep every family, so only the copy can leave the settled one out.
|
|
const result = yield* runMigrateDevDb(
|
|
{ baseDir: destDir, source, projects: 5, threadsPerProject: 100 },
|
|
{ sharedHome: sourceDir },
|
|
);
|
|
|
|
const copied = yield* withDatabase(
|
|
result.databasePath,
|
|
Effect.gen(function* () {
|
|
const sql = yield* SqlClient.SqlClient;
|
|
const [settledRows] = yield* sql<{ count: number }>`
|
|
SELECT
|
|
(SELECT COUNT(*) FROM orchestration_v2_projection_runs WHERE thread_id = 'settled-thread')
|
|
+ (SELECT COUNT(*) FROM orchestration_events WHERE stream_id = 'settled-thread')
|
|
AS count`;
|
|
const [sequence] = yield* sql<{ seq: number }>`
|
|
SELECT seq FROM sqlite_sequence WHERE name = 'orchestration_events'`;
|
|
return { settledRows: settledRows?.count, sequence: sequence?.seq };
|
|
}),
|
|
);
|
|
assert.equal(copied.settledRows, 0);
|
|
assert.equal(copied.sequence, sourceSequence?.seq);
|
|
}),
|
|
);
|
|
|
|
it.effect("upgrades a source from before the V2 thread tables", () =>
|
|
Effect.gen(function* () {
|
|
const fs = yield* FileSystem.FileSystem;
|
|
const path = yield* Path.Path;
|
|
const sourceDir = yield* fs.makeTempDirectoryScoped({ prefix: "migrate-dev-db-v1-" });
|
|
const destDir = yield* fs.makeTempDirectoryScoped({ prefix: "migrate-dev-db-v1-dest-" });
|
|
const stateDir = path.join(sourceDir, "userdata");
|
|
const source = path.join(stateDir, "statev2.sqlite");
|
|
yield* fs.makeDirectory(stateDir, { recursive: true });
|
|
yield* withDatabase(source, runMigrations({ toMigrationInclusive: 54 }));
|
|
|
|
const result = yield* runMigrateDevDb(
|
|
{ baseDir: destDir, source, projects: 5, threadsPerProject: 10 },
|
|
{ sharedHome: sourceDir },
|
|
);
|
|
assert.include(result.executedMigrations, "55_OrchestrationV2");
|
|
}),
|
|
);
|
|
|
|
it.effect("fails loudly on a migration slot collision", () =>
|
|
Effect.gen(function* () {
|
|
const fs = yield* FileSystem.FileSystem;
|
|
const sourceDir = yield* fs.makeTempDirectoryScoped({ prefix: "migrate-dev-db-slot-" });
|
|
const destDir = yield* fs.makeTempDirectoryScoped({ prefix: "migrate-dev-db-slot-dest-" });
|
|
const source = yield* createFixtureSource(sourceDir);
|
|
// Simulate another branch having claimed slot 1 first: the id is
|
|
// recorded, so this checkout's migration 1 silently never runs.
|
|
yield* withDatabase(
|
|
source,
|
|
Effect.gen(function* () {
|
|
const sql = yield* SqlClient.SqlClient;
|
|
yield* sql`UPDATE effect_sql_migrations
|
|
SET name = 'SomebodyElsesMigration' WHERE migration_id = 1`;
|
|
}),
|
|
);
|
|
|
|
const error = yield* runMigrateDevDb(
|
|
{ baseDir: destDir, source, projects: 5, threadsPerProject: 10 },
|
|
{ sharedHome: sourceDir },
|
|
).pipe(Effect.flip);
|
|
assert.equal(error._tag, "MigrateDevDbSlotCollisionError");
|
|
if (error._tag === "MigrateDevDbSlotCollisionError") {
|
|
assert.equal(error.slot, 1);
|
|
assert.equal(error.appliedName, "SomebodyElsesMigration");
|
|
}
|
|
}),
|
|
);
|
|
|
|
it.effect("refuses while a dev server holds the destination", () =>
|
|
Effect.gen(function* () {
|
|
const fs = yield* FileSystem.FileSystem;
|
|
const path = yield* Path.Path;
|
|
const sourceDir = yield* fs.makeTempDirectoryScoped({ prefix: "migrate-dev-db-busy-" });
|
|
const destDir = yield* fs.makeTempDirectoryScoped({ prefix: "migrate-dev-db-busy-dest-" });
|
|
const source = yield* createFixtureSource(sourceDir);
|
|
// This test process stands in for a live dev server.
|
|
const stateDir = path.join(destDir, "userdata");
|
|
yield* fs.makeDirectory(stateDir, { recursive: true });
|
|
yield* fs.writeFileString(
|
|
path.join(stateDir, "server-runtime.json"),
|
|
`{"version":1,"pid":${process.pid}}`,
|
|
);
|
|
|
|
const error = yield* runMigrateDevDb(
|
|
{ baseDir: destDir, source, projects: 5, threadsPerProject: 10 },
|
|
{ sharedHome: sourceDir },
|
|
).pipe(Effect.flip);
|
|
assert.equal(error._tag, "MigrateDevDbServerRunningError");
|
|
if (error._tag === "MigrateDevDbServerRunningError") {
|
|
assert.equal(error.pid, process.pid);
|
|
}
|
|
}),
|
|
);
|
|
|
|
it.effect("refuses a source that resolves to a destination path", () =>
|
|
Effect.gen(function* () {
|
|
const fs = yield* FileSystem.FileSystem;
|
|
const path = yield* Path.Path;
|
|
const sharedDir = yield* fs.makeTempDirectoryScoped({ prefix: "migrate-dev-db-overlap-" });
|
|
const destDir = yield* fs.makeTempDirectoryScoped({ prefix: "migrate-dev-db-overlap-dest-" });
|
|
// A leftover snapshot from a prior failed run, passed as --source: it
|
|
// must not be deleted before it is read.
|
|
const leftoverSnapshot = path.join(destDir, "userdata", "statev2.sqlite.migrate-dev-db-tmp");
|
|
yield* fs.makeDirectory(path.dirname(leftoverSnapshot), { recursive: true });
|
|
yield* fs.writeFileString(leftoverSnapshot, "not a real db");
|
|
|
|
const error = yield* runMigrateDevDb(
|
|
{ baseDir: destDir, source: leftoverSnapshot, projects: 5, threadsPerProject: 10 },
|
|
{ sharedHome: sharedDir },
|
|
).pipe(Effect.flip);
|
|
assert.equal(error._tag, "MigrateDevDbSourceIsDestinationError");
|
|
assert.equal(yield* fs.exists(leftoverSnapshot), true);
|
|
}),
|
|
);
|
|
|
|
it.effect("refuses to rebuild the shared home", () =>
|
|
Effect.gen(function* () {
|
|
const fs = yield* FileSystem.FileSystem;
|
|
const sourceDir = yield* fs.makeTempDirectoryScoped({ prefix: "migrate-dev-db-shared-" });
|
|
const source = yield* createFixtureSource(sourceDir);
|
|
|
|
const error = yield* runMigrateDevDb(
|
|
{ baseDir: sourceDir, source, projects: 5, threadsPerProject: 10 },
|
|
{ sharedHome: sourceDir },
|
|
).pipe(Effect.flip);
|
|
assert.equal(error._tag, "MigrateDevDbSharedHomeError");
|
|
}),
|
|
);
|
|
});
|