mirror of
https://github.com/stablyai/orca.git
synced 2026-10-02 16:02:15 +00:00
feat(session-search): index OpenCode SQLite sessions
OpenCode sessions live in one SQLite database read on a worker thread, and the worker only ever answered with the newest few messages for the panel preview. The parser therefore published nothing over the transcript channel, so the search index wrote a placeholder row for every OpenCode candidate and no OpenCode message was ever searchable. Adds a `capture` request to the worker protocol that returns the session and every text part of every user/assistant turn from one open of the database. The agent parser asks for it whenever a sink is listening, so OpenCode joins the whole-document sources on the same path as Grok, Cursor and Gemini. The placeholder path (`parserPublishesMessages`, `noteUnreachableParser`) is gone; an OpenCode read that fails now fails like any other file. Bumps the index schema so existing indexes drop their placeholder rows, and adds `sessionsByAgent` to the index status, which is the count that made this bug visible.
This commit is contained in:
@@ -1,7 +1,6 @@
|
||||
import { chmod, mkdir, rm, writeFile } from 'node:fs/promises'
|
||||
import { join } from 'node:path'
|
||||
import { afterEach, beforeEach, expect, it } from 'vitest'
|
||||
import { parserPublishesMessages } from '../ai-vault/session-scanner-agent-parser'
|
||||
import { resetTranscriptConsumersForTests } from '../ai-vault/session-transcript-consumers'
|
||||
import { retireDeletedSessionSearchSources } from './session-search-deleted-sources'
|
||||
import {
|
||||
@@ -280,22 +279,6 @@ it('retires a synthetic row when the container it came from is gone', async () =
|
||||
})
|
||||
})
|
||||
|
||||
// Nothing in this PR can hold a synthetic row: the index pass refuses a source
|
||||
// whose parser decodes its messages where the message channel cannot reach
|
||||
// them, and OpenCode's SQLite sessions are read on a worker thread. The rule
|
||||
// above is the guard for the day that changes -- without it the walk would read
|
||||
// `<db>#<id>` as a filename and retire every such row the moment it appeared.
|
||||
it('does not index a source whose messages the channel cannot reach', () => {
|
||||
const db = join(harness.root, 'opencode.db')
|
||||
expect(
|
||||
parserPublishesMessages({
|
||||
agent: 'opencode',
|
||||
codexHome: null,
|
||||
file: { path: `${db}#session-1`, mtimeMs: 1, modifiedAt: '', sizeBytes: 0 }
|
||||
})
|
||||
).toBe(false)
|
||||
})
|
||||
|
||||
// Round 12, F1. The cap counts directories because that is what costs: rows
|
||||
// sharing one are a single read and then map lookups.
|
||||
it('caps the directories one pass reads, not the rows it answers', async () => {
|
||||
|
||||
@@ -1,4 +1,3 @@
|
||||
import { parserPublishesMessages } from '../ai-vault/session-scanner-agent-parser'
|
||||
import {
|
||||
registerTranscriptConsumer,
|
||||
type TranscriptConsumer,
|
||||
@@ -7,7 +6,6 @@ import {
|
||||
type TranscriptReadOutcome,
|
||||
type TranscriptReadStart
|
||||
} from '../ai-vault/session-transcript-consumers'
|
||||
import type { SessionFileCandidate } from '../ai-vault/session-scanner-types'
|
||||
import { fileIdentity } from './session-search-file-cursor'
|
||||
import type { SessionSearchFileWrite } from './session-search-index-writer'
|
||||
import type { SessionSearchStore } from './session-search-store'
|
||||
@@ -31,10 +29,6 @@ export class SessionSearchIndexConsumer implements TranscriptConsumer {
|
||||
|
||||
beginRead(start: TranscriptReadStart): TranscriptReadConsumer | null {
|
||||
const { candidate } = start
|
||||
if (!parserPublishesMessages(candidate)) {
|
||||
this.noteUnreachableParser(candidate)
|
||||
return null
|
||||
}
|
||||
if (start.mode === 'append') {
|
||||
const cursor = this.store.indexedFile(candidate.file.path, fileIdentity(candidate.file))
|
||||
if (!cursor || cursor.byteOffset !== start.previousByteOffset) {
|
||||
@@ -60,35 +54,6 @@ export class SessionSearchIndexConsumer implements TranscriptConsumer {
|
||||
}
|
||||
return new SessionSearchReadConsumer(this.store, start, write)
|
||||
}
|
||||
|
||||
/**
|
||||
* A source no read can ever index, recorded as one this index has seen.
|
||||
*
|
||||
* A parser that decodes where the message channel cannot reach it -- OpenCode's
|
||||
* SQLite sessions today -- publishes nothing, so no read of it will ever
|
||||
* commit a row. Leaving the file table silent about it is not free: the next
|
||||
* pass sees a path the index holds nothing for, asks for a read, and asking
|
||||
* over a warm cache drops the session list's own resume point. The sidebar's
|
||||
* fold is thrown away and the whole database is decoded again, on every pass,
|
||||
* for ever.
|
||||
*
|
||||
* The row written is the shape the store already has for a read that went
|
||||
* through and decoded no session: cursor at the file's size, no session row.
|
||||
* The decide step then skips it until its stat moves, and the retirement walk
|
||||
* retires it like any other row when it goes.
|
||||
*/
|
||||
private noteUnreachableParser(candidate: SessionFileCandidate): void {
|
||||
const write = this.store.beginWrite(candidate, 'replace', 0)
|
||||
const committed =
|
||||
write?.commit({
|
||||
session: null,
|
||||
byteOffset: candidate.file.sizeBytes ?? 0,
|
||||
incomplete: false
|
||||
}) === true
|
||||
if (committed) {
|
||||
this.store.writeCommitted(candidate)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
class SessionSearchReadConsumer implements TranscriptReadConsumer {
|
||||
|
||||
@@ -40,6 +40,8 @@ export type SessionSearchIndexStatus = {
|
||||
lastReconcileAt: number | null
|
||||
/** When a whole-machine sweep last finished; null until one has. */
|
||||
lastSweepCompletedAt: number | null
|
||||
/** Indexed sessions per agent; an agent with files and none is unsearchable. */
|
||||
sessionsByAgent: Record<string, number>
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -187,7 +189,8 @@ export class SessionSearchIndexer {
|
||||
const settled = (this.closed ? this.lastCounts : this.readCounts()) ?? {
|
||||
current: 0,
|
||||
due: 0,
|
||||
failed: 0
|
||||
failed: 0,
|
||||
sessionsByAgent: {}
|
||||
}
|
||||
return {
|
||||
phase: this.phase(settled),
|
||||
@@ -196,7 +199,8 @@ export class SessionSearchIndexer {
|
||||
filesFailed: settled.failed,
|
||||
degradedRoots: this.degradedRoots.map((root) => ({ ...root })),
|
||||
lastReconcileAt: this.lastReconcileAt,
|
||||
lastSweepCompletedAt: this.lastSweepCompletedAt
|
||||
lastSweepCompletedAt: this.lastSweepCompletedAt,
|
||||
sessionsByAgent: { ...settled.sessionsByAgent }
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -1,177 +0,0 @@
|
||||
import { mkdirSync } from 'node:fs'
|
||||
import { join } from 'node:path'
|
||||
import { afterEach, beforeEach, expect, it, vi } from 'vitest'
|
||||
|
||||
// Only the thread hop is replaced: both implementations below are the repo's
|
||||
// own in-process readers, which the worker entry calls on the other side.
|
||||
export const openCodeParseCalls: string[] = []
|
||||
vi.mock('../ai-vault/session-scanner-opencode-sqlite-worker-spawn', async () => {
|
||||
const list = await import('../ai-vault/session-scanner-opencode-sqlite-list')
|
||||
const parse = await import('../ai-vault/session-scanner-opencode-sqlite')
|
||||
const own = await import('./session-search-opencode-decline.test')
|
||||
return {
|
||||
resolveOpenCodeSqliteWorkerEntryPath: () => null,
|
||||
listOpenCodeSqliteSessionsViaWorker: (
|
||||
args: Parameters<typeof list.listOpenCodeSqliteSessions>[0]
|
||||
) => list.listOpenCodeSqliteSessions(args),
|
||||
parseOpenCodeSqliteSessionViaWorker: (
|
||||
args: Parameters<typeof parse.parseOpenCodeSqliteSession>[0]
|
||||
) => {
|
||||
own.openCodeParseCalls.push(args.sessionId)
|
||||
return parse.parseOpenCodeSqliteSession(args)
|
||||
}
|
||||
}
|
||||
})
|
||||
import Database from '../sqlite/sync-database'
|
||||
import { getSessionParseCacheEntry } from '../ai-vault/session-parse-cache-store'
|
||||
import { resetSessionParseCacheForTests } from '../ai-vault/session-scanner-parse-cache'
|
||||
import { resetTranscriptConsumersForTests } from '../ai-vault/session-transcript-consumers'
|
||||
import { buildOpenCodeSqliteCandidatePath } from '../ai-vault/session-scanner-opencode-sqlite-paths'
|
||||
import { SessionSearchIndexer } from './session-search-indexer'
|
||||
import {
|
||||
FakeSessionSearchClock,
|
||||
openSessionSearchIndexerHarness,
|
||||
writeClaudeTranscript,
|
||||
type SessionSearchIndexerHarness
|
||||
} from './session-search-indexer-test-fixture'
|
||||
|
||||
/*
|
||||
* Round 12, F3. An OpenCode SQLite session decodes where the message channel
|
||||
* cannot reach it, so no read of one will ever commit a row. The consumer
|
||||
* declined it and wrote nothing, which left the file table silent about a
|
||||
* source discovery returns on every pass: the decide step saw a path the index
|
||||
* held nothing for, asked for a read, and asking for one over a warm cache
|
||||
* drops the session list's own resume point. Every OpenCode session was fully
|
||||
* decoded on every pass and the sidebar's fold was thrown away with it, which
|
||||
* is the cache STA-1278 and STA-1417 added.
|
||||
*/
|
||||
|
||||
const SESSION = 'ses_r12'
|
||||
const CLAUDE_SESSION = 'aaaaaaaa-bbbb-4ccc-8ddd-eeeeeeeeeeee'
|
||||
|
||||
let harness: SessionSearchIndexerHarness
|
||||
let clock: FakeSessionSearchClock
|
||||
let indexer: SessionSearchIndexer | null = null
|
||||
|
||||
beforeEach(async () => {
|
||||
resetSessionParseCacheForTests()
|
||||
resetTranscriptConsumersForTests()
|
||||
clock = new FakeSessionSearchClock()
|
||||
harness = await openSessionSearchIndexerHarness('ss-opencode-decline')
|
||||
indexer = null
|
||||
openCodeParseCalls.length = 0
|
||||
})
|
||||
|
||||
afterEach(async () => {
|
||||
indexer?.close()
|
||||
resetTranscriptConsumersForTests()
|
||||
resetSessionParseCacheForTests()
|
||||
await harness.cleanup()
|
||||
})
|
||||
|
||||
function writeOpenCodeDb(path: string, sessionId: string): void {
|
||||
const db = new Database(path)
|
||||
db.exec(`
|
||||
CREATE TABLE session (
|
||||
id TEXT PRIMARY KEY, project_id TEXT NOT NULL, parent_id TEXT, slug TEXT NOT NULL,
|
||||
directory TEXT NOT NULL, title TEXT NOT NULL, version TEXT NOT NULL, share_url TEXT,
|
||||
summary_additions INTEGER, summary_deletions INTEGER, summary_files INTEGER,
|
||||
summary_diffs TEXT, revert TEXT, permission TEXT,
|
||||
time_created INTEGER NOT NULL, time_updated INTEGER NOT NULL, time_compacting INTEGER,
|
||||
time_archived INTEGER, workspace_id TEXT, path TEXT, agent TEXT, model TEXT,
|
||||
cost REAL DEFAULT 0 NOT NULL, tokens_input INTEGER DEFAULT 0 NOT NULL,
|
||||
tokens_output INTEGER DEFAULT 0 NOT NULL, tokens_reasoning INTEGER DEFAULT 0 NOT NULL,
|
||||
tokens_cache_read INTEGER DEFAULT 0 NOT NULL, tokens_cache_write INTEGER DEFAULT 0 NOT NULL,
|
||||
metadata TEXT
|
||||
);
|
||||
CREATE TABLE message (
|
||||
id TEXT PRIMARY KEY, session_id TEXT NOT NULL, time_created INTEGER NOT NULL,
|
||||
time_updated INTEGER NOT NULL, data TEXT NOT NULL
|
||||
);
|
||||
CREATE TABLE project (
|
||||
id TEXT PRIMARY KEY, worktree TEXT NOT NULL, vcs TEXT, name TEXT, icon_url TEXT,
|
||||
icon_color TEXT, time_created INTEGER NOT NULL, time_updated INTEGER NOT NULL,
|
||||
time_initialized INTEGER, sandboxes TEXT NOT NULL, commands TEXT, icon_url_override TEXT
|
||||
);
|
||||
CREATE TABLE part (
|
||||
id TEXT PRIMARY KEY, message_id TEXT NOT NULL, session_id TEXT NOT NULL,
|
||||
time_created INTEGER NOT NULL, time_updated INTEGER NOT NULL, data TEXT NOT NULL
|
||||
);
|
||||
`)
|
||||
db.prepare(
|
||||
`INSERT INTO session (id, project_id, parent_id, slug, directory, title, version,
|
||||
time_created, time_updated, agent, model, cost, tokens_input, tokens_output,
|
||||
tokens_reasoning, tokens_cache_read, tokens_cache_write)
|
||||
VALUES (?, 'proj-1', NULL, 'slug-1', '/tmp/opencode', 'OpenCode title', '1.0.0',
|
||||
?, ?, 'build', '{"id":"glm"}', 0, 1, 1, 0, 0, 0)`
|
||||
).run(sessionId, 1_740_000_000_000, 1_740_000_100_000)
|
||||
db.prepare(
|
||||
`INSERT INTO message (id, session_id, time_created, time_updated, data) VALUES (?, ?, ?, ?, ?)`
|
||||
).run(
|
||||
'msg-1',
|
||||
sessionId,
|
||||
1_740_000_000_000,
|
||||
1_740_000_000_000,
|
||||
JSON.stringify({ role: 'user', time: { created: 1_740_000_000_000 } })
|
||||
)
|
||||
db.prepare(
|
||||
`INSERT INTO part (id, message_id, session_id, time_created, time_updated, data) VALUES (?, ?, ?, ?, ?, ?)`
|
||||
).run(
|
||||
'part-1',
|
||||
'msg-1',
|
||||
sessionId,
|
||||
1_740_000_000_000,
|
||||
1_740_000_000_000,
|
||||
JSON.stringify({ type: 'text', text: 'hello opencode' })
|
||||
)
|
||||
db.prepare(
|
||||
`INSERT INTO project (id, worktree, name, time_created, time_updated, sandboxes)
|
||||
VALUES ('proj-1', '/tmp/opencode', 'proj', ?, ?, '[]')`
|
||||
).run(1_740_000_000_000, 1_740_000_000_000)
|
||||
db.close()
|
||||
}
|
||||
|
||||
it('reads an OpenCode session once, not on every pass', async () => {
|
||||
const dbPath = join(harness.root, 'opencode-db', 'opencode.db')
|
||||
mkdirSync(join(harness.root, 'opencode-db'), { recursive: true })
|
||||
writeOpenCodeDb(dbPath, SESSION)
|
||||
const claudePath = join(harness.claudeProjectDir, 'control.jsonl')
|
||||
await writeClaudeTranscript(claudePath, ['control turn'], CLAUDE_SESSION)
|
||||
|
||||
indexer = new SessionSearchIndexer({
|
||||
databasePath: harness.databasePath,
|
||||
roots: { ...harness.roots, opencodeDbPaths: [dbPath] },
|
||||
historyDays: null,
|
||||
clock,
|
||||
reconcileIntervalMs: 20_000,
|
||||
onError: () => undefined
|
||||
})
|
||||
await indexer.start()
|
||||
|
||||
const syntheticPath = buildOpenCodeSqliteCandidatePath(dbPath, SESSION)
|
||||
const openCodeAfterFirst = getSessionParseCacheEntry(syntheticPath)
|
||||
const claudeAfterFirst = getSessionParseCacheEntry(claudePath)
|
||||
|
||||
await indexer.reconcile()
|
||||
await indexer.reconcile()
|
||||
|
||||
// One decode across three passes, and the session list's cached fold for it
|
||||
// is the same object it was after the first: nothing invalidated it.
|
||||
expect(openCodeParseCalls).toHaveLength(1)
|
||||
expect(getSessionParseCacheEntry(syntheticPath)).toBe(openCodeAfterFirst)
|
||||
// The control, which the index really does hold, is untouched either way.
|
||||
expect(getSessionParseCacheEntry(claudePath)).toBe(claudeAfterFirst)
|
||||
|
||||
// What makes it skippable: a row saying the index has seen this source and
|
||||
// holds no session for it, which is the shape a read-through-with-no-session
|
||||
// already leaves.
|
||||
const rows = harness.read((db) =>
|
||||
db.prepare('SELECT path, state, session_row_id FROM files ORDER BY path').all()
|
||||
) as { path: string; state: string; session_row_id: number | null }[]
|
||||
expect(rows).toHaveLength(2)
|
||||
expect(rows.find((row) => row.path === syntheticPath)).toMatchObject({
|
||||
state: 'current',
|
||||
session_row_id: null
|
||||
})
|
||||
expect(indexer.status()).toMatchObject({ filesDue: 0, filesFailed: 0, phase: 'current' })
|
||||
})
|
||||
@@ -0,0 +1,203 @@
|
||||
import { join } from 'node:path'
|
||||
import { afterEach, beforeEach, expect, it, vi } from 'vitest'
|
||||
|
||||
// Only the thread hop is replaced: all three implementations below are the
|
||||
// repo's own in-process readers, which the worker entry calls on the other side.
|
||||
export const openCodeReadCalls: string[] = []
|
||||
vi.mock('../ai-vault/session-scanner-opencode-sqlite-worker-spawn', async () => {
|
||||
const list = await import('../ai-vault/session-scanner-opencode-sqlite-list')
|
||||
const parse = await import('../ai-vault/session-scanner-opencode-sqlite')
|
||||
const capture = await import('../ai-vault/session-scanner-opencode-sqlite-capture')
|
||||
const own = await import('./session-search-opencode-index.test')
|
||||
return {
|
||||
resolveOpenCodeSqliteWorkerEntryPath: () => null,
|
||||
listOpenCodeSqliteSessionsViaWorker: (
|
||||
args: Parameters<typeof list.listOpenCodeSqliteSessions>[0]
|
||||
) => list.listOpenCodeSqliteSessions(args),
|
||||
parseOpenCodeSqliteSessionViaWorker: (
|
||||
args: Parameters<typeof parse.parseOpenCodeSqliteSession>[0]
|
||||
) => {
|
||||
own.openCodeReadCalls.push(`parse:${args.sessionId}`)
|
||||
return parse.parseOpenCodeSqliteSession(args)
|
||||
},
|
||||
captureOpenCodeSqliteSessionViaWorker: (
|
||||
args: Parameters<typeof capture.captureOpenCodeSqliteSession>[0]
|
||||
) => {
|
||||
own.openCodeReadCalls.push(`capture:${args.sessionId}`)
|
||||
return capture.captureOpenCodeSqliteSession(args)
|
||||
}
|
||||
}
|
||||
})
|
||||
import { resetSessionParseCacheForTests } from '../ai-vault/session-scanner-parse-cache'
|
||||
import { resetTranscriptConsumersForTests } from '../ai-vault/session-transcript-consumers'
|
||||
import { buildOpenCodeSqliteCandidatePath } from '../ai-vault/session-scanner-opencode-sqlite-paths'
|
||||
import {
|
||||
appendOpenCodeSqliteTurn,
|
||||
writeOpenCodeSqliteDatabase
|
||||
} from '../ai-vault/session-scanner-opencode-sqlite-fixture'
|
||||
import type SyncDatabase from '../sqlite/sync-database'
|
||||
import { SessionSearchEngine } from './session-search-engine'
|
||||
import { SessionSearchIndexer } from './session-search-indexer'
|
||||
import { openSessionSearchDatabase } from './session-search-schema'
|
||||
import {
|
||||
FakeSessionSearchClock,
|
||||
openSessionSearchIndexerHarness,
|
||||
writeClaudeTranscript,
|
||||
type SessionSearchIndexerHarness
|
||||
} from './session-search-indexer-test-fixture'
|
||||
|
||||
/*
|
||||
* Nothing any OpenCode agent said used to be searchable. Its sessions live in
|
||||
* one SQLite database read on a worker thread, and the worker only ever
|
||||
* returned the newest few messages for the panel preview, so the index recorded
|
||||
* a placeholder row and moved on. This is the end-to-end proof that a sentence
|
||||
* an OpenCode assistant wrote comes back from a real search over a real index.
|
||||
*/
|
||||
|
||||
// Literal-looking on purpose: the `phrase` route is the one a user quoting a
|
||||
// remembered sentence takes, and only a literal query reaches it.
|
||||
const ANSWER = 'the quokkaTelemetry harness reindexes every shard'
|
||||
const OTHER = 'a completely unrelated conversation about typography'
|
||||
const SESSION = 'ses_capture'
|
||||
const SECOND_SESSION = 'ses_second'
|
||||
const CLAUDE_SESSION = 'aaaaaaaa-bbbb-4ccc-8ddd-eeeeeeeeeeee'
|
||||
|
||||
let harness: SessionSearchIndexerHarness
|
||||
let clock: FakeSessionSearchClock
|
||||
let indexer: SessionSearchIndexer | null = null
|
||||
let engineDbs: SyncDatabase[] = []
|
||||
|
||||
beforeEach(async () => {
|
||||
resetSessionParseCacheForTests()
|
||||
resetTranscriptConsumersForTests()
|
||||
clock = new FakeSessionSearchClock()
|
||||
harness = await openSessionSearchIndexerHarness('ss-opencode-index')
|
||||
indexer = null
|
||||
engineDbs = []
|
||||
openCodeReadCalls.length = 0
|
||||
})
|
||||
|
||||
afterEach(async () => {
|
||||
indexer?.close()
|
||||
for (const db of engineDbs) {
|
||||
db.close()
|
||||
}
|
||||
resetTranscriptConsumersForTests()
|
||||
resetSessionParseCacheForTests()
|
||||
await harness.cleanup()
|
||||
})
|
||||
|
||||
function dbPath(): string {
|
||||
return join(harness.root, 'opencode-db', 'opencode.db')
|
||||
}
|
||||
|
||||
async function startIndexer(): Promise<SessionSearchIndexer> {
|
||||
const started = new SessionSearchIndexer({
|
||||
databasePath: harness.databasePath,
|
||||
roots: { ...harness.roots, opencodeDbPaths: [dbPath()] },
|
||||
historyDays: null,
|
||||
clock,
|
||||
reconcileIntervalMs: 20_000,
|
||||
onError: (error) => {
|
||||
throw error
|
||||
}
|
||||
})
|
||||
indexer = started
|
||||
await started.start()
|
||||
return started
|
||||
}
|
||||
|
||||
/** A second connection on the index file, the way the live instance pairs them. */
|
||||
function openEngine(): SessionSearchEngine {
|
||||
const db = openSessionSearchDatabase(harness.databasePath)
|
||||
engineDbs.push(db)
|
||||
return new SessionSearchEngine(db)
|
||||
}
|
||||
|
||||
function writeVault(): void {
|
||||
writeOpenCodeSqliteDatabase(dbPath(), [
|
||||
{
|
||||
id: SESSION,
|
||||
title: 'Telemetry work',
|
||||
directory: '/tmp/opencode',
|
||||
turns: [
|
||||
{ role: 'user', parts: ['how do I reindex the shards'] },
|
||||
{ role: 'assistant', parts: ['Here is the plan.', ANSWER] }
|
||||
]
|
||||
},
|
||||
{
|
||||
id: SECOND_SESSION,
|
||||
title: 'Typography',
|
||||
directory: '/tmp/opencode-two',
|
||||
turns: [{ role: 'assistant', parts: [OTHER] }]
|
||||
}
|
||||
])
|
||||
}
|
||||
|
||||
it('finds a sentence an OpenCode assistant wrote, through the real indexer', async () => {
|
||||
writeVault()
|
||||
await startIndexer()
|
||||
|
||||
const response = openEngine().search({ query: ANSWER })
|
||||
|
||||
expect(response.planner.route).toBe('phrase')
|
||||
expect(response.hits).toHaveLength(1)
|
||||
const hit = response.hits[0]
|
||||
expect(hit).toMatchObject({
|
||||
agent: 'opencode',
|
||||
sessionId: SESSION,
|
||||
cwd: '/tmp/opencode'
|
||||
})
|
||||
expect(hit?.evidence?.role).toBe('assistant')
|
||||
expect(hit?.evidence?.snippet).toContain('quokkaTelemetry')
|
||||
// The whole-session read, not the preview window: the user turn is indexed too.
|
||||
expect(openEngine().search({ query: 'reindex the shards' }).hits).toHaveLength(1)
|
||||
// And the sibling session is a session of its own, not folded into this one.
|
||||
expect(openEngine().search({ query: OTHER }).hits[0]?.sessionId).toBe(SECOND_SESSION)
|
||||
})
|
||||
|
||||
it('reads an OpenCode session once, not on every pass', async () => {
|
||||
writeVault()
|
||||
const claudePath = join(harness.claudeProjectDir, 'control.jsonl')
|
||||
await writeClaudeTranscript(claudePath, ['control turn'], CLAUDE_SESSION)
|
||||
const started = await startIndexer()
|
||||
|
||||
await started.reconcile()
|
||||
await started.reconcile()
|
||||
|
||||
// One capture per session across three passes; nothing re-decodes a session
|
||||
// whose `time_updated` has not moved.
|
||||
expect(openCodeReadCalls).toEqual([`capture:${SESSION}`, `capture:${SECOND_SESSION}`])
|
||||
const rows = harness.read((db) =>
|
||||
db.prepare('SELECT path, state, session_row_id FROM files ORDER BY path').all()
|
||||
) as { path: string; state: string; session_row_id: number | null }[]
|
||||
expect(
|
||||
rows.find((row) => row.path === buildOpenCodeSqliteCandidatePath(dbPath(), SESSION))
|
||||
).toMatchObject({ state: 'current' })
|
||||
expect(
|
||||
rows.find((row) => row.path === buildOpenCodeSqliteCandidatePath(dbPath(), SESSION))
|
||||
?.session_row_id
|
||||
).not.toBeNull()
|
||||
expect(started.status()).toMatchObject({ filesDue: 0, filesFailed: 0, phase: 'current' })
|
||||
// The count that made this bug visible: two OpenCode files, two OpenCode
|
||||
// sessions. Before the capture channel it read two files and zero sessions.
|
||||
expect(started.status().sessionsByAgent).toMatchObject({ opencode: 2, claude: 1 })
|
||||
})
|
||||
|
||||
it('re-reads a session that gained a message and replaces its rows', async () => {
|
||||
writeVault()
|
||||
const started = await startIndexer()
|
||||
expect(openEngine().search({ query: 'orthogonal vestibule' }).hits).toHaveLength(0)
|
||||
|
||||
appendOpenCodeSqliteTurn(dbPath(), SESSION, {
|
||||
role: 'assistant',
|
||||
parts: ['an orthogonal vestibule appeared']
|
||||
})
|
||||
await started.reconcile()
|
||||
|
||||
const engine = openEngine()
|
||||
expect(engine.search({ query: 'orthogonal vestibule' }).hits[0]?.sessionId).toBe(SESSION)
|
||||
// Replaced whole, not appended twice: the original turn is still one hit.
|
||||
expect(engine.search({ query: ANSWER }).hits).toHaveLength(1)
|
||||
expect(openCodeReadCalls.filter((call) => call === `capture:${SESSION}`)).toHaveLength(2)
|
||||
})
|
||||
@@ -10,7 +10,7 @@ import { removeTreeSync } from '../../shared/windows-transient-lock-removal'
|
||||
// policy, decided where the wire is.
|
||||
|
||||
// Bump to drop and rebuild: the index is a cache over the transcripts, never a source.
|
||||
export const SESSION_SEARCH_SCHEMA_VERSION = 5
|
||||
export const SESSION_SEARCH_SCHEMA_VERSION = 6
|
||||
|
||||
// unicode61 keeps `_ . - /` inside tokens so paths and identifiers match exactly;
|
||||
// the `identifiers` column carries the split form (see session-search-identifier-split).
|
||||
|
||||
@@ -24,7 +24,7 @@ async function fixture() {
|
||||
// degradedRoots is re-stated because the contract type leaves `root` optional
|
||||
// for relay redaction, while the indexer always names the root it degraded.
|
||||
const indexer = {
|
||||
status: () => ({ ...status, degradedRoots: [] }),
|
||||
status: () => ({ ...status, degradedRoots: [], sessionsByAgent: {} }),
|
||||
reconcile: vi.fn(async () => {})
|
||||
}
|
||||
const service = createSessionSearchService({ engine: harness.engine, indexer })
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
import type SyncDatabase from '../sqlite/sync-database'
|
||||
import { asRecord } from '../ai-vault/session-scanner-record-value'
|
||||
import type { SessionFileCandidate } from '../ai-vault/session-scanner-types'
|
||||
import type { TranscriptSessionIdentity } from '../ai-vault/session-transcript-consumers'
|
||||
import type {
|
||||
@@ -42,7 +43,20 @@ export type SessionSearchFileRow = {
|
||||
}
|
||||
|
||||
/** How many rows are in each state; the whole of the indexer's progress report. */
|
||||
export type SessionSearchStateCounts = { current: number; due: number; failed: number }
|
||||
export type SessionSearchStateCounts = {
|
||||
current: number
|
||||
due: number
|
||||
failed: number
|
||||
/**
|
||||
* Indexed sessions per agent.
|
||||
*
|
||||
* The one number that distinguishes an agent the index has read from one it
|
||||
* has only listed: OpenCode's 606 rows in `files` with nothing in `sessions`
|
||||
* was the shape of a whole source being silently unsearchable, and no
|
||||
* file-state count could show it.
|
||||
*/
|
||||
sessionsByAgent: Record<string, number>
|
||||
}
|
||||
|
||||
/**
|
||||
* Owns the index database. PR 2 scope: the write half only — the transcript
|
||||
@@ -272,13 +286,34 @@ export class SessionSearchStore {
|
||||
state: SessionSearchFileState
|
||||
n: number
|
||||
}[]
|
||||
const counts: SessionSearchStateCounts = { current: 0, due: 0, failed: 0 }
|
||||
const counts: SessionSearchStateCounts = {
|
||||
current: 0,
|
||||
due: 0,
|
||||
failed: 0,
|
||||
sessionsByAgent: this.sessionsByAgent()
|
||||
}
|
||||
for (const row of rows) {
|
||||
counts[row.state] = Number(row.n)
|
||||
}
|
||||
return counts
|
||||
}
|
||||
|
||||
// Grouped on `sessions_agent`, over one row per indexed session. Deliberately
|
||||
// not the message count beside it: that would scan every indexed row on a call
|
||||
// the panel polls, and it answers the same question one table later.
|
||||
private sessionsByAgent(): Record<string, number> {
|
||||
const rows = this.db.prepare('SELECT agent, count(*) AS n FROM sessions GROUP BY agent').all()
|
||||
const counts: Record<string, number> = {}
|
||||
for (const row of rows) {
|
||||
const agent = asRecord(row)?.agent
|
||||
const total = asRecord(row)?.n
|
||||
if (typeof agent === 'string' && typeof total === 'number') {
|
||||
counts[agent] = total
|
||||
}
|
||||
}
|
||||
return counts
|
||||
}
|
||||
|
||||
/**
|
||||
* Drops a source's rows. Only a proven deletion may call this: an unreadable
|
||||
* source is `unverifiable`, not `missing`, and keeps its rows
|
||||
|
||||
@@ -6,11 +6,11 @@ import { parseClineSessionFile } from './session-scanner-cline-parser'
|
||||
import { parseGrokSessionFile } from './session-scanner-grok-parser'
|
||||
import { parseMessageGraphSessionFile, parseRovoSessionFile } from './session-scanner-graph-parsers'
|
||||
import { parseKimiSessionFile } from './session-scanner-kimi-parser'
|
||||
import { splitOpenCodeSqliteCandidate } from './session-scanner-opencode-sqlite-paths'
|
||||
import {
|
||||
looksLikeOpenCodeSqliteCandidate,
|
||||
splitOpenCodeSqliteCandidate
|
||||
} from './session-scanner-opencode-sqlite-paths'
|
||||
import { parseOpenCodeSqliteSessionViaWorker } from './session-scanner-opencode-sqlite-worker-spawn'
|
||||
captureOpenCodeSqliteSessionViaWorker,
|
||||
parseOpenCodeSqliteSessionViaWorker
|
||||
} from './session-scanner-opencode-sqlite-worker-spawn'
|
||||
import { parseClaudeSessionFile } from './session-scanner-primary-parsers'
|
||||
import { parseGeminiSessionFile } from './session-scanner-gemini-parsers'
|
||||
import { parseCodexSessionFile } from './session-scanner-codex-parser'
|
||||
@@ -22,12 +22,27 @@ import type { SessionFileCandidate } from './session-scanner-types'
|
||||
import type { TranscriptMessageSink } from './session-transcript-consumers'
|
||||
|
||||
/**
|
||||
* False when a parser decodes its messages somewhere the channel cannot reach.
|
||||
* OpenCode's SQLite sessions are read on a worker thread, so their messages
|
||||
* never come back over the sink and the read must not be reported as complete.
|
||||
* Read an OpenCode SQLite session on the worker thread.
|
||||
*
|
||||
* Two request kinds rather than one, chosen by whether anyone is listening: a
|
||||
* list scan wants the newest few messages for the panel preview, so asking for
|
||||
* the whole transcript would read every part of every session on every refresh.
|
||||
* A read with a sink is the search index's, and that one needs all of it.
|
||||
*/
|
||||
export function parserPublishesMessages(candidate: SessionFileCandidate): boolean {
|
||||
return candidate.agent !== 'opencode' || !looksLikeOpenCodeSqliteCandidate(candidate.file.path)
|
||||
async function readOpenCodeSqliteCandidate(
|
||||
sqliteCandidate: { dbPath: string; sessionId: string },
|
||||
platform: NodeJS.Platform,
|
||||
messages?: TranscriptMessageSink
|
||||
): Promise<AiVaultSession | null> {
|
||||
const request = { ...sqliteCandidate, platform }
|
||||
if (!messages?.active) {
|
||||
return parseOpenCodeSqliteSessionViaWorker(request)
|
||||
}
|
||||
const capture = await captureOpenCodeSqliteSessionViaWorker(request)
|
||||
for (const message of capture.messages) {
|
||||
messages.push(message)
|
||||
}
|
||||
return capture.session
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -70,11 +85,7 @@ export async function parseAgentSessionFile(
|
||||
// real filesystem paths and fall through to the JSON parser.
|
||||
const sqliteCandidate = splitOpenCodeSqliteCandidate(candidate.file.path)
|
||||
if (sqliteCandidate) {
|
||||
return parseOpenCodeSqliteSessionViaWorker({
|
||||
dbPath: sqliteCandidate.dbPath,
|
||||
sessionId: sqliteCandidate.sessionId,
|
||||
platform
|
||||
})
|
||||
return readOpenCodeSqliteCandidate(sqliteCandidate, platform, messages)
|
||||
}
|
||||
return parseOpenCodeSessionFile(candidate.file, platform, messages)
|
||||
}
|
||||
|
||||
@@ -0,0 +1,222 @@
|
||||
import { mkdir, writeFile } from 'node:fs/promises'
|
||||
import { join } from 'node:path'
|
||||
import { writeAntigravityScannerFixture } from './session-scanner-test-fixtures'
|
||||
import { jsonlBody, type AgentVaultRoots } from './session-scanner-vault-roots'
|
||||
|
||||
// The agents whose session is a JSON document, or a directory of them, rewritten
|
||||
// in place rather than appended to. Antigravity rides along here because its
|
||||
// fixture writer already owns the layout.
|
||||
|
||||
/**
|
||||
* Write one session per document-shaped agent.
|
||||
* @param root - The vault root, which Kimi's session index lives directly in.
|
||||
* @param roots - The scan roots to write under.
|
||||
* @param antigravitySessionId - The conversation id Antigravity resumes by.
|
||||
*/
|
||||
export async function writeDocumentAgentFixtures(
|
||||
root: string,
|
||||
roots: AgentVaultRoots,
|
||||
antigravitySessionId: string
|
||||
): Promise<void> {
|
||||
await mkdir(roots.geminiSessionsDir, { recursive: true })
|
||||
await writeFile(
|
||||
join(roots.geminiSessionsDir, 'gemini-session.json'),
|
||||
JSON.stringify({
|
||||
sessionId: 'gemini-session',
|
||||
startTime: '2026-05-01T10:02:00.000Z',
|
||||
lastUpdated: '2026-05-01T10:02:01.000Z',
|
||||
messages: [
|
||||
{
|
||||
type: 'user',
|
||||
timestamp: '2026-05-01T10:02:00.000Z',
|
||||
content: [{ text: 'Gemini title' }]
|
||||
},
|
||||
{
|
||||
type: 'gemini',
|
||||
timestamp: '2026-05-01T10:02:01.000Z',
|
||||
model: 'gemini-2.5-pro',
|
||||
tokens: { input: 10, output: 5 }
|
||||
}
|
||||
]
|
||||
})
|
||||
)
|
||||
|
||||
await writeAntigravityScannerFixture(roots.antigravityBrainDir, antigravitySessionId)
|
||||
|
||||
await mkdir(join(roots.opencodeStorageDir, 'session', 'project'), { recursive: true })
|
||||
await mkdir(join(roots.opencodeStorageDir, 'message', 'opencode-session'), { recursive: true })
|
||||
await writeFile(
|
||||
join(roots.opencodeStorageDir, 'session', 'project', 'ses_opencode.json'),
|
||||
JSON.stringify({
|
||||
id: 'opencode-session',
|
||||
directory: '/tmp/opencode',
|
||||
title: 'OpenCode title',
|
||||
time: { created: 1_777_634_000_000, updated: 1_777_634_001_000 }
|
||||
})
|
||||
)
|
||||
await writeFile(
|
||||
join(roots.opencodeStorageDir, 'message', 'opencode-session', 'msg_1.json'),
|
||||
JSON.stringify({
|
||||
role: 'user',
|
||||
summary: { title: 'OpenCode title' },
|
||||
time: { created: 1_777_634_000_000 },
|
||||
tokens: { input: 7, output: 3 }
|
||||
})
|
||||
)
|
||||
|
||||
await mkdir(join(roots.grokSessionsDir, encodeURIComponent('/tmp/grok'), 'grok-session'), {
|
||||
recursive: true
|
||||
})
|
||||
await writeFile(
|
||||
join(roots.grokSessionsDir, encodeURIComponent('/tmp/grok'), 'grok-session', 'summary.json'),
|
||||
JSON.stringify({
|
||||
info: { id: 'grok-session', cwd: '/tmp/grok' },
|
||||
session_summary: '',
|
||||
created_at: '2026-05-01T10:04:00.000Z',
|
||||
updated_at: '2026-05-01T10:04:01.000Z',
|
||||
num_chat_messages: 2,
|
||||
current_model_id: 'grok-build',
|
||||
head_branch: 'feature/grok-vault'
|
||||
})
|
||||
)
|
||||
await writeFile(
|
||||
join(
|
||||
roots.grokSessionsDir,
|
||||
encodeURIComponent('/tmp/grok'),
|
||||
'grok-session',
|
||||
'chat_history.jsonl'
|
||||
),
|
||||
jsonlBody([
|
||||
{
|
||||
type: 'user',
|
||||
content: [
|
||||
{
|
||||
type: 'text',
|
||||
text: '<user_info>context</user_info><user_query>Grok title</user_query>'
|
||||
}
|
||||
]
|
||||
},
|
||||
{ type: 'assistant', content: 'Done' }
|
||||
])
|
||||
)
|
||||
|
||||
await mkdir(roots.hermesSessionsDir, { recursive: true })
|
||||
await writeFile(
|
||||
join(roots.hermesSessionsDir, 'session_hermes-session.json'),
|
||||
JSON.stringify({
|
||||
session_id: 'hermes-session',
|
||||
model: 'hermes-1',
|
||||
cwd: '/tmp/hermes',
|
||||
session_start: '2026-05-01T10:05:00.000Z',
|
||||
last_updated: '2026-05-01T10:05:01.000Z',
|
||||
messages: [{ role: 'user', content: 'Hermes title' }]
|
||||
})
|
||||
)
|
||||
|
||||
await mkdir(join(roots.rovoSessionsDir, 'rovo-session'), { recursive: true })
|
||||
await writeFile(
|
||||
join(roots.rovoSessionsDir, 'rovo-session', 'metadata.json'),
|
||||
JSON.stringify({ title: 'Rovo title', workspace_path: '/tmp/rovo' })
|
||||
)
|
||||
await writeFile(
|
||||
join(roots.rovoSessionsDir, 'rovo-session', 'session_context.json'),
|
||||
JSON.stringify({
|
||||
message_history: [
|
||||
{
|
||||
kind: 'request',
|
||||
timestamp: '2026-05-01T10:06:00.000Z',
|
||||
parts: [{ part_kind: 'user-prompt', content: 'Rovo title' }]
|
||||
}
|
||||
]
|
||||
})
|
||||
)
|
||||
|
||||
await mkdir(roots.devinTranscriptsDir, { recursive: true })
|
||||
await writeFile(
|
||||
join(roots.devinTranscriptsDir, 'devin-session.json'),
|
||||
JSON.stringify({
|
||||
session_id: 'devin-session',
|
||||
working_directory: '/tmp/devin',
|
||||
agent: { model_name: 'swe-1-6-fast' },
|
||||
steps: [
|
||||
{
|
||||
metadata: {
|
||||
created_at: '2026-05-01T10:10:00.000Z',
|
||||
is_user_input: true,
|
||||
metrics: { input_tokens: 1, output_tokens: 2 }
|
||||
},
|
||||
text: 'Devin vault title'
|
||||
}
|
||||
]
|
||||
})
|
||||
)
|
||||
|
||||
const clineSessionId = 'cline-session'
|
||||
const clineSessionDir = join(roots.clineSessionsDir, clineSessionId)
|
||||
await mkdir(clineSessionDir, { recursive: true })
|
||||
await writeFile(
|
||||
join(clineSessionDir, `${clineSessionId}.json`),
|
||||
JSON.stringify({
|
||||
session_id: clineSessionId,
|
||||
started_at: '2026-05-01T10:10:30.000Z',
|
||||
model: 'cline-model',
|
||||
cwd: '/tmp/cline'
|
||||
})
|
||||
)
|
||||
await writeFile(
|
||||
join(clineSessionDir, `${clineSessionId}.messages.json`),
|
||||
JSON.stringify({
|
||||
updated_at: '2026-05-01T10:10:31.000Z',
|
||||
messages: [{ role: 'user', content: [{ type: 'text', text: 'Cline vault title' }] }]
|
||||
})
|
||||
)
|
||||
|
||||
// Kimi: <sessions>/wd_*/session_*/state.json + sibling agents/main/wire.jsonl,
|
||||
// with the work dir resolved from the top-level session_index.jsonl.
|
||||
const kimiSessionDir = join(roots.kimiSessionsDir, 'wd_app_abc', 'session_kimi-session')
|
||||
await mkdir(join(kimiSessionDir, 'agents', 'main'), { recursive: true })
|
||||
await writeFile(
|
||||
join(kimiSessionDir, 'state.json'),
|
||||
JSON.stringify({
|
||||
createdAt: '2026-05-01T10:11:00.000Z',
|
||||
updatedAt: '2026-05-01T10:11:05.000Z',
|
||||
title: 'Kimi vault title',
|
||||
lastPrompt: 'Kimi vault title',
|
||||
agents: { main: { type: 'main', parentAgentId: null } }
|
||||
})
|
||||
)
|
||||
await writeFile(
|
||||
join(root, 'session_index.jsonl'),
|
||||
jsonlBody([
|
||||
{ sessionId: 'session_kimi-session', sessionDir: kimiSessionDir, workDir: '/tmp/kimi' }
|
||||
])
|
||||
)
|
||||
await writeFile(
|
||||
join(kimiSessionDir, 'agents', 'main', 'wire.jsonl'),
|
||||
jsonlBody([
|
||||
{ type: 'config.update', modelAlias: 'kimi-k2.6', time: 1781853559132 },
|
||||
{
|
||||
type: 'context.append_message',
|
||||
message: {
|
||||
role: 'user',
|
||||
content: [{ type: 'text', text: 'Kimi vault title' }],
|
||||
origin: { kind: 'user' }
|
||||
},
|
||||
time: 1781853559164
|
||||
},
|
||||
{
|
||||
type: 'context.append_loop_event',
|
||||
event: { type: 'content.part', part: { type: 'text', text: 'Kimi reply' } },
|
||||
time: 1781853559177
|
||||
},
|
||||
{ type: 'context.append_loop_event', event: { type: 'step.end' }, time: 1781853559178 },
|
||||
{
|
||||
type: 'usage.record',
|
||||
model: 'kimi-k2.6',
|
||||
usage: { inputOther: 4, output: 6, inputCacheRead: 0, inputCacheCreation: 0 },
|
||||
usageScope: 'turn',
|
||||
time: 1781853559178
|
||||
}
|
||||
])
|
||||
)
|
||||
}
|
||||
@@ -0,0 +1,32 @@
|
||||
import { isolatedScanRoots } from './session-scanner-test-fixtures'
|
||||
import { writeDocumentAgentFixtures } from './session-scanner-document-agent-fixtures'
|
||||
import { writeLogAgentFixtures } from './session-scanner-log-agent-fixtures'
|
||||
|
||||
// Why this is shared rather than inline in one test: it is the only place that
|
||||
// writes one transcript in every supported agent's own format. A scan test and
|
||||
// the search index's capture guard both need exactly that, and a second copy
|
||||
// would drift the moment an agent's layout changed.
|
||||
|
||||
export type EveryAgentVault = {
|
||||
roots: ReturnType<typeof isolatedScanRoots>
|
||||
/** Ids the caller asserts resume commands against. */
|
||||
antigravitySessionId: string
|
||||
/** OMP and Prime Agent resume by absolute transcript path, not by id. */
|
||||
ompSessionFile: string
|
||||
primeAgentSessionFile: string
|
||||
}
|
||||
|
||||
/**
|
||||
* Write one session per supported agent under `root`, each in that agent's own
|
||||
* on-disk layout. OpenCode gets its legacy JSON layout here; its SQLite layout
|
||||
* has its own builder, because it needs a database rather than a tree.
|
||||
* @param root - An empty temporary directory to build the vault in.
|
||||
* @returns The scan roots for `root`, and the ids a caller asserts against.
|
||||
*/
|
||||
export async function writeEveryAgentVault(root: string): Promise<EveryAgentVault> {
|
||||
const roots = isolatedScanRoots(root)
|
||||
const antigravitySessionId = 'aaaaaaaa-bbbb-4ccc-8ddd-eeeeeeeeeeee'
|
||||
const { ompSessionFile, primeAgentSessionFile } = await writeLogAgentFixtures(roots)
|
||||
await writeDocumentAgentFixtures(root, roots, antigravitySessionId)
|
||||
return { roots, antigravitySessionId, ompSessionFile, primeAgentSessionFile }
|
||||
}
|
||||
@@ -27,8 +27,12 @@ const ALLOWLIST = new Set([
|
||||
// On-demand IPC readers, gated in the STA-4049 follow-up.
|
||||
'session-scanner-claude-subagents.ts',
|
||||
'session-scanner-omp-subagent-listing.ts',
|
||||
// Test-only fixture builder.
|
||||
'session-scanner-test-fixtures.ts'
|
||||
// Test-only fixture builders: they create the vault a test reads, so the
|
||||
// paths they touch are temp directories this process just made.
|
||||
'session-scanner-test-fixtures.ts',
|
||||
'session-scanner-document-agent-fixtures.ts',
|
||||
'session-scanner-log-agent-fixtures.ts',
|
||||
'session-scanner-opencode-sqlite-fixture.ts'
|
||||
])
|
||||
|
||||
const FS_IMPORT = /import\s+([\s\S]*?)\s+from\s+['"]node:fs(?:\/promises)?['"]/g
|
||||
|
||||
@@ -0,0 +1,160 @@
|
||||
import { mkdir, writeFile } from 'node:fs/promises'
|
||||
import { join } from 'node:path'
|
||||
import {
|
||||
writeOmpScannerFixture,
|
||||
writePrimeAgentScannerFixture
|
||||
} from './session-scanner-test-fixtures'
|
||||
import { jsonlBody, type AgentVaultRoots } from './session-scanner-vault-roots'
|
||||
|
||||
// The agents whose session is an append-only JSONL log the CLI writes a record
|
||||
// at a time. Split from the document-shaped agents purely by file size; the two
|
||||
// halves are called together and neither is meaningful alone.
|
||||
|
||||
/**
|
||||
* Write one append-only transcript per log-shaped agent.
|
||||
* @param roots - The scan roots to write under.
|
||||
* @returns The transcript paths OMP and Prime Agent resume by.
|
||||
*/
|
||||
export async function writeLogAgentFixtures(
|
||||
roots: AgentVaultRoots
|
||||
): Promise<{ ompSessionFile: string; primeAgentSessionFile: string }> {
|
||||
await mkdir(join(roots.claudeProjectsDir, 'project'), { recursive: true })
|
||||
await writeFile(
|
||||
join(roots.claudeProjectsDir, 'project', 'claude-session.jsonl'),
|
||||
jsonlBody([
|
||||
{
|
||||
type: 'user',
|
||||
sessionId: 'claude-session',
|
||||
timestamp: '2026-05-01T10:00:00.000Z',
|
||||
cwd: '/tmp/claude',
|
||||
message: { role: 'user', content: 'Claude title' }
|
||||
}
|
||||
])
|
||||
)
|
||||
|
||||
await mkdir(join(roots.codexSessionsDir, '2026', '05', '01'), { recursive: true })
|
||||
await writeFile(
|
||||
join(roots.codexSessionsDir, '2026', '05', '01', 'rollout-2026-codex-session.jsonl'),
|
||||
jsonlBody([
|
||||
{
|
||||
timestamp: '2026-05-01T10:01:00.000Z',
|
||||
type: 'session_meta',
|
||||
payload: { id: 'codex-session', cwd: '/tmp/codex' }
|
||||
},
|
||||
{
|
||||
timestamp: '2026-05-01T10:01:01.000Z',
|
||||
type: 'response_item',
|
||||
payload: {
|
||||
type: 'message',
|
||||
role: 'user',
|
||||
content: [{ type: 'text', text: 'Codex title' }]
|
||||
}
|
||||
}
|
||||
])
|
||||
)
|
||||
|
||||
await mkdir(roots.copilotSessionsDir, { recursive: true })
|
||||
await writeFile(
|
||||
join(roots.copilotSessionsDir, 'copilot-session.jsonl'),
|
||||
jsonlBody([
|
||||
{
|
||||
type: 'session.start',
|
||||
data: { sessionId: 'copilot-session', startTime: '2026-05-01T10:03:00.000Z' },
|
||||
timestamp: '2026-05-01T10:03:00.000Z'
|
||||
},
|
||||
{
|
||||
type: 'session.info',
|
||||
data: {
|
||||
infoType: 'folder_trust',
|
||||
message: 'Folder /tmp/copilot has been added to trusted folders.'
|
||||
},
|
||||
timestamp: '2026-05-01T10:03:01.000Z'
|
||||
},
|
||||
{
|
||||
type: 'user.message',
|
||||
data: { transformedContent: 'Copilot title' },
|
||||
timestamp: '2026-05-01T10:03:02.000Z'
|
||||
}
|
||||
])
|
||||
)
|
||||
|
||||
await mkdir(join(roots.cursorProjectsDir, 'project', 'agent-transcripts'), { recursive: true })
|
||||
await writeFile(
|
||||
join(roots.cursorProjectsDir, 'project', 'agent-transcripts', 'cursor-session.jsonl'),
|
||||
jsonlBody([
|
||||
{
|
||||
role: 'user',
|
||||
message: { content: [{ type: 'text', text: 'Cursor title' }] }
|
||||
},
|
||||
{ role: 'assistant', message: { content: [{ type: 'text', text: 'Done' }] } }
|
||||
])
|
||||
)
|
||||
|
||||
await mkdir(join(roots.openclawStateDir, 'agents', 'default', 'sessions'), { recursive: true })
|
||||
await writeFile(
|
||||
join(roots.openclawStateDir, 'agents', 'default', 'sessions', 'openclaw-session.jsonl'),
|
||||
jsonlBody([
|
||||
{
|
||||
type: 'session',
|
||||
id: 'openclaw-session',
|
||||
timestamp: '2026-05-01T10:07:00.000Z',
|
||||
cwd: '/tmp/openclaw'
|
||||
},
|
||||
{
|
||||
type: 'message',
|
||||
timestamp: '2026-05-01T10:07:01.000Z',
|
||||
message: { role: 'user', content: [{ type: 'text', text: 'OpenClaw title' }] }
|
||||
}
|
||||
])
|
||||
)
|
||||
|
||||
await mkdir(roots.piSessionsDir, { recursive: true })
|
||||
await writeFile(
|
||||
join(roots.piSessionsDir, 'pi-session.jsonl'),
|
||||
jsonlBody([
|
||||
{
|
||||
type: 'session',
|
||||
id: 'pi-session',
|
||||
timestamp: '2026-05-01T10:08:00.000Z',
|
||||
cwd: '/tmp/pi'
|
||||
},
|
||||
{
|
||||
type: 'message',
|
||||
timestamp: '2026-05-01T10:08:01.000Z',
|
||||
message: { role: 'user', content: [{ type: 'text', text: 'Pi title' }] }
|
||||
}
|
||||
])
|
||||
)
|
||||
|
||||
const ompSessionFile = await writeOmpScannerFixture(roots.ompSessionsDir)
|
||||
const primeAgentSessionFile = await writePrimeAgentScannerFixture(roots.primeAgentSessionsDir)
|
||||
|
||||
await mkdir(roots.droidSessionsDir, { recursive: true })
|
||||
await writeFile(
|
||||
join(roots.droidSessionsDir, 'droid-session.jsonl'),
|
||||
jsonlBody([
|
||||
{
|
||||
type: 'system',
|
||||
session_id: 'droid-session',
|
||||
timestamp: '2026-05-01T10:09:00.000Z',
|
||||
model: 'droid-model',
|
||||
cwd: '/tmp/droid'
|
||||
},
|
||||
{
|
||||
type: 'message',
|
||||
session_id: 'droid-session',
|
||||
timestamp: '2026-05-01T10:09:01.000Z',
|
||||
role: 'user',
|
||||
text: 'Droid title'
|
||||
},
|
||||
{
|
||||
type: 'completion',
|
||||
session_id: 'droid-session',
|
||||
timestamp: '2026-05-01T10:09:02.000Z',
|
||||
usage: { input_tokens: 2, output_tokens: 3 }
|
||||
}
|
||||
])
|
||||
)
|
||||
|
||||
return { ompSessionFile, primeAgentSessionFile }
|
||||
}
|
||||
@@ -0,0 +1,161 @@
|
||||
import type { AiVaultSession } from '../../shared/ai-vault-types'
|
||||
import { timestampIso } from './session-scanner-accumulator'
|
||||
import { asRecord } from './session-scanner-record-value'
|
||||
import { extractPartText, readOpenCodeSqliteSession } from './session-scanner-opencode-sqlite'
|
||||
import { readOpenCodeDatabase } from './session-scanner-opencode-sqlite-open'
|
||||
import { canReadOpenCodeMessageParts } from './session-scanner-opencode-sqlite-schema'
|
||||
import type { TranscriptMessage, TranscriptMessageRole } from './session-transcript-consumers'
|
||||
import { boundedText } from './session-transcript-message-content'
|
||||
import type SyncDatabase from '../sqlite/sync-database'
|
||||
|
||||
// Why: the session list needs the newest few messages, and the search index
|
||||
// needs every one of them. That is the only difference between this read and
|
||||
// `parseOpenCodeSqliteSession`, so the decoding is shared and only the query
|
||||
// that selects the rows differs.
|
||||
|
||||
/**
|
||||
* How many text parts one session may hold before this read gives up.
|
||||
*
|
||||
* A safety valve on memory, not a policy: the rows are materialized and then
|
||||
* posted across the worker boundary, so an unbounded session would be held
|
||||
* twice. Exceeding it throws rather than returning a prefix, because a prefix
|
||||
* committed under a complete-read cursor would leave the tail unsearchable with
|
||||
* nothing on the row to say so. A failed read is retried and surfaces; a silent
|
||||
* truncation does neither. 20,000 text parts is roughly 10,000 conversation
|
||||
* turns, far past any real session.
|
||||
*/
|
||||
const OPENCODE_CAPTURE_PART_LIMIT = 20_000
|
||||
|
||||
type CaptureRow = {
|
||||
messageId: string
|
||||
role: string | null
|
||||
partData: string
|
||||
messageTimeMs: number
|
||||
}
|
||||
|
||||
// A row this build cannot read is dropped rather than failing the session: the
|
||||
// schema probe only proves the columns exist, not what any one row holds.
|
||||
function toCaptureRow(value: unknown): CaptureRow | null {
|
||||
const record = asRecord(value)
|
||||
if (!record) {
|
||||
return null
|
||||
}
|
||||
const { message_id: messageId, role, part_data: partData, message_time: messageTime } = record
|
||||
if (
|
||||
typeof messageId !== 'string' ||
|
||||
typeof partData !== 'string' ||
|
||||
typeof messageTime !== 'number'
|
||||
) {
|
||||
return null
|
||||
}
|
||||
return {
|
||||
messageId,
|
||||
role: typeof role === 'string' ? role : null,
|
||||
partData,
|
||||
messageTimeMs: messageTime
|
||||
}
|
||||
}
|
||||
|
||||
/** The session the panel renders, and every message the index folds. */
|
||||
export type OpenCodeSqliteCapture = {
|
||||
session: AiVaultSession | null
|
||||
messages: TranscriptMessage[]
|
||||
}
|
||||
|
||||
function captureRole(role: string | null): TranscriptMessageRole | null {
|
||||
return role === 'user' || role === 'assistant' ? role : null
|
||||
}
|
||||
|
||||
function buildCaptureQuery(): string {
|
||||
// Message order, then part order within a message: the same key the preview
|
||||
// read uses, run forwards and without the newest-N window.
|
||||
return `SELECT m.id AS message_id,
|
||||
json_extract(m.data, '$.role') AS role,
|
||||
p.data AS part_data,
|
||||
m.time_created AS message_time
|
||||
FROM message m
|
||||
JOIN part p ON p.message_id = m.id
|
||||
WHERE m.session_id = ?
|
||||
AND json_extract(m.data, '$.role') IN ('user','assistant')
|
||||
AND json_extract(p.data, '$.type') = 'text'
|
||||
ORDER BY m.time_created ASC, m.id ASC, p.time_created ASC, p.rowid ASC
|
||||
LIMIT ?`
|
||||
}
|
||||
|
||||
/**
|
||||
* Decode one session's whole transcript, one message per `message` row.
|
||||
*
|
||||
* Parts are joined rather than emitted separately because every other provider
|
||||
* hands a consumer one message per turn; a phrase that runs across two blocks of
|
||||
* the same turn is then still one indexable row.
|
||||
*/
|
||||
export function readOpenCodeSessionMessages(
|
||||
db: SyncDatabase,
|
||||
sessionId: string
|
||||
): TranscriptMessage[] {
|
||||
if (!canReadOpenCodeMessageParts(db)) {
|
||||
return []
|
||||
}
|
||||
const rows = db.prepare(buildCaptureQuery()).all(sessionId, OPENCODE_CAPTURE_PART_LIMIT + 1)
|
||||
if (rows.length > OPENCODE_CAPTURE_PART_LIMIT) {
|
||||
throw new Error(
|
||||
`OpenCode session ${sessionId} holds more than ${OPENCODE_CAPTURE_PART_LIMIT} text parts; its transcript was not read.`
|
||||
)
|
||||
}
|
||||
|
||||
const messages: TranscriptMessage[] = []
|
||||
let openMessageId: string | null = null
|
||||
let openParts: string[] = []
|
||||
let openRole: TranscriptMessageRole | null = null
|
||||
let openTimestamp: string | null = null
|
||||
|
||||
const flush = (): void => {
|
||||
const text = openRole && openParts.length > 0 ? boundedText(openParts.join('\n')) : null
|
||||
if (openRole && text) {
|
||||
messages.push({ role: openRole, text, timestamp: openTimestamp })
|
||||
}
|
||||
openParts = []
|
||||
}
|
||||
|
||||
for (const value of rows) {
|
||||
const row = toCaptureRow(value)
|
||||
if (!row) {
|
||||
continue
|
||||
}
|
||||
if (row.messageId !== openMessageId) {
|
||||
flush()
|
||||
openMessageId = row.messageId
|
||||
openRole = captureRole(row.role)
|
||||
openTimestamp = timestampIso(row.messageTimeMs)
|
||||
}
|
||||
const text = extractPartText(row.partData)
|
||||
if (text) {
|
||||
openParts.push(text)
|
||||
}
|
||||
}
|
||||
flush()
|
||||
return messages
|
||||
}
|
||||
|
||||
/**
|
||||
* Read one OpenCode session and its whole transcript from a single open of the
|
||||
* database, so the two can never describe different generations of the session.
|
||||
* @param args.dbPath - Absolute path to the opencode.db file.
|
||||
* @param args.sessionId - Primary key in the `session` table.
|
||||
* @param args.platform - Platform used for resume-command generation.
|
||||
* @returns The parsed session (null when it does not exist) and its messages.
|
||||
*/
|
||||
export async function captureOpenCodeSqliteSession(args: {
|
||||
dbPath: string
|
||||
sessionId: string
|
||||
platform: NodeJS.Platform
|
||||
}): Promise<OpenCodeSqliteCapture> {
|
||||
return readOpenCodeDatabase({
|
||||
dbPath: args.dbPath,
|
||||
read: (db) => {
|
||||
const session = readOpenCodeSqliteSession({ db, ...args })
|
||||
// No session row is no transcript: the id names nothing in this database.
|
||||
return { session, messages: session ? readOpenCodeSessionMessages(db, args.sessionId) : [] }
|
||||
}
|
||||
})
|
||||
}
|
||||
@@ -0,0 +1,174 @@
|
||||
import { mkdirSync } from 'node:fs'
|
||||
import { dirname } from 'node:path'
|
||||
import SyncDatabase from '../sqlite/sync-database'
|
||||
|
||||
// The OpenCode 1.17.x schema, as the app itself creates it. Written out in full
|
||||
// rather than trimmed to the columns a reader names, because every read probes
|
||||
// for its columns and a trimmed fixture would pass a probe the real database
|
||||
// fails (or the reverse) without the test being able to tell.
|
||||
|
||||
const OPENCODE_SCHEMA = `
|
||||
CREATE TABLE session (
|
||||
id TEXT PRIMARY KEY, project_id TEXT NOT NULL, parent_id TEXT, slug TEXT NOT NULL,
|
||||
directory TEXT NOT NULL, title TEXT NOT NULL, version TEXT NOT NULL, share_url TEXT,
|
||||
summary_additions INTEGER, summary_deletions INTEGER, summary_files INTEGER,
|
||||
summary_diffs TEXT, revert TEXT, permission TEXT,
|
||||
time_created INTEGER NOT NULL, time_updated INTEGER NOT NULL, time_compacting INTEGER,
|
||||
time_archived INTEGER, workspace_id TEXT, path TEXT, agent TEXT, model TEXT,
|
||||
cost REAL DEFAULT 0 NOT NULL, tokens_input INTEGER DEFAULT 0 NOT NULL,
|
||||
tokens_output INTEGER DEFAULT 0 NOT NULL, tokens_reasoning INTEGER DEFAULT 0 NOT NULL,
|
||||
tokens_cache_read INTEGER DEFAULT 0 NOT NULL, tokens_cache_write INTEGER DEFAULT 0 NOT NULL,
|
||||
metadata TEXT
|
||||
);
|
||||
CREATE TABLE message (
|
||||
id TEXT PRIMARY KEY, session_id TEXT NOT NULL, time_created INTEGER NOT NULL,
|
||||
time_updated INTEGER NOT NULL, data TEXT NOT NULL
|
||||
);
|
||||
CREATE TABLE project (
|
||||
id TEXT PRIMARY KEY, worktree TEXT NOT NULL, vcs TEXT, name TEXT, icon_url TEXT,
|
||||
icon_color TEXT, time_created INTEGER NOT NULL, time_updated INTEGER NOT NULL,
|
||||
time_initialized INTEGER, sandboxes TEXT NOT NULL, commands TEXT, icon_url_override TEXT
|
||||
);
|
||||
CREATE TABLE part (
|
||||
id TEXT PRIMARY KEY, message_id TEXT NOT NULL, session_id TEXT NOT NULL,
|
||||
time_created INTEGER NOT NULL, time_updated INTEGER NOT NULL, data TEXT NOT NULL
|
||||
);
|
||||
`
|
||||
|
||||
export const OPENCODE_FIXTURE_EPOCH_MS = 1_740_000_000_000
|
||||
|
||||
export type OpenCodeSqliteFixtureTurn = {
|
||||
role: 'user' | 'assistant'
|
||||
/** One text part per string, in the order the session recorded them. */
|
||||
parts: readonly string[]
|
||||
}
|
||||
|
||||
export type OpenCodeSqliteFixtureSession = {
|
||||
id: string
|
||||
title?: string
|
||||
directory?: string
|
||||
turns: readonly OpenCodeSqliteFixtureTurn[]
|
||||
}
|
||||
|
||||
/**
|
||||
* Create an OpenCode SQLite database holding `sessions`.
|
||||
*
|
||||
* Each turn's parts are written as separate `part` rows, which is the shape a
|
||||
* reader has to reassemble; a fixture with one part per turn would never
|
||||
* exercise it. A session's `time_updated` is its last turn's timestamp, the
|
||||
* same stat the real database moves when a session gains a message.
|
||||
* @param dbPath - Where to create the database; parent directories are created.
|
||||
* @param sessions - The sessions to write, in the order they were created.
|
||||
*/
|
||||
export function writeOpenCodeSqliteDatabase(
|
||||
dbPath: string,
|
||||
sessions: readonly OpenCodeSqliteFixtureSession[]
|
||||
): void {
|
||||
mkdirSync(dirname(dbPath), { recursive: true })
|
||||
const db = new SyncDatabase(dbPath)
|
||||
try {
|
||||
if (!tableAlreadyThere(db)) {
|
||||
db.exec(OPENCODE_SCHEMA)
|
||||
db.prepare(
|
||||
`INSERT INTO project (id, worktree, name, time_created, time_updated, sandboxes)
|
||||
VALUES ('proj-1', '/tmp/opencode', 'proj', ?, ?, '[]')`
|
||||
).run(OPENCODE_FIXTURE_EPOCH_MS, OPENCODE_FIXTURE_EPOCH_MS)
|
||||
}
|
||||
for (const session of sessions) {
|
||||
writeSession(db, session)
|
||||
}
|
||||
} finally {
|
||||
db.close()
|
||||
}
|
||||
}
|
||||
|
||||
function tableAlreadyThere(db: SyncDatabase): boolean {
|
||||
return (
|
||||
db.prepare(`SELECT name FROM sqlite_master WHERE type='table' AND name='session'`).get() !==
|
||||
undefined
|
||||
)
|
||||
}
|
||||
|
||||
function writeSession(db: SyncDatabase, session: OpenCodeSqliteFixtureSession): void {
|
||||
const created = OPENCODE_FIXTURE_EPOCH_MS
|
||||
const updated = created + Math.max(1, session.turns.length) * 60_000
|
||||
db.prepare(
|
||||
`INSERT INTO session (id, project_id, parent_id, slug, directory, title, version,
|
||||
time_created, time_updated, agent, model, cost, tokens_input, tokens_output,
|
||||
tokens_reasoning, tokens_cache_read, tokens_cache_write)
|
||||
VALUES (?, 'proj-1', NULL, 'slug-1', ?, ?, '1.0.0', ?, ?, 'build', '{"id":"glm"}',
|
||||
0, 1, 1, 0, 0, 0)
|
||||
ON CONFLICT(id) DO UPDATE SET time_updated = excluded.time_updated`
|
||||
).run(
|
||||
session.id,
|
||||
session.directory ?? '/tmp/opencode',
|
||||
session.title ?? 'OpenCode title',
|
||||
created,
|
||||
updated
|
||||
)
|
||||
appendTurns(db, session, created)
|
||||
}
|
||||
|
||||
/** Appends `turns` after whatever the session already holds. */
|
||||
export function appendTurns(
|
||||
db: SyncDatabase,
|
||||
session: OpenCodeSqliteFixtureSession,
|
||||
startMs: number
|
||||
): void {
|
||||
const insertMessage = db.prepare(
|
||||
`INSERT INTO message (id, session_id, time_created, time_updated, data) VALUES (?, ?, ?, ?, ?)`
|
||||
)
|
||||
const insertPart = db.prepare(
|
||||
`INSERT INTO part (id, message_id, session_id, time_created, time_updated, data)
|
||||
VALUES (?, ?, ?, ?, ?, ?)`
|
||||
)
|
||||
session.turns.forEach((turn, turnIndex) => {
|
||||
const at = startMs + (turnIndex + 1) * 60_000
|
||||
const messageId = `${session.id}-msg-${turnIndex}-${at}`
|
||||
insertMessage.run(
|
||||
messageId,
|
||||
session.id,
|
||||
at,
|
||||
at,
|
||||
JSON.stringify({ role: turn.role, time: { created: at } })
|
||||
)
|
||||
turn.parts.forEach((text, partIndex) => {
|
||||
insertPart.run(
|
||||
`${messageId}-part-${partIndex}`,
|
||||
messageId,
|
||||
session.id,
|
||||
at + partIndex,
|
||||
at + partIndex,
|
||||
JSON.stringify({ type: 'text', text })
|
||||
)
|
||||
})
|
||||
})
|
||||
}
|
||||
|
||||
/**
|
||||
* Add one turn to an existing session and move its `time_updated`, the way
|
||||
* OpenCode does when a session continues.
|
||||
* @param dbPath - The fixture database to append to.
|
||||
* @param sessionId - The session to continue.
|
||||
* @param turn - The turn to append.
|
||||
*/
|
||||
export function appendOpenCodeSqliteTurn(
|
||||
dbPath: string,
|
||||
sessionId: string,
|
||||
turn: OpenCodeSqliteFixtureTurn
|
||||
): void {
|
||||
const db = new SyncDatabase(dbPath)
|
||||
try {
|
||||
const updated = currentUpdatedMs(db, sessionId)
|
||||
appendTurns(db, { id: sessionId, turns: [turn] }, updated)
|
||||
db.prepare('UPDATE session SET time_updated = ? WHERE id = ?').run(updated + 60_000, sessionId)
|
||||
} finally {
|
||||
db.close()
|
||||
}
|
||||
}
|
||||
|
||||
function currentUpdatedMs(db: SyncDatabase, sessionId: string): number {
|
||||
const row = db.prepare('SELECT time_updated FROM session WHERE id = ?').get(sessionId)
|
||||
const updated = row === undefined ? undefined : Object.values(row)[0]
|
||||
return typeof updated === 'number' ? updated : OPENCODE_FIXTURE_EPOCH_MS
|
||||
}
|
||||
@@ -1,12 +1,15 @@
|
||||
import { LazyWorkerThreadHost, type WorkerThreadFactory } from '../lazy-worker-thread-host'
|
||||
import type { AiVaultScanIssue, AiVaultSession } from '../../shared/ai-vault-types'
|
||||
import type {
|
||||
OpenCodeSqliteCaptureRequest,
|
||||
OpenCodeSqliteCaptureValue,
|
||||
OpenCodeSqliteListRequest,
|
||||
OpenCodeSqliteListValue,
|
||||
OpenCodeSqliteParseRequest,
|
||||
OpenCodeSqliteWorkerRequest,
|
||||
OpenCodeSqliteWorkerResponse
|
||||
} from './session-scanner-opencode-sqlite-worker-protocol'
|
||||
import { parseOpenCodeSqliteCaptureValue } from './session-scanner-opencode-sqlite-worker-response'
|
||||
import type { SessionFileCandidate } from './session-scanner-types'
|
||||
import { errorMessage } from './session-scanner-values'
|
||||
|
||||
@@ -19,6 +22,9 @@ import { errorMessage } from './session-scanner-values'
|
||||
|
||||
export const LIST_TIMEOUT_MS = 30_000
|
||||
export const PARSE_TIMEOUT_MS = 15_000
|
||||
// Longer than a parse because it reads every part of the session rather than
|
||||
// the newest window, and shorter than nothing at all because the queue is FIFO.
|
||||
export const CAPTURE_TIMEOUT_MS = 30_000
|
||||
export const IDLE_TEARDOWN_MS = 30_000
|
||||
// After this many consecutive worker deaths, fail the remaining queued calls to
|
||||
// scan issues instead of respawning so a DB that reliably kills the worker can't
|
||||
@@ -31,6 +37,7 @@ export const MAX_CONSECUTIVE_DEATHS = 3
|
||||
type OpenCodeSqliteRequestBody =
|
||||
| Omit<OpenCodeSqliteListRequest, 'id'>
|
||||
| Omit<OpenCodeSqliteParseRequest, 'id'>
|
||||
| Omit<OpenCodeSqliteCaptureRequest, 'id'>
|
||||
|
||||
type PendingCall = {
|
||||
request: OpenCodeSqliteWorkerRequest
|
||||
@@ -44,6 +51,15 @@ type PendingCall = {
|
||||
// can surface a precise issue while keeping synchronous SQLite off the main thread.
|
||||
class OpenCodeSqliteWorkerUnavailableError extends Error {}
|
||||
|
||||
// One session failed, not the whole source: the scanner turns this throw into a
|
||||
// per-session scan issue and the search index records a failed read.
|
||||
function sessionReadFailure(err: unknown): Error {
|
||||
if (err instanceof OpenCodeSqliteWorkerUnavailableError) {
|
||||
return new Error('OpenCode SQLite background scanner could not start.')
|
||||
}
|
||||
return err instanceof Error ? err : new Error(String(err))
|
||||
}
|
||||
|
||||
/**
|
||||
* Main-thread bridge that runs OpenCode SQLite reads on a persistent worker
|
||||
* thread. Dispatches one request at a time (FIFO), times each request out from
|
||||
@@ -140,13 +156,43 @@ export class OpenCodeSqliteWorkerClient {
|
||||
{ kind: 'parse', dbPath: args.dbPath, sessionId: args.sessionId, platform: args.platform },
|
||||
PARSE_TIMEOUT_MS
|
||||
)
|
||||
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: the worker's parse leg returns exactly this, built by the repo's own reader on the other side of a structured clone.
|
||||
return value as AiVaultSession | null
|
||||
} catch (err) {
|
||||
if (err instanceof OpenCodeSqliteWorkerUnavailableError) {
|
||||
throw new Error('OpenCode SQLite background scanner could not start.')
|
||||
}
|
||||
// Reject only this session; the scanner turns the throw into a scan issue.
|
||||
throw err instanceof Error ? err : new Error(String(err))
|
||||
throw sessionReadFailure(err)
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Read one OpenCode session and its whole transcript on the worker.
|
||||
*
|
||||
* One request rather than a parse plus a second read: both halves then come
|
||||
* from a single open of the database, so the messages the index folds cannot
|
||||
* belong to a different generation of the session than the panel shows.
|
||||
* @param args.dbPath - Absolute path to the opencode.db file.
|
||||
* @param args.sessionId - Primary key in the `session` table.
|
||||
* @param args.platform - Platform used for resume-command generation.
|
||||
* @returns The session (null when it does not exist) and its messages;
|
||||
* rejects on worker timeout/crash so the read is recorded as failed.
|
||||
*/
|
||||
async capture(args: {
|
||||
dbPath: string
|
||||
sessionId: string
|
||||
platform: NodeJS.Platform
|
||||
}): Promise<OpenCodeSqliteCaptureValue> {
|
||||
try {
|
||||
const value = await this.dispatch(
|
||||
{
|
||||
kind: 'capture',
|
||||
dbPath: args.dbPath,
|
||||
sessionId: args.sessionId,
|
||||
platform: args.platform
|
||||
},
|
||||
CAPTURE_TIMEOUT_MS
|
||||
)
|
||||
return parseOpenCodeSqliteCaptureValue(value)
|
||||
} catch (err) {
|
||||
throw sessionReadFailure(err)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
import { parentPort } from 'node:worker_threads'
|
||||
import type { AiVaultScanIssue } from '../../shared/ai-vault-types'
|
||||
import { captureOpenCodeSqliteSession } from './session-scanner-opencode-sqlite-capture'
|
||||
import { listOpenCodeSqliteSessions } from './session-scanner-opencode-sqlite-list'
|
||||
import { parseOpenCodeSqliteSession } from './session-scanner-opencode-sqlite'
|
||||
import type {
|
||||
@@ -30,6 +31,14 @@ async function handleRequest(
|
||||
})
|
||||
return { id: request.id, ok: true, value: { candidates, issues } }
|
||||
}
|
||||
if (request.kind === 'capture') {
|
||||
const capture = await captureOpenCodeSqliteSession({
|
||||
dbPath: request.dbPath,
|
||||
sessionId: request.sessionId,
|
||||
platform: request.platform
|
||||
})
|
||||
return { id: request.id, ok: true, value: capture }
|
||||
}
|
||||
const session = await parseOpenCodeSqliteSession({
|
||||
dbPath: request.dbPath,
|
||||
sessionId: request.sessionId,
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
import type { AiVaultScanIssue } from '../../shared/ai-vault-types'
|
||||
import type { AiVaultScanIssue, AiVaultSession } from '../../shared/ai-vault-types'
|
||||
import type { SessionFileCandidate } from './session-scanner-types'
|
||||
import type { TranscriptMessage } from './session-transcript-consumers'
|
||||
|
||||
// Why: request/response shapes shared by the worker entry and the main-thread
|
||||
// client. Kept type-only (and electron-free) so importing it into the worker
|
||||
@@ -20,7 +21,21 @@ export type OpenCodeSqliteParseRequest = {
|
||||
platform: NodeJS.Platform
|
||||
}
|
||||
|
||||
export type OpenCodeSqliteWorkerRequest = OpenCodeSqliteListRequest | OpenCodeSqliteParseRequest
|
||||
// Same arguments as `parse`, different answer: the session plus every message
|
||||
// the session holds. Its own kind rather than a flag on `parse` so the two
|
||||
// response shapes stay distinguishable at the type level on both sides.
|
||||
export type OpenCodeSqliteCaptureRequest = {
|
||||
id: number
|
||||
kind: 'capture'
|
||||
dbPath: string
|
||||
sessionId: string
|
||||
platform: NodeJS.Platform
|
||||
}
|
||||
|
||||
export type OpenCodeSqliteWorkerRequest =
|
||||
| OpenCodeSqliteListRequest
|
||||
| OpenCodeSqliteParseRequest
|
||||
| OpenCodeSqliteCaptureRequest
|
||||
|
||||
// The list leg returns candidates plus the issues it accumulated; the worker
|
||||
// mutates a local array and hands it back so the caller can merge it into the
|
||||
@@ -30,6 +45,14 @@ export type OpenCodeSqliteListValue = {
|
||||
issues: AiVaultScanIssue[]
|
||||
}
|
||||
|
||||
// The session the panel shows, and the transcript the search index folds. Both
|
||||
// come from one open of the database, so the two can never disagree about which
|
||||
// generation of the session they describe.
|
||||
export type OpenCodeSqliteCaptureValue = {
|
||||
session: AiVaultSession | null
|
||||
messages: TranscriptMessage[]
|
||||
}
|
||||
|
||||
export type OpenCodeSqliteWorkerResponse =
|
||||
| { id: number; ok: true; value: unknown }
|
||||
| { id: number; ok: false; error: string }
|
||||
|
||||
@@ -0,0 +1,41 @@
|
||||
import type { AiVaultSession } from '../../shared/ai-vault-types'
|
||||
import { asRecord } from './session-scanner-record-value'
|
||||
import type { OpenCodeSqliteCaptureValue } from './session-scanner-opencode-sqlite-worker-protocol'
|
||||
import type { TranscriptMessage } from './session-transcript-consumers'
|
||||
|
||||
// Why: a worker posts back a structured clone, which arrives as `unknown`. The
|
||||
// messages are checked one by one because they are written into the index as
|
||||
// rows keyed by role, and a value with no role at all would land under none.
|
||||
|
||||
function isTranscriptMessage(value: unknown): value is TranscriptMessage {
|
||||
const record = asRecord(value)
|
||||
return (
|
||||
record !== null &&
|
||||
(record.role === 'user' || record.role === 'assistant' || record.role === 'tool') &&
|
||||
typeof record.text === 'string' &&
|
||||
(record.timestamp === null || typeof record.timestamp === 'string')
|
||||
)
|
||||
}
|
||||
|
||||
/**
|
||||
* Read a `capture` response from the OpenCode SQLite worker.
|
||||
*
|
||||
* A message that is not one is dropped rather than failing the read: the rest
|
||||
* of the session is still worth indexing, and a row with a role the index has
|
||||
* no column for would be written under an empty one.
|
||||
* @param value - The worker's response value.
|
||||
* @returns The session and the messages the response carried.
|
||||
*/
|
||||
export function parseOpenCodeSqliteCaptureValue(value: unknown): OpenCodeSqliteCaptureValue {
|
||||
const record = asRecord(value)
|
||||
if (!record) {
|
||||
return { session: null, messages: [] }
|
||||
}
|
||||
const messages = Array.isArray(record.messages) ? record.messages.filter(isTranscriptMessage) : []
|
||||
// Held to the same standard as the parse leg rather than validated harder: a
|
||||
// session this build dropped here but kept there would be in the panel and
|
||||
// absent from the index, which is worse than trusting our own worker.
|
||||
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: the worker builds this with the repo's own reader; only the structured clone sits between.
|
||||
const session = (record.session ?? null) as AiVaultSession | null
|
||||
return { session, messages }
|
||||
}
|
||||
@@ -3,6 +3,7 @@ import { join } from 'node:path'
|
||||
import { Worker } from 'node:worker_threads'
|
||||
import type { AiVaultScanIssue, AiVaultSession } from '../../shared/ai-vault-types'
|
||||
import type { SessionFileCandidate } from './session-scanner-types'
|
||||
import type { OpenCodeSqliteCaptureValue } from './session-scanner-opencode-sqlite-worker-protocol'
|
||||
import { OpenCodeSqliteWorkerClient } from './session-scanner-opencode-sqlite-worker-client'
|
||||
|
||||
// Why: resolve the built worker entry + own the process-wide shared client so
|
||||
@@ -70,3 +71,19 @@ export function parseOpenCodeSqliteSessionViaWorker(args: {
|
||||
}): Promise<AiVaultSession | null> {
|
||||
return getSharedClient().parse(args)
|
||||
}
|
||||
|
||||
/**
|
||||
* Read one OpenCode SQLite session and its whole transcript through the shared
|
||||
* worker client.
|
||||
* @param args.dbPath - Absolute path to the opencode.db file.
|
||||
* @param args.sessionId - Primary key in the `session` table.
|
||||
* @param args.platform - Platform used for resume-command generation.
|
||||
* @returns The session and every message it holds.
|
||||
*/
|
||||
export function captureOpenCodeSqliteSessionViaWorker(args: {
|
||||
dbPath: string
|
||||
sessionId: string
|
||||
platform: NodeJS.Platform
|
||||
}): Promise<OpenCodeSqliteCaptureValue> {
|
||||
return getSharedClient().capture(args)
|
||||
}
|
||||
|
||||
@@ -130,7 +130,8 @@ function mapPreviewRole(role: string | null): AiVaultSessionPreviewMessage['role
|
||||
return 'unknown'
|
||||
}
|
||||
|
||||
function extractPartText(partData: string): string | null {
|
||||
/** The text a `type: 'text'` part carries; null for every other part shape. */
|
||||
export function extractPartText(partData: string): string | null {
|
||||
try {
|
||||
const parsed = JSON.parse(partData) as unknown
|
||||
const record =
|
||||
@@ -233,12 +234,13 @@ export async function parseOpenCodeSqliteSession(args: {
|
||||
}): Promise<AiVaultSession | null> {
|
||||
return readOpenCodeDatabase({
|
||||
dbPath: args.dbPath,
|
||||
read: (db) => readSession({ db, ...args })
|
||||
read: (db) => readOpenCodeSqliteSession({ db, ...args })
|
||||
})
|
||||
}
|
||||
|
||||
// Extracted so the open wrapper owns the handle's lifetime.
|
||||
function readSession(args: {
|
||||
// Exported so a capture read can take the session and its whole transcript from
|
||||
// one open of the database rather than opening it twice.
|
||||
export function readOpenCodeSqliteSession(args: {
|
||||
db: SyncDatabase
|
||||
dbPath: string
|
||||
sessionId: string
|
||||
|
||||
@@ -33,9 +33,12 @@ export function jsonLines(records: unknown[]): string {
|
||||
return records.map((record) => JSON.stringify(record)).join('\n')
|
||||
}
|
||||
|
||||
// Newline-terminated, the way an agent writes each record: a file whose last
|
||||
// line has no break is a transcript mid-write, and the reader deliberately
|
||||
// withholds that line from consumers until it is complete.
|
||||
export async function writeJsonlFile(filePath: string, records: unknown[]): Promise<void> {
|
||||
await mkdir(dirname(filePath), { recursive: true })
|
||||
await writeFile(filePath, jsonLines(records))
|
||||
await writeFile(filePath, `${jsonLines(records)}\n`)
|
||||
}
|
||||
|
||||
export async function writeAntigravityTranscript(
|
||||
|
||||
@@ -0,0 +1,11 @@
|
||||
import { jsonLines, type isolatedScanRoots } from './session-scanner-test-fixtures'
|
||||
|
||||
// Shared by the two halves of the every-agent vault, which are split only
|
||||
// because one file of every agent's layout is past the line ceiling.
|
||||
|
||||
export type AgentVaultRoots = ReturnType<typeof isolatedScanRoots>
|
||||
|
||||
/** Records as a file body: newline-terminated, the way an agent writes them. */
|
||||
export function jsonlBody(records: unknown[]): string {
|
||||
return `${jsonLines(records)}\n`
|
||||
}
|
||||
@@ -4,13 +4,8 @@ import { join } from 'node:path'
|
||||
import { afterEach, describe, expect, it, vi } from 'vitest'
|
||||
import { AI_VAULT_AGENTS } from '../../shared/ai-vault-types'
|
||||
import { scanAiVaultSessions } from './session-scanner'
|
||||
import {
|
||||
isolatedScanRoots,
|
||||
jsonLines,
|
||||
writeAntigravityScannerFixture,
|
||||
writeOmpScannerFixture,
|
||||
writePrimeAgentScannerFixture
|
||||
} from './session-scanner-test-fixtures'
|
||||
import { isolatedScanRoots, jsonLines } from './session-scanner-test-fixtures'
|
||||
import { writeEveryAgentVault } from './session-scanner-every-agent-fixture'
|
||||
|
||||
let tempRoots: string[] = []
|
||||
|
||||
@@ -373,348 +368,8 @@ describe('scanAiVaultSessions', () => {
|
||||
it('indexes every supported agent transcript format with native resume commands', async () => {
|
||||
const root = await mkdtemp(join(tmpdir(), 'orca-ai-vault-all-agents-'))
|
||||
tempRoots.push(root)
|
||||
const roots = isolatedScanRoots(root)
|
||||
|
||||
await mkdir(join(roots.claudeProjectsDir, 'project'), { recursive: true })
|
||||
await writeFile(
|
||||
join(roots.claudeProjectsDir, 'project', 'claude-session.jsonl'),
|
||||
jsonLines([
|
||||
{
|
||||
type: 'user',
|
||||
sessionId: 'claude-session',
|
||||
timestamp: '2026-05-01T10:00:00.000Z',
|
||||
cwd: '/tmp/claude',
|
||||
message: { role: 'user', content: 'Claude title' }
|
||||
}
|
||||
])
|
||||
)
|
||||
|
||||
await mkdir(join(roots.codexSessionsDir, '2026', '05', '01'), { recursive: true })
|
||||
await writeFile(
|
||||
join(roots.codexSessionsDir, '2026', '05', '01', 'rollout-2026-codex-session.jsonl'),
|
||||
jsonLines([
|
||||
{
|
||||
timestamp: '2026-05-01T10:01:00.000Z',
|
||||
type: 'session_meta',
|
||||
payload: { id: 'codex-session', cwd: '/tmp/codex' }
|
||||
},
|
||||
{
|
||||
timestamp: '2026-05-01T10:01:01.000Z',
|
||||
type: 'response_item',
|
||||
payload: {
|
||||
type: 'message',
|
||||
role: 'user',
|
||||
content: [{ type: 'text', text: 'Codex title' }]
|
||||
}
|
||||
}
|
||||
])
|
||||
)
|
||||
|
||||
await mkdir(roots.geminiSessionsDir, { recursive: true })
|
||||
await writeFile(
|
||||
join(roots.geminiSessionsDir, 'gemini-session.json'),
|
||||
JSON.stringify({
|
||||
sessionId: 'gemini-session',
|
||||
startTime: '2026-05-01T10:02:00.000Z',
|
||||
lastUpdated: '2026-05-01T10:02:01.000Z',
|
||||
messages: [
|
||||
{
|
||||
type: 'user',
|
||||
timestamp: '2026-05-01T10:02:00.000Z',
|
||||
content: [{ text: 'Gemini title' }]
|
||||
},
|
||||
{
|
||||
type: 'gemini',
|
||||
timestamp: '2026-05-01T10:02:01.000Z',
|
||||
model: 'gemini-2.5-pro',
|
||||
tokens: { input: 10, output: 5 }
|
||||
}
|
||||
]
|
||||
})
|
||||
)
|
||||
|
||||
const antigravitySessionId = 'aaaaaaaa-bbbb-4ccc-8ddd-eeeeeeeeeeee'
|
||||
await writeAntigravityScannerFixture(roots.antigravityBrainDir, antigravitySessionId)
|
||||
|
||||
await mkdir(roots.copilotSessionsDir, { recursive: true })
|
||||
await writeFile(
|
||||
join(roots.copilotSessionsDir, 'copilot-session.jsonl'),
|
||||
jsonLines([
|
||||
{
|
||||
type: 'session.start',
|
||||
data: { sessionId: 'copilot-session', startTime: '2026-05-01T10:03:00.000Z' },
|
||||
timestamp: '2026-05-01T10:03:00.000Z'
|
||||
},
|
||||
{
|
||||
type: 'session.info',
|
||||
data: {
|
||||
infoType: 'folder_trust',
|
||||
message: 'Folder /tmp/copilot has been added to trusted folders.'
|
||||
},
|
||||
timestamp: '2026-05-01T10:03:01.000Z'
|
||||
},
|
||||
{
|
||||
type: 'user.message',
|
||||
data: { transformedContent: 'Copilot title' },
|
||||
timestamp: '2026-05-01T10:03:02.000Z'
|
||||
}
|
||||
])
|
||||
)
|
||||
|
||||
await mkdir(join(roots.cursorProjectsDir, 'project', 'agent-transcripts'), { recursive: true })
|
||||
await writeFile(
|
||||
join(roots.cursorProjectsDir, 'project', 'agent-transcripts', 'cursor-session.jsonl'),
|
||||
jsonLines([
|
||||
{
|
||||
role: 'user',
|
||||
message: { content: [{ type: 'text', text: 'Cursor title' }] }
|
||||
},
|
||||
{ role: 'assistant', message: { content: [{ type: 'text', text: 'Done' }] } }
|
||||
])
|
||||
)
|
||||
|
||||
await mkdir(join(roots.opencodeStorageDir, 'session', 'project'), { recursive: true })
|
||||
await mkdir(join(roots.opencodeStorageDir, 'message', 'opencode-session'), { recursive: true })
|
||||
await writeFile(
|
||||
join(roots.opencodeStorageDir, 'session', 'project', 'ses_opencode.json'),
|
||||
JSON.stringify({
|
||||
id: 'opencode-session',
|
||||
directory: '/tmp/opencode',
|
||||
title: 'OpenCode title',
|
||||
time: { created: 1_777_634_000_000, updated: 1_777_634_001_000 }
|
||||
})
|
||||
)
|
||||
await writeFile(
|
||||
join(roots.opencodeStorageDir, 'message', 'opencode-session', 'msg_1.json'),
|
||||
JSON.stringify({
|
||||
role: 'user',
|
||||
summary: { title: 'OpenCode title' },
|
||||
time: { created: 1_777_634_000_000 },
|
||||
tokens: { input: 7, output: 3 }
|
||||
})
|
||||
)
|
||||
|
||||
await mkdir(join(roots.grokSessionsDir, encodeURIComponent('/tmp/grok'), 'grok-session'), {
|
||||
recursive: true
|
||||
})
|
||||
await writeFile(
|
||||
join(roots.grokSessionsDir, encodeURIComponent('/tmp/grok'), 'grok-session', 'summary.json'),
|
||||
JSON.stringify({
|
||||
info: { id: 'grok-session', cwd: '/tmp/grok' },
|
||||
session_summary: '',
|
||||
created_at: '2026-05-01T10:04:00.000Z',
|
||||
updated_at: '2026-05-01T10:04:01.000Z',
|
||||
num_chat_messages: 2,
|
||||
current_model_id: 'grok-build',
|
||||
head_branch: 'feature/grok-vault'
|
||||
})
|
||||
)
|
||||
await writeFile(
|
||||
join(
|
||||
roots.grokSessionsDir,
|
||||
encodeURIComponent('/tmp/grok'),
|
||||
'grok-session',
|
||||
'chat_history.jsonl'
|
||||
),
|
||||
jsonLines([
|
||||
{
|
||||
type: 'user',
|
||||
content: [
|
||||
{
|
||||
type: 'text',
|
||||
text: '<user_info>context</user_info><user_query>Grok title</user_query>'
|
||||
}
|
||||
]
|
||||
},
|
||||
{ type: 'assistant', content: 'Done' }
|
||||
])
|
||||
)
|
||||
|
||||
await mkdir(roots.hermesSessionsDir, { recursive: true })
|
||||
await writeFile(
|
||||
join(roots.hermesSessionsDir, 'session_hermes-session.json'),
|
||||
JSON.stringify({
|
||||
session_id: 'hermes-session',
|
||||
model: 'hermes-1',
|
||||
cwd: '/tmp/hermes',
|
||||
session_start: '2026-05-01T10:05:00.000Z',
|
||||
last_updated: '2026-05-01T10:05:01.000Z',
|
||||
messages: [{ role: 'user', content: 'Hermes title' }]
|
||||
})
|
||||
)
|
||||
|
||||
await mkdir(join(roots.rovoSessionsDir, 'rovo-session'), { recursive: true })
|
||||
await writeFile(
|
||||
join(roots.rovoSessionsDir, 'rovo-session', 'metadata.json'),
|
||||
JSON.stringify({ title: 'Rovo title', workspace_path: '/tmp/rovo' })
|
||||
)
|
||||
await writeFile(
|
||||
join(roots.rovoSessionsDir, 'rovo-session', 'session_context.json'),
|
||||
JSON.stringify({
|
||||
message_history: [
|
||||
{
|
||||
kind: 'request',
|
||||
timestamp: '2026-05-01T10:06:00.000Z',
|
||||
parts: [{ part_kind: 'user-prompt', content: 'Rovo title' }]
|
||||
}
|
||||
]
|
||||
})
|
||||
)
|
||||
|
||||
await mkdir(join(roots.openclawStateDir, 'agents', 'default', 'sessions'), { recursive: true })
|
||||
await writeFile(
|
||||
join(roots.openclawStateDir, 'agents', 'default', 'sessions', 'openclaw-session.jsonl'),
|
||||
jsonLines([
|
||||
{
|
||||
type: 'session',
|
||||
id: 'openclaw-session',
|
||||
timestamp: '2026-05-01T10:07:00.000Z',
|
||||
cwd: '/tmp/openclaw'
|
||||
},
|
||||
{
|
||||
type: 'message',
|
||||
timestamp: '2026-05-01T10:07:01.000Z',
|
||||
message: { role: 'user', content: [{ type: 'text', text: 'OpenClaw title' }] }
|
||||
}
|
||||
])
|
||||
)
|
||||
|
||||
await mkdir(roots.piSessionsDir, { recursive: true })
|
||||
await writeFile(
|
||||
join(roots.piSessionsDir, 'pi-session.jsonl'),
|
||||
jsonLines([
|
||||
{
|
||||
type: 'session',
|
||||
id: 'pi-session',
|
||||
timestamp: '2026-05-01T10:08:00.000Z',
|
||||
cwd: '/tmp/pi'
|
||||
},
|
||||
{
|
||||
type: 'message',
|
||||
timestamp: '2026-05-01T10:08:01.000Z',
|
||||
message: { role: 'user', content: [{ type: 'text', text: 'Pi title' }] }
|
||||
}
|
||||
])
|
||||
)
|
||||
|
||||
const ompSessionFile = await writeOmpScannerFixture(roots.ompSessionsDir)
|
||||
const primeAgentSessionFile = await writePrimeAgentScannerFixture(roots.primeAgentSessionsDir)
|
||||
|
||||
await mkdir(roots.devinTranscriptsDir, { recursive: true })
|
||||
await writeFile(
|
||||
join(roots.devinTranscriptsDir, 'devin-session.json'),
|
||||
JSON.stringify({
|
||||
session_id: 'devin-session',
|
||||
working_directory: '/tmp/devin',
|
||||
agent: { model_name: 'swe-1-6-fast' },
|
||||
steps: [
|
||||
{
|
||||
metadata: {
|
||||
created_at: '2026-05-01T10:10:00.000Z',
|
||||
is_user_input: true,
|
||||
metrics: { input_tokens: 1, output_tokens: 2 }
|
||||
},
|
||||
text: 'Devin vault title'
|
||||
}
|
||||
]
|
||||
})
|
||||
)
|
||||
|
||||
await mkdir(roots.droidSessionsDir, { recursive: true })
|
||||
await writeFile(
|
||||
join(roots.droidSessionsDir, 'droid-session.jsonl'),
|
||||
jsonLines([
|
||||
{
|
||||
type: 'system',
|
||||
session_id: 'droid-session',
|
||||
timestamp: '2026-05-01T10:09:00.000Z',
|
||||
model: 'droid-model',
|
||||
cwd: '/tmp/droid'
|
||||
},
|
||||
{
|
||||
type: 'message',
|
||||
session_id: 'droid-session',
|
||||
timestamp: '2026-05-01T10:09:01.000Z',
|
||||
role: 'user',
|
||||
text: 'Droid title'
|
||||
},
|
||||
{
|
||||
type: 'completion',
|
||||
session_id: 'droid-session',
|
||||
timestamp: '2026-05-01T10:09:02.000Z',
|
||||
usage: { input_tokens: 2, output_tokens: 3 }
|
||||
}
|
||||
])
|
||||
)
|
||||
|
||||
const clineSessionId = 'cline-session'
|
||||
const clineSessionDir = join(roots.clineSessionsDir, clineSessionId)
|
||||
await mkdir(clineSessionDir, { recursive: true })
|
||||
await writeFile(
|
||||
join(clineSessionDir, `${clineSessionId}.json`),
|
||||
JSON.stringify({
|
||||
session_id: clineSessionId,
|
||||
started_at: '2026-05-01T10:10:30.000Z',
|
||||
model: 'cline-model',
|
||||
cwd: '/tmp/cline'
|
||||
})
|
||||
)
|
||||
await writeFile(
|
||||
join(clineSessionDir, `${clineSessionId}.messages.json`),
|
||||
JSON.stringify({
|
||||
updated_at: '2026-05-01T10:10:31.000Z',
|
||||
messages: [{ role: 'user', content: [{ type: 'text', text: 'Cline vault title' }] }]
|
||||
})
|
||||
)
|
||||
|
||||
// Kimi: <sessions>/wd_*/session_*/state.json + sibling agents/main/wire.jsonl,
|
||||
// with the work dir resolved from the top-level session_index.jsonl.
|
||||
const kimiSessionDir = join(roots.kimiSessionsDir, 'wd_app_abc', 'session_kimi-session')
|
||||
await mkdir(join(kimiSessionDir, 'agents', 'main'), { recursive: true })
|
||||
await writeFile(
|
||||
join(kimiSessionDir, 'state.json'),
|
||||
JSON.stringify({
|
||||
createdAt: '2026-05-01T10:11:00.000Z',
|
||||
updatedAt: '2026-05-01T10:11:05.000Z',
|
||||
title: 'Kimi vault title',
|
||||
lastPrompt: 'Kimi vault title',
|
||||
agents: { main: { type: 'main', parentAgentId: null } }
|
||||
})
|
||||
)
|
||||
await writeFile(
|
||||
join(root, 'session_index.jsonl'),
|
||||
jsonLines([
|
||||
{ sessionId: 'session_kimi-session', sessionDir: kimiSessionDir, workDir: '/tmp/kimi' }
|
||||
])
|
||||
)
|
||||
await writeFile(
|
||||
join(kimiSessionDir, 'agents', 'main', 'wire.jsonl'),
|
||||
jsonLines([
|
||||
{ type: 'config.update', modelAlias: 'kimi-k2.6', time: 1781853559132 },
|
||||
{
|
||||
type: 'context.append_message',
|
||||
message: {
|
||||
role: 'user',
|
||||
content: [{ type: 'text', text: 'Kimi vault title' }],
|
||||
origin: { kind: 'user' }
|
||||
},
|
||||
time: 1781853559164
|
||||
},
|
||||
{
|
||||
type: 'context.append_loop_event',
|
||||
event: { type: 'content.part', part: { type: 'text', text: 'Kimi reply' } },
|
||||
time: 1781853559177
|
||||
},
|
||||
{ type: 'context.append_loop_event', event: { type: 'step.end' }, time: 1781853559178 },
|
||||
{
|
||||
type: 'usage.record',
|
||||
model: 'kimi-k2.6',
|
||||
usage: { inputOther: 4, output: 6, inputCacheRead: 0, inputCacheCreation: 0 },
|
||||
usageScope: 'turn',
|
||||
time: 1781853559178
|
||||
}
|
||||
])
|
||||
)
|
||||
const { roots, antigravitySessionId, ompSessionFile, primeAgentSessionFile } =
|
||||
await writeEveryAgentVault(root)
|
||||
|
||||
const result = await scanAiVaultSessions({ ...roots, platform: 'darwin', limit: 20 })
|
||||
|
||||
|
||||
@@ -16,12 +16,18 @@ const OPENCODE_SQLITE_SESSION = {
|
||||
agent: 'opencode' as const,
|
||||
sessionId: 'sqlite-session'
|
||||
}
|
||||
const OPENCODE_SQLITE_MESSAGES = [
|
||||
{ role: 'user' as const, text: 'ask sqlite', timestamp: null },
|
||||
{ role: 'assistant' as const, text: 'reply sqlite', timestamp: null }
|
||||
]
|
||||
|
||||
// Stands in for the worker thread: the point is that its messages never come
|
||||
// back over the channel, not what the SQLite read returns.
|
||||
// Stands in for the worker thread: the point is which leg the reader asks for
|
||||
// and that what comes back reaches the channel, not what the SQLite read returns.
|
||||
vi.mock('./session-scanner-opencode-sqlite-worker-spawn', async (importOriginal) => ({
|
||||
...(await importOriginal<typeof OpenCodeSqliteWorkerSpawn>()),
|
||||
parseOpenCodeSqliteSessionViaWorker: () => Promise.resolve(OPENCODE_SQLITE_SESSION)
|
||||
parseOpenCodeSqliteSessionViaWorker: () => Promise.resolve(OPENCODE_SQLITE_SESSION),
|
||||
captureOpenCodeSqliteSessionViaWorker: () =>
|
||||
Promise.resolve({ session: OPENCODE_SQLITE_SESSION, messages: OPENCODE_SQLITE_MESSAGES })
|
||||
}))
|
||||
import type * as OpenCodeSqliteWorkerSpawn from './session-scanner-opencode-sqlite-worker-spawn'
|
||||
import {
|
||||
@@ -304,7 +310,7 @@ it('serializes overlapping parses of one path so no consumer read is orphaned',
|
||||
expect(second?.messageCount).toBe(10)
|
||||
})
|
||||
|
||||
it('reports a read whose parser cannot publish its messages as not complete', async () => {
|
||||
it('publishes an OpenCode SQLite session over the channel and reports it complete', async () => {
|
||||
const root = await mkdtemp(join(tmpdir(), 'orca-transcript-opencode-'))
|
||||
tempRoots.push(root)
|
||||
const dbPath = join(root, 'opencode.db')
|
||||
@@ -327,8 +333,9 @@ it('reports a read whose parser cannot publish its messages as not complete', as
|
||||
|
||||
expect(session).toEqual(OPENCODE_SQLITE_SESSION)
|
||||
expect(consumer.reads).toHaveLength(1)
|
||||
expect(consumer.reads[0].messages).toEqual([])
|
||||
expect(consumer.reads[0].outcome?.incomplete).toBe(true)
|
||||
expect(consumer.reads[0].messages).toEqual(OPENCODE_SQLITE_MESSAGES)
|
||||
expect(consumer.reads[0].outcome?.incomplete).toBe(false)
|
||||
consumer.unregister()
|
||||
})
|
||||
|
||||
it('reports the transcript size, not the cache key, as a whole-file read offset', async () => {
|
||||
|
||||
@@ -0,0 +1,130 @@
|
||||
import { mkdtemp, rm } from 'node:fs/promises'
|
||||
import { tmpdir } from 'node:os'
|
||||
import { join } from 'node:path'
|
||||
import { afterEach, expect, it, vi } from 'vitest'
|
||||
|
||||
// Only the thread hop is replaced: the implementations below are the repo's own
|
||||
// in-process readers, which the worker entry calls on the other side.
|
||||
vi.mock('./session-scanner-opencode-sqlite-worker-spawn', async () => {
|
||||
const list = await import('./session-scanner-opencode-sqlite-list')
|
||||
const parse = await import('./session-scanner-opencode-sqlite')
|
||||
const capture = await import('./session-scanner-opencode-sqlite-capture')
|
||||
return {
|
||||
resolveOpenCodeSqliteWorkerEntryPath: () => null,
|
||||
listOpenCodeSqliteSessionsViaWorker: (
|
||||
args: Parameters<typeof list.listOpenCodeSqliteSessions>[0]
|
||||
) => list.listOpenCodeSqliteSessions(args),
|
||||
parseOpenCodeSqliteSessionViaWorker: (
|
||||
args: Parameters<typeof parse.parseOpenCodeSqliteSession>[0]
|
||||
) => parse.parseOpenCodeSqliteSession(args),
|
||||
captureOpenCodeSqliteSessionViaWorker: (
|
||||
args: Parameters<typeof capture.captureOpenCodeSqliteSession>[0]
|
||||
) => capture.captureOpenCodeSqliteSession(args)
|
||||
}
|
||||
})
|
||||
import { AI_VAULT_AGENTS, type AiVaultAgent } from '../../shared/ai-vault-types'
|
||||
import { scanAiVaultSessions } from './session-scanner'
|
||||
import { writeEveryAgentVault } from './session-scanner-every-agent-fixture'
|
||||
import { resetSessionParseCacheForTests } from './session-scanner-parse-cache'
|
||||
import { writeOpenCodeSqliteDatabase } from './session-scanner-opencode-sqlite-fixture'
|
||||
import { splitOpenCodeSqliteCandidate } from './session-scanner-opencode-sqlite-paths'
|
||||
import {
|
||||
registerTranscriptConsumer,
|
||||
resetTranscriptConsumersForTests,
|
||||
type TranscriptMessage
|
||||
} from './session-transcript-consumers'
|
||||
|
||||
/*
|
||||
* The guard the OpenCode capture gap needed.
|
||||
*
|
||||
* Every consumer of the transcript reader -- the search index today, a digest
|
||||
* tomorrow -- sees an agent only through the messages its parser publishes. A
|
||||
* parser can list a session, show a preview and resume it correctly while
|
||||
* publishing nothing at all, which is exactly how 606 OpenCode sessions came to
|
||||
* hold zero indexed messages. Nothing above this layer can tell the difference,
|
||||
* so the assertion has to live here: one fixture per supported agent, read the
|
||||
* way the app reads it, and every agent has to say something.
|
||||
*/
|
||||
|
||||
const OPENCODE_SQLITE_SESSION = 'ses_capture_guard'
|
||||
|
||||
let tempRoots: string[] = []
|
||||
|
||||
afterEach(async () => {
|
||||
resetTranscriptConsumersForTests()
|
||||
resetSessionParseCacheForTests()
|
||||
await Promise.all(tempRoots.map((root) => rm(root, { recursive: true, force: true })))
|
||||
tempRoots = []
|
||||
})
|
||||
|
||||
type CapturedRead = { agent: AiVaultAgent; path: string; messages: TranscriptMessage[] }
|
||||
|
||||
async function readEveryAgentVault(): Promise<CapturedRead[]> {
|
||||
const root = await mkdtemp(join(tmpdir(), 'orca-transcript-every-agent-'))
|
||||
tempRoots.push(root)
|
||||
const { roots } = await writeEveryAgentVault(root)
|
||||
const dbPath = join(root, 'opencode-db', 'opencode.db')
|
||||
writeOpenCodeSqliteDatabase(dbPath, [
|
||||
{
|
||||
id: OPENCODE_SQLITE_SESSION,
|
||||
turns: [
|
||||
{ role: 'user', parts: ['what does the sqlite reader publish'] },
|
||||
{ role: 'assistant', parts: ['Every part of every turn.'] }
|
||||
]
|
||||
}
|
||||
])
|
||||
|
||||
const reads: CapturedRead[] = []
|
||||
registerTranscriptConsumer({
|
||||
beginRead: (start) => {
|
||||
const read: CapturedRead = {
|
||||
agent: start.candidate.agent,
|
||||
path: start.candidate.file.path,
|
||||
messages: []
|
||||
}
|
||||
reads.push(read)
|
||||
return { message: (message) => read.messages.push(message), finish: () => undefined }
|
||||
}
|
||||
})
|
||||
const result = await scanAiVaultSessions({
|
||||
...roots,
|
||||
opencodeDbPaths: [dbPath],
|
||||
platform: 'darwin',
|
||||
limit: 40
|
||||
})
|
||||
expect(result.issues).toEqual([])
|
||||
return reads
|
||||
}
|
||||
|
||||
function spokeIn(read: CapturedRead): boolean {
|
||||
return read.messages.some((message) => message.role === 'user' || message.role === 'assistant')
|
||||
}
|
||||
|
||||
it('publishes at least one user or assistant message for every supported agent', async () => {
|
||||
const reads = await readEveryAgentVault()
|
||||
|
||||
const silent = AI_VAULT_AGENTS.filter(
|
||||
(agent) => !reads.some((read) => read.agent === agent && spokeIn(read))
|
||||
)
|
||||
expect(silent).toEqual([])
|
||||
})
|
||||
|
||||
it('publishes an OpenCode SQLite session through the same channel as every file source', async () => {
|
||||
const reads = await readEveryAgentVault()
|
||||
|
||||
const sqliteRead = reads.find(
|
||||
(read) => splitOpenCodeSqliteCandidate(read.path)?.sessionId === OPENCODE_SQLITE_SESSION
|
||||
)
|
||||
expect(sqliteRead?.messages).toEqual([
|
||||
{
|
||||
role: 'user',
|
||||
text: 'what does the sqlite reader publish',
|
||||
timestamp: expect.any(String)
|
||||
},
|
||||
{
|
||||
role: 'assistant',
|
||||
text: 'Every part of every turn.',
|
||||
timestamp: expect.any(String)
|
||||
}
|
||||
])
|
||||
})
|
||||
@@ -1,6 +1,6 @@
|
||||
import { readTranscriptSlice } from '../native-chat/wsl-transcript-fs-access'
|
||||
import type { AiVaultSession } from '../../shared/ai-vault-types'
|
||||
import { parseAgentSessionFile, parserPublishesMessages } from './session-scanner-agent-parser'
|
||||
import { parseAgentSessionFile } from './session-scanner-agent-parser'
|
||||
import { consumeCompleteJsonlLines } from './session-scanner-jsonl-reader'
|
||||
import type { ResumableSessionParseState, SessionFileCandidate } from './session-scanner-types'
|
||||
import {
|
||||
@@ -159,12 +159,11 @@ export async function readWholeTranscript(args: {
|
||||
args.stats.fullParses++
|
||||
args.stats.bytesRead += file.sizeBytes ?? 0
|
||||
}
|
||||
const publishes = parserPublishesMessages(args.candidate)
|
||||
const channel = new TranscriptMessageChannel()
|
||||
channel.beginRead({ candidate: args.candidate, mode: 'replace', previousByteOffset: 0 })
|
||||
try {
|
||||
const session = await parseAgentSessionFile(args.candidate, args.platform, channel)
|
||||
channel.finishRead({ session, byteOffset: file.sizeBytes ?? 0, incomplete: !publishes })
|
||||
channel.finishRead({ session, byteOffset: file.sizeBytes ?? 0, incomplete: false })
|
||||
return session
|
||||
} catch (error) {
|
||||
channel.finishRead({ session: null, byteOffset: 0, incomplete: true })
|
||||
|
||||
@@ -104,5 +104,8 @@ export const AiVaultSearchStatusSchema = z.object({
|
||||
degradedRoots: z.array(z.object({ root: z.string().optional(), reason: z.string() })),
|
||||
lastReconcileAt: z.number().nullable(),
|
||||
lastSweepCompletedAt: z.number().nullable(),
|
||||
// Optional: an older host answers without it, and a reader that has none
|
||||
// should show no breakdown rather than a breakdown of zeroes.
|
||||
sessionsByAgent: z.record(z.string(), z.number().int().nonnegative()).optional(),
|
||||
generation: z.number().int().nonnegative()
|
||||
})
|
||||
|
||||
Reference in New Issue
Block a user