File size: 11,687 Bytes
bf21785 7323d61 bf21785 4360e71 7323d61 bf21785 dfd5fc7 bf21785 abec589 bf21785 7323d61 bf21785 4360e71 e720a94 bf21785 7323d61 bf21785 dfd5fc7 7323d61 bf21785 4360e71 bf21785 e720a94 abec589 bf21785 abec589 bf21785 abec589 bf21785 abec589 fce80d1 abec589 bf21785 e720a94 4360e71 bf21785 e720a94 bf21785 dfd5fc7 bf21785 dfd5fc7 bf21785 abec589 bf21785 abec589 0040641 abec589 0040641 abec589 0040641 abec589 bf21785 4360e71 bf21785 e720a94 7323d61 bf21785 7323d61 bf21785 abec589 e720a94 bf21785 e720a94 bf21785 abec589 bf21785 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 | // WebSocket control channel β the device dials OUT to agent.mona.expert.
// The WEBSITE is the controller: it sends commands down; the device executes
// them and streams metrics, steps, tokens, and results back up.
// There is NO local server and NO local UI served over HTTP.
import { EventEmitter } from 'node:events';
import { WebSocket } from 'ws';
import os from 'node:os';
import { statfsSync } from 'node:fs';
import { CLOUD, DEFAULTS } from './config.js';
import { log } from './log.js';
/**
* Close codes the cloud uses to say "this credential is no longer valid".
* Reconnecting would just fail again β the daemon stops and asks for re-login.
* 4001 β unauthorized (bad / expired API key)
* 4003 β forbidden (device revoked, agent disabled)
*/
const TERMINAL_CLOSE_CODES = new Set([4001, 4003]);
/** Sampled CPU busy ratio β two os.cpus() readings 100ms apart. */
async function cpuPercent() {
const a = os.cpus();
await new Promise((r) => setTimeout(r, 100));
const b = os.cpus();
let idle = 0, total = 0;
for (let i = 0; i < a.length; i++) {
const ta = a[i].times, tb = b[i].times;
idle += tb.idle - ta.idle;
for (const k of Object.keys(tb)) total += tb[k] - ta[k];
}
return total > 0 ? Math.round((1 - idle / total) * 1000) / 10 : 0;
}
/** Used disk % on the device's home volume. */
function diskPercent() {
try {
const s = statfsSync(os.homedir());
const total = s.blocks * s.bsize, free = s.bavail * s.bsize;
return total > 0 ? Math.round((1 - free / total) * 1000) / 10 : 0;
} catch { return null; }
}
export class ControlChannel extends EventEmitter {
#apiKey;
#agentId;
#capabilities;
#ws = null;
#queue = [];
#llmPending = new Map();
#metricsTimer = null;
#metricsIntervalMs;
#backoff = DEFAULTS.reconnectMinMs;
#reconnectTimer = null;
#closing = false;
#stopped = false;
#wsSkipped = false;
constructor(apiKey, agentId, capabilities = null, { metricsIntervalMs } = {}) {
super();
this.#apiKey = apiKey;
this.#agentId = agentId;
this.#capabilities = capabilities;
this.#metricsIntervalMs = metricsIntervalMs || DEFAULTS.metricsIntervalMs;
}
/** Connect (or reconnect) to the cloud. Returns this for chaining. */
connect() {
if (this.#closing || this.#stopped) return this;
// Device metrics stream over HTTPS, independent of the WS link.
// This keeps the dashboard live even when the WS relay is down
// (shared hosting has no Node.js).
this.#startMetrics();
let url = CLOUD.wsUrl;
// Docker platform registers agents by agentId in the query string.
if (CLOUD.platform === 'docker') {
url += (url.includes('?') ? '&' : '?') + `agentId=${encodeURIComponent(this.#agentId || 'agent-1')}`;
}
log.debug(`Connecting to ${url}`);
this.#ws = new WebSocket(url, {
headers: {
'authorization': `Bearer ${this.#apiKey}`,
'x-mona-agent-id': this.#agentId || '',
'user-agent': `mona-agent/${DEFAULTS.version}`,
},
});
this.#ws.on('open', () => {
this.#backoff = DEFAULTS.reconnectMinMs;
log.info(`Connected to ${new URL(url).host}`);
if (CLOUD.platform === 'docker') {
// Docker platform protocol: flat register message, no hello handshake.
this.#sendFlat('register', { name: os.hostname(), model: `mona-agent/${DEFAULTS.version}` });
} else {
this.#send('hello', {
agentId: this.#agentId,
host: os.hostname(),
platform: os.platform(),
arch: os.arch(),
cpus: os.cpus().length,
mem: os.totalmem(),
version: DEFAULTS.version,
capabilities: this.#capabilities,
});
}
this.#flush();
this.emit('connected');
});
this.#ws.on('message', (raw) => {
let msg;
try { msg = JSON.parse(raw.toString()); } catch { return; }
if (msg.type === 'llm:response' || msg.type === 'llm:error') {
this.#resolveLlm(msg);
} else if (msg.type === 'command') {
this.emit('command', msg);
} else if (CLOUD.platform === 'docker' && msg.type === 'chat') {
// Docker dashboard chat same run flow as a command.
this.emit('command', { action: 'run', runId: msg.requestId, payload: { task: msg.message } });
} else if (msg.type === 'ping') {
this.#send('pong', {});
} else {
this.emit('message', msg);
}
});
this.#ws.on('close', (code) => {
if (this.#closing) return;
// WS relay absent (HTTP fallback active): no reconnect loop.
if (this.#wsSkipped) return;
// Terminal close: the credential itself was rejected β do not loop.
if (TERMINAL_CLOSE_CODES.has(code)) {
this.#stopped = true;
clearTimeout(this.#reconnectTimer);
log.error(`Cloud rejected credentials (code ${code}) β stopping. Run: mona-agent login`);
this.emit('auth-failed', code);
this.emit('disconnected', code);
return;
}
// Exponential backoff with jitter
const jitter = Math.random() * this.#backoff * 0.3;
const wait = Math.min(this.#backoff + jitter, DEFAULTS.reconnectMaxMs);
this.#backoff = Math.min(this.#backoff * 2, DEFAULTS.reconnectMaxMs);
log.warn(`Disconnected (code=${code}), reconnecting in ${(wait / 1000).toFixed(1)}s`);
this.emit('disconnected', code);
this.#reconnectTimer = setTimeout(() => this.connect(), wait);
});
this.#ws.on('error', (err) => {
// Shared hosting has no WS relay: LiteSpeed answers the upgrade
// with a normal HTTP response. Skip WS and keep the HTTPS
// metrics fallback β the dashboard stays live.
if (CLOUD.platform === 'sngine' && /unexpected server response/i.test(err.message)) {
if (!this.#wsSkipped) {
this.#wsSkipped = true;
log.info('WS relay unavailable on this control plane β device metrics stream via HTTPS');
}
return;
}
log.error(`WebSocket error: ${err.message}`);
this.emit('error', err);
});
return this;
}
/** Send a typed message upstream. Every envelope carries a protocol version. */
#send(type, data) {
const msg = JSON.stringify({
v: 1,
type,
ts: Date.now(),
agentId: this.#agentId,
data,
});
if (this.#ws?.readyState === WebSocket.OPEN) {
this.#ws.send(msg);
} else {
this.#queue.push(msg);
}
}
/** Send a flat protocol message (docker platform: fields at top level). */
#sendFlat(type, obj) {
const msg = JSON.stringify({ type, ts: Date.now(), ...obj });
if (this.#ws?.readyState === WebSocket.OPEN) {
this.#ws.send(msg);
} else {
this.#queue.push(msg);
}
}
#flush() {
while (this.#queue.length && this.#ws?.readyState === WebSocket.OPEN) {
this.#ws.send(this.#queue.shift());
}
}
// ββ Public emitters (all consumed by the website dashboard) βββββ
step(name, detail) { if (CLOUD.platform !== 'docker') this.#send('agent.step', { name, detail }); }
token(delta, runId) { if (CLOUD.platform !== 'docker') this.#send('agent.token', { delta, runId }); }
result(runId, output) { if (CLOUD.platform !== 'docker') this.#send('agent.result', { runId, output }); }
log(level, message) { if (CLOUD.platform !== 'docker') this.#send('agent.log', { level, message }); }
// ββ Docker platform protocol ββββββββββββββββββββββββββββββββββββ
/** Reply to a dashboard chat message (docker platform). */
chatResponse(requestId, message) {
this.#sendFlat('chat:response', { requestId, message });
}
/**
* Proxy an LLM call through the docker platform (request/response RPC).
* No provider or model is named β the control plane decides those.
* @returns {Promise<{content:string, usage?:object, model?:string, finishReason?:string}>}
*/
llmRequest({ messages, temperature = 0.7 }) {
const requestId = `req_${Date.now()}_${Math.random().toString(36).slice(2, 8)}`;
return new Promise((resolve, reject) => {
const timeout = setTimeout(() => {
this.#llmPending.delete(requestId);
reject(new Error('LLM request timed out after 120s'));
}, 120_000);
this.#llmPending.set(requestId, { resolve, reject, timeout });
this.#sendFlat('llm:request', { requestId, messages, temperature });
});
}
#resolveLlm(msg) {
const p = this.#llmPending.get(msg.requestId);
if (!p) return;
clearTimeout(p.timeout);
this.#llmPending.delete(msg.requestId);
if (msg.type === 'llm:error') p.reject(new Error(msg.error));
else p.resolve({ content: msg.content, usage: msg.usage, model: msg.model, finishReason: msg.finishReason });
}
/** Current connection state. */
get connected() {
return this.#ws?.readyState === WebSocket.OPEN;
}
/** True once the cloud has rejected this credential β the daemon is done. */
get stopped() {
return this.#stopped;
}
// ββ Device metrics stream ββββββββββββββββββββββββββββββββββββββ
#startMetrics() {
if (this.#metricsTimer) return;
const tick = async () => {
const totalMem = os.totalmem(), freeMem = os.freemem();
const cpus = os.cpus();
const metrics = {
cpuLoad: os.loadavg(),
cpuPercent: await cpuPercent(),
cpuModel: cpus[0]?.model || 'unknown',
mem: {
total: totalMem,
free: freeMem,
used: totalMem - freeMem,
percent: Math.round((1 - freeMem / totalMem) * 1000) / 10,
},
diskPercent: diskPercent(),
uptime: os.uptime(),
uptimeSeconds: os.uptime(),
cpus: cpus.length,
};
if (CLOUD.platform === 'docker') {
// Docker dashboard shows live agent status from these broadcasts.
this.#sendFlat('status', { status: 'online', details: metrics });
} else {
this.#send('device.metrics', metrics);
}
// HTTP fallback: shared hosting cannot run the Node WS relay,
// so also push metrics straight to the Sngine PHP API.
this.#httpStats(metrics);
this.emit('metrics', metrics);
};
tick();
this.#metricsTimer = setInterval(tick, this.#metricsIntervalMs);
}
/** Push metrics + host info to the Sngine PHP API over HTTPS. */
#httpStats(payload) {
fetch(`${CLOUD.base}/api/v1/agent/stats`, {
method: 'POST',
headers: {
'content-type': 'application/json',
'authorization': `Bearer ${this.#apiKey}`,
},
body: JSON.stringify({
agentId: this.#agentId,
host: os.hostname(),
platform: os.platform(),
arch: os.arch(),
version: DEFAULTS.version,
...payload,
}),
}).catch(() => {});
}
#stopMetrics() {
if (this.#metricsTimer) {
clearInterval(this.#metricsTimer);
this.#metricsTimer = null;
}
}
// ββ Lifecycle βββββββββββββββββββββββββββββββββββββββββββββββββββ
close() {
this.#closing = true;
clearTimeout(this.#reconnectTimer);
this.#stopMetrics();
for (const [, p] of this.#llmPending) { clearTimeout(p.timeout); p.reject(new Error('closed')); }
this.#llmPending.clear();
if (this.#ws) {
this.#ws.removeAllListeners('close');
this.#ws.close(1000, 'agent shutdown');
}
}
}
|