Refactor runtime around isolated bot instances
This commit is contained in:
@@ -1,23 +1,12 @@
|
||||
import type { AcpBackendConfig } from "../config.js";
|
||||
import type { AgentConfig } from "../config.js";
|
||||
import type { AcpBackendSpec } from "./types.js";
|
||||
|
||||
// Legacy compatibility wrapper; Config v3 runtime uses bot.agent directly.
|
||||
export class AcpBackendRegistry {
|
||||
private readonly backends = new Map<string, AcpBackendSpec>();
|
||||
|
||||
constructor(configs: AcpBackendConfig[]) {
|
||||
for (const config of configs) this.register({ ...config });
|
||||
}
|
||||
|
||||
register(backend: AcpBackendSpec): void {
|
||||
if (this.backends.has(backend.id)) throw new Error(`Duplicate ACP backend: ${backend.id}`);
|
||||
this.backends.set(backend.id, backend);
|
||||
}
|
||||
|
||||
constructor(private readonly agent: AgentConfig) {}
|
||||
get(id: string): AcpBackendSpec {
|
||||
const backend = this.backends.get(id);
|
||||
if (!backend) throw new Error(`Unknown ACP backend: ${id}`);
|
||||
return backend;
|
||||
if (id !== this.agent.id) throw new Error(`Unknown ACP agent: ${id}`);
|
||||
return { ...this.agent };
|
||||
}
|
||||
|
||||
list(): string[] { return [...this.backends.keys()].sort(); }
|
||||
list(): string[] { return [this.agent.id]; }
|
||||
}
|
||||
|
||||
@@ -1,8 +1,8 @@
|
||||
import type { AcpBackendConfig } from "../../config.js";
|
||||
import type { AgentConfig } from "../../config.js";
|
||||
import type { AcpBackendSpec } from "../types.js";
|
||||
|
||||
export function kimiBackend(config: AcpBackendConfig): AcpBackendSpec {
|
||||
if (config.id !== "kimi") throw new Error(`Expected kimi backend, got '${config.id}'`);
|
||||
if (!config.args.includes("acp")) throw new Error("Kimi ACP backend args must include 'acp'");
|
||||
export function kimiBackend(config: AgentConfig): AcpBackendSpec {
|
||||
if (config.id !== "kimi") throw new Error(`Expected kimi agent, got '${config.id}'`);
|
||||
if (!config.args.includes("acp")) throw new Error("Kimi ACP agent args must include 'acp'");
|
||||
return { ...config };
|
||||
}
|
||||
|
||||
+32
-14
@@ -2,11 +2,11 @@ import type { ChildProcessWithoutNullStreams } from "node:child_process";
|
||||
import { Readable, Writable } from "node:stream";
|
||||
import * as acp from "@agentclientprotocol/sdk";
|
||||
import type { AgentCapabilities, InitializeResponse, RequestPermissionRequest, RequestPermissionResponse, SessionNotification } from "@agentclientprotocol/sdk";
|
||||
import type { RolePolicy } from "../config.js";
|
||||
import type { PermissionPolicy } from "../config.js";
|
||||
|
||||
export interface AcpClientOptions {
|
||||
initializeTimeoutMs: number;
|
||||
policy: RolePolicy;
|
||||
policy: PermissionPolicy;
|
||||
}
|
||||
|
||||
export class AcpClient {
|
||||
@@ -48,7 +48,12 @@ export class AcpClient {
|
||||
this.collecting = false;
|
||||
this.chunks = [];
|
||||
if (this.capabilities.sessionCapabilities?.resume) {
|
||||
await this.connection.agent.request(acp.methods.agent.session.resume, { sessionId, cwd, mcpServers: [] });
|
||||
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 {
|
||||
@@ -92,24 +97,37 @@ export class AcpClient {
|
||||
}
|
||||
}
|
||||
|
||||
export function decidePermission(request: RequestPermissionRequest, policy: RolePolicy): RequestPermissionResponse {
|
||||
if (policy.permissionMode === "deny") return reject(request);
|
||||
export function decidePermission(request: RequestPermissionRequest, policy: PermissionPolicy): RequestPermissionResponse {
|
||||
if (policy.mode === "deny") return reject(request);
|
||||
const allowOption = request.options.find((option) => option.kind === "allow_once") || request.options.find((option) => option.kind === "allow_always");
|
||||
if (!allowOption) return reject(request);
|
||||
if (policy.permissionMode === "auto") return { outcome: { outcome: "selected", optionId: allowOption.optionId } };
|
||||
if (policy.mode === "auto") return { outcome: { outcome: "selected", optionId: allowOption.optionId } };
|
||||
|
||||
const name = String(request.toolCall.name || request.toolCall.kind || "").toLowerCase();
|
||||
const title = String(request.toolCall.title || "").toLowerCase();
|
||||
const allowedTool = policy.allowedTools.some((tool) => name === tool.toLowerCase() || title.startsWith(tool.toLowerCase()));
|
||||
if (!allowedTool) return reject(request);
|
||||
if (name === "bash" || name === "terminal" || title.startsWith("bash") || title.startsWith("terminal")) {
|
||||
if (request.toolCall.rawInput === undefined || policy.allowedCommandPatterns.length === 0) return reject(request);
|
||||
const input = typeof request.toolCall.rawInput === "string" ? request.toolCall.rawInput : JSON.stringify(request.toolCall.rawInput);
|
||||
if (!policy.allowedCommandPatterns.some((pattern) => new RegExp(pattern).test(input))) return reject(request);
|
||||
const name = typeof request.toolCall.name === "string" ? request.toolCall.name.trim().toLowerCase() : "";
|
||||
if (!name || !policy.allowedTools.some((tool) => name === tool.trim().toLowerCase())) return reject(request);
|
||||
if (name === "bash" || name === "terminal") {
|
||||
if (policy.allowedCommandPatterns.length === 0) return reject(request);
|
||||
const input = commandInput(request.toolCall.rawInput);
|
||||
if (!input || !policy.allowedCommandPatterns.some((pattern) => fullMatch(pattern, input))) return reject(request);
|
||||
}
|
||||
return { outcome: { outcome: "selected", optionId: allowOption.optionId } };
|
||||
}
|
||||
|
||||
function commandInput(rawInput: unknown): string | undefined {
|
||||
if (typeof rawInput === "string") return rawInput;
|
||||
if (typeof rawInput === "object" && rawInput !== null && !Array.isArray(rawInput)) {
|
||||
const record = rawInput as Record<string, unknown>;
|
||||
if (Object.keys(record).some((key) => !["command", "timeout", "timeoutMs"].includes(key))) return undefined;
|
||||
return typeof record.command === "string" ? record.command : undefined;
|
||||
}
|
||||
return undefined;
|
||||
}
|
||||
|
||||
function fullMatch(pattern: string, input: string): boolean {
|
||||
const match = new RegExp(pattern).exec(input);
|
||||
return match?.index === 0 && match[0] === input;
|
||||
}
|
||||
|
||||
function reject(request: RequestPermissionRequest): RequestPermissionResponse {
|
||||
const option = request.options.find((item) => item.kind === "reject_once") || request.options.find((item) => item.kind === "reject_always");
|
||||
return option ? { outcome: { outcome: "selected", optionId: option.optionId } } : { outcome: { outcome: "cancelled" } };
|
||||
|
||||
+4
-4
@@ -1,10 +1,10 @@
|
||||
import { spawn } from "node:child_process";
|
||||
import type { AcpBackendConfig } from "../config.js";
|
||||
import type { AgentConfig } from "../config.js";
|
||||
import type { InitializeResponse } from "@agentclientprotocol/sdk";
|
||||
import { AcpClient } from "./client.js";
|
||||
|
||||
export async function probeAcpBackend(backend: AcpBackendConfig, timeoutMs = 10_000): Promise<InitializeResponse> {
|
||||
const child = spawn(backend.command, backend.args, { stdio: ["pipe", "pipe", "pipe"], env: { ...process.env, ...backend.env } });
|
||||
const client = new AcpClient(child, { initializeTimeoutMs: timeoutMs, policy: { permissionMode: "deny", allowedTools: [], allowedCommandPatterns: [] } });
|
||||
export async function probeAcpBackend(agent: AgentConfig, timeoutMs = 10_000): Promise<InitializeResponse> {
|
||||
const child = spawn(agent.command, agent.args, { stdio: ["pipe", "pipe", "pipe"], env: { ...process.env, ...agent.env } });
|
||||
const client = new AcpClient(child, { initializeTimeoutMs: timeoutMs, policy: { mode: "deny", allowedTools: [], allowedCommandPatterns: [] } });
|
||||
try { return await client.initialize(); } finally { client.close(); child.kill("SIGTERM"); }
|
||||
}
|
||||
|
||||
+39
-41
@@ -1,9 +1,8 @@
|
||||
import type { AcpConfig } from "../config.js";
|
||||
import { chatKeyFor, type DurableSessionStore, type SessionBinding } from "../core/durable-session-store.js";
|
||||
import type { RoleRegistry } from "../roles/role-registry.js";
|
||||
import type { AcpBackendRegistry } from "./backend-registry.js";
|
||||
import type { ResolvedBot } from "../roles/role-registry.js";
|
||||
import type { ConversationRequest, ConversationResponse, ConversationRuntime, RuntimeStats } from "./types.js";
|
||||
import { AcpWorker } from "./worker.js";
|
||||
import { AcpSessionRestoreError, AcpWorker } from "./worker.js";
|
||||
|
||||
export class AcpSessionManager implements ConversationRuntime {
|
||||
private readonly workers = new Map<string, AcpWorker>();
|
||||
@@ -14,8 +13,7 @@ export class AcpSessionManager implements ConversationRuntime {
|
||||
|
||||
constructor(
|
||||
private readonly config: AcpConfig,
|
||||
private readonly backends: AcpBackendRegistry,
|
||||
private readonly roles: RoleRegistry,
|
||||
private readonly bot: ResolvedBot,
|
||||
private readonly store: DurableSessionStore
|
||||
) {
|
||||
this.sweeper = setInterval(() => void this.sweep(), config.sweepIntervalMs);
|
||||
@@ -25,30 +23,36 @@ export class AcpSessionManager implements ConversationRuntime {
|
||||
async prompt(request: ConversationRequest): Promise<ConversationResponse> {
|
||||
if (this.shuttingDown) throw new Error("ACP runtime is shutting down");
|
||||
const chatKey = chatKeyFor(request.platform, request.chatId);
|
||||
const role = this.roles.get(this.store.getSelectedRole(chatKey, this.roles.defaultId()));
|
||||
let binding = this.store.getBinding(chatKey, role.id);
|
||||
if (binding && (binding.roleFingerprint !== role.fingerprint || binding.backendId !== role.backend || binding.workspace !== role.workspace)) {
|
||||
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;
|
||||
}
|
||||
|
||||
const worker = await this.acquireWorker(role.id, binding);
|
||||
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(role.bootstrap);
|
||||
await worker.prompt(this.bot.bootstrap);
|
||||
binding = {
|
||||
chatKey, roleId: role.id, backendId: role.backend, nativeSessionId: worker.nativeSessionId!,
|
||||
workspace: role.workspace, roleFingerprint: role.fingerprint, createdAt: now, updatedAt: now
|
||||
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, role.id);
|
||||
return { text: text || "(ACP agent returned no text)", roleId: role.id, backendId: role.backend };
|
||||
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(worker.backend.id, worker.nativeSessionId));
|
||||
if (worker.nativeSessionId) this.workers.delete(workerKey(this.bot.agent.id, worker.nativeSessionId));
|
||||
await worker.terminate();
|
||||
throw error;
|
||||
} finally {
|
||||
@@ -64,25 +68,20 @@ export class AcpSessionManager implements ConversationRuntime {
|
||||
async reset(platform: string, chatId: string): Promise<void> {
|
||||
const chatKey = chatKeyFor(platform, chatId);
|
||||
await this.cancel(platform, chatId);
|
||||
const roleId = this.store.getSelectedRole(chatKey, this.roles.defaultId());
|
||||
const binding = await this.store.deleteBinding(chatKey, roleId);
|
||||
if (binding) await this.stopWorker(binding.backendId, binding.nativeSessionId);
|
||||
}
|
||||
|
||||
async selectRole(platform: string, chatId: string, roleId: string): Promise<void> {
|
||||
this.roles.get(roleId);
|
||||
await this.store.setSelectedRole(chatKeyFor(platform, chatId), roleId);
|
||||
}
|
||||
|
||||
selectedRole(platform: string, chatId: string): string {
|
||||
return this.store.getSelectedRole(chatKeyFor(platform, chatId), this.roles.defaultId());
|
||||
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 roleId = this.store.getSelectedRole(chatKey, this.roles.defaultId());
|
||||
const binding = this.store.getBinding(chatKey, roleId);
|
||||
return { role: roleId, backend: this.roles.get(roleId).backend, persisted: Boolean(binding), running: this.inFlight.has(chatKey) };
|
||||
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 {
|
||||
@@ -99,20 +98,19 @@ export class AcpSessionManager implements ConversationRuntime {
|
||||
this.inFlight.clear();
|
||||
}
|
||||
|
||||
private async acquireWorker(roleId: string, binding?: SessionBinding): Promise<AcpWorker> {
|
||||
const role = this.roles.get(roleId);
|
||||
private async acquireWorker(binding?: SessionBinding): Promise<AcpWorker> {
|
||||
if (binding) {
|
||||
const existing = this.workers.get(workerKey(binding.backendId, binding.nativeSessionId));
|
||||
const existing = this.workers.get(workerKey(binding.agentId, binding.nativeSessionId));
|
||||
if (existing) return existing;
|
||||
}
|
||||
await this.ensureCapacity();
|
||||
const worker = new AcpWorker(this.backends.get(role.backend), role, this.config, (crashed, error) => {
|
||||
const worker = new AcpWorker(this.bot, this.config, (crashed, error) => {
|
||||
this.crashes++;
|
||||
if (crashed.nativeSessionId) this.workers.delete(workerKey(crashed.backend.id, crashed.nativeSessionId));
|
||||
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(role.backend, nativeSessionId), worker);
|
||||
this.workers.set(workerKey(this.bot.agent.id, nativeSessionId), worker);
|
||||
return worker;
|
||||
}
|
||||
|
||||
@@ -134,12 +132,12 @@ export class AcpSessionManager implements ConversationRuntime {
|
||||
}
|
||||
|
||||
private async dropBinding(binding: SessionBinding): Promise<void> {
|
||||
await this.store.deleteBinding(binding.chatKey, binding.roleId);
|
||||
await this.stopWorker(binding.backendId, binding.nativeSessionId);
|
||||
await this.store.deleteBinding(binding.chatKey);
|
||||
await this.stopWorker(binding.agentId, binding.nativeSessionId);
|
||||
}
|
||||
|
||||
private async stopWorker(backendId: string, nativeSessionId: string): Promise<void> {
|
||||
const key = workerKey(backendId, 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);
|
||||
@@ -147,4 +145,4 @@ export class AcpSessionManager implements ConversationRuntime {
|
||||
}
|
||||
}
|
||||
|
||||
function workerKey(backendId: string, nativeSessionId: string): string { return `${backendId}:${nativeSessionId}`; }
|
||||
function workerKey(agentId: string, nativeSessionId: string): string { return `${agentId}:${nativeSessionId}`; }
|
||||
|
||||
+4
-6
@@ -1,4 +1,4 @@
|
||||
import type { RolePolicy } from "../config.js";
|
||||
import type { PermissionPolicy } from "../config.js";
|
||||
|
||||
export interface AcpBackendSpec {
|
||||
id: string;
|
||||
@@ -17,8 +17,8 @@ export interface ConversationRequest {
|
||||
|
||||
export interface ConversationResponse {
|
||||
text: string;
|
||||
roleId: string;
|
||||
backendId: string;
|
||||
botId: string;
|
||||
agentId: string;
|
||||
}
|
||||
|
||||
export interface RuntimeStats {
|
||||
@@ -32,13 +32,11 @@ export interface ConversationRuntime {
|
||||
prompt(request: ConversationRequest): Promise<ConversationResponse>;
|
||||
cancel(platform: string, chatId: string): Promise<boolean>;
|
||||
reset(platform: string, chatId: string): Promise<void>;
|
||||
selectRole(platform: string, chatId: string, roleId: string): Promise<void>;
|
||||
selectedRole(platform: string, chatId: string): string;
|
||||
status(platform: string, chatId: string): Record<string, string | number | boolean>;
|
||||
stats(): RuntimeStats;
|
||||
shutdown(): Promise<void>;
|
||||
}
|
||||
|
||||
export interface PermissionContext {
|
||||
policy: RolePolicy;
|
||||
policy: PermissionPolicy;
|
||||
}
|
||||
|
||||
+12
-11
@@ -1,8 +1,9 @@
|
||||
import { spawn, type ChildProcessWithoutNullStreams } from "node:child_process";
|
||||
import type { AcpConfig } from "../config.js";
|
||||
import type { ResolvedRole } from "../roles/role-registry.js";
|
||||
import type { ResolvedBot } from "../roles/role-registry.js";
|
||||
import { AcpClient } from "./client.js";
|
||||
import type { AcpBackendSpec } from "./types.js";
|
||||
|
||||
export class AcpSessionRestoreError extends Error {}
|
||||
|
||||
export class AcpWorker {
|
||||
private child?: ChildProcessWithoutNullStreams;
|
||||
@@ -16,16 +17,15 @@ export class AcpWorker {
|
||||
inFlight = false;
|
||||
|
||||
constructor(
|
||||
readonly backend: AcpBackendSpec,
|
||||
readonly role: ResolvedRole,
|
||||
readonly bot: ResolvedBot,
|
||||
private readonly config: AcpConfig,
|
||||
private readonly onCrash: (worker: AcpWorker, error: Error) => void
|
||||
) {}
|
||||
|
||||
async start(nativeSessionId?: string): Promise<string> {
|
||||
this.child = spawn(this.backend.command, this.backend.args, {
|
||||
cwd: this.role.workspace,
|
||||
env: { ...process.env, ...this.backend.env },
|
||||
this.child = spawn(this.bot.agent.command, this.bot.agent.args, {
|
||||
cwd: this.bot.workspace,
|
||||
env: { ...process.env, ...this.bot.agent.env },
|
||||
shell: false,
|
||||
stdio: ["pipe", "pipe", "pipe"]
|
||||
});
|
||||
@@ -35,18 +35,19 @@ 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.role.policy });
|
||||
this.client = new AcpClient(this.child, { initializeTimeoutMs: this.config.initializeTimeoutMs, policy: this.bot.permissions });
|
||||
try {
|
||||
await this.client.initialize();
|
||||
if (nativeSessionId) {
|
||||
await this.client.resumeSession(nativeSessionId, this.role.workspace);
|
||||
await this.client.resumeSession(nativeSessionId, this.bot.workspace);
|
||||
this.nativeSessionId = nativeSessionId;
|
||||
} else {
|
||||
this.nativeSessionId = await this.client.newSession(this.role.workspace);
|
||||
this.nativeSessionId = await this.client.newSession(this.bot.workspace);
|
||||
}
|
||||
return this.nativeSessionId;
|
||||
} catch (error) {
|
||||
await this.terminate();
|
||||
if (nativeSessionId) throw new AcpSessionRestoreError(`Cannot restore ACP session '${nativeSessionId}': ${error instanceof Error ? error.message : String(error)}`);
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
@@ -109,7 +110,7 @@ export class AcpWorker {
|
||||
if (!remaining) return;
|
||||
const text = chunk.subarray(0, remaining).toString("utf8").trimEnd();
|
||||
this.stderrBytes += Buffer.byteLength(text);
|
||||
if (text) console.error(`[acp:${this.backend.id}] ${text}`);
|
||||
if (text) console.error(`[acp:${this.bot.agent.id}] ${text}`);
|
||||
}
|
||||
|
||||
private crashed(error: Error): void {
|
||||
|
||||
+49
-12
@@ -1,8 +1,9 @@
|
||||
#!/usr/bin/env node
|
||||
import process from "node:process";
|
||||
import { discoverBackends } from "./acp/discovery.js";
|
||||
import { loadConfigFile } from "./cli/config-file.js";
|
||||
import { assertOperationalConfig, loadConfigFile } from "./cli/config-file.js";
|
||||
import { runDoctor } from "./cli/doctor.js";
|
||||
import { runInstanceCommand } from "./cli/instance.js";
|
||||
import { localBaseUrl } from "./cli/net.js";
|
||||
import { printFeishu } from "./cli/print.js";
|
||||
import { runSetup } from "./cli/setup.js";
|
||||
@@ -14,6 +15,8 @@ interface ParsedArgs { command?: string; rest: string[]; configPath?: string; js
|
||||
async function main(argv: string[]): Promise<number> {
|
||||
const parsed = parseArgs(argv);
|
||||
if (parsed.help || !parsed.command) { printHelp(); return 0; }
|
||||
validateCommandArgs(parsed);
|
||||
if (parsed.command === "instance") return runInstanceCommand(parsed.rest[0], parsed.rest.slice(1));
|
||||
if (parsed.command === "setup") { await runSetup(parsed.configPath); return 0; }
|
||||
if (parsed.command === "discover-backends" || parsed.command === "discover-agents") {
|
||||
const backends = await discoverBackends();
|
||||
@@ -27,6 +30,7 @@ async function main(argv: string[]): Promise<number> {
|
||||
}
|
||||
if (parsed.command === "start") {
|
||||
const loaded = loadConfigFile(parsed.configPath);
|
||||
assertOperationalConfig(loaded.config, loaded.path);
|
||||
const running = await startServer(loaded.config);
|
||||
return await new Promise<number>((resolve) => {
|
||||
let stopping = false;
|
||||
@@ -44,10 +48,10 @@ async function main(argv: string[]): Promise<number> {
|
||||
if (parsed.command === "doctor") { const loaded = loadConfigFile(parsed.configPath); return runDoctor(loaded.config, loaded.path); }
|
||||
if (parsed.command === "print") {
|
||||
const loaded = loadConfigFile(parsed.configPath);
|
||||
if (parsed.rest[0] === "feishu") { await printFeishu(loaded.config); return 0; }
|
||||
console.error(`Unknown print topic: ${parsed.rest[0] || "(missing)"}`); return 1;
|
||||
await printFeishu(loaded.config);
|
||||
return 0;
|
||||
}
|
||||
console.error(`Unknown command: ${parsed.command}`); printHelp(); return 1;
|
||||
throw new Error(`Unknown command: ${parsed.command}`);
|
||||
}
|
||||
|
||||
function parseArgs(argv: string[]): ParsedArgs {
|
||||
@@ -56,19 +60,52 @@ function parseArgs(argv: string[]): ParsedArgs {
|
||||
const arg = argv[index];
|
||||
if (arg === "--help" || arg === "-h") help = true;
|
||||
else if (arg === "--json") json = true;
|
||||
else if (arg === "--config") { configPath = argv[++index]; if (!configPath) throw new Error("--config requires a path"); }
|
||||
else if (!command) command = arg; else rest.push(arg);
|
||||
else if (arg === "--config") {
|
||||
if (configPath !== undefined) throw new Error("--config may only be specified once");
|
||||
const value = argv[++index];
|
||||
if (!value || value.startsWith("-")) throw new Error("--config requires a path");
|
||||
configPath = value;
|
||||
} else if (arg.startsWith("-")) throw new Error(`Unknown option: ${arg}`);
|
||||
else if (!command) command = arg;
|
||||
else rest.push(arg);
|
||||
}
|
||||
return { command, rest, configPath, json, help };
|
||||
}
|
||||
|
||||
function validateCommandArgs(parsed: ParsedArgs): void {
|
||||
const { command, rest, configPath, json } = parsed;
|
||||
if (command === "instance") {
|
||||
if (configPath || json) throw new Error("instance commands do not accept --config or --json");
|
||||
return;
|
||||
}
|
||||
if (command === "setup") {
|
||||
if (json || rest.length > 0) throw new Error("setup accepts only --config <path>");
|
||||
return;
|
||||
}
|
||||
if (command === "discover-backends" || command === "discover-agents") {
|
||||
if (configPath || rest.length > 0) throw new Error(`${command} accepts only --json`);
|
||||
return;
|
||||
}
|
||||
if (["start", "status", "doctor"].includes(command || "")) {
|
||||
if (json || rest.length > 0) throw new Error(`${command} accepts only --config <path>`);
|
||||
return;
|
||||
}
|
||||
if (command === "print") {
|
||||
if (json || rest.length !== 1 || rest[0] !== "feishu") throw new Error("print requires exactly the topic 'feishu'");
|
||||
return;
|
||||
}
|
||||
throw new Error(`Unknown command: ${command}`);
|
||||
}
|
||||
|
||||
async function printStatus(configPath?: string): Promise<void> {
|
||||
const loaded = loadConfigFile(configPath); const baseUrl = localBaseUrl(loaded.config);
|
||||
console.log(`Config: ${loaded.path}${loaded.exists ? "" : " (seeded from config.example.json)"}`);
|
||||
console.log(`Server: ${loaded.config.server.host}:${loaded.config.server.port}`);
|
||||
console.log(`Default role: ${loaded.config.defaultRole}`);
|
||||
console.log(`Roles: ${loaded.config.roles.map(({ id }) => id).join(", ")}`);
|
||||
console.log(`Backends: ${loaded.config.backends.map(({ id }) => id).join(", ")}`);
|
||||
console.log(`Config: ${loaded.path}`);
|
||||
console.log(`Config version: ${loaded.config.configVersion}`);
|
||||
console.log(`Bot: ${loaded.config.bot.id}`);
|
||||
console.log(`Agent: ${loaded.config.bot.agent.id}`);
|
||||
console.log(`Platform: ${loaded.config.gateway.platform.type}`);
|
||||
console.log(`Server: ${loaded.config.gateway.server.host}:${loaded.config.gateway.server.port}`);
|
||||
console.log(`Workspace: ${loaded.config.bot.workspace}`);
|
||||
console.log(`State file: ${defaultStateFile(loaded.config)}`);
|
||||
try {
|
||||
const response = await fetch(`${baseUrl}/health`, { signal: AbortSignal.timeout(1_000) });
|
||||
@@ -77,7 +114,7 @@ async function printStatus(configPath?: string): Promise<void> {
|
||||
}
|
||||
|
||||
function printHelp(): void {
|
||||
console.log(`gori-agent - multi-IM ACP agent gateway\n\nUsage:\n gori-agent setup [--config path]\n gori-agent discover-backends [--json]\n gori-agent start [--config path]\n gori-agent status [--config path]\n gori-agent doctor [--config path]\n gori-agent print feishu [--config path]\n\nDeprecated alias: discover-agents`);
|
||||
console.log(`gori-agent - single-Bot ACP gateway\n\nUsage:\n gori-agent instance init <bot-id>\n gori-agent instance start <bot-id>\n gori-agent instance stop <bot-id>\n gori-agent instance restart <bot-id>\n gori-agent instance status <bot-id>\n gori-agent instance logs <bot-id>\n gori-agent instance doctor <bot-id>\n gori-agent instance list\n\nDebug commands:\n gori-agent setup [--config path]\n gori-agent discover-backends [--json]\n gori-agent start --config path\n gori-agent status --config path\n gori-agent doctor --config path\n gori-agent print feishu --config path\n\nInstances root: \${GORI_AGENT_ROOT:-$HOME/.gori-agent}/instances`);
|
||||
}
|
||||
|
||||
main(process.argv.slice(2)).then((code) => { if (Number.isInteger(code)) process.exitCode = code; }).catch((error) => {
|
||||
|
||||
+98
-36
@@ -7,60 +7,122 @@ import { loadConfigFromPath, parseConfig } from "../config.js";
|
||||
export interface LoadedConfigFile {
|
||||
path: string;
|
||||
config: AppConfig;
|
||||
exists: boolean;
|
||||
source: "explicit" | "env" | "local" | "example";
|
||||
source: "explicit" | "env" | "local";
|
||||
}
|
||||
|
||||
export function projectRoot(): string {
|
||||
return path.resolve(new URL("../..", import.meta.url).pathname);
|
||||
}
|
||||
|
||||
export function resolveCliConfigPath(configPath?: string): { path: string; exists: boolean; source: LoadedConfigFile["source"] } {
|
||||
if (configPath) {
|
||||
const resolved = path.resolve(configPath);
|
||||
return { path: resolved, exists: fs.existsSync(resolved), source: "explicit" };
|
||||
}
|
||||
|
||||
if (process.env.GORI_GATEWAY_CONFIG) {
|
||||
const resolved = path.resolve(process.env.GORI_GATEWAY_CONFIG);
|
||||
return { path: resolved, exists: fs.existsSync(resolved), source: "env" };
|
||||
}
|
||||
|
||||
const local = path.resolve("config.json");
|
||||
if (fs.existsSync(local)) return { path: local, exists: true, source: "local" };
|
||||
|
||||
return { path: path.resolve("config.json"), exists: false, source: "example" };
|
||||
export function resolveCliConfigPath(configPath?: string): { path: string; source: LoadedConfigFile["source"] } {
|
||||
if (configPath) return { path: path.resolve(configPath), source: "explicit" };
|
||||
if (process.env.GORI_GATEWAY_CONFIG) return { path: path.resolve(process.env.GORI_GATEWAY_CONFIG), source: "env" };
|
||||
return { path: path.resolve("config.json"), source: "local" };
|
||||
}
|
||||
|
||||
export function loadConfigFile(configPath?: string): LoadedConfigFile {
|
||||
const resolved = resolveCliConfigPath(configPath);
|
||||
if (resolved.exists) {
|
||||
return {
|
||||
path: resolved.path,
|
||||
config: loadConfigFromPath(resolved.path),
|
||||
exists: true,
|
||||
source: resolved.source
|
||||
};
|
||||
if (!fs.existsSync(resolved.path)) {
|
||||
throw new Error(`Runtime config does not exist: ${resolved.path}; config.example.json is documentation only`);
|
||||
}
|
||||
return { path: resolved.path, config: loadConfigFromPath(resolved.path), source: resolved.source };
|
||||
}
|
||||
|
||||
export function loadExampleConfig(): AppConfig {
|
||||
const examplePath = path.join(projectRoot(), "config.example.json");
|
||||
const raw = fs.readFileSync(examplePath, "utf8");
|
||||
const config = parseConfig(JSON.parse(raw) as unknown);
|
||||
return {
|
||||
path: resolved.path,
|
||||
config,
|
||||
exists: false,
|
||||
source: "example"
|
||||
};
|
||||
return parseConfig(JSON.parse(fs.readFileSync(examplePath, "utf8")) as unknown);
|
||||
}
|
||||
|
||||
export function writeConfigFile(configPath: string, config: AppConfig): void {
|
||||
const parsed = parseConfig(config);
|
||||
fs.mkdirSync(path.dirname(configPath), { recursive: true });
|
||||
fs.writeFileSync(configPath, `${JSON.stringify(parsed, null, 2)}\n`, "utf8");
|
||||
const directory = path.dirname(configPath);
|
||||
fs.mkdirSync(directory, { recursive: true, mode: 0o700 });
|
||||
const temp = path.join(directory, `.${path.basename(configPath)}.${process.pid}.${Date.now()}.tmp`);
|
||||
let fd: number | undefined;
|
||||
try {
|
||||
fd = fs.openSync(temp, "wx", 0o600);
|
||||
fs.writeFileSync(fd, `${JSON.stringify(parsed, null, 2)}\n`, "utf8");
|
||||
fs.fsyncSync(fd);
|
||||
fs.closeSync(fd);
|
||||
fd = undefined;
|
||||
fs.renameSync(temp, configPath);
|
||||
fs.chmodSync(configPath, 0o600);
|
||||
const dirFd = fs.openSync(directory, "r");
|
||||
try { fs.fsyncSync(dirFd); } finally { fs.closeSync(dirFd); }
|
||||
} catch (error) {
|
||||
if (fd !== undefined) try { fs.closeSync(fd); } catch { /* ignore cleanup failure */ }
|
||||
try { fs.unlinkSync(temp); } catch { /* ignore cleanup failure */ }
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
export function operationalConfigProblems(config: AppConfig): string[] {
|
||||
const problems: string[] = [];
|
||||
if (config.bot.id === "BOT_ID") problems.push("bot.id is still the example placeholder");
|
||||
if (isPlaceholder(config.bot.workspace) || config.bot.workspace.includes("/absolute/path/")) problems.push("bot.workspace is still a placeholder");
|
||||
if (isPlaceholder(config.bot.agent.command) || config.bot.agent.command.includes("/home/USER/")) problems.push("bot.agent.command is still a placeholder");
|
||||
else if (!isExecutable(config.bot.agent.command)) problems.push("bot.agent.command is not executable");
|
||||
try {
|
||||
if (!fs.statSync(config.bot.workspace).isDirectory()) problems.push("bot.workspace must be an existing directory");
|
||||
} catch { problems.push("bot.workspace must be an existing directory"); }
|
||||
|
||||
const platform = config.gateway.platform;
|
||||
switch (platform.type) {
|
||||
case "qq":
|
||||
if (missingOrPlaceholder(platform.appId)) problems.push("QQ appId is missing or placeholder");
|
||||
if (missingOrPlaceholder(platform.clientSecret)) problems.push("QQ clientSecret is missing or placeholder");
|
||||
if (platform.botNames.some(isPlaceholder)) problems.push("QQ botNames contains a placeholder");
|
||||
break;
|
||||
case "feishu":
|
||||
if (missingOrPlaceholder(platform.appId)) problems.push("Feishu appId is missing or placeholder");
|
||||
if (missingOrPlaceholder(platform.appSecret)) problems.push("Feishu appSecret is missing or placeholder");
|
||||
break;
|
||||
case "wecom":
|
||||
if (missingOrPlaceholder(platform.corpId)) problems.push("WeCom corpId is missing or placeholder");
|
||||
if (missingOrPlaceholder(platform.agentId)) problems.push("WeCom agentId is missing or placeholder");
|
||||
if (missingOrPlaceholder(platform.secret)) problems.push("WeCom secret is missing or placeholder");
|
||||
break;
|
||||
case "webhook":
|
||||
if (missingOrPlaceholder(platform.secret)) problems.push("Generic webhook secret is missing or placeholder");
|
||||
break;
|
||||
case "weixin":
|
||||
if (missingOrPlaceholder(platform.secret)) problems.push("Weixin bridge secret is missing or placeholder");
|
||||
break;
|
||||
}
|
||||
return problems;
|
||||
}
|
||||
|
||||
export function assertOperationalConfig(config: AppConfig, configFile?: string): void {
|
||||
const problems = operationalConfigProblems(config);
|
||||
if (configFile) {
|
||||
const stat = fs.lstatSync(configFile);
|
||||
if (!stat.isFile() || stat.isSymbolicLink()) problems.push("config must be a regular file");
|
||||
const mode = stat.mode & 0o777;
|
||||
if (mode !== 0o600) problems.push(`config file mode must be 0600, got 0${mode.toString(8)}`);
|
||||
}
|
||||
if (problems.length > 0) throw new Error(`Refusing to start: ${problems.join("; ")}`);
|
||||
}
|
||||
|
||||
function missingOrPlaceholder(value: string): boolean { return !value || isPlaceholder(value); }
|
||||
|
||||
function isExecutable(command: string): boolean {
|
||||
try {
|
||||
if (command.includes(path.sep)) {
|
||||
fs.accessSync(command, fs.constants.X_OK);
|
||||
return fs.statSync(command).isFile();
|
||||
}
|
||||
return (process.env.PATH || "").split(path.delimiter).some((directory) => {
|
||||
try { fs.accessSync(path.join(directory, command), fs.constants.X_OK); return fs.statSync(path.join(directory, command)).isFile(); }
|
||||
catch { return false; }
|
||||
});
|
||||
} catch { return false; }
|
||||
}
|
||||
|
||||
function isPlaceholder(value: string): boolean {
|
||||
return /^(?:BOT_ID|QQ_(?:APP_ID|CLIENT_SECRET|BOT_NAME)|(?:FEISHU|WECOM|WEBHOOK|WEIXIN)_[A-Z0-9_]+)$/.test(value)
|
||||
|| value.includes("<") || value.includes(">");
|
||||
}
|
||||
|
||||
export function readConfigJson(configPath?: string): unknown {
|
||||
const loaded = loadConfigFile(configPath);
|
||||
return loaded.config;
|
||||
return loadConfigFile(configPath).config;
|
||||
}
|
||||
|
||||
+27
-33
@@ -2,61 +2,55 @@ import fs from "node:fs";
|
||||
import path from "node:path";
|
||||
import { probeAcpBackend } from "../acp/probe.js";
|
||||
import { defaultStateFile, type AppConfig } from "../config.js";
|
||||
import { RoleRegistry } from "../roles/role-registry.js";
|
||||
import { BotProfileResolver } from "../roles/role-registry.js";
|
||||
import { operationalConfigProblems } from "./config-file.js";
|
||||
import { printUrlHints } from "./net.js";
|
||||
|
||||
export async function runDoctor(config: AppConfig, configPath: string): Promise<number> {
|
||||
let problems = 0;
|
||||
console.log(`Config: ${configPath}`);
|
||||
console.log(`Server: ${config.server.host}:${config.server.port}`);
|
||||
console.log(`Config version: ${config.configVersion}`);
|
||||
console.log(`Bot: ${config.bot.id}`);
|
||||
console.log(`Server: ${config.gateway.server.host}:${config.gateway.server.port}`);
|
||||
console.log(`Platform: ${config.gateway.platform.type}`);
|
||||
|
||||
for (const problem of operationalConfigProblems(config)) problems += reportError(problem);
|
||||
try {
|
||||
new RoleRegistry(config);
|
||||
console.log(`Default role: ${config.defaultRole}`);
|
||||
const mode = fs.statSync(configPath).mode & 0o777;
|
||||
if (mode !== 0o600) problems += reportError(`config file mode must be 0600, got 0${mode.toString(8)}`);
|
||||
} catch (error) {
|
||||
problems += reportError(`cannot stat config file: ${error instanceof Error ? error.message : String(error)}`);
|
||||
}
|
||||
|
||||
try { new BotProfileResolver(config); } catch (error) {
|
||||
problems += reportError(error instanceof Error ? error.message : String(error));
|
||||
}
|
||||
|
||||
for (const role of config.roles) {
|
||||
if (!fs.existsSync(role.workspace)) problems += reportError(`role '${role.id}' workspace does not exist: ${role.workspace}`);
|
||||
if (role.policy.permissionMode === "auto") console.log(`WARN: role '${role.id}' auto-approves every permission request.`);
|
||||
if (role.policy.permissionMode === "allowlist" && role.policy.allowedTools.includes("bash") && role.policy.allowedCommandPatterns.length === 0) {
|
||||
console.log(`WARN: role '${role.id}' allows bash by name but has no command patterns; bash requests will be denied.`);
|
||||
}
|
||||
if (config.gateway.platform.type === "qq") console.log(`QQ connection mode: ${config.gateway.platform.connectionMode}`);
|
||||
if (config.bot.permissions.mode === "auto") console.log("WARN: bot auto-approves every permission request.");
|
||||
if (config.bot.permissions.mode === "allowlist" && config.bot.permissions.allowedTools.some((tool) => ["bash", "terminal"].includes(tool.toLowerCase())) && config.bot.permissions.allowedCommandPatterns.length === 0) {
|
||||
console.log("WARN: bot allows bash/terminal by name but has no command patterns; command requests will be denied.");
|
||||
}
|
||||
|
||||
for (const backend of config.backends) {
|
||||
try {
|
||||
const initialized = await probeAcpBackend(backend, config.acp.initializeTimeoutMs);
|
||||
const capabilities = initialized.agentCapabilities;
|
||||
console.log(`ACP '${backend.id}': ready (${initialized.agentInfo?.name || "unknown"} ${initialized.agentInfo?.version || ""})`);
|
||||
console.log(` load=${Boolean(capabilities?.loadSession)} resume=${Boolean(capabilities?.sessionCapabilities?.resume)} list=${Boolean(capabilities?.sessionCapabilities?.list)} close=${Boolean(capabilities?.sessionCapabilities?.close)}`);
|
||||
if (!capabilities?.loadSession && !capabilities?.sessionCapabilities?.resume) problems += reportError(`backend '${backend.id}' cannot restore sessions`);
|
||||
console.log(" model config: not advertised by initialize; using the Kimi default model");
|
||||
} catch (error) {
|
||||
problems += reportError(`ACP '${backend.id}' initialize failed: ${error instanceof Error ? error.message : String(error)}`);
|
||||
}
|
||||
try {
|
||||
const initialized = await probeAcpBackend(config.bot.agent, config.runtime.acp.initializeTimeoutMs);
|
||||
const capabilities = initialized.agentCapabilities;
|
||||
console.log(`ACP '${config.bot.agent.id}': ready (${initialized.agentInfo?.name || "unknown"} ${initialized.agentInfo?.version || ""})`);
|
||||
console.log(` load=${Boolean(capabilities?.loadSession)} resume=${Boolean(capabilities?.sessionCapabilities?.resume)} list=${Boolean(capabilities?.sessionCapabilities?.list)} close=${Boolean(capabilities?.sessionCapabilities?.close)}`);
|
||||
if (!capabilities?.loadSession && !capabilities?.sessionCapabilities?.resume) problems += reportError(`agent '${config.bot.agent.id}' cannot restore sessions`);
|
||||
} catch (error) {
|
||||
problems += reportError(`ACP '${config.bot.agent.id}' initialize failed: ${error instanceof Error ? error.message : String(error)}`);
|
||||
}
|
||||
|
||||
const stateFile = defaultStateFile(config);
|
||||
try {
|
||||
const stateDirectory = path.dirname(stateFile);
|
||||
fs.mkdirSync(stateDirectory, { recursive: true });
|
||||
fs.mkdirSync(stateDirectory, { recursive: true, mode: 0o700 });
|
||||
fs.accessSync(stateDirectory, fs.constants.R_OK | fs.constants.W_OK);
|
||||
console.log(`State file: ${stateFile}`);
|
||||
} catch (error) {
|
||||
problems += reportError(`state directory is unavailable: ${error instanceof Error ? error.message : String(error)}`);
|
||||
}
|
||||
|
||||
const enabledPlatforms = Object.entries(config.platforms).filter(([, value]) => value.enabled).map(([name]) => name);
|
||||
console.log(`Enabled platforms: ${enabledPlatforms.join(", ") || "none"}`);
|
||||
if (config.platforms.qq.enabled) {
|
||||
console.log(`QQ connection mode: ${config.platforms.qq.connectionMode}`);
|
||||
if (!config.platforms.qq.appId) problems += reportError("QQ appId is empty.");
|
||||
if (!config.platforms.qq.clientSecret || config.platforms.qq.clientSecret === "replace-me") console.log("WARN: QQ clientSecret is missing or placeholder.");
|
||||
if (config.platforms.qq.connectionMode === "webhook" && config.platforms.qq.verifySignature && !(config.platforms.qq.botSecret || config.platforms.qq.clientSecret)) {
|
||||
problems += reportError("QQ botSecret or clientSecret is required for webhook signature verification.");
|
||||
}
|
||||
}
|
||||
await printUrlHints(config);
|
||||
return problems > 0 ? 1 : 0;
|
||||
}
|
||||
|
||||
@@ -0,0 +1,348 @@
|
||||
import { spawn } from "node:child_process";
|
||||
import fs from "node:fs";
|
||||
import net from "node:net";
|
||||
import os from "node:os";
|
||||
import path from "node:path";
|
||||
import process from "node:process";
|
||||
import { defaultStateFile } from "../config.js";
|
||||
import { assertOperationalConfig, loadConfigFile } from "./config-file.js";
|
||||
import { runDoctor } from "./doctor.js";
|
||||
import { localBaseUrl } from "./net.js";
|
||||
import { runSetup } from "./setup.js";
|
||||
|
||||
const BOT_ID_PATTERN = /^[a-z0-9](?:[a-z0-9-]{0,62})$/;
|
||||
|
||||
export function agentRoot(): string {
|
||||
return path.resolve(process.env.GORI_AGENT_ROOT || path.join(os.homedir(), ".gori-agent"));
|
||||
}
|
||||
|
||||
export function instancesRoot(): string { return path.join(agentRoot(), "instances"); }
|
||||
|
||||
export function instanceDirectory(botId: string): string {
|
||||
validateBotId(botId);
|
||||
return path.join(instancesRoot(), botId);
|
||||
}
|
||||
|
||||
export function instanceConfigPath(botId: string): string { return path.join(instanceDirectory(botId), "config.json"); }
|
||||
|
||||
export async function runInstanceCommand(action: string | undefined, args: string[]): Promise<number> {
|
||||
if (action === "list") {
|
||||
if (args.length !== 0) throw new Error("instance list does not accept <bot-id>");
|
||||
return listInstances();
|
||||
}
|
||||
if (!action || !["init", "start", "stop", "restart", "status", "logs", "doctor"].includes(action)) {
|
||||
throw new Error(`Unknown instance command: ${action || "(missing)"}`);
|
||||
}
|
||||
if (args.length !== 1) throw new Error(`instance ${action} requires exactly one <bot-id>`);
|
||||
const botId = args[0];
|
||||
validateBotId(botId);
|
||||
switch (action) {
|
||||
case "init": return initInstance(botId);
|
||||
case "start": return startInstance(botId);
|
||||
case "stop": return stopInstance(botId);
|
||||
case "restart": await stopInstance(botId); return startInstance(botId);
|
||||
case "status": return statusInstance(botId);
|
||||
case "logs": return logsInstance(botId);
|
||||
case "doctor": return doctorInstance(botId);
|
||||
default: throw new Error(`Unknown instance command: ${action}`);
|
||||
}
|
||||
}
|
||||
|
||||
async function initInstance(botId: string): Promise<number> {
|
||||
const directory = instanceDirectory(botId);
|
||||
if (fs.existsSync(directory)) throw new Error(`Instance already exists: ${botId}`);
|
||||
ensureInstanceDirectories(directory, true);
|
||||
try {
|
||||
await runSetup(path.join(directory, "config.json"), { botId, requireNew: true, writeWithoutConfirmation: true });
|
||||
return 0;
|
||||
} catch (error) {
|
||||
cleanupFailedInit(directory);
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
async function startInstance(botId: string): Promise<number> {
|
||||
const loaded = loadInstance(botId);
|
||||
assertOperationalConfig(loaded.config);
|
||||
assertConfigMode(loaded.path);
|
||||
const paths = runtimePaths(botId);
|
||||
ensureInstanceDirectories(paths.home, false);
|
||||
const current = inspectPid(paths.pidFile, loaded.path);
|
||||
if (current.kind === "running") throw new Error(`Instance '${botId}' is already running (pid ${current.pid})`);
|
||||
if (current.kind === "foreign") throw new Error(`Refusing to start: PID file references unrelated live process ${current.pid}`);
|
||||
if (current.kind === "stale") removeStalePid(paths.pidFile, current);
|
||||
await assertPortAvailable(loaded.config.gateway.server.host, loaded.config.gateway.server.port);
|
||||
|
||||
const logFd = fs.openSync(paths.logFile, "a", 0o600);
|
||||
fs.chmodSync(paths.logFile, 0o600);
|
||||
const cliFile = path.join(path.resolve(new URL("../..", import.meta.url).pathname), "dist", "cli.js");
|
||||
let child: ReturnType<typeof spawn> | undefined;
|
||||
try {
|
||||
child = spawn(process.execPath, [cliFile, "start", "--config", loaded.path], {
|
||||
cwd: paths.home,
|
||||
env: { ...process.env, GORI_AGENT_HOME: paths.home, GORI_GATEWAY_CONFIG: loaded.path },
|
||||
detached: true,
|
||||
stdio: ["ignore", logFd, logFd]
|
||||
});
|
||||
await new Promise<void>((resolve, reject) => {
|
||||
child!.once("spawn", resolve);
|
||||
child!.once("error", reject);
|
||||
});
|
||||
if (!child.pid) throw new Error("Instance child process has no PID");
|
||||
atomicWriteMode(paths.pidFile, `${child.pid}\n`, 0o600);
|
||||
await waitForHealthyInstance(child.pid, loaded.config.gateway.platform.type, botId, loaded.config);
|
||||
child.unref();
|
||||
console.log(`Instance '${botId}' started (pid ${child.pid})`);
|
||||
console.log(`Log: ${paths.logFile}`);
|
||||
return 0;
|
||||
} catch (error) {
|
||||
if (child?.pid) {
|
||||
await terminateStartedChild(child.pid);
|
||||
removePidIfOwned(paths.pidFile, child.pid);
|
||||
}
|
||||
throw new Error(`Instance '${botId}' failed to start; inspect ${paths.logFile}: ${error instanceof Error ? error.message : String(error)}`);
|
||||
} finally { fs.closeSync(logFd); }
|
||||
}
|
||||
|
||||
async function stopInstance(botId: string): Promise<number> {
|
||||
const loaded = loadInstance(botId);
|
||||
const paths = runtimePaths(botId);
|
||||
const current = inspectPid(paths.pidFile, loaded.path);
|
||||
if (current.kind === "foreign") throw new Error(`Refusing to stop: PID file references unrelated live process ${current.pid}`);
|
||||
if (current.kind !== "running") {
|
||||
if (current.kind === "stale") removeStalePid(paths.pidFile, current);
|
||||
console.log(`Instance '${botId}' is not running`);
|
||||
return 0;
|
||||
}
|
||||
process.kill(current.pid, "SIGTERM");
|
||||
for (let attempt = 0; attempt < 100 && isPidRunning(current.pid); attempt++) await delay(100);
|
||||
if (isPidRunning(current.pid)) throw new Error(`Instance '${botId}' did not stop within 10 seconds (pid ${current.pid})`);
|
||||
removeStalePid(paths.pidFile);
|
||||
console.log(`Instance '${botId}' stopped (pid ${current.pid})`);
|
||||
return 0;
|
||||
}
|
||||
|
||||
async function statusInstance(botId: string): Promise<number> {
|
||||
const loaded = loadInstance(botId);
|
||||
const paths = runtimePaths(botId);
|
||||
const current = inspectPid(paths.pidFile, loaded.path);
|
||||
const pid = current.kind === "running" ? current.pid : undefined;
|
||||
console.log(`Instance: ${botId}`);
|
||||
console.log(`Process: ${pid ? `running (pid ${pid})` : current.kind === "foreign" ? `PID identity mismatch (${current.pid})` : "not running"}`);
|
||||
console.log(`Config: ${loaded.path}`);
|
||||
console.log(`Platform: ${loaded.config.gateway.platform.type}`);
|
||||
console.log(`Agent: ${loaded.config.bot.agent.id}`);
|
||||
console.log(`Workspace: ${loaded.config.bot.workspace}`);
|
||||
console.log(`State file: ${withInstanceHome(paths.home, () => defaultStateFile(loaded.config))}`);
|
||||
if (current.kind === "foreign") return 1;
|
||||
if (pid) {
|
||||
try {
|
||||
const response = await fetch(`${localBaseUrl(loaded.config)}/health`, { signal: AbortSignal.timeout(1_000) });
|
||||
const health = await response.json() as { botId?: string; platform?: string };
|
||||
const matches = response.ok && health.botId === botId && health.platform === loaded.config.gateway.platform.type;
|
||||
console.log(`Health: HTTP ${response.status}, identity ${matches ? "verified" : "mismatch"}`);
|
||||
if (!matches) return 1;
|
||||
} catch { console.log("Health: not reachable on local URL"); return 1; }
|
||||
}
|
||||
return 0;
|
||||
}
|
||||
|
||||
async function logsInstance(botId: string): Promise<number> {
|
||||
loadInstance(botId);
|
||||
const paths = runtimePaths(botId);
|
||||
ensurePrivateDirectory(path.dirname(paths.logFile), true);
|
||||
const fd = fs.openSync(paths.logFile, "a", 0o600);
|
||||
fs.closeSync(fd);
|
||||
fs.chmodSync(paths.logFile, 0o600);
|
||||
return new Promise<number>((resolve, reject) => {
|
||||
const child = spawn("tail", ["-f", paths.logFile], { stdio: "inherit" });
|
||||
child.once("error", reject);
|
||||
child.once("close", (code) => resolve(code || 0));
|
||||
});
|
||||
}
|
||||
|
||||
async function doctorInstance(botId: string): Promise<number> {
|
||||
const loaded = loadInstance(botId);
|
||||
return withInstanceHomeAsync(instanceDirectory(botId), () => runDoctor(loaded.config, loaded.path));
|
||||
}
|
||||
|
||||
function listInstances(): number {
|
||||
const root = instancesRoot();
|
||||
if (!fs.existsSync(root)) return 0;
|
||||
ensurePrivateDirectory(root, false);
|
||||
for (const entry of fs.readdirSync(root, { withFileTypes: true }).filter((item) => item.isDirectory()).sort((a, b) => a.name.localeCompare(b.name))) {
|
||||
if (!BOT_ID_PATTERN.test(entry.name)) continue;
|
||||
const configFile = path.join(root, entry.name, "config.json");
|
||||
let state = "missing-config";
|
||||
if (fs.existsSync(configFile)) {
|
||||
const current = inspectPid(path.join(root, entry.name, "state", "gori-agent.pid"), configFile);
|
||||
state = current.kind === "running" ? "running" : current.kind === "foreign" ? "pid-mismatch" : "stopped";
|
||||
}
|
||||
console.log(`${entry.name}\t${state}`);
|
||||
}
|
||||
return 0;
|
||||
}
|
||||
|
||||
function loadInstance(botId: string): ReturnType<typeof loadConfigFile> {
|
||||
const configFile = instanceConfigPath(botId);
|
||||
if (fs.existsSync(path.dirname(configFile))) ensurePrivateDirectory(path.dirname(configFile), false);
|
||||
const loaded = loadConfigFile(configFile);
|
||||
if (loaded.config.bot.id !== botId) throw new Error(`Instance directory '${botId}' does not match config bot.id '${loaded.config.bot.id}'`);
|
||||
return loaded;
|
||||
}
|
||||
|
||||
function runtimePaths(botId: string): { home: string; pidFile: string; logFile: string } {
|
||||
const home = instanceDirectory(botId);
|
||||
return { home, pidFile: path.join(home, "state", "gori-agent.pid"), logFile: path.join(home, "logs", "gori-agent.log") };
|
||||
}
|
||||
|
||||
function ensureInstanceDirectories(home: string, createHome: boolean): void {
|
||||
ensurePrivateDirectory(agentRoot(), createHome);
|
||||
ensurePrivateDirectory(instancesRoot(), createHome);
|
||||
ensurePrivateDirectory(home, createHome);
|
||||
ensurePrivateDirectory(path.join(home, "state"), true);
|
||||
ensurePrivateDirectory(path.join(home, "logs"), true);
|
||||
}
|
||||
|
||||
function ensurePrivateDirectory(directory: string, create: boolean): void {
|
||||
if (!fs.existsSync(directory)) {
|
||||
if (!create) throw new Error(`Required directory does not exist: ${directory}`);
|
||||
fs.mkdirSync(directory, { mode: 0o700 });
|
||||
}
|
||||
const stat = fs.lstatSync(directory);
|
||||
if (!stat.isDirectory() || stat.isSymbolicLink()) throw new Error(`Refusing unsafe instance path: ${directory}`);
|
||||
const mode = stat.mode & 0o777;
|
||||
if (mode !== 0o700) throw new Error(`Directory mode must be 0700: ${directory} has 0${mode.toString(8)}`);
|
||||
if (typeof process.getuid === "function" && stat.uid !== process.getuid()) throw new Error(`Directory is not owned by the current user: ${directory}`);
|
||||
}
|
||||
|
||||
function cleanupFailedInit(directory: string): void {
|
||||
if (fs.existsSync(path.join(directory, "config.json"))) return;
|
||||
for (const name of ["state", "logs"]) {
|
||||
try { fs.rmdirSync(path.join(directory, name)); } catch { /* preserve non-empty or absent directories */ }
|
||||
}
|
||||
try { fs.rmdirSync(directory); } catch { /* preserve non-empty directory */ }
|
||||
}
|
||||
|
||||
function validateBotId(botId: string): void {
|
||||
if (!BOT_ID_PATTERN.test(botId)) throw new Error(`Invalid bot ID '${botId}'; use lowercase letters, digits, and hyphens`);
|
||||
}
|
||||
|
||||
function assertConfigMode(configFile: string): void {
|
||||
const stat = fs.lstatSync(configFile);
|
||||
if (!stat.isFile() || stat.isSymbolicLink()) throw new Error("Refusing to start: config must be a regular file");
|
||||
const mode = stat.mode & 0o777;
|
||||
if (mode !== 0o600) throw new Error(`Refusing to start: config file mode must be 0600, got 0${mode.toString(8)}`);
|
||||
}
|
||||
|
||||
type PidInspection = { kind: "absent" } | { kind: "stale"; dev?: number; ino?: number } | { kind: "running" | "foreign"; pid: number };
|
||||
|
||||
function inspectPid(pidFile: string, configFile: string): PidInspection {
|
||||
let text: string;
|
||||
let stat: fs.Stats;
|
||||
try { stat = fs.lstatSync(pidFile); text = fs.readFileSync(pidFile, "utf8").trim(); }
|
||||
catch (error) { return (error as NodeJS.ErrnoException).code === "ENOENT" ? { kind: "absent" } : { kind: "stale" }; }
|
||||
const pid = Number(text);
|
||||
if (!Number.isSafeInteger(pid) || pid <= 1 || !isPidRunning(pid)) return { kind: "stale", dev: stat.dev, ino: stat.ino };
|
||||
try {
|
||||
const commandLine = fs.readFileSync(`/proc/${pid}/cmdline`, "utf8").split("\0").filter(Boolean);
|
||||
const configIndex = commandLine.indexOf("--config");
|
||||
const matches = commandLine[2] === "start"
|
||||
&& commandLine[1]?.endsWith("/dist/cli.js")
|
||||
&& configIndex >= 0
|
||||
&& path.resolve(commandLine[configIndex + 1] || "") === path.resolve(configFile);
|
||||
return { kind: matches ? "running" : "foreign", pid };
|
||||
} catch { return isPidRunning(pid) ? { kind: "foreign", pid } : { kind: "stale", dev: stat.dev, ino: stat.ino }; }
|
||||
}
|
||||
|
||||
function isPidRunning(pid: number | undefined): boolean {
|
||||
if (!pid) return false;
|
||||
try { process.kill(pid, 0); return true; }
|
||||
catch (error) { return (error as NodeJS.ErrnoException).code === "EPERM"; }
|
||||
}
|
||||
|
||||
function removeStalePid(pidFile: string, inspection?: Extract<PidInspection, { kind: "stale" }>): void {
|
||||
try {
|
||||
if (inspection?.dev !== undefined && inspection.ino !== undefined) {
|
||||
const current = fs.lstatSync(pidFile);
|
||||
if (current.dev !== inspection.dev || current.ino !== inspection.ino) return;
|
||||
}
|
||||
fs.unlinkSync(pidFile);
|
||||
} catch (error) { if ((error as NodeJS.ErrnoException).code !== "ENOENT") throw error; }
|
||||
}
|
||||
|
||||
function removePidIfOwned(pidFile: string, pid: number): void {
|
||||
try {
|
||||
if (fs.readFileSync(pidFile, "utf8").trim() === String(pid)) fs.unlinkSync(pidFile);
|
||||
} catch (error) { if ((error as NodeJS.ErrnoException).code !== "ENOENT") throw error; }
|
||||
}
|
||||
|
||||
function atomicWriteMode(file: string, content: string, mode: number): void {
|
||||
const temp = `${file}.${process.pid}.${Date.now()}.tmp`;
|
||||
let fd: number | undefined;
|
||||
try {
|
||||
fd = fs.openSync(temp, "wx", mode);
|
||||
fs.writeFileSync(fd, content, "utf8");
|
||||
fs.fsyncSync(fd);
|
||||
fs.closeSync(fd);
|
||||
fd = undefined;
|
||||
fs.linkSync(temp, file);
|
||||
fs.unlinkSync(temp);
|
||||
fs.chmodSync(file, mode);
|
||||
const dirFd = fs.openSync(path.dirname(file), "r");
|
||||
try { fs.fsyncSync(dirFd); } finally { fs.closeSync(dirFd); }
|
||||
} catch (error) {
|
||||
if (fd !== undefined) try { fs.closeSync(fd); } catch { /* ignore cleanup failure */ }
|
||||
try { fs.unlinkSync(temp); } catch { /* ignore cleanup failure */ }
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
function assertPortAvailable(host: string, port: number): Promise<void> {
|
||||
return new Promise((resolve, reject) => {
|
||||
const server = net.createServer();
|
||||
server.once("error", (error: NodeJS.ErrnoException) => reject(new Error(`Cannot start: ${host}:${port} is unavailable (${error.code || error.message})`)));
|
||||
server.listen(port, host, () => server.close((error) => error ? reject(error) : resolve()));
|
||||
});
|
||||
}
|
||||
|
||||
async function waitForHealthyInstance(pid: number, platform: string, botId: string, config: Parameters<typeof localBaseUrl>[0]): Promise<void> {
|
||||
let lastError = "health endpoint not ready";
|
||||
for (let attempt = 0; attempt < 30; attempt++) {
|
||||
if (!isPidRunning(pid)) throw new Error("child process exited before becoming healthy");
|
||||
try {
|
||||
const response = await fetch(`${localBaseUrl(config)}/health`, { signal: AbortSignal.timeout(500) });
|
||||
const health = await response.json() as { botId?: string; platform?: string };
|
||||
if (response.ok && health.botId === botId && health.platform === platform) return;
|
||||
lastError = `health identity mismatch (HTTP ${response.status})`;
|
||||
} catch (error) { lastError = error instanceof Error ? error.message : String(error); }
|
||||
await delay(100);
|
||||
}
|
||||
throw new Error(`health check timed out: ${lastError}`);
|
||||
}
|
||||
|
||||
function delay(ms: number): Promise<void> { return new Promise((resolve) => setTimeout(resolve, ms)); }
|
||||
|
||||
async function terminateStartedChild(pid: number): Promise<void> {
|
||||
if (!isPidRunning(pid)) return;
|
||||
try { process.kill(pid, "SIGTERM"); } catch { return; }
|
||||
for (let attempt = 0; attempt < 20 && isPidRunning(pid); attempt++) await delay(50);
|
||||
if (isPidRunning(pid)) try { process.kill(pid, "SIGKILL"); } catch { /* process already exited */ }
|
||||
}
|
||||
|
||||
function withInstanceHome<T>(home: string, fn: () => T): T {
|
||||
const previous = process.env.GORI_AGENT_HOME;
|
||||
process.env.GORI_AGENT_HOME = home;
|
||||
try { return fn(); } finally { restoreInstanceHome(previous); }
|
||||
}
|
||||
|
||||
async function withInstanceHomeAsync<T>(home: string, fn: () => Promise<T>): Promise<T> {
|
||||
const previous = process.env.GORI_AGENT_HOME;
|
||||
process.env.GORI_AGENT_HOME = home;
|
||||
try { return await fn(); } finally { restoreInstanceHome(previous); }
|
||||
}
|
||||
|
||||
function restoreInstanceHome(previous: string | undefined): void {
|
||||
if (previous === undefined) delete process.env.GORI_AGENT_HOME;
|
||||
else process.env.GORI_AGENT_HOME = previous;
|
||||
}
|
||||
+10
-17
@@ -2,48 +2,41 @@ import os from "node:os";
|
||||
import type { AppConfig } from "../config.js";
|
||||
|
||||
export function localBaseUrl(config: AppConfig): string {
|
||||
return `http://localhost:${config.server.port}`;
|
||||
const configured = config.gateway.server.host;
|
||||
const host = configured === "0.0.0.0" ? "127.0.0.1" : configured === "::" ? "[::1]" : configured.includes(":") && !configured.startsWith("[") ? `[${configured}]` : configured;
|
||||
return `http://${host}:${config.gateway.server.port}`;
|
||||
}
|
||||
|
||||
export function lanBaseUrls(config: AppConfig): string[] {
|
||||
const urls: string[] = [];
|
||||
for (const interfaces of Object.values(os.networkInterfaces())) {
|
||||
for (const item of interfaces || []) {
|
||||
if (item.family === "IPv4" && !item.internal) {
|
||||
urls.push(`http://${item.address}:${config.server.port}`);
|
||||
}
|
||||
if (item.family === "IPv4" && !item.internal) urls.push(`http://${item.address}:${config.gateway.server.port}`);
|
||||
}
|
||||
}
|
||||
return urls;
|
||||
}
|
||||
|
||||
export async function publicBaseUrlHint(config: AppConfig): Promise<string | undefined> {
|
||||
if (config.server.publicBaseUrl) return config.server.publicBaseUrl.replace(/\/$/, "");
|
||||
|
||||
if (config.gateway.server.publicBaseUrl) return config.gateway.server.publicBaseUrl.replace(/\/$/, "");
|
||||
const controller = new AbortController();
|
||||
const timer = setTimeout(() => controller.abort(), 1_500);
|
||||
try {
|
||||
const response = await fetch("https://api.ipify.org?format=text", { signal: controller.signal });
|
||||
if (!response.ok) return undefined;
|
||||
const ip = (await response.text()).trim();
|
||||
if (!ip) return undefined;
|
||||
return `http://${ip}:${config.server.port}`;
|
||||
} catch {
|
||||
return undefined;
|
||||
} finally {
|
||||
clearTimeout(timer);
|
||||
}
|
||||
return ip ? `http://${ip}:${config.gateway.server.port}` : undefined;
|
||||
} catch { return undefined; } finally { clearTimeout(timer); }
|
||||
}
|
||||
|
||||
export async function printUrlHints(config: AppConfig): Promise<void> {
|
||||
console.log(`Local: ${localBaseUrl(config)}`);
|
||||
const lanUrls = lanBaseUrls(config);
|
||||
if (lanUrls.length > 0) console.log(`LAN: ${lanUrls.join(", ")}`);
|
||||
const publicHint = await publicBaseUrlHint(config);
|
||||
if (publicHint) console.log(`Public hint: ${publicHint}`);
|
||||
if (!config.server.publicBaseUrl) console.log("Set server.publicBaseUrl when exposing through HTTPS/reverse proxy.");
|
||||
if (config.gateway.server.publicBaseUrl) console.log(`Public: ${config.gateway.server.publicBaseUrl.replace(/\/$/, "")}`);
|
||||
else console.log("Set gateway.server.publicBaseUrl when exposing through HTTPS/reverse proxy.");
|
||||
}
|
||||
|
||||
export async function webhookBaseUrl(config: AppConfig): Promise<string> {
|
||||
return config.server.publicBaseUrl?.replace(/\/$/, "") || await publicBaseUrlHint(config) || localBaseUrl(config);
|
||||
return config.gateway.server.publicBaseUrl?.replace(/\/$/, "") || await publicBaseUrlHint(config) || localBaseUrl(config);
|
||||
}
|
||||
|
||||
+3
-7
@@ -2,17 +2,13 @@ import type { AppConfig } from "../config.js";
|
||||
import { webhookBaseUrl } from "./net.js";
|
||||
|
||||
export async function printFeishu(config: AppConfig): Promise<void> {
|
||||
if (config.gateway.platform.type !== "feishu") throw new Error("Current instance platform is not Feishu");
|
||||
const baseUrl = await webhookBaseUrl(config);
|
||||
console.log("Feishu/Lark setup instructions");
|
||||
console.log("");
|
||||
console.log(`Webhook URL: ${baseUrl}/webhook/feishu`);
|
||||
console.log("Event subscription: im.message.receive_v1");
|
||||
console.log("Required config fields: platforms.feishu.appId, appSecret, verificationToken, botNames");
|
||||
console.log("Required config fields: gateway.platform.appId, appSecret, verificationToken, botNames");
|
||||
console.log("Do not paste appSecret into chats or logs.");
|
||||
if (!config.platforms.feishu.enabled) {
|
||||
console.log("Warning: platforms.feishu.enabled is false in this config.");
|
||||
}
|
||||
if (!config.server.publicBaseUrl) {
|
||||
console.log("Warning: server.publicBaseUrl is empty; configure your public HTTPS URL before production use.");
|
||||
}
|
||||
if (!config.gateway.server.publicBaseUrl) console.log("Warning: gateway.server.publicBaseUrl is empty; configure your public HTTPS URL before production use.");
|
||||
}
|
||||
|
||||
+20
-1
@@ -1,8 +1,10 @@
|
||||
import readline from "node:readline/promises";
|
||||
import { stdin as input, stdout as output } from "node:process";
|
||||
import { Writable } from "node:stream";
|
||||
|
||||
export interface PromptSession {
|
||||
ask(question: string, defaultValue?: string): Promise<string>;
|
||||
askSecret(question: string, hasExisting?: boolean): Promise<string>;
|
||||
askBoolean(question: string, defaultValue?: boolean): Promise<boolean>;
|
||||
askList(question: string, defaultValues?: string[]): Promise<string[]>;
|
||||
choose<T>(question: string, choices: Choice<T>[], defaultIndex?: number): Promise<T>;
|
||||
@@ -16,9 +18,14 @@ export interface Choice<T> {
|
||||
}
|
||||
|
||||
export function createPromptSession(): PromptSession {
|
||||
const rl = readline.createInterface({ input, output });
|
||||
let rl = readline.createInterface({ input, output });
|
||||
return {
|
||||
ask: (question, defaultValue) => ask(rl, question, defaultValue),
|
||||
askSecret: async (question, hasExisting) => {
|
||||
rl.close();
|
||||
try { return await askSecret(question, hasExisting); }
|
||||
finally { rl = readline.createInterface({ input, output }); }
|
||||
},
|
||||
askBoolean: (question, defaultValue) => askBoolean(rl, question, defaultValue),
|
||||
askList: (question, defaultValues) => askList(rl, question, defaultValues),
|
||||
choose: (question, choices, defaultIndex) => choose(rl, question, choices, defaultIndex),
|
||||
@@ -26,6 +33,18 @@ export function createPromptSession(): PromptSession {
|
||||
};
|
||||
}
|
||||
|
||||
async function askSecret(question: string, hasExisting = false): Promise<string> {
|
||||
if (!input.isTTY || !output.isTTY) throw new Error("Secret input requires an interactive terminal");
|
||||
const muted = new Writable({ write(_chunk, _encoding, callback) { callback(); } });
|
||||
const secretRl = readline.createInterface({ input, output: muted, terminal: true });
|
||||
output.write(`${question}${hasExisting ? " (leave blank to keep existing)" : ""}: `);
|
||||
try {
|
||||
const answer = (await secretRl.question("")).trim();
|
||||
output.write("\n");
|
||||
return answer;
|
||||
} finally { secretRl.close(); }
|
||||
}
|
||||
|
||||
export async function ask(rl: readline.Interface, question: string, defaultValue?: string): Promise<string> {
|
||||
const suffix = defaultValue !== undefined && defaultValue !== "" ? ` [${defaultValue}]` : "";
|
||||
const answer = (await rl.question(`${question}${suffix}: `)).trim();
|
||||
|
||||
+7
-21
@@ -1,26 +1,12 @@
|
||||
import type { AppConfig } from "../config.js";
|
||||
import type { FeishuConfig } from "../config.js";
|
||||
import type { PromptSession } from "./prompt.js";
|
||||
|
||||
export async function configureFeishu(
|
||||
prompt: PromptSession,
|
||||
existing: AppConfig["platforms"]["feishu"]
|
||||
): Promise<AppConfig["platforms"]["feishu"]> {
|
||||
console.log("\nFeishu/Lark setup");
|
||||
console.log("Create a Feishu/Lark app, enable bot messaging, and subscribe to im.message.receive_v1.");
|
||||
|
||||
const appId = await prompt.ask("App ID", existing.appId && existing.appId !== "cli_xxx" ? existing.appId : undefined);
|
||||
const appSecretAnswer = await prompt.ask(existing.appSecret ? "App Secret (leave blank to keep existing)" : "App Secret");
|
||||
const verificationToken = await prompt.ask(
|
||||
"Verification token",
|
||||
existing.verificationToken && existing.verificationToken !== "replace-me" ? existing.verificationToken : undefined
|
||||
);
|
||||
const botNames = await prompt.askList("Bot display names, comma-separated", existing.botNames);
|
||||
|
||||
export async function configureFeishu(prompt: PromptSession, existing: FeishuConfig): Promise<FeishuConfig> {
|
||||
return {
|
||||
enabled: true,
|
||||
appId,
|
||||
appSecret: appSecretAnswer || existing.appSecret,
|
||||
verificationToken,
|
||||
botNames
|
||||
type: "feishu",
|
||||
appId: await prompt.ask("App ID", existing.appId),
|
||||
appSecret: await prompt.askSecret("App Secret", Boolean(existing.appSecret)) || existing.appSecret,
|
||||
verificationToken: await prompt.askSecret("Verification token", Boolean(existing.verificationToken)) || existing.verificationToken,
|
||||
botNames: await prompt.askList("Bot display names, comma-separated", existing.botNames)
|
||||
};
|
||||
}
|
||||
|
||||
@@ -0,0 +1,52 @@
|
||||
import fs from "node:fs";
|
||||
import os from "node:os";
|
||||
import path from "node:path";
|
||||
import process from "node:process";
|
||||
import { isValidServerHost, parseConfig } from "../config.js";
|
||||
|
||||
export function parseServerHost(input: string): string {
|
||||
const host = input.trim();
|
||||
if (!isValidServerHost(host)) throw new Error("server host must be a hostname, IPv4 address, or unbracketed IPv6 address without a port");
|
||||
return host;
|
||||
}
|
||||
|
||||
export function parseServerPort(input: string): number {
|
||||
const value = input.trim();
|
||||
if (!/^[0-9]+$/.test(value)) throw new Error("server port must be a decimal integer from 1 to 65535");
|
||||
const port = Number(value);
|
||||
if (!Number.isSafeInteger(port) || port < 1 || port > 65_535) throw new Error("server port must be a decimal integer from 1 to 65535");
|
||||
return port;
|
||||
}
|
||||
|
||||
export function suggestAvailablePort(preferred: number, usedPorts: ReadonlySet<number>): number {
|
||||
for (let port = preferred; port <= 65_535; port++) {
|
||||
if (!usedPorts.has(port)) return port;
|
||||
}
|
||||
throw new Error(`No unassigned instance port is available at or above ${preferred}`);
|
||||
}
|
||||
|
||||
export function defaultInstancesDirectory(): string {
|
||||
const root = path.resolve(process.env.GORI_AGENT_ROOT || path.join(os.homedir(), ".gori-agent"));
|
||||
return path.join(root, "instances");
|
||||
}
|
||||
|
||||
export function collectUsedInstancePorts(instancesDirectory: string, excludeConfigPath?: string): Set<number> {
|
||||
const used = new Set<number>();
|
||||
const excluded = excludeConfigPath ? path.resolve(excludeConfigPath) : undefined;
|
||||
let entries: fs.Dirent[];
|
||||
try { entries = fs.readdirSync(instancesDirectory, { withFileTypes: true }); }
|
||||
catch { return used; }
|
||||
|
||||
for (const entry of entries) {
|
||||
if (!entry.isDirectory() || entry.isSymbolicLink()) continue;
|
||||
const configFile = path.join(instancesDirectory, entry.name, "config.json");
|
||||
if (excluded && path.resolve(configFile) === excluded) continue;
|
||||
try {
|
||||
const stat = fs.lstatSync(configFile);
|
||||
if (!stat.isFile() || stat.isSymbolicLink()) continue;
|
||||
const config = parseConfig(JSON.parse(fs.readFileSync(configFile, "utf8")) as unknown);
|
||||
used.add(config.gateway.server.port);
|
||||
} catch { /* Ignore invalid or unreadable peer instances. */ }
|
||||
}
|
||||
return used;
|
||||
}
|
||||
@@ -1,15 +1,8 @@
|
||||
import crypto from "node:crypto";
|
||||
import type { AppConfig } from "../config.js";
|
||||
import type { WebhookConfig } from "../config.js";
|
||||
import type { PromptSession } from "./prompt.js";
|
||||
|
||||
export async function configureGenericWebhook(
|
||||
prompt: PromptSession,
|
||||
existing: AppConfig["platforms"]["webhook"]
|
||||
): Promise<AppConfig["platforms"]["webhook"]> {
|
||||
console.log("\nGeneric webhook setup");
|
||||
console.log("Use POST /webhook/generic with optional X-Gori-Signature HMAC-SHA256 authentication.");
|
||||
|
||||
const generated = existing.secret && existing.secret !== "replace-me" ? existing.secret : crypto.randomBytes(24).toString("hex");
|
||||
const secret = await prompt.ask("Webhook secret", generated);
|
||||
return { enabled: true, secret };
|
||||
export async function configureGenericWebhook(prompt: PromptSession, existing: WebhookConfig): Promise<WebhookConfig> {
|
||||
const answer = await prompt.askSecret("Webhook secret (leave blank to generate)", Boolean(existing.secret));
|
||||
return { type: "webhook", secret: answer || existing.secret || crypto.randomBytes(24).toString("hex") };
|
||||
}
|
||||
|
||||
+4
-11
@@ -1,15 +1,8 @@
|
||||
import crypto from "node:crypto";
|
||||
import type { AppConfig } from "../config.js";
|
||||
import type { WeixinConfig } from "../config.js";
|
||||
import type { PromptSession } from "./prompt.js";
|
||||
|
||||
export async function configureWeixin(
|
||||
prompt: PromptSession,
|
||||
existing: AppConfig["platforms"]["weixin"]
|
||||
): Promise<AppConfig["platforms"]["weixin"]> {
|
||||
console.log("\nPersonal WeChat external webhook setup");
|
||||
console.log("Native personal WeChat integration is not included; use an external bridge that POSTs to /webhook/weixin.");
|
||||
|
||||
const generated = existing.secret && existing.secret !== "replace-me" ? existing.secret : crypto.randomBytes(24).toString("hex");
|
||||
const secret = await prompt.ask("Bridge secret/reference", generated);
|
||||
return { enabled: true, mode: "external-webhook", secret };
|
||||
export async function configureWeixin(prompt: PromptSession, existing: WeixinConfig): Promise<WeixinConfig> {
|
||||
const answer = await prompt.askSecret("Bridge secret (leave blank to generate)", Boolean(existing.secret));
|
||||
return { type: "weixin", mode: "external-webhook", secret: answer || existing.secret || crypto.randomBytes(24).toString("hex") };
|
||||
}
|
||||
|
||||
+161
-53
@@ -1,66 +1,174 @@
|
||||
import crypto from "node:crypto";
|
||||
import fs from "node:fs";
|
||||
import path from "node:path";
|
||||
import { discoverBackends } from "../acp/discovery.js";
|
||||
import type { AppConfig } from "../config.js";
|
||||
import { createPromptSession } from "./prompt.js";
|
||||
import { loadConfigFile, projectRoot, writeConfigFile } from "./config-file.js";
|
||||
import type { AppConfig, PlatformConfig } from "../config.js";
|
||||
import { createPromptSession, type PromptSession } from "./prompt.js";
|
||||
import { loadConfigFile, loadExampleConfig, projectRoot, writeConfigFile } from "./config-file.js";
|
||||
import { collectUsedInstancePorts, defaultInstancesDirectory, parseServerHost, parseServerPort, suggestAvailablePort } from "./setup-server.js";
|
||||
|
||||
const GORI_SKILL = "/home/ubuntu/gori-space/gori-deploy/.kimi-code/skills/gori-update/SKILL.md";
|
||||
export interface SetupOptions {
|
||||
botId?: string;
|
||||
requireNew?: boolean;
|
||||
writeWithoutConfirmation?: boolean;
|
||||
prompt?: PromptSession;
|
||||
discoverAgents?: typeof discoverBackends;
|
||||
instancesDirectory?: string;
|
||||
log?: (message: string) => void;
|
||||
}
|
||||
|
||||
export async function runSetup(configPath?: string): Promise<void> {
|
||||
const loaded = loadConfigFile(configPath);
|
||||
const prompt = createPromptSession();
|
||||
export async function runSetup(configPath?: string, options: SetupOptions = {}): Promise<void> {
|
||||
const target = path.resolve(configPath || "config.json");
|
||||
const configExists = fs.existsSync(target);
|
||||
if (options.requireNew && configExists) throw new Error(`Config already exists: ${target}`);
|
||||
const current = configExists ? loadConfigFile(target).config : loadExampleConfig();
|
||||
const prompt = options.prompt || createPromptSession();
|
||||
const discoverAgents = options.discoverAgents || discoverBackends;
|
||||
const log = options.log || console.log;
|
||||
try {
|
||||
console.log("gori-agent ACP setup");
|
||||
console.log(`Config target: ${loaded.path}`);
|
||||
const discovered = await discoverBackends();
|
||||
for (const backend of discovered) console.log(`- ${backend.id}: ${backend.status}${backend.version ? ` (${backend.version})` : ""}${backend.reason ? ` - ${backend.reason}` : ""}`);
|
||||
const kimi = discovered.find((backend) => backend.id === "kimi" && backend.status === "ready");
|
||||
if (!kimi) throw new Error("Kimi ACP backend is not available");
|
||||
log("gori-agent Config v3 setup");
|
||||
log(`Config target: ${target}`);
|
||||
const discovered = await discoverAgents();
|
||||
for (const agent of discovered) log(`- ${agent.id}: ${agent.status}${agent.version ? ` (${agent.version})` : ""}${agent.reason ? ` - ${agent.reason}` : ""}`);
|
||||
const ready = discovered.filter((agent) => agent.status === "ready");
|
||||
if (ready.length === 0) throw new Error("No ACP agent is available");
|
||||
const selected = ready.find((agent) => agent.id === current.bot.agent.id) || ready[0];
|
||||
|
||||
const existingRole = loaded.config.roles.find((role) => role.id === loaded.config.defaultRole);
|
||||
const workspace = path.resolve(await prompt.ask("Assistant workspace", existingRole?.workspace || projectRoot()));
|
||||
const includeOps = await prompt.askBoolean("Include ops role with the existing gori-update skill", loaded.config.roles.some((role) => role.id === "ops"));
|
||||
const roles: AppConfig["roles"] = [{
|
||||
id: "assistant", backend: "kimi", workspace, persona: existingRole?.persona || "", skills: [],
|
||||
policy: { permissionMode: "deny", allowedTools: [], allowedCommandPatterns: [] }
|
||||
}];
|
||||
const skills: AppConfig["skills"] = [];
|
||||
if (includeOps) {
|
||||
skills.push({ id: "gori-update", file: GORI_SKILL, maxBytes: 256_000 });
|
||||
roles.push(opsRole());
|
||||
}
|
||||
const botId = options.botId || await prompt.ask("Bot ID", current.bot.id === "BOT_ID" ? undefined : current.bot.id);
|
||||
const workspace = path.resolve(await prompt.ask("Bot workspace", current.bot.workspace.startsWith("/absolute/") ? projectRoot() : current.bot.workspace));
|
||||
const persona = await prompt.ask("Bot persona", current.bot.persona);
|
||||
const usedPorts = collectUsedInstancePorts(options.instancesDirectory || defaultInstancesDirectory(), target);
|
||||
const server = await configureServer(prompt, current.gateway.server, usedPorts, !configExists, log);
|
||||
const platformType = await prompt.choose("Platform", [
|
||||
{ label: "QQ", value: "qq" as const },
|
||||
{ label: "Feishu/Lark", value: "feishu" as const },
|
||||
{ label: "WeCom", value: "wecom" as const },
|
||||
{ label: "Generic webhook", value: "webhook" as const },
|
||||
{ label: "Weixin bridge", value: "weixin" as const }
|
||||
], ["qq", "feishu", "wecom", "webhook", "weixin"].indexOf(current.gateway.platform.type));
|
||||
const platform = await configurePlatform(prompt, platformType, current.gateway.platform);
|
||||
const isTemplate = current.bot.id === "BOT_ID";
|
||||
const next: AppConfig = {
|
||||
...loaded.config,
|
||||
configVersion: 2,
|
||||
backends: [{ id: "kimi", command: kimi.command, args: ["acp"], env: {} }],
|
||||
skills,
|
||||
defaultRole: "assistant",
|
||||
roles,
|
||||
platforms: structuredClone(loaded.config.platforms)
|
||||
configVersion: 3,
|
||||
bot: {
|
||||
id: botId,
|
||||
workspace,
|
||||
persona,
|
||||
agent: {
|
||||
id: selected.id,
|
||||
command: selected.command,
|
||||
args: selected.args,
|
||||
env: !isTemplate && selected.id === current.bot.agent.id ? current.bot.agent.env : {}
|
||||
},
|
||||
skills: isTemplate ? [] : current.bot.skills,
|
||||
permissions: isTemplate ? { mode: "deny", allowedTools: [], allowedCommandPatterns: [] } : current.bot.permissions
|
||||
},
|
||||
gateway: { server, policy: current.gateway.policy, platform },
|
||||
runtime: current.runtime
|
||||
};
|
||||
console.log("\nMigration summary:");
|
||||
console.log(`- default role: ${next.defaultRole}`);
|
||||
console.log(`- roles: ${next.roles.map(({ id }) => id).join(", ")}`);
|
||||
console.log(`- backend: ${kimi.command} acp`);
|
||||
console.log("- all existing platform settings and credentials are preserved");
|
||||
if (await prompt.askBoolean(`Write config to ${loaded.path}`, false)) {
|
||||
writeConfigFile(loaded.path, next);
|
||||
console.log(`Wrote ${loaded.path}`);
|
||||
} else console.log("No changes written.");
|
||||
const shouldWrite = options.writeWithoutConfirmation || await prompt.askBoolean(`Write Config v3 to ${target}`, false);
|
||||
if (!shouldWrite) { log("No changes written."); return; }
|
||||
writeConfigFile(target, next);
|
||||
log(`Wrote ${target} with mode 0600`);
|
||||
} finally { prompt.close(); }
|
||||
}
|
||||
|
||||
function opsRole(): AppConfig["roles"][number] {
|
||||
return {
|
||||
id: "ops", backend: "kimi", workspace: "/home/ubuntu/gori-space",
|
||||
persona: "你是 Gori 团队运维角色。严格遵循 gori-update skill;有风险或需要外部确认时停止并报告。",
|
||||
skills: ["gori-update"],
|
||||
policy: {
|
||||
permissionMode: "allowlist",
|
||||
allowedTools: ["read", "grep", "glob", "bash"],
|
||||
allowedCommandPatterns: [
|
||||
"^(?:.*\\\"command\\\":\\\")?(?:git (?:status|log|diff|pull --ff-only)|bash gori-deploy/(?:build\\.sh|dist/deploy-[a-z-]+\\.sh)|docker (?:ps|logs))"
|
||||
]
|
||||
async function configureServer(
|
||||
prompt: PromptSession,
|
||||
existing: AppConfig["gateway"]["server"],
|
||||
usedPorts: ReadonlySet<number>,
|
||||
isNew: boolean,
|
||||
log: (message: string) => void
|
||||
): Promise<AppConfig["gateway"]["server"]> {
|
||||
let host: string;
|
||||
for (;;) {
|
||||
try {
|
||||
host = parseServerHost(await prompt.ask("Gateway server host", existing.host));
|
||||
break;
|
||||
} catch (error) {
|
||||
log(`ERROR: ${error instanceof Error ? error.message : String(error)}`);
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
const preferredPort = existing.port;
|
||||
let suggestedPort = preferredPort;
|
||||
if (usedPorts.has(preferredPort)) {
|
||||
try {
|
||||
suggestedPort = suggestAvailablePort(preferredPort, usedPorts);
|
||||
log(`WARNING: Gateway server port ${preferredPort} is already declared by another instance; suggested port: ${suggestedPort}.`);
|
||||
} catch (error) {
|
||||
if (isNew) throw error;
|
||||
log(`WARNING: Gateway server port ${preferredPort} is already declared by another instance, and no higher unassigned port is available.`);
|
||||
}
|
||||
}
|
||||
|
||||
let port: number;
|
||||
for (;;) {
|
||||
try {
|
||||
port = parseServerPort(await prompt.ask("Gateway server port", String(isNew ? suggestedPort : preferredPort)));
|
||||
break;
|
||||
} catch (error) {
|
||||
log(`ERROR: ${error instanceof Error ? error.message : String(error)}`);
|
||||
}
|
||||
}
|
||||
if (usedPorts.has(port)) {
|
||||
log(`WARNING: Selected gateway server port ${port} is already declared by another instance; start will fail if that port is in use.`);
|
||||
}
|
||||
|
||||
return { host, port, publicBaseUrl: existing.publicBaseUrl };
|
||||
}
|
||||
|
||||
async function configurePlatform(
|
||||
prompt: PromptSession,
|
||||
type: PlatformConfig["type"],
|
||||
existing: PlatformConfig
|
||||
): Promise<PlatformConfig> {
|
||||
switch (type) {
|
||||
case "qq": {
|
||||
const previous = existing.type === "qq" ? existing : undefined;
|
||||
return {
|
||||
type,
|
||||
connectionMode: await prompt.choose("QQ connection mode", [
|
||||
{ label: "WebSocket", value: "websocket" as const },
|
||||
{ label: "Webhook", value: "webhook" as const }
|
||||
], previous?.connectionMode === "webhook" ? 1 : 0),
|
||||
appId: await prompt.ask("QQ App ID", previous?.appId),
|
||||
clientSecret: await prompt.askSecret("QQ Client Secret", Boolean(previous?.clientSecret)) || previous?.clientSecret || "",
|
||||
botSecret: await prompt.askSecret("QQ Bot Secret", Boolean(previous?.botSecret)) || previous?.botSecret || "",
|
||||
verifySignature: previous?.verifySignature ?? true,
|
||||
botNames: await prompt.askList("QQ Bot names", previous?.botNames),
|
||||
intents: previous?.intents ?? 33_554_432,
|
||||
shard: previous?.shard ?? [0, 1]
|
||||
};
|
||||
}
|
||||
case "feishu": {
|
||||
const previous = existing.type === "feishu" ? existing : undefined;
|
||||
return {
|
||||
type,
|
||||
appId: await prompt.ask("Feishu App ID", previous?.appId),
|
||||
appSecret: await prompt.askSecret("Feishu App Secret", Boolean(previous?.appSecret)) || previous?.appSecret || "",
|
||||
verificationToken: await prompt.askSecret("Feishu verification token", Boolean(previous?.verificationToken)) || previous?.verificationToken || "",
|
||||
botNames: await prompt.askList("Feishu Bot names", previous?.botNames)
|
||||
};
|
||||
}
|
||||
case "wecom": {
|
||||
const previous = existing.type === "wecom" ? existing : undefined;
|
||||
return {
|
||||
type,
|
||||
corpId: await prompt.ask("WeCom Corp ID", previous?.corpId),
|
||||
agentId: await prompt.ask("WeCom Agent ID", previous?.agentId),
|
||||
secret: await prompt.askSecret("WeCom Secret", Boolean(previous?.secret)) || previous?.secret || ""
|
||||
};
|
||||
}
|
||||
case "webhook": {
|
||||
const previous = existing.type === "webhook" ? existing : undefined;
|
||||
const secret = await prompt.askSecret("Webhook secret", Boolean(previous?.secret));
|
||||
return { type, secret: secret || previous?.secret || crypto.randomBytes(24).toString("hex") };
|
||||
}
|
||||
case "weixin": {
|
||||
const previous = existing.type === "weixin" ? existing : undefined;
|
||||
const secret = await prompt.askSecret("Weixin bridge secret", Boolean(previous?.secret));
|
||||
return { type, mode: "external-webhook", secret: secret || previous?.secret || crypto.randomBytes(24).toString("hex") };
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+103
-106
@@ -1,95 +1,120 @@
|
||||
import fs from "node:fs";
|
||||
import net from "node:net";
|
||||
import path from "node:path";
|
||||
import process from "node:process";
|
||||
import { z } from "zod";
|
||||
|
||||
const policySchema = z.object({
|
||||
const gatewayPolicySchema = z.object({
|
||||
allowedUsers: z.array(z.string()).default([]),
|
||||
allowedChats: z.array(z.string()).default([]),
|
||||
requireMentionInGroup: z.boolean().default(false)
|
||||
});
|
||||
requireMentionInGroup: z.boolean().default(true)
|
||||
}).strict();
|
||||
|
||||
const rolePolicySchema = z.object({
|
||||
permissionMode: z.enum(["deny", "allowlist", "auto"]).default("deny"),
|
||||
const permissionPolicySchema = z.object({
|
||||
mode: z.enum(["deny", "allowlist", "auto"]).default("deny"),
|
||||
allowedTools: z.array(z.string()).default([]),
|
||||
allowedCommandPatterns: z.array(z.string()).default([])
|
||||
});
|
||||
}).strict();
|
||||
|
||||
const backendSchema = z.object({
|
||||
const agentSchema = z.object({
|
||||
id: z.string().min(1),
|
||||
command: z.string().min(1),
|
||||
args: z.array(z.string()).default([]),
|
||||
env: z.record(z.string()).default({})
|
||||
});
|
||||
}).strict();
|
||||
|
||||
const skillSchema = z.object({
|
||||
id: z.string().min(1),
|
||||
file: z.string().min(1),
|
||||
maxBytes: z.number().int().positive().default(256_000)
|
||||
});
|
||||
}).strict();
|
||||
|
||||
const roleSchema = z.object({
|
||||
id: z.string().min(1),
|
||||
backend: z.string().min(1),
|
||||
const botIdSchema = z.union([
|
||||
z.literal("BOT_ID"),
|
||||
z.string().regex(/^[a-z0-9](?:[a-z0-9-]{0,62})$/, "bot.id must contain only lowercase letters, digits, and hyphens")
|
||||
]);
|
||||
const botSchema = z.object({
|
||||
id: botIdSchema,
|
||||
workspace: z.string().min(1),
|
||||
persona: z.string().default(""),
|
||||
skills: z.array(z.string()).default([]),
|
||||
policy: rolePolicySchema.default({})
|
||||
});
|
||||
agent: agentSchema,
|
||||
skills: z.array(skillSchema).default([]),
|
||||
permissions: permissionPolicySchema.default({})
|
||||
}).strict();
|
||||
|
||||
const serverSchema = z.object({
|
||||
host: z.string().refine(isValidServerHost, "server host must be a hostname, IPv4 address, or unbracketed IPv6 address without a port").default("0.0.0.0"),
|
||||
port: z.number().int().positive().max(65_535).default(8787),
|
||||
publicBaseUrl: z.string().default("")
|
||||
}).strict();
|
||||
|
||||
const feishuSchema = z.object({
|
||||
type: z.literal("feishu"),
|
||||
appId: z.string().default(""),
|
||||
appSecret: z.string().default(""),
|
||||
verificationToken: z.string().default(""),
|
||||
botNames: z.array(z.string()).default([])
|
||||
}).strict();
|
||||
const wecomSchema = z.object({
|
||||
type: z.literal("wecom"),
|
||||
corpId: z.string().default(""),
|
||||
agentId: z.string().default(""),
|
||||
secret: z.string().default("")
|
||||
}).strict();
|
||||
const qqSchema = z.object({
|
||||
type: z.literal("qq"),
|
||||
connectionMode: z.enum(["websocket", "webhook"]).default("websocket"),
|
||||
appId: z.string().default(""),
|
||||
clientSecret: z.string().default(""),
|
||||
botSecret: z.string().default(""),
|
||||
verifySignature: z.boolean().default(true),
|
||||
botNames: z.array(z.string()).default([]),
|
||||
intents: z.number().int().positive().default(1 << 25),
|
||||
shard: z.tuple([z.number().int().nonnegative(), z.number().int().positive()]).default([0, 1])
|
||||
}).strict();
|
||||
const webhookSchema = z.object({ type: z.literal("webhook"), secret: z.string().default("") }).strict();
|
||||
const weixinSchema = z.object({
|
||||
type: z.literal("weixin"),
|
||||
mode: z.enum(["external-webhook", "not-implemented"]).default("external-webhook"),
|
||||
secret: z.string().default("")
|
||||
}).strict();
|
||||
const platformSchema = z.discriminatedUnion("type", [feishuSchema, wecomSchema, qqSchema, webhookSchema, weixinSchema]);
|
||||
|
||||
const acpSchema = z.object({
|
||||
stateFile: z.string().default(""),
|
||||
initializeTimeoutMs: z.number().int().positive().default(10_000),
|
||||
promptTimeoutMs: z.number().int().positive().default(600_000),
|
||||
promptTimeoutMs: z.number().int().positive().default(7_200_000),
|
||||
cancelGraceMs: z.number().int().positive().default(5_000),
|
||||
idleTimeoutMs: z.number().int().positive().default(1_800_000),
|
||||
sweepIntervalMs: z.number().int().positive().default(60_000),
|
||||
maxProcesses: z.number().int().positive().default(8)
|
||||
});
|
||||
|
||||
const feishuSchema = z.object({
|
||||
enabled: z.boolean().default(false), appId: z.string().default(""), appSecret: z.string().default(""),
|
||||
verificationToken: z.string().default(""), botNames: z.array(z.string()).default([])
|
||||
});
|
||||
const wecomSchema = z.object({
|
||||
enabled: z.boolean().default(false), corpId: z.string().default(""), agentId: z.string().default(""), secret: z.string().default("")
|
||||
});
|
||||
const qqSchema = z.object({
|
||||
enabled: z.boolean().default(false), connectionMode: z.enum(["websocket", "webhook"]).default("websocket"),
|
||||
appId: z.string().default(""), clientSecret: z.string().default(""), botSecret: z.string().default(""),
|
||||
verifySignature: z.boolean().default(true), botNames: z.array(z.string()).default([]),
|
||||
intents: z.number().int().positive().default(1 << 25),
|
||||
shard: z.tuple([z.number().int().nonnegative(), z.number().int().positive()]).default([0, 1])
|
||||
});
|
||||
const webhookSchema = z.object({ enabled: z.boolean().default(true), secret: z.string().default("") });
|
||||
const weixinSchema = z.object({
|
||||
enabled: z.boolean().default(false), mode: z.enum(["external-webhook", "not-implemented"]).default("external-webhook"), secret: z.string().default("")
|
||||
});
|
||||
}).strict();
|
||||
|
||||
export const configSchema = z.object({
|
||||
configVersion: z.literal(2),
|
||||
server: z.object({
|
||||
host: z.string().default("0.0.0.0"), port: z.number().int().positive().max(65_535).default(3000), publicBaseUrl: z.string().default("")
|
||||
}).default({}),
|
||||
policy: policySchema.default({}),
|
||||
acp: acpSchema.default({}),
|
||||
backends: z.array(backendSchema).min(1),
|
||||
skills: z.array(skillSchema).default([]),
|
||||
defaultRole: z.string().min(1),
|
||||
roles: z.array(roleSchema).min(1),
|
||||
platforms: z.object({
|
||||
feishu: feishuSchema.default({}), wecom: wecomSchema.default({}), qq: qqSchema.default({}),
|
||||
webhook: webhookSchema.default({}), weixin: weixinSchema.default({})
|
||||
}).default({})
|
||||
});
|
||||
configVersion: z.literal(3),
|
||||
bot: botSchema,
|
||||
gateway: z.object({
|
||||
server: serverSchema.default({}),
|
||||
policy: gatewayPolicySchema.default({}),
|
||||
platform: platformSchema
|
||||
}).strict(),
|
||||
runtime: z.object({ acp: acpSchema.default({}) }).strict().default({})
|
||||
}).strict();
|
||||
|
||||
export type AppConfig = z.infer<typeof configSchema>;
|
||||
export type GatewayPolicy = z.infer<typeof policySchema>;
|
||||
export type AcpConfig = AppConfig["acp"];
|
||||
export type AcpBackendConfig = AppConfig["backends"][number];
|
||||
export type RoleConfig = AppConfig["roles"][number];
|
||||
export type RolePolicy = RoleConfig["policy"];
|
||||
export type SkillConfig = AppConfig["skills"][number];
|
||||
export type BotConfig = AppConfig["bot"];
|
||||
export type AgentConfig = BotConfig["agent"];
|
||||
export type PermissionPolicy = BotConfig["permissions"];
|
||||
export type SkillConfig = BotConfig["skills"][number];
|
||||
export type GatewayPolicy = AppConfig["gateway"]["policy"];
|
||||
export type AcpConfig = AppConfig["runtime"]["acp"];
|
||||
export type PlatformConfig = AppConfig["gateway"]["platform"];
|
||||
export type PlatformType = PlatformConfig["type"];
|
||||
export type QqConfig = Extract<PlatformConfig, { type: "qq" }>;
|
||||
export type FeishuConfig = Extract<PlatformConfig, { type: "feishu" }>;
|
||||
export type WeComConfig = Extract<PlatformConfig, { type: "wecom" }>;
|
||||
export type WebhookConfig = Extract<PlatformConfig, { type: "webhook" }>;
|
||||
export type WeixinConfig = Extract<PlatformConfig, { type: "weixin" }>;
|
||||
|
||||
// Kept only so legacy CliAgent source remains type-checkable; it is not used by the runtime.
|
||||
export interface CliAgentConfig {
|
||||
@@ -102,64 +127,29 @@ export interface CliAgentConfig {
|
||||
cwd?: string;
|
||||
}
|
||||
|
||||
interface V1AgentConfig {
|
||||
name?: string;
|
||||
command?: string;
|
||||
args?: string[];
|
||||
cwd?: string;
|
||||
}
|
||||
|
||||
export function resolveConfigPath(configPath = process.env.GORI_GATEWAY_CONFIG): string {
|
||||
return path.resolve(configPath || "config.example.json");
|
||||
return path.resolve(configPath || "config.json");
|
||||
}
|
||||
|
||||
export function parseConfig(rawConfig: unknown): AppConfig {
|
||||
const candidate = isRecord(rawConfig) && rawConfig.configVersion === 2 ? rawConfig : migrateV1Config(rawConfig);
|
||||
const config = configSchema.parse(candidate);
|
||||
validateUnique(config.backends.map((item) => item.id), "backend");
|
||||
validateUnique(config.roles.map((item) => item.id), "role");
|
||||
validateUnique(config.skills.map((item) => item.id), "skill");
|
||||
|
||||
const backendIds = new Set(config.backends.map((item) => item.id));
|
||||
const skillIds = new Set(config.skills.map((item) => item.id));
|
||||
if (!config.roles.some((role) => role.id === config.defaultRole)) throw new Error(`defaultRole '${config.defaultRole}' is not present in roles`);
|
||||
for (const role of config.roles) {
|
||||
if (!path.isAbsolute(role.workspace)) throw new Error(`role '${role.id}' workspace must be absolute`);
|
||||
if (!backendIds.has(role.backend)) throw new Error(`role '${role.id}' references unknown backend '${role.backend}'`);
|
||||
for (const skill of role.skills) if (!skillIds.has(skill)) throw new Error(`role '${role.id}' references unknown skill '${skill}'`);
|
||||
for (const pattern of role.policy.allowedCommandPatterns) {
|
||||
try { new RegExp(pattern); } catch { throw new Error(`role '${role.id}' has invalid command pattern '${pattern}'`); }
|
||||
}
|
||||
if (!isRecord(rawConfig)) throw new Error("configuration must be an object");
|
||||
if (rawConfig.configVersion !== 3) {
|
||||
throw new Error(`unsupported configVersion '${String(rawConfig.configVersion)}'; gori-agent requires Config v3`);
|
||||
}
|
||||
const config = configSchema.parse(rawConfig);
|
||||
if (!path.isAbsolute(config.bot.workspace)) throw new Error("bot.workspace must be absolute");
|
||||
validateUnique(config.bot.skills.map((item) => item.id), "skill");
|
||||
for (const skill of config.bot.skills) {
|
||||
if (!path.isAbsolute(skill.file)) throw new Error(`skill '${skill.id}' file must be absolute`);
|
||||
}
|
||||
for (const pattern of config.bot.permissions.allowedCommandPatterns) {
|
||||
try { new RegExp(pattern); } catch { throw new Error(`bot has invalid command pattern '${pattern}'`); }
|
||||
}
|
||||
return config;
|
||||
}
|
||||
|
||||
export function migrateV1Config(rawConfig: unknown): unknown {
|
||||
if (!isRecord(rawConfig)) throw new Error("configuration must be an object");
|
||||
const agents = Array.isArray(rawConfig.agents) ? rawConfig.agents.filter(isRecord) as V1AgentConfig[] : [];
|
||||
const defaultAgent = typeof rawConfig.defaultAgent === "string" ? rawConfig.defaultAgent : "";
|
||||
const selected = agents.find((agent) => agent.name === defaultAgent);
|
||||
if (!selected || selected.name !== "kimi" || !selected.command || path.basename(selected.command) !== "kimi") {
|
||||
throw new Error("v1 migration only supports a Kimi default agent; configure an ACP backend explicitly for other agents");
|
||||
}
|
||||
return {
|
||||
configVersion: 2,
|
||||
server: rawConfig.server,
|
||||
policy: rawConfig.policy,
|
||||
acp: {},
|
||||
backends: [{ id: "kimi", command: selected.command, args: ["acp"] }],
|
||||
skills: [],
|
||||
defaultRole: "assistant",
|
||||
roles: [{
|
||||
id: "assistant", backend: "kimi", workspace: path.resolve(selected.cwd || process.cwd()), persona: "", skills: [],
|
||||
policy: { permissionMode: "deny", allowedTools: [], allowedCommandPatterns: [] }
|
||||
}],
|
||||
platforms: rawConfig.platforms
|
||||
};
|
||||
}
|
||||
|
||||
export function defaultStateFile(config: AppConfig): string {
|
||||
if (config.acp.stateFile) return path.resolve(config.acp.stateFile);
|
||||
if (config.runtime.acp.stateFile) return path.resolve(config.runtime.acp.stateFile);
|
||||
const home = process.env.GORI_AGENT_HOME || process.cwd();
|
||||
return path.join(path.resolve(home), "state", "acp-sessions.json");
|
||||
}
|
||||
@@ -172,6 +162,13 @@ export function loadConfigFromPath(configPath: string): AppConfig {
|
||||
return parseConfig(JSON.parse(fs.readFileSync(configPath, "utf8")) as unknown);
|
||||
}
|
||||
|
||||
export function isValidServerHost(value: string): boolean {
|
||||
if (!value || value !== value.trim() || /[\s\u0000-\u001f\u007f/\\\[\]]/.test(value) || value.includes("://")) return false;
|
||||
if (net.isIP(value) !== 0) return true;
|
||||
if (value.length > 253 || value.includes(":")) return false;
|
||||
return value.split(".").every((label) => /^(?=.{1,63}$)[A-Za-z0-9](?:[A-Za-z0-9-]*[A-Za-z0-9])?$/.test(label));
|
||||
}
|
||||
|
||||
function validateUnique(ids: string[], label: string): void {
|
||||
if (new Set(ids).size !== ids.length) throw new Error(`${label} IDs must be unique`);
|
||||
}
|
||||
|
||||
@@ -1,25 +1,23 @@
|
||||
export type CommandKind = "help" | "roles" | "role" | "status" | "cancel" | "new";
|
||||
export type CommandKind = "help" | "status" | "cancel" | "new" | "retired-role";
|
||||
|
||||
export interface ParsedCommand {
|
||||
kind: CommandKind;
|
||||
argument?: string;
|
||||
deprecatedAlias?: boolean;
|
||||
}
|
||||
|
||||
export class CommandRouter {
|
||||
parse(text: string): ParsedCommand | undefined {
|
||||
const trimmed = text.trim();
|
||||
if (!trimmed.startsWith("/")) return undefined;
|
||||
const [command, ...args] = trimmed.split(/\s+/);
|
||||
const [command] = trimmed.split(/\s+/);
|
||||
switch (command.toLowerCase()) {
|
||||
case "/help": return { kind: "help" };
|
||||
case "/roles": return { kind: "roles" };
|
||||
case "/role": return { kind: "role", argument: args[0] };
|
||||
case "/status": return { kind: "status" };
|
||||
case "/cancel": return { kind: "cancel" };
|
||||
case "/new": return { kind: "new" };
|
||||
case "/agents": return { kind: "roles", deprecatedAlias: true };
|
||||
case "/agent": return { kind: "role", argument: args[0], deprecatedAlias: true };
|
||||
case "/roles":
|
||||
case "/role":
|
||||
case "/agents":
|
||||
case "/agent": return { kind: "retired-role" };
|
||||
default: return undefined;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,53 +1,72 @@
|
||||
import fs from "node:fs";
|
||||
import path from "node:path";
|
||||
|
||||
export interface ChatState {
|
||||
selectedRole: string;
|
||||
createdAt: number;
|
||||
updatedAt: number;
|
||||
export interface StoreIdentity {
|
||||
botId: string;
|
||||
platform: string;
|
||||
}
|
||||
|
||||
export interface SessionBinding {
|
||||
chatKey: string;
|
||||
roleId: string;
|
||||
backendId: string;
|
||||
agentId: string;
|
||||
nativeSessionId: string;
|
||||
workspace: string;
|
||||
roleFingerprint: string;
|
||||
botFingerprint: string;
|
||||
createdAt: number;
|
||||
updatedAt: number;
|
||||
}
|
||||
|
||||
interface StoreData {
|
||||
version: 1;
|
||||
chats: Record<string, ChatState>;
|
||||
version: 2;
|
||||
botId: string;
|
||||
platform: string;
|
||||
bindings: Record<string, SessionBinding>;
|
||||
}
|
||||
|
||||
const EMPTY: StoreData = { version: 1, chats: {}, bindings: {} };
|
||||
|
||||
export class DurableSessionStore {
|
||||
private data: StoreData = structuredClone(EMPTY);
|
||||
private data: StoreData;
|
||||
private queue: Promise<void> = Promise.resolve();
|
||||
private lockFd?: fs.promises.FileHandle;
|
||||
private closed = false;
|
||||
|
||||
constructor(readonly file: string) {}
|
||||
constructor(readonly file: string, readonly identity: StoreIdentity) {
|
||||
this.data = emptyData(identity);
|
||||
}
|
||||
|
||||
async open(): Promise<void> {
|
||||
await fs.promises.mkdir(path.dirname(this.file), { recursive: true });
|
||||
const directory = path.dirname(this.file);
|
||||
await fs.promises.mkdir(directory, { recursive: true, mode: 0o700 });
|
||||
const directoryStat = await fs.promises.lstat(directory);
|
||||
if (!directoryStat.isDirectory() || directoryStat.isSymbolicLink()) throw new Error(`Refusing unsafe state directory: ${directory}`);
|
||||
const directoryMode = directoryStat.mode & 0o777;
|
||||
if (directoryMode !== 0o700) throw new Error(`State directory mode must be 0700, got 0${directoryMode.toString(8)}`);
|
||||
if (typeof process.getuid === "function" && directoryStat.uid !== process.getuid()) throw new Error("State directory is not owned by the current user");
|
||||
const lockFile = `${this.file}.lock`;
|
||||
try {
|
||||
this.lockFd = await acquireLock(lockFile);
|
||||
await this.lockFd.writeFile(`${process.pid}\n`);
|
||||
} catch (error) {
|
||||
if ((error as NodeJS.ErrnoException).code === "EEXIST") throw new Error(`Session state is locked by another gori-agent instance: ${lockFile}`);
|
||||
if ((error as NodeJS.ErrnoException).code !== "EEXIST" || !await removeStaleLock(lockFile)) {
|
||||
if ((error as NodeJS.ErrnoException).code === "EEXIST") throw new Error(`Session state is locked by another gori-agent instance: ${lockFile}`);
|
||||
throw error;
|
||||
}
|
||||
this.lockFd = await acquireLock(lockFile).catch((retryError) => {
|
||||
if ((retryError as NodeJS.ErrnoException).code === "EEXIST") throw new Error(`Session state is locked by another gori-agent instance: ${lockFile}`);
|
||||
throw retryError;
|
||||
});
|
||||
}
|
||||
try {
|
||||
await this.lockFd.writeFile(`${process.pid}\n`);
|
||||
await this.lockFd.sync();
|
||||
} catch (error) {
|
||||
await this.releaseLock();
|
||||
throw error;
|
||||
}
|
||||
try {
|
||||
const raw = await fs.promises.readFile(this.file, "utf8");
|
||||
const parsed = JSON.parse(raw) as StoreData;
|
||||
if (parsed.version !== 1 || !parsed.chats || !parsed.bindings) throw new Error("unsupported or malformed state data");
|
||||
const parsed = parseStoreData(JSON.parse(raw) as unknown);
|
||||
if (parsed.botId !== this.identity.botId || parsed.platform !== this.identity.platform) {
|
||||
throw new Error(`state identity mismatch: expected bot '${this.identity.botId}' platform '${this.identity.platform}'`);
|
||||
}
|
||||
this.data = parsed;
|
||||
} catch (error) {
|
||||
if ((error as NodeJS.ErrnoException).code !== "ENOENT") {
|
||||
@@ -57,44 +76,32 @@ export class DurableSessionStore {
|
||||
}
|
||||
}
|
||||
|
||||
getSelectedRole(chatKey: string, defaultRole: string): string {
|
||||
return this.data.chats[chatKey]?.selectedRole || defaultRole;
|
||||
}
|
||||
|
||||
async setSelectedRole(chatKey: string, roleId: string): Promise<void> {
|
||||
const now = Date.now();
|
||||
const current = this.data.chats[chatKey];
|
||||
this.data.chats[chatKey] = { selectedRole: roleId, createdAt: current?.createdAt || now, updatedAt: now };
|
||||
await this.persist();
|
||||
}
|
||||
|
||||
getBinding(chatKey: string, roleId: string): SessionBinding | undefined {
|
||||
const value = this.data.bindings[bindingKey(chatKey, roleId)];
|
||||
getBinding(chatKey: string): SessionBinding | undefined {
|
||||
const value = this.data.bindings[chatKey];
|
||||
return value ? structuredClone(value) : undefined;
|
||||
}
|
||||
|
||||
async setBinding(binding: SessionBinding): Promise<void> {
|
||||
this.data.bindings[bindingKey(binding.chatKey, binding.roleId)] = structuredClone(binding);
|
||||
this.data.bindings[binding.chatKey] = structuredClone(binding);
|
||||
await this.persist();
|
||||
}
|
||||
|
||||
async touchBinding(chatKey: string, roleId: string): Promise<void> {
|
||||
const binding = this.data.bindings[bindingKey(chatKey, roleId)];
|
||||
async touchBinding(chatKey: string): Promise<void> {
|
||||
const binding = this.data.bindings[chatKey];
|
||||
if (!binding) return;
|
||||
binding.updatedAt = Date.now();
|
||||
await this.persist();
|
||||
}
|
||||
|
||||
async deleteBinding(chatKey: string, roleId: string): Promise<SessionBinding | undefined> {
|
||||
const key = bindingKey(chatKey, roleId);
|
||||
const existing = this.data.bindings[key];
|
||||
delete this.data.bindings[key];
|
||||
async deleteBinding(chatKey: string): Promise<SessionBinding | undefined> {
|
||||
const existing = this.data.bindings[chatKey];
|
||||
delete this.data.bindings[chatKey];
|
||||
if (existing) await this.persist();
|
||||
return existing;
|
||||
return existing ? structuredClone(existing) : undefined;
|
||||
}
|
||||
|
||||
stats(): { chats: number; bindings: number } {
|
||||
return { chats: Object.keys(this.data.chats).length, bindings: Object.keys(this.data.bindings).length };
|
||||
stats(): { bindings: number } {
|
||||
return { bindings: Object.keys(this.data.bindings).length };
|
||||
}
|
||||
|
||||
async flush(): Promise<void> { await this.queue; }
|
||||
@@ -102,15 +109,15 @@ export class DurableSessionStore {
|
||||
async close(): Promise<void> {
|
||||
if (this.closed) return;
|
||||
this.closed = true;
|
||||
await this.flush();
|
||||
await this.releaseLock();
|
||||
try { await this.flush(); } finally { await this.releaseLock(); }
|
||||
}
|
||||
|
||||
private persist(): Promise<void> {
|
||||
if (this.closed) return Promise.reject(new Error("Session store is closed"));
|
||||
const snapshot = JSON.stringify(this.data, null, 2) + "\n";
|
||||
this.queue = this.queue.then(() => atomicWrite(this.file, snapshot));
|
||||
return this.queue;
|
||||
const operation = this.queue.then(() => atomicWrite(this.file, snapshot));
|
||||
this.queue = operation.catch(() => undefined);
|
||||
return operation;
|
||||
}
|
||||
|
||||
private async releaseLock(): Promise<void> {
|
||||
@@ -124,26 +131,73 @@ export class DurableSessionStore {
|
||||
}
|
||||
|
||||
export function chatKeyFor(platform: string, chatId: string): string { return `${platform}:${chatId}`; }
|
||||
export function bindingKey(chatKey: string, roleId: string): string { return `${chatKey}\u0000${roleId}`; }
|
||||
|
||||
function emptyData(identity: StoreIdentity): StoreData {
|
||||
return { version: 2, botId: identity.botId, platform: identity.platform, bindings: {} };
|
||||
}
|
||||
|
||||
function parseStoreData(value: unknown): StoreData {
|
||||
if (!isRecord(value) || value.version !== 2 || typeof value.botId !== "string" || typeof value.platform !== "string" || !isRecord(value.bindings)) {
|
||||
throw new Error("unsupported state version; Config v3 requires state v2 in a new GORI_AGENT_HOME");
|
||||
}
|
||||
for (const [chatKey, binding] of Object.entries(value.bindings)) {
|
||||
if (!isRecord(binding)
|
||||
|| binding.chatKey !== chatKey
|
||||
|| typeof binding.agentId !== "string"
|
||||
|| typeof binding.nativeSessionId !== "string"
|
||||
|| typeof binding.workspace !== "string"
|
||||
|| typeof binding.botFingerprint !== "string"
|
||||
|| typeof binding.createdAt !== "number"
|
||||
|| typeof binding.updatedAt !== "number") {
|
||||
throw new Error(`invalid state v2 binding '${chatKey}'`);
|
||||
}
|
||||
}
|
||||
return value as unknown as StoreData;
|
||||
}
|
||||
|
||||
function isRecord(value: unknown): value is Record<string, unknown> {
|
||||
return typeof value === "object" && value !== null && !Array.isArray(value);
|
||||
}
|
||||
|
||||
async function acquireLock(lockFile: string): Promise<fs.promises.FileHandle> {
|
||||
return fs.promises.open(lockFile, "wx", 0o600);
|
||||
}
|
||||
|
||||
async function removeStaleLock(lockFile: string): Promise<boolean> {
|
||||
let handle: fs.promises.FileHandle | undefined;
|
||||
try {
|
||||
handle = await fs.promises.open(lockFile, "r");
|
||||
const pid = Number((await handle.readFile("utf8")).trim());
|
||||
if (!Number.isSafeInteger(pid) || pid <= 1 || processExists(pid)) return false;
|
||||
const opened = await handle.stat();
|
||||
const current = await fs.promises.lstat(lockFile);
|
||||
if (opened.dev !== current.dev || opened.ino !== current.ino) return false;
|
||||
await fs.promises.unlink(lockFile);
|
||||
return true;
|
||||
} catch { return false; }
|
||||
finally { await handle?.close().catch(() => undefined); }
|
||||
}
|
||||
|
||||
function processExists(pid: number): boolean {
|
||||
try { process.kill(pid, 0); return true; }
|
||||
catch (error) { return (error as NodeJS.ErrnoException).code === "EPERM"; }
|
||||
}
|
||||
|
||||
async function atomicWrite(file: string, content: string): Promise<void> {
|
||||
const temp = `${file}.${process.pid}.${Date.now()}.tmp`;
|
||||
const handle = await fs.promises.open(temp, "wx", 0o600);
|
||||
let handle: fs.promises.FileHandle | undefined;
|
||||
try {
|
||||
handle = await fs.promises.open(temp, "wx", 0o600);
|
||||
await handle.writeFile(content, "utf8");
|
||||
await handle.sync();
|
||||
} finally {
|
||||
await handle.close();
|
||||
}
|
||||
try {
|
||||
handle = undefined;
|
||||
await fs.promises.rename(temp, file);
|
||||
await fs.promises.chmod(file, 0o600);
|
||||
const dir = await fs.promises.open(path.dirname(file), "r");
|
||||
try { await dir.sync(); } finally { await dir.close(); }
|
||||
} catch (error) {
|
||||
if (handle) await handle.close().catch(() => undefined);
|
||||
await fs.promises.unlink(temp).catch(() => undefined);
|
||||
throw error;
|
||||
}
|
||||
|
||||
+7
-15
@@ -1,6 +1,5 @@
|
||||
import type { ConversationRuntime } from "../acp/types.js";
|
||||
import type { GatewayPolicy } from "../config.js";
|
||||
import type { RoleRegistry } from "../roles/role-registry.js";
|
||||
import type { PlatformAdapter } from "./adapter.js";
|
||||
import { CommandRouter, type ParsedCommand } from "./command-router.js";
|
||||
import { chatKeyFor } from "./durable-session-store.js";
|
||||
@@ -12,12 +11,12 @@ export class Gateway {
|
||||
private readonly locks = new Map<string, Promise<void>>();
|
||||
readonly commandRouter = new CommandRouter();
|
||||
|
||||
constructor(private readonly policy: GatewayPolicy, private readonly runtime: ConversationRuntime, private readonly roles: RoleRegistry) {}
|
||||
constructor(private readonly policy: GatewayPolicy, private readonly runtime: ConversationRuntime) {}
|
||||
|
||||
async receive(message: IncomingMessage, adapter: PlatformAdapter, options: { synchronous?: boolean } = {}): Promise<GatewayResult> {
|
||||
const policyError = this.checkPolicy(message);
|
||||
if (policyError) {
|
||||
console.log(`Message ignored by policy: ${policyError} (${message.platform} ${message.chatId} ${message.userId})`);
|
||||
console.log(`Message ignored by policy: ${policyError} (platform=${message.platform})`);
|
||||
return { ok: true, ignored: true, error: policyError };
|
||||
}
|
||||
|
||||
@@ -43,24 +42,17 @@ export class Gateway {
|
||||
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 "help": return ["Commands:", "/status", "/cancel", "/new", "/help"].join("\n");
|
||||
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.";
|
||||
return "Started a new native ACP session for this chat.";
|
||||
case "cancel": return this.cancelText(message);
|
||||
case "retired-role": return "Role switching was removed in Config v3; this instance has one fixed Bot and ACP agent.";
|
||||
}
|
||||
}
|
||||
|
||||
@@ -78,8 +70,8 @@ export class Gateway {
|
||||
}
|
||||
|
||||
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.allowedUsers.length > 0 && !this.policy.allowedUsers.includes(message.userId)) return "User not allowed";
|
||||
if (this.policy.allowedChats.length > 0 && !this.policy.allowedChats.includes(message.chatId)) return "Chat not allowed";
|
||||
if (this.policy.requireMentionInGroup && message.isGroup && !message.mentionsBot) return "Mention required in group chat";
|
||||
return undefined;
|
||||
}
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
import type { AppConfig } from "../../config.js";
|
||||
import type { FeishuConfig } from "../../config.js";
|
||||
import { jsonResponse } from "../../core/adapter.js";
|
||||
import type { PlatformAdapter } from "../../core/adapter.js";
|
||||
import type { Gateway } from "../../core/gateway.js";
|
||||
@@ -11,7 +11,7 @@ export class FeishuAdapter implements PlatformAdapter {
|
||||
private tenantToken?: { token: string; expiresAt: number };
|
||||
|
||||
constructor(
|
||||
private readonly config: AppConfig["platforms"]["feishu"],
|
||||
private readonly config: FeishuConfig,
|
||||
private readonly gateway: Gateway
|
||||
) {}
|
||||
|
||||
@@ -72,9 +72,7 @@ export class FeishuAdapter implements PlatformAdapter {
|
||||
})
|
||||
});
|
||||
|
||||
if (!response.ok) {
|
||||
throw new Error(`Feishu reply failed: ${response.status} ${await response.text()}`);
|
||||
}
|
||||
if (!response.ok) throw new Error(`Feishu reply failed: HTTP ${response.status}`);
|
||||
const data = await response.json() as { code?: number; msg?: string };
|
||||
if (data.code && data.code !== 0) {
|
||||
throw new Error(`Feishu reply failed: ${data.code} ${data.msg || ""}`.trim());
|
||||
@@ -90,9 +88,7 @@ export class FeishuAdapter implements PlatformAdapter {
|
||||
headers: { "Content-Type": "application/json; charset=utf-8" },
|
||||
body: JSON.stringify({ app_id: this.config.appId, app_secret: this.config.appSecret })
|
||||
});
|
||||
if (!response.ok) {
|
||||
throw new Error(`Feishu token request failed: ${response.status} ${await response.text()}`);
|
||||
}
|
||||
if (!response.ok) throw new Error(`Feishu token request failed: HTTP ${response.status}`);
|
||||
const data = await response.json() as FeishuTenantTokenResponse;
|
||||
if (data.code !== 0 || !data.tenant_access_token) {
|
||||
throw new Error(`Feishu token request failed: ${data.code} ${data.msg || ""}`.trim());
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
import type { AppConfig } from "../../config.js";
|
||||
import type { QqConfig } from "../../config.js";
|
||||
import { jsonResponse } from "../../core/adapter.js";
|
||||
import type { PlatformAdapter } from "../../core/adapter.js";
|
||||
import type { Gateway } from "../../core/gateway.js";
|
||||
@@ -18,7 +18,7 @@ export class QqAdapter implements PlatformAdapter {
|
||||
private accessToken?: { token: string; expiresAt: number };
|
||||
|
||||
constructor(
|
||||
private readonly config: AppConfig["platforms"]["qq"],
|
||||
private readonly config: QqConfig,
|
||||
private readonly gateway: Gateway
|
||||
) {}
|
||||
|
||||
@@ -48,7 +48,7 @@ export class QqAdapter implements PlatformAdapter {
|
||||
console.log(`QQ recv ${payload.t} ignored (empty text or missing ids)`);
|
||||
return false;
|
||||
}
|
||||
console.log(`QQ recv ${payload.t} chat=${message.chatId} user=${message.userId} msgId=${message.messageId || ""} text=${JSON.stringify(message.text)}`);
|
||||
console.log(`QQ recv ${payload.t} accepted length=${message.text.length}`);
|
||||
void this.gateway.receive(message, this).catch((error) => {
|
||||
console.error("QQ gateway error", error);
|
||||
});
|
||||
@@ -64,7 +64,7 @@ export class QqAdapter implements PlatformAdapter {
|
||||
const isGroup = Boolean(groupOpenId) || message.target.chatId.startsWith("group:");
|
||||
const targetId = groupOpenId || userOpenId || message.target.chatId.replace(/^group:/, "").replace(/^user:/, "");
|
||||
const path = isGroup ? `/v2/groups/${encodeURIComponent(targetId)}/messages` : `/v2/users/${encodeURIComponent(targetId)}/messages`;
|
||||
console.log(`QQ send -> ${path} msgId=${message.replyTo || ""} length=${message.text.length}`);
|
||||
console.log(`QQ send ${isGroup ? "group" : "user"} message length=${message.text.length}`);
|
||||
|
||||
const response = await fetch(`https://api.sgroup.qq.com${path}`, {
|
||||
method: "POST",
|
||||
@@ -74,7 +74,7 @@ export class QqAdapter implements PlatformAdapter {
|
||||
},
|
||||
body: JSON.stringify({ content: message.text, msg_id: message.replyTo })
|
||||
});
|
||||
if (!response.ok) throw new Error(`QQ send failed: ${response.status} ${await response.text()}`);
|
||||
if (!response.ok) throw new Error(`QQ send failed: HTTP ${response.status}`);
|
||||
const data = await response.json() as QqSendMessageResponse;
|
||||
if (data.code && data.code !== 0) throw new Error(`QQ send failed: ${data.code} ${data.message || ""}`.trim());
|
||||
}
|
||||
@@ -128,7 +128,7 @@ export class QqAdapter implements PlatformAdapter {
|
||||
headers: { "Content-Type": "application/json; charset=utf-8" },
|
||||
body: JSON.stringify({ appId: this.config.appId, clientSecret: this.config.clientSecret })
|
||||
});
|
||||
if (!response.ok) throw new Error(`QQ token request failed: ${response.status} ${await response.text()}`);
|
||||
if (!response.ok) throw new Error(`QQ token request failed: HTTP ${response.status}`);
|
||||
const data = await response.json() as QqAccessTokenResponse;
|
||||
if (!data.access_token) {
|
||||
throw new Error(`QQ token request failed: ${data.error || "unknown"} ${data.error_description || ""}`.trim());
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
import os from "node:os";
|
||||
import WebSocket from "ws";
|
||||
import type { AppConfig } from "../../config.js";
|
||||
import type { QqConfig } from "../../config.js";
|
||||
import type { QqAdapter } from "./adapter.js";
|
||||
import type { QqGatewayResponse, QqWebhookPayload } from "./types.js";
|
||||
|
||||
@@ -19,7 +19,7 @@ export class QqGatewayClient {
|
||||
private stopped = false;
|
||||
|
||||
constructor(
|
||||
private readonly config: AppConfig["platforms"]["qq"],
|
||||
private readonly config: QqConfig,
|
||||
private readonly adapter: QqAdapter
|
||||
) {}
|
||||
|
||||
@@ -113,7 +113,7 @@ export class QqGatewayClient {
|
||||
const response = await fetch(QQ_GATEWAY_API, {
|
||||
headers: { Authorization: `QQBot ${token}` }
|
||||
});
|
||||
if (!response.ok) throw new Error(`QQ gateway request failed: ${response.status} ${await response.text()}`);
|
||||
if (!response.ok) throw new Error(`QQ gateway request failed: HTTP ${response.status}`);
|
||||
return response.json() as Promise<QqGatewayResponse>;
|
||||
}
|
||||
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
import crypto from "node:crypto";
|
||||
import type { AppConfig } from "../../config.js";
|
||||
import type { WebhookConfig } from "../../config.js";
|
||||
import { jsonResponse } from "../../core/adapter.js";
|
||||
import type { PlatformAdapter } from "../../core/adapter.js";
|
||||
import type { Gateway } from "../../core/gateway.js";
|
||||
@@ -32,7 +32,7 @@ export class GenericWebhookAdapter implements PlatformAdapter {
|
||||
readonly name = "webhook";
|
||||
|
||||
constructor(
|
||||
private readonly config: AppConfig["platforms"]["webhook"],
|
||||
private readonly config: WebhookConfig,
|
||||
private readonly gateway: Gateway
|
||||
) {}
|
||||
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
import type { AppConfig } from "../../config.js";
|
||||
import type { WeComConfig } from "../../config.js";
|
||||
import { notImplemented } from "../../core/adapter.js";
|
||||
import type { PlatformAdapter } from "../../core/adapter.js";
|
||||
import type { OutgoingMessage, WebhookRequestContext, WebhookResponse } from "../../core/types.js";
|
||||
@@ -8,7 +8,7 @@ export class WeComAdapter implements PlatformAdapter {
|
||||
readonly name = "wecom";
|
||||
private accessToken?: { token: string; expiresAt: number };
|
||||
|
||||
constructor(private readonly config: AppConfig["platforms"]["wecom"]) {}
|
||||
constructor(private readonly config: WeComConfig) {}
|
||||
|
||||
async handleWebhook(_context: WebhookRequestContext): Promise<WebhookResponse> {
|
||||
return notImplemented("WeCom");
|
||||
@@ -27,7 +27,7 @@ export class WeComAdapter implements PlatformAdapter {
|
||||
safe: 0
|
||||
})
|
||||
});
|
||||
if (!response.ok) throw new Error(`WeCom send failed: ${response.status} ${await response.text()}`);
|
||||
if (!response.ok) throw new Error(`WeCom send failed: HTTP ${response.status}`);
|
||||
const data = await response.json() as WeComSendMessageResponse;
|
||||
if (data.errcode !== 0) throw new Error(`WeCom send failed: ${data.errcode} ${data.errmsg || ""}`.trim());
|
||||
}
|
||||
@@ -39,7 +39,7 @@ export class WeComAdapter implements PlatformAdapter {
|
||||
url.searchParams.set("corpid", this.config.corpId);
|
||||
url.searchParams.set("corpsecret", this.config.secret);
|
||||
const response = await fetch(url);
|
||||
if (!response.ok) throw new Error(`WeCom token request failed: ${response.status} ${await response.text()}`);
|
||||
if (!response.ok) throw new Error(`WeCom token request failed: HTTP ${response.status}`);
|
||||
const data = await response.json() as WeComAccessTokenResponse;
|
||||
if (data.errcode !== 0 || !data.access_token) {
|
||||
throw new Error(`WeCom token request failed: ${data.errcode} ${data.errmsg || ""}`.trim());
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
import type { AppConfig } from "../../config.js";
|
||||
import type { WeixinConfig } from "../../config.js";
|
||||
import { jsonResponse, notImplemented } from "../../core/adapter.js";
|
||||
import type { PlatformAdapter } from "../../core/adapter.js";
|
||||
import type { Gateway } from "../../core/gateway.js";
|
||||
@@ -17,7 +17,7 @@ export class WeixinAdapter implements PlatformAdapter {
|
||||
readonly name = "weixin";
|
||||
|
||||
constructor(
|
||||
private readonly config: AppConfig["platforms"]["weixin"],
|
||||
private readonly config: WeixinConfig,
|
||||
private readonly gateway: Gateway
|
||||
) {}
|
||||
|
||||
|
||||
+29
-35
@@ -1,51 +1,45 @@
|
||||
import crypto from "node:crypto";
|
||||
import type { AppConfig, RoleConfig } from "../config.js";
|
||||
import type { AppConfig, BotConfig } from "../config.js";
|
||||
import { SkillLoader, type LoadedSkill } from "./skill-loader.js";
|
||||
|
||||
export interface ResolvedRole extends RoleConfig {
|
||||
const BOOTSTRAP_SCHEMA_VERSION = 1;
|
||||
|
||||
export interface ResolvedBot extends BotConfig {
|
||||
loadedSkills: LoadedSkill[];
|
||||
fingerprint: string;
|
||||
bootstrap: string;
|
||||
}
|
||||
|
||||
export class RoleRegistry {
|
||||
private readonly roles = new Map<string, ResolvedRole>();
|
||||
export class BotProfileResolver {
|
||||
readonly bot: ResolvedBot;
|
||||
|
||||
constructor(private readonly config: AppConfig) {
|
||||
const loader = new SkillLoader(config.skills);
|
||||
for (const role of config.roles) {
|
||||
const loadedSkills = role.skills.map((id) => loader.load(id));
|
||||
const fingerprint = crypto.createHash("sha256").update(JSON.stringify({
|
||||
id: role.id,
|
||||
backend: role.backend,
|
||||
workspace: role.workspace,
|
||||
persona: role.persona,
|
||||
policy: role.policy,
|
||||
skills: loadedSkills.map(({ id, file, hash }) => ({ id, file, hash }))
|
||||
})).digest("hex");
|
||||
this.roles.set(role.id, { ...role, loadedSkills, fingerprint, bootstrap: buildBootstrap(role, loadedSkills) });
|
||||
}
|
||||
constructor(config: AppConfig) {
|
||||
const loadedSkills = config.bot.skills.map((skill) => new SkillLoader(config.bot.skills).load(skill.id));
|
||||
const fingerprint = crypto.createHash("sha256").update(JSON.stringify({
|
||||
bootstrapSchemaVersion: BOOTSTRAP_SCHEMA_VERSION,
|
||||
id: config.bot.id,
|
||||
workspace: config.bot.workspace,
|
||||
persona: config.bot.persona,
|
||||
agent: config.bot.agent,
|
||||
permissions: config.bot.permissions,
|
||||
skills: loadedSkills.map(({ id, file, hash }) => ({ id, file, hash }))
|
||||
})).digest("hex");
|
||||
this.bot = {
|
||||
...config.bot,
|
||||
loadedSkills,
|
||||
fingerprint,
|
||||
bootstrap: buildBootstrap(config.bot, loadedSkills)
|
||||
};
|
||||
}
|
||||
|
||||
get(id?: string): ResolvedRole {
|
||||
const roleId = id || this.config.defaultRole;
|
||||
const role = this.roles.get(roleId);
|
||||
if (!role) throw new Error(`Unknown role: ${roleId}`);
|
||||
return role;
|
||||
}
|
||||
|
||||
has(id: string): boolean { return this.roles.has(id); }
|
||||
list(): string[] { return [...this.roles.keys()].sort(); }
|
||||
defaultId(): string { return this.config.defaultRole; }
|
||||
}
|
||||
|
||||
function buildBootstrap(role: RoleConfig, skills: LoadedSkill[]): string {
|
||||
function buildBootstrap(bot: BotConfig, skills: LoadedSkill[]): string {
|
||||
return [
|
||||
"Initialize this ACP session with the following role. Treat these instructions as persistent context. Reply only with READY.",
|
||||
`Role: ${role.id}`,
|
||||
`Workspace: ${role.workspace}`,
|
||||
role.persona ? `Persona:\n${role.persona}` : "Persona: general coding assistant",
|
||||
`Permission policy enforced by the ACP client: ${JSON.stringify(role.policy)}`,
|
||||
`Initialize this ACP session with gori-agent bootstrap schema v${BOOTSTRAP_SCHEMA_VERSION}. Treat these instructions as persistent context. Reply only with READY.`,
|
||||
`Bot: ${bot.id}`,
|
||||
`Workspace: ${bot.workspace}`,
|
||||
bot.persona ? `Persona:\n${bot.persona}` : "Persona: general coding assistant",
|
||||
`Permission policy enforced by the ACP client: ${JSON.stringify(bot.permissions)}`,
|
||||
...skills.map((skill) => `Skill ${skill.id} (${skill.file}):\n${skill.content}`)
|
||||
].join("\n\n");
|
||||
}
|
||||
|
||||
+77
-49
@@ -1,103 +1,131 @@
|
||||
import express from "express";
|
||||
import type { Server } from "node:http";
|
||||
import { fileURLToPath } from "node:url";
|
||||
import { AcpBackendRegistry } from "./acp/backend-registry.js";
|
||||
import { AcpSessionManager } from "./acp/session-manager.js";
|
||||
import { defaultStateFile, loadConfig, type AppConfig } from "./config.js";
|
||||
import type { PlatformAdapter } from "./core/adapter.js";
|
||||
import { DurableSessionStore } from "./core/durable-session-store.js";
|
||||
import { Gateway } from "./core/gateway.js";
|
||||
import { PlatformRegistry } from "./core/platform-registry.js";
|
||||
import { FeishuAdapter } from "./platforms/feishu/adapter.js";
|
||||
import { QqAdapter } from "./platforms/qq/adapter.js";
|
||||
import { QqGatewayClient } from "./platforms/qq/gateway-client.js";
|
||||
import { WeComAdapter } from "./platforms/wecom/adapter.js";
|
||||
import { GenericWebhookAdapter } from "./platforms/webhook/adapter.js";
|
||||
import { WeixinAdapter } from "./platforms/weixin/adapter.js";
|
||||
import { RoleRegistry } from "./roles/role-registry.js";
|
||||
import { BotProfileResolver } from "./roles/role-registry.js";
|
||||
|
||||
export interface GatewayRuntime {
|
||||
app: express.Express;
|
||||
gateway: Gateway;
|
||||
sessionManager: AcpSessionManager;
|
||||
store: DurableSessionStore;
|
||||
platformAdapter: PlatformAdapter;
|
||||
qqGatewayClient?: QqGatewayClient;
|
||||
shutdown(): Promise<void>;
|
||||
}
|
||||
|
||||
export async function createGatewayRuntime(config: AppConfig): Promise<GatewayRuntime> {
|
||||
const store = new DurableSessionStore(defaultStateFile(config));
|
||||
const platformType = config.gateway.platform.type;
|
||||
const store = new DurableSessionStore(defaultStateFile(config), { botId: config.bot.id, platform: platformType });
|
||||
await store.open();
|
||||
const roles = new RoleRegistry(config);
|
||||
const sessionManager = new AcpSessionManager(config.acp, new AcpBackendRegistry(config.backends), roles, store);
|
||||
const gateway = new Gateway(config.policy, sessionManager, roles);
|
||||
const platforms = new PlatformRegistry();
|
||||
const qqAdapter = new QqAdapter(config.platforms.qq, gateway);
|
||||
const adapters: PlatformAdapter[] = [
|
||||
new FeishuAdapter(config.platforms.feishu, gateway), new WeComAdapter(config.platforms.wecom), qqAdapter,
|
||||
new GenericWebhookAdapter(config.platforms.webhook, gateway), new WeixinAdapter(config.platforms.weixin, gateway)
|
||||
];
|
||||
for (const adapter of adapters) platforms.register(adapter);
|
||||
try {
|
||||
const bot = new BotProfileResolver(config).bot;
|
||||
const sessionManager = new AcpSessionManager(config.runtime.acp, bot, store);
|
||||
const gateway = new Gateway(config.gateway.policy, sessionManager);
|
||||
const platformAdapter = createPlatformAdapter(config, gateway);
|
||||
const qqGatewayClient = config.gateway.platform.type === "qq" && config.gateway.platform.connectionMode === "websocket"
|
||||
? new QqGatewayClient(config.gateway.platform, platformAdapter as QqAdapter) : undefined;
|
||||
|
||||
const app = express();
|
||||
app.use(express.json({ limit: "1mb", verify: (req, _res, buf) => {
|
||||
(req as express.Request & { rawBody?: Buffer }).rawBody = Buffer.from(buf);
|
||||
} }));
|
||||
app.get("/health", (_req, res) => res.json({ ok: true, acp: gateway.stats() }));
|
||||
app.get("/platforms", (_req, res) => res.json({ ok: true, platforms: platforms.list(), roles: roles.list(), backends: config.backends.map(({ id }) => id) }));
|
||||
const app = express();
|
||||
app.use(express.json({ limit: "1mb", verify: (req, _res, buf) => {
|
||||
(req as express.Request & { rawBody?: Buffer }).rawBody = Buffer.from(buf);
|
||||
} }));
|
||||
app.get("/health", (_req, res) => res.json({
|
||||
ok: true,
|
||||
configVersion: config.configVersion,
|
||||
botId: config.bot.id,
|
||||
platform: platformType,
|
||||
acp: gateway.stats()
|
||||
}));
|
||||
app.get("/platforms", (_req, res) => res.json({ ok: true, botId: config.bot.id, platform: platformType }));
|
||||
|
||||
function mountWebhook(routeName: string, adapterName = routeName): void {
|
||||
app.post(`/webhook/${routeName}`, async (req, res) => {
|
||||
const routeName = platformType === "webhook" ? "generic" : platformType;
|
||||
const needsWebhook = config.gateway.platform.type !== "qq" || config.gateway.platform.connectionMode === "webhook";
|
||||
if (needsWebhook) app.post(`/webhook/${routeName}`, async (req, res) => {
|
||||
try {
|
||||
const response = await platforms.get(adapterName).handleWebhook({
|
||||
const response = await platformAdapter.handleWebhook({
|
||||
req, body: req.body, headers: req.headers, query: req.query,
|
||||
rawBody: (req as express.Request & { rawBody?: Buffer }).rawBody
|
||||
});
|
||||
if (response.headers) for (const [key, value] of Object.entries(response.headers)) res.setHeader(key, value);
|
||||
res.status(response.status || 200).json(response.body ?? { ok: true });
|
||||
} catch (error) {
|
||||
const message = error instanceof Error ? error.message : String(error);
|
||||
console.error(`Webhook ${routeName} failed`, error);
|
||||
res.status(500).json({ ok: false, error: message });
|
||||
console.error(`Webhook ${routeName} failed (${error instanceof Error ? error.name : "unknown error"})`);
|
||||
res.status(500).json({ ok: false, error: "Internal webhook error" });
|
||||
}
|
||||
});
|
||||
}
|
||||
for (const name of ["feishu", "wecom", "qq", "weixin"]) mountWebhook(name);
|
||||
mountWebhook("generic", "webhook");
|
||||
|
||||
const qqGatewayClient = config.platforms.qq.enabled && config.platforms.qq.connectionMode === "websocket"
|
||||
? new QqGatewayClient(config.platforms.qq, qqAdapter) : undefined;
|
||||
let closing: Promise<void> | undefined;
|
||||
const runtime: GatewayRuntime = {
|
||||
app, gateway, sessionManager, store, qqGatewayClient,
|
||||
shutdown: () => closing ||= (async () => {
|
||||
qqGatewayClient?.stop();
|
||||
await sessionManager.shutdown();
|
||||
await store.close();
|
||||
})()
|
||||
};
|
||||
return runtime;
|
||||
let closing: Promise<void> | undefined;
|
||||
return {
|
||||
app, gateway, sessionManager, store, platformAdapter, qqGatewayClient,
|
||||
shutdown: () => closing ||= (async () => {
|
||||
qqGatewayClient?.stop();
|
||||
const results = await Promise.allSettled([sessionManager.shutdown(), store.close()]);
|
||||
const failed = results.find((result): result is PromiseRejectedResult => result.status === "rejected");
|
||||
if (failed) throw failed.reason;
|
||||
})()
|
||||
};
|
||||
} catch (error) {
|
||||
await store.close();
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
export async function createApp(config: AppConfig): Promise<express.Express> { return (await createGatewayRuntime(config)).app; }
|
||||
function createPlatformAdapter(config: AppConfig, gateway: Gateway): PlatformAdapter {
|
||||
const platform = config.gateway.platform;
|
||||
switch (platform.type) {
|
||||
case "feishu": return new FeishuAdapter(platform, gateway);
|
||||
case "wecom": return new WeComAdapter(platform);
|
||||
case "qq": return new QqAdapter(platform, gateway);
|
||||
case "webhook": return new GenericWebhookAdapter(platform, gateway);
|
||||
case "weixin": return new WeixinAdapter(platform, gateway);
|
||||
}
|
||||
}
|
||||
|
||||
export interface RunningServer { server: Server; runtime: GatewayRuntime; shutdown(): Promise<void> }
|
||||
|
||||
export async function startServer(config: AppConfig): Promise<RunningServer> {
|
||||
const runtime = await createGatewayRuntime(config);
|
||||
const server = runtime.app.listen(config.server.port, config.server.host, () => {
|
||||
console.log(`gori-agent listening on ${config.server.host}:${config.server.port}`);
|
||||
if (runtime.qqGatewayClient) { console.log("QQ websocket gateway enabled; connecting to QQ..."); runtime.qqGatewayClient.start(); }
|
||||
});
|
||||
const server = runtime.app.listen(config.gateway.server.port, config.gateway.server.host);
|
||||
try {
|
||||
await new Promise<void>((resolve, reject) => {
|
||||
server.once("listening", resolve);
|
||||
server.once("error", reject);
|
||||
});
|
||||
console.log(`gori-agent '${config.bot.id}' listening on ${config.gateway.server.host}:${config.gateway.server.port}`);
|
||||
if (runtime.qqGatewayClient) {
|
||||
console.log("QQ websocket gateway enabled; connecting to QQ...");
|
||||
runtime.qqGatewayClient.start();
|
||||
}
|
||||
} catch (error) {
|
||||
if (server.listening) await closeServer(server).catch(() => undefined);
|
||||
await runtime.shutdown().catch(() => undefined);
|
||||
throw error;
|
||||
}
|
||||
let closing: Promise<void> | undefined;
|
||||
return { server, runtime, shutdown: () => closing ||= (async () => {
|
||||
runtime.qqGatewayClient?.stop();
|
||||
const runtimeShutdown = runtime.shutdown();
|
||||
await new Promise<void>((resolve, reject) => server.close((error) => error ? reject(error) : resolve()));
|
||||
await runtimeShutdown;
|
||||
const results = await Promise.allSettled([closeServer(server), runtime.shutdown()]);
|
||||
const failed = results.find((result): result is PromiseRejectedResult => result.status === "rejected");
|
||||
if (failed) throw failed.reason;
|
||||
})() };
|
||||
}
|
||||
|
||||
function closeServer(server: Server): Promise<void> {
|
||||
if (!server.listening) return Promise.resolve();
|
||||
return new Promise<void>((resolve, reject) => server.close((error) => error ? reject(error) : resolve()));
|
||||
}
|
||||
|
||||
if (process.argv[1] && fileURLToPath(import.meta.url) === process.argv[1]) {
|
||||
void startServer(loadConfig()).catch((error) => { console.error(error); process.exitCode = 1; });
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user