/** * Digital twin demo server. * * Owns the telemetry source, broadcasts frames over WebSocket at 2 Hz, exposes a * small control API for the what-if panel, and proxies the copilot so the API key * never reaches the browser. */ import fs from 'node:fs'; import path from 'node:path'; import { fileURLToPath } from 'node:url'; import dotenv from 'dotenv'; import http from 'node:http'; import express from 'express'; import { WebSocketServer } from 'ws'; // Load .env from the repo root explicitly. `dotenv/config` resolves relative to // process.cwd(), and `npm run dev -w server` sets cwd to server/, so the root // .env was silently ignored - the symptom being a copilot that quietly falls back // to a different provider than the one you configured. const HERE = path.dirname(fileURLToPath(import.meta.url)); dotenv.config({ path: path.join(HERE, '..', '.env') }); dotenv.config({ path: path.join(HERE, '.env') }); import { STATION_SPECS, BUFFER_CAPACITY } from './sim/stations.js'; import { FAULT_PROFILES } from './sim/faults.js'; import { KPI_WINDOW_SEC } from './sim/kpi.js'; import { SimulatedSource, SPEEDS } from './ingest/simulatedSource.js'; import { MqttSource } from './ingest/mqttSource.js'; import { OpcUaSource } from './ingest/opcuaSource.js'; import { Copilot } from './ai/copilot.js'; const PORT = Number(process.env.PORT || 8787); const SOURCE_KIND = (process.env.TELEMETRY_SOURCE || 'simulated').toLowerCase(); function createSource() { switch (SOURCE_KIND) { case 'mqtt': return new MqttSource({ url: process.env.MQTT_URL, username: process.env.MQTT_USERNAME, password: process.env.MQTT_PASSWORD, }); case 'opcua': return new OpcUaSource({ endpoint: process.env.OPCUA_ENDPOINT, username: process.env.OPCUA_USERNAME, password: process.env.OPCUA_PASSWORD, }); case 'simulated': return new SimulatedSource({ seed: process.env.SIM_SEED ? Number(process.env.SIM_SEED) : undefined, }); default: throw new Error(`Unknown TELEMETRY_SOURCE "${SOURCE_KIND}". Use simulated, mqtt or opcua.`); } } const source = createSource(); const copilot = new Copilot(); const app = express(); app.use(express.json({ limit: '256kb' })); // --- static metadata the UI needs once, not on every frame ------------------- app.get('/api/meta', (_req, res) => { res.json({ lineId: 'LINE-1', stations: STATION_SPECS, faults: FAULT_PROFILES, bufferCapacity: BUFFER_CAPACITY, kpiWindowSec: KPI_WINDOW_SEC, speeds: SPEEDS, source: { kind: source.name, capabilities: source.capabilities }, copilot: copilot.status(), }); }); app.get('/api/health', (_req, res) => { res.json({ ok: true, source: source.name, hasFrame: Boolean(source.latest), simTime: source.latest ? source.latest.t : 0, copilot: copilot.status(), }); }); app.get('/api/state', (_req, res) => { if (!source.latest) return res.status(503).json({ error: 'No telemetry yet.' }); res.json(source.latest); }); // --- control surface -------------------------------------------------------- /** Wrap a control action so an unsupported source returns 400, not a 500. */ function control(handler) { return (req, res) => { try { const result = handler(req); res.json({ ok: true, result }); } catch (err) { res.status(400).json({ ok: false, error: err.message }); } }; } app.post('/api/control/speed', control((req) => source.setSpeed(req.body.speed))); app.post('/api/control/setpoint', control((req) => source.setSetpoint(req.body.value))); app.post('/api/control/line-speed', control((req) => source.setLineSpeed(req.body.value))); app.post('/api/control/tool-change', control(() => source.toolChange())); app.post('/api/control/reset', control(() => source.reset())); app.post('/api/control/fault', control((req) => { const { id, action } = req.body; if (!FAULT_PROFILES.some((f) => f.id === id)) throw new Error(`Unknown fault "${id}".`); return action === 'clear' ? source.clearFault(id) : source.injectFault(id); })); // --- copilot --------------------------------------------------------------- app.get('/api/copilot/status', async (_req, res) => { res.json(await copilot.detect()); }); // --- settings --------------------------------------------------------------- app.get('/api/settings', (_req, res) => { res.json({ settings: copilot.publicSettings(), status: copilot.status() }); }); app.post('/api/settings', async (req, res) => { const result = await copilot.configure(req.body || {}); // A rejected patch changes nothing, so report it as a client error. res.status(result.ok ? 200 : 400).json(result); }); app.post('/api/settings/reset', async (_req, res) => { res.json(await copilot.reset()); }); app.get('/api/copilot/models', async (req, res) => { const provider = String(req.query.provider || '').toLowerCase(); if (provider !== 'openrouter' && provider !== 'ollama') { return res.status(400).json({ error: 'provider must be openrouter or ollama' }); } res.json(await copilot.listModels(provider, { force: req.query.force === '1' })); }); app.post('/api/copilot/test', async (req, res) => { const { provider, model } = req.body || {}; res.json(await copilot.testProvider({ provider, model })); }); app.post('/api/copilot/chat', async (req, res) => { const { question, history, provider } = req.body || {}; res.writeHead(200, { 'Content-Type': 'text/event-stream', 'Cache-Control': 'no-cache, no-transform', Connection: 'keep-alive', 'X-Accel-Buffering': 'no', }); const send = (obj) => res.write(`data: ${JSON.stringify(obj)}\n\n`); // Detect client disconnect on the RESPONSE, not the request. Express has // already consumed the POST body by this point, which ends the request stream // and fires req 'close' immediately - watching that aborts the stream before // the first token is ever written. let closed = false; res.on('close', () => { closed = true; }); try { for await (const chunk of copilot.stream(question, history, source.latest, provider)) { if (closed) break; send(chunk); } } catch (err) { // copilot.stream() degrades internally, so reaching here means something // unexpected broke. Still return prose rather than a broken stream. console.error('[copilot] unexpected failure:', err); if (!closed) send({ type: 'token', text: `The copilot failed unexpectedly: ${err.message}` }); } if (!res.writableEnded) { send({ type: 'end' }); res.end(); } }); // --- built front end ------------------------------------------------------- /** * Serve web/dist when it exists, so `npm run build` plus this process is a * complete, single-port deployment with no dev server and no second origin. * * This is what makes the "one Node process, works with no internet" claim true. * In development Vite serves the front end on :5173 and proxies here instead, so * this block simply finds no dist directory and does nothing. */ const DIST = path.join(HERE, '..', 'web', 'dist'); if (fs.existsSync(path.join(DIST, 'index.html'))) { app.use(express.static(DIST)); // SPA fallback for anything that is not an API call. Registered as middleware // rather than app.get('*') because Express 5 rejects a bare '*' path. app.use((req, res, next) => { if (req.method !== 'GET' || req.path.startsWith('/api/')) return next(); res.sendFile(path.join(DIST, 'index.html')); }); console.log(`[server] serving built front end from ${DIST}`); } else { console.log('[server] no web/dist found - run "npm run build" for single-port mode'); } // --- websocket broadcast --------------------------------------------------- const server = http.createServer(app); const wss = new WebSocketServer({ server, path: '/ws' }); wss.on('connection', (ws) => { // Replay recent history first so the trend charts are populated on load, then // the current frame so the dashboard renders without waiting for the next tick. if (source.replay.length) { ws.send(JSON.stringify({ type: 'history', frames: source.replay })); } if (source.latest) { ws.send(JSON.stringify({ type: 'frame', frame: source.latest })); } ws.on('error', (err) => console.error('[ws] client error:', err.message)); }); source.onFrame((frame) => { const msg = JSON.stringify({ type: 'frame', frame }); for (const client of wss.clients) { if (client.readyState === 1) client.send(msg); } }); // --- startup --------------------------------------------------------------- async function main() { const status = await copilot.detect(); console.log(`[copilot] ${status.label} - ${status.detail}`); try { await source.start(); console.log(`[source] ${source.name} started`); } catch (err) { // A stub source throws a long explanatory message. Print it and stop, rather // than serving a dashboard with no data behind it. console.error(`\n[source] failed to start "${SOURCE_KIND}":\n${err.message}\n`); process.exit(1); } server.listen(PORT, () => { console.log(`[server] http://localhost:${PORT} (ws://localhost:${PORT}/ws)`); }); } for (const sig of ['SIGINT', 'SIGTERM']) { process.on(sig, async () => { await source.stop(); server.close(() => process.exit(0)); }); } main();