Initial gori agent gateway
This commit is contained in:
@@ -0,0 +1,33 @@
|
||||
import type { AgentRunner } from "../core/types.js";
|
||||
|
||||
export class AgentRegistry {
|
||||
private readonly agents = new Map<string, AgentRunner>();
|
||||
|
||||
constructor(private readonly defaultAgentName: string) {}
|
||||
|
||||
register(agent: AgentRunner): void {
|
||||
if (this.agents.has(agent.name)) {
|
||||
throw new Error(`Agent already registered: ${agent.name}`);
|
||||
}
|
||||
this.agents.set(agent.name, agent);
|
||||
}
|
||||
|
||||
get(name?: string): AgentRunner {
|
||||
const requested = name || this.defaultAgentName;
|
||||
const agent = this.agents.get(requested);
|
||||
if (!agent) throw new Error(`Unknown agent: ${requested}`);
|
||||
return agent;
|
||||
}
|
||||
|
||||
has(name: string): boolean {
|
||||
return this.agents.has(name);
|
||||
}
|
||||
|
||||
defaultName(): string {
|
||||
return this.defaultAgentName;
|
||||
}
|
||||
|
||||
list(): string[] {
|
||||
return [...this.agents.keys()].sort();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,80 @@
|
||||
import { spawn } from "node:child_process";
|
||||
import { once } from "node:events";
|
||||
import type { CliAgentConfig } from "../config.js";
|
||||
import type { AgentRequest, AgentResponse, AgentRunner } from "../core/types.js";
|
||||
|
||||
export class CliAgent implements AgentRunner {
|
||||
readonly name: string;
|
||||
|
||||
constructor(private readonly config: CliAgentConfig) {
|
||||
this.name = config.name;
|
||||
}
|
||||
|
||||
async run(request: AgentRequest): Promise<AgentResponse> {
|
||||
const args = this.config.inputMode === "arg"
|
||||
? [...this.config.args, request.input]
|
||||
: this.config.args;
|
||||
|
||||
const child = spawn(this.config.command, args, {
|
||||
cwd: this.config.cwd,
|
||||
shell: false,
|
||||
stdio: ["pipe", "pipe", "pipe"],
|
||||
env: {
|
||||
...process.env,
|
||||
GORI_AGENT_NAME: this.name,
|
||||
GORI_AGENT_SESSION_ID: request.sessionId,
|
||||
GORI_AGENT_PLATFORM: request.platform,
|
||||
GORI_AGENT_CHAT_ID: request.chatId,
|
||||
GORI_AGENT_USER_ID: request.userId,
|
||||
GORI_AGENT_MESSAGE_ID: request.messageId || ""
|
||||
}
|
||||
});
|
||||
|
||||
let stdout = "";
|
||||
let stderr = "";
|
||||
let killedByTimeout = false;
|
||||
const maxBytes = this.config.outputMaxBytes;
|
||||
|
||||
const appendLimited = (current: string, chunk: Buffer): string => {
|
||||
if (Buffer.byteLength(current) >= maxBytes) return current;
|
||||
const combined = Buffer.concat([Buffer.from(current), chunk]);
|
||||
return combined.subarray(0, maxBytes).toString("utf8");
|
||||
};
|
||||
|
||||
child.stdout.on("data", (chunk: Buffer) => {
|
||||
stdout = appendLimited(stdout, chunk);
|
||||
});
|
||||
child.stderr.on("data", (chunk: Buffer) => {
|
||||
stderr = appendLimited(stderr, chunk);
|
||||
});
|
||||
|
||||
const timer = setTimeout(() => {
|
||||
killedByTimeout = true;
|
||||
child.kill("SIGTERM");
|
||||
setTimeout(() => {
|
||||
child.kill("SIGKILL");
|
||||
}, 2_000).unref();
|
||||
}, this.config.timeoutMs);
|
||||
|
||||
if (this.config.inputMode === "stdin") {
|
||||
child.stdin.end(request.input);
|
||||
} else {
|
||||
child.stdin.end();
|
||||
}
|
||||
|
||||
const [code, signal] = await once(child, "close") as [number | null, NodeJS.Signals | null];
|
||||
clearTimeout(timer);
|
||||
|
||||
if (killedByTimeout) {
|
||||
throw new Error(`Agent '${this.name}' timed out after ${this.config.timeoutMs}ms`);
|
||||
}
|
||||
if (code !== 0) {
|
||||
throw new Error(`Agent '${this.name}' exited with code ${code ?? signal}: ${stderr || stdout}`.trim());
|
||||
}
|
||||
|
||||
return {
|
||||
agentName: this.name,
|
||||
text: stdout.trim() || stderr.trim() || "(empty response)"
|
||||
};
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,85 @@
|
||||
import fs from "node:fs";
|
||||
import path from "node:path";
|
||||
import { z } from "zod";
|
||||
|
||||
const cliAgentSchema = z.object({
|
||||
name: z.string().min(1),
|
||||
command: z.string().min(1),
|
||||
args: z.array(z.string()).default([]),
|
||||
inputMode: z.enum(["stdin", "arg"]).default("stdin"),
|
||||
timeoutMs: z.number().int().positive().default(120_000),
|
||||
outputMaxBytes: z.number().int().positive().default(64_000),
|
||||
cwd: z.string().optional()
|
||||
});
|
||||
|
||||
const policySchema = z.object({
|
||||
allowedUsers: z.array(z.string()).default([]),
|
||||
allowedChats: z.array(z.string()).default([]),
|
||||
requireMentionInGroup: z.boolean().default(false)
|
||||
});
|
||||
|
||||
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),
|
||||
appId: z.string().default(""),
|
||||
clientSecret: z.string().default("")
|
||||
});
|
||||
|
||||
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("")
|
||||
});
|
||||
|
||||
export const configSchema = z.object({
|
||||
server: z.object({
|
||||
host: z.string().default("0.0.0.0"),
|
||||
port: z.number().int().positive().max(65_535).default(3000)
|
||||
}).default({}),
|
||||
policy: policySchema.default({}),
|
||||
defaultAgent: z.string().min(1).default("echo"),
|
||||
agents: z.array(cliAgentSchema).min(1),
|
||||
platforms: z.object({
|
||||
feishu: feishuSchema.default({}),
|
||||
wecom: wecomSchema.default({}),
|
||||
qq: qqSchema.default({}),
|
||||
webhook: webhookSchema.default({}),
|
||||
weixin: weixinSchema.default({})
|
||||
}).default({})
|
||||
});
|
||||
|
||||
export type AppConfig = z.infer<typeof configSchema>;
|
||||
export type CliAgentConfig = z.infer<typeof cliAgentSchema>;
|
||||
export type GatewayPolicy = z.infer<typeof policySchema>;
|
||||
|
||||
export function loadConfig(configPath = process.env.GORI_GATEWAY_CONFIG): AppConfig {
|
||||
const resolvedPath = path.resolve(configPath || "config.example.json");
|
||||
const raw = fs.readFileSync(resolvedPath, "utf8");
|
||||
const parsed = JSON.parse(raw) as unknown;
|
||||
const config = configSchema.parse(parsed);
|
||||
|
||||
if (!config.agents.some((agent) => agent.name === config.defaultAgent)) {
|
||||
throw new Error(`defaultAgent '${config.defaultAgent}' is not present in agents`);
|
||||
}
|
||||
|
||||
return config;
|
||||
}
|
||||
@@ -0,0 +1,15 @@
|
||||
import type { OutgoingMessage, WebhookRequestContext, WebhookResponse } from "./types.js";
|
||||
|
||||
export interface PlatformAdapter {
|
||||
name: string;
|
||||
handleWebhook(context: WebhookRequestContext): Promise<WebhookResponse>;
|
||||
sendMessage(message: OutgoingMessage): Promise<void>;
|
||||
}
|
||||
|
||||
export function jsonResponse(body: unknown, status = 200): WebhookResponse {
|
||||
return { status, body };
|
||||
}
|
||||
|
||||
export function notImplemented(platform: string): WebhookResponse {
|
||||
return jsonResponse({ ok: false, error: `${platform} inbound webhook is not implemented in v1` }, 501);
|
||||
}
|
||||
@@ -0,0 +1,57 @@
|
||||
import { AgentRegistry } from "../agents/agent-registry.js";
|
||||
import { SessionStore } from "./session-store.js";
|
||||
import type { IncomingMessage } from "./types.js";
|
||||
|
||||
export interface CommandResult {
|
||||
handled: boolean;
|
||||
text?: string;
|
||||
}
|
||||
|
||||
export class CommandRouter {
|
||||
constructor(
|
||||
private readonly agents: AgentRegistry,
|
||||
private readonly sessions: SessionStore
|
||||
) {}
|
||||
|
||||
route(message: IncomingMessage, sessionId: string): CommandResult {
|
||||
const text = message.text.trim();
|
||||
if (!text.startsWith("/")) return { handled: false };
|
||||
|
||||
const [command, ...args] = text.split(/\s+/);
|
||||
switch (command) {
|
||||
case "/help":
|
||||
return {
|
||||
handled: true,
|
||||
text: [
|
||||
"Commands:",
|
||||
"/help - show this help",
|
||||
"/status - show gateway/session status",
|
||||
"/agents - list available agents",
|
||||
"/agent <name> - select an agent for this chat",
|
||||
"/new - reset this chat session"
|
||||
].join("\n")
|
||||
};
|
||||
case "/status": {
|
||||
const session = this.sessions.getSession(sessionId);
|
||||
return {
|
||||
handled: true,
|
||||
text: `OK\nplatform=${message.platform}\nchat=${message.chatId}\nagent=${session.selectedAgent || this.agents.defaultName()}`
|
||||
};
|
||||
}
|
||||
case "/agents":
|
||||
return { handled: true, text: `Available agents: ${this.agents.list().join(", ")}` };
|
||||
case "/agent": {
|
||||
const agentName = args[0];
|
||||
if (!agentName) return { handled: true, text: "Usage: /agent <name>" };
|
||||
if (!this.agents.has(agentName)) return { handled: true, text: `Unknown agent: ${agentName}` };
|
||||
this.sessions.setSelectedAgent(sessionId, agentName);
|
||||
return { handled: true, text: `Selected agent: ${agentName}` };
|
||||
}
|
||||
case "/new":
|
||||
this.sessions.reset(sessionId);
|
||||
return { handled: true, text: "Started a new session for this chat." };
|
||||
default:
|
||||
return { handled: false };
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,112 @@
|
||||
import type { GatewayPolicy } from "../config.js";
|
||||
import { AgentRegistry } from "../agents/agent-registry.js";
|
||||
import type { PlatformAdapter } from "./adapter.js";
|
||||
import { CommandRouter } from "./command-router.js";
|
||||
import { SessionStore, sessionIdFor } from "./session-store.js";
|
||||
import type { IncomingMessage } from "./types.js";
|
||||
|
||||
export interface GatewayResult {
|
||||
ok: boolean;
|
||||
reply?: string;
|
||||
ignored?: boolean;
|
||||
error?: string;
|
||||
}
|
||||
|
||||
export class Gateway {
|
||||
private readonly locks = new Map<string, Promise<void>>();
|
||||
readonly commandRouter: CommandRouter;
|
||||
|
||||
constructor(
|
||||
private readonly policy: GatewayPolicy,
|
||||
private readonly agents: AgentRegistry,
|
||||
private readonly sessions: SessionStore
|
||||
) {
|
||||
this.commandRouter = new CommandRouter(agents, sessions);
|
||||
}
|
||||
|
||||
async receive(message: IncomingMessage, adapter: PlatformAdapter, options: { synchronous?: boolean } = {}): Promise<GatewayResult> {
|
||||
const policyError = this.checkPolicy(message);
|
||||
if (policyError) return { ok: true, ignored: true, error: policyError };
|
||||
|
||||
const sessionId = sessionIdFor(message.platform, message.chatId);
|
||||
return this.withChatLock(sessionId, async () => {
|
||||
const command = this.commandRouter.route(message, sessionId);
|
||||
if (command.handled) {
|
||||
const reply = command.text || "";
|
||||
await adapter.sendMessage({
|
||||
target: { platform: message.platform, chatId: message.chatId, userId: message.userId, raw: message.raw },
|
||||
text: reply,
|
||||
replyTo: message.messageId
|
||||
});
|
||||
return { ok: true, reply };
|
||||
}
|
||||
|
||||
const session = this.sessions.getSession(sessionId);
|
||||
const agent = this.agents.get(session.selectedAgent);
|
||||
try {
|
||||
const response = await agent.run({
|
||||
input: message.text,
|
||||
sessionId,
|
||||
platform: message.platform,
|
||||
chatId: message.chatId,
|
||||
userId: message.userId,
|
||||
messageId: message.messageId
|
||||
});
|
||||
await adapter.sendMessage({
|
||||
target: { platform: message.platform, chatId: message.chatId, userId: message.userId, raw: message.raw },
|
||||
text: response.text,
|
||||
replyTo: message.messageId
|
||||
});
|
||||
return { ok: true, reply: response.text };
|
||||
} catch (error) {
|
||||
const errorText = error instanceof Error ? error.message : String(error);
|
||||
const reply = `Agent error: ${errorText}`;
|
||||
if (!options.synchronous) {
|
||||
await adapter.sendMessage({
|
||||
target: { platform: message.platform, chatId: message.chatId, userId: message.userId, raw: message.raw },
|
||||
text: reply,
|
||||
replyTo: message.messageId
|
||||
});
|
||||
}
|
||||
return { ok: false, error: errorText, reply };
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
stats(): { sessions: number; lockedChats: number } {
|
||||
return { ...this.sessions.stats(), lockedChats: this.locks.size };
|
||||
}
|
||||
|
||||
private checkPolicy(message: IncomingMessage): string | undefined {
|
||||
if (this.policy.allowedUsers.length > 0 && !this.policy.allowedUsers.includes(message.userId)) {
|
||||
return `User not allowed: ${message.userId}`;
|
||||
}
|
||||
if (this.policy.allowedChats.length > 0 && !this.policy.allowedChats.includes(message.chatId)) {
|
||||
return `Chat not allowed: ${message.chatId}`;
|
||||
}
|
||||
if (this.policy.requireMentionInGroup && message.isGroup && !message.mentionsBot) {
|
||||
return "Mention required in group chat";
|
||||
}
|
||||
return undefined;
|
||||
}
|
||||
|
||||
private async withChatLock<T>(key: string, fn: () => Promise<T>): Promise<T> {
|
||||
const previous = this.locks.get(key) || Promise.resolve();
|
||||
let release!: () => void;
|
||||
const current = new Promise<void>((resolve) => {
|
||||
release = resolve;
|
||||
});
|
||||
const queued = previous.then(() => current);
|
||||
this.locks.set(key, queued);
|
||||
|
||||
await previous;
|
||||
try {
|
||||
return await fn();
|
||||
} finally {
|
||||
release();
|
||||
if (this.locks.get(key) === queued) {
|
||||
this.locks.delete(key);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,22 @@
|
||||
import type { PlatformAdapter } from "./adapter.js";
|
||||
|
||||
export class PlatformRegistry {
|
||||
private readonly adapters = new Map<string, PlatformAdapter>();
|
||||
|
||||
register(adapter: PlatformAdapter): void {
|
||||
if (this.adapters.has(adapter.name)) {
|
||||
throw new Error(`Platform adapter already registered: ${adapter.name}`);
|
||||
}
|
||||
this.adapters.set(adapter.name, adapter);
|
||||
}
|
||||
|
||||
get(name: string): PlatformAdapter {
|
||||
const adapter = this.adapters.get(name);
|
||||
if (!adapter) throw new Error(`Unknown platform adapter: ${name}`);
|
||||
return adapter;
|
||||
}
|
||||
|
||||
list(): string[] {
|
||||
return [...this.adapters.keys()].sort();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,37 @@
|
||||
export interface ChatSession {
|
||||
selectedAgent?: string;
|
||||
createdAt: number;
|
||||
updatedAt: number;
|
||||
}
|
||||
|
||||
export class SessionStore {
|
||||
private readonly sessions = new Map<string, ChatSession>();
|
||||
|
||||
getSession(sessionId: string): ChatSession {
|
||||
const existing = this.sessions.get(sessionId);
|
||||
if (existing) return existing;
|
||||
|
||||
const now = Date.now();
|
||||
const session: ChatSession = { createdAt: now, updatedAt: now };
|
||||
this.sessions.set(sessionId, session);
|
||||
return session;
|
||||
}
|
||||
|
||||
setSelectedAgent(sessionId: string, agentName: string): void {
|
||||
const session = this.getSession(sessionId);
|
||||
session.selectedAgent = agentName;
|
||||
session.updatedAt = Date.now();
|
||||
}
|
||||
|
||||
reset(sessionId: string): void {
|
||||
this.sessions.delete(sessionId);
|
||||
}
|
||||
|
||||
stats(): { sessions: number } {
|
||||
return { sessions: this.sessions.size };
|
||||
}
|
||||
}
|
||||
|
||||
export function sessionIdFor(platform: string, chatId: string): string {
|
||||
return `${platform}:${chatId}`;
|
||||
}
|
||||
@@ -0,0 +1,36 @@
|
||||
export function stripBotMentions(text: string): string {
|
||||
return text
|
||||
.replace(/<at[^>]*>.*?<\/at>/g, " ")
|
||||
.replace(/<@!?[A-Za-z0-9_:-]+>/g, " ")
|
||||
.replace(/@[\w\-\u4e00-\u9fa5]+/g, " ")
|
||||
.replace(/\s+/g, " ")
|
||||
.trim();
|
||||
}
|
||||
|
||||
export function flattenText(value: unknown): string {
|
||||
if (typeof value === "string") return value;
|
||||
if (Array.isArray(value)) return value.map(flattenText).filter(Boolean).join(" ");
|
||||
if (!value || typeof value !== "object") return "";
|
||||
|
||||
const record = value as Record<string, unknown>;
|
||||
if (typeof record.text === "string") return record.text;
|
||||
if (typeof record.content === "string") return record.content;
|
||||
if (Array.isArray(record.content)) return flattenText(record.content);
|
||||
if (Array.isArray(record.elements)) return flattenText(record.elements);
|
||||
if (Array.isArray(record.children)) return flattenText(record.children);
|
||||
if (record.tag === "text" && typeof record.un_escape_text === "string") return record.un_escape_text;
|
||||
|
||||
return Object.values(record).map(flattenText).filter(Boolean).join(" ");
|
||||
}
|
||||
|
||||
export function parseMaybeJson(text: string): unknown {
|
||||
try {
|
||||
return JSON.parse(text);
|
||||
} catch {
|
||||
return text;
|
||||
}
|
||||
}
|
||||
|
||||
export function firstHeader(value: string | string[] | undefined): string | undefined {
|
||||
return Array.isArray(value) ? value[0] : value;
|
||||
}
|
||||
@@ -0,0 +1,61 @@
|
||||
import type { Request } from "express";
|
||||
|
||||
export type PlatformName = "feishu" | "wecom" | "qq" | "webhook" | "weixin" | string;
|
||||
|
||||
export interface MessageTarget {
|
||||
platform: PlatformName;
|
||||
chatId: string;
|
||||
userId?: string;
|
||||
threadId?: string;
|
||||
raw?: unknown;
|
||||
}
|
||||
|
||||
export interface IncomingMessage {
|
||||
platform: PlatformName;
|
||||
chatId: string;
|
||||
userId: string;
|
||||
text: string;
|
||||
messageId?: string;
|
||||
isGroup?: boolean;
|
||||
mentionsBot?: boolean;
|
||||
raw?: unknown;
|
||||
}
|
||||
|
||||
export interface OutgoingMessage {
|
||||
target: MessageTarget;
|
||||
text: string;
|
||||
replyTo?: string;
|
||||
}
|
||||
|
||||
export interface WebhookRequestContext {
|
||||
req: Request;
|
||||
body: unknown;
|
||||
headers: Request["headers"];
|
||||
query: Request["query"];
|
||||
rawBody?: Buffer;
|
||||
}
|
||||
|
||||
export interface WebhookResponse {
|
||||
status?: number;
|
||||
headers?: Record<string, string>;
|
||||
body?: unknown;
|
||||
}
|
||||
|
||||
export interface AgentRequest {
|
||||
input: string;
|
||||
sessionId: string;
|
||||
platform: PlatformName;
|
||||
chatId: string;
|
||||
userId: string;
|
||||
messageId?: string;
|
||||
}
|
||||
|
||||
export interface AgentResponse {
|
||||
text: string;
|
||||
agentName: string;
|
||||
}
|
||||
|
||||
export interface AgentRunner {
|
||||
name: string;
|
||||
run(request: AgentRequest): Promise<AgentResponse>;
|
||||
}
|
||||
@@ -0,0 +1,114 @@
|
||||
import type { AppConfig } from "../../config.js";
|
||||
import { jsonResponse } from "../../core/adapter.js";
|
||||
import type { PlatformAdapter } from "../../core/adapter.js";
|
||||
import type { Gateway } from "../../core/gateway.js";
|
||||
import { flattenText, parseMaybeJson, stripBotMentions } from "../../core/text.js";
|
||||
import type { IncomingMessage, OutgoingMessage, WebhookRequestContext, WebhookResponse } from "../../core/types.js";
|
||||
import type { FeishuTenantTokenResponse, FeishuWebhookBody } from "./types.js";
|
||||
|
||||
export class FeishuAdapter implements PlatformAdapter {
|
||||
readonly name = "feishu";
|
||||
private tenantToken?: { token: string; expiresAt: number };
|
||||
|
||||
constructor(
|
||||
private readonly config: AppConfig["platforms"]["feishu"],
|
||||
private readonly gateway: Gateway
|
||||
) {}
|
||||
|
||||
async handleWebhook(context: WebhookRequestContext): Promise<WebhookResponse> {
|
||||
const body = context.body as FeishuWebhookBody;
|
||||
if (body.challenge) return jsonResponse({ challenge: body.challenge });
|
||||
|
||||
const token = body.header?.token || body.token;
|
||||
if (this.config.verificationToken && token !== this.config.verificationToken) {
|
||||
return jsonResponse({ ok: false, error: "invalid verification token" }, 401);
|
||||
}
|
||||
|
||||
if (body.header?.event_type !== "im.message.receive_v1") {
|
||||
return jsonResponse({ ok: true, ignored: true });
|
||||
}
|
||||
|
||||
const message = body.event?.message;
|
||||
const sender = body.event?.sender?.sender_id;
|
||||
if (!message?.chat_id || !sender || !message.content) {
|
||||
return jsonResponse({ ok: false, error: "missing Feishu message fields" }, 400);
|
||||
}
|
||||
|
||||
const content = parseMaybeJson(message.content);
|
||||
const rawText = flattenText(content);
|
||||
const mentionsBot = this.hasBotMention(rawText, message.mentions);
|
||||
const normalized: IncomingMessage = {
|
||||
platform: this.name,
|
||||
chatId: message.chat_id,
|
||||
userId: sender.open_id || sender.user_id || sender.union_id || "unknown",
|
||||
text: stripBotMentions(rawText),
|
||||
messageId: message.message_id,
|
||||
isGroup: message.chat_type !== "p2p",
|
||||
mentionsBot,
|
||||
raw: body
|
||||
};
|
||||
|
||||
void this.gateway.receive(normalized, this).catch((error) => {
|
||||
console.error("Feishu gateway error", error);
|
||||
});
|
||||
return jsonResponse({ ok: true });
|
||||
}
|
||||
|
||||
async sendMessage(message: OutgoingMessage): Promise<void> {
|
||||
if (!message.replyTo) {
|
||||
console.warn("Feishu sendMessage requires replyTo message_id; skipping");
|
||||
return;
|
||||
}
|
||||
const token = await this.getTenantAccessToken();
|
||||
const response = await fetch(`https://open.feishu.cn/open-apis/im/v1/messages/${encodeURIComponent(message.replyTo)}/reply`, {
|
||||
method: "POST",
|
||||
headers: {
|
||||
Authorization: `Bearer ${token}`,
|
||||
"Content-Type": "application/json; charset=utf-8"
|
||||
},
|
||||
body: JSON.stringify({
|
||||
msg_type: "text",
|
||||
content: JSON.stringify({ text: message.text })
|
||||
})
|
||||
});
|
||||
|
||||
if (!response.ok) {
|
||||
throw new Error(`Feishu reply failed: ${response.status} ${await response.text()}`);
|
||||
}
|
||||
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());
|
||||
}
|
||||
}
|
||||
|
||||
private async getTenantAccessToken(): Promise<string> {
|
||||
if (this.tenantToken && this.tenantToken.expiresAt > Date.now() + 60_000) {
|
||||
return this.tenantToken.token;
|
||||
}
|
||||
const response = await fetch("https://open.feishu.cn/open-apis/auth/v3/tenant_access_token/internal", {
|
||||
method: "POST",
|
||||
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()}`);
|
||||
}
|
||||
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());
|
||||
}
|
||||
this.tenantToken = {
|
||||
token: data.tenant_access_token,
|
||||
expiresAt: Date.now() + Math.max(1, (data.expire || 7200) - 120) * 1000
|
||||
};
|
||||
return this.tenantToken.token;
|
||||
}
|
||||
|
||||
private hasBotMention(
|
||||
text: string,
|
||||
mentions: NonNullable<NonNullable<FeishuWebhookBody["event"]>["message"]>["mentions"] | undefined
|
||||
): boolean {
|
||||
if (mentions && mentions.length > 0) return true;
|
||||
return this.config.botNames.some((name) => text.includes(`@${name}`) || text.includes(name));
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,38 @@
|
||||
export interface FeishuTenantTokenResponse {
|
||||
code: number;
|
||||
msg?: string;
|
||||
tenant_access_token?: string;
|
||||
expire?: number;
|
||||
}
|
||||
|
||||
export interface FeishuWebhookBody {
|
||||
challenge?: string;
|
||||
token?: string;
|
||||
header?: {
|
||||
token?: string;
|
||||
event_type?: string;
|
||||
event_id?: string;
|
||||
};
|
||||
event?: {
|
||||
message?: {
|
||||
message_id?: string;
|
||||
chat_id?: string;
|
||||
chat_type?: string;
|
||||
content?: string;
|
||||
message_type?: string;
|
||||
mentions?: Array<{
|
||||
key?: string;
|
||||
name?: string;
|
||||
id?: { open_id?: string; user_id?: string; union_id?: string };
|
||||
}>;
|
||||
};
|
||||
sender?: {
|
||||
sender_id?: { open_id?: string; user_id?: string; union_id?: string };
|
||||
};
|
||||
};
|
||||
}
|
||||
|
||||
export interface FeishuReplyResponse {
|
||||
code: number;
|
||||
msg?: string;
|
||||
}
|
||||
@@ -0,0 +1,58 @@
|
||||
import type { AppConfig } 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";
|
||||
import type { QqAccessTokenResponse, QqSendMessageResponse } from "./types.js";
|
||||
|
||||
export class QqAdapter implements PlatformAdapter {
|
||||
readonly name = "qq";
|
||||
private accessToken?: { token: string; expiresAt: number };
|
||||
|
||||
constructor(private readonly config: AppConfig["platforms"]["qq"]) {}
|
||||
|
||||
async handleWebhook(_context: WebhookRequestContext): Promise<WebhookResponse> {
|
||||
return notImplemented("QQ");
|
||||
}
|
||||
|
||||
async sendMessage(message: OutgoingMessage): Promise<void> {
|
||||
const token = await this.getAccessToken();
|
||||
const raw = (message.target.raw || {}) as Record<string, unknown>;
|
||||
const groupOpenId = typeof raw.group_openid === "string" ? raw.group_openid : undefined;
|
||||
const userOpenId = typeof raw.user_openid === "string" ? raw.user_openid : undefined;
|
||||
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`;
|
||||
|
||||
const response = await fetch(`https://api.sgroup.qq.com${path}`, {
|
||||
method: "POST",
|
||||
headers: {
|
||||
Authorization: `QQBot ${token}`,
|
||||
"Content-Type": "application/json; charset=utf-8"
|
||||
},
|
||||
body: JSON.stringify({ content: message.text, msg_id: message.replyTo })
|
||||
});
|
||||
if (!response.ok) throw new Error(`QQ send failed: ${response.status} ${await response.text()}`);
|
||||
const data = await response.json() as QqSendMessageResponse;
|
||||
if (data.code && data.code !== 0) throw new Error(`QQ send failed: ${data.code} ${data.message || ""}`.trim());
|
||||
}
|
||||
|
||||
private async getAccessToken(): Promise<string> {
|
||||
if (this.accessToken && this.accessToken.expiresAt > Date.now() + 60_000) return this.accessToken.token;
|
||||
|
||||
const response = await fetch("https://bots.qq.com/app/getAppAccessToken", {
|
||||
method: "POST",
|
||||
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()}`);
|
||||
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());
|
||||
}
|
||||
this.accessToken = {
|
||||
token: data.access_token,
|
||||
expiresAt: Date.now() + Math.max(1, (data.expires_in || 7200) - 120) * 1000
|
||||
};
|
||||
return this.accessToken.token;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,13 @@
|
||||
export interface QqAccessTokenResponse {
|
||||
access_token?: string;
|
||||
expires_in?: number;
|
||||
token_type?: string;
|
||||
error?: string;
|
||||
error_description?: string;
|
||||
}
|
||||
|
||||
export interface QqSendMessageResponse {
|
||||
id?: string;
|
||||
code?: number;
|
||||
message?: string;
|
||||
}
|
||||
@@ -0,0 +1,68 @@
|
||||
import crypto from "node:crypto";
|
||||
import type { AppConfig } from "../../config.js";
|
||||
import { jsonResponse } from "../../core/adapter.js";
|
||||
import type { PlatformAdapter } from "../../core/adapter.js";
|
||||
import type { Gateway } from "../../core/gateway.js";
|
||||
import { firstHeader } from "../../core/text.js";
|
||||
import type { IncomingMessage, OutgoingMessage, WebhookRequestContext, WebhookResponse } from "../../core/types.js";
|
||||
|
||||
interface GenericWebhookPayload {
|
||||
chat_id?: string;
|
||||
user_id?: string;
|
||||
text?: string;
|
||||
message_id?: string;
|
||||
is_group?: boolean;
|
||||
mentions_bot?: boolean;
|
||||
}
|
||||
|
||||
export function verifyHmacSha256(rawBody: string, signature: string | undefined, secret: string): boolean {
|
||||
if (!secret) return true;
|
||||
if (!signature) return false;
|
||||
|
||||
const expectedHex = crypto.createHmac("sha256", secret).update(rawBody).digest("hex");
|
||||
const actualHex = signature.startsWith("sha256=") ? signature.slice("sha256=".length) : signature;
|
||||
|
||||
if (!/^[a-fA-F0-9]+$/.test(actualHex)) return false;
|
||||
const expected = Buffer.from(expectedHex, "hex");
|
||||
const actual = Buffer.from(actualHex, "hex");
|
||||
return expected.length === actual.length && crypto.timingSafeEqual(expected, actual);
|
||||
}
|
||||
|
||||
export class GenericWebhookAdapter implements PlatformAdapter {
|
||||
readonly name = "webhook";
|
||||
|
||||
constructor(
|
||||
private readonly config: AppConfig["platforms"]["webhook"],
|
||||
private readonly gateway: Gateway
|
||||
) {}
|
||||
|
||||
async handleWebhook(context: WebhookRequestContext): Promise<WebhookResponse> {
|
||||
const rawBody = context.rawBody?.toString("utf8") ?? JSON.stringify(context.body ?? {});
|
||||
const signature = firstHeader(context.headers["x-gori-signature"] as string | string[] | undefined);
|
||||
if (!verifyHmacSha256(rawBody, signature, this.config.secret)) {
|
||||
return jsonResponse({ ok: false, error: "invalid signature" }, 401);
|
||||
}
|
||||
|
||||
const payload = context.body as GenericWebhookPayload;
|
||||
if (!payload.chat_id || !payload.user_id || typeof payload.text !== "string") {
|
||||
return jsonResponse({ ok: false, error: "chat_id, user_id and text are required" }, 400);
|
||||
}
|
||||
|
||||
const message: IncomingMessage = {
|
||||
platform: this.name,
|
||||
chatId: payload.chat_id,
|
||||
userId: payload.user_id,
|
||||
text: payload.text,
|
||||
messageId: payload.message_id,
|
||||
isGroup: payload.is_group,
|
||||
mentionsBot: payload.mentions_bot,
|
||||
raw: payload
|
||||
};
|
||||
const result = await this.gateway.receive(message, this, { synchronous: true });
|
||||
return jsonResponse(result.ok ? { ok: true, output: result.reply } : { ok: false, error: result.error, output: result.reply }, result.ok ? 200 : 500);
|
||||
}
|
||||
|
||||
async sendMessage(_message: OutgoingMessage): Promise<void> {
|
||||
// Generic webhook responses are returned synchronously from handleWebhook.
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,53 @@
|
||||
import type { AppConfig } 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";
|
||||
import type { WeComAccessTokenResponse, WeComSendMessageResponse } from "./types.js";
|
||||
|
||||
export class WeComAdapter implements PlatformAdapter {
|
||||
readonly name = "wecom";
|
||||
private accessToken?: { token: string; expiresAt: number };
|
||||
|
||||
constructor(private readonly config: AppConfig["platforms"]["wecom"]) {}
|
||||
|
||||
async handleWebhook(_context: WebhookRequestContext): Promise<WebhookResponse> {
|
||||
return notImplemented("WeCom");
|
||||
}
|
||||
|
||||
async sendMessage(message: OutgoingMessage): Promise<void> {
|
||||
const token = await this.getAccessToken();
|
||||
const response = await fetch(`https://qyapi.weixin.qq.com/cgi-bin/message/send?access_token=${encodeURIComponent(token)}`, {
|
||||
method: "POST",
|
||||
headers: { "Content-Type": "application/json; charset=utf-8" },
|
||||
body: JSON.stringify({
|
||||
touser: message.target.userId || message.target.chatId,
|
||||
msgtype: "text",
|
||||
agentid: Number(this.config.agentId) || this.config.agentId,
|
||||
text: { content: message.text },
|
||||
safe: 0
|
||||
})
|
||||
});
|
||||
if (!response.ok) throw new Error(`WeCom send failed: ${response.status} ${await response.text()}`);
|
||||
const data = await response.json() as WeComSendMessageResponse;
|
||||
if (data.errcode !== 0) throw new Error(`WeCom send failed: ${data.errcode} ${data.errmsg || ""}`.trim());
|
||||
}
|
||||
|
||||
private async getAccessToken(): Promise<string> {
|
||||
if (this.accessToken && this.accessToken.expiresAt > Date.now() + 60_000) return this.accessToken.token;
|
||||
|
||||
const url = new URL("https://qyapi.weixin.qq.com/cgi-bin/gettoken");
|
||||
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()}`);
|
||||
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());
|
||||
}
|
||||
this.accessToken = {
|
||||
token: data.access_token,
|
||||
expiresAt: Date.now() + Math.max(1, (data.expires_in || 7200) - 120) * 1000
|
||||
};
|
||||
return this.accessToken.token;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,11 @@
|
||||
export interface WeComAccessTokenResponse {
|
||||
errcode: number;
|
||||
errmsg?: string;
|
||||
access_token?: string;
|
||||
expires_in?: number;
|
||||
}
|
||||
|
||||
export interface WeComSendMessageResponse {
|
||||
errcode: number;
|
||||
errmsg?: string;
|
||||
}
|
||||
@@ -0,0 +1,49 @@
|
||||
import type { AppConfig } from "../../config.js";
|
||||
import { jsonResponse, notImplemented } from "../../core/adapter.js";
|
||||
import type { PlatformAdapter } from "../../core/adapter.js";
|
||||
import type { Gateway } from "../../core/gateway.js";
|
||||
import type { IncomingMessage, OutgoingMessage, WebhookRequestContext, WebhookResponse } from "../../core/types.js";
|
||||
|
||||
interface ExternalWeixinPayload {
|
||||
chat_id?: string;
|
||||
user_id?: string;
|
||||
text?: string;
|
||||
message_id?: string;
|
||||
is_group?: boolean;
|
||||
mentions_bot?: boolean;
|
||||
}
|
||||
|
||||
export class WeixinAdapter implements PlatformAdapter {
|
||||
readonly name = "weixin";
|
||||
|
||||
constructor(
|
||||
private readonly config: AppConfig["platforms"]["weixin"],
|
||||
private readonly gateway: Gateway
|
||||
) {}
|
||||
|
||||
async handleWebhook(context: WebhookRequestContext): Promise<WebhookResponse> {
|
||||
if (this.config.mode === "not-implemented") return notImplemented("Weixin");
|
||||
|
||||
const payload = context.body as ExternalWeixinPayload;
|
||||
if (!payload.chat_id || !payload.user_id || typeof payload.text !== "string") {
|
||||
return jsonResponse({ ok: false, error: "external Weixin webhook requires chat_id, user_id and text" }, 400);
|
||||
}
|
||||
|
||||
const message: IncomingMessage = {
|
||||
platform: this.name,
|
||||
chatId: payload.chat_id,
|
||||
userId: payload.user_id,
|
||||
text: payload.text,
|
||||
messageId: payload.message_id,
|
||||
isGroup: payload.is_group,
|
||||
mentionsBot: payload.mentions_bot,
|
||||
raw: payload
|
||||
};
|
||||
const result = await this.gateway.receive(message, this, { synchronous: true });
|
||||
return jsonResponse(result.ok ? { ok: true, output: result.reply } : { ok: false, error: result.error, output: result.reply }, result.ok ? 200 : 500);
|
||||
}
|
||||
|
||||
async sendMessage(_message: OutgoingMessage): Promise<void> {
|
||||
// Personal WeChat has no official native iLink adapter in v1. External bridges should consume the synchronous webhook response.
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,78 @@
|
||||
import express from "express";
|
||||
import { loadConfig } from "./config.js";
|
||||
import { AgentRegistry } from "./agents/agent-registry.js";
|
||||
import { CliAgent } from "./agents/cli-agent.js";
|
||||
import { Gateway } from "./core/gateway.js";
|
||||
import { PlatformRegistry } from "./core/platform-registry.js";
|
||||
import { SessionStore } from "./core/session-store.js";
|
||||
import type { PlatformAdapter } from "./core/adapter.js";
|
||||
import { FeishuAdapter } from "./platforms/feishu/adapter.js";
|
||||
import { WeComAdapter } from "./platforms/wecom/adapter.js";
|
||||
import { QqAdapter } from "./platforms/qq/adapter.js";
|
||||
import { GenericWebhookAdapter } from "./platforms/webhook/adapter.js";
|
||||
import { WeixinAdapter } from "./platforms/weixin/adapter.js";
|
||||
|
||||
const config = loadConfig();
|
||||
const sessions = new SessionStore();
|
||||
const agents = new AgentRegistry(config.defaultAgent);
|
||||
for (const agentConfig of config.agents) {
|
||||
agents.register(new CliAgent(agentConfig));
|
||||
}
|
||||
|
||||
const gateway = new Gateway(config.policy, agents, sessions);
|
||||
const platforms = new PlatformRegistry();
|
||||
const adapters: PlatformAdapter[] = [
|
||||
new FeishuAdapter(config.platforms.feishu, gateway),
|
||||
new WeComAdapter(config.platforms.wecom),
|
||||
new QqAdapter(config.platforms.qq),
|
||||
new GenericWebhookAdapter(config.platforms.webhook, gateway),
|
||||
new WeixinAdapter(config.platforms.weixin, gateway)
|
||||
];
|
||||
for (const adapter of adapters) platforms.register(adapter);
|
||||
|
||||
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, gateway: gateway.stats() });
|
||||
});
|
||||
|
||||
app.get("/platforms", (_req, res) => {
|
||||
res.json({ ok: true, platforms: platforms.list(), agents: agents.list() });
|
||||
});
|
||||
|
||||
function mountWebhook(routeName: string, adapterName = routeName): void {
|
||||
app.post(`/webhook/${routeName}`, async (req, res) => {
|
||||
try {
|
||||
const response = await platforms.get(adapterName).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 });
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
for (const platformName of ["feishu", "wecom", "qq", "weixin"]) {
|
||||
mountWebhook(platformName);
|
||||
}
|
||||
mountWebhook("generic", "webhook");
|
||||
|
||||
app.listen(config.server.port, config.server.host, () => {
|
||||
console.log(`gori-agent-gateway listening on ${config.server.host}:${config.server.port}`);
|
||||
});
|
||||
Reference in New Issue
Block a user