Add QQ gateway runtime support
This commit is contained in:
+2
-1
@@ -32,9 +32,10 @@ export async function runDoctor(config: AppConfig, configPath: string): Promise<
|
||||
if (!config.platforms.feishu.verificationToken || config.platforms.feishu.verificationToken === "replace-me") warn("Feishu verificationToken is missing or placeholder.");
|
||||
}
|
||||
if (config.platforms.qq.enabled) {
|
||||
console.log(`QQ connection mode: ${config.platforms.qq.connectionMode}`);
|
||||
if (!config.platforms.qq.appId) problems += error("QQ appId is empty.");
|
||||
if (!config.platforms.qq.clientSecret || config.platforms.qq.clientSecret === "replace-me") warn("QQ clientSecret is missing or placeholder.");
|
||||
if (config.platforms.qq.verifySignature) {
|
||||
if (config.platforms.qq.connectionMode === "webhook" && config.platforms.qq.verifySignature) {
|
||||
const callbackSecret = config.platforms.qq.botSecret || config.platforms.qq.clientSecret;
|
||||
if (!callbackSecret || callbackSecret === "replace-me") problems += error("QQ botSecret or clientSecret is required when verifySignature is true.");
|
||||
}
|
||||
|
||||
+38
-9
@@ -202,15 +202,39 @@ async function buildNextConfig(
|
||||
secret: await prompt.ask(existing.platforms.wecom.secret ? "WeCom secret (leave blank to keep existing)" : "WeCom secret") || existing.platforms.wecom.secret
|
||||
};
|
||||
} else if (selectedPlatform === "qq") {
|
||||
console.log("\nQQ Bot HTTP callback endpoint: /webhook/qq");
|
||||
platforms.qq = {
|
||||
enabled: true,
|
||||
appId: await prompt.ask("QQ appId", existing.platforms.qq.appId),
|
||||
clientSecret: await prompt.ask(existing.platforms.qq.clientSecret ? "QQ clientSecret (leave blank to keep existing)" : "QQ clientSecret") || existing.platforms.qq.clientSecret,
|
||||
botSecret: await prompt.ask(existing.platforms.qq.botSecret ? "QQ botSecret for callback signing (leave blank to keep existing)" : "QQ botSecret for callback signing, blank to reuse clientSecret", "") || existing.platforms.qq.botSecret,
|
||||
verifySignature: await prompt.askBoolean("Verify QQ callback signatures", existing.platforms.qq.verifySignature),
|
||||
botNames: splitCommaList(await prompt.ask("QQ bot names, comma separated", existing.platforms.qq.botNames.join(", ")))
|
||||
};
|
||||
const connectionMode = await prompt.choose("QQ connection mode", [
|
||||
{ label: "WebSocket gateway", value: "websocket" as const, hint: "recommended; no public callback URL needed" },
|
||||
{ label: "HTTP callback webhook", value: "webhook" as const, hint: "requires public HTTPS callback URL" }
|
||||
], existing.platforms.qq.connectionMode === "webhook" ? 1 : 0);
|
||||
|
||||
if (connectionMode === "websocket") {
|
||||
console.log("\nQQ WebSocket gateway will actively connect to QQ; no public domain is needed.");
|
||||
const intentsText = await prompt.ask("QQ gateway intents", String(existing.platforms.qq.intents));
|
||||
platforms.qq = {
|
||||
enabled: true,
|
||||
connectionMode,
|
||||
appId: await prompt.ask("QQ appId", existing.platforms.qq.appId),
|
||||
clientSecret: await prompt.ask(existing.platforms.qq.clientSecret ? "QQ clientSecret (leave blank to keep existing)" : "QQ clientSecret") || existing.platforms.qq.clientSecret,
|
||||
botSecret: existing.platforms.qq.botSecret,
|
||||
verifySignature: existing.platforms.qq.verifySignature,
|
||||
botNames: splitCommaList(await prompt.ask("QQ bot names, comma separated", existing.platforms.qq.botNames.join(", "))),
|
||||
intents: parsePositiveInt(intentsText, existing.platforms.qq.intents),
|
||||
shard: existing.platforms.qq.shard
|
||||
};
|
||||
} else {
|
||||
console.log("\nQQ Bot HTTP callback endpoint: /webhook/qq");
|
||||
platforms.qq = {
|
||||
enabled: true,
|
||||
connectionMode,
|
||||
appId: await prompt.ask("QQ appId", existing.platforms.qq.appId),
|
||||
clientSecret: await prompt.ask(existing.platforms.qq.clientSecret ? "QQ clientSecret (leave blank to keep existing)" : "QQ clientSecret") || existing.platforms.qq.clientSecret,
|
||||
botSecret: await prompt.ask(existing.platforms.qq.botSecret ? "QQ botSecret for callback signing (leave blank to keep existing)" : "QQ botSecret for callback signing, blank to reuse clientSecret", "") || existing.platforms.qq.botSecret,
|
||||
verifySignature: await prompt.askBoolean("Verify QQ callback signatures", existing.platforms.qq.verifySignature),
|
||||
botNames: splitCommaList(await prompt.ask("QQ bot names, comma separated", existing.platforms.qq.botNames.join(", "))),
|
||||
intents: existing.platforms.qq.intents,
|
||||
shard: existing.platforms.qq.shard
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
const echo = echoConfig();
|
||||
@@ -263,6 +287,11 @@ function splitArgs(value: string): string[] {
|
||||
return value.split(" ").map((part) => part.trim()).filter(Boolean);
|
||||
}
|
||||
|
||||
function parsePositiveInt(value: string, fallback: number): number {
|
||||
const parsed = Number.parseInt(value, 10);
|
||||
return Number.isInteger(parsed) && parsed > 0 ? parsed : fallback;
|
||||
}
|
||||
|
||||
function splitCommaList(value: string): string[] {
|
||||
return value.split(",").map((part) => part.trim()).filter(Boolean);
|
||||
}
|
||||
|
||||
+4
-1
@@ -36,11 +36,14 @@ const wecomSchema = z.object({
|
||||
|
||||
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([])
|
||||
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({
|
||||
|
||||
@@ -38,13 +38,17 @@ export class QqAdapter implements PlatformAdapter {
|
||||
return jsonResponse(QQ_CALLBACK_ACK);
|
||||
}
|
||||
|
||||
const message = this.normalizeMessage(payload);
|
||||
if (!message) return jsonResponse(QQ_CALLBACK_ACK);
|
||||
if (!this.handleDispatch(payload)) return jsonResponse(QQ_CALLBACK_ACK);
|
||||
return jsonResponse(QQ_CALLBACK_ACK);
|
||||
}
|
||||
|
||||
handleDispatch(payload: QqWebhookPayload): boolean {
|
||||
const message = this.normalizeMessage(payload);
|
||||
if (!message) return false;
|
||||
void this.gateway.receive(message, this).catch((error) => {
|
||||
console.error("QQ gateway error", error);
|
||||
});
|
||||
return jsonResponse(QQ_CALLBACK_ACK);
|
||||
return true;
|
||||
}
|
||||
|
||||
async sendMessage(message: OutgoingMessage): Promise<void> {
|
||||
@@ -110,7 +114,7 @@ export class QqAdapter implements PlatformAdapter {
|
||||
return this.config.botSecret || this.config.clientSecret;
|
||||
}
|
||||
|
||||
private async getAccessToken(): Promise<string> {
|
||||
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", {
|
||||
|
||||
@@ -0,0 +1,188 @@
|
||||
import os from "node:os";
|
||||
import WebSocket from "ws";
|
||||
import type { AppConfig } from "../../config.js";
|
||||
import type { QqAdapter } from "./adapter.js";
|
||||
import type { QqGatewayResponse, QqWebhookPayload } from "./types.js";
|
||||
|
||||
const QQ_GATEWAY_API = "https://api.sgroup.qq.com/gateway/bot";
|
||||
const RECONNECT_BASE_MS = 2_000;
|
||||
const RECONNECT_MAX_MS = 60_000;
|
||||
|
||||
export class QqGatewayClient {
|
||||
private socket?: WebSocket;
|
||||
private heartbeat?: NodeJS.Timeout;
|
||||
private reconnectTimer?: NodeJS.Timeout;
|
||||
private seq: number | null = null;
|
||||
private sessionId?: string;
|
||||
private token = "";
|
||||
private reconnectAttempts = 0;
|
||||
private stopped = false;
|
||||
|
||||
constructor(
|
||||
private readonly config: AppConfig["platforms"]["qq"],
|
||||
private readonly adapter: QqAdapter
|
||||
) {}
|
||||
|
||||
start(): void {
|
||||
this.stopped = false;
|
||||
void this.connect().catch((error) => {
|
||||
console.error("QQ websocket connect failed", error);
|
||||
this.scheduleReconnect();
|
||||
});
|
||||
}
|
||||
|
||||
stop(): void {
|
||||
this.stopped = true;
|
||||
if (this.heartbeat) clearInterval(this.heartbeat);
|
||||
if (this.reconnectTimer) clearTimeout(this.reconnectTimer);
|
||||
this.heartbeat = undefined;
|
||||
this.reconnectTimer = undefined;
|
||||
this.socket?.close();
|
||||
this.socket = undefined;
|
||||
}
|
||||
|
||||
private async connect(): Promise<void> {
|
||||
if (this.stopped) return;
|
||||
|
||||
const token = await this.adapter.getAccessToken();
|
||||
this.token = token;
|
||||
const gateway = await this.fetchGateway(token);
|
||||
if (!gateway.url) throw new Error("QQ gateway response did not include url");
|
||||
|
||||
const socket = new WebSocket(gateway.url);
|
||||
this.socket = socket;
|
||||
|
||||
socket.on("open", () => {
|
||||
console.log(`QQ websocket connected: ${gateway.url}`);
|
||||
});
|
||||
|
||||
socket.on("message", (data) => {
|
||||
const payload = parsePayload(data);
|
||||
if (payload) this.handlePayload(payload);
|
||||
});
|
||||
|
||||
socket.on("close", (code, reason) => {
|
||||
if (this.stopped) return;
|
||||
this.clearHeartbeat();
|
||||
console.warn(`QQ websocket closed: ${code} ${reason.toString()}`.trim());
|
||||
this.scheduleReconnect();
|
||||
});
|
||||
|
||||
socket.on("error", (error) => {
|
||||
if (!this.stopped) console.error("QQ websocket error", error);
|
||||
});
|
||||
}
|
||||
|
||||
handlePayload(payload: QqWebhookPayload): void {
|
||||
if (typeof payload.s === "number") this.seq = payload.s;
|
||||
|
||||
if (payload.op === 10) {
|
||||
const heartbeatInterval = payload.d?.heartbeat_interval || 45_000;
|
||||
this.sendIdentifyOrResume();
|
||||
this.startHeartbeat(heartbeatInterval);
|
||||
return;
|
||||
}
|
||||
|
||||
if (payload.op === 0) {
|
||||
if (payload.t === "READY") {
|
||||
this.sessionId = payload.d?.session_id;
|
||||
this.reconnectAttempts = 0;
|
||||
console.log("QQ websocket ready");
|
||||
return;
|
||||
}
|
||||
if (payload.t === "GROUP_AT_MESSAGE_CREATE" || payload.t === "C2C_MESSAGE_CREATE") {
|
||||
this.adapter.handleDispatch(payload);
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
if (payload.op === 7) {
|
||||
console.warn("QQ websocket requested reconnect");
|
||||
this.reconnectNow();
|
||||
return;
|
||||
}
|
||||
|
||||
if (payload.op === 9) {
|
||||
this.sessionId = undefined;
|
||||
console.warn("QQ websocket invalid session; identifying again");
|
||||
this.sendIdentifyOrResume();
|
||||
}
|
||||
}
|
||||
|
||||
private async fetchGateway(token: string): Promise<QqGatewayResponse> {
|
||||
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()}`);
|
||||
return response.json() as Promise<QqGatewayResponse>;
|
||||
}
|
||||
|
||||
private sendIdentifyOrResume(): void {
|
||||
if (this.sessionId && this.seq !== null) {
|
||||
this.send({ op: 6, d: { token: this.tokenHeader(), session_id: this.sessionId, seq: this.seq } });
|
||||
return;
|
||||
}
|
||||
|
||||
this.send({
|
||||
op: 2,
|
||||
d: {
|
||||
token: this.tokenHeader(),
|
||||
intents: this.config.intents,
|
||||
shard: this.config.shard,
|
||||
properties: {
|
||||
$os: process.platform,
|
||||
$browser: "gori-agent-gateway",
|
||||
$device: os.hostname()
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
private tokenHeader(): string {
|
||||
return `QQBot ${this.token}`;
|
||||
}
|
||||
|
||||
private startHeartbeat(intervalMs: number): void {
|
||||
this.clearHeartbeat();
|
||||
this.heartbeat = setInterval(() => {
|
||||
this.send({ op: 1, d: this.seq });
|
||||
}, intervalMs);
|
||||
}
|
||||
|
||||
private clearHeartbeat(): void {
|
||||
if (this.heartbeat) clearInterval(this.heartbeat);
|
||||
this.heartbeat = undefined;
|
||||
}
|
||||
|
||||
private send(payload: Record<string, unknown>): void {
|
||||
if (!this.socket || this.socket.readyState !== WebSocket.OPEN) return;
|
||||
this.socket.send(JSON.stringify(payload));
|
||||
}
|
||||
|
||||
private reconnectNow(): void {
|
||||
if (this.socket) this.socket.close();
|
||||
else this.scheduleReconnect();
|
||||
}
|
||||
|
||||
private scheduleReconnect(): void {
|
||||
if (this.stopped || this.reconnectTimer) return;
|
||||
const delay = Math.min(RECONNECT_BASE_MS * 2 ** this.reconnectAttempts, RECONNECT_MAX_MS);
|
||||
this.reconnectAttempts += 1;
|
||||
this.reconnectTimer = setTimeout(() => {
|
||||
this.reconnectTimer = undefined;
|
||||
void this.connect().catch((error) => {
|
||||
console.error("QQ websocket reconnect failed", error);
|
||||
this.scheduleReconnect();
|
||||
});
|
||||
}, delay);
|
||||
}
|
||||
}
|
||||
|
||||
function parsePayload(data: WebSocket.RawData): QqWebhookPayload | undefined {
|
||||
try {
|
||||
return JSON.parse(data.toString("utf8")) as QqWebhookPayload;
|
||||
} catch (error) {
|
||||
console.error("QQ websocket received invalid payload", error);
|
||||
return undefined;
|
||||
}
|
||||
}
|
||||
@@ -20,9 +20,22 @@ export interface QqWebhookPayload {
|
||||
t?: string;
|
||||
}
|
||||
|
||||
export interface QqGatewayResponse {
|
||||
url?: string;
|
||||
shards?: number;
|
||||
session_start_limit?: {
|
||||
total?: number;
|
||||
remaining?: number;
|
||||
reset_after?: number;
|
||||
max_concurrency?: number;
|
||||
};
|
||||
}
|
||||
|
||||
export interface QqWebhookEventData {
|
||||
plain_token?: string;
|
||||
event_ts?: string;
|
||||
session_id?: string;
|
||||
heartbeat_interval?: number;
|
||||
id?: string;
|
||||
msg_id?: string;
|
||||
message_id?: string;
|
||||
|
||||
+26
-5
@@ -11,10 +11,16 @@ 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 { QqGatewayClient } from "./platforms/qq/gateway-client.js";
|
||||
import { GenericWebhookAdapter } from "./platforms/webhook/adapter.js";
|
||||
import { WeixinAdapter } from "./platforms/weixin/adapter.js";
|
||||
|
||||
export function createApp(config: AppConfig): express.Express {
|
||||
export interface GatewayRuntime {
|
||||
app: express.Express;
|
||||
qqGatewayClient?: QqGatewayClient;
|
||||
}
|
||||
|
||||
export function createGatewayRuntime(config: AppConfig): GatewayRuntime {
|
||||
const sessions = new SessionStore();
|
||||
const agents = new AgentRegistry(config.defaultAgent);
|
||||
for (const agentConfig of config.agents) {
|
||||
@@ -23,10 +29,11 @@ export function createApp(config: AppConfig): express.Express {
|
||||
|
||||
const gateway = new Gateway(config.policy, agents, sessions);
|
||||
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),
|
||||
new QqAdapter(config.platforms.qq, gateway),
|
||||
qqAdapter,
|
||||
new GenericWebhookAdapter(config.platforms.webhook, gateway),
|
||||
new WeixinAdapter(config.platforms.weixin, gateway)
|
||||
];
|
||||
@@ -75,14 +82,28 @@ export function createApp(config: AppConfig): express.Express {
|
||||
}
|
||||
mountWebhook("generic", "webhook");
|
||||
|
||||
return app;
|
||||
const qqGatewayClient = config.platforms.qq.enabled && config.platforms.qq.connectionMode === "websocket"
|
||||
? new QqGatewayClient(config.platforms.qq, qqAdapter)
|
||||
: undefined;
|
||||
|
||||
return { app, qqGatewayClient };
|
||||
}
|
||||
|
||||
export function createApp(config: AppConfig): express.Express {
|
||||
return createGatewayRuntime(config).app;
|
||||
}
|
||||
|
||||
export function startServer(config: AppConfig): Server {
|
||||
const app = createApp(config);
|
||||
return app.listen(config.server.port, config.server.host, () => {
|
||||
const runtime = createGatewayRuntime(config);
|
||||
const server = runtime.app.listen(config.server.port, config.server.host, () => {
|
||||
console.log(`gori-agent-gateway listening on ${config.server.host}:${config.server.port}`);
|
||||
if (runtime.qqGatewayClient) {
|
||||
console.log("QQ websocket gateway enabled; connecting to QQ...");
|
||||
runtime.qqGatewayClient.start();
|
||||
}
|
||||
});
|
||||
server.on("close", () => runtime.qqGatewayClient?.stop());
|
||||
return server;
|
||||
}
|
||||
|
||||
if (process.argv[1] && fileURLToPath(import.meta.url) === process.argv[1]) {
|
||||
|
||||
Reference in New Issue
Block a user