File size: 6,356 Bytes
51c026d | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 | import type { JsonValue, ServiceCall } from "@earendil-works/chord";
import type { Context, Session, SessionMetadata } from "@earendil-works/pi-agent-core";
import { BACKGROUND_CONTEXT, MemorySessionRepo } from "@earendil-works/pi-agent-core";
import { SessionAmbiguousError, SessionNotFoundError } from "../errors.ts";
import type { RoutedServerServiceHost, RoutedSessionHandle, ServerHost } from "../types.ts";
export class Deferred<T> {
readonly promise: Promise<T>;
private resolvePromise!: (value: T) => void;
constructor() {
this.promise = new Promise<T>((resolve) => {
this.resolvePromise = resolve;
});
}
resolve(value: T): void {
this.resolvePromise(value);
}
}
interface OpenGate {
entered: Deferred<void>;
release: Deferred<void>;
}
export class TestHarness {
readonly session: Session;
readonly closed = new Deferred<void>();
readonly #termination = new Deferred<Error | undefined>();
readonly terminated = this.#termination.promise;
attachedClients = 0;
attachmentReleaseCount = 0;
closeCount = 0;
readonly serviceCalls: ServiceCall[] = [];
failAttachmentRelease?: Error;
failClose?: Error;
nextServiceError?: Error;
nextServiceResult: JsonValue | undefined = { ok: true };
private nextCloseGate?: OpenGate;
private nextServiceGate?: OpenGate;
constructor(session: Session) {
this.session = session;
}
attachClient(_context: Context): {
invokeService: TestHarness["invokeService"];
release(context: Context): void;
} {
this.attachedClients += 1;
let released = false;
return {
invokeService: (call) => this.invokeService(call),
release: (_context) => {
if (released) return;
this.attachmentReleaseCount += 1;
if (this.failAttachmentRelease) throw this.failAttachmentRelease;
released = true;
this.attachedClients -= 1;
},
};
}
async invokeService(call: ServiceCall): Promise<JsonValue | undefined> {
this.serviceCalls.push(call);
if (this.nextServiceError) {
const error = this.nextServiceError;
this.nextServiceError = undefined;
throw error;
}
const gate = this.nextServiceGate;
if (gate) {
this.nextServiceGate = undefined;
gate.entered.resolve(undefined);
await gate.release.promise;
}
const result = this.nextServiceResult;
this.nextServiceResult = { ok: true };
return result;
}
async close(context: Context): Promise<void> {
this.closeCount += 1;
const gate = this.nextCloseGate;
if (gate) {
this.nextCloseGate = undefined;
gate.entered.resolve(undefined);
await gate.release.promise;
}
if (this.failClose) {
const error = this.failClose;
this.failClose = undefined;
throw error;
}
await this.session.close(context);
this.closed.resolve(undefined);
this.#termination.resolve(undefined);
}
async terminate(error: Error): Promise<void> {
await this.session.close(BACKGROUND_CONTEXT);
this.#termination.resolve(error);
}
gateNextClose(): OpenGate {
const gate = { entered: new Deferred<void>(), release: new Deferred<void>() };
this.nextCloseGate = gate;
return gate;
}
gateNextServiceCall(): OpenGate {
const gate = { entered: new Deferred<void>(), release: new Deferred<void>() };
this.nextServiceGate = gate;
return gate;
}
}
export function createTestServerServices(): RoutedServerServiceHost {
return {
attachClient(presentation) {
return {
async invokeService(call, _publish, context) {
if (
call.instance === undefined &&
call.serviceId === "pi.session-management" &&
call.member === "attach" &&
call.args.length === 1 &&
typeof call.args[0] === "string"
) {
await presentation.attachSession(call.args[0], context);
return null;
}
if (
call.instance === undefined &&
call.serviceId === "pi.session-management" &&
call.member === "detach" &&
call.args.length === 0
) {
await presentation.detachSession(context);
return null;
}
throw new Error(`Unsupported test server service ${call.serviceId}.${call.member}`);
},
release() {},
};
},
};
}
export class TestServerHost implements ServerHost {
readonly serverServices = createTestServerServices();
readonly repo = new MemorySessionRepo({ now: () => 1 });
readonly harnesses = new Map<string, TestHarness[]>();
openSessionCount = 0;
nextOpenSessionError?: Error;
nextHarnessCloseError?: Error;
private nextOpenSessionGate?: OpenGate;
async resolveSession(sessionId: string, context: Context): Promise<SessionMetadata> {
const matches = (await this.repo.list(undefined, context)).filter(({ id }) => id === sessionId);
if (matches.length === 0) throw new SessionNotFoundError(`Unknown session: ${sessionId}`);
if (matches.length > 1) throw new SessionAmbiguousError();
return matches[0]!;
}
async openSession(metadata: SessionMetadata, context: Context): Promise<RoutedSessionHandle> {
this.openSessionCount += 1;
const gate = this.nextOpenSessionGate;
if (gate) {
this.nextOpenSessionGate = undefined;
gate.entered.resolve(undefined);
await gate.release.promise;
}
const session = await this.repo.open(metadata, context);
try {
if (this.nextOpenSessionError) {
const error = this.nextOpenSessionError;
this.nextOpenSessionError = undefined;
throw error;
}
const harness = new TestHarness(session);
if (this.nextHarnessCloseError) {
harness.failClose = this.nextHarnessCloseError;
this.nextHarnessCloseError = undefined;
}
const harnesses = this.harnesses.get(metadata.id) ?? [];
harnesses.push(harness);
this.harnesses.set(metadata.id, harnesses);
return harness;
} catch (error) {
await session.close(context);
throw error;
}
}
async seed(id = "session-1", parentSessionId?: string): Promise<SessionMetadata> {
const session = await this.repo.create({ id, parentSessionId }, BACKGROUND_CONTEXT);
const metadata = session.metadata;
await session.close(BACKGROUND_CONTEXT);
return metadata;
}
gateNextOpenSession(): OpenGate {
const gate = { entered: new Deferred<void>(), release: new Deferred<void>() };
this.nextOpenSessionGate = gate;
return gate;
}
latestHarness(id: string): TestHarness {
const harnesses = this.harnesses.get(id);
if (!harnesses?.length) throw new Error(`No harness for ${id}`);
return harnesses.at(-1)!;
}
}
|