Add ACP-backed role sessions
This commit is contained in:
+62
-75
@@ -1,112 +1,99 @@
|
||||
import type { ConversationRuntime } from "../acp/types.js";
|
||||
import type { GatewayPolicy } from "../config.js";
|
||||
import { AgentRegistry } from "../agents/agent-registry.js";
|
||||
import type { RoleRegistry } from "../roles/role-registry.js";
|
||||
import type { PlatformAdapter } from "./adapter.js";
|
||||
import { CommandRouter } from "./command-router.js";
|
||||
import { SessionStore, sessionIdFor } from "./session-store.js";
|
||||
import { CommandRouter, type ParsedCommand } from "./command-router.js";
|
||||
import { chatKeyFor } from "./durable-session-store.js";
|
||||
import type { IncomingMessage } from "./types.js";
|
||||
|
||||
export interface GatewayResult {
|
||||
ok: boolean;
|
||||
reply?: string;
|
||||
ignored?: boolean;
|
||||
error?: string;
|
||||
}
|
||||
export interface GatewayResult { ok: boolean; reply?: string; ignored?: boolean; error?: string }
|
||||
|
||||
export class Gateway {
|
||||
private readonly locks = new Map<string, Promise<void>>();
|
||||
readonly commandRouter: CommandRouter;
|
||||
readonly commandRouter = new CommandRouter();
|
||||
|
||||
constructor(
|
||||
private readonly policy: GatewayPolicy,
|
||||
private readonly agents: AgentRegistry,
|
||||
private readonly sessions: SessionStore
|
||||
) {
|
||||
this.commandRouter = new CommandRouter(agents, sessions);
|
||||
}
|
||||
constructor(private readonly policy: GatewayPolicy, private readonly runtime: ConversationRuntime, private readonly roles: RoleRegistry) {}
|
||||
|
||||
async receive(message: IncomingMessage, adapter: PlatformAdapter, options: { synchronous?: boolean } = {}): Promise<GatewayResult> {
|
||||
const policyError = this.checkPolicy(message);
|
||||
if (policyError) return { ok: true, ignored: true, error: policyError };
|
||||
if (policyError) {
|
||||
console.log(`Message ignored by policy: ${policyError} (${message.platform} ${message.chatId} ${message.userId})`);
|
||||
return { ok: true, ignored: true, error: policyError };
|
||||
}
|
||||
|
||||
const sessionId = sessionIdFor(message.platform, message.chatId);
|
||||
return this.withChatLock(sessionId, async () => {
|
||||
const command = this.commandRouter.route(message, sessionId);
|
||||
if (command.handled) {
|
||||
const reply = command.text || "";
|
||||
await adapter.sendMessage({
|
||||
target: { platform: message.platform, chatId: message.chatId, userId: message.userId, raw: message.raw },
|
||||
text: reply,
|
||||
replyTo: message.messageId
|
||||
});
|
||||
return { ok: true, reply };
|
||||
}
|
||||
|
||||
const session = this.sessions.getSession(sessionId);
|
||||
const agent = this.agents.get(session.selectedAgent);
|
||||
const command = this.commandRouter.parse(message.text);
|
||||
if (command?.kind === "cancel") return this.reply(message, adapter, await this.cancelText(message), options);
|
||||
if (command?.kind === "new") await this.runtime.cancel(message.platform, message.chatId);
|
||||
const chatKey = chatKeyFor(message.platform, message.chatId);
|
||||
return this.withChatLock(chatKey, async () => {
|
||||
try {
|
||||
const response = await agent.run({
|
||||
input: message.text,
|
||||
sessionId,
|
||||
platform: message.platform,
|
||||
chatId: message.chatId,
|
||||
userId: message.userId,
|
||||
messageId: message.messageId
|
||||
});
|
||||
await adapter.sendMessage({
|
||||
target: { platform: message.platform, chatId: message.chatId, userId: message.userId, raw: message.raw },
|
||||
text: response.text,
|
||||
replyTo: message.messageId
|
||||
});
|
||||
return { ok: true, reply: response.text };
|
||||
const reply = command ? await this.executeCommand(command, message) : (await this.runtime.prompt({
|
||||
platform: message.platform, chatId: message.chatId, userId: message.userId, text: message.text, messageId: message.messageId
|
||||
})).text;
|
||||
return this.reply(message, adapter, reply, options);
|
||||
} catch (error) {
|
||||
const errorText = error instanceof Error ? error.message : String(error);
|
||||
const reply = `Agent error: ${errorText}`;
|
||||
if (!options.synchronous) {
|
||||
await adapter.sendMessage({
|
||||
target: { platform: message.platform, chatId: message.chatId, userId: message.userId, raw: message.raw },
|
||||
text: reply,
|
||||
replyTo: message.messageId
|
||||
});
|
||||
}
|
||||
if (!options.synchronous) await this.send(message, adapter, reply);
|
||||
return { ok: false, error: errorText, reply };
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
stats(): { sessions: number; lockedChats: number } {
|
||||
return { ...this.sessions.stats(), lockedChats: this.locks.size };
|
||||
stats(): ReturnType<ConversationRuntime["stats"]> & { lockedChats: number } { return { ...this.runtime.stats(), lockedChats: this.locks.size }; }
|
||||
|
||||
private async executeCommand(command: ParsedCommand, message: IncomingMessage): Promise<string> {
|
||||
const prefix = command.deprecatedAlias ? "Deprecated alias; use /role or /roles.\n" : "";
|
||||
switch (command.kind) {
|
||||
case "help": return ["Commands:", "/roles", "/role <id>", "/status", "/cancel", "/new", "/help"].join("\n");
|
||||
case "roles": return `${prefix}Available roles: ${this.roles.list().join(", ")}`;
|
||||
case "role":
|
||||
if (!command.argument) return `${prefix}Usage: /role <id>`;
|
||||
if (!this.roles.has(command.argument)) return `${prefix}Unknown role: ${command.argument}`;
|
||||
await this.runtime.cancel(message.platform, message.chatId);
|
||||
await this.runtime.selectRole(message.platform, message.chatId, command.argument);
|
||||
return `${prefix}Selected role: ${command.argument}`;
|
||||
case "status": {
|
||||
const status = this.runtime.status(message.platform, message.chatId);
|
||||
return `OK\n${Object.entries(status).map(([key, value]) => `${key}=${value}`).join("\n")}`;
|
||||
}
|
||||
case "new":
|
||||
await this.runtime.reset(message.platform, message.chatId);
|
||||
return "Started a new native ACP session for this chat and role.";
|
||||
case "cancel": return this.cancelText(message);
|
||||
}
|
||||
}
|
||||
|
||||
private async cancelText(message: IncomingMessage): Promise<string> {
|
||||
return await this.runtime.cancel(message.platform, message.chatId) ? "Cancellation requested." : "No active turn to cancel.";
|
||||
}
|
||||
|
||||
private async reply(message: IncomingMessage, adapter: PlatformAdapter, reply: string, options: { synchronous?: boolean }): Promise<GatewayResult> {
|
||||
if (!options.synchronous) await this.send(message, adapter, reply);
|
||||
return { ok: true, reply };
|
||||
}
|
||||
|
||||
private send(message: IncomingMessage, adapter: PlatformAdapter, text: string): Promise<void> {
|
||||
return adapter.sendMessage({ target: { platform: message.platform, chatId: message.chatId, userId: message.userId, raw: message.raw }, text, replyTo: message.messageId });
|
||||
}
|
||||
|
||||
private checkPolicy(message: IncomingMessage): string | undefined {
|
||||
if (this.policy.allowedUsers.length > 0 && !this.policy.allowedUsers.includes(message.userId)) {
|
||||
return `User not allowed: ${message.userId}`;
|
||||
}
|
||||
if (this.policy.allowedChats.length > 0 && !this.policy.allowedChats.includes(message.chatId)) {
|
||||
return `Chat not allowed: ${message.chatId}`;
|
||||
}
|
||||
if (this.policy.requireMentionInGroup && message.isGroup && !message.mentionsBot) {
|
||||
return "Mention required in group chat";
|
||||
}
|
||||
if (this.policy.allowedUsers.length > 0 && !this.policy.allowedUsers.includes(message.userId)) return `User not allowed: ${message.userId}`;
|
||||
if (this.policy.allowedChats.length > 0 && !this.policy.allowedChats.includes(message.chatId)) return `Chat not allowed: ${message.chatId}`;
|
||||
if (this.policy.requireMentionInGroup && message.isGroup && !message.mentionsBot) return "Mention required in group chat";
|
||||
return undefined;
|
||||
}
|
||||
|
||||
private async withChatLock<T>(key: string, fn: () => Promise<T>): Promise<T> {
|
||||
const previous = this.locks.get(key) || Promise.resolve();
|
||||
let release!: () => void;
|
||||
const current = new Promise<void>((resolve) => {
|
||||
release = resolve;
|
||||
});
|
||||
const current = new Promise<void>((resolve) => { release = resolve; });
|
||||
const queued = previous.then(() => current);
|
||||
this.locks.set(key, queued);
|
||||
|
||||
await previous;
|
||||
try {
|
||||
return await fn();
|
||||
} finally {
|
||||
try { return await fn(); } finally {
|
||||
release();
|
||||
if (this.locks.get(key) === queued) {
|
||||
this.locks.delete(key);
|
||||
}
|
||||
if (this.locks.get(key) === queued) this.locks.delete(key);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user