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
9 changes: 9 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -191,6 +191,15 @@ write, Maka batch-idempotently imports legacy RuntimeEvent JSONL without
rewriting it. Legacy-only workspaces remain available to read-only inspection
until that first write.

Session metadata and Agent Graph control tables now use the same process-local
operational database owner and the same `runtime.sqlite` transaction authority.
The first operational open copies a WAL-consistent `sessions.sqlite` source
into `runtime.sqlite`, validates every source row, and records the source digest
and result in `cutover_journal`. An interrupted copy resumes without partial
rows; a legacy database changed after cutover fails closed. The old database is
retained as migration evidence but is no longer a production writer. Session
transcript bodies remain append-only JSONL.

Runtime continuation remains opt-in:

- `MAKA_RUNTIME_SAFE_BOUNDARY_RESUME=1` enables the Desktop interrupted-turn
Expand Down
258 changes: 258 additions & 0 deletions packages/storage/src/__tests__/operational-state-store.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,258 @@
import assert from 'node:assert/strict';
import { mkdtemp, rm, stat, writeFile } from 'node:fs/promises';
import { tmpdir } from 'node:os';
import { join } from 'node:path';
import { DatabaseSync } from 'node:sqlite';
import { describe, test } from 'node:test';
import type { SessionHeader } from '@maka/core';
import {
acquireOperationalStateDatabase,
type OperationalStateCutoverFailpoint,
} from '../operational-state-store.js';
import { createSqliteRuntimeStore } from '../sqlite-runtime-store.js';
import { createSqliteSessionMetadataStore } from '../sqlite-session-metadata-store.js';

describe('operational state database cutover', () => {
test('imports sessions.sqlite into runtime.sqlite once and rejects later legacy changes', async () => {
const root = await mkdtemp(join(tmpdir(), 'maka-operational-cutover-'));
const legacyPath = join(root, 'sessions.sqlite');
const runtimePath = join(root, 'runtime.sqlite');
try {
const legacy = createSqliteSessionMetadataStore(legacyPath, { now: () => 10 });
await legacy.create(sessionHeader());
legacy.close();

const lease = acquireOperationalStateDatabase(root, { now: () => 20 });
const secondLease = acquireOperationalStateDatabase(root);
assert.equal(secondLease.database, lease.database);
secondLease.close();
const metadata = createSqliteSessionMetadataStore(runtimePath, {
databaseLease: lease,
now: () => 30,
});
assert.equal((await metadata.read('session-1')).header.name, 'Legacy session');
assert.deepEqual(
lease.database
.prepare(`SELECT scope, version FROM operational_schema_migrations ORDER BY scope`)
.all()
.map((row) => ({ ...row })),
[
{ scope: 'operational', version: 1 },
{ scope: 'runtime', version: 5 },
{ scope: 'session_metadata', version: 14 },
],
);
assert.deepEqual(
{
...(lease.database
.prepare(`
SELECT state, source_path, validation_json
FROM cutover_journal
WHERE store_name = 'session_metadata'
`)
.get() as Record<string, unknown>),
},
{
state: 'completed',
source_path: legacyPath,
validation_json: JSON.stringify({
session_metadata: 1,
session_metadata_labels: 1,
session_metadata_import_sources: 0,
session_metadata_tombstones: 0,
subagent_spawns: 0,
agent_graph_intent_claims: 0,
agent_graph_schedule_updates: 0,
agent_graph_operator_provisions: 0,
agent_graph_client_projections: 0,
agent_graph_client_operator_projections: 0,
agent_graph_client_terminal_activity: 0,
agent_graph_client_applied_records: 0,
agent_graph_supervisor_wakes: 0,
agent_graph_supervisor_wake_attempts: 0,
sandbox_boundary_log: 1,
}),
},
);
metadata.close();
assert.ok((await stat(legacyPath)).isFile(), 'legacy database remains as cutover evidence');

const reopenedLease = acquireOperationalStateDatabase(root);
reopenedLease.close();

const changedLegacy = createSqliteSessionMetadataStore(legacyPath);
await changedLegacy.update('session-1', { name: 'Changed by an old binary' });
changedLegacy.close();
assert.throws(
() => acquireOperationalStateDatabase(root),
/changed after session metadata cutover completed/,
);
} finally {
await rm(root, { recursive: true, force: true });
}
});

test('fails closed instead of merging conflicting canonical and legacy rows', async () => {
const root = await mkdtemp(join(tmpdir(), 'maka-operational-cutover-conflict-'));
const legacyPath = join(root, 'sessions.sqlite');
const runtimePath = join(root, 'runtime.sqlite');
try {
const legacy = createSqliteSessionMetadataStore(legacyPath);
await legacy.create(sessionHeader());
legacy.close();

const runtime = createSqliteRuntimeStore(runtimePath);
runtime.close();
const canonical = createSqliteSessionMetadataStore(runtimePath);
await canonical.create(sessionHeader({ name: 'Canonical session' }));
canonical.close();

assert.throws(
() => acquireOperationalStateDatabase(root),
/Session metadata cutover conflict in table session_metadata/,
);
const verified = createSqliteSessionMetadataStore(runtimePath);
try {
assert.equal((await verified.read('session-1')).header.name, 'Canonical session');
} finally {
verified.close();
}
} finally {
await rm(root, { recursive: true, force: true });
}
});

test('rejects an empty legacy database instead of recording an empty cutover', async () => {
const root = await mkdtemp(join(tmpdir(), 'maka-operational-cutover-empty-'));
try {
await writeFile(join(root, 'sessions.sqlite'), '');
assert.throws(
() => acquireOperationalStateDatabase(root),
/not a non-empty regular database file/,
);
const target = new DatabaseSync(join(root, 'runtime.sqlite'));
try {
assert.equal(
(
target.prepare(`SELECT COUNT(*) AS count FROM cutover_journal`).get() as {
count: number;
}
).count,
0,
);
} finally {
target.close();
}
} finally {
await rm(root, { recursive: true, force: true });
}
});

for (const failpoint of [
'after_cutover_started',
'after_cutover_rows_copied',
'after_cutover_validated',
] satisfies OperationalStateCutoverFailpoint[]) {
test(`resumes atomically after ${failpoint}`, async () => {
const root = await mkdtemp(join(tmpdir(), 'maka-operational-cutover-crash-'));
const legacyPath = join(root, 'sessions.sqlite');
const runtimePath = join(root, 'runtime.sqlite');
try {
const legacy = createSqliteSessionMetadataStore(legacyPath);
await legacy.create(sessionHeader());
legacy.close();

assert.throws(
() =>
acquireOperationalStateDatabase(root, {
failpoint: (point) => {
if (point === failpoint) throw new Error(`failpoint:${point}`);
},
}),
new RegExp(`failpoint:${failpoint}`),
);

const interrupted = new DatabaseSync(runtimePath);
try {
assert.deepEqual(
{
...(interrupted
.prepare(`
SELECT state, completed_at
FROM cutover_journal
WHERE store_name = 'session_metadata'
`)
.get() as Record<string, unknown>),
},
{ state: 'started', completed_at: null },
);
assert.equal(
(
interrupted.prepare(`SELECT COUNT(*) AS count FROM session_metadata`).get() as {
count: number;
}
).count,
0,
);
} finally {
interrupted.close();
}

const resumed = acquireOperationalStateDatabase(root);
try {
assert.equal(
(
resumed.database.prepare(`SELECT COUNT(*) AS count FROM session_metadata`).get() as {
count: number;
}
).count,
1,
);
assert.equal(
(
resumed.database
.prepare(`
SELECT state
FROM cutover_journal
WHERE store_name = 'session_metadata'
`)
.get() as { state: string }
).state,
'completed',
);
} finally {
resumed.close();
}
} finally {
await rm(root, { recursive: true, force: true });
}
});
}
});

function sessionHeader(overrides: Partial<SessionHeader> = {}): SessionHeader {
return {
id: 'session-1',
workspaceRoot: '/workspace',
cwd: '/workspace',
createdAt: 1,
lastUsedAt: 2,
name: 'Legacy session',
titleIsManual: true,
isFlagged: false,
labels: ['legacy'],
isArchived: false,
status: 'active',
hasUnread: false,
backend: 'fake',
llmConnectionSlug: 'test',
connectionLocked: true,
model: 'test-model',
thinkingLevel: 'medium',
permissionMode: 'ask',
collaborationMode: 'agent',
orchestrationMode: 'swarm',
schemaVersion: 1,
...overrides,
};
}
5 changes: 2 additions & 3 deletions packages/storage/src/__tests__/session-bundle-policy.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -114,7 +114,6 @@ test('exports one session only and excludes credential/config canaries', async (
'.maka_cli_claude_device_id',
'credentials.json',
'llm-connections.json',
'sessions.sqlite',
`artifacts/${other.id}`,
`sessions/${other.id}`,
].sort(),
Expand Down Expand Up @@ -182,7 +181,7 @@ test('exports one session only and excludes credential/config canaries', async (
});
});

test('exports while the session metadata WAL remains open and protects its sidecars', async () => {
test('exports while the operational DB remains open and protects its sidecars', async () => {
await withBundleRoots(async ({ stateRoot, configRoot, destinationRoot }) => {
const sessions = createSessionStore(stateRoot);
try {
Expand All @@ -195,7 +194,7 @@ test('exports while the session metadata WAL remains open and protects its sidec
text: 'selected transcript',
});

const sidecars = ['sessions.sqlite-wal', 'sessions.sqlite-shm'];
const sidecars = ['runtime.sqlite-wal', 'runtime.sqlite-shm'];
const liveEntries = await readdir(stateRoot);
for (const sidecar of sidecars) {
assert.ok(liveEntries.includes(sidecar), liveEntries.join(', '));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -97,7 +97,7 @@ describe('default SQLite session metadata store', () => {
assert.equal('name' in (marker ?? {}), false);
assert.equal(message?.type, 'user');
await stat(join(root, SQLITE_SESSION_METADATA_DATABASE_NAME));
await assert.rejects(() => stat(join(root, 'runtime.sqlite')), { code: 'ENOENT' });
await assert.rejects(() => stat(join(root, 'sessions.sqlite')), { code: 'ENOENT' });

await store.close?.();
const reopened = createSessionStore(root);
Expand Down
18 changes: 11 additions & 7 deletions packages/storage/src/agent-graph-control-store.ts
Original file line number Diff line number Diff line change
@@ -1,17 +1,21 @@
import { join } from 'node:path';
import { SQLITE_SESSION_METADATA_DATABASE_NAME } from './session-store.js';
import {
acquireOperationalStateDatabase,
OPERATIONAL_STATE_DATABASE_NAME,
} from './operational-state-store.js';
import {
createSqliteSessionMetadataStore,
type SqliteSessionMetadataStore,
} from './sqlite-session-metadata-store.js';

/**
* Open the SQLite metadata control plane used by one host-owned agent graph
* coordinator. Session JSONL remains the transcript authority; graph
* schedule/topology/claim relationships are queried from SQLite.
* Open the metadata repository owned by the operational database. Session
* JSONL remains the transcript-body authority; graph schedule/topology/claim
* relationships are canonical in runtime.sqlite.
*/
export function createAgentGraphControlStore(workspaceRoot: string): SqliteSessionMetadataStore {
return createSqliteSessionMetadataStore(
join(workspaceRoot, SQLITE_SESSION_METADATA_DATABASE_NAME),
);
const databaseLease = acquireOperationalStateDatabase(workspaceRoot);
return createSqliteSessionMetadataStore(join(workspaceRoot, OPERATIONAL_STATE_DATABASE_NAME), {
databaseLease,
});
}
1 change: 1 addition & 0 deletions packages/storage/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,7 @@ export * from './config-transfer.js';
export * from './automation-store.js';
export * from './sqlite-runtime-store.js';
export * from './runtime-event-transfer.js';
export * from './operational-state-store.js';
export * from './mcp-config-store.js';
export * from './workspace-identity.js';
export * from './memory-bundle-store.js';
Expand Down
Loading
Loading