agmission/server/controllers/api_export.js

607 lines
28 KiB
JavaScript
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

'use strict';
/**
* Async Export controller — /api/v1/jobs/:jobId/export and /api/v1/exports/:exportId
*
* Flow:
* 1. POST /api/v1/jobs/:jobId/export → creates ExportJob record (status=pending),
* kicks off async generation, returns { exportId, status: 'pending' }.
* 2. GET /api/v1/exports/:exportId → poll status; when ready returns { status: 'ready', downloadUrl }.
* 3. GET /api/v1/exports/:exportId/download → streams the file, schedules cleanup.
*
* FE / integration notes:
* - For the daily 17:00 batch: POST export after previous day's jobs are confirmed sprayed,
* poll every 1030 s, then download the CSV when ready.
* - interval applies GPS point thinning; records where sprayStat changes are always preserved.
* - Use interval=0 (or omit interval) to export all points without thinning.
* - CSV has all raw trace fields + job/session header columns repeated per row for
* direct Power BI / data-warehouse import without joins.
*/
const path = require('path');
const fs = require('fs');
const { Transform, pipeline } = require('stream');
const { promisify } = require('util');
const pipelineAsync = promisify(pipeline);
const ObjectId = require('mongodb').ObjectId;
const moment = require('moment');
const { Job, App, AppFile, AppDetail } = require('../model');
const ExportJob = require('../model/export_job');
const { AppParamError, AppAuthError } = require('../helpers/app_error');
const { Errors, HttpStatus, ExportUnits, ExportJobStatus, RateUnits } = require('../helpers/constants');
const utils = require('../helpers/utils');
const env = require('../helpers/env');
const { computeAppRateApplied, flowRateFromAppRate, isPositiveNumber, inferRateUnitCode, isLikelyLiquidMaterial, resolveTargetRatePerHa } = require('../helpers/record_utils');
const EXPORT_TTL_HOURS = env.EXPORT_TTL_HOURS || 24;
const EXPORT_DEDUP_MINS = env.EXPORT_DEDUP_MINS ?? 5;
/**
* On startup: delete orphaned export files whose ExportJob is expired or missing.
* Runs fire-and-forget so it never blocks server startup.
*/
setImmediate(async () => {
try {
const pattern = /^export_[a-f0-9]+\.(csv|json)$/;
const files = await fs.promises.readdir(env.TEMP_DIR).catch(() => []);
for (const file of files) {
if (!pattern.test(file)) continue;
const id = file.replace(/^export_/, '').replace(/\.(csv|json)$/, '');
const exists = ObjectId.isValid(id) && await ExportJob.exists({ _id: id, expiresAt: { $gt: new Date() } });
if (!exists) {
fs.unlink(path.join(env.TEMP_DIR, file), () => {});
}
}
} catch { /* non-fatal */ }
});
// Re-use the same helpers from api_pub (inline to avoid a shared helper module for now)
function parseInterval(raw) {
if (raw == null || raw === '') return null;
const v = parseFloat(raw);
return isFinite(v) && v > 0 ? v : null;
}
function getLaserAlt(detail) {
return detail?.laserAlt ?? detail?.raserAlt ?? '';
}
/**
* Convert AppDetail.gpsTime to an ISO UTC timestamp.
* Supports both epoch-seconds and legacy seconds-of-day values.
*/
function toRecordTimeUtc(gpsTime, appStartDateTime) {
if (!utils.isNumber(gpsTime)) return null;
// Epoch seconds (>= year 2000-01-01 UTC) can be converted directly.
if (gpsTime >= 946684800) {
return moment.unix(gpsTime).utc().toISOString();
}
// Legacy format: seconds-of-day, anchor to app start date when available.
const base = moment.utc(appStartDateTime, [moment.ISO_8601, 'YYYYMMDDTHHmmss'], true);
if (base.isValid()) {
const dayOffset = Math.floor(gpsTime / 86400);
const secOfDay = ((gpsTime % 86400) + 86400) % 86400;
return base.clone().startOf('day').add(dayOffset, 'days').add(secOfDay, 'seconds').toISOString();
}
// Fallback for malformed app start datetime.
return moment.unix(gpsTime).utc().toISOString();
}
/** Verify job ownership — throws on mismatch. */
async function ownerJob(jobId, ownerId) {
const job = await Job.findOne({ _id: jobId, markedDelete: { $ne: true } }).lean();
if (!job) AppParamError.throw(Errors.JOB_NOT_FOUND);
if (!job.byPuid || job.byPuid.toString() !== ownerId.toString()) AppAuthError.throw();
return job;
}
// ─── Unit conversion helpers ─────────────────────────────────────────────────
// All raw AppDetail values are stored in SI/metric units.
// When units='us', these factors convert to US customary equivalents.
const CONV = {
msToMph: v => utils.roundTo(v * 2.23694, 2), // m/s → mph
msToKt: v => utils.roundTo(v * 1.94384, 2), // m/s → kt (knots, matches playback display)
mToFt: v => utils.roundTo(v * 3.28084, 2), // m → ft
cToF: v => utils.roundTo(v * 9 / 5 + 32, 1), // °C → °F
LminToGmin: v => utils.roundTo(v * 0.264172, 4), // L/min → gal/min
LhaToGac: v => utils.roundTo(v * 0.10694, 4), // L/ha → gal/ac
};
function applyConv(v, fn) {
return (v != null && v !== '') ? fn(Number(v)) : v;
}
/**
* Returns CSV column definitions for the requested unit system.
* Each entry: { key (row-object property), header (CSV column name) }.
*/
function getCsvColumns(units, includeFm = false) {
const us = units === ExportUnits.US;
const cols = [
// Job / session metadata — no unit conversion
{ key: 'jobId' }, { key: 'orderNumber' }, { key: 'jobName' },
{ key: 'clientId' }, { key: 'clientName' },
{ key: 'sessionId' }, { key: 'fileName' }, { key: 'pilotName' },
// GPS data
{ key: 'timeUtc' }, { key: 'gpsTime' }, { key: 'lat' }, { key: 'lon' },
{ key: 'utmX' }, { key: 'utmY' },
{ key: 'alt', header: us ? 'alt_ft' : 'alt_m' },
{ key: 'grSpeed', header: us ? 'groundSpeed_mph' : 'groundSpeed_ms' },
{ key: 'heading' },
{ key: 'xTrack', header: us ? 'crossTrackError_ft' : 'crossTrackError_m' },
{ key: 'lockedLine' }, { key: 'hdop' }, { key: 'satsIn' },
{ key: 'tslu' }, { key: 'calcodeFreq' },
{ key: 'sprayStat' },
// Application data
{ key: 'flowRateApplied', header: us ? 'flowRateApplied_galMin' : 'flowRateApplied_Lmin' },
{ key: 'flowRateRequired', header: us ? 'flowRateRequired_galMin' : 'flowRateRequired_Lmin' },
{ key: 'appRateRequired', header: us ? 'appRateRequired_galAc' : 'appRateRequired_Lha' },
{ key: 'appRateApplied', header: us ? 'appRateApplied_galAc' : 'appRateApplied_Lha' },
{ key: 'swathWidth', header: us ? 'swathWidth_ft' : 'swathWidth_m' },
{ key: 'boomPressure_psi' },
{ key: 'flowController' },
{ key: 'sprayOnLag_s' }, { key: 'sprayOffLag_s' }, { key: 'pulsesPerLiter' },
{ key: 'rpm' },
// MET — wind in knots (metric) or mph (US) to match playback display
{ key: 'windSpeed_kt', header: us ? 'windSpeed_mph' : 'windSpeed_kt' },
{ key: 'windDir_deg' },
{ key: 'temp_c', header: us ? 'temp_f' : 'temp_c' },
{ key: 'humidity_pct' },
];
if (includeFm) {
// Flight Master / AgDisp fields — only when fm=true requested
cols.push(
{ key: 'sprayHeight_m' },
{ key: 'driftX_m' }, { key: 'driftY_m' },
{ key: 'depositX_m' }, { key: 'depositY_m' },
{ key: 'radarAlt_m' },
{ key: 'laserAlt_m' } // DB field is raserAlt (schema typo); exposed as laserAlt_m
);
}
return cols;
}
function escapeCsv(val) {
if (val == null) return '';
const s = String(val);
if (s.includes(',') || s.includes('"') || s.includes('\n')) return `"${s.replace(/"/g, '""')}"`;
return s;
}
function recordToRow(d, sessionMeta, jobHeader, units, includeFm = false) {
const us = units === ExportUnits.US;
// sprayStat: 0=off, 1=on, 3=segment marker. Only compute rates for actual spray-on records (1)
const sprayOn = d.sprayStat === 1;
const rateUnitCode = inferRateUnitCode(sessionMeta.meta, sessionMeta.job);
const liquidMaterial = isLikelyLiquidMaterial(sessionMeta.meta, rateUnitCode);
const targetRateMetric = resolveTargetRatePerHa(sessionMeta.meta, sessionMeta.job);
const metaAppRate = sessionMeta.meta?.appRate;
// flowRate fields: raw stored values, no fallback (matches frontend)
const flowRateAppliedRaw = d.lminApp ?? null;
const flowRateRequiredRaw = d.lminReq ?? null;
// appRateRequired: matches frontend applicRate display (metric before unit conversion below)
// Priority 1: meta.appRate (converted to metric) when present
// Priority 2: per-point lhaReq when present
// Priority 3: job.appRate (converted to metric) as fallback
const appRateRequiredRaw = (utils.isNumber(metaAppRate) && metaAppRate !== 0)
? targetRateMetric
: (utils.isNumber(d.lhaReq) ? d.lhaReq : targetRateMetric);
// appRateApplied: matches frontend appRateAp — only meaningful when spraying
// Priority 1: meta.appRate (metric) when no FC or FC has no reading
// Priority 2: liquid — compute from measured flow rate (L/min → L/ha)
// Priority 3: dry/granular — lminApp stores kg/ha directly
const appRateApplied = (() => {
if (!sprayOn) return null;
const useFC = sessionMeta.meta?.useFC;
if (metaAppRate && (!useFC || !d.lminApp)) return targetRateMetric;
if (liquidMaterial) return utils.appRateFromFlowRate(d.lminApp, d.swath, d.grSpeed);
return utils.isNumber(d.lminApp) ? d.lminApp : null;
})();
const fcName = sessionMeta.meta?.fcName;
const row = {
...jobHeader,
sessionId: sessionMeta.appId,
fileName: sessionMeta.fileName,
pilotName: sessionMeta.operator ?? '',
timeUtc: toRecordTimeUtc(d.gpsTime, sessionMeta.appStartDateTime),
gpsTime: d.gpsTime ?? '',
lat: utils.isNumber(d.lat) ? utils.roundTo(d.lat, 7) : (d.lat ?? ''),
lon: utils.isNumber(d.lon) ? utils.roundTo(d.lon, 7) : (d.lon ?? ''),
utmX: utils.isNumber(d.utmX) ? utils.roundTo(d.utmX, 1) : (d.utmX ?? ''),
utmY: utils.isNumber(d.utmY) ? utils.roundTo(d.utmY, 1) : (d.utmY ?? ''),
alt: us ? applyConv(d.alt, CONV.mToFt) : (utils.isNumber(d.alt) ? utils.roundTo(d.alt, 2) : (d.alt ?? '')),
grSpeed: us ? applyConv(d.grSpeed, CONV.msToMph) : (utils.isNumber(d.grSpeed) ? utils.roundTo(d.grSpeed, 2) : (d.grSpeed ?? '')),
heading: utils.isNumber(d.head) ? utils.roundTo(d.head, 2) : (d.head ?? ''),
xTrack: us ? applyConv(d.xTrack, CONV.mToFt) : (utils.isNumber(d.xTrack) ? utils.roundTo(d.xTrack, 2) : (d.xTrack ?? '')),
lockedLine: d.llnum ?? '', hdop: utils.isNumber(d.stdHdop) ? utils.roundTo(d.stdHdop, 2) : (d.stdHdop ?? ''),
satsIn: d.satsIn ?? '',
tslu: d.tslu ?? '', calcodeFreq: d.calcodeFreq ?? '',
sprayStat: d.sprayStat ?? '',
flowRateApplied: us ? applyConv(flowRateAppliedRaw, CONV.LminToGmin) : (utils.isNumber(flowRateAppliedRaw) ? utils.roundTo(flowRateAppliedRaw, 4) : (flowRateAppliedRaw ?? '')),
flowRateRequired: us ? applyConv(flowRateRequiredRaw, CONV.LminToGmin) : (utils.isNumber(flowRateRequiredRaw) ? utils.roundTo(flowRateRequiredRaw, 4) : (flowRateRequiredRaw ?? '')),
appRateRequired: us ? applyConv(appRateRequiredRaw, CONV.LhaToGac) : (utils.isNumber(appRateRequiredRaw) ? utils.roundTo(appRateRequiredRaw, 4) : (appRateRequiredRaw ?? '')),
appRateApplied: us ? applyConv(appRateApplied, CONV.LhaToGac) : (utils.isNumber(appRateApplied) ? utils.roundTo(appRateApplied, 4) : (appRateApplied ?? '')),
swathWidth: us ? applyConv(d.swath, CONV.mToFt) : (d.swath ?? ''),
boomPressure_psi: utils.isNumber(d.psi) ? utils.roundTo(d.psi, 2) : (d.psi ?? ''),
flowController: (fcName && !/none/i.test(fcName)) ? fcName : 'No FC',
sprayOnLag_s: sessionMeta.meta?.sprOnLag ?? '',
sprayOffLag_s: sessionMeta.meta?.sprOffLag ?? '',
pulsesPerLiter: sessionMeta.meta?.pulsesPerLit ?? '',
rpm: (Array.isArray(d.rpm) && d.rpm.length) ? JSON.stringify(d.rpm) : '',
// Wind speed in knots (metric) or mph (US) — matches playback display
windSpeed_kt: us ? applyConv(d.windSpd, CONV.msToMph) : applyConv(d.windSpd, CONV.msToKt),
windDir_deg: utils.isNumber(d.windDir) ? utils.roundTo(d.windDir, 1) : (d.windDir ?? ''),
temp_c: us ? applyConv(d.temp, CONV.cToF) : (utils.isNumber(d.temp) ? utils.roundTo(d.temp, 1) : (d.temp ?? '')),
humidity_pct: utils.isNumber(d.humid) ? utils.roundTo(d.humid, 1) : (d.humid ?? '')
};
if (includeFm) {
row.sprayHeight_m = utils.isNumber(d.sprayHeight) ? utils.roundTo(d.sprayHeight, 2) : (d.sprayHeight ?? '');
row.driftX_m = utils.isNumber(d.driftX) ? utils.roundTo(d.driftX, 2) : (d.driftX ?? '');
row.driftY_m = utils.isNumber(d.driftY) ? utils.roundTo(d.driftY, 2) : (d.driftY ?? '');
row.depositX_m = utils.isNumber(d.depositX) ? utils.roundTo(d.depositX, 2) : (d.depositX ?? '');
row.depositY_m = utils.isNumber(d.depositY) ? utils.roundTo(d.depositY, 2) : (d.depositY ?? '');
row.radarAlt_m = utils.isNumber(d.radarAlt) ? utils.roundTo(d.radarAlt, 2) : (d.radarAlt ?? '');
row.laserAlt_m = getLaserAlt(d);
}
const cols = getCsvColumns(units, includeFm);
return cols.map(c => escapeCsv(row[c.key])).join(',') + '\n';
}
// ─── Async generation ─────────────────────────────────────────────────────────
async function generateExport(exportJobId) {
const exportJob = await ExportJob.findById(exportJobId);
if (!exportJob) return;
try {
exportJob.status = ExportJobStatus.PROCESSING;
await exportJob.save();
const job = await Job.findById(exportJob.jobId, 'name orderNumber client')
.select('name orderNumber client appRate appRateUnit')
.populate('client', '_id name')
.lean();
const jobHeader = {
jobId: exportJob.jobId,
orderNumber: job?.orderNumber ?? '',
jobName: job?.name ?? '',
clientId: job?.client?._id?.toString() ?? '',
clientName: job?.client?.name ?? ''
};
const apps = await App.find({ jobId: exportJob.jobId, markedDelete: { $ne: true } }).lean();
const appFiles = await AppFile.find(
{ appId: { $in: apps.map(a => a._id) }, markedDelete: { $ne: true } }
).lean();
const filesByAppId = {};
for (const f of appFiles) {
const key = f.appId.toString();
if (!filesByAppId[key]) filesByAppId[key] = [];
filesByAppId[key].push(f);
}
const interval = exportJob.interval;
const includeFm = !!exportJob.fm;
const outPath = path.join(env.TEMP_DIR, `export_${exportJobId}.${exportJob.format}`);
const writeStream = fs.createWriteStream(outPath);
const units = exportJob.units || 'metric';
if (exportJob.format === 'csv') {
// Write header row (unit-aware column names)
const cols = getCsvColumns(units, includeFm);
writeStream.write(cols.map(c => c.header || c.key).join(',') + '\n');
for (const app of apps) {
const files = filesByAppId[app._id.toString()] || [];
for (const appFile of files) {
const sessionMeta = {
appId: app._id,
fileName: app.fileName,
operator: appFile.meta?.operator,
meta: appFile.meta,
job,
appStartDateTime: app.startDateTime
};
const cursor = AppDetail.find(
{ fileId: appFile._id },
null,
{ sort: { _id: 1 }, lean: true }
).cursor();
let prevGpsTime = null;
let prevSprayStat = null;
for await (const record of cursor) {
if (interval) {
const sprayStatChanged = prevSprayStat !== null && record.sprayStat !== prevSprayStat;
if (prevGpsTime !== null && (record.gpsTime - prevGpsTime) < interval && !sprayStatChanged) continue;
prevGpsTime = record.gpsTime;
}
prevSprayStat = record.sprayStat;
writeStream.write(recordToRow(record, sessionMeta, jobHeader, units, includeFm));
}
}
}
} else if (exportJob.format === 'json') {
// JSON array of records — one object per GPS point with all fields.
const records = [];
for (const app of apps) {
const files = filesByAppId[app._id.toString()] || [];
for (const appFile of files) {
const sessionMeta = {
appId: app._id,
fileName: app.fileName,
operator: appFile.meta?.operator,
meta: appFile.meta,
job,
appStartDateTime: app.startDateTime
};
const cursor = AppDetail.find(
{ fileId: appFile._id },
null,
{ sort: { _id: 1 }, lean: true }
).cursor();
let prevGpsTime = null;
let prevSprayStat = null;
for await (const record of cursor) {
if (interval) {
const sprayStatChanged = prevSprayStat !== null && record.sprayStat !== prevSprayStat;
if (prevGpsTime !== null && (record.gpsTime - prevGpsTime) < interval && !sprayStatChanged) continue;
prevGpsTime = record.gpsTime;
}
prevSprayStat = record.sprayStat;
// Build record object directly with formatted values
const us = units === ExportUnits.US;
const sprayOn = record.sprayStat === 1 || record.sprayStat === 2;
const rateUnitCode = inferRateUnitCode(sessionMeta.meta, sessionMeta.job);
const liquidMaterial = isLikelyLiquidMaterial(sessionMeta.meta, rateUnitCode);
const targetRateMetric = resolveTargetRatePerHa(sessionMeta.meta, sessionMeta.job);
const metaAppRate = sessionMeta.meta?.appRate;
const flowRateAppliedRaw = record.lminApp ?? null;
const flowRateRequiredRaw = record.lminReq ?? null;
const appRateRequiredRaw = (utils.isNumber(metaAppRate) && metaAppRate !== 0)
? targetRateMetric
: (utils.isNumber(record.lhaReq) ? record.lhaReq : targetRateMetric);
const appRateApplied = (() => {
if (!sprayOn) return null;
const useFC = sessionMeta.meta?.useFC;
if (metaAppRate && (!useFC || !record.lminApp)) return targetRateMetric;
if (liquidMaterial) return utils.appRateFromFlowRate(record.lminApp, record.swath, record.grSpeed);
return utils.isNumber(record.lminApp) ? record.lminApp : null;
})();
const recordObj = {
...jobHeader,
sessionId: sessionMeta.appId,
fileName: sessionMeta.fileName,
pilotName: sessionMeta.operator ?? '',
timeUtc: toRecordTimeUtc(record.gpsTime, sessionMeta.appStartDateTime),
gpsTime: record.gpsTime ?? '',
lat: utils.isNumber(record.lat) ? utils.roundTo(record.lat, 7) : (record.lat ?? ''),
lon: utils.isNumber(record.lon) ? utils.roundTo(record.lon, 7) : (record.lon ?? ''),
utmX: utils.isNumber(record.utmX) ? utils.roundTo(record.utmX, 1) : (record.utmX ?? ''),
utmY: utils.isNumber(record.utmY) ? utils.roundTo(record.utmY, 1) : (record.utmY ?? ''),
alt: us ? applyConv(record.alt, CONV.mToFt) : (utils.isNumber(record.alt) ? utils.roundTo(record.alt, 2) : (record.alt ?? '')),
grSpeed: us ? applyConv(record.grSpeed, CONV.msToMph) : (utils.isNumber(record.grSpeed) ? utils.roundTo(record.grSpeed, 2) : (record.grSpeed ?? '')),
heading: utils.isNumber(record.head) ? utils.roundTo(record.head, 2) : (record.head ?? ''),
xTrack: us ? applyConv(record.xTrack, CONV.mToFt) : (utils.isNumber(record.xTrack) ? utils.roundTo(record.xTrack, 2) : (record.xTrack ?? '')),
lockedLine: record.llnum ?? '',
hdop: utils.isNumber(record.stdHdop) ? utils.roundTo(record.stdHdop, 2) : (record.stdHdop ?? ''),
satsIn: record.satsIn ?? '',
tslu: record.tslu ?? '',
calcodeFreq: record.calcodeFreq ?? '',
sprayStat: record.sprayStat ?? '',
flowRateApplied: us ? applyConv(flowRateAppliedRaw, CONV.LminToGmin) : (utils.isNumber(flowRateAppliedRaw) ? utils.roundTo(flowRateAppliedRaw, 4) : (flowRateAppliedRaw ?? '')),
flowRateRequired: us ? applyConv(flowRateRequiredRaw, CONV.LminToGmin) : (utils.isNumber(flowRateRequiredRaw) ? utils.roundTo(flowRateRequiredRaw, 4) : (flowRateRequiredRaw ?? '')),
appRateRequired: us ? applyConv(appRateRequiredRaw, CONV.LhaToGac) : (utils.isNumber(appRateRequiredRaw) ? utils.roundTo(appRateRequiredRaw, 4) : (appRateRequiredRaw ?? '')),
appRateApplied: us ? applyConv(appRateApplied, CONV.LhaToGac) : (utils.isNumber(appRateApplied) ? utils.roundTo(appRateApplied, 4) : (appRateApplied ?? '')),
swathWidth: us ? applyConv(record.swath, CONV.mToFt) : (record.swath ?? ''),
boomPressure_psi: utils.isNumber(record.psi) ? utils.roundTo(record.psi, 2) : (record.psi ?? ''),
flowController: (sessionMeta.meta?.fcName && !/none/i.test(sessionMeta.meta?.fcName)) ? sessionMeta.meta?.fcName : 'No FC',
sprayOnLag_s: sessionMeta.meta?.sprOnLag ?? '',
sprayOffLag_s: sessionMeta.meta?.sprOffLag ?? '',
pulsesPerLiter: sessionMeta.meta?.pulsesPerLit ?? '',
rpm: (Array.isArray(record.rpm) && record.rpm.length) ? record.rpm : '',
windSpeed: us ? applyConv(record.windSpd, CONV.msToMph) : applyConv(record.windSpd, CONV.msToKt),
windDir_deg: utils.isNumber(record.windDir) ? utils.roundTo(record.windDir, 1) : (record.windDir ?? ''),
temp: us ? applyConv(record.temp, CONV.cToF) : (utils.isNumber(record.temp) ? utils.roundTo(record.temp, 1) : (record.temp ?? '')),
humidity_pct: utils.isNumber(record.humid) ? utils.roundTo(record.humid, 1) : (record.humid ?? '')
};
if (includeFm) {
recordObj.sprayHeight_m = utils.isNumber(record.sprayHeight) ? utils.roundTo(record.sprayHeight, 2) : (record.sprayHeight ?? '');
recordObj.driftX_m = utils.isNumber(record.driftX) ? utils.roundTo(record.driftX, 2) : (record.driftX ?? '');
recordObj.driftY_m = utils.isNumber(record.driftY) ? utils.roundTo(record.driftY, 2) : (record.driftY ?? '');
recordObj.depositX_m = utils.isNumber(record.depositX) ? utils.roundTo(record.depositX, 2) : (record.depositX ?? '');
recordObj.depositY_m = utils.isNumber(record.depositY) ? utils.roundTo(record.depositY, 2) : (record.depositY ?? '');
recordObj.radarAlt_m = utils.isNumber(record.radarAlt) ? utils.roundTo(record.radarAlt, 2) : (record.radarAlt ?? '');
recordObj.laserAlt_m = getLaserAlt(record);
}
records.push(recordObj);
}
}
}
writeStream.write(JSON.stringify(records, null, 2));
}
await new Promise((resolve, reject) => {
writeStream.end();
writeStream.on('finish', resolve);
writeStream.on('error', reject);
});
const expiresAt = new Date(Date.now() + EXPORT_TTL_HOURS * 3600 * 1000);
exportJob.status = ExportJobStatus.READY;
exportJob.filePath = outPath;
exportJob.expiresAt = expiresAt;
await exportJob.save();
} catch (err) {
exportJob.status = ExportJobStatus.ERROR;
exportJob.errorMsg = err.message;
await exportJob.save();
console.error('[export] generation failed', err);
}
}
// ─── Route handlers ───────────────────────────────────────────────────────────
/**
* POST /api/v1/jobs/:jobId/export
* Body: { format: 'csv' | 'json', interval?: number, units?: 'metric' | 'us', fm?: boolean }
* interval thins by GPS time window and preserves sprayStat transition points.
*/
async function triggerExport(req, res) {
const jobId = parseInt(req.params.jobId, 10);
if (!isFinite(jobId)) AppParamError.throw('invalid jobId');
await ownerJob(jobId, req.uid);
const format = req.body?.format;
if (!['csv', 'json'].includes(format)) {
return res.status(HttpStatus.BAD_REQUEST).json({ error: 'format must be csv or json' });
}
const interval = parseInterval(req.body?.interval);
const rawUnits = req.body?.units;
const units = rawUnits === ExportUnits.US ? ExportUnits.US : ExportUnits.METRIC;
const fm = req.body?.fm === true; // opt-in: include Flight Master / AgDisp fields
// Deduplication: reuse an existing export for the same params within the dedup window.
// - ready + not yet expired → can be re-downloaded immediately
// - pending/processing + created within dedup window → generation already in flight
const dedupSince = new Date(Date.now() - EXPORT_DEDUP_MINS * 60 * 1000);
const existing = await ExportJob.findOne({
owner: ObjectId(req.uid),
jobId,
format,
interval: interval ?? null,
units,
fm: fm || false,
$or: [
{ status: ExportJobStatus.READY, expiresAt: { $gt: new Date() } },
{ status: { $in: [ExportJobStatus.PENDING, ExportJobStatus.PROCESSING] }, createdAt: { $gte: dedupSince } }
]
}).sort({ createdAt: -1 }).lean();
if (existing) {
const statusCode = existing.status === ExportJobStatus.READY ? HttpStatus.OK : HttpStatus.ACCEPTED;
const payload = {
exportId: existing._id,
status: existing.status,
format: existing.format,
units: existing.units,
createdAt: existing.createdAt,
reused: true
};
if (existing.status === ExportJobStatus.READY) payload.downloadUrl = `/api/v1/exports/${existing._id}/download`;
return res.status(statusCode).json(payload);
}
const exportJob = await ExportJob.create({
owner: ObjectId(req.uid),
jobId,
format,
interval,
units,
fm,
status: ExportJobStatus.PENDING
});
// Kick off async generation — do not await
setImmediate(() => generateExport(exportJob._id));
res.status(HttpStatus.ACCEPTED).json({
exportId: exportJob._id,
status: exportJob.status,
format: exportJob.format,
units: exportJob.units,
createdAt: exportJob.createdAt
});
}
/**
* GET /api/v1/exports/:exportId
* Poll for export status. When ready, includes downloadUrl.
*/
async function getExportStatus(req, res) {
const exportId = req.params.exportId;
if (!ObjectId.isValid(exportId)) AppParamError.throw('invalid exportId');
const exportJob = await ExportJob.findOne({
_id: ObjectId(exportId),
owner: ObjectId(req.uid)
}).lean();
if (!exportJob) return res.status(HttpStatus.NOT_FOUND).json({ error: Errors.NOT_FOUND });
const payload = {
exportId: exportJob._id,
status: exportJob.status,
format: exportJob.format,
units: exportJob.units,
createdAt: exportJob.createdAt,
expiresAt: exportJob.expiresAt ?? null,
error: exportJob.errorMsg ?? null
};
if (exportJob.status === ExportJobStatus.READY) {
// Provide a download URL — the frontend calls this to stream the file
payload.downloadUrl = `/api/v1/exports/${exportId}/download`;
}
res.json(payload);
}
/**
* GET /api/v1/exports/:exportId/download
* Streams the generated export file. Schedules file deletion after streaming.
*/
async function downloadExport(req, res) {
const exportId = req.params.exportId;
if (!ObjectId.isValid(exportId)) AppParamError.throw('invalid exportId');
const exportJob = await ExportJob.findOne({
_id: ObjectId(exportId),
owner: ObjectId(req.uid),
status: ExportJobStatus.READY
}).lean();
if (!exportJob || !exportJob.filePath) {
return res.status(HttpStatus.NOT_FOUND).json({ error: Errors.NOT_FOUND });
}
const ext = exportJob.format === 'json' ? 'json' : 'csv';
const contentType = exportJob.format === 'json' ? 'application/json' : 'text/csv';
const filename = `export_job${exportJob.jobId}_${exportJob._id}.${ext}`;
res.setHeader('Content-Type', contentType);
res.setHeader('Content-Disposition', `attachment; filename="${filename}"`);
const readStream = fs.createReadStream(exportJob.filePath);
readStream.pipe(res);
readStream.on('error', (err) => {
console.error('[export] stream error', err);
res.end();
});
}
module.exports = { triggerExport, getExportStatus, downloadExport };