Pi cannot be read as another linear JSONL stream. Its history is a tree with a durable leaf, orphan roots, branch summaries, and two compaction forms, so the active context is something the format states rather than something line order implies. The adapter keeps those semantics inside itself and projects the result into the existing canonical tables. Sessions are keyed by (normalized header cwd, header id) rather than by path, because Pi's --session-id lookup is project-local: two projects may reuse an id, while a move or an identical copy is still one session. Discovery covers both layouts Pi writes and fingerprints each file by mtime, ctime, size and inode, so a rewrite that preserves mtime is not read as unchanged. Abandoned branches are preserved rather than dropped. Visibility becomes three-state -- visible, inactive, hidden -- and helpers return only visible rows until includeInactive asks for the superseded path, labeling every row so a caller knows which it holds. Usage counts all three, because an abandoned call still spent tokens; message_count reports only the visible transcript. A committed MIT-licensed oracle transcribed from Pi 0.83.0 pins the context algorithms, and a fixed-seed differential runs 512 generated sessions against it on every test run. Schema changes are additive.
478 lines
13 KiB
JavaScript
478 lines
13 KiB
JavaScript
import { test } from 'node:test';
|
|
import assert from 'node:assert/strict';
|
|
import { createRequire } from 'node:module';
|
|
import { mkdirSync, mkdtempSync, writeFileSync } from 'node:fs';
|
|
import { tmpdir } from 'node:os';
|
|
import { join } from 'node:path';
|
|
|
|
const require = createRequire(import.meta.url);
|
|
import { createIndexerService } from '../app/src/main/indexer-service.ts';
|
|
|
|
function manualTimers() {
|
|
const timers = new Set();
|
|
return {
|
|
setTimeout(fn) {
|
|
timers.add(fn);
|
|
return fn;
|
|
},
|
|
clearTimeout(fn) {
|
|
timers.delete(fn);
|
|
},
|
|
flush() {
|
|
const pending = [...timers];
|
|
timers.clear();
|
|
for (const fn of pending) fn();
|
|
},
|
|
};
|
|
}
|
|
|
|
test('indexer service debounces repeated build requests', async () => {
|
|
const timers = manualTimers();
|
|
const calls = [];
|
|
const service = createIndexerService({
|
|
buildIndex: async ({ reason }) => calls.push(reason),
|
|
watchProjects: () => null,
|
|
writeHeartbeat: () => {},
|
|
timers,
|
|
stabilityMs: 0,
|
|
});
|
|
|
|
service.scheduleBuild('first');
|
|
service.scheduleBuild('second');
|
|
service.scheduleBuild('third');
|
|
timers.flush();
|
|
await service.idle();
|
|
|
|
assert.deepEqual(calls, ['third']);
|
|
});
|
|
|
|
test('indexer service runs one pending build after an in-flight build finishes', async () => {
|
|
const timers = manualTimers();
|
|
const calls = [];
|
|
let finishFirst;
|
|
const service = createIndexerService({
|
|
buildIndex: async ({ reason }) => {
|
|
calls.push(reason);
|
|
if (reason === 'first') await new Promise(resolve => { finishFirst = resolve; });
|
|
},
|
|
watchProjects: () => null,
|
|
writeHeartbeat: () => {},
|
|
timers,
|
|
stabilityMs: 0,
|
|
});
|
|
|
|
const first = service.runBuildNow('first');
|
|
service.scheduleBuild('second');
|
|
timers.flush();
|
|
|
|
assert.deepEqual(calls, ['first']);
|
|
finishFirst();
|
|
await first;
|
|
await service.idle();
|
|
|
|
assert.deepEqual(calls, ['first', 'pending']);
|
|
});
|
|
|
|
test('indexer service reschedules a writer-lease deferral without publishing a heartbeat', async () => {
|
|
const timers = manualTimers();
|
|
const calls = [];
|
|
let heartbeats = 0;
|
|
const service = createIndexerService({
|
|
buildIndex: async ({ reason, changedPaths }) => {
|
|
calls.push({ reason, changedPaths });
|
|
return calls.length === 1 ? { deferred: true, reason: 'writer_busy' } : { deferred: false };
|
|
},
|
|
watchProjects: () => null,
|
|
writeHeartbeat: () => { heartbeats += 1; },
|
|
timers,
|
|
stabilityMs: 0,
|
|
});
|
|
|
|
await service.runBuildNow('watch', ['project/session.jsonl']);
|
|
assert.equal(heartbeats, 0);
|
|
assert.equal(calls.length, 1);
|
|
|
|
timers.flush();
|
|
await service.idle();
|
|
assert.deepEqual(calls, [
|
|
{ reason: 'watch', changedPaths: ['project/session.jsonl'] },
|
|
{ reason: 'writer-lease', changedPaths: ['project/session.jsonl'] },
|
|
]);
|
|
assert.equal(heartbeats, 1);
|
|
});
|
|
|
|
test('a deferred full-inventory build stays full when retried', async () => {
|
|
const timers = manualTimers();
|
|
const calls = [];
|
|
const service = createIndexerService({
|
|
buildIndex: async (args) => {
|
|
calls.push(args);
|
|
return calls.length === 1 ? { deferred: true } : { deferred: false };
|
|
},
|
|
watchProjects: () => null,
|
|
writeHeartbeat: () => {},
|
|
timers,
|
|
stabilityMs: 0,
|
|
});
|
|
|
|
service.scheduleBuild('root-appeared');
|
|
timers.flush();
|
|
await service.idle();
|
|
service.scheduleBuild('ordinary-change', '/tmp/later.jsonl');
|
|
timers.flush();
|
|
await service.idle();
|
|
|
|
assert.deepEqual(calls, [
|
|
{ reason: 'root-appeared', changedPaths: undefined },
|
|
{ reason: 'ordinary-change', changedPaths: undefined },
|
|
]);
|
|
});
|
|
|
|
test('indexer service does not log a build cancelled by a service stop', async () => {
|
|
const timers = manualTimers();
|
|
const warnings = [];
|
|
let rejectBuild;
|
|
const service = createIndexerService({
|
|
buildIndex: () => new Promise((_resolve, reject) => { rejectBuild = reject; }),
|
|
watchProjects: () => null,
|
|
writeHeartbeat: () => {},
|
|
timers,
|
|
stabilityMs: 0,
|
|
logger: { warn: (msg) => warnings.push(msg) },
|
|
});
|
|
|
|
const build = service.runBuildNow('startup');
|
|
service.stop(); // manual rebuild path tears the worker down mid-build
|
|
rejectBuild(new Error('Indexer worker stopped'));
|
|
await build;
|
|
|
|
assert.deepEqual(warnings, []);
|
|
});
|
|
|
|
test('indexer service logs a build that fails while running', async () => {
|
|
const timers = manualTimers();
|
|
const warnings = [];
|
|
let rejectBuild;
|
|
const service = createIndexerService({
|
|
buildIndex: () => new Promise((_resolve, reject) => { rejectBuild = reject; }),
|
|
watchProjects: () => null,
|
|
writeHeartbeat: () => {},
|
|
timers,
|
|
stabilityMs: 0,
|
|
logger: { warn: (msg) => warnings.push(msg) },
|
|
});
|
|
|
|
const build = service.runBuildNow('watch');
|
|
rejectBuild(new Error('disk on fire'));
|
|
await build;
|
|
|
|
assert.equal(warnings.length, 1);
|
|
assert.match(warnings[0], /Obelisk index build failed: disk on fire/);
|
|
});
|
|
|
|
test('indexer service reports partial inventory paths on ordinary builds', async () => {
|
|
const warnings = [];
|
|
const service = createIndexerService({
|
|
buildIndex: async () => ({
|
|
deferred: false,
|
|
complete: false,
|
|
inventoryIssues: [{
|
|
provider: 'pi',
|
|
path: '/tmp/pi/locked',
|
|
error: 'EACCES: permission denied',
|
|
}],
|
|
}),
|
|
watchProjects: () => null,
|
|
writeHeartbeat: () => {},
|
|
logger: { warn: (msg) => warnings.push(msg) },
|
|
stabilityMs: 0,
|
|
});
|
|
|
|
await service.runBuildNow('startup');
|
|
|
|
assert.deepEqual(warnings, [
|
|
'Obelisk indexed a partial pi inventory at /tmp/pi/locked: EACCES: permission denied',
|
|
]);
|
|
});
|
|
|
|
test('indexer service reports partial inventory paths before a deferred retry', async () => {
|
|
const timers = manualTimers();
|
|
const warnings = [];
|
|
const service = createIndexerService({
|
|
buildIndex: async () => ({
|
|
deferred: true,
|
|
complete: false,
|
|
inventoryIssues: [{
|
|
provider: 'pi',
|
|
path: '/tmp/pi/locked',
|
|
error: 'EACCES: permission denied',
|
|
}],
|
|
}),
|
|
watchProjects: () => null,
|
|
writeHeartbeat: () => {},
|
|
logger: { warn: (msg) => warnings.push(msg) },
|
|
timers,
|
|
stabilityMs: 0,
|
|
});
|
|
|
|
await service.runBuildNow('startup');
|
|
service.stop();
|
|
|
|
assert.deepEqual(warnings, [
|
|
'Obelisk indexed a partial pi inventory at /tmp/pi/locked: EACCES: permission denied',
|
|
]);
|
|
});
|
|
|
|
test('indexer service waits for a stability window before building', async () => {
|
|
const timers = manualTimers();
|
|
const calls = [];
|
|
const service = createIndexerService({
|
|
buildIndex: async ({ reason }) => calls.push(reason),
|
|
watchProjects: () => null,
|
|
writeHeartbeat: () => {},
|
|
timers,
|
|
stabilityMs: 500,
|
|
});
|
|
|
|
service.scheduleBuild('jsonl-change');
|
|
timers.flush();
|
|
await service.idle();
|
|
assert.deepEqual(calls, []);
|
|
|
|
timers.flush();
|
|
await service.idle();
|
|
assert.deepEqual(calls, ['jsonl-change']);
|
|
});
|
|
|
|
test('indexer service retries watcher setup when the projects directory is missing', () => {
|
|
const timers = manualTimers();
|
|
let attempts = 0;
|
|
const service = createIndexerService({
|
|
buildIndex: async () => {},
|
|
watchProjects: () => {
|
|
attempts++;
|
|
return attempts === 1 ? null : { close() {} };
|
|
},
|
|
writeHeartbeat: () => {},
|
|
timers,
|
|
stabilityMs: 0,
|
|
});
|
|
|
|
service.start({ buildOnStart: false });
|
|
assert.equal(attempts, 1);
|
|
|
|
timers.flush();
|
|
assert.equal(attempts, 2);
|
|
|
|
timers.flush();
|
|
assert.equal(attempts, 2);
|
|
});
|
|
|
|
test('indexer service publishes daemon ownership as soon as it starts', () => {
|
|
const timers = manualTimers();
|
|
let heartbeats = 0;
|
|
const service = createIndexerService({
|
|
buildIndex: async () => ({ deferred: false }),
|
|
watchProjects: () => null,
|
|
writeHeartbeat: () => { heartbeats += 1; },
|
|
timers,
|
|
stabilityMs: 0,
|
|
});
|
|
|
|
service.start({ buildOnStart: false });
|
|
assert.equal(heartbeats, 1);
|
|
service.stop();
|
|
});
|
|
|
|
test('indexer service watches Claude JSON files through chokidar', async () => {
|
|
const projectsDir = mkdtempSync(join(tmpdir(), 'obelisk-chokidar-projects-'));
|
|
const timers = manualTimers();
|
|
const calls = [];
|
|
let watchArgs = null;
|
|
const handlers = {};
|
|
const watcher = {
|
|
on(event, handler) {
|
|
handlers[event] = handler;
|
|
return watcher;
|
|
},
|
|
closeCalled: false,
|
|
close() {
|
|
watcher.closeCalled = true;
|
|
},
|
|
};
|
|
const chokidar = {
|
|
watch(paths, options) {
|
|
watchArgs = { paths, options };
|
|
return watcher;
|
|
},
|
|
};
|
|
|
|
const service = createIndexerService({
|
|
projectsDir,
|
|
buildIndex: async ({ reason }) => calls.push(reason),
|
|
chokidar,
|
|
writeHeartbeat: () => {},
|
|
timers,
|
|
stabilityMs: 0,
|
|
debounceMs: 0,
|
|
});
|
|
|
|
try {
|
|
service.start({ buildOnStart: false });
|
|
assert.equal(watchArgs.paths, projectsDir);
|
|
assert.equal(watchArgs.options.cwd, projectsDir);
|
|
assert.equal(watchArgs.options.ignoreInitial, true);
|
|
assert.ok(watchArgs.options.awaitWriteFinish);
|
|
|
|
handlers.change('session.jsonl');
|
|
timers.flush();
|
|
await service.idle();
|
|
assert.deepEqual(calls, ['watch']);
|
|
} finally {
|
|
service.stop();
|
|
}
|
|
|
|
assert.equal(watcher.closeCalled, true);
|
|
});
|
|
|
|
test('indexer service passes changed JSONL paths to the build worker', async () => {
|
|
const projectsDir = mkdtempSync(join(tmpdir(), 'obelisk-changed-paths-'));
|
|
const timers = manualTimers();
|
|
const calls = [];
|
|
const handlers = {};
|
|
const watcher = {
|
|
on(event, handler) {
|
|
handlers[event] = handler;
|
|
return watcher;
|
|
},
|
|
close() {},
|
|
};
|
|
const chokidar = {
|
|
watch() {
|
|
return watcher;
|
|
},
|
|
};
|
|
|
|
const service = createIndexerService({
|
|
projectsDir,
|
|
buildIndex: async (args) => calls.push(args),
|
|
chokidar,
|
|
writeHeartbeat: () => {},
|
|
timers,
|
|
stabilityMs: 0,
|
|
debounceMs: 0,
|
|
});
|
|
|
|
service.start({ buildOnStart: false });
|
|
handlers.change('project-a/session-1.jsonl');
|
|
handlers.add('project-a/session-2.json');
|
|
timers.flush();
|
|
await service.idle();
|
|
|
|
assert.equal(calls.length, 1);
|
|
assert.equal(calls[0].reason, 'watch');
|
|
assert.deepEqual(calls[0].changedPaths, [
|
|
join(projectsDir, 'project-a/session-1.jsonl'),
|
|
join(projectsDir, 'project-a/session-2.json'),
|
|
]);
|
|
});
|
|
|
|
test('indexer service watches Claude projects and Codex sessions for app-side indexing', async () => {
|
|
const claudeProjectsDir = mkdtempSync(join(tmpdir(), 'obelisk-watch-claude-'));
|
|
const codexSessionsDir = mkdtempSync(join(tmpdir(), 'obelisk-watch-codex-sessions-'));
|
|
const timers = manualTimers();
|
|
const calls = [];
|
|
const watchers = [];
|
|
const watchArgs = [];
|
|
const chokidar = {
|
|
watch(paths, options) {
|
|
const handlers = {};
|
|
const watcher = {
|
|
handlers,
|
|
on(event, handler) {
|
|
handlers[event] = handler;
|
|
return watcher;
|
|
},
|
|
close() {},
|
|
};
|
|
watchers.push(watcher);
|
|
watchArgs.push({ paths, options });
|
|
return watcher;
|
|
},
|
|
};
|
|
|
|
const service = createIndexerService({
|
|
projectsDir: claudeProjectsDir,
|
|
watchDirs: [claudeProjectsDir, codexSessionsDir],
|
|
buildIndex: async (args) => calls.push(args),
|
|
chokidar,
|
|
writeHeartbeat: () => {},
|
|
timers,
|
|
stabilityMs: 0,
|
|
debounceMs: 0,
|
|
});
|
|
|
|
service.start({ buildOnStart: false });
|
|
assert.deepEqual(watchArgs.map(arg => arg.paths), [claudeProjectsDir, codexSessionsDir]);
|
|
assert.deepEqual(watchArgs.map(arg => arg.options.cwd), [claudeProjectsDir, codexSessionsDir]);
|
|
|
|
watchers[1].handlers.change('2026/06/15/rollout-2026-06-15T00-00-00-codex.jsonl');
|
|
timers.flush();
|
|
await service.idle();
|
|
|
|
assert.equal(calls.length, 1);
|
|
assert.deepEqual(calls[0].changedPaths, [
|
|
join(codexSessionsDir, '2026/06/15/rollout-2026-06-15T00-00-00-codex.jsonl'),
|
|
]);
|
|
});
|
|
|
|
test('indexer service starts watching a configured root that appears after startup', async () => {
|
|
const existingRoot = mkdtempSync(join(tmpdir(), 'obelisk-watch-existing-'));
|
|
const parent = mkdtempSync(join(tmpdir(), 'obelisk-watch-late-parent-'));
|
|
const lateRoot = join(parent, 'nested', 'sessions');
|
|
const timers = manualTimers();
|
|
const calls = [];
|
|
const watchArgs = [];
|
|
const chokidar = {
|
|
watch(root) {
|
|
watchArgs.push(root);
|
|
const watcher = {
|
|
on() {
|
|
return watcher;
|
|
},
|
|
close() {},
|
|
};
|
|
return watcher;
|
|
},
|
|
};
|
|
const service = createIndexerService({
|
|
watchDirs: [existingRoot, lateRoot],
|
|
buildIndex: async (args) => calls.push(args),
|
|
chokidar,
|
|
writeHeartbeat: () => {},
|
|
timers,
|
|
stabilityMs: 0,
|
|
debounceMs: 0,
|
|
watchRetryMs: 0,
|
|
});
|
|
|
|
try {
|
|
service.start({ buildOnStart: false });
|
|
assert.deepEqual(watchArgs, [existingRoot]);
|
|
|
|
mkdirSync(lateRoot, { recursive: true });
|
|
writeFileSync(join(lateRoot, 'pre-existing.jsonl'), '{}\n');
|
|
timers.flush();
|
|
assert.deepEqual(watchArgs, [existingRoot, lateRoot]);
|
|
|
|
timers.flush();
|
|
await service.idle();
|
|
assert.deepEqual(calls, [{
|
|
reason: 'watch',
|
|
changedPaths: undefined,
|
|
}]);
|
|
} finally {
|
|
service.stop();
|
|
}
|
|
});
|