mirror of
https://github.com/earendil-works/pi.git
synced 2026-10-02 00:35:27 +08:00
feat(agent): add pico memory storage foundation
This commit is contained in:
@@ -1,5 +1,8 @@
|
||||
export * from "./addresses.ts";
|
||||
export * from "./core.ts";
|
||||
export * from "./entries.ts";
|
||||
export * from "./memory-storage.ts";
|
||||
export * from "./runtime.ts";
|
||||
export * from "./session.ts";
|
||||
export * from "./storage.ts";
|
||||
export * from "./tasks.ts";
|
||||
|
||||
@@ -0,0 +1,427 @@
|
||||
import type { Context } from "@earendil-works/chord";
|
||||
import type { Address, Element, List, Value } from "./addresses.ts";
|
||||
import type { Id, JsonValue, Seq } from "./core.ts";
|
||||
import type { Conversation, Entry } from "./entries.ts";
|
||||
import {
|
||||
type ConversationScan,
|
||||
type EntryScan,
|
||||
InvalidHistoryPosition,
|
||||
type Page,
|
||||
type PageQuery,
|
||||
type Storage,
|
||||
type TaskScan,
|
||||
type Write,
|
||||
} from "./storage.ts";
|
||||
import type { Task } from "./tasks.ts";
|
||||
|
||||
type ValueHistoryWrite =
|
||||
| { readonly seq: Seq; readonly type: "set"; readonly value: JsonValue }
|
||||
| { readonly seq: Seq; readonly type: "delete" };
|
||||
|
||||
type ListHistoryWrite =
|
||||
| { readonly seq: Seq; readonly type: "append"; readonly element: Element<JsonValue> }
|
||||
| { readonly seq: Seq; readonly type: "remove"; readonly elementId: Id }
|
||||
| { readonly seq: Seq; readonly type: "clear" };
|
||||
|
||||
type StoredState =
|
||||
| { readonly kind: "value"; readonly rewind: false; value?: JsonValue }
|
||||
| { readonly kind: "list"; readonly rewind: false; readonly elements: Map<Id, Element<JsonValue>> }
|
||||
| { readonly kind: "value"; readonly rewind: true; readonly writes: ValueHistoryWrite[] }
|
||||
| { readonly kind: "list"; readonly rewind: true; readonly writes: ListHistoryWrite[] };
|
||||
|
||||
type RewindListState = Extract<StoredState, { kind: "list"; rewind: true }>;
|
||||
|
||||
interface AddressState {
|
||||
value?: Extract<StoredState, { kind: "value" }>;
|
||||
list?: Extract<StoredState, { kind: "list" }>;
|
||||
}
|
||||
|
||||
interface NamespaceState {
|
||||
unkeyed?: AddressState;
|
||||
readonly keyed: Map<string, AddressState>;
|
||||
}
|
||||
|
||||
type ScopeState = Map<string, NamespaceState>;
|
||||
|
||||
interface CappedListState {
|
||||
readonly state: RewindListState | undefined;
|
||||
readonly cap: Seq | undefined;
|
||||
}
|
||||
|
||||
interface StoredEntry {
|
||||
readonly seq: Seq;
|
||||
readonly entry: Entry;
|
||||
}
|
||||
|
||||
function rewindable(
|
||||
address: Address,
|
||||
): address is Address & { readonly scope: Extract<Address["scope"], { type: "conversation" }>; readonly rewind: true } {
|
||||
return address.scope.type === "conversation" && address.rewind;
|
||||
}
|
||||
|
||||
function scopeId(address: Address): Id | undefined {
|
||||
if (address.scope.type === "conversation") return address.scope.conversationId;
|
||||
if (address.scope.type === "task") return address.scope.taskId;
|
||||
if (address.scope.type === "shared") return address.scope.id;
|
||||
return undefined;
|
||||
}
|
||||
|
||||
function page<T>(
|
||||
values: Iterable<T>,
|
||||
query: PageQuery,
|
||||
id: (value: T) => Id,
|
||||
matches: (value: T) => boolean,
|
||||
ascending: boolean,
|
||||
): Page<T> {
|
||||
const items: T[] = [];
|
||||
let more = false;
|
||||
for (const value of values) {
|
||||
const valueId = id(value);
|
||||
if (
|
||||
!matches(value) ||
|
||||
(query.cursor !== undefined && (ascending ? valueId <= query.cursor : valueId >= query.cursor))
|
||||
) {
|
||||
continue;
|
||||
}
|
||||
if (items.length === query.limit) {
|
||||
more = true;
|
||||
break;
|
||||
}
|
||||
items.push(value);
|
||||
}
|
||||
return { items, ...(more && items.length > 0 ? { next: id(items.at(-1)!) } : {}) };
|
||||
}
|
||||
|
||||
export class MemoryStorage implements Storage {
|
||||
private nextObjectId: Id = 1;
|
||||
private nextSeq: Seq = 1;
|
||||
private readonly conversations = new Map<Id, Conversation>();
|
||||
private readonly entries = new Map<Id, StoredEntry>();
|
||||
private readonly entriesByConversation = new Map<Id, StoredEntry[]>();
|
||||
private readonly tasks = new Map<Id, Task>();
|
||||
private readonly outputReferences = new Map<Id, number>();
|
||||
private readonly sessionState: ScopeState = new Map();
|
||||
private readonly conversationStickyState = new Map<Id, ScopeState>();
|
||||
private readonly conversationRewindableState = new Map<Id, ScopeState>();
|
||||
private readonly taskState = new Map<Id, ScopeState>();
|
||||
private readonly sharedState = new Map<Id, ScopeState>();
|
||||
|
||||
nextId(): Id {
|
||||
return this.nextObjectId++;
|
||||
}
|
||||
|
||||
async commit(writes: readonly Write[], _ctx: Context): Promise<readonly Seq[]> {
|
||||
const seqs = writes.map((_, index) => this.nextSeq + index);
|
||||
const affectedOutputs = new Set<Id>();
|
||||
let nextObjectId = this.nextObjectId;
|
||||
for (const write of writes) {
|
||||
if (write.type === "conversation.create") nextObjectId = Math.max(nextObjectId, write.conversation.id + 1);
|
||||
else if (write.type === "entry.append") nextObjectId = Math.max(nextObjectId, write.entry.id + 1);
|
||||
else if (write.type === "task.create") nextObjectId = Math.max(nextObjectId, write.task.id + 1);
|
||||
else if (write.type === "list.append") nextObjectId = Math.max(nextObjectId, write.element.id + 1);
|
||||
if (write.type !== "task.create" && write.type !== "task.set") continue;
|
||||
const currentOutput = this.tasks.get(write.task.id)?.output?.id;
|
||||
if (currentOutput !== undefined) affectedOutputs.add(currentOutput);
|
||||
if (write.task.output !== undefined) affectedOutputs.add(write.task.output.id);
|
||||
}
|
||||
for (let index = 0; index < writes.length; index++) this.apply(writes[index]!, seqs[index]!);
|
||||
for (const write of writes) {
|
||||
if (write.type === "task.set" && write.task.status === "terminal") this.taskState.delete(write.task.id);
|
||||
}
|
||||
for (const id of affectedOutputs) {
|
||||
if (!this.outputReferences.has(id)) this.sharedState.delete(id);
|
||||
}
|
||||
this.nextObjectId = nextObjectId;
|
||||
this.nextSeq += writes.length;
|
||||
return seqs;
|
||||
}
|
||||
|
||||
async getConversations(ids: readonly Id[], _ctx: Context): Promise<ReadonlyMap<Id, Conversation>> {
|
||||
const result = new Map<Id, Conversation>();
|
||||
for (const id of ids) {
|
||||
const conversation = this.conversations.get(id);
|
||||
if (conversation !== undefined) result.set(id, conversation);
|
||||
}
|
||||
return result;
|
||||
}
|
||||
|
||||
async scanConversations(query: ConversationScan, _ctx: Context): Promise<Page<Conversation>> {
|
||||
return page(
|
||||
this.conversations.values(),
|
||||
query,
|
||||
(conversation) => conversation.id,
|
||||
(conversation) =>
|
||||
(query.parent === undefined || conversation.parent?.conversationId === query.parent) &&
|
||||
(query.owner === undefined || conversation.owner === query.owner),
|
||||
true,
|
||||
);
|
||||
}
|
||||
|
||||
async getEntries(ids: readonly Id[], _ctx: Context): Promise<ReadonlyMap<Id, Entry>> {
|
||||
const result = new Map<Id, Entry>();
|
||||
for (const id of ids) {
|
||||
const stored = this.entries.get(id);
|
||||
if (stored !== undefined) result.set(id, stored.entry);
|
||||
}
|
||||
return result;
|
||||
}
|
||||
|
||||
async scanEntries(query: EntryScan, _ctx: Context): Promise<Page<Entry>> {
|
||||
return page(
|
||||
this.visibleEntries(query.conversationId, query.through),
|
||||
query,
|
||||
(entry) => entry.id,
|
||||
(entry) => query.kind === undefined || entry.kind === query.kind,
|
||||
false,
|
||||
);
|
||||
}
|
||||
|
||||
async newestHead(conversationId: Id, at: Id, _ctx: Context): Promise<Entry | undefined> {
|
||||
for (const entry of this.visibleEntries(conversationId, at)) {
|
||||
if (entry.head !== undefined) return entry;
|
||||
}
|
||||
return undefined;
|
||||
}
|
||||
|
||||
async getTasks(ids: readonly Id[], _ctx: Context): Promise<ReadonlyMap<Id, Task>> {
|
||||
const result = new Map<Id, Task>();
|
||||
for (const id of ids) {
|
||||
const task = this.tasks.get(id);
|
||||
if (task !== undefined) result.set(id, task);
|
||||
}
|
||||
return result;
|
||||
}
|
||||
|
||||
async scanTasks(query: TaskScan, _ctx: Context): Promise<Page<Task>> {
|
||||
return page(
|
||||
this.tasks.values(),
|
||||
query,
|
||||
(task) => task.id,
|
||||
(task) =>
|
||||
(query.conversationIds === undefined || query.conversationIds.includes(task.conversationId)) &&
|
||||
(query.statuses === undefined || query.statuses.includes(task.status)) &&
|
||||
(query.kind === undefined || task.kind === query.kind) &&
|
||||
(query.abort === undefined || (task.abort === true) === query.abort) &&
|
||||
(query.outputId === undefined || task.output?.id === query.outputId),
|
||||
true,
|
||||
);
|
||||
}
|
||||
|
||||
async getValue<T extends JsonValue>(address: Value<T>, at: Id | undefined, _ctx: Context): Promise<T | undefined> {
|
||||
if (at !== undefined) {
|
||||
if (!rewindable(address)) throw new InvalidHistoryPosition(at);
|
||||
return this.rewindValue(address, this.historyPosition(address.scope.conversationId, at));
|
||||
}
|
||||
if (rewindable(address)) return this.rewindValue(address, undefined);
|
||||
const state = this.storedState(address);
|
||||
return state?.kind === "value" && !state.rewind ? (state.value as T | undefined) : undefined;
|
||||
}
|
||||
|
||||
async readList<T extends JsonValue>(
|
||||
address: List<T>,
|
||||
at: Id | undefined,
|
||||
_ctx: Context,
|
||||
): Promise<readonly Element<T>[]> {
|
||||
if (at !== undefined) {
|
||||
if (!rewindable(address)) throw new InvalidHistoryPosition(at);
|
||||
return this.rewindListElements(address, this.historyPosition(address.scope.conversationId, at));
|
||||
}
|
||||
if (rewindable(address)) return this.rewindListElements(address, undefined);
|
||||
const state = this.storedState(address);
|
||||
if (state?.kind !== "list" || state.rewind) return [];
|
||||
return [...(state.elements.values() as Iterable<Element<T>>)];
|
||||
}
|
||||
|
||||
private historyPosition(conversationId: Id, at: Id): Seq {
|
||||
const stored = this.entries.get(at);
|
||||
if (stored === undefined) throw new InvalidHistoryPosition(at);
|
||||
|
||||
let currentId = conversationId;
|
||||
let cap: Seq | undefined;
|
||||
while (true) {
|
||||
if (stored.entry.conversationId === currentId && (cap === undefined || stored.seq <= cap)) return stored.seq;
|
||||
const conversation = this.conversations.get(currentId);
|
||||
if (conversation?.parent === undefined) throw new InvalidHistoryPosition(at);
|
||||
const parentAtSeq = this.entries.get(conversation.parent.at)?.seq;
|
||||
if (parentAtSeq === undefined) throw new InvalidHistoryPosition(at);
|
||||
cap = Math.min(cap ?? parentAtSeq, parentAtSeq);
|
||||
currentId = conversation.parent.conversationId;
|
||||
}
|
||||
}
|
||||
|
||||
private rewindValue<T extends JsonValue>(address: Value<T>, cap: Seq | undefined): T | undefined {
|
||||
let conversationId = address.scope.type === "conversation" ? address.scope.conversationId : undefined;
|
||||
while (conversationId !== undefined) {
|
||||
const state = this.conversationStoredState(conversationId, address);
|
||||
if (state?.kind === "value" && state.rewind) {
|
||||
for (let index = state.writes.length - 1; index >= 0; index--) {
|
||||
const write = state.writes[index]!;
|
||||
if (cap !== undefined && write.seq > cap) continue;
|
||||
return write.type === "set" ? (write.value as T) : undefined;
|
||||
}
|
||||
}
|
||||
const conversation = this.conversations.get(conversationId);
|
||||
if (conversation?.parent === undefined) return undefined;
|
||||
const parentAtSeq = this.historyPosition(conversation.parent.conversationId, conversation.parent.at);
|
||||
cap = Math.min(cap ?? parentAtSeq, parentAtSeq);
|
||||
conversationId = conversation.parent.conversationId;
|
||||
}
|
||||
return undefined;
|
||||
}
|
||||
|
||||
private rewindListElements<T extends JsonValue>(
|
||||
address: List<T> & { readonly scope: Extract<Address["scope"], { type: "conversation" }> },
|
||||
cap: Seq | undefined,
|
||||
): readonly Element<T>[] {
|
||||
const ancestry: CappedListState[] = [];
|
||||
let conversationId = address.scope.conversationId;
|
||||
while (true) {
|
||||
const stored = this.conversationStoredState(conversationId, address);
|
||||
ancestry.push({ state: stored?.kind === "list" && stored.rewind ? stored : undefined, cap });
|
||||
const conversation = this.conversations.get(conversationId);
|
||||
if (conversation?.parent === undefined) break;
|
||||
const parentAtSeq = this.historyPosition(conversation.parent.conversationId, conversation.parent.at);
|
||||
cap = Math.min(cap ?? parentAtSeq, parentAtSeq);
|
||||
conversationId = conversation.parent.conversationId;
|
||||
}
|
||||
|
||||
const elements = new Map<Id, Element<T>>();
|
||||
for (const segment of ancestry.reverse()) {
|
||||
for (const write of segment.state?.writes ?? []) {
|
||||
if (segment.cap !== undefined && write.seq > segment.cap) continue;
|
||||
if (write.type === "append") elements.set(write.element.id, write.element as Element<T>);
|
||||
else if (write.type === "remove") elements.delete(write.elementId);
|
||||
else elements.clear();
|
||||
}
|
||||
}
|
||||
return [...elements.values()];
|
||||
}
|
||||
|
||||
private *visibleEntries(conversationId: Id, through: Id | undefined): Iterable<Entry> {
|
||||
let conversation = this.conversations.get(conversationId);
|
||||
let cap = through;
|
||||
while (conversation !== undefined) {
|
||||
const entries = this.entriesByConversation.get(conversation.id) ?? [];
|
||||
for (let index = entries.length - 1; index >= 0; index--) {
|
||||
const entry = entries[index]!.entry;
|
||||
if (cap === undefined || entry.id <= cap) yield entry;
|
||||
}
|
||||
if (conversation.parent === undefined) return;
|
||||
cap = Math.min(cap ?? conversation.parent.at, conversation.parent.at);
|
||||
conversation = this.conversations.get(conversation.parent.conversationId);
|
||||
}
|
||||
}
|
||||
|
||||
private apply(write: Write, seq: Seq): void {
|
||||
if (write.type === "conversation.create") {
|
||||
this.conversations.set(write.conversation.id, write.conversation);
|
||||
} else if (write.type === "entry.append") {
|
||||
const entry = write.entry;
|
||||
const stored = { seq, entry };
|
||||
this.entries.set(entry.id, stored);
|
||||
const entries = this.entriesByConversation.get(entry.conversationId);
|
||||
if (entries === undefined) this.entriesByConversation.set(entry.conversationId, [stored]);
|
||||
else entries.push(stored);
|
||||
} else if (write.type === "task.create" || write.type === "task.set") {
|
||||
const current = this.tasks.get(write.task.id);
|
||||
if (current?.status !== "terminal" && current?.output !== undefined)
|
||||
this.removeOutputReference(current.output.id);
|
||||
this.tasks.set(write.task.id, write.task);
|
||||
if (write.task.status !== "terminal" && write.task.output !== undefined) {
|
||||
this.outputReferences.set(write.task.output.id, (this.outputReferences.get(write.task.output.id) ?? 0) + 1);
|
||||
}
|
||||
} else if (write.type === "value.set") {
|
||||
const state = this.valueState(write.address);
|
||||
if (state.rewind) state.writes.push({ seq, type: "set", value: write.value });
|
||||
else state.value = write.value;
|
||||
} else if (write.type === "value.delete") {
|
||||
const state = this.valueState(write.address);
|
||||
if (state.rewind) state.writes.push({ seq, type: "delete" });
|
||||
else state.value = undefined;
|
||||
} else {
|
||||
const state = this.listState(write.address);
|
||||
if (state.rewind) {
|
||||
if (write.type === "list.append") state.writes.push({ seq, type: "append", element: write.element });
|
||||
else if (write.type === "list.remove")
|
||||
state.writes.push({ seq, type: "remove", elementId: write.elementId });
|
||||
else state.writes.push({ seq, type: "clear" });
|
||||
} else if (write.type === "list.append") state.elements.set(write.element.id, write.element);
|
||||
else if (write.type === "list.remove") state.elements.delete(write.elementId);
|
||||
else state.elements.clear();
|
||||
}
|
||||
}
|
||||
|
||||
private removeOutputReference(id: Id): void {
|
||||
const count = this.outputReferences.get(id)!;
|
||||
if (count === 1) this.outputReferences.delete(id);
|
||||
else this.outputReferences.set(id, count - 1);
|
||||
}
|
||||
|
||||
private storedState(address: Address): StoredState | undefined {
|
||||
return this.stateInScope(this.scopeState(address, false), address);
|
||||
}
|
||||
|
||||
private conversationStoredState(conversationId: Id, address: Address): StoredState | undefined {
|
||||
const scopes = address.rewind ? this.conversationRewindableState : this.conversationStickyState;
|
||||
return this.stateInScope(scopes.get(conversationId), address);
|
||||
}
|
||||
|
||||
private stateInScope(scope: ScopeState | undefined, address: Address): StoredState | undefined {
|
||||
const namespace = scope?.get(address.namespace);
|
||||
const state = address.key === undefined ? namespace?.unkeyed : namespace?.keyed.get(address.key);
|
||||
return address.kind === "value" ? state?.value : state?.list;
|
||||
}
|
||||
|
||||
private valueState(address: Value<JsonValue>): Extract<StoredState, { kind: "value" }> {
|
||||
const initial: StoredState = rewindable(address)
|
||||
? { kind: "value", rewind: true, writes: [] }
|
||||
: { kind: "value", rewind: false };
|
||||
return this.ensureState(address, initial) as Extract<StoredState, { kind: "value" }>;
|
||||
}
|
||||
|
||||
private listState(address: List<JsonValue>): Extract<StoredState, { kind: "list" }> {
|
||||
const initial: StoredState = rewindable(address)
|
||||
? { kind: "list", rewind: true, writes: [] }
|
||||
: { kind: "list", rewind: false, elements: new Map() };
|
||||
return this.ensureState(address, initial) as Extract<StoredState, { kind: "list" }>;
|
||||
}
|
||||
|
||||
private ensureState(address: Address, initial: StoredState): StoredState {
|
||||
const scope = this.scopeState(address, true)!;
|
||||
let namespace = scope.get(address.namespace);
|
||||
if (namespace === undefined) {
|
||||
namespace = { keyed: new Map() };
|
||||
scope.set(address.namespace, namespace);
|
||||
}
|
||||
let state = address.key === undefined ? namespace.unkeyed : namespace.keyed.get(address.key);
|
||||
if (state === undefined) {
|
||||
state = {};
|
||||
if (address.key === undefined) namespace.unkeyed = state;
|
||||
else namespace.keyed.set(address.key, state);
|
||||
}
|
||||
const stored = address.kind === "value" ? state.value : state.list;
|
||||
if (stored !== undefined) return stored;
|
||||
if (initial.kind === "value") state.value = initial;
|
||||
else state.list = initial;
|
||||
return initial;
|
||||
}
|
||||
|
||||
private scopeState(address: Address, create: boolean): ScopeState | undefined {
|
||||
if (address.scope.type === "session") return this.sessionState;
|
||||
const id = scopeId(address)!;
|
||||
const scopes =
|
||||
address.scope.type === "conversation"
|
||||
? address.rewind
|
||||
? this.conversationRewindableState
|
||||
: this.conversationStickyState
|
||||
: address.scope.type === "task"
|
||||
? this.taskState
|
||||
: this.sharedState;
|
||||
let state = scopes.get(id);
|
||||
if (state === undefined && create) {
|
||||
state = new Map();
|
||||
scopes.set(id, state);
|
||||
}
|
||||
return state;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,310 @@
|
||||
import type { Context } from "@earendil-works/chord";
|
||||
import type { Address, Element, List, Value } from "./addresses.ts";
|
||||
import type { Id, JsonValue } from "./core.ts";
|
||||
import type { Conversation, Entry } from "./entries.ts";
|
||||
import type { ConversationScan, EntryScan, NewTask, Page, Storage, TaskScan, Write } from "./storage.ts";
|
||||
import type { Task } from "./tasks.ts";
|
||||
|
||||
export class ReadAfterWrite extends Error {
|
||||
constructor() {
|
||||
super("Transaction reads must precede writes");
|
||||
this.name = "ReadAfterWrite";
|
||||
}
|
||||
}
|
||||
|
||||
export class ScratchRetired extends Error {
|
||||
readonly taskId: Id;
|
||||
|
||||
constructor(taskId: Id) {
|
||||
super(`Scratch for task ${taskId} is retired`);
|
||||
this.name = "ScratchRetired";
|
||||
this.taskId = taskId;
|
||||
}
|
||||
}
|
||||
|
||||
export class SharedOutputRetired extends Error {
|
||||
readonly outputId: Id;
|
||||
|
||||
constructor(outputId: Id) {
|
||||
super(`Shared task output ${outputId} is retired`);
|
||||
this.name = "SharedOutputRetired";
|
||||
this.outputId = outputId;
|
||||
}
|
||||
}
|
||||
|
||||
export type TaskCreate = Omit<NewTask, "id" | "output"> & {
|
||||
readonly output?: { readonly id?: Id; readonly kind: string };
|
||||
};
|
||||
|
||||
export interface Transaction {
|
||||
getConversation(id: Id): Promise<Conversation | undefined>;
|
||||
getEntry(id: Id): Promise<Entry | undefined>;
|
||||
getEntries(ids: readonly Id[]): Promise<ReadonlyMap<Id, Entry>>;
|
||||
getTask(id: Id): Promise<Task | undefined>;
|
||||
getTasks(ids: readonly Id[]): Promise<ReadonlyMap<Id, Task>>;
|
||||
createConversation(conversation: Omit<Conversation, "id">): Id;
|
||||
entry(entry: Omit<Entry, "id">): Id;
|
||||
task(task: TaskCreate): Id;
|
||||
setTask(task: Task): void;
|
||||
value<T extends JsonValue>(
|
||||
address: Value<T>,
|
||||
): {
|
||||
get(at?: Id): Promise<T | undefined>;
|
||||
set(value: T): void;
|
||||
delete(): void;
|
||||
};
|
||||
list<T extends JsonValue>(
|
||||
address: List<T>,
|
||||
): {
|
||||
read(at?: Id): Promise<readonly Element<T>[]>;
|
||||
append(value: T): Id;
|
||||
remove(id: Id): void;
|
||||
clear(): void;
|
||||
};
|
||||
}
|
||||
|
||||
class TransactionState {
|
||||
private readonly storage: Storage;
|
||||
private readonly ctx: Context;
|
||||
private active = true;
|
||||
private writing = false;
|
||||
private pendingReads = 0;
|
||||
private readonly writes: Write[] = [];
|
||||
|
||||
constructor(storage: Storage, ctx: Context) {
|
||||
this.storage = storage;
|
||||
this.ctx = ctx;
|
||||
}
|
||||
|
||||
facade(): Transaction {
|
||||
return Object.freeze({
|
||||
getConversation: (id: Id) => this.getConversation(id),
|
||||
getEntry: (id: Id) => this.getEntry(id),
|
||||
getEntries: (ids: readonly Id[]) => this.getEntries(ids),
|
||||
getTask: (id: Id) => this.getTask(id),
|
||||
getTasks: (ids: readonly Id[]) => this.getTasks(ids),
|
||||
createConversation: (conversation: Omit<Conversation, "id">) => this.createConversation(conversation),
|
||||
entry: (entry: Omit<Entry, "id">) => this.entry(entry),
|
||||
task: (task: TaskCreate) => this.task(task),
|
||||
setTask: (task: Task) => this.setTask(task),
|
||||
value: <T extends JsonValue>(address: Value<T>) => this.value(address),
|
||||
list: <T extends JsonValue>(address: List<T>) => this.list(address),
|
||||
});
|
||||
}
|
||||
|
||||
finish(): readonly Write[] {
|
||||
this.assertActive();
|
||||
if (this.pendingReads !== 0) throw new Error("Transaction has pending reads");
|
||||
this.active = false;
|
||||
return this.writes;
|
||||
}
|
||||
|
||||
seal(): void {
|
||||
this.active = false;
|
||||
}
|
||||
|
||||
async getConversation(id: Id): Promise<Conversation | undefined> {
|
||||
return (await this.read(() => this.storage.getConversations([id], this.ctx))).get(id);
|
||||
}
|
||||
|
||||
async getEntry(id: Id): Promise<Entry | undefined> {
|
||||
return (await this.getEntries([id])).get(id);
|
||||
}
|
||||
|
||||
getEntries(ids: readonly Id[]): Promise<ReadonlyMap<Id, Entry>> {
|
||||
return this.read(() => this.storage.getEntries(ids, this.ctx));
|
||||
}
|
||||
|
||||
async getTask(id: Id): Promise<Task | undefined> {
|
||||
return (await this.getTasks([id])).get(id);
|
||||
}
|
||||
|
||||
getTasks(ids: readonly Id[]): Promise<ReadonlyMap<Id, Task>> {
|
||||
return this.read(() => this.storage.getTasks(ids, this.ctx));
|
||||
}
|
||||
|
||||
createConversation(conversation: Omit<Conversation, "id">): Id {
|
||||
const id = this.allocate();
|
||||
this.writes.push({ type: "conversation.create", conversation: { ...conversation, id } });
|
||||
return id;
|
||||
}
|
||||
|
||||
entry(entry: Omit<Entry, "id">): Id {
|
||||
const id = this.allocate();
|
||||
this.writes.push({ type: "entry.append", entry: { ...entry, id } });
|
||||
return id;
|
||||
}
|
||||
|
||||
task(task: TaskCreate): Id {
|
||||
const id = this.allocate();
|
||||
const { output: requestedOutput, ...base } = task;
|
||||
const output =
|
||||
requestedOutput === undefined ? undefined : { id: requestedOutput.id ?? id, kind: requestedOutput.kind };
|
||||
this.writes.push({
|
||||
type: "task.create",
|
||||
task: { ...base, id, ...(output === undefined ? {} : { output }) },
|
||||
});
|
||||
return id;
|
||||
}
|
||||
|
||||
setTask(task: Task): void {
|
||||
this.write({ type: "task.set", task });
|
||||
}
|
||||
|
||||
value<T extends JsonValue>(
|
||||
address: Value<T>,
|
||||
): {
|
||||
get(at?: Id): Promise<T | undefined>;
|
||||
set(value: T): void;
|
||||
delete(): void;
|
||||
} {
|
||||
this.assertActive();
|
||||
return Object.freeze({
|
||||
get: (at?: Id) => this.read(() => this.storage.getValue(address, at, this.ctx)),
|
||||
set: (value: T) => this.write({ type: "value.set", address, value }),
|
||||
delete: () => this.write({ type: "value.delete", address }),
|
||||
});
|
||||
}
|
||||
|
||||
list<T extends JsonValue>(
|
||||
address: List<T>,
|
||||
): {
|
||||
read(at?: Id): Promise<readonly Element<T>[]>;
|
||||
append(value: T): Id;
|
||||
remove(id: Id): void;
|
||||
clear(): void;
|
||||
} {
|
||||
this.assertActive();
|
||||
return Object.freeze({
|
||||
read: (at?: Id) => this.read(() => this.storage.readList(address, at, this.ctx)),
|
||||
append: (value: T) => {
|
||||
const id = this.allocate();
|
||||
this.writes.push({ type: "list.append", address, element: { id, value } });
|
||||
return id;
|
||||
},
|
||||
remove: (elementId: Id) => this.write({ type: "list.remove", address, elementId }),
|
||||
clear: () => this.write({ type: "list.clear", address }),
|
||||
});
|
||||
}
|
||||
|
||||
private allocate(): Id {
|
||||
this.assertWritable();
|
||||
return this.storage.nextId();
|
||||
}
|
||||
|
||||
private write(write: Write): void {
|
||||
this.assertWritable();
|
||||
this.writes.push(write);
|
||||
}
|
||||
|
||||
private async read<T>(read: () => Promise<T>): Promise<T> {
|
||||
this.assertReadable();
|
||||
this.pendingReads++;
|
||||
try {
|
||||
return await read();
|
||||
} finally {
|
||||
this.pendingReads--;
|
||||
}
|
||||
}
|
||||
|
||||
private assertReadable(): void {
|
||||
this.assertActive();
|
||||
if (this.writing) throw new ReadAfterWrite();
|
||||
}
|
||||
|
||||
private assertWritable(): void {
|
||||
this.assertActive();
|
||||
if (this.pendingReads !== 0) throw new Error("Transaction writes must await reads");
|
||||
this.writing = true;
|
||||
}
|
||||
|
||||
private assertActive(): void {
|
||||
if (!this.active) throw new Error("Transaction is closed");
|
||||
}
|
||||
}
|
||||
|
||||
export class Session {
|
||||
private readonly storage: Storage;
|
||||
private tail: Promise<void> = Promise.resolve();
|
||||
|
||||
constructor(storage: Storage) {
|
||||
this.storage = storage;
|
||||
}
|
||||
|
||||
async commit<T>(build: (tx: Transaction) => T | Promise<T>, ctx: Context): Promise<T> {
|
||||
const previous = this.tail;
|
||||
let release!: () => void;
|
||||
this.tail = new Promise((resolve) => {
|
||||
release = resolve;
|
||||
});
|
||||
await previous;
|
||||
const tx = new TransactionState(this.storage, ctx);
|
||||
try {
|
||||
const result = await build(tx.facade());
|
||||
const writes = tx.finish();
|
||||
if (writes.length !== 0) await this.storage.commit(writes, ctx);
|
||||
return result;
|
||||
} finally {
|
||||
tx.seal();
|
||||
release();
|
||||
}
|
||||
}
|
||||
|
||||
getConversations(ids: readonly Id[], ctx: Context): Promise<ReadonlyMap<Id, Conversation>> {
|
||||
return this.storage.getConversations(ids, ctx);
|
||||
}
|
||||
|
||||
scanConversations(query: ConversationScan, ctx: Context): Promise<Page<Conversation>> {
|
||||
return this.storage.scanConversations(query, ctx);
|
||||
}
|
||||
|
||||
getEntries(ids: readonly Id[], ctx: Context): Promise<ReadonlyMap<Id, Entry>> {
|
||||
return this.storage.getEntries(ids, ctx);
|
||||
}
|
||||
|
||||
scanEntries(query: EntryScan, ctx: Context): Promise<Page<Entry>> {
|
||||
return this.storage.scanEntries(query, ctx);
|
||||
}
|
||||
|
||||
newestHead(conversationId: Id, at: Id, ctx: Context): Promise<Entry | undefined> {
|
||||
return this.storage.newestHead(conversationId, at, ctx);
|
||||
}
|
||||
|
||||
getTasks(ids: readonly Id[], ctx: Context): Promise<ReadonlyMap<Id, Task>> {
|
||||
return this.storage.getTasks(ids, ctx);
|
||||
}
|
||||
|
||||
scanTasks(query: TaskScan, ctx: Context): Promise<Page<Task>> {
|
||||
return this.storage.scanTasks(query, ctx);
|
||||
}
|
||||
|
||||
async getValue<T extends JsonValue>(address: Value<T>, at: Id | undefined, ctx: Context): Promise<T | undefined> {
|
||||
await this.assertReadable(address, ctx);
|
||||
return this.storage.getValue(address, at, ctx);
|
||||
}
|
||||
|
||||
async readList<T extends JsonValue>(
|
||||
address: List<T>,
|
||||
at: Id | undefined,
|
||||
ctx: Context,
|
||||
): Promise<readonly Element<T>[]> {
|
||||
await this.assertReadable(address, ctx);
|
||||
return this.storage.readList(address, at, ctx);
|
||||
}
|
||||
|
||||
private async outputIsLive(id: Id, ctx: Context): Promise<boolean> {
|
||||
return (
|
||||
(await this.storage.scanTasks({ outputId: id, statuses: ["pending", "running"], limit: 1 }, ctx)).items
|
||||
.length > 0
|
||||
);
|
||||
}
|
||||
|
||||
private async assertReadable(address: Address, ctx: Context): Promise<void> {
|
||||
if (address.scope.type === "task") {
|
||||
const task = (await this.storage.getTasks([address.scope.taskId], ctx)).get(address.scope.taskId);
|
||||
if (task === undefined || task.status === "terminal") throw new ScratchRetired(address.scope.taskId);
|
||||
} else if (address.scope.type === "shared" && !(await this.outputIsLive(address.scope.id, ctx))) {
|
||||
throw new SharedOutputRetired(address.scope.id);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,81 @@
|
||||
import type { Context } from "@earendil-works/chord";
|
||||
import type { Element, List, Value } from "./addresses.ts";
|
||||
import type { Id, JsonValue, Seq } from "./core.ts";
|
||||
import type { Conversation, Entry } from "./entries.ts";
|
||||
import type { Task } from "./tasks.ts";
|
||||
|
||||
export interface PageQuery {
|
||||
/** Last item returned by the same scan. */
|
||||
readonly cursor?: Id;
|
||||
readonly limit: number;
|
||||
}
|
||||
|
||||
export interface Page<T> {
|
||||
readonly items: readonly T[];
|
||||
readonly next?: Id;
|
||||
}
|
||||
|
||||
export class InvalidHistoryPosition extends Error {
|
||||
readonly at: Id;
|
||||
|
||||
constructor(at: Id) {
|
||||
super(`Entry ${at} is not a valid history position for this address`);
|
||||
this.name = "InvalidHistoryPosition";
|
||||
this.at = at;
|
||||
}
|
||||
}
|
||||
|
||||
export type NewTask = Omit<Task, "abort" | "checkpoint" | "outcome" | "owns" | "status"> & {
|
||||
readonly status: "pending";
|
||||
readonly owns: readonly [];
|
||||
readonly abort?: never;
|
||||
readonly checkpoint?: never;
|
||||
readonly outcome?: never;
|
||||
};
|
||||
|
||||
export type StateWrite =
|
||||
| { readonly type: "value.set"; readonly address: Value<JsonValue>; readonly value: JsonValue }
|
||||
| { readonly type: "value.delete"; readonly address: Value<JsonValue> }
|
||||
| { readonly type: "list.append"; readonly address: List<JsonValue>; readonly element: Element<JsonValue> }
|
||||
| { readonly type: "list.remove"; readonly address: List<JsonValue>; readonly elementId: Id }
|
||||
| { readonly type: "list.clear"; readonly address: List<JsonValue> };
|
||||
|
||||
export type Write =
|
||||
| StateWrite
|
||||
| { readonly type: "conversation.create"; readonly conversation: Conversation }
|
||||
| { readonly type: "entry.append"; readonly entry: Entry }
|
||||
| { readonly type: "task.create"; readonly task: NewTask }
|
||||
| { readonly type: "task.set"; readonly task: Task };
|
||||
|
||||
export type ConversationScan = PageQuery & {
|
||||
readonly parent?: Id;
|
||||
readonly owner?: Id;
|
||||
};
|
||||
|
||||
export type EntryScan = PageQuery & {
|
||||
readonly conversationId: Id;
|
||||
readonly kind?: string;
|
||||
readonly through?: Id;
|
||||
};
|
||||
|
||||
export type TaskScan = PageQuery & {
|
||||
readonly conversationIds?: readonly Id[];
|
||||
readonly statuses?: readonly Task["status"][];
|
||||
readonly kind?: string;
|
||||
readonly abort?: boolean;
|
||||
readonly outputId?: Id;
|
||||
};
|
||||
|
||||
export interface Storage {
|
||||
nextId(): Id;
|
||||
commit(writes: readonly Write[], ctx: Context): Promise<readonly Seq[]>;
|
||||
getConversations(ids: readonly Id[], ctx: Context): Promise<ReadonlyMap<Id, Conversation>>;
|
||||
scanConversations(query: ConversationScan, ctx: Context): Promise<Page<Conversation>>;
|
||||
getEntries(ids: readonly Id[], ctx: Context): Promise<ReadonlyMap<Id, Entry>>;
|
||||
scanEntries(query: EntryScan, ctx: Context): Promise<Page<Entry>>;
|
||||
newestHead(conversationId: Id, at: Id, ctx: Context): Promise<Entry | undefined>;
|
||||
getTasks(ids: readonly Id[], ctx: Context): Promise<ReadonlyMap<Id, Task>>;
|
||||
scanTasks(query: TaskScan, ctx: Context): Promise<Page<Task>>;
|
||||
getValue<T extends JsonValue>(address: Value<T>, at: Id | undefined, ctx: Context): Promise<T | undefined>;
|
||||
readList<T extends JsonValue>(address: List<T>, at: Id | undefined, ctx: Context): Promise<readonly Element<T>[]>;
|
||||
}
|
||||
@@ -0,0 +1,571 @@
|
||||
import { TODO_CONTEXT } from "@earendil-works/chord/context";
|
||||
import { describe, expect, it, vi } from "vitest";
|
||||
import {
|
||||
defineList,
|
||||
defineValue,
|
||||
InvalidHistoryPosition,
|
||||
type JsonValue,
|
||||
MemoryStorage,
|
||||
ReadAfterWrite,
|
||||
ScratchRetired,
|
||||
Session,
|
||||
SharedOutputRetired,
|
||||
type Task,
|
||||
type TaskCreate,
|
||||
type Transaction,
|
||||
} from "../../../src/harness/pico/index.ts";
|
||||
|
||||
const ctx = TODO_CONTEXT;
|
||||
|
||||
function taskSpec(
|
||||
conversationId: number,
|
||||
extra: Partial<Pick<TaskCreate, "kind" | "output" | "background">> = {},
|
||||
): TaskCreate {
|
||||
return {
|
||||
conversationId,
|
||||
kind: "test.task",
|
||||
input: { value: conversationId },
|
||||
after: [],
|
||||
owns: [],
|
||||
status: "pending",
|
||||
...extra,
|
||||
};
|
||||
}
|
||||
|
||||
type LiveTask = Task & { readonly status: "pending" | "running"; readonly outcome?: never };
|
||||
|
||||
async function setTask(session: Session, id: number, update: (task: LiveTask) => Task): Promise<void> {
|
||||
await session.commit(async (tx) => {
|
||||
const task = await tx.getTask(id);
|
||||
if (task === undefined || task.status === "terminal") throw new Error(`Task ${id} is not live`);
|
||||
tx.setTask(update(task));
|
||||
}, ctx);
|
||||
}
|
||||
|
||||
describe("Pico Session over MemoryStorage", () => {
|
||||
it("separates IDs from sequences and burns IDs from failed callbacks", async () => {
|
||||
const storage = new MemoryStorage();
|
||||
expect(await storage.commit([{ type: "conversation.create", conversation: { id: 10 } }], ctx)).toEqual([1]);
|
||||
const state = defineValue<string>({ type: "session" }, "test.state");
|
||||
expect(await storage.commit([{ type: "value.set", address: state, value: "set" }], ctx)).toEqual([2]);
|
||||
expect(await new Session(storage).commit((tx) => tx.createConversation({}), ctx)).toBe(11);
|
||||
|
||||
const sessionStorage = new MemoryStorage();
|
||||
const commit = vi.spyOn(sessionStorage, "commit");
|
||||
const session = new Session(sessionStorage);
|
||||
await expect(
|
||||
session.commit((tx) => {
|
||||
tx.createConversation({});
|
||||
throw new Error("reject");
|
||||
}, ctx),
|
||||
).rejects.toThrow("reject");
|
||||
expect(commit).not.toHaveBeenCalled();
|
||||
expect(await sessionStorage.getConversations([1], ctx)).toEqual(new Map());
|
||||
expect(await session.commit((tx) => tx.createConversation({}), ctx)).toBe(2);
|
||||
expect(commit).toHaveBeenCalledTimes(1);
|
||||
await expect(commit.mock.results[0]!.value).resolves.toEqual([1]);
|
||||
await expect(
|
||||
session.commit(async (tx) => {
|
||||
tx.createConversation({});
|
||||
await tx.getConversation(1);
|
||||
}, ctx),
|
||||
).rejects.toBeInstanceOf(ReadAfterWrite);
|
||||
expect(await session.commit((tx) => tx.createConversation({}), ctx)).toBe(4);
|
||||
expect(
|
||||
await session.commit(async (tx) => {
|
||||
const read = tx.getConversation(1);
|
||||
expect(() => tx.createConversation({})).toThrow("Transaction writes must await reads");
|
||||
await read;
|
||||
return tx.createConversation({});
|
||||
}, ctx),
|
||||
).toBe(5);
|
||||
|
||||
let leaked!: Transaction;
|
||||
await session.commit((tx) => {
|
||||
leaked = tx;
|
||||
}, ctx);
|
||||
expect(() => leaked.createConversation({})).toThrow("Transaction is closed");
|
||||
|
||||
const concurrent = await Promise.all([
|
||||
session.commit(async (tx) => {
|
||||
await tx.getConversation(1);
|
||||
return tx.createConversation({});
|
||||
}, ctx),
|
||||
session.commit((tx) => tx.createConversation({}), ctx),
|
||||
]);
|
||||
expect(concurrent).toEqual([6, 7]);
|
||||
expect(commit).toHaveBeenCalledTimes(5);
|
||||
});
|
||||
|
||||
it("leaves state and sequence unchanged when persistence rejects", async () => {
|
||||
const storage = new MemoryStorage();
|
||||
const commit = vi.spyOn(storage, "commit").mockRejectedValueOnce(new Error("persistence failed"));
|
||||
await expect(new Session(storage).commit((tx) => tx.createConversation({}), ctx)).rejects.toThrow(
|
||||
"persistence failed",
|
||||
);
|
||||
expect(commit).toHaveBeenCalledTimes(1);
|
||||
expect(await storage.getConversations([1], ctx)).toEqual(new Map());
|
||||
commit.mockRestore();
|
||||
|
||||
const state = defineValue<string>({ type: "session" }, "test.after-failure");
|
||||
expect(await storage.commit([{ type: "value.set", address: state, value: "stored" }], ctx)).toEqual([1]);
|
||||
expect(await storage.getValue(state, undefined, ctx)).toBe("stored");
|
||||
expect(storage.nextId()).toBe(2);
|
||||
});
|
||||
|
||||
it("uses write sequences for historical state inside one atomic storage batch", async () => {
|
||||
const storage = new MemoryStorage();
|
||||
const conversation = storage.nextId();
|
||||
const entry = storage.nextId();
|
||||
const beforeElement = storage.nextId();
|
||||
const afterElement = storage.nextId();
|
||||
const value = defineValue<string>({ type: "conversation", conversationId: conversation }, "test.value", {
|
||||
rewind: true,
|
||||
});
|
||||
const list = defineList<string>({ type: "conversation", conversationId: conversation }, "test.list", {
|
||||
rewind: true,
|
||||
});
|
||||
|
||||
expect(
|
||||
await storage.commit(
|
||||
[
|
||||
{ type: "conversation.create", conversation: { id: conversation } },
|
||||
{ type: "value.set", address: value, value: "before" },
|
||||
{ type: "list.append", address: list, element: { id: beforeElement, value: "before" } },
|
||||
{ type: "entry.append", entry: { id: entry, conversationId: conversation, kind: "test.entry" } },
|
||||
{ type: "value.set", address: value, value: "after" },
|
||||
{ type: "list.append", address: list, element: { id: afterElement, value: "after" } },
|
||||
],
|
||||
ctx,
|
||||
),
|
||||
).toEqual([1, 2, 3, 4, 5, 6]);
|
||||
expect(entry).toBe(2);
|
||||
expect(await storage.getValue(value, entry, ctx)).toBe("before");
|
||||
expect(await storage.getValue(value, undefined, ctx)).toBe("after");
|
||||
expect((await storage.readList(list, entry, ctx)).map((item) => item.value)).toEqual(["before"]);
|
||||
expect((await storage.readList(list, undefined, ctx)).map((item) => item.value)).toEqual(["before", "after"]);
|
||||
});
|
||||
|
||||
it("pages conversations and fork-visible entries", async () => {
|
||||
const session = new Session(new MemoryStorage());
|
||||
const ids = await session.commit((tx) => {
|
||||
const root = tx.createConversation({});
|
||||
const first = tx.entry({ conversationId: root, kind: "test.first" });
|
||||
const head = tx.entry({ conversationId: root, kind: "test.head", head: first });
|
||||
const fork = tx.createConversation({ parent: { conversationId: root, at: first } });
|
||||
const forkEntry = tx.entry({ conversationId: fork, kind: "fork.entry" });
|
||||
const late = tx.entry({ conversationId: root, kind: "root.late" });
|
||||
const ownedFork = tx.createConversation({ parent: { conversationId: root, at: first }, owner: 99 });
|
||||
return { root, first, head, fork, forkEntry, late, ownedFork };
|
||||
}, ctx);
|
||||
expect(ids).toEqual({ root: 1, first: 2, head: 3, fork: 4, forkEntry: 5, late: 6, ownedFork: 7 });
|
||||
|
||||
const conversations = await session.scanConversations({ limit: 1 }, ctx);
|
||||
expect(conversations.items.map(({ id }) => id)).toEqual([1]);
|
||||
expect(
|
||||
(await session.scanConversations({ cursor: conversations.next, limit: 1 }, ctx)).items.map(({ id }) => id),
|
||||
).toEqual([4]);
|
||||
expect((await session.getConversations([ids.ownedFork, 999_999], ctx)).has(ids.ownedFork)).toBe(true);
|
||||
expect((await session.scanConversations({ parent: ids.root, limit: 10 }, ctx)).items.map(({ id }) => id)).toEqual(
|
||||
[ids.fork, ids.ownedFork],
|
||||
);
|
||||
expect((await session.scanConversations({ owner: 99, limit: 10 }, ctx)).items.map(({ id }) => id)).toEqual([
|
||||
ids.ownedFork,
|
||||
]);
|
||||
expect(
|
||||
(await session.scanConversations({ parent: ids.root, owner: 99, limit: 10 }, ctx)).items.map(({ id }) => id),
|
||||
).toEqual([ids.ownedFork]);
|
||||
|
||||
const fork = await session.scanEntries({ conversationId: 4, limit: 1 }, ctx);
|
||||
expect(fork.items.map(({ id }) => id)).toEqual([5]);
|
||||
expect(
|
||||
(await session.scanEntries({ conversationId: 4, cursor: fork.next, limit: 1 }, ctx)).items.map(({ id }) => id),
|
||||
).toEqual([2]);
|
||||
expect(
|
||||
(await session.scanEntries({ conversationId: 1, through: 3, limit: 10 }, ctx)).items.map(({ id }) => id),
|
||||
).toEqual([3, 2]);
|
||||
expect(
|
||||
(await session.scanEntries({ conversationId: ids.root, kind: "test.head", limit: 10 }, ctx)).items.map(
|
||||
({ id }) => id,
|
||||
),
|
||||
).toEqual([ids.head]);
|
||||
expect((await session.getEntries([ids.first, 999_999], ctx)).has(ids.first)).toBe(true);
|
||||
expect((await session.newestHead(1, 3, ctx))?.id).toBe(3);
|
||||
expect(await session.newestHead(4, 5, ctx)).toBeUndefined();
|
||||
});
|
||||
|
||||
it("buffers values and list elements into one commit", async () => {
|
||||
const session = new Session(new MemoryStorage());
|
||||
const value = defineValue<JsonValue>({ type: "conversation", conversationId: 1 }, "test.value", {
|
||||
rewind: false,
|
||||
});
|
||||
const list = defineList<string>({ type: "conversation", conversationId: 1 }, "test.list", { rewind: false });
|
||||
const ids = await session.commit((tx) => {
|
||||
const conversation = tx.createConversation({});
|
||||
tx.value(value).set({ version: 1 });
|
||||
const first = tx.list(list).append("a");
|
||||
const second = tx.list(list).append("b");
|
||||
return { conversation, first, second };
|
||||
}, ctx);
|
||||
expect(ids).toEqual({ conversation: 1, first: 2, second: 3 });
|
||||
expect(await session.getValue(value, undefined, ctx)).toEqual({ version: 1 });
|
||||
|
||||
expect((await session.readList(list, undefined, ctx)).map(({ value }) => value)).toEqual(["a", "b"]);
|
||||
|
||||
await session.commit((tx) => {
|
||||
tx.value(value).delete();
|
||||
tx.list(list).remove(ids.first);
|
||||
tx.list(list).clear();
|
||||
}, ctx);
|
||||
expect(await session.getValue(value, undefined, ctx)).toBeUndefined();
|
||||
expect(await session.readList(list, undefined, ctx)).toEqual([]);
|
||||
});
|
||||
|
||||
it("rewinds values through fork caps and preserves local tombstones", async () => {
|
||||
const storage = new MemoryStorage();
|
||||
const session = new Session(storage);
|
||||
const root = await session.commit((tx) => tx.createConversation({}), ctx);
|
||||
const rootValue = defineValue<string>({ type: "conversation", conversationId: root }, "test.history", {
|
||||
rewind: true,
|
||||
});
|
||||
const first = await session.commit((tx) => {
|
||||
tx.value(rootValue).set("root-first");
|
||||
return tx.entry({ conversationId: root, kind: "root.first" });
|
||||
}, ctx);
|
||||
const deleted = await session.commit((tx) => {
|
||||
tx.value(rootValue).delete();
|
||||
return tx.entry({ conversationId: root, kind: "root.deleted" });
|
||||
}, ctx);
|
||||
const late = await session.commit((tx) => {
|
||||
tx.value(rootValue).set("root-late");
|
||||
return tx.entry({ conversationId: root, kind: "root.late" });
|
||||
}, ctx);
|
||||
|
||||
expect(await session.getValue(rootValue, undefined, ctx)).toBe("root-late");
|
||||
expect(await session.getValue(rootValue, first, ctx)).toBe("root-first");
|
||||
expect(await session.getValue(rootValue, deleted, ctx)).toBeUndefined();
|
||||
expect(await session.getValue(rootValue, late, ctx)).toBe("root-late");
|
||||
|
||||
const child = await session.commit(
|
||||
(tx) => tx.createConversation({ parent: { conversationId: root, at: first } }),
|
||||
ctx,
|
||||
);
|
||||
const childValue = defineValue<string>({ type: "conversation", conversationId: child }, "test.history", {
|
||||
rewind: true,
|
||||
});
|
||||
expect(await session.getValue(childValue, undefined, ctx)).toBe("root-first");
|
||||
const childDeleted = await session.commit((tx) => {
|
||||
tx.value(childValue).delete();
|
||||
return tx.entry({ conversationId: child, kind: "child.deleted" });
|
||||
}, ctx);
|
||||
const grandchild = await session.commit(
|
||||
(tx) => tx.createConversation({ parent: { conversationId: child, at: childDeleted } }),
|
||||
ctx,
|
||||
);
|
||||
const grandchildValue = defineValue<string>(
|
||||
{ type: "conversation", conversationId: grandchild },
|
||||
"test.history",
|
||||
{ rewind: true },
|
||||
);
|
||||
await session.commit((tx) => {
|
||||
tx.value(childValue).set("child-late");
|
||||
tx.entry({ conversationId: child, kind: "child.late" });
|
||||
}, ctx);
|
||||
|
||||
expect(await session.getValue(childValue, undefined, ctx)).toBe("child-late");
|
||||
expect(await session.getValue(grandchildValue, undefined, ctx)).toBeUndefined();
|
||||
expect(await session.getValue(grandchildValue, first, ctx)).toBe("root-first");
|
||||
expect(await session.getValue(grandchildValue, childDeleted, ctx)).toBeUndefined();
|
||||
|
||||
const sibling = await session.commit(
|
||||
(tx) => tx.createConversation({ parent: { conversationId: root, at: first } }),
|
||||
ctx,
|
||||
);
|
||||
const siblingValue = defineValue<string>({ type: "conversation", conversationId: sibling }, "test.history", {
|
||||
rewind: true,
|
||||
});
|
||||
expect(await session.getValue(siblingValue, undefined, ctx)).toBe("root-first");
|
||||
await session.commit((tx) => tx.value(siblingValue).set("sibling"), ctx);
|
||||
expect(await session.getValue(siblingValue, undefined, ctx)).toBe("sibling");
|
||||
expect(await session.getValue(childValue, undefined, ctx)).toBe("child-late");
|
||||
expect(await session.getValue(grandchildValue, undefined, ctx)).toBeUndefined();
|
||||
|
||||
const deepAncestorFork = await session.commit(
|
||||
(tx) => tx.createConversation({ parent: { conversationId: grandchild, at: first } }),
|
||||
ctx,
|
||||
);
|
||||
const deepAncestorValue = defineValue<string>(
|
||||
{ type: "conversation", conversationId: deepAncestorFork },
|
||||
"test.history",
|
||||
{ rewind: true },
|
||||
);
|
||||
expect(await session.getValue(deepAncestorValue, undefined, ctx)).toBe("root-first");
|
||||
await expect(session.getValue(childValue, deleted, ctx)).rejects.toBeInstanceOf(InvalidHistoryPosition);
|
||||
await expect(session.getValue(childValue, 999_999, ctx)).rejects.toBeInstanceOf(InvalidHistoryPosition);
|
||||
|
||||
const unrelated = await session.commit((tx) => {
|
||||
const conversation = tx.createConversation({});
|
||||
return tx.entry({ conversationId: conversation, kind: "unrelated" });
|
||||
}, ctx);
|
||||
await expect(session.getValue(childValue, unrelated, ctx)).rejects.toBeInstanceOf(InvalidHistoryPosition);
|
||||
|
||||
const sticky = defineValue<string>({ type: "conversation", conversationId: root }, "test.sticky", {
|
||||
rewind: false,
|
||||
});
|
||||
const sessionValue = defineValue<string>({ type: "session" }, "test.history");
|
||||
const taskList = defineList<string>({ type: "task", taskId: 100 }, "test.history");
|
||||
const sharedValue = defineValue<string>({ type: "shared", id: 200 }, "test.history");
|
||||
await expect(storage.getValue(sticky, first, ctx)).rejects.toBeInstanceOf(InvalidHistoryPosition);
|
||||
await expect(storage.getValue(sessionValue, first, ctx)).rejects.toBeInstanceOf(InvalidHistoryPosition);
|
||||
await expect(storage.readList(taskList, first, ctx)).rejects.toBeInstanceOf(InvalidHistoryPosition);
|
||||
await expect(storage.getValue(sharedValue, first, ctx)).rejects.toBeInstanceOf(InvalidHistoryPosition);
|
||||
});
|
||||
|
||||
it("rewinds complete lists through append, remove, clear, and nested fork caps", async () => {
|
||||
const session = new Session(new MemoryStorage());
|
||||
const root = await session.commit((tx) => tx.createConversation({}), ctx);
|
||||
const rootList = defineList<string>({ type: "conversation", conversationId: root }, "test.history-list", {
|
||||
rewind: true,
|
||||
});
|
||||
const first = await session.commit((tx) => {
|
||||
const a = tx.list(rootList).append("a");
|
||||
const b = tx.list(rootList).append("b");
|
||||
const entry = tx.entry({ conversationId: root, kind: "root.first" });
|
||||
return { a, b, entry };
|
||||
}, ctx);
|
||||
const second = await session.commit((tx) => {
|
||||
tx.list(rootList).remove(first.a);
|
||||
const c = tx.list(rootList).append("c");
|
||||
const entry = tx.entry({ conversationId: root, kind: "root.second" });
|
||||
return { c, entry };
|
||||
}, ctx);
|
||||
const third = await session.commit((tx) => {
|
||||
tx.list(rootList).clear();
|
||||
const d = tx.list(rootList).append("d");
|
||||
const entry = tx.entry({ conversationId: root, kind: "root.third" });
|
||||
return { d, entry };
|
||||
}, ctx);
|
||||
await session.commit((tx) => tx.list(rootList).append("e"), ctx);
|
||||
|
||||
expect((await session.readList(rootList, first.entry, ctx)).map((item) => item.value)).toEqual(["a", "b"]);
|
||||
expect((await session.readList(rootList, second.entry, ctx)).map((item) => item.value)).toEqual(["b", "c"]);
|
||||
expect((await session.readList(rootList, third.entry, ctx)).map((item) => item.value)).toEqual(["d"]);
|
||||
expect((await session.readList(rootList, undefined, ctx)).map((item) => item.value)).toEqual(["d", "e"]);
|
||||
|
||||
const child = await session.commit(
|
||||
(tx) => tx.createConversation({ parent: { conversationId: root, at: second.entry } }),
|
||||
ctx,
|
||||
);
|
||||
const childList = defineList<string>({ type: "conversation", conversationId: child }, "test.history-list", {
|
||||
rewind: true,
|
||||
});
|
||||
const childAppend = await session.commit((tx) => {
|
||||
tx.list(childList).append("x");
|
||||
return tx.entry({ conversationId: child, kind: "child.append" });
|
||||
}, ctx);
|
||||
const childRemove = await session.commit((tx) => {
|
||||
tx.list(childList).remove(first.b);
|
||||
return tx.entry({ conversationId: child, kind: "child.remove" });
|
||||
}, ctx);
|
||||
const grandchild = await session.commit(
|
||||
(tx) => tx.createConversation({ parent: { conversationId: child, at: childRemove } }),
|
||||
ctx,
|
||||
);
|
||||
const grandchildList = defineList<string>(
|
||||
{ type: "conversation", conversationId: grandchild },
|
||||
"test.history-list",
|
||||
{ rewind: true },
|
||||
);
|
||||
const childClear = await session.commit((tx) => {
|
||||
tx.list(childList).clear();
|
||||
tx.list(childList).append("y");
|
||||
return tx.entry({ conversationId: child, kind: "child.clear" });
|
||||
}, ctx);
|
||||
await session.commit((tx) => {
|
||||
tx.list(grandchildList).append("z");
|
||||
tx.entry({ conversationId: grandchild, kind: "grandchild.append" });
|
||||
}, ctx);
|
||||
|
||||
expect((await session.readList(childList, childAppend, ctx)).map((item) => item.value)).toEqual(["b", "c", "x"]);
|
||||
expect((await session.readList(childList, childRemove, ctx)).map((item) => item.value)).toEqual(["c", "x"]);
|
||||
expect((await session.readList(childList, childClear, ctx)).map((item) => item.value)).toEqual(["y"]);
|
||||
expect((await session.readList(childList, undefined, ctx)).map((item) => item.value)).toEqual(["y"]);
|
||||
expect((await session.readList(grandchildList, undefined, ctx)).map((item) => item.value)).toEqual([
|
||||
"c",
|
||||
"x",
|
||||
"z",
|
||||
]);
|
||||
expect((await session.readList(grandchildList, childRemove, ctx)).map((item) => item.value)).toEqual(["c", "x"]);
|
||||
|
||||
const sibling = await session.commit(
|
||||
(tx) => tx.createConversation({ parent: { conversationId: root, at: second.entry } }),
|
||||
ctx,
|
||||
);
|
||||
const siblingList = defineList<string>({ type: "conversation", conversationId: sibling }, "test.history-list", {
|
||||
rewind: true,
|
||||
});
|
||||
expect((await session.readList(siblingList, undefined, ctx)).map((item) => item.value)).toEqual(["b", "c"]);
|
||||
await session.commit((tx) => tx.list(siblingList).append("s"), ctx);
|
||||
expect((await session.readList(siblingList, undefined, ctx)).map((item) => item.value)).toEqual(["b", "c", "s"]);
|
||||
expect((await session.readList(childList, undefined, ctx)).map((item) => item.value)).toEqual(["y"]);
|
||||
|
||||
const deepAncestorFork = await session.commit(
|
||||
(tx) => tx.createConversation({ parent: { conversationId: grandchild, at: second.entry } }),
|
||||
ctx,
|
||||
);
|
||||
const deepAncestorList = defineList<string>(
|
||||
{ type: "conversation", conversationId: deepAncestorFork },
|
||||
"test.history-list",
|
||||
{ rewind: true },
|
||||
);
|
||||
expect((await session.readList(deepAncestorList, undefined, ctx)).map((item) => item.value)).toEqual(["b", "c"]);
|
||||
|
||||
const unrelated = await session.commit((tx) => {
|
||||
const conversation = tx.createConversation({});
|
||||
return tx.entry({ conversationId: conversation, kind: "unrelated" });
|
||||
}, ctx);
|
||||
const stickyList = defineList<string>({ type: "conversation", conversationId: child }, "test.sticky", {
|
||||
rewind: false,
|
||||
});
|
||||
await expect(session.readList(childList, third.entry, ctx)).rejects.toBeInstanceOf(InvalidHistoryPosition);
|
||||
await expect(session.readList(childList, unrelated, ctx)).rejects.toBeInstanceOf(InvalidHistoryPosition);
|
||||
await expect(session.readList(childList, 999_999, ctx)).rejects.toBeInstanceOf(InvalidHistoryPosition);
|
||||
await expect(session.readList(stickyList, childRemove, ctx)).rejects.toBeInstanceOf(InvalidHistoryPosition);
|
||||
});
|
||||
|
||||
it("keeps scopes, kinds, optional keys, and rewind stores separate", async () => {
|
||||
const session = new Session(new MemoryStorage());
|
||||
const conversation = await session.commit((tx) => tx.createConversation({}), ctx);
|
||||
const scope = { type: "conversation", conversationId: conversation } as const;
|
||||
const stickyValue = defineValue<string>(scope, "test.identity", { key: "same", rewind: false });
|
||||
const rewindValue = defineValue<string>(scope, "test.identity", { key: "same", rewind: true });
|
||||
const stickyList = defineList<string>(scope, "test.identity", { key: "same", rewind: false });
|
||||
const rewindList = defineList<string>(scope, "test.identity", { key: "same", rewind: true });
|
||||
const absentKey = defineValue<string>(scope, "test.optional-key", { rewind: false });
|
||||
const undefinedKey = defineValue<string>(scope, "test.optional-key", { key: undefined, rewind: false });
|
||||
const emptyKey = defineValue<string>(scope, "test.optional-key", { key: "", rewind: false });
|
||||
const sessionValue = defineValue<string>({ type: "session" }, "test.identity", { key: "same" });
|
||||
|
||||
await session.commit((tx) => {
|
||||
tx.value(stickyValue).set("sticky-value");
|
||||
tx.value(rewindValue).set("rewind-value");
|
||||
tx.list(stickyList).append("sticky-list");
|
||||
tx.list(rewindList).append("rewind-list");
|
||||
tx.value(absentKey).set("absent");
|
||||
tx.value(emptyKey).set("empty");
|
||||
tx.value(sessionValue).set("session");
|
||||
}, ctx);
|
||||
|
||||
expect(await session.getValue(stickyValue, undefined, ctx)).toBe("sticky-value");
|
||||
expect(await session.getValue(rewindValue, undefined, ctx)).toBe("rewind-value");
|
||||
expect((await session.readList(stickyList, undefined, ctx)).map((item) => item.value)).toEqual(["sticky-list"]);
|
||||
expect((await session.readList(rewindList, undefined, ctx)).map((item) => item.value)).toEqual(["rewind-list"]);
|
||||
expect(await session.getValue(absentKey, undefined, ctx)).toBe("absent");
|
||||
expect(await session.getValue(undefinedKey, undefined, ctx)).toBe("absent");
|
||||
expect(await session.getValue(emptyKey, undefined, ctx)).toBe("empty");
|
||||
expect(await session.getValue(sessionValue, undefined, ctx)).toBe("session");
|
||||
});
|
||||
|
||||
it("stores task snapshots and retires scratch and shared output", async () => {
|
||||
const storage = new MemoryStorage();
|
||||
const session = new Session(storage);
|
||||
const ids = await session.commit((tx) => {
|
||||
const conversation = tx.createConversation({});
|
||||
const producer = tx.task(taskSpec(conversation, { output: { kind: "test.output" } }));
|
||||
const output = defineList<JsonValue>({ type: "shared", id: producer }, "pi.output");
|
||||
tx.list(output).append([["r", { text: "" }]]);
|
||||
const consumer = tx.task(taskSpec(conversation, { output: { id: producer, kind: "test.output" } }));
|
||||
return { conversation, producer, consumer };
|
||||
}, ctx);
|
||||
expect(ids).toEqual({ conversation: 1, producer: 2, consumer: 4 });
|
||||
|
||||
const output = defineList<JsonValue>({ type: "shared", id: ids.producer }, "pi.output");
|
||||
const scratch = defineValue<JsonValue>({ type: "task", taskId: ids.producer }, "test.scratch");
|
||||
await session.commit((tx) => {
|
||||
tx.value(scratch).set({ n: 1 });
|
||||
tx.list(output).append([["s", ["text"], "done"]]);
|
||||
tx.list(output).append([["s", ["status"], "complete"]]);
|
||||
}, ctx);
|
||||
|
||||
await setTask(session, ids.producer, (task) => ({ ...task, status: "running" }));
|
||||
await setTask(session, ids.producer, (task) => ({ ...task, checkpoint: { phase: "working" } }));
|
||||
await setTask(session, ids.producer, (task) => ({
|
||||
...task,
|
||||
status: "terminal",
|
||||
outcome: { status: "completed", result: null },
|
||||
}));
|
||||
await expect(session.getValue(scratch, undefined, ctx)).rejects.toBeInstanceOf(ScratchRetired);
|
||||
expect(await storage.getValue(scratch, undefined, ctx)).toBeUndefined();
|
||||
|
||||
await setTask(session, ids.consumer, (task) => ({ ...task, status: "running" }));
|
||||
const replacement = await session.commit(async (tx) => {
|
||||
const consumer = await tx.getTask(ids.consumer);
|
||||
if (consumer === undefined || consumer.status === "terminal") throw new Error("Consumer is not live");
|
||||
tx.setTask({ ...consumer, status: "terminal", outcome: { status: "completed", result: null } });
|
||||
return tx.task(taskSpec(ids.conversation, { output: { id: ids.producer, kind: "test.output" } }));
|
||||
}, ctx);
|
||||
expect(await session.readList(output, undefined, ctx)).toHaveLength(3);
|
||||
await setTask(session, replacement, (task) => ({ ...task, status: "running" }));
|
||||
await setTask(session, replacement, (task) => ({
|
||||
...task,
|
||||
status: "terminal",
|
||||
outcome: { status: "completed", result: null },
|
||||
}));
|
||||
await expect(session.readList(output, undefined, ctx)).rejects.toBeInstanceOf(SharedOutputRetired);
|
||||
expect(await storage.readList(output, undefined, ctx)).toEqual([]);
|
||||
});
|
||||
|
||||
it("filters and pages tasks", async () => {
|
||||
const session = new Session(new MemoryStorage());
|
||||
const ids = await session.commit((tx) => {
|
||||
const conversation = tx.createConversation({});
|
||||
const otherConversation = tx.createConversation({});
|
||||
return {
|
||||
conversation,
|
||||
otherConversation,
|
||||
first: tx.task(taskSpec(conversation)),
|
||||
second: tx.task(
|
||||
taskSpec(conversation, { kind: "other", background: true, output: { kind: "test.output" } }),
|
||||
),
|
||||
third: tx.task(taskSpec(otherConversation)),
|
||||
};
|
||||
}, ctx);
|
||||
await setTask(session, ids.first, (task) => ({ ...task, abort: true }));
|
||||
await setTask(session, ids.first, (task) => ({ ...task, status: "running" }));
|
||||
await setTask(session, ids.first, (task) => ({
|
||||
...task,
|
||||
status: "terminal",
|
||||
outcome: { status: "aborted", result: null },
|
||||
}));
|
||||
|
||||
const first = await session.scanTasks({ limit: 1 }, ctx);
|
||||
expect(first.items.map(({ id }) => id)).toEqual([ids.first]);
|
||||
expect((await session.scanTasks({ cursor: first.next, limit: 1 }, ctx)).items.map(({ id }) => id)).toEqual([
|
||||
ids.second,
|
||||
]);
|
||||
expect((await session.getTasks([ids.first, 999_999], ctx)).has(ids.first)).toBe(true);
|
||||
expect(
|
||||
(await session.scanTasks({ conversationIds: [ids.conversation], limit: 10 }, ctx)).items.map(({ id }) => id),
|
||||
).toEqual([ids.first, ids.second]);
|
||||
expect((await session.scanTasks({ kind: "other", limit: 10 }, ctx)).items.map(({ id }) => id)).toEqual([
|
||||
ids.second,
|
||||
]);
|
||||
expect((await session.scanTasks({ outputId: ids.second, limit: 10 }, ctx)).items.map(({ id }) => id)).toEqual([
|
||||
ids.second,
|
||||
]);
|
||||
expect(
|
||||
(
|
||||
await session.scanTasks(
|
||||
{
|
||||
conversationIds: [ids.conversation],
|
||||
statuses: ["terminal"],
|
||||
kind: "test.task",
|
||||
abort: true,
|
||||
limit: 10,
|
||||
},
|
||||
ctx,
|
||||
)
|
||||
).items.map(({ id }) => id),
|
||||
).toEqual([ids.first]);
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user