'use strict';
const IM_TYPE = process.env.AUTCLAW_ADAPTER_IM_TYPE || 'demo';
const CONTROL_URL = process.env.AUTCLAW_ADAPTER_CONNECT_URL
|| `ws://127.0.0.1:${process.env.AUTCLAW_PORT || '8200'}/api/adapters/${IM_TYPE}/connect`;
const CALL_TIMEOUT_MS = 15000;
// ---- 账号快照:每次 adapter_init 重建 ----
let accounts = new Map(); // account_id -> {instanceID, token, enabled}
function applyInit(params) {
accounts = new Map();
for (const config of params.account_configs || []) {
const values = config.schema_values || {};
accounts.set(String(config.account_id), {
instanceID: String(config.instance_id || ''), // 标准事件必须回填
token: String(values.bot_token || ''),
enabled: config.enabled !== false,
});
}
}
// ---- 每次连接一个 session:pending 调用与平台循环都属于当前连接 ----
let generation = 0;
function connect() {
const session = {
generation: ++generation,
ws: new WebSocket(CONTROL_URL),
pending: new Map(), // aut_echo -> {resolve, reject, timer}
closed: false,
seq: 0,
};
// autClaw 会周期性发 ws ping,标准 WebSocket 客户端自动回 pong,无需自行实现心跳
session.ws.addEventListener('message', (event) => {
let packet;
try { packet = JSON.parse(String(event.data)); } catch { return; }
onPacket(session, packet).catch((err) => console.error('[demo] handle failed:', err));
});
const reopen = () => {
if (session.closed) return;
session.closed = true;
for (const entry of session.pending.values()) { // 断线:所有 pending 立即失败
clearTimeout(entry.timer);
entry.reject(new Error('control websocket closed'));
}
session.pending.clear();
setTimeout(connect, 3000);
};
session.ws.addEventListener('close', reopen);
session.ws.addEventListener('error', reopen);
}
function send(session, payload) {
if (session.closed || session.ws.readyState !== WebSocket.OPEN) {
throw new Error('control websocket is not open');
}
session.ws.send(JSON.stringify(payload));
}
function reply(session, echo, params, error) {
const response = { aut_echo: echo, ok: !error }; // 响应不带 aut_action
if (error) response.aut_error = String(error.message || error);
else response.aut_params = params || {};
send(session, response);
}
// worker 主动调用 autClaw:注册 pending,按 aut_echo 匹配响应,必须设超时
function call(session, action, params, timeoutMS = CALL_TIMEOUT_MS) {
return new Promise((resolve, reject) => {
const echo = `js_worker_${session.generation}_${++session.seq}`;
const timer = setTimeout(() => {
session.pending.delete(echo);
reject(new Error(`${action} timed out`));
}, timeoutMS);
session.pending.set(echo, { resolve, reject, timer });
try {
send(session, { aut_action: action, aut_echo: echo, aut_params: params || {} });
} catch (err) {
clearTimeout(timer);
session.pending.delete(echo);
reject(err);
}
});
}
async function onPacket(session, packet) {
const echo = typeof packet.aut_echo === 'string' ? packet.aut_echo : '';
// 没有 aut_action 且带 aut_echo 的是响应:完成 pending 调用
if (!packet.aut_action && echo) {
const entry = session.pending.get(echo);
if (!entry) return;
session.pending.delete(echo);
clearTimeout(entry.timer);
if (packet.ok === false) entry.reject(new Error(packet.aut_error || 'control call failed'));
else entry.resolve(packet.aut_params || {});
return;
}
const action = String(packet.aut_action || '');
if (action === 'adapter_init') {
applyInit(packet.aut_params || {});
send(session, {
aut_action: 'adapter_worker_ready',
aut_params: {
contract: 'v3', // 必须;缺失时 autClaw 会直接关闭连接
resolve_account: true,
normalize_inbound: true,
build_outbound: true,
decode_receipt: true,
execute_outbound: true,
emit_inbound_events: true,
},
});
startPolling(session); // ready 之后启动绑定当前 session 的平台循环
return;
}
if (!echo) return;
try {
reply(session, echo, await handleAction(session, action, packet.aut_params || {}));
} catch (err) {
reply(session, echo, null, err); // 异常必须转为失败响应,不能让连接崩溃
}
}
async function handleAction(session, action, params) {
switch (action) {
case 'adapter_codec_resolve_account': {
// 共享路由上报进来时从 headers / query / raw_payload 解析账号;示例取第一个账号
const first = accounts.keys().next();
if (first.done) throw new Error('no account configured');
return { account_id: first.value };
}
case 'adapter_codec_normalize_inbound':
// 共享路由的 raw_payload -> 标准事件数组;本示例事件走主动上报,这里返回空
return { events: [] };
case 'adapter_codec_build_outbound': {
const account = accounts.get(String(params.account_id));
if (!account) throw new Error(`unknown account: ${params.account_id}`);
if (params.action !== 'sendText') throw new Error(`unsupported action: ${params.action}`);
const payload = params.payload || {};
return {
executor: 'worker', // 由本 worker 通过 adapter_execute_outbound 发送
await_receipt: false,
raw_payload: {
chat_id: params.session_type === 'group' ? params.chat_id : params.user_id,
text: String(payload.text || payload.content || payload.message || ''),
},
};
}
case 'adapter_codec_decode_receipt':
return { matched: false }; // 自发送模式没有发送连接回执
case 'adapter_execute_outbound': {
const account = accounts.get(String(params.account_id));
if (!account) throw new Error(`unknown account: ${params.account_id}`);
const response = await fetch('https://api.demo.example/send', {
method: 'POST',
headers: { 'content-type': 'application/json', authorization: `Bearer ${account.token}` },
body: JSON.stringify(params.raw_payload || {}),
signal: AbortSignal.timeout(10000),
});
const body = await response.json().catch(() => ({}));
if (!response.ok) return { ok: false, error: `demo api status ${response.status}` };
return { ok: true, data: { message_id: String(body.message_id || '') } };
}
default:
throw new Error(`unsupported action: ${action}`);
}
}
// ---- 平台长轮询:随连接断开结束,重连后随新 adapter_init 重新启动 ----
function startPolling(session) {
for (const [accountID, account] of accounts) {
if (!account.enabled || !account.instanceID) continue;
pollAccount(session, accountID, account).catch((err) => console.error(`[demo] poll ${accountID}:`, err));
}
}
async function pollAccount(session, accountID, account) {
let cursor = '';
while (!session.closed) {
let updates = [];
try {
const response = await fetch(`https://api.demo.example/poll?cursor=${encodeURIComponent(cursor)}`, {
headers: { authorization: `Bearer ${account.token}` },
signal: AbortSignal.timeout(30000),
});
updates = (await response.json()).events || [];
} catch {
await sleep(3000);
continue;
}
if (!updates.length) continue;
const events = updates.map((update) => ({
kind: 'message',
event_type: 'message',
event_id: String(update.id),
dedupe_key: `update:${update.id}`,
im_type: IM_TYPE,
instance_id: account.instanceID,
self_id: accountID,
session: update.chat_id ? { type: 'group', chat_id: String(update.chat_id) } : { type: 'private' },
actor: { id: String(update.user_id), name: String(update.user_name || '') },
message: {
message_id: String(update.id),
chain: { items: [{ type: 'text', text: String(update.text || '') }] },
},
raw: { account_id: accountID, raw_message: String(update.text || '') },
}));
try {
const result = await call(session, 'adapter_emit_inbound_events', { im_type: IM_TYPE, events });
const advance = advanceCount(result, events.length);
if (advance > 0) cursor = String(updates[advance - 1].id);
if (advance < events.length) await sleep(1000); // rejected 缺口:游标停在缺口前,下轮重发
} catch {
await sleep(1000); // 调用失败:整批未确认,不推进游标
}
}
}
// 平台游标只允许推进到第一个 rejected 之前
function advanceCount(result, total) {
const rejected = new Set(
(result.outcomes || []).filter((o) => o.status === 'rejected').map((o) => o.index),
);
for (let index = 0; index < total; index++) {
if (rejected.has(index)) return index;
}
return total;
}
const sleep = (ms) => new Promise((resolve) => setTimeout(resolve, ms));
connect();