import { Injectable, Logger } from '@nestjs/common'; import { Device } from '@prisma/client'; import { plainToInstance } from 'class-transformer'; import { validate } from 'class-validator'; import { PrismaService } from '../prisma/prisma.service'; import { MqttJobMessageDto } from './dto/mqtt-messages.dto'; import { MqttClientService } from './mqtt-client.service'; @Injectable() export class JobsIngestService { private readonly logger = new Logger(JobsIngestService.name); constructor( private readonly prisma: PrismaService, private readonly mqttClient: MqttClientService, ) {} async handle(device: Device, rawPayload: string) { let msg: MqttJobMessageDto; try { msg = plainToInstance(MqttJobMessageDto, JSON.parse(rawPayload) as object); } catch { this.ack(device, { ticket: null, status: 'error', reason: 'invalid JSON' }); return; } const errors = await validate(msg, { whitelist: true }); if (errors.length > 0) { this.ack(device, { ticket: msg.ticket ?? null, status: 'error', reason: 'validation failed' }); this.logger.warn(`Invalid job message from ${device.mqttUsername}: ${errors}`); return; } const existing = await this.prisma.job.findUnique({ where: { orgId_ticketNumber: { orgId: device.orgId, ticketNumber: msg.ticket } }, }); if (existing) { this.ack(device, { ticket: msg.ticket, jobId: existing.id, status: 'exists' }); return; } const job = await this.prisma.job.create({ data: { orgId: device.orgId, ticketNumber: msg.ticket, title: msg.title || `Ticket ${msg.ticket}`, description: msg.description, address: msg.address, source: 'DEVICE', createdByDeviceId: device.id, }, }); this.ack(device, { ticket: msg.ticket, jobId: job.id, status: 'created' }); this.logger.log(`Device ${device.mqttUsername} created job ${job.ticketNumber}`); } private ack(device: Device, payload: object) { this.mqttClient.publish(`devices/${device.mqttUsername}/jobs/ack`, payload); } }