Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
50 changes: 40 additions & 10 deletions packages/pi/src/actor-sqlite.ts
Original file line number Diff line number Diff line change
@@ -1,10 +1,39 @@
import type {
SqliteDatabase,
SqliteExecutor,
SqliteValue,
import {
SQLITE_MIGRATIONS,
type SqliteDatabase,
type SqliteExecutor,
type SqliteValue,
} from "@earendil-works/pi-durable/storage/sqlite";
import type { RawAccess } from "rivetkit/db";

/**
* Pi Durable's tables: the ones its migrations create, and `durable_schema`,
* which its migration runner creates. The actor's database also holds the
* app's tables, so Pi's tables are stored with a `pi_` prefix.
*/
export const PI_DURABLE_TABLES: readonly string[] = [
"durable_schema",
...SQLITE_MIGRATIONS.flatMap((migration) =>
migration.statements.flatMap((statement) => {
const table = /^\s*CREATE TABLE (?:IF NOT EXISTS )?(\w+)/i.exec(
statement,
)?.[1];
return table === undefined ? [] : [table];
}),
),
];

/** Pi Durable's SQL names its tables only as table references, so whole words are table names. */
const PI_TABLE_NAMES = new RegExp(
`\\b(${PI_DURABLE_TABLES.join("|")})\\b`,
"g",
);

/** Pi Durable's SQL with its table names prefixed. */
function prefixTables(sql: string): string {
return sql.replace(PI_TABLE_NAMES, "pi_$1");
}

/**
* Pi Durable's `SqliteDatabase` over the actor's SQLite database.
*
Expand All @@ -23,20 +52,21 @@ export function actorSqlite(db: RawAccess): SqliteDatabase {
}

/**
* Pi Durable names the row type of each query it sends. SQLite returns plain
* rows, so `get` and `all` take Pi's word for their shape.
* Pi Durable's SQL against `pi_` tables. Pi Durable names the row type of each
* query it sends. SQLite returns plain rows, so `get` and `all` take Pi's word
* for their shape.
*/
function actorSqliteExecutor(db: Pick<RawAccess, "execute">): SqliteExecutor {
return {
exec: async (sql) => {
await db.execute(sql);
await db.execute(prefixTables(sql));
},
run: async (sql, ...params) => {
await db.execute(sql, ...params);
await db.execute(prefixTables(sql), ...params);
},
get: async <T extends object>(sql: string, ...params: SqliteValue[]) =>
((await db.execute(sql, ...params)) as T[]).at(0),
((await db.execute(prefixTables(sql), ...params)) as T[]).at(0),
all: async <T extends object>(sql: string, ...params: SqliteValue[]) =>
(await db.execute(sql, ...params)) as T[],
(await db.execute(prefixTables(sql), ...params)) as T[],
};
}
18 changes: 17 additions & 1 deletion packages/pi/src/storage.ts
Original file line number Diff line number Diff line change
@@ -1,12 +1,14 @@
import type { ConversationId } from "@earendil-works/pi-durable";
import type { RawAccess } from "rivetkit/db";
import { migrations } from "rivetkit/unstable/migrations";
import { PI_DURABLE_TABLES } from "./actor-sqlite.js";

export type PiDatabase = Pick<RawAccess, "execute">;

/**
* `pi()`'s own tables. Pi Durable creates its tables in the same database
* itself. The `pi_durable` names are stored in actors' databases, so they stay.
* itself, with the `pi_` prefix `actorSqlite` adds. The `pi_durable` names are
* stored in actors' databases, so they stay.
*/
export const migratePiTables = migrations({
tableName: "pi_durable_schema_version",
Expand Down Expand Up @@ -45,6 +47,20 @@ export const migratePiTables = migrations({
) STRICT;
`,
},
{
version: 4,
// Pi Durable's tables of agents from 0.5.1 and earlier have no prefix. Renaming them keeps their data.
up: async (db) => {
for (const table of PI_DURABLE_TABLES) {
const [found] = await db.execute(
"SELECT 1 FROM sqlite_master WHERE type = 'table' AND name = ?",
table,
);
if (found)
await db.execute(`ALTER TABLE ${table} RENAME TO pi_${table}`);
}
},
},
],
});

Expand Down
46 changes: 43 additions & 3 deletions packages/pi/tests/actor.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,19 @@ const LOG = Array.from({ length: 20_000 }, (_, i) => `line ${i + 1}`).join(
"\n",
);

/** Pi Durable's tables as 0.5.1 stored them, without a prefix. */
const PI_DURABLE_TABLES_051 = [
"durable_schema",
"durable_metadata",
"record_ids",
"conversations",
"entries",
"tasks",
"submissions",
"documents",
"document_revisions",
];

/** Tool calls of the crash test. A tool's first run blocks until the stop aborts it; a rerun returns at once. */
const toolRuns = { safe: 0, unsafe: 0 };

Expand Down Expand Up @@ -147,7 +160,7 @@ function buildRegistry(mock: MockModel, root: string) {
},
failCommitsOfPoison: async (c) => {
await c.db.execute(
`CREATE TRIGGER IF NOT EXISTS test_fail_poison BEFORE INSERT ON documents
`CREATE TRIGGER IF NOT EXISTS test_fail_poison BEFORE INSERT ON pi_documents
WHEN NEW.kind = '"app.poison"' BEGIN SELECT RAISE(ABORT, 'injected storage failure'); END`,
);
},
Expand Down Expand Up @@ -179,12 +192,24 @@ function buildRegistry(mock: MockModel, root: string) {
nap: (c) => {
c.sleep();
},
// What 0.5.1 stored: `pi_file` with its rows, at schema version 2.
// What 0.5.1 stored: Pi Durable's tables without a prefix, and
// `pi_file` with its rows, at schema version 2.
storeAs051: async (c) => {
for (const table of PI_DURABLE_TABLES_051) {
await c.db.execute(`ALTER TABLE pi_${table} RENAME TO ${table}`);
}
await c.db.execute(
"UPDATE pi_durable_schema_version SET schema_version = 2",
);
},
// The app's own table, named like one of Pi Durable's.
addAppTask: async (c, title: string) => {
await c.db.execute(
"CREATE TABLE IF NOT EXISTS tasks (title TEXT NOT NULL) STRICT",
);
await c.db.execute("INSERT INTO tasks (title) VALUES (?)", title);
return c.db.execute<{ title: string }>("SELECT title FROM tasks");
},
},
});
const backoff = pi({
Expand Down Expand Up @@ -609,7 +634,7 @@ describe("pi actor", () => {
});
});

test("files an agent without a sandbox stored on 0.5.1 stay readable after the upgrade", async (c) => {
test("an agent without a sandbox stored on 0.5.1 keeps its conversation and its files after the upgrade", async (c) => {
const { client } = await setupTest(c, registry);
const key = ["files-051", randomUUID()];
const handle = client.files.getOrCreate(key);
Expand All @@ -622,13 +647,28 @@ describe("pi actor", () => {

const root = await handle.harness.root();
const { messages } = await handle.conversation.context(root.id);
expect(JSON.stringify(messages)).toContain("create hello.txt");
const read = messages.findLast(
(message) => message.role === "toolResult" && message.toolName === "read",
);
expect(read).toMatchObject({ isError: false });
expect(JSON.stringify(read?.content)).toContain("line two");
});

test("the app's own table named like one of Pi Durable's lives next to the agent's conversations", async (c) => {
const { client } = await setupTest(c, registry);
const handle = client.files.getOrCreate(["app-tables", randomUUID()]);
await handle.prompt("say hello");

expect(await handle.addAppTask("ship 0.5.2")).toEqual([
{ title: "ship 0.5.2" },
]);
expect(await handle.prompt("say hello")).toMatchObject({
status: "done",
text: "Hi there!",
});
});

test("sandbox tools cannot write outside the sandbox or read the actor host's environment", async (c) => {
const { client } = await setupTest(c, registry);
const handle = client.coder.getOrCreate(["escape", randomUUID()]);
Expand Down
2 changes: 1 addition & 1 deletion packages/pi/tests/upgrade.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -217,7 +217,7 @@ function buildRegistry(mock: MockModel, root: string) {
bumpSchema: async (c: {
db: { execute: (sql: string) => Promise<unknown> };
}) => {
await c.db.execute("UPDATE durable_schema SET version = version + 1");
await c.db.execute("UPDATE pi_durable_schema SET version = version + 1");
},
// What the earlier session-based pi() stored: its own session table and
// `pi_sandbox`, and none of this pi()'s tables.
Expand Down
Loading