feat: expose bot task progress
This commit is contained in:
+21
-13
@@ -7,6 +7,7 @@ import type { PermissionPolicy } from "../config.js";
|
||||
export interface AcpClientOptions {
|
||||
initializeTimeoutMs: number;
|
||||
policy: PermissionPolicy;
|
||||
onSessionActivity?: () => void;
|
||||
}
|
||||
|
||||
export class AcpClient {
|
||||
@@ -47,19 +48,24 @@ export class AcpClient {
|
||||
async resumeSession(sessionId: string, cwd: string): Promise<void> {
|
||||
this.collecting = false;
|
||||
this.chunks = [];
|
||||
if (this.capabilities.sessionCapabilities?.resume) {
|
||||
try {
|
||||
await this.connection.agent.request(acp.methods.agent.session.resume, { sessionId, cwd, mcpServers: [] });
|
||||
} catch (error) {
|
||||
if (!this.capabilities.loadSession) throw error;
|
||||
await this.connection.agent.request(acp.methods.agent.session.load, { sessionId, cwd, mcpServers: [] });
|
||||
}
|
||||
} else if (this.capabilities.loadSession) {
|
||||
await this.connection.agent.request(acp.methods.agent.session.load, { sessionId, cwd, mcpServers: [] });
|
||||
} else {
|
||||
throw new Error("ACP backend cannot resume or load sessions");
|
||||
}
|
||||
this.activeSessionId = sessionId;
|
||||
try {
|
||||
if (this.capabilities.sessionCapabilities?.resume) {
|
||||
try {
|
||||
await this.connection.agent.request(acp.methods.agent.session.resume, { sessionId, cwd, mcpServers: [] });
|
||||
} catch (error) {
|
||||
if (!this.capabilities.loadSession) throw error;
|
||||
await this.connection.agent.request(acp.methods.agent.session.load, { sessionId, cwd, mcpServers: [] });
|
||||
}
|
||||
} else if (this.capabilities.loadSession) {
|
||||
await this.connection.agent.request(acp.methods.agent.session.load, { sessionId, cwd, mcpServers: [] });
|
||||
} else {
|
||||
throw new Error("ACP backend cannot resume or load sessions");
|
||||
}
|
||||
} catch (error) {
|
||||
this.activeSessionId = undefined;
|
||||
throw error;
|
||||
}
|
||||
this.chunks = [];
|
||||
}
|
||||
|
||||
@@ -91,7 +97,9 @@ export class AcpClient {
|
||||
close(error?: unknown): void { this.connection.close(error); }
|
||||
|
||||
private handleUpdate(notification: SessionNotification): void {
|
||||
if (!this.collecting || notification.sessionId !== this.activeSessionId) return;
|
||||
if (notification.sessionId !== this.activeSessionId) return;
|
||||
this.options.onSessionActivity?.();
|
||||
if (!this.collecting) return;
|
||||
const update = notification.update;
|
||||
if (update.sessionUpdate === "agent_message_chunk" && update.content.type === "text") this.chunks.push(update.content.text);
|
||||
}
|
||||
|
||||
@@ -41,7 +41,7 @@ export class AcpSessionManager implements ConversationRuntime {
|
||||
try {
|
||||
if (!binding) {
|
||||
const now = Date.now();
|
||||
await worker.prompt(this.bot.bootstrap);
|
||||
await worker.prompt(this.bot.bootstrap, "initializing");
|
||||
binding = {
|
||||
chatKey, agentId: this.bot.agent.id, nativeSessionId: worker.nativeSessionId!,
|
||||
workspace: this.bot.workspace, botFingerprint: this.bot.fingerprint, createdAt: now, updatedAt: now
|
||||
@@ -75,12 +75,17 @@ export class AcpSessionManager implements ConversationRuntime {
|
||||
status(platform: string, chatId: string): Record<string, string | number | boolean> {
|
||||
const chatKey = chatKeyFor(platform, chatId);
|
||||
const binding = this.store.getBinding(chatKey);
|
||||
const worker = this.inFlight.get(chatKey);
|
||||
const now = Date.now();
|
||||
return {
|
||||
bot: this.bot.id,
|
||||
agent: this.bot.agent.id,
|
||||
workspace: this.bot.workspace,
|
||||
persisted: Boolean(binding),
|
||||
running: this.inFlight.has(chatKey)
|
||||
running: Boolean(worker),
|
||||
phase: worker?.phase || "idle",
|
||||
runningSeconds: worker?.turnStartedAt ? Math.max(0, Math.floor((now - worker.turnStartedAt) / 1000)) : 0,
|
||||
idleSeconds: worker?.lastActivityAt ? Math.max(0, Math.floor((now - worker.lastActivityAt) / 1000)) : 0
|
||||
};
|
||||
}
|
||||
|
||||
|
||||
+19
-3
@@ -5,6 +5,8 @@ import { AcpClient } from "./client.js";
|
||||
|
||||
export class AcpSessionRestoreError extends Error {}
|
||||
|
||||
export type AcpWorkerPhase = "idle" | "initializing" | "processing";
|
||||
|
||||
export class AcpWorker {
|
||||
private child?: ChildProcessWithoutNullStreams;
|
||||
private client?: AcpClient;
|
||||
@@ -14,6 +16,9 @@ export class AcpWorker {
|
||||
private stderrBytes = 0;
|
||||
nativeSessionId?: string;
|
||||
lastUsedAt = Date.now();
|
||||
turnStartedAt?: number;
|
||||
lastActivityAt?: number;
|
||||
phase: AcpWorkerPhase = "idle";
|
||||
inFlight = false;
|
||||
|
||||
constructor(
|
||||
@@ -35,7 +40,11 @@ export class AcpWorker {
|
||||
this.exited = true;
|
||||
if (!this.stopping) this.crashed(new Error(`ACP worker exited (code=${code ?? "null"}, signal=${signal ?? "null"})`));
|
||||
});
|
||||
this.client = new AcpClient(this.child, { initializeTimeoutMs: this.config.initializeTimeoutMs, policy: this.bot.permissions });
|
||||
this.client = new AcpClient(this.child, {
|
||||
initializeTimeoutMs: this.config.initializeTimeoutMs,
|
||||
policy: this.bot.permissions,
|
||||
onSessionActivity: () => { this.lastActivityAt = Date.now(); }
|
||||
});
|
||||
try {
|
||||
await this.client.initialize();
|
||||
if (nativeSessionId) {
|
||||
@@ -52,11 +61,15 @@ export class AcpWorker {
|
||||
}
|
||||
}
|
||||
|
||||
async prompt(text: string): Promise<string> {
|
||||
async prompt(text: string, phase: Exclude<AcpWorkerPhase, "idle"> = "processing"): Promise<string> {
|
||||
if (!this.client || !this.nativeSessionId || this.exited) throw new Error("ACP worker is not available");
|
||||
if (this.inFlight) throw new Error("ACP worker already has an in-flight turn");
|
||||
const now = Date.now();
|
||||
this.inFlight = true;
|
||||
this.lastUsedAt = Date.now();
|
||||
this.phase = phase;
|
||||
this.turnStartedAt = now;
|
||||
this.lastActivityAt = now;
|
||||
this.lastUsedAt = now;
|
||||
this.abort = new AbortController();
|
||||
let timeout: NodeJS.Timeout | undefined;
|
||||
const timeoutPromise = new Promise<never>((_resolve, reject) => {
|
||||
@@ -72,6 +85,9 @@ export class AcpWorker {
|
||||
} finally {
|
||||
if (timeout) clearTimeout(timeout);
|
||||
this.inFlight = false;
|
||||
this.phase = "idle";
|
||||
this.turnStartedAt = undefined;
|
||||
this.lastActivityAt = undefined;
|
||||
this.abort = undefined;
|
||||
this.lastUsedAt = Date.now();
|
||||
}
|
||||
|
||||
+59
-12
@@ -7,8 +7,15 @@ import type { IncomingMessage } from "./types.js";
|
||||
|
||||
export interface GatewayResult { ok: boolean; reply?: string; ignored?: boolean; error?: string }
|
||||
|
||||
interface ChatWorkState {
|
||||
running: boolean;
|
||||
queued: number;
|
||||
startedAt?: number;
|
||||
}
|
||||
|
||||
export class Gateway {
|
||||
private readonly locks = new Map<string, Promise<void>>();
|
||||
private readonly chatWork = new Map<string, ChatWorkState>();
|
||||
readonly commandRouter = new CommandRouter();
|
||||
|
||||
constructor(private readonly policy: GatewayPolicy, private readonly runtime: ConversationRuntime) {}
|
||||
@@ -21,19 +28,28 @@ export class Gateway {
|
||||
}
|
||||
|
||||
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);
|
||||
if (command?.kind === "help" || command?.kind === "status" || command?.kind === "cancel") {
|
||||
return this.reply(message, adapter, await this.executeCommand(command, message), options);
|
||||
}
|
||||
|
||||
const chatKey = chatKeyFor(message.platform, message.chatId);
|
||||
const existing = this.chatWork.get(chatKey);
|
||||
const wasBusy = Boolean(existing?.running || existing?.queued);
|
||||
const acknowledgement = !command && !options.synchronous
|
||||
? this.sendAcknowledgement(message, adapter, wasBusy)
|
||||
: Promise.resolve();
|
||||
return this.withChatLock(chatKey, async () => {
|
||||
await acknowledgement;
|
||||
const replySequence = !command && !options.synchronous ? 2 : 1;
|
||||
try {
|
||||
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);
|
||||
return this.reply(message, adapter, reply, options, replySequence);
|
||||
} catch (error) {
|
||||
const errorText = error instanceof Error ? error.message : String(error);
|
||||
const reply = `Agent error: ${errorText}`;
|
||||
if (!options.synchronous) await this.send(message, adapter, reply);
|
||||
if (!options.synchronous) await this.send(message, adapter, reply, replySequence);
|
||||
return { ok: false, error: errorText, reply };
|
||||
}
|
||||
});
|
||||
@@ -45,7 +61,14 @@ export class Gateway {
|
||||
switch (command.kind) {
|
||||
case "help": return ["Commands:", "/status", "/cancel", "/new", "/help"].join("\n");
|
||||
case "status": {
|
||||
const status = this.runtime.status(message.platform, message.chatId);
|
||||
const chatKey = chatKeyFor(message.platform, message.chatId);
|
||||
const work = this.chatWork.get(chatKey);
|
||||
const status = {
|
||||
...this.runtime.status(message.platform, message.chatId),
|
||||
gatewayRunning: Boolean(work?.running),
|
||||
queued: work?.queued || 0,
|
||||
gatewayRunningSeconds: work?.startedAt ? Math.max(0, Math.floor((Date.now() - work.startedAt) / 1000)) : 0
|
||||
};
|
||||
return `OK\n${Object.entries(status).map(([key, value]) => `${key}=${value}`).join("\n")}`;
|
||||
}
|
||||
case "new":
|
||||
@@ -60,13 +83,26 @@ export class Gateway {
|
||||
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);
|
||||
private async reply(message: IncomingMessage, adapter: PlatformAdapter, reply: string, options: { synchronous?: boolean }, replySequence = 1): Promise<GatewayResult> {
|
||||
if (!options.synchronous) await this.send(message, adapter, reply, replySequence);
|
||||
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 send(message: IncomingMessage, adapter: PlatformAdapter, text: string, replySequence: number): Promise<void> {
|
||||
return adapter.sendMessage({
|
||||
target: { platform: message.platform, chatId: message.chatId, userId: message.userId, raw: message.raw },
|
||||
text,
|
||||
replyTo: message.messageId,
|
||||
replySequence
|
||||
});
|
||||
}
|
||||
|
||||
private async sendAcknowledgement(message: IncomingMessage, adapter: PlatformAdapter, queued: boolean): Promise<void> {
|
||||
const text = queued
|
||||
? "已收到,已排队。可随时发送 /status 查看状态。"
|
||||
: "已收到,正在处理。可随时发送 /status 查看状态。";
|
||||
try { await this.send(message, adapter, text, 1); }
|
||||
catch { console.error(`Gateway acknowledgement send failed (platform=${message.platform})`); }
|
||||
}
|
||||
|
||||
private checkPolicy(message: IncomingMessage): string | undefined {
|
||||
@@ -78,14 +114,25 @@ export class Gateway {
|
||||
|
||||
private async withChatLock<T>(key: string, fn: () => Promise<T>): Promise<T> {
|
||||
const previous = this.locks.get(key) || Promise.resolve();
|
||||
const state = this.chatWork.get(key) || { running: false, queued: 0 };
|
||||
state.queued++;
|
||||
this.chatWork.set(key, state);
|
||||
let release!: () => void;
|
||||
const current = new Promise<void>((resolve) => { release = resolve; });
|
||||
const queued = previous.then(() => current);
|
||||
this.locks.set(key, queued);
|
||||
const lock = previous.then(() => current);
|
||||
this.locks.set(key, lock);
|
||||
await previous;
|
||||
state.queued--;
|
||||
state.running = true;
|
||||
state.startedAt = Date.now();
|
||||
try { return await fn(); } finally {
|
||||
state.running = false;
|
||||
state.startedAt = undefined;
|
||||
release();
|
||||
if (this.locks.get(key) === queued) this.locks.delete(key);
|
||||
if (this.locks.get(key) === lock) {
|
||||
this.locks.delete(key);
|
||||
this.chatWork.delete(key);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -25,6 +25,7 @@ export interface OutgoingMessage {
|
||||
target: MessageTarget;
|
||||
text: string;
|
||||
replyTo?: string;
|
||||
replySequence?: number;
|
||||
}
|
||||
|
||||
export interface WebhookRequestContext {
|
||||
|
||||
@@ -72,7 +72,7 @@ export class QqAdapter implements PlatformAdapter {
|
||||
Authorization: `QQBot ${token}`,
|
||||
"Content-Type": "application/json; charset=utf-8"
|
||||
},
|
||||
body: JSON.stringify({ content: message.text, msg_id: message.replyTo })
|
||||
body: JSON.stringify({ content: message.text, msg_id: message.replyTo, msg_seq: message.replySequence ?? 1 })
|
||||
});
|
||||
if (!response.ok) throw new Error(`QQ send failed: HTTP ${response.status}`);
|
||||
const data = await response.json() as QqSendMessageResponse;
|
||||
|
||||
Reference in New Issue
Block a user