mirror of
https://github.com/earendil-works/pi.git
synced 2026-09-28 05:54:43 +08:00
fix(agent): expose session stores through factories
This commit is contained in:
@@ -3,23 +3,24 @@ import type {
|
||||
JsonlSessionCreateOptions,
|
||||
JsonlSessionListOptions,
|
||||
JsonlSessionMetadata,
|
||||
JsonlSessionStoreApi,
|
||||
LeafEntry,
|
||||
SessionEntryCursorOptions,
|
||||
SessionSnapshot,
|
||||
SessionStorage,
|
||||
SessionStore,
|
||||
SessionTreeEntry,
|
||||
} from "../types.ts";
|
||||
import { SessionError, toError } from "../types.ts";
|
||||
import { JsonlSessionStorage, loadJsonlSessionMetadata } from "./jsonl-storage.ts";
|
||||
import {
|
||||
createSessionId,
|
||||
createSessionRepository,
|
||||
createTimestamp,
|
||||
getEntriesToFork,
|
||||
getFileSystemResultOrThrow,
|
||||
SessionRepository,
|
||||
type SessionRepository,
|
||||
} from "./repo-utils.ts";
|
||||
import { ScanningSessionSearch } from "./search-backend.ts";
|
||||
import { createScanningSessionSearch } from "./search-backend.ts";
|
||||
|
||||
export type JsonlSessionStoreOptions = { fs: JsonlSessionStoreFileSystem; sessionsRoot: string };
|
||||
|
||||
@@ -42,7 +43,9 @@ function encodeCwd(cwd: string): string {
|
||||
return `--${cwd.replace(/^[/\\]/, "").replace(/[/\\:]/g, "-")}--`;
|
||||
}
|
||||
|
||||
export class JsonlSessionStore implements JsonlSessionStoreApi {
|
||||
class JsonlSessionStore
|
||||
implements SessionStore<JsonlSessionMetadata, JsonlSessionCreateOptions, JsonlSessionListOptions>
|
||||
{
|
||||
private readonly fs: JsonlSessionStoreFileSystem;
|
||||
private readonly sessionsRootInput: string;
|
||||
private sessionsRoot: string | undefined;
|
||||
@@ -209,7 +212,9 @@ export class JsonlSessionStore implements JsonlSessionStoreApi {
|
||||
}
|
||||
}
|
||||
|
||||
export function createJsonlSessionStore(options: JsonlSessionStoreOptions): JsonlSessionStore {
|
||||
export function createJsonlSessionStore(
|
||||
options: JsonlSessionStoreOptions,
|
||||
): SessionStore<JsonlSessionMetadata, JsonlSessionCreateOptions, JsonlSessionListOptions> {
|
||||
return new JsonlSessionStore(options);
|
||||
}
|
||||
|
||||
@@ -217,5 +222,5 @@ export function createJsonlSessionRepository(
|
||||
options: JsonlSessionStoreOptions,
|
||||
): SessionRepository<JsonlSessionMetadata, JsonlSessionCreateOptions, JsonlSessionListOptions> {
|
||||
const store = createJsonlSessionStore(options);
|
||||
return new SessionRepository({ store, search: new ScanningSessionSearch(store) });
|
||||
return createSessionRepository({ store, search: createScanningSessionSearch(store) });
|
||||
}
|
||||
|
||||
@@ -9,12 +9,18 @@ import {
|
||||
type SessionTreeEntry,
|
||||
} from "../types.ts";
|
||||
import { InMemorySessionStorage } from "./memory-storage.ts";
|
||||
import { createSessionId, createTimestamp, getEntriesToFork, SessionRepository } from "./repo-utils.ts";
|
||||
import { ScanningSessionSearch } from "./search-backend.ts";
|
||||
import {
|
||||
createSessionId,
|
||||
createSessionRepository,
|
||||
createTimestamp,
|
||||
getEntriesToFork,
|
||||
type SessionRepository,
|
||||
} from "./repo-utils.ts";
|
||||
import { createScanningSessionSearch } from "./search-backend.ts";
|
||||
|
||||
export type InMemorySessionCreateOptions = { id?: string };
|
||||
|
||||
export class InMemorySessionStore implements SessionStore<SessionMetadata, InMemorySessionCreateOptions, void> {
|
||||
class InMemorySessionStore implements SessionStore<SessionMetadata, InMemorySessionCreateOptions, void> {
|
||||
private sessions = new Map<string, InMemorySessionStorage<SessionMetadata>>();
|
||||
|
||||
async create(options: InMemorySessionCreateOptions = {}): Promise<SessionMetadata> {
|
||||
@@ -82,7 +88,7 @@ export class InMemorySessionStore implements SessionStore<SessionMetadata, InMem
|
||||
}
|
||||
}
|
||||
|
||||
export function createInMemorySessionStore(): InMemorySessionStore {
|
||||
export function createInMemorySessionStore(): SessionStore<SessionMetadata, InMemorySessionCreateOptions, void> {
|
||||
return new InMemorySessionStore();
|
||||
}
|
||||
|
||||
@@ -92,5 +98,5 @@ export function createInMemorySessionRepository(): SessionRepository<
|
||||
void
|
||||
> {
|
||||
const store = createInMemorySessionStore();
|
||||
return new SessionRepository({ store, search: new ScanningSessionSearch(store) });
|
||||
return createSessionRepository({ store, search: createScanningSessionSearch(store) });
|
||||
}
|
||||
|
||||
@@ -13,9 +13,7 @@ type SessionSearchSource<TMetadata extends SessionMetadata> = {
|
||||
};
|
||||
|
||||
/** Searches canonical sessions directly and therefore has no index to maintain. */
|
||||
export class ScanningSessionSearch<TMetadata extends SessionMetadata = SessionMetadata>
|
||||
implements SessionSearch<TMetadata>
|
||||
{
|
||||
class ScanningSessionSearch<TMetadata extends SessionMetadata = SessionMetadata> implements SessionSearch<TMetadata> {
|
||||
private readonly source: SessionSearchSource<TMetadata>;
|
||||
|
||||
constructor(source: SessionSearchSource<TMetadata>) {
|
||||
@@ -33,3 +31,9 @@ export class ScanningSessionSearch<TMetadata extends SessionMetadata = SessionMe
|
||||
return hits;
|
||||
}
|
||||
}
|
||||
|
||||
export function createScanningSessionSearch<TMetadata extends SessionMetadata>(
|
||||
source: SessionSearchSource<TMetadata>,
|
||||
): SessionSearch<TMetadata> {
|
||||
return new ScanningSessionSearch(source);
|
||||
}
|
||||
|
||||
@@ -582,9 +582,6 @@ export interface JsonlSessionListOptions {
|
||||
cwd?: string;
|
||||
}
|
||||
|
||||
export interface JsonlSessionStoreApi
|
||||
extends SessionStore<JsonlSessionMetadata, JsonlSessionCreateOptions, JsonlSessionListOptions> {}
|
||||
|
||||
export type AgentHarnessPhase = "idle" | "turn" | "compaction" | "branch_summary" | "retry";
|
||||
|
||||
export type PendingSessionWrite = SessionTreeEntry extends infer TEntry
|
||||
|
||||
@@ -1,23 +1,43 @@
|
||||
import { existsSync } from "node:fs";
|
||||
import { describe, expect, it } from "vitest";
|
||||
import { NodeExecutionEnv } from "../../src/harness/env/nodejs.ts";
|
||||
import { JsonlSessionStore } from "../../src/harness/session/jsonl-repo.ts";
|
||||
import { InMemorySessionStore } from "../../src/harness/session/memory-repo.ts";
|
||||
import { SessionRepository } from "../../src/harness/session/repo-utils.ts";
|
||||
import { createJsonlSessionStore } from "../../src/harness/session/jsonl-repo.ts";
|
||||
import {
|
||||
createInMemorySessionStore,
|
||||
type InMemorySessionCreateOptions,
|
||||
} from "../../src/harness/session/memory-repo.ts";
|
||||
import { createSessionRepository } from "../../src/harness/session/repo-utils.ts";
|
||||
import type { SessionMetadata, SessionStore } from "../../src/harness/types.ts";
|
||||
import { createAssistantMessage, createTempDir, createUserMessage } from "./session-test-utils.ts";
|
||||
|
||||
class CountingInMemorySessionStore extends InMemorySessionStore {
|
||||
loadCount = 0;
|
||||
|
||||
override async load(...args: Parameters<InMemorySessionStore["load"]>) {
|
||||
this.loadCount += 1;
|
||||
return super.load(...args);
|
||||
}
|
||||
function createCountingInMemorySessionStore(): {
|
||||
store: SessionStore<SessionMetadata, InMemorySessionCreateOptions, void>;
|
||||
counter: { loadCount: number };
|
||||
} {
|
||||
const source = createInMemorySessionStore();
|
||||
const counter = { loadCount: 0 };
|
||||
return {
|
||||
counter,
|
||||
store: {
|
||||
create: (options) => source.create(options),
|
||||
async load(metadata) {
|
||||
counter.loadCount += 1;
|
||||
return source.load(metadata);
|
||||
},
|
||||
list: (options) => source.list(options),
|
||||
getEntries: (metadata, options) => source.getEntries(metadata, options),
|
||||
createEntryId: (metadata) => source.createEntryId(metadata),
|
||||
appendEntry: (metadata, entry) => source.appendEntry(metadata, entry),
|
||||
setLeafId: (metadata, leafId) => source.setLeafId(metadata, leafId),
|
||||
delete: (metadata) => source.delete(metadata),
|
||||
fork: (metadata, options) => source.fork(metadata, options),
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
describe("InMemorySessionStore", () => {
|
||||
it("opens, deletes, and forks by metadata", async () => {
|
||||
const repo = new SessionRepository({ store: new InMemorySessionStore() });
|
||||
const repo = createSessionRepository({ store: createInMemorySessionStore() });
|
||||
const session = await repo.create({ id: "session-1" });
|
||||
const metadata = await session.getMetadata();
|
||||
const user1 = await session.appendMessage(createUserMessage("one"));
|
||||
@@ -34,21 +54,21 @@ describe("InMemorySessionStore", () => {
|
||||
});
|
||||
|
||||
it("does not repeatedly load full snapshots for scoped reads", async () => {
|
||||
const store = new CountingInMemorySessionStore();
|
||||
const repo = new SessionRepository({ store });
|
||||
const { store, counter } = createCountingInMemorySessionStore();
|
||||
const repo = createSessionRepository({ store });
|
||||
const session = await repo.create({ id: "session-1" });
|
||||
const entryId = await session.appendMessage(createUserMessage("one"));
|
||||
|
||||
store.loadCount = 0;
|
||||
counter.loadCount = 0;
|
||||
await session.getMetadata();
|
||||
expect(store.loadCount).toBe(0);
|
||||
expect(counter.loadCount).toBe(0);
|
||||
|
||||
await session.getLeafId();
|
||||
expect(store.loadCount).toBe(1);
|
||||
expect(counter.loadCount).toBe(1);
|
||||
|
||||
store.loadCount = 0;
|
||||
counter.loadCount = 0;
|
||||
await session.getEntry(entryId);
|
||||
expect(store.loadCount).toBe(1);
|
||||
expect(counter.loadCount).toBe(1);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -58,7 +78,7 @@ describe("JsonlSessionStore", () => {
|
||||
const env = new NodeExecutionEnv({ cwd: root });
|
||||
const cwd = "/tmp/my-project";
|
||||
const otherCwd = "/tmp/other-project";
|
||||
const repo = new SessionRepository({ store: new JsonlSessionStore({ fs: env, sessionsRoot: root }) });
|
||||
const repo = createSessionRepository({ store: createJsonlSessionStore({ fs: env, sessionsRoot: root }) });
|
||||
const session = await repo.create({ cwd, id: "019de8c2-de29-73e9-ae0c-e134db34c447" });
|
||||
const otherSession = await repo.create({ cwd: otherCwd, id: "other-session" });
|
||||
const metadata = await session.getMetadata();
|
||||
@@ -75,7 +95,7 @@ describe("JsonlSessionStore", () => {
|
||||
it("opens, deletes, and forks by metadata", async () => {
|
||||
const root = createTempDir();
|
||||
const env = new NodeExecutionEnv({ cwd: root });
|
||||
const repo = new SessionRepository({ store: new JsonlSessionStore({ fs: env, sessionsRoot: root }) });
|
||||
const repo = createSessionRepository({ store: createJsonlSessionStore({ fs: env, sessionsRoot: root }) });
|
||||
const source = await repo.create({ cwd: "/tmp/source", id: "source-session" });
|
||||
const sourceMetadata = await source.getMetadata();
|
||||
const user1 = await source.appendMessage(createUserMessage("one"));
|
||||
@@ -97,7 +117,7 @@ describe("JsonlSessionStore", () => {
|
||||
it("persists header metadata through create, list, and fork", async () => {
|
||||
const root = createTempDir();
|
||||
const env = new NodeExecutionEnv({ cwd: root });
|
||||
const repo = new SessionRepository({ store: new JsonlSessionStore({ fs: env, sessionsRoot: root }) });
|
||||
const repo = createSessionRepository({ store: createJsonlSessionStore({ fs: env, sessionsRoot: root }) });
|
||||
const source = await repo.create({
|
||||
cwd: "/tmp/source",
|
||||
id: "source-session",
|
||||
|
||||
@@ -5,16 +5,16 @@ import { describe, expect, it } from "vitest";
|
||||
import {
|
||||
applyMigrations,
|
||||
createNodeSqliteFactory,
|
||||
createSqliteSessionStore,
|
||||
type SqliteDatabase,
|
||||
type SqliteDatabaseFactory,
|
||||
type SqliteRunResult,
|
||||
type SqliteSessionMetadata,
|
||||
SqliteSessionStorage,
|
||||
SqliteSessionStore,
|
||||
type SqliteStatement,
|
||||
} from "../../../storage/sqlite-node/src/index.ts";
|
||||
import { NodeExecutionEnv } from "../../src/harness/env/nodejs.ts";
|
||||
import { SessionRepository } from "../../src/harness/session/repo-utils.ts";
|
||||
import { createSessionRepository } from "../../src/harness/session/repo-utils.ts";
|
||||
import { createAssistantMessage, createUserMessage } from "./session-test-utils.ts";
|
||||
|
||||
function createTempDir(): string {
|
||||
@@ -64,13 +64,39 @@ class CountingDatabase implements SqliteDatabase {
|
||||
}
|
||||
}
|
||||
|
||||
function createCloseCountingSqliteFactory(): {
|
||||
sqlite: SqliteDatabaseFactory;
|
||||
counts: { opens: number; closes: number };
|
||||
} {
|
||||
const source = createNodeSqliteFactory();
|
||||
const counts = { opens: 0, closes: 0 };
|
||||
return {
|
||||
counts,
|
||||
sqlite: {
|
||||
async open(path) {
|
||||
const db = await source.open(path);
|
||||
counts.opens += 1;
|
||||
return {
|
||||
exec: (sql) => db.exec(sql),
|
||||
prepare: (sql) => db.prepare(sql),
|
||||
transaction: (fn) => db.transaction(fn),
|
||||
async close() {
|
||||
counts.closes += 1;
|
||||
await db.close();
|
||||
},
|
||||
};
|
||||
},
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
describe("SQLite migrations", () => {
|
||||
it("applies file-based migrations and records them", async () => {
|
||||
const root = createTempDir();
|
||||
const databasePath = join(root, "sessions.sqlite");
|
||||
const env = new NodeExecutionEnv({ cwd: root });
|
||||
const sqlite = createNodeSqliteFactory();
|
||||
const repo = new SessionRepository({ store: new SqliteSessionStore({ env, sqlite, databasePath }) });
|
||||
const repo = createSessionRepository({ store: createSqliteSessionStore({ env, sqlite, databasePath }) });
|
||||
await repo.create({ cwd: root, id: "session-1" });
|
||||
|
||||
const db = await sqlite.open(databasePath);
|
||||
@@ -112,8 +138,8 @@ describe("SQLite migrations", () => {
|
||||
const root = createTempDir();
|
||||
const databasePath = join(root, "sessions.sqlite");
|
||||
const env = new NodeExecutionEnv({ cwd: root });
|
||||
const repo = new SessionRepository({
|
||||
store: new SqliteSessionStore({ env, sqlite: createNodeSqliteFactory(), databasePath }),
|
||||
const repo = createSessionRepository({
|
||||
store: createSqliteSessionStore({ env, sqlite: createNodeSqliteFactory(), databasePath }),
|
||||
});
|
||||
const source = await repo.create({
|
||||
cwd: root,
|
||||
@@ -139,7 +165,7 @@ describe("SQLite migrations", () => {
|
||||
const databasePath = join(root, "sessions.sqlite");
|
||||
const env = new NodeExecutionEnv({ cwd: root });
|
||||
const sqlite = createNodeSqliteFactory();
|
||||
const repo = new SessionRepository({ store: new SqliteSessionStore({ env, sqlite, databasePath }) });
|
||||
const repo = createSessionRepository({ store: createSqliteSessionStore({ env, sqlite, databasePath }) });
|
||||
const session = await repo.create({ cwd: root, id: "session-1" });
|
||||
const rootId = await session.appendMessage(createUserMessage("root"));
|
||||
const childId = await session.appendMessage(createAssistantMessage("child"));
|
||||
@@ -175,7 +201,7 @@ describe("SQLite migrations", () => {
|
||||
const databasePath = join(root, "sessions.sqlite");
|
||||
const env = new NodeExecutionEnv({ cwd: root });
|
||||
const sqlite = createNodeSqliteFactory();
|
||||
const repo = new SessionRepository({ store: new SqliteSessionStore({ env, sqlite, databasePath }) });
|
||||
const repo = createSessionRepository({ store: createSqliteSessionStore({ env, sqlite, databasePath }) });
|
||||
const session = await repo.create({ cwd: root, id: "session-1" });
|
||||
const rootId = await session.appendMessage(createUserMessage("root"));
|
||||
const firstChildId = await session.appendMessage(createAssistantMessage("first child"));
|
||||
@@ -203,8 +229,8 @@ describe("SQLite migrations", () => {
|
||||
const root = createTempDir();
|
||||
const databasePath = join(root, "sessions.sqlite");
|
||||
const env = new NodeExecutionEnv({ cwd: root });
|
||||
const repo = new SessionRepository({
|
||||
store: new SqliteSessionStore({ env, sqlite: createNodeSqliteFactory(), databasePath }),
|
||||
const repo = createSessionRepository({
|
||||
store: createSqliteSessionStore({ env, sqlite: createNodeSqliteFactory(), databasePath }),
|
||||
});
|
||||
const session = await repo.create({ cwd: root, id: "session-1" });
|
||||
const rootId = await session.appendMessage(createUserMessage("root"));
|
||||
@@ -225,8 +251,8 @@ describe("SQLite migrations", () => {
|
||||
const root = createTempDir();
|
||||
const databasePath = join(root, "sessions.sqlite");
|
||||
const env = new NodeExecutionEnv({ cwd: root });
|
||||
const repo = new SessionRepository({
|
||||
store: new SqliteSessionStore({ env, sqlite: createNodeSqliteFactory(), databasePath }),
|
||||
const repo = createSessionRepository({
|
||||
store: createSqliteSessionStore({ env, sqlite: createNodeSqliteFactory(), databasePath }),
|
||||
});
|
||||
const session = await repo.create({ cwd: root, id: "session-1" });
|
||||
await session.appendMessage(createUserMessage("one"));
|
||||
@@ -254,8 +280,8 @@ describe("SQLite migrations", () => {
|
||||
open: async () => db,
|
||||
};
|
||||
const env = new NodeExecutionEnv({ cwd: root });
|
||||
const repo = new SessionRepository({
|
||||
store: new SqliteSessionStore({ env, sqlite, databasePath: join(root, "sessions.sqlite") }),
|
||||
const repo = createSessionRepository({
|
||||
store: createSqliteSessionStore({ env, sqlite, databasePath: join(root, "sessions.sqlite") }),
|
||||
});
|
||||
|
||||
await expect(repo.create({ cwd: root, id: "session-1" })).rejects.toThrow("insert failed");
|
||||
@@ -274,8 +300,8 @@ describe("SQLite migrations", () => {
|
||||
open: async () => db,
|
||||
};
|
||||
const env = new NodeExecutionEnv({ cwd: root });
|
||||
const repo = new SessionRepository({
|
||||
store: new SqliteSessionStore({ env, sqlite, databasePath: join(root, "sessions.sqlite") }),
|
||||
const repo = createSessionRepository({
|
||||
store: createSqliteSessionStore({ env, sqlite, databasePath: join(root, "sessions.sqlite") }),
|
||||
});
|
||||
const metadata: SqliteSessionMetadata = {
|
||||
id: "missing",
|
||||
@@ -293,38 +319,15 @@ describe("SQLite migrations", () => {
|
||||
const root = createTempDir();
|
||||
const databasePath = join(root, "sessions.sqlite");
|
||||
const env = new NodeExecutionEnv({ cwd: root });
|
||||
const storage = new SqliteSessionStore({ env, sqlite: createNodeSqliteFactory(), databasePath });
|
||||
const repo = new SessionRepository({ store: storage });
|
||||
let cleanupCount = 0;
|
||||
const sourceStorage = {
|
||||
async getEntries() {
|
||||
return [];
|
||||
},
|
||||
async getPathToRootOrCompaction() {
|
||||
return [];
|
||||
},
|
||||
async cleanup() {
|
||||
cleanupCount += 1;
|
||||
},
|
||||
} as const;
|
||||
const originalOpen = storage.open.bind(storage);
|
||||
storage.open = async () => sourceStorage as never;
|
||||
const { sqlite, counts } = createCloseCountingSqliteFactory();
|
||||
const repo = createSessionRepository({ store: createSqliteSessionStore({ env, sqlite, databasePath }) });
|
||||
const source = await repo.create({ cwd: root, id: "session-1" });
|
||||
|
||||
try {
|
||||
await repo.fork(
|
||||
{
|
||||
id: "session-1",
|
||||
createdAt: new Date().toISOString(),
|
||||
cwd: root,
|
||||
path: databasePath,
|
||||
},
|
||||
{ cwd: root, id: "session-2" },
|
||||
);
|
||||
} finally {
|
||||
storage.open = originalOpen;
|
||||
}
|
||||
counts.opens = 0;
|
||||
counts.closes = 0;
|
||||
await repo.fork(await source.getMetadata(), { cwd: root, id: "session-2" });
|
||||
|
||||
expect(cleanupCount).toBe(1);
|
||||
expect(counts).toEqual({ opens: 2, closes: 2 });
|
||||
});
|
||||
|
||||
it("restores in-memory state when appendEntry fails after mutating caches", async () => {
|
||||
@@ -367,7 +370,7 @@ describe("SQLite migrations", () => {
|
||||
const databasePath = join(root, "sessions.sqlite");
|
||||
const env = new NodeExecutionEnv({ cwd: root });
|
||||
const sqlite = createNodeSqliteFactory();
|
||||
const repo = new SessionRepository({ store: new SqliteSessionStore({ env, sqlite, databasePath }) });
|
||||
const repo = createSessionRepository({ store: createSqliteSessionStore({ env, sqlite, databasePath }) });
|
||||
const session = await repo.create({ cwd: root, id: "session-1" });
|
||||
const userId = await session.appendMessage(createUserMessage("one"));
|
||||
await session.appendThinkingLevelChange("high");
|
||||
|
||||
@@ -3,17 +3,15 @@ import { describe, expect, it } from "vitest";
|
||||
import {
|
||||
createNodeSqliteFactory,
|
||||
createSqliteSessionRepository,
|
||||
createSqliteSessionSearch,
|
||||
createSqliteSessionStore,
|
||||
type SqliteSessionMetadata,
|
||||
SqliteSessionSearch,
|
||||
SqliteSessionStore,
|
||||
type SqliteSessionStoreApi,
|
||||
} from "../../../storage/sqlite-node/src/index.ts";
|
||||
import { NodeExecutionEnv } from "../../src/harness/env/nodejs.ts";
|
||||
import { createJsonlSessionRepository, JsonlSessionStore } from "../../src/harness/session/jsonl-repo.ts";
|
||||
import { SessionRepository } from "../../src/harness/session/repo-utils.ts";
|
||||
import { createJsonlSessionRepository, createJsonlSessionStore } from "../../src/harness/session/jsonl-repo.ts";
|
||||
import { createSessionRepository } from "../../src/harness/session/repo-utils.ts";
|
||||
import type {
|
||||
JsonlSessionMetadata,
|
||||
JsonlSessionStoreApi,
|
||||
SessionSearch,
|
||||
SessionSearchHit,
|
||||
SessionSearchIndex,
|
||||
@@ -100,12 +98,12 @@ describe("JsonlSessionStore with SQLite search index", () => {
|
||||
const root = createTempDir();
|
||||
const env = new NodeExecutionEnv({ cwd: root });
|
||||
const sqlite = createNodeSqliteFactory();
|
||||
const search = new SqliteSessionSearch<JsonlSessionMetadata>({
|
||||
const search = createSqliteSessionSearch<JsonlSessionMetadata>({
|
||||
env,
|
||||
sqlite,
|
||||
databasePath: join(root, "search.sqlite"),
|
||||
});
|
||||
const canonical = new JsonlSessionStore({ fs: env, sessionsRoot: join(root, "sessions") });
|
||||
const canonical = createJsonlSessionStore({ fs: env, sessionsRoot: join(root, "sessions") });
|
||||
const store = {
|
||||
load: (metadata) => canonical.load(metadata),
|
||||
list: (options) => canonical.list(options),
|
||||
@@ -134,8 +132,8 @@ describe("JsonlSessionStore with SQLite search index", () => {
|
||||
await search.replaceSession(metadata, (await canonical.load(metadata)).entries);
|
||||
return metadata;
|
||||
},
|
||||
} satisfies JsonlSessionStoreApi;
|
||||
const repo = new SessionRepository({ store, search });
|
||||
} satisfies ReturnType<typeof createJsonlSessionStore>;
|
||||
const repo = createSessionRepository({ store, search });
|
||||
const session = await repo.create({ cwd: root, id: "jsonl-session" });
|
||||
const entryId = await session.appendMessage(createUserMessage("Find the auth defect"));
|
||||
|
||||
@@ -150,17 +148,17 @@ describe("JsonlSessionStore with multiple search indexes", () => {
|
||||
const root = createTempDir();
|
||||
const env = new NodeExecutionEnv({ cwd: root });
|
||||
const sqlite = createNodeSqliteFactory();
|
||||
const primary = new SqliteSessionSearch<JsonlSessionMetadata>({
|
||||
const primary = createSqliteSessionSearch<JsonlSessionMetadata>({
|
||||
env,
|
||||
sqlite,
|
||||
databasePath: join(root, "primary-search.sqlite"),
|
||||
});
|
||||
const secondary = new SqliteSessionSearch<JsonlSessionMetadata>({
|
||||
const secondary = createSqliteSessionSearch<JsonlSessionMetadata>({
|
||||
env,
|
||||
sqlite,
|
||||
databasePath: join(root, "secondary-search.sqlite"),
|
||||
});
|
||||
const canonical = new JsonlSessionStore({ fs: env, sessionsRoot: join(root, "sessions") });
|
||||
const canonical = createJsonlSessionStore({ fs: env, sessionsRoot: join(root, "sessions") });
|
||||
const store = {
|
||||
load: (metadata) => canonical.load(metadata),
|
||||
list: (options) => canonical.list(options),
|
||||
@@ -195,8 +193,8 @@ describe("JsonlSessionStore with multiple search indexes", () => {
|
||||
await secondary.replaceSession(metadata, entries);
|
||||
return metadata;
|
||||
},
|
||||
} satisfies JsonlSessionStoreApi;
|
||||
const repo = new SessionRepository({ store, search: primary });
|
||||
} satisfies ReturnType<typeof createJsonlSessionStore>;
|
||||
const repo = createSessionRepository({ store, search: primary });
|
||||
const session = await repo.create({ cwd: root, id: "jsonl-session" });
|
||||
const entryId = await session.appendMessage(createUserMessage("indexed in both places"));
|
||||
|
||||
@@ -235,7 +233,7 @@ describe("SqliteSessionStore with custom search", () => {
|
||||
const env = new NodeExecutionEnv({ cwd: root });
|
||||
const sqlite = createNodeSqliteFactory();
|
||||
const databasePath = join(root, "sessions.sqlite");
|
||||
const canonical = new SqliteSessionStore({ env, sqlite, databasePath });
|
||||
const canonical = createSqliteSessionStore({ env, sqlite, databasePath });
|
||||
const store = {
|
||||
load: (metadata) => canonical.load(metadata),
|
||||
list: (options) => canonical.list(options),
|
||||
@@ -264,8 +262,8 @@ describe("SqliteSessionStore with custom search", () => {
|
||||
await index.replaceSession(metadata, (await canonical.load(metadata)).entries);
|
||||
return metadata;
|
||||
},
|
||||
} satisfies SqliteSessionStoreApi;
|
||||
const repo = new SessionRepository({ store, search });
|
||||
} satisfies ReturnType<typeof createSqliteSessionStore>;
|
||||
const repo = createSessionRepository({ store, search });
|
||||
const session = await repo.create({ cwd: root, id: "session-1" });
|
||||
const metadata = await session.getMetadata();
|
||||
const entryId = await session.appendMessage(createUserMessage("indexed remotely"));
|
||||
|
||||
@@ -2,4 +2,4 @@
|
||||
|
||||
Node sqlite storage backend for `@earendil-works/pi-agent-core` sessions. Provides the
|
||||
`node:sqlite` adapter (`SqliteDatabase` implementation) and the SQLite session
|
||||
store/storage implementation (`SqliteSessionStore`, migrations, materialized views).
|
||||
store/storage implementation (`createSqliteSessionStore`, migrations, materialized views).
|
||||
|
||||
@@ -7,13 +7,15 @@ import type {
|
||||
} from "@earendil-works/pi-agent-core";
|
||||
import {
|
||||
createSessionId,
|
||||
createSessionRepository,
|
||||
getEntriesToFork,
|
||||
getFileSystemResultOrThrow,
|
||||
SessionError,
|
||||
SessionRepository,
|
||||
type SessionRepository,
|
||||
type SessionStore,
|
||||
} from "@earendil-works/pi-agent-core";
|
||||
import { applyMigrations } from "./migrations.ts";
|
||||
import { SqliteSessionSearch } from "./search-backend.ts";
|
||||
import { createSqliteSessionSearch } from "./search-backend.ts";
|
||||
import { SqliteSessionStorage } from "./storage/index.ts";
|
||||
import { rowToMetadata, type SessionRow } from "./storage/sessions.ts";
|
||||
import type {
|
||||
@@ -22,7 +24,6 @@ import type {
|
||||
SqliteSessionCreateOptions,
|
||||
SqliteSessionListOptions,
|
||||
SqliteSessionMetadata,
|
||||
SqliteSessionStoreApi,
|
||||
SqliteSessionStoreEnv,
|
||||
} from "./types.ts";
|
||||
|
||||
@@ -51,7 +52,9 @@ export type SqliteSessionStoreOptions = {
|
||||
databasePath: string;
|
||||
};
|
||||
|
||||
export class SqliteSessionStore implements SqliteSessionStoreApi {
|
||||
class SqliteSessionStore
|
||||
implements SessionStore<SqliteSessionMetadata, SqliteSessionCreateOptions, SqliteSessionListOptions>
|
||||
{
|
||||
private readonly env: SqliteSessionStoreEnv;
|
||||
private readonly sqlite: SqliteDatabaseFactory;
|
||||
private readonly databasePathInput: string;
|
||||
@@ -242,7 +245,9 @@ export class SqliteSessionStore implements SqliteSessionStoreApi {
|
||||
}
|
||||
}
|
||||
|
||||
export function createSqliteSessionStore(options: SqliteSessionStoreOptions): SqliteSessionStore {
|
||||
export function createSqliteSessionStore(
|
||||
options: SqliteSessionStoreOptions,
|
||||
): SessionStore<SqliteSessionMetadata, SqliteSessionCreateOptions, SqliteSessionListOptions> {
|
||||
return new SqliteSessionStore(options);
|
||||
}
|
||||
|
||||
@@ -250,8 +255,8 @@ export function createSqliteSessionRepository(
|
||||
options: SqliteSessionStoreOptions,
|
||||
): SessionRepository<SqliteSessionMetadata, SqliteSessionCreateOptions, SqliteSessionListOptions> {
|
||||
const store = createSqliteSessionStore(options);
|
||||
return new SessionRepository({
|
||||
return createSessionRepository({
|
||||
store,
|
||||
search: new SqliteSessionSearch<SqliteSessionMetadata>({ ...options, mode: "canonical" }),
|
||||
search: createSqliteSessionSearch<SqliteSessionMetadata>({ ...options, mode: "canonical" }),
|
||||
});
|
||||
}
|
||||
|
||||
@@ -83,7 +83,7 @@ CREATE VIRTUAL TABLE IF NOT EXISTS session_search_fts USING fts5(
|
||||
* Storage-independent SQLite FTS search. Its database may be separate from,
|
||||
* or shared with, the canonical session backend.
|
||||
*/
|
||||
export class SqliteSessionSearch<TMetadata extends SessionMetadata = SessionMetadata>
|
||||
class SqliteSessionSearch<TMetadata extends SessionMetadata = SessionMetadata>
|
||||
implements SessionSearch<TMetadata>, SessionSearchIndex<TMetadata>
|
||||
{
|
||||
private readonly options: {
|
||||
@@ -236,3 +236,9 @@ export class SqliteSessionSearch<TMetadata extends SessionMetadata = SessionMeta
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
export function createSqliteSessionSearch<TMetadata extends SessionMetadata = SessionMetadata>(
|
||||
options: SqliteSessionSearchOptions,
|
||||
): SessionSearch<TMetadata> & SessionSearchIndex<TMetadata> {
|
||||
return new SqliteSessionSearch<TMetadata>(options);
|
||||
}
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
import type { FileSystem, SessionCreateOptions, SessionMetadata, SessionStore } from "@earendil-works/pi-agent-core";
|
||||
import type { FileSystem, SessionCreateOptions, SessionMetadata } from "@earendil-works/pi-agent-core";
|
||||
|
||||
/** Result of a prepared SQLite statement execution. */
|
||||
export interface SqliteRunResult {
|
||||
@@ -44,7 +44,4 @@ export interface SqliteSessionListOptions {
|
||||
cwd?: string;
|
||||
}
|
||||
|
||||
export interface SqliteSessionStoreApi
|
||||
extends SessionStore<SqliteSessionMetadata, SqliteSessionCreateOptions, SqliteSessionListOptions> {}
|
||||
|
||||
export type SqliteSessionStoreEnv = Pick<FileSystem, "absolutePath" | "createDir" | "exists">;
|
||||
|
||||
@@ -5,15 +5,14 @@ import {
|
||||
bashExecutionToText,
|
||||
convertToLlm,
|
||||
createCustomMessage,
|
||||
createInMemorySessionRepository,
|
||||
FileError,
|
||||
formatPromptTemplateInvocation,
|
||||
formatSkillInvocation,
|
||||
formatSkillsForSystemPrompt,
|
||||
getOrThrow,
|
||||
InMemorySessionStore,
|
||||
ok,
|
||||
parseCommandArgs,
|
||||
SessionRepository,
|
||||
streamProxy,
|
||||
toError,
|
||||
truncateHead,
|
||||
@@ -28,7 +27,7 @@ const stream = createAssistantMessageEventStream();
|
||||
|
||||
const agent = new Agent({ initialState: { model }, streamFn: streamSimple });
|
||||
agent.steer({ role: "user", content: [{ type: "text", text: "queued" }], timestamp: 0 });
|
||||
const repo = new SessionRepository({ store: new InMemorySessionStore() });
|
||||
const repo = createInMemorySessionRepository();
|
||||
const result = getOrThrow(ok({ value: 1 }));
|
||||
const customMessage = createCustomMessage("note", "hello", true, undefined, "2026-01-01T00:00:00.000Z");
|
||||
const llmMessages = convertToLlm([customMessage]);
|
||||
|
||||
Reference in New Issue
Block a user