cb5bbbf382
- Eliminated core/ directory (252 files, 23 subdirs → distributed)
- Reduced max nesting from 9 levels to 5 levels
- Promoted tool subdirs to top level: exa/, lsp/, patch/, task/, web/
- Merged web-scrapers/ + web-search/ into web/{scrapers,search}/
- Flattened modes/interactive/ to modes/
- Split execution/ into: ipy/ (python), exec/ (bash), ssh/
- Renamed ipy/python-*.ts to ipy/*.ts (executor, kernel, etc.)
- Flattened cursor/exec-bridge.ts to cursor.ts
- Created logical groupings: config/, session/, extensibility/, export/
- Updated all imports across 500+ files
214 lines
6.0 KiB
TypeScript
214 lines
6.0 KiB
TypeScript
/**
|
|
* Tests for AgentSession concurrent prompt guard.
|
|
*/
|
|
|
|
import { afterEach, beforeEach, describe, expect, it } from "bun:test";
|
|
import { existsSync, mkdirSync, rmSync } from "node:fs";
|
|
import { tmpdir } from "node:os";
|
|
import { join } from "node:path";
|
|
import { Agent } from "@oh-my-pi/pi-agent-core";
|
|
import { type AssistantMessage, type AssistantMessageEvent, EventStream, getModel } from "@oh-my-pi/pi-ai";
|
|
import { nanoid } from "nanoid";
|
|
import { ModelRegistry } from "$c/config/model-registry";
|
|
import { SettingsManager } from "$c/config/settings-manager";
|
|
import { AgentSession } from "$c/session/agent-session";
|
|
import { AuthStorage } from "$c/session/auth-storage";
|
|
import { SessionManager } from "$c/session/session-manager";
|
|
|
|
// Mock stream that mimics AssistantMessageEventStream
|
|
class MockAssistantStream extends EventStream<AssistantMessageEvent, AssistantMessage> {
|
|
constructor() {
|
|
super(
|
|
(event) => event.type === "done" || event.type === "error",
|
|
(event) => {
|
|
if (event.type === "done") return event.message;
|
|
if (event.type === "error") return event.error;
|
|
throw new Error("Unexpected event type");
|
|
},
|
|
);
|
|
}
|
|
}
|
|
|
|
function createAssistantMessage(text: string): AssistantMessage {
|
|
return {
|
|
role: "assistant",
|
|
content: [{ type: "text", text }],
|
|
api: "anthropic-messages",
|
|
provider: "anthropic",
|
|
model: "mock",
|
|
usage: {
|
|
input: 0,
|
|
output: 0,
|
|
cacheRead: 0,
|
|
cacheWrite: 0,
|
|
totalTokens: 0,
|
|
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
|
|
},
|
|
stopReason: "stop",
|
|
timestamp: Date.now(),
|
|
};
|
|
}
|
|
|
|
describe("AgentSession concurrent prompt guard", () => {
|
|
let session: AgentSession;
|
|
let tempDir: string;
|
|
|
|
beforeEach(() => {
|
|
tempDir = join(tmpdir(), `pi-concurrent-test-${nanoid()}`);
|
|
mkdirSync(tempDir, { recursive: true });
|
|
});
|
|
|
|
afterEach(async () => {
|
|
if (session) {
|
|
session.dispose();
|
|
}
|
|
if (tempDir && existsSync(tempDir)) {
|
|
rmSync(tempDir, { recursive: true });
|
|
}
|
|
});
|
|
|
|
async function createSession() {
|
|
const model = getModel("anthropic", "claude-sonnet-4-5")!;
|
|
let abortSignal: AbortSignal | undefined;
|
|
|
|
// Use a stream function that responds to abort
|
|
const agent = new Agent({
|
|
getApiKey: () => "test-key",
|
|
initialState: {
|
|
model,
|
|
systemPrompt: "Test",
|
|
tools: [],
|
|
},
|
|
streamFn: (_model, _context, options) => {
|
|
abortSignal = options?.signal;
|
|
const stream = new MockAssistantStream();
|
|
queueMicrotask(() => {
|
|
stream.push({ type: "start", partial: createAssistantMessage("") });
|
|
const checkAbort = () => {
|
|
if (abortSignal?.aborted) {
|
|
stream.push({ type: "error", reason: "aborted", error: createAssistantMessage("Aborted") });
|
|
} else {
|
|
setTimeout(checkAbort, 5);
|
|
}
|
|
};
|
|
checkAbort();
|
|
});
|
|
return stream;
|
|
},
|
|
});
|
|
|
|
const sessionManager = SessionManager.inMemory();
|
|
const settingsManager = await SettingsManager.create(tempDir, tempDir);
|
|
const authStorage = new AuthStorage(join(tempDir, "auth.json"));
|
|
const modelRegistry = new ModelRegistry(authStorage, tempDir);
|
|
// Set a runtime API key so validation passes
|
|
authStorage.setRuntimeApiKey("anthropic", "test-key");
|
|
|
|
session = new AgentSession({
|
|
agent,
|
|
sessionManager,
|
|
settingsManager,
|
|
modelRegistry,
|
|
});
|
|
|
|
return session;
|
|
}
|
|
|
|
it("should throw when prompt() called while streaming", async () => {
|
|
await createSession();
|
|
|
|
// Start first prompt (don't await, it will block until abort)
|
|
const firstPrompt = session.prompt("First message");
|
|
|
|
// Wait a tick for isStreaming to be set
|
|
await Bun.sleep(10);
|
|
|
|
// Verify we're streaming
|
|
expect(session.isStreaming).toBe(true);
|
|
|
|
// Second prompt should reject
|
|
await expect(session.prompt("Second message")).rejects.toThrow(
|
|
"Agent is already processing. Specify streamingBehavior ('steer' or 'followUp') to queue the message.",
|
|
);
|
|
|
|
// Cleanup
|
|
await session.abort();
|
|
await firstPrompt.catch(() => {}); // Ignore abort error
|
|
});
|
|
|
|
it("should allow steer() while streaming", async () => {
|
|
await createSession();
|
|
|
|
// Start first prompt
|
|
const firstPrompt = session.prompt("First message");
|
|
await Bun.sleep(10);
|
|
|
|
// steer should work while streaming
|
|
expect(() => session.steer("Steering message")).not.toThrow();
|
|
expect(session.queuedMessageCount).toBe(1);
|
|
|
|
// Cleanup
|
|
await session.abort();
|
|
await firstPrompt.catch(() => {});
|
|
});
|
|
|
|
it("should allow followUp() while streaming", async () => {
|
|
await createSession();
|
|
|
|
// Start first prompt
|
|
const firstPrompt = session.prompt("First message");
|
|
await Bun.sleep(10);
|
|
|
|
// followUp should work while streaming
|
|
expect(() => session.followUp("Follow-up message")).not.toThrow();
|
|
expect(session.queuedMessageCount).toBe(1);
|
|
|
|
// Cleanup
|
|
await session.abort();
|
|
await firstPrompt.catch(() => {});
|
|
});
|
|
|
|
it("should allow prompt() after previous completes", async () => {
|
|
// Create session with a stream that completes immediately
|
|
const model = getModel("anthropic", "claude-sonnet-4-5")!;
|
|
const agent = new Agent({
|
|
getApiKey: () => "test-key",
|
|
initialState: {
|
|
model,
|
|
systemPrompt: "Test",
|
|
tools: [],
|
|
},
|
|
streamFn: () => {
|
|
const stream = new MockAssistantStream();
|
|
queueMicrotask(() => {
|
|
stream.push({ type: "start", partial: createAssistantMessage("") });
|
|
stream.push({ type: "done", reason: "stop", message: createAssistantMessage("Done") });
|
|
});
|
|
return stream;
|
|
},
|
|
});
|
|
|
|
const sessionManager = SessionManager.inMemory();
|
|
const settingsManager = await SettingsManager.create(tempDir, tempDir);
|
|
const authStorage = new AuthStorage(join(tempDir, "auth.json"));
|
|
const modelRegistry = new ModelRegistry(authStorage, tempDir);
|
|
authStorage.setRuntimeApiKey("anthropic", "test-key");
|
|
|
|
session = new AgentSession({
|
|
agent,
|
|
sessionManager,
|
|
settingsManager,
|
|
modelRegistry,
|
|
});
|
|
|
|
// First prompt completes
|
|
await session.prompt("First message");
|
|
|
|
// Should not be streaming anymore
|
|
expect(session.isStreaming).toBe(false);
|
|
|
|
// Second prompt should work
|
|
await expect(session.prompt("Second message")).resolves.toBeUndefined();
|
|
});
|
|
});
|