Spaces:
Paused
Paused
| const { withTransaction } = require("../db/transaction"); | |
| function createTelemetryIngestService(deps) { | |
| const { | |
| pool, | |
| logger, | |
| runtimeState, | |
| telemetryRepository, | |
| assetRepository, | |
| alertEngineService, | |
| telemetryValidator, | |
| } = deps; | |
| async function handleIncomingTelemetry(topicInfo, payload) { | |
| const validation = telemetryValidator.validateAndNormalizeTelemetry(topicInfo, payload); | |
| if (!validation.valid) { | |
| runtimeState.markMqttMessageRejected(`invalid_payload:${validation.errors.join("|")}`); | |
| logger.warn("Telemetry payload rejected", { | |
| tenantCode: topicInfo.tenantCode, | |
| truckCode: topicInfo.truckCode, | |
| containerCode: topicInfo.containerCode, | |
| errors: validation.errors, | |
| }); | |
| return; | |
| } | |
| const context = await assetRepository.resolveAssetContextByCodes(pool, { | |
| tenantCode: topicInfo.tenantCode, | |
| truckCode: topicInfo.truckCode, | |
| containerCode: topicInfo.containerCode, | |
| }); | |
| if (!context) { | |
| runtimeState.markMqttMessageRejected("unknown_asset_reference"); | |
| logger.warn("Telemetry rejected due to unknown tenant/truck/container mapping", { | |
| tenantCode: topicInfo.tenantCode, | |
| truckCode: topicInfo.truckCode, | |
| containerCode: topicInfo.containerCode, | |
| }); | |
| return; | |
| } | |
| const receivedAt = new Date().toISOString(); | |
| const normalized = validation.normalized; | |
| const telemetry = { | |
| tenantId: context.tenant_id, | |
| fleetId: context.fleet_id, | |
| truckId: context.truck_id, | |
| containerId: context.container_id, | |
| tripId: context.trip_id, | |
| gatewayDeviceId: null, | |
| sensorDeviceId: null, | |
| mqttTopic: topicInfo.topic, | |
| seq: normalized.seq, | |
| sourceTs: normalized.sourceTs, | |
| receivedAt, | |
| gpsLat: normalized.gpsLat, | |
| gpsLon: normalized.gpsLon, | |
| speedKph: normalized.speedKph, | |
| temperatureC: normalized.temperatureC, | |
| humidityPct: normalized.humidityPct, | |
| pressureHpa: normalized.pressureHpa, | |
| tiltDeg: normalized.tiltDeg, | |
| shock: normalized.shock, | |
| gasRaw: normalized.gasRaw, | |
| gasAlert: normalized.gasAlert, | |
| sdOk: normalized.sdOk, | |
| gpsFix: normalized.gpsFix, | |
| uplink: normalized.uplink, | |
| rawPayload: normalized.rawPayload, | |
| }; | |
| await withTransaction(pool, async (client) => { | |
| await telemetryRepository.insertTelemetryHistory(client, telemetry); | |
| await telemetryRepository.upsertTelemetryLatest(client, telemetry); | |
| await alertEngineService.evaluateTelemetryInTransaction(client, { | |
| tenantId: context.tenant_id, | |
| fleetId: context.fleet_id, | |
| truckId: context.truck_id, | |
| containerId: context.container_id, | |
| tripId: context.trip_id, | |
| }, telemetry); | |
| }); | |
| runtimeState.markMqttMessageAccepted(); | |
| } | |
| return { | |
| handleIncomingTelemetry, | |
| }; | |
| } | |
| module.exports = { | |
| createTelemetryIngestService, | |
| }; | |