cc / server.js
Selinore's picture
Update server.js
68042f2 verified
Raw History Blame Contribute Delete
38.1 kB
'use strict';
require('dotenv').config();
const crypto = require('node:crypto');
const express = require('express');
const path = require('path');
const WebSocket = require('ws');
const axios = require('axios');
const {
S3Client,
PutObjectCommand,
GetObjectCommand,
ListObjectsV2Command,
DeleteObjectCommand
} = require('@aws-sdk/client-s3');
// ============================================================
// 0) Config
// ============================================================
const PORT = process.env.PORT || 7860;
const {
OPENWEBUI_BASE_URL,
OPENWEBUI_API_KEY,
OPENCLAW_WS_URL,
OPENCLAW_TOKEN,
BOT_USERNAME = 'openman',
POLL_INTERVAL = 5000,
DEBUG_LOGS = 'true',
CONNECT_PASSWORD = '',
FRONT_STREAM_INTERVAL = '1500',
MAX_UPLOAD_MB = '30',
DEFAULT_MESSAGE_LIMIT = '120',
PUBLIC_URL = '',
STORAGE_EXTERNAL_ENDPOINT,
STORAGE_PUBLIC_BUCKET,
STORAGE_REGION,
STORAGE_ACCESS_KEY_ID,
STORAGE_SECRET_ACCESS_KEY,
STORAGE_S3_FORCE_PATH_STYLE
} = process.env;
// 运行时配置(可被 /api/runtime-settings 更新并立即生效)
const runtimeConfig = {
botUsername: BOT_USERNAME,
debugLogs: String(DEBUG_LOGS).toLowerCase() !== 'false',
streamInterval: Math.max(parseInt(FRONT_STREAM_INTERVAL, 10) || 1500, 800),
defaultMessageLimit: Math.max(parseInt(DEFAULT_MESSAGE_LIMIT, 10) || 120, 20),
uiProfile: {
OPENCLAW_WS_URL: OPENCLAW_WS_URL || '',
// 注意:token / secret 不会通过 config 返回给前端,但会在运行时保存并用于连接
OPENCLAW_TOKEN: OPENCLAW_TOKEN || '',
PUBLIC_URL: PUBLIC_URL || '',
STORAGE_EXTERNAL_ENDPOINT: STORAGE_EXTERNAL_ENDPOINT || '',
STORAGE_PUBLIC_BUCKET: STORAGE_PUBLIC_BUCKET || '',
STORAGE_REGION: STORAGE_REGION || '',
STORAGE_ACCESS_KEY_ID: STORAGE_ACCESS_KEY_ID || '',
STORAGE_SECRET_ACCESS_KEY: STORAGE_SECRET_ACCESS_KEY || '',
STORAGE_S3_FORCE_PATH_STYLE: STORAGE_S3_FORCE_PATH_STYLE || 'false'
}
};
function safeRuntimeConfig() {
const ui = runtimeConfig.uiProfile || {};
const safeUi = {
OPENCLAW_WS_URL: ui.OPENCLAW_WS_URL || '',
PUBLIC_URL: ui.PUBLIC_URL || '',
STORAGE_EXTERNAL_ENDPOINT: ui.STORAGE_EXTERNAL_ENDPOINT || '',
STORAGE_PUBLIC_BUCKET: ui.STORAGE_PUBLIC_BUCKET || '',
STORAGE_REGION: ui.STORAGE_REGION || '',
STORAGE_S3_FORCE_PATH_STYLE: ui.STORAGE_S3_FORCE_PATH_STYLE || 'false',
// 不泄露敏感字段,仅提示是否已配置
OPENCLAW_TOKEN_CONFIGURED: !!(ui.OPENCLAW_TOKEN || OPENCLAW_TOKEN),
STORAGE_ACCESS_KEY_CONFIGURED: !!(ui.STORAGE_ACCESS_KEY_ID || STORAGE_ACCESS_KEY_ID),
STORAGE_SECRET_KEY_CONFIGURED: !!(ui.STORAGE_SECRET_ACCESS_KEY || STORAGE_SECRET_ACCESS_KEY)
};
return {
botUsername: runtimeConfig.botUsername,
debugLogs: runtimeConfig.debugLogs,
streamInterval: runtimeConfig.streamInterval,
defaultMessageLimit: runtimeConfig.defaultMessageLimit,
uiProfile: safeUi
};
}
const uploadMaxMB = Math.max(parseInt(MAX_UPLOAD_MB, 10) || 30, 5);
const uploadMaxBytes = uploadMaxMB * 1024 * 1024;
// ============================================================
// S3 dynamic client (runtime apply)
// ============================================================
let currentS3Config = {
region: runtimeConfig.uiProfile.STORAGE_REGION || 'us-east-1',
endpoint: runtimeConfig.uiProfile.STORAGE_EXTERNAL_ENDPOINT || '',
forcePathStyle: runtimeConfig.uiProfile.STORAGE_S3_FORCE_PATH_STYLE === 'true',
accessKeyId: runtimeConfig.uiProfile.STORAGE_ACCESS_KEY_ID || '',
secretAccessKey: runtimeConfig.uiProfile.STORAGE_SECRET_ACCESS_KEY || ''
};
function createS3Client(conf) {
const region = conf.region || 'us-east-1';
const endpoint = conf.endpoint || undefined;
const forcePathStyle = !!conf.forcePathStyle;
const accessKeyId = conf.accessKeyId || '';
const secretAccessKey = conf.secretAccessKey || '';
return new S3Client({
region,
endpoint,
forcePathStyle,
credentials:
accessKeyId && secretAccessKey
? { accessKeyId, secretAccessKey }
: undefined
});
}
let s3 = createS3Client(currentS3Config);
function rebuildS3ClientIfNeeded() {
const ui = runtimeConfig.uiProfile || {};
const next = {
region: ui.STORAGE_REGION || STORAGE_REGION || 'us-east-1',
endpoint: ui.STORAGE_EXTERNAL_ENDPOINT || STORAGE_EXTERNAL_ENDPOINT || '',
forcePathStyle: (ui.STORAGE_S3_FORCE_PATH_STYLE || STORAGE_S3_FORCE_PATH_STYLE || 'false') === 'true',
accessKeyId: ui.STORAGE_ACCESS_KEY_ID || STORAGE_ACCESS_KEY_ID || '',
secretAccessKey: ui.STORAGE_SECRET_ACCESS_KEY || STORAGE_SECRET_ACCESS_KEY || ''
};
const changed =
next.region !== currentS3Config.region ||
next.endpoint !== currentS3Config.endpoint ||
next.forcePathStyle !== currentS3Config.forcePathStyle ||
next.accessKeyId !== currentS3Config.accessKeyId ||
next.secretAccessKey !== currentS3Config.secretAccessKey;
if (changed) {
log('🔄 S3 配置变更,重建 S3Client');
currentS3Config = next;
s3 = createS3Client(currentS3Config);
}
}
const S3_PREFIX_CONFIG = 'openclaw/channels.json';
const S3_PREFIX_UPLOADS = 'openclaw/uploads/';
// ============================================================
// logs
// ============================================================
const logBuffer = [];
function addLog(type, msg) {
if (type === 'info' && !runtimeConfig.debugLogs) return;
const time = new Date().toLocaleTimeString('zh-CN', { hour12: false });
const line = `[${time}] ${type === 'err' ? '❌ ' : ''}${msg}`;
console[type === 'err' ? 'error' : 'log'](line);
logBuffer.unshift(line);
if (logBuffer.length > 300) logBuffer.pop();
}
function log(msg) { addLog('info', msg); }
function err(msg) { addLog('err', msg); }
// ============================================================
// runtime state
// ============================================================
const channelRuntime = new Map();
let ws = null;
let isClawReady = false;
let isReconnecting = false;
let pingInterval = null;
let reconnectTimer = null;
const pendingRequests = new Map();
// OpenWebUI client (固定使用 env;不做运行时改动)
const webuiApi = axios.create({
baseURL: OPENWEBUI_BASE_URL,
headers: {
'Content-Type': 'application/json',
Authorization: `Bearer ${OPENWEBUI_API_KEY}`
},
timeout: 10000
});
const fileDownloader = axios.create({
baseURL: OPENWEBUI_BASE_URL,
headers: { Authorization: `Bearer ${OPENWEBUI_API_KEY}` },
responseType: 'arraybuffer',
timeout: 30000
});
// ============================================================
// 1) S3 helpers
// ============================================================
async function streamToString(stream) {
return new Promise((resolve, reject) => {
const chunks = [];
stream.on('data', c => chunks.push(c));
stream.on('error', reject);
stream.on('end', () => resolve(Buffer.concat(chunks).toString('utf8')));
});
}
function getS3Bucket() {
return runtimeConfig.uiProfile.STORAGE_PUBLIC_BUCKET || STORAGE_PUBLIC_BUCKET || '';
}
function getPublicBaseUrl() {
const ui = runtimeConfig.uiProfile || {};
return ui.PUBLIC_URL || PUBLIC_URL || ui.STORAGE_EXTERNAL_ENDPOINT || STORAGE_EXTERNAL_ENDPOINT || '';
}
async function loadChannelsFromS3() {
const bucket = getS3Bucket();
if (!bucket) {
log('ℹ️ [S3] STORAGE_PUBLIC_BUCKET 未配置,跳过读取频道配置');
return;
}
try {
log(`☁️ [S3] 读取频道配置... bucket=${bucket}`);
const rsp = await s3.send(
new GetObjectCommand({
Bucket: bucket,
Key: S3_PREFIX_CONFIG
})
);
const str = await streamToString(rsp.Body);
const data = JSON.parse(str);
const channels = Array.isArray(data?.channels) ? data.channels : [];
// 规范化 + 去重
const normalized = [];
const seen = new Set();
for (const ch of channels) {
const id = String(ch?.id || '').trim();
if (!id || seen.has(id)) continue;
seen.add(id);
normalized.push({ id, name: String(ch?.name || id).trim() || id });
}
// 以 S3 为准“替换”本地 channelRuntime
const nextIds = new Set(normalized.map(x => x.id));
// 1) 删除本地存在但 S3 不存在的频道
for (const id of Array.from(channelRuntime.keys())) {
if (!nextIds.has(id)) channelRuntime.delete(id);
}
// 2) 新增/更新 S3 中的频道
for (const ch of normalized) {
if (!channelRuntime.has(ch.id)) {
channelRuntime.set(ch.id, { id: ch.id, name: ch.name, lastMsgId: null, syncing: false });
initChannelBaseline(ch.id);
} else {
const old = channelRuntime.get(ch.id);
old.name = ch.name || old.name || ch.id;
}
}
log(`✅ [S3] 频道配置已应用:${normalized.length} 个`);
} catch (e) {
if (e.name === 'NoSuchKey') {
log('ℹ️ [S3] 无 channels.json(将清空本地频道列表)');
channelRuntime.clear();
} else {
err(`S3读取错误: ${e.message}`);
}
}
}
async function saveChannelsToS3() {
const bucket = getS3Bucket();
if (!bucket) throw new Error('STORAGE_PUBLIC_BUCKET 未配置');
try {
const channels = Array.from(channelRuntime.values()).map(c => ({ id: c.id, name: c.name || c.id }));
const body = JSON.stringify({ version: '2.0', savedAt: new Date().toISOString(), channels }, null, 2);
await s3.send(new PutObjectCommand({
Bucket: bucket,
Key: S3_PREFIX_CONFIG,
Body: body,
ContentType: 'application/json'
}));
log(`☁️ [S3] 已保存频道配置 (${channels.length} 个)`);
} catch (e) {
err(`S3保存错误: ${e.message}`);
throw e;
}
}
async function uploadToS3(filename, buffer, mimeType) {
const bucket = getS3Bucket();
if (!bucket) {
err('S3 Bucket 未配置,无法上传');
return null;
}
const safeName = `${Date.now()}_${String(filename).replace(/[^a-zA-Z0-9._-]/g, '_')}`;
const key = S3_PREFIX_UPLOADS + safeName;
try {
await s3.send(new PutObjectCommand({
Bucket: bucket,
Key: key,
Body: buffer,
ContentType: mimeType || 'application/octet-stream',
ACL: 'public-read'
}));
let base = getPublicBaseUrl();
if (base && !base.endsWith('/')) base += '/';
return `${base}${bucket}/${key}`;
} catch (e) {
err(`S3上传失败 [${filename}]: ${e.message}`);
return null;
}
}
async function cleanOldS3Files() {
const bucket = getS3Bucket();
if (!bucket) return;
try {
const threeDaysAgo = Date.now() - 3 * 24 * 60 * 60 * 1000;
let continuationToken = null;
let deleted = 0;
do {
const result = await s3.send(new ListObjectsV2Command({
Bucket: bucket,
Prefix: S3_PREFIX_UPLOADS,
ContinuationToken: continuationToken
}));
if (result.Contents) {
for (const item of result.Contents) {
if (item.LastModified && item.LastModified.getTime() < threeDaysAgo) {
await s3.send(new DeleteObjectCommand({ Bucket: bucket, Key: item.Key }));
deleted++;
}
}
}
continuationToken = result.NextContinuationToken;
} while (continuationToken);
if (deleted > 0) log(`🧹 [S3] 清理旧文件 ${deleted} 个`);
} catch (e) {
err(`清理失败: ${e.message}`);
}
}
// ============================================================
// 2) OpenClaw request helpers
// ============================================================
function clawRequest(method, params = {}, timeout = 5000) {
return new Promise((resolve, reject) => {
if (!isClawReady || !ws || ws.readyState !== WebSocket.OPEN) {
return reject(new Error('OpenClaw未连接'));
}
const id = crypto.randomUUID();
const timer = setTimeout(() => {
pendingRequests.delete(id);
reject(new Error(`OpenClaw请求超时: ${method}`));
}, timeout);
pendingRequests.set(id, { resolve, reject, timer, method });
try {
ws.send(JSON.stringify({ type: 'req', id, method, params }));
} catch (e) {
clearTimeout(timer);
pendingRequests.delete(id);
reject(e);
}
});
}
function parseClawContent(content) {
const result = { text: '', toolCalls: [], toolResults: [], files: [] };
if (typeof content === 'string') { result.text = content; return result; }
if (Array.isArray(content)) {
const lines = [];
for (const item of content) {
if (!item) continue;
if (typeof item === 'string') { lines.push(item); continue; }
const type = item.type || '';
if (type === 'text' && item.text) lines.push(item.text);
else if (type === 'image' && item.url) { lines.push(`![图片](${item.url})`); result.files.push({ url: item.url, type: 'image' }); }
else if (type === 'audio' && item.url) { lines.push(`[语音](${item.url}#audio)`); result.files.push({ url: item.url, type: 'audio' }); }
else if (type === 'video' && item.url) { lines.push(`[视频](${item.url})`); result.files.push({ url: item.url, type: 'video' }); }
else if (type === 'toolCall') result.toolCalls.push(item);
else if (type === 'toolResult') result.toolResults.push(item);
}
result.text = lines.join('\n').trim();
return result;
}
if (content && typeof content === 'object') {
if (typeof content.text === 'string') result.text = content.text;
if (Array.isArray(content.parts)) {
const p = parseClawContent(content.parts);
result.text = p.text;
result.toolCalls = p.toolCalls;
result.toolResults = p.toolResults;
result.files = p.files;
}
}
return result;
}
function stableMsgId(m, idx = 0) {
if (m.id) return String(m.id);
const tsRaw = m.timestamp ?? m.created_at ?? 0;
const ts = Number(tsRaw) || 0;
const p = parseClawContent(m.content);
const hash = crypto
.createHash('sha1')
.update(`${ts}|${m.role || ''}|${p.text || ''}|${JSON.stringify(p.toolCalls || [])}|${JSON.stringify(p.toolResults || [])}`)
.digest('hex')
.slice(0, 12);
return `claw_${ts}_${idx}_${hash}`;
}
function normalizeClawMessages(rawList = []) {
return rawList.map((m, idx) => {
const role = m.role || 'assistant';
const parsed = parseClawContent(m.content);
const ts = Number(m.timestamp ?? m.created_at ?? Date.now()) || Date.now();
return {
id: stableMsgId(m, idx),
role,
content: parsed.text || '',
timestamp: ts,
user: {
name: role === 'assistant' ? runtimeConfig.botUsername : role === 'toolResult' ? 'Tool' : 'User'
},
data: {
files: parsed.files || [],
toolCalls: parsed.toolCalls || [],
toolResults: parsed.toolResults || [],
rawRole: role
}
};
});
}
async function getClawMessagesByChannel(cid, limit = 120) {
const payload = await clawRequest(
'chat.history',
{ sessionKey: `agent:main:${cid}`, limit: Math.max(10, limit) },
7000
);
const raw = Array.isArray(payload?.messages) ? payload.messages : [];
const normalized = normalizeClawMessages(raw);
normalized.sort((a, b) => Number(b.timestamp || 0) - Number(a.timestamp || 0)); // new->old
return normalized;
}
// ============================================================
// 3) API
// ============================================================
const app = express();
app.use(express.json({ limit: `${uploadMaxMB}mb` }));
app.use(express.static(path.join(__dirname, 'public')));
function passOK(req) {
if (!CONNECT_PASSWORD) return true;
const p = req.headers['x-connect-password'] || req.query.p || '';
return String(p) === String(CONNECT_PASSWORD);
}
function apiGuard(req, res, next) {
if (req.path === '/config') return next();
if (passOK(req)) return next();
return res.status(401).json({ error: '连接密码错误或未提供' });
}
app.use('/api', apiGuard);
app.get('/api/config', (req, res) => {
res.json({
needPassword: !!CONNECT_PASSWORD,
maxUploadMB: uploadMaxMB,
runtimeConfig: safeRuntimeConfig(),
envKeys: [
'CONNECT_PASSWORD',
'OPENCLAW_WS_URL',
'OPENCLAW_TOKEN',
'BOT_USERNAME',
'FRONT_STREAM_INTERVAL',
'MAX_UPLOAD_MB',
'DEFAULT_MESSAGE_LIMIT',
'DEBUG_LOGS',
'STORAGE_EXTERNAL_ENDPOINT',
'STORAGE_PUBLIC_BUCKET',
'STORAGE_REGION',
'STORAGE_ACCESS_KEY_ID',
'STORAGE_SECRET_ACCESS_KEY',
'STORAGE_S3_FORCE_PATH_STYLE',
'PUBLIC_URL'
]
});
});
app.get('/api/runtime-settings', (req, res) => res.json(safeRuntimeConfig()));
app.post('/api/runtime-settings', async (req, res) => {
const body = req.body || {};
// 记录旧值用于判断是否需要“立即应用”
const oldWs = runtimeConfig.uiProfile.OPENCLAW_WS_URL || '';
const oldToken = runtimeConfig.uiProfile.OPENCLAW_TOKEN || '';
const oldBot = runtimeConfig.botUsername;
const oldDebug = runtimeConfig.debugLogs;
const oldS3Fingerprint = JSON.stringify({
bucket: runtimeConfig.uiProfile.STORAGE_PUBLIC_BUCKET || STORAGE_PUBLIC_BUCKET || '',
endpoint: runtimeConfig.uiProfile.STORAGE_EXTERNAL_ENDPOINT || STORAGE_EXTERNAL_ENDPOINT || '',
region: runtimeConfig.uiProfile.STORAGE_REGION || STORAGE_REGION || '',
forcePathStyle: runtimeConfig.uiProfile.STORAGE_S3_FORCE_PATH_STYLE || STORAGE_S3_FORCE_PATH_STYLE || 'false',
accessKeyId: runtimeConfig.uiProfile.STORAGE_ACCESS_KEY_ID || STORAGE_ACCESS_KEY_ID || '',
secretAccessKey: runtimeConfig.uiProfile.STORAGE_SECRET_ACCESS_KEY || STORAGE_SECRET_ACCESS_KEY || ''
});
// ---- 普通运行时项 ----
if (typeof body.botUsername === 'string' && body.botUsername.trim()) {
runtimeConfig.botUsername = body.botUsername.trim();
}
if (typeof body.debugLogs === 'boolean') runtimeConfig.debugLogs = body.debugLogs;
if (Number.isFinite(Number(body.streamInterval))) {
runtimeConfig.streamInterval = Math.max(800, parseInt(body.streamInterval, 10));
}
if (Number.isFinite(Number(body.defaultMessageLimit))) {
runtimeConfig.defaultMessageLimit = Math.max(20, parseInt(body.defaultMessageLimit, 10));
}
// ---- uiProfile 合并(OpenClaw/S3 都在这里)----
if (body.uiProfile && typeof body.uiProfile === 'object') {
runtimeConfig.uiProfile = { ...runtimeConfig.uiProfile, ...body.uiProfile };
}
const newS3Fingerprint = JSON.stringify({
bucket: runtimeConfig.uiProfile.STORAGE_PUBLIC_BUCKET || STORAGE_PUBLIC_BUCKET || '',
endpoint: runtimeConfig.uiProfile.STORAGE_EXTERNAL_ENDPOINT || STORAGE_EXTERNAL_ENDPOINT || '',
region: runtimeConfig.uiProfile.STORAGE_REGION || STORAGE_REGION || '',
forcePathStyle: runtimeConfig.uiProfile.STORAGE_S3_FORCE_PATH_STYLE || STORAGE_S3_FORCE_PATH_STYLE || 'false',
accessKeyId: runtimeConfig.uiProfile.STORAGE_ACCESS_KEY_ID || STORAGE_ACCESS_KEY_ID || '',
secretAccessKey: runtimeConfig.uiProfile.STORAGE_SECRET_ACCESS_KEY || STORAGE_SECRET_ACCESS_KEY || ''
});
// 立即应用:S3 client
const s3Changed = oldS3Fingerprint !== newS3Fingerprint;
if (s3Changed) {
rebuildS3ClientIfNeeded();
// 关键:S3 配置变了就重新拉 channels.json,刷新频道列表
await loadChannelsFromS3();
}
// botUsername / debugLogs 打日志(可选)
if (runtimeConfig.botUsername !== oldBot) log(`👤 BOT_USERNAME 更新: ${oldBot} -> ${runtimeConfig.botUsername}`);
if (runtimeConfig.debugLogs !== oldDebug) log(`🐞 DEBUG_LOGS 更新: ${oldDebug} -> ${runtimeConfig.debugLogs}`);
// 立即应用:OpenClaw 连接变更 → 重连
const newWs = runtimeConfig.uiProfile.OPENCLAW_WS_URL || '';
const newToken = runtimeConfig.uiProfile.OPENCLAW_TOKEN || '';
if (newWs !== oldWs || newToken !== oldToken) {
log('🔄 OpenClaw 配置变更,触发重连');
forceReconnectOpenClaw();
}
log('⚙️ 运行设置更新');
res.json({ success: true, runtimeConfig: safeRuntimeConfig() });
});
app.get('/api/channels', (req, res) => {
res.json(Array.from(channelRuntime.values()));
});
app.post('/api/channels', async (req, res) => {
try {
const id = String(req.body?.id || '').trim();
const name = String(req.body?.name || id).trim();
if (!id) return res.status(400).json({ error: 'id不能为空' });
if (channelRuntime.has(id)) return res.status(409).json({ error: '频道已存在' });
channelRuntime.set(id, { id, name: name || id, lastMsgId: null, syncing: false });
await saveChannelsToS3();
initChannelBaseline(id);
log(`➕ 新增频道: ${name} (${id})`);
res.json({ success: true, channel: { id, name: name || id } });
} catch (e) {
err(`新增频道失败: ${e.message}`);
res.status(500).json({ error: '新增频道失败' });
}
});
app.delete('/api/channels/:id', async (req, res) => {
try {
const id = String(req.params.id || '').trim();
if (!id) return res.status(400).json({ error: 'id不能为空' });
if (!channelRuntime.has(id)) return res.status(404).json({ error: '频道不存在' });
channelRuntime.delete(id);
await saveChannelsToS3();
log(`➖ 删除频道: ${id}`);
res.json({ success: true });
} catch (e) {
err(`删除频道失败: ${e.message}`);
res.status(500).json({ error: '删除频道失败' });
}
});
app.get('/api/logs', (req, res) => res.json(logBuffer));
app.get('/api/messages/:cid', async (req, res) => {
const cid = req.params.cid;
try {
const limit = Math.max(parseInt(req.query.limit, 10) || runtimeConfig.defaultMessageLimit, 20);
const list = await getClawMessagesByChannel(cid, limit);
res.json(list);
} catch (e) {
err(`[前端] Claw消息获取失败 ${cid}: ${e.message}`);
res.json([]);
}
});
app.get('/api/messages/:cid/stream', async (req, res) => {
const cid = req.params.cid;
const limit = Math.max(parseInt(req.query.limit, 10) || runtimeConfig.defaultMessageLimit, 20);
const interval = Math.max(parseInt(req.query.interval, 10) || runtimeConfig.streamInterval, 800);
if (!passOK(req)) return res.status(401).end('unauthorized');
res.setHeader('Content-Type', 'text/event-stream; charset=utf-8');
res.setHeader('Cache-Control', 'no-cache, no-transform');
res.setHeader('Connection', 'keep-alive');
res.flushHeaders();
let closed = false;
let lastSeen = req.query.last_id || null;
const sendEvent = (event, payload) => {
if (closed) return;
res.write(`event: ${event}\n`);
res.write(`data: ${JSON.stringify(payload)}\n\n`);
};
const heartbeat = setInterval(() => {
if (!closed) res.write(': ping\n\n');
}, 15000);
const tick = async () => {
if (closed) return;
try {
const latest = await getClawMessagesByChannel(cid, limit); // new->old
const sorted = [...latest].sort((a, b) => Number(a.timestamp || 0) - Number(b.timestamp || 0)); // old->new
if (!lastSeen) {
if (sorted.length > 0) lastSeen = sorted[sorted.length - 1].id;
sendEvent('ready', { ok: true, count: sorted.length, lastSeen });
return;
}
const idx = sorted.findIndex(m => m.id === lastSeen);
const fresh = idx === -1 ? sorted : sorted.slice(idx + 1);
if (fresh.length > 0) {
lastSeen = fresh[fresh.length - 1].id;
sendEvent('messages', { items: fresh, lastSeen });
}
} catch (e) {
sendEvent('warn', { message: e.message || 'stream tick error' });
}
};
const timer = setInterval(tick, interval);
tick();
req.on('close', () => {
closed = true;
clearInterval(timer);
clearInterval(heartbeat);
try { res.end(); } catch {}
});
});
app.post('/api/upload', async (req, res) => {
const { filename, base64, mimeType } = req.body || {};
if (!base64 || !filename) return res.status(400).json({ error: '参数缺失' });
try {
const buffer = Buffer.from(base64, 'base64');
if (buffer.length > uploadMaxBytes) {
return res.status(413).json({ error: `文件过大,超过 ${uploadMaxMB}MB` });
}
const url = await uploadToS3(filename, buffer, mimeType);
if (!url) return res.status(500).json({ error: 'S3上传失败' });
res.json({ url, mimeType: mimeType || '' });
} catch (e) {
err(`/api/upload失败: ${e.message}`);
res.status(500).json({ error: '上传异常' });
}
});
app.post('/api/send/:cid', async (req, res) => {
const cid = req.params.cid;
const { text = '' } = req.body || {};
if (!channelRuntime.has(cid)) return res.status(404).json({ error: '频道不存在' });
const clean = String(text).trim();
if (!clean) return res.status(400).json({ error: 'text不能为空' });
const isCommand = clean.startsWith('/');
triggerAgent(cid, null, 'admin', clean, isCommand);
log(`[Web发送] ${cid}: ${clean.slice(0, 120)}...`);
res.json({ success: true });
});
app.get('/health', (req, res) => {
res.json({ status: 'ok', wsConnected: isClawReady, channels: channelRuntime.size });
});
app.get('/', (req, res) => res.sendFile(path.join(__dirname, 'public', 'index.html')));
// ============================================================
// 4) OpenClaw WS + sync logic
// ============================================================
app.listen(PORT, async () => {
log(`🚀 服务启动: ${PORT}`);
await init();
});
async function init() {
try {
rebuildS3ClientIfNeeded();
await loadChannelsFromS3();
connectOpenClaw();
setInterval(cleanOldS3Files, 3600 * 1000);
cleanOldS3Files();
} catch (e) {
err(e.message);
}
}
async function initChannelBaseline(id) {
const state = channelRuntime.get(id);
if (!state) return;
try {
const res = await webuiApi.get(`/api/v1/channels/${id}/messages?limit=1`);
if (res.data && res.data.length > 0) {
state.lastMsgId = res.data[0].id;
log(`✅ [${state.name}] 锚点初始化: ${state.lastMsgId}`);
} else {
state.lastMsgId = 'EMPTY_START';
log(`✅ [${state.name}] 空频道,等待第一条消息`);
}
} catch {
setTimeout(() => initChannelBaseline(id), 5000);
}
}
const { publicKey, privateKey } = crypto.generateKeyPairSync('ed25519');
const rawPublicKey = publicKey.export({ type: 'spki', format: 'der' }).subarray(12);
const deviceId = crypto.createHash('sha256').update(rawPublicKey).digest('hex');
const toB64Url = (buf) => buf.toString('base64').replace(/\+/g, '-').replace(/\//g, '_').replace(/=+$/, '');
function cleanupConnection() {
if (pingInterval) clearInterval(pingInterval);
if (ws) {
ws.removeAllListeners();
try { ws.terminate(); } catch {}
ws = null;
}
isClawReady = false;
for (const s of channelRuntime.values()) s.syncing = false;
}
function getEffectiveWsUrl() {
return (runtimeConfig.uiProfile.OPENCLAW_WS_URL || OPENCLAW_WS_URL || '').trim();
}
function getEffectiveToken() {
return (runtimeConfig.uiProfile.OPENCLAW_TOKEN || OPENCLAW_TOKEN || '').trim();
}
function forceReconnectOpenClaw() {
cleanupConnection();
isReconnecting = false;
if (reconnectTimer) {
clearTimeout(reconnectTimer);
reconnectTimer = null;
}
connectOpenClaw(0);
}
function connectOpenClaw(retryCount = 0) {
const wsUrl = getEffectiveWsUrl();
if (!wsUrl) {
err('OPENCLAW_WS_URL 未配置,无法连接 OpenClaw');
return;
}
if (isReconnecting || ws) return;
isReconnecting = true;
log(`🔌 连接 OpenClaw (${retryCount + 1})...`);
try {
// ws endpoint
const url = new URL(wsUrl.replace(/\/$/, '') + '/ws');
// 关键:Origin header(通常需要用 https/http,而不是 wss/ws)
// wss://xxx -> https://xxx
// ws://xxx -> http://xxx
const origin = url.origin.replace(/^wss:/i, 'https:').replace(/^ws:/i, 'http:');
ws = new WebSocket(url, {
headers: {
Origin: origin
}
// 也可以用 ws 的 origin 选项(等价):
// origin
});
ws.on('open', () => {
log(`✅ WS 握手成功 (Origin=${origin})`);
isReconnecting = false;
ws.isAlive = true;
ws.on('pong', () => { ws.isAlive = true; });
pingInterval = setInterval(() => {
if (ws.isAlive === false) {
err('心跳超时,终止连接');
return ws.terminate();
}
ws.isAlive = false;
ws.ping();
}, 10000);
});
ws.on('message', handleClawMessage);
ws.on('close', (code) => {
log(`🔌 连接断开 (${code})`);
handleDisconnect(retryCount);
});
ws.on('error', (e) => err(`WS错误: ${e.message}`));
} catch (e) {
err(e.message);
handleDisconnect(retryCount);
}
}
function handleDisconnect(retryCount) {
cleanupConnection();
isReconnecting = false;
const delay = Math.min(5000 * Math.pow(1.5, retryCount), 60000);
if (reconnectTimer) clearTimeout(reconnectTimer);
reconnectTimer = setTimeout(() => connectOpenClaw(retryCount + 1), delay);
}
function handleClawMessage(data) {
let msg;
try { msg = JSON.parse(data.toString()); } catch { return; }
if (msg.event === 'connect.challenge') {
const { nonce } = msg.payload || {};
const signedAt = Date.now();
const token = getEffectiveToken();
// scopes:排序,保证签名一致性
const scopesArr = ['operator.admin', 'operator.approvals', 'operator.pairing'].slice().sort();
const scopesStr = scopesArr.join(',');
// v3 新增字段:platform + deviceFamily(必须与 connect.params 里保持一致)
const platform =
process.platform === 'win32' ? 'Win32' :
process.platform === 'darwin' ? 'Darwin' :
process.platform === 'linux' ? 'Linux' :
String(process.platform || 'unknown');
const deviceFamily = 'desktop';
// v3 payload
const payload = [
'v3',
deviceId,
'openclaw-control-ui', // clientId
'webchat', // mode
'operator', // role
scopesStr, // scopes (string)
String(signedAt),
token,
nonce,
platform, // NEW
deviceFamily // NEW
].join('|');
const signature = crypto.sign(null, Buffer.from(payload), privateKey);
ws.send(JSON.stringify({
type: 'req',
id: crypto.randomUUID(),
method: 'connect',
params: {
minProtocol: 3,
maxProtocol: 3,
client: {
id: 'openclaw-control-ui',
version: 'dev',
platform: platform,
deviceFamily: deviceFamily, // NEW: 让网关有字段可取,避免“签名绑定字段缺失”
mode: 'webchat'
},
device: {
id: deviceId,
publicKey: toB64Url(rawPublicKey),
signature: toB64Url(signature),
signedAt,
nonce // 保留:避免 DEVICE_AUTH_NONCE_REQUIRED / MISMATCH
},
auth: { token },
role: 'operator',
scopes: scopesArr
}
}));
}
if (msg.payload?.type === 'hello-ok') {
log('🎉 OpenClaw 认证成功');
isClawReady = true;
}
if (msg.type === 'res' && msg.id && pendingRequests.has(msg.id)) {
const req = pendingRequests.get(msg.id);
pendingRequests.delete(msg.id);
if (req?.timer) clearTimeout(req.timer);
try { req?.resolve?.(msg.payload); } catch {}
}
}
async function fetchClawHistory(cid, uid) {
try {
const payload = await clawRequest('chat.history', { sessionKey: `agent:main:${cid}`, limit: 5 }, 5000);
return payload?.messages || [];
} catch {
return [];
}
}
function triggerAgent(cid, uid, name, text, isCommand) {
if (!isClawReady || !ws || ws.readyState !== WebSocket.OPEN) return;
const finalMsg = isCommand ? text : `${name} 说: ${text}`;
ws.send(JSON.stringify({
type: 'req',
id: crypto.randomUUID(),
method: isCommand ? 'chat.send' : 'agent',
params: {
sessionKey: `agent:main:${cid}`,
message: finalMsg,
idempotencyKey: crypto.randomUUID(),
deliver: isCommand ? false : undefined
}
}));
log(`🚀 [发送] ${name}: ${text.slice(0, 30)}...`);
}
// ============================================================
// sync logic(保留你的同步,但把 BOT_USERNAME 改为 runtimeConfig.botUsername 生效)
// ============================================================
async function syncChannel(cid) {
const state = channelRuntime.get(cid);
if (!state || !state.lastMsgId || !isClawReady) return;
if (state.syncing) {
if (Date.now() - (state.syncStartTime || 0) > 60000) {
log(`⚠️ [${state.name}] 同步超时复位`);
state.syncing = false;
} else return;
}
state.syncing = true;
state.syncStartTime = Date.now();
try {
const { data: msgs } = await webuiApi.get(`/api/v1/channels/${cid}/messages?skip=0&limit=10`);
if (!msgs?.length) throw new Error('No messages');
const botName = (runtimeConfig.botUsername || 'openman').toLowerCase();
const webuiFingerprints = new Set(
msgs.map(m => (m.content || m.data?.content || '').trim()).filter(Boolean)
);
const activeUsers = new Map();
msgs.slice(0, 5).forEach(m => {
if (m.user_id) activeUsers.set(m.user_id, { id: m.user_id, name: m.user?.name || 'User' });
});
for (const u of activeUsers.values()) {
await new Promise(r => setTimeout(r, 200));
const history = await fetchClawHistory(cid, u.id);
const replies = history.filter(h => h.role === 'assistant').slice(-2);
for (const r of replies) {
const txt =
typeof r.content === 'string'
? r.content.trim()
: Array.isArray(r.content)
? (r.content[0]?.text || '').trim()
: '';
if (txt && !webuiFingerprints.has(txt)) {
log(`📥 [回填回复] -> ${u.name}`);
await webuiApi.post(`/api/v1/channels/${cid}/messages/post`, { content: txt, data: { files: [] }, reply_to_id: null });
webuiFingerprints.add(txt);
}
}
}
const sorted = msgs.reverse();
let newMsgs = [];
const startIdx = sorted.findIndex(m => m.id === state.lastMsgId);
if (state.lastMsgId === 'EMPTY_START') newMsgs = sorted;
else if (startIdx !== -1) newMsgs = sorted.slice(startIdx + 1);
else {
if (sorted.length > 0) {
const newestId = sorted[sorted.length - 1].id;
log(`⚠️ [${state.name}] 锚点丢失,重置 -> ${newestId}`);
state.lastMsgId = newestId;
}
newMsgs = [];
}
if (newMsgs.length > 0) {
state.lastMsgId = sorted[sorted.length - 1].id;
for (const m of newMsgs) {
const sender = m.user?.name || '';
if (sender.toLowerCase() === botName) continue;
let txt = (m.content || m.data?.content || '')
.replace(/<@U:[^>]+>/g, '')
.replace(new RegExp(`@${runtimeConfig.botUsername}`, 'gi'), '')
.trim();
const atts = [];
let rawFiles = [];
if (m.data === true) {
try {
const { data: detailData } = await webuiApi.get(`/api/v1/channels/${cid}/messages/${m.id}/data`);
if (detailData && detailData.files) rawFiles = detailData.files;
} catch (e) {
err(`详情获取失败: ${e.message}`);
}
} else if (typeof m.data === 'object' && m.data?.files) {
rawFiles = m.data.files;
}
if (rawFiles.length > 0) {
for (const fItem of rawFiles) {
try {
let buffer, filename, mimeType;
const f = fItem.file || fItem;
if (fItem.url && fItem.url.startsWith('data:')) {
const matches = fItem.url.match(/^data:([a-zA-Z0-9/+-]+);base64,(.+)$/);
if (matches) {
mimeType = matches[1];
const extMap = {
'image/jpeg': 'jpg',
'image/png': 'png',
'image/gif': 'gif',
'image/webp': 'webp',
'audio/ogg': 'ogg',
'audio/webm': 'webm'
};
const ext = extMap[mimeType] || 'bin';
buffer = Buffer.from(matches[2], 'base64');
filename = `upload_${Date.now()}_${Math.random().toString(36).slice(2)}.${ext}`;
}
} else if (f.id) {
const { data } = await fileDownloader.get(`/api/v1/files/${f.id}/content`);
buffer = data;
filename = `${Date.now()}_${(f.filename || 'file').replace(/[^a-zA-Z0-9.-]/g, '_')}`;
mimeType = f.content_type || fItem.content_type || 'application/octet-stream';
}
if (buffer && filename) {
const s3Url = await uploadToS3(filename, buffer, mimeType);
if (s3Url) {
const isImage = mimeType?.startsWith('image/');
const isAudio = mimeType?.startsWith('audio/');
atts.push(
`[${isImage ? '图片' : isAudio ? '语音' : '文件'}: ${filename}](${s3Url}${isAudio ? '#audio' : ''}) (必须下载附件,了解内容后才能回复)`
);
}
}
} catch (e) {
err(`附件处理失败: ${e.message}`);
}
}
}
if (atts.length) txt += (txt ? '\n\n' : '') + atts.join('\n');
const isCommand = txt.trim().startsWith('/');
if (txt || atts.length > 0) {
if (!txt && atts.length > 0) txt = '发送了附件。';
await new Promise(r => setTimeout(r, 500));
triggerAgent(cid, m.user_id, sender, txt, isCommand);
}
}
}
} catch (e) {
if (e.message !== 'No messages') err(`[${state.name}] 跳过: ${e.message}`);
} finally {
state.syncing = false;
}
}
async function loop() {
const delay = Math.max(parseInt(POLL_INTERVAL || 5000, 10), 2000);
if (isClawReady) {
for (const cid of channelRuntime.keys()) {
await syncChannel(cid);
await new Promise(r => setTimeout(r, 200));
}
}
setTimeout(loop, delay);
}
loop();