'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();