Files
digitaltwin/server/ingest/mqttSource.js
T
2026-08-24 15:35:32 +05:30

101 lines
4.3 KiB
JavaScript

/**
* MQTT telemetry source - DOCUMENTED STUB.
*
* This is not implemented, and it says so honestly rather than pretending. What
* it does provide is the exact shape of the work: the tag map, the frame
* assembly, and where analytics plugs in. Wiring this to a real broker is a
* day's work, not a rewrite, because everything downstream consumes frames.
*
* To implement:
* 1. npm i mqtt
* 2. Fill TAG_MAP with the customer's actual topic names.
* 3. Implement start() as marked below.
* 4. Set TELEMETRY_SOURCE=mqtt and MQTT_URL in .env
*
* Sparkplug B note: most industrial brokers publish Sparkplug B protobuf rather
* than plain JSON on flat topics. If so, add `sparkplug-payload` to decode
* NBIRTH/NDATA messages and map metric aliases instead of topic strings.
*/
import { TelemetrySource } from './source.js';
import { AnalyticsEngine } from '../analytics/alarms.js';
/**
* Maps a broker topic to a station signal.
*
* The demo model expects the signals declared in server/sim/stations.js. Any
* topic not mapped here is ignored; any signal not supplied by the broker simply
* has no data, and the UI shows it as such rather than inventing a value.
*/
export const TAG_MAP = {
// 'plant/line1/conveyor01/belt_speed': { station: 'CONV-01', signal: 'beltSpeed' },
// 'plant/line1/conveyor01/motor_current': { station: 'CONV-01', signal: 'motorAmps' },
// 'plant/line1/cnc02/vibration_rms': { station: 'CNC-02', signal: 'vibration' },
// 'plant/line1/cnc02/spindle_load': { station: 'CNC-02', signal: 'spindleLoad' },
// 'plant/line1/oven03/zone2_pv': { station: 'OVN-03', signal: 'zone2Temp' },
// 'plant/line1/oven03/zone2_sp': { station: 'OVN-03', signal: 'setpoint' },
// 'plant/line1/ins04/reject_rate': { station: 'INS-04', signal: 'rejectRate' },
// 'plant/line1/pkg05/units_per_min': { station: 'PKG-05', signal: 'unitsPerMin' },
};
export class MqttSource extends TelemetrySource {
constructor({ url, username, password, topicPrefix } = {}) {
super('mqtt');
this.url = url;
this.username = username;
this.password = password;
this.topicPrefix = topicPrefix;
this.analytics = new AnalyticsEngine();
}
/** A real broker feed is read-only: you observe the plant, you do not drive it. */
get capabilities() {
return { timeControl: false, faultInjection: false, setpointControl: false };
}
async start() {
throw new Error(
'MQTT source is not configured. This is a documented stub.\n' +
'To enable it: npm i mqtt, populate TAG_MAP in server/ingest/mqttSource.js ' +
'with your topic names, implement start(), then set TELEMETRY_SOURCE=mqtt ' +
'and MQTT_URL in .env.\n' +
'Run with TELEMETRY_SOURCE=simulated for the demo.',
);
/* Implementation outline:
*
* const mqtt = await import('mqtt');
* this.client = mqtt.connect(this.url, { username: this.username, password: this.password });
* this.client.on('connect', () => this.client.subscribe(Object.keys(TAG_MAP)));
*
* // Accumulate the latest value per tag. Industrial tags publish on change,
* // at wildly different rates, so you assemble a frame on a timer rather
* // than trying to emit one per message.
* this.client.on('message', (topic, payload) => {
* const tag = TAG_MAP[topic];
* if (!tag) return;
* this.values[`${tag.station}.${tag.signal}`] = Number(payload.toString());
* });
*
* this.timer = setInterval(() => this.assembleAndEmit(), 500);
*
* assembleAndEmit() builds the same frame shape SimulatedSource emits:
* stations with their signals, KPI rollups (see server/sim/kpi.js - the
* KpiTracker works on any counter source, not just the simulator), then
* this.analytics.update(snapshot) and this.emit(frame).
*
* Two things that bite in the real world:
* - Staleness. Track a per-tag last-seen timestamp and mark a station
* offline when its tags go quiet, exactly as the F4 dropout fault does.
* Never let a stale value render as if it were live.
* - Units. Vibration in in/s, temperature in F, and pressure in psi are all
* common. Convert at the boundary here, not downstream.
*/
}
async stop() {
if (this.timer) clearInterval(this.timer);
if (this.client) this.client.end();
}
}