149 lines
6.0 KiB
TypeScript
149 lines
6.0 KiB
TypeScript
import type { AcpConfig } from "../config.js";
|
|
import { chatKeyFor, type DurableSessionStore, type SessionBinding } from "../core/durable-session-store.js";
|
|
import type { ResolvedBot } from "../roles/role-registry.js";
|
|
import type { ConversationRequest, ConversationResponse, ConversationRuntime, RuntimeStats } from "./types.js";
|
|
import { AcpSessionRestoreError, AcpWorker } from "./worker.js";
|
|
|
|
export class AcpSessionManager implements ConversationRuntime {
|
|
private readonly workers = new Map<string, AcpWorker>();
|
|
private readonly inFlight = new Map<string, AcpWorker>();
|
|
private readonly sweeper: NodeJS.Timeout;
|
|
private crashes = 0;
|
|
private shuttingDown = false;
|
|
|
|
constructor(
|
|
private readonly config: AcpConfig,
|
|
private readonly bot: ResolvedBot,
|
|
private readonly store: DurableSessionStore
|
|
) {
|
|
this.sweeper = setInterval(() => void this.sweep(), config.sweepIntervalMs);
|
|
this.sweeper.unref();
|
|
}
|
|
|
|
async prompt(request: ConversationRequest): Promise<ConversationResponse> {
|
|
if (this.shuttingDown) throw new Error("ACP runtime is shutting down");
|
|
const chatKey = chatKeyFor(request.platform, request.chatId);
|
|
let binding = this.store.getBinding(chatKey);
|
|
if (binding && (binding.botFingerprint !== this.bot.fingerprint || binding.agentId !== this.bot.agent.id || binding.workspace !== this.bot.workspace)) {
|
|
await this.dropBinding(binding);
|
|
binding = undefined;
|
|
}
|
|
|
|
let worker: AcpWorker;
|
|
try { worker = await this.acquireWorker(binding); }
|
|
catch (error) {
|
|
if (!(error instanceof AcpSessionRestoreError) || !binding) throw error;
|
|
await this.store.deleteBinding(chatKey);
|
|
binding = undefined;
|
|
worker = await this.acquireWorker();
|
|
}
|
|
this.inFlight.set(chatKey, worker);
|
|
try {
|
|
if (!binding) {
|
|
const now = Date.now();
|
|
await worker.prompt(this.bot.bootstrap);
|
|
binding = {
|
|
chatKey, agentId: this.bot.agent.id, nativeSessionId: worker.nativeSessionId!,
|
|
workspace: this.bot.workspace, botFingerprint: this.bot.fingerprint, createdAt: now, updatedAt: now
|
|
};
|
|
await this.store.setBinding(binding);
|
|
}
|
|
const text = await worker.prompt(request.text);
|
|
await this.store.touchBinding(chatKey);
|
|
return { text: text || "(ACP agent returned no text)", botId: this.bot.id, agentId: this.bot.agent.id };
|
|
} catch (error) {
|
|
if (worker.nativeSessionId) this.workers.delete(workerKey(this.bot.agent.id, worker.nativeSessionId));
|
|
await worker.terminate();
|
|
throw error;
|
|
} finally {
|
|
if (this.inFlight.get(chatKey) === worker) this.inFlight.delete(chatKey);
|
|
}
|
|
}
|
|
|
|
async cancel(platform: string, chatId: string): Promise<boolean> {
|
|
const worker = this.inFlight.get(chatKeyFor(platform, chatId));
|
|
return worker ? worker.cancel() : false;
|
|
}
|
|
|
|
async reset(platform: string, chatId: string): Promise<void> {
|
|
const chatKey = chatKeyFor(platform, chatId);
|
|
await this.cancel(platform, chatId);
|
|
const binding = await this.store.deleteBinding(chatKey);
|
|
if (binding) await this.stopWorker(binding.agentId, binding.nativeSessionId);
|
|
}
|
|
|
|
status(platform: string, chatId: string): Record<string, string | number | boolean> {
|
|
const chatKey = chatKeyFor(platform, chatId);
|
|
const binding = this.store.getBinding(chatKey);
|
|
return {
|
|
bot: this.bot.id,
|
|
agent: this.bot.agent.id,
|
|
workspace: this.bot.workspace,
|
|
persisted: Boolean(binding),
|
|
running: this.inFlight.has(chatKey)
|
|
};
|
|
}
|
|
|
|
stats(): RuntimeStats {
|
|
return { activeWorkers: this.workers.size, inFlight: this.inFlight.size, crashes: this.crashes, persistedBindings: this.store.stats().bindings };
|
|
}
|
|
|
|
async shutdown(): Promise<void> {
|
|
if (this.shuttingDown) return;
|
|
this.shuttingDown = true;
|
|
clearInterval(this.sweeper);
|
|
await Promise.all([...this.inFlight.values()].map((worker) => worker.cancel().catch(() => false)));
|
|
await Promise.all([...this.workers.values()].map((worker) => worker.terminate()));
|
|
this.workers.clear();
|
|
this.inFlight.clear();
|
|
}
|
|
|
|
private async acquireWorker(binding?: SessionBinding): Promise<AcpWorker> {
|
|
if (binding) {
|
|
const existing = this.workers.get(workerKey(binding.agentId, binding.nativeSessionId));
|
|
if (existing) return existing;
|
|
}
|
|
await this.ensureCapacity();
|
|
const worker = new AcpWorker(this.bot, this.config, (crashed, error) => {
|
|
this.crashes++;
|
|
if (crashed.nativeSessionId) this.workers.delete(workerKey(this.bot.agent.id, crashed.nativeSessionId));
|
|
console.error(`ACP worker crash: ${error.message}`);
|
|
});
|
|
const nativeSessionId = await worker.start(binding?.nativeSessionId);
|
|
this.workers.set(workerKey(this.bot.agent.id, nativeSessionId), worker);
|
|
return worker;
|
|
}
|
|
|
|
private async ensureCapacity(): Promise<void> {
|
|
if (this.workers.size < this.config.maxProcesses) return;
|
|
const candidate = [...this.workers.entries()].filter(([, worker]) => !worker.inFlight).sort((a, b) => a[1].lastUsedAt - b[1].lastUsedAt)[0];
|
|
if (!candidate) throw new Error(`ACP worker limit reached (${this.config.maxProcesses})`);
|
|
this.workers.delete(candidate[0]);
|
|
await candidate[1].terminate();
|
|
}
|
|
|
|
private async sweep(): Promise<void> {
|
|
const cutoff = Date.now() - this.config.idleTimeoutMs;
|
|
const expired = [...this.workers.entries()].filter(([, worker]) => !worker.inFlight && worker.lastUsedAt < cutoff);
|
|
for (const [key, worker] of expired) {
|
|
this.workers.delete(key);
|
|
await worker.terminate();
|
|
}
|
|
}
|
|
|
|
private async dropBinding(binding: SessionBinding): Promise<void> {
|
|
await this.store.deleteBinding(binding.chatKey);
|
|
await this.stopWorker(binding.agentId, binding.nativeSessionId);
|
|
}
|
|
|
|
private async stopWorker(agentId: string, nativeSessionId: string): Promise<void> {
|
|
const key = workerKey(agentId, nativeSessionId);
|
|
const worker = this.workers.get(key);
|
|
if (!worker) return;
|
|
this.workers.delete(key);
|
|
await worker.terminate();
|
|
}
|
|
}
|
|
|
|
function workerKey(agentId: string, nativeSessionId: string): string { return `${agentId}:${nativeSessionId}`; }
|