From d6549c07f84d452ea7314bb6164e4e13afcc66af Mon Sep 17 00:00:00 2001 From: jmcqueen Date: Wed, 26 Aug 2026 07:06:48 -0400 Subject: [PATCH] Add FedEx shipment tracking with polling, ops delivery alerts, and backfill. Auto-register tracking numbers from SC notes, poll FedEx on a business-hours cron, and expose /trackStatus, /trackTest, and admin /track-backfill. Co-authored-by: Cursor --- .env.example | 17 +- index.js | 36 +- src/bot/index.js | 8 + src/commands/help.js | 4 + src/commands/woTrackStatus.js | 61 +++ src/commands/woTrackTest.js | 83 +++ src/config/secrets.js | 7 + src/integrations/fedex/client.js | 199 +++++++ src/server/adminAuth.js | 2 +- src/server/app.js | 105 +++- src/services/shipmentTrackingService.js | 654 ++++++++++++++++++++++++ src/services/webhookProcessor.js | 14 + 12 files changed, 1186 insertions(+), 4 deletions(-) create mode 100644 src/commands/woTrackStatus.js create mode 100644 src/commands/woTrackTest.js create mode 100644 src/integrations/fedex/client.js create mode 100644 src/services/shipmentTrackingService.js diff --git a/.env.example b/.env.example index 7484657..a4a4c97 100644 --- a/.env.example +++ b/.env.example @@ -50,7 +50,7 @@ CS_API_BASE=https://bot.joesjavajoint.com/CollabSupport # Cron for daily consolidated digest (default 14:00 UTC). Digest only — no auto-delete. # SPACE_CLEANUP_REMINDER_CRON=0 0 14 * * * -# --- Admin endpoints (/cleanup-test, /stale-workorders) --- +# --- Admin endpoints (/cleanup-test, /stale-workorders, /track-backfill) --- # Required in production. If unset in NODE_ENV=production the endpoints refuse # requests with 503. In dev (NODE_ENV!=production) unset means "allow" with a # warning in the log. @@ -96,6 +96,21 @@ SC_WEBHOOK_AUTH_MODE=off # off | log | enforce # If unset, ServChan picks the first reason matching "revised"/"superseded"/etc. # SC_PROPOSAL_REJECT_REASON_ID=7 +# SPACE_CLEANUP_REMINDER_CRON=0 0 14 * * * + +# --- FedEx shipment tracking (optional) --- +# Register at https://developer.fedex.com → create project → add Track API v1. +# Sandbox: FEDEX_API_BASE=https://apis-sandbox.fedex.com +# Production: FEDEX_API_BASE=https://apis.fedex.com +# FEDEX_CLIENT_ID=your-api-key +# FEDEX_CLIENT_SECRET=your-secret-key +# FEDEX_ACCOUNT_NUMBER=your-fedex-account-number +# FEDEX_API_BASE=https://apis.fedex.com +# SHIPMENT_TRACKING_ENABLED=true +# Every 2 hours 8am–6pm Eastern (override cron/timezone if needed) +# SHIPMENT_TRACKING_CRON=0 0 8,10,12,14,16,18 * * * +# SHIPMENT_TRACKING_TIMEZONE=America/New_York + # --- Attachment auto-post (optional) --- # When true (default), ServChan posts SC photos/invoices to the Webex WO room # on room creation and when WorkOrderNoteAdded webhooks include AttachmentIds. diff --git a/index.js b/index.js index 680e609..700e195 100644 --- a/index.js +++ b/index.js @@ -10,6 +10,7 @@ import { createApp, setupCron } from './src/server/app.js'; import { summarizeTicketDescription } from './src/integrations/xai/client.js'; import { WebexService } from './src/services/webexService.js'; import { getStaleWorkOrdersReport } from './src/services/staleWorkOrderReportService.js'; +import { logTrackingBootStatus } from './src/services/shipmentTrackingService.js'; // ──────────────────────────────────────────────────────────────────────────────── // CONFIGURATION @@ -109,6 +110,37 @@ db.serialize(() => { noteText TEXT ) `); + db.run(` + CREATE TABLE IF NOT EXISTS shipment_tracking ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + workOrderId INTEGER NOT NULL, + roomId TEXT NOT NULL, + carrier TEXT NOT NULL DEFAULT 'FEDEX', + trackingNumber TEXT NOT NULL, + statusCode TEXT, + statusDescription TEXT, + estimatedDelivery TEXT, + deliveredAt TEXT, + lastCheckedAt TEXT, + lastPostedStatus TEXT, + detectedAt TEXT NOT NULL, + sourceNote TEXT, + opsDeliveredNotifiedAt TEXT, + UNIQUE(workOrderId, trackingNumber) + ) + `); + db.run(` + CREATE INDEX IF NOT EXISTS idx_shipment_tracking_active + ON shipment_tracking(deliveredAt) + `); + db.run( + `ALTER TABLE shipment_tracking ADD COLUMN opsDeliveredNotifiedAt TEXT`, + (err) => { + if (err && !/duplicate column name/i.test(err.message)) { + console.warn(`[DB] shipment_tracking migration: ${err.message}`); + } + } + ); console.log(`[DB] Connected to ${DB_PATH}`); }); @@ -155,7 +187,9 @@ const Framework = initializeBot({ // ──────────────────────────────────────────────────────────────────────────────── // Cron + Express App (extracted) // ──────────────────────────────────────────────────────────────────────────────── -setupCron({ db, runSpaceCleanup }); +setupCron({ db, runSpaceCleanup, webex: webexService }); + +logTrackingBootStatus(); const app = createApp({ db, diff --git a/src/bot/index.js b/src/bot/index.js index e3bf8d3..64584c4 100644 --- a/src/bot/index.js +++ b/src/bot/index.js @@ -18,6 +18,8 @@ import { handleAvStatus } from '../commands/avStatus.js'; import { handleWoApprove } from '../commands/woApprove.js'; import { handleWoConfirmed } from '../commands/woConfirmed.js'; import { handleWoAddNote } from '../commands/woAddNote.js'; +import { handleWoTrackStatus } from '../commands/woTrackStatus.js'; +import { handleWoTrackTest } from '../commands/woTrackTest.js'; import approvalService from '../services/approvalService.js'; import invoiceApprovalService from '../services/invoiceApprovalService.js'; import { installMercuryGuard } from './mercuryGuard.js'; @@ -139,6 +141,12 @@ export function initializeBot({ webexConfig }) { case 'addnote': await handleWoAddNote(bot, trigger); break; + case 'trackstatus': + await handleWoTrackStatus(bot, trigger); + break; + case 'tracktest': + await handleWoTrackTest(bot, trigger); + break; default: await handleUnknown(bot, trigger); break; diff --git a/src/commands/help.js b/src/commands/help.js index 6ebae26..1b7e8d9 100644 --- a/src/commands/help.js +++ b/src/commands/help.js @@ -12,12 +12,16 @@ export async function handleHelp(bot, trigger) { text += `- **/woApprove** — (re)post proposal approval card for current WO (when WAITING FOR APPROVAL)\n`; text += `- **/confirmed** — confirm WO resolution (SC CONFIRMED if allowed, else note + ops)\n`; text += `- **/addNote** — add a note to this work order in ServiceChannel (include your text after the command)\n`; + text += `- **/trackStatus** — refresh FedEx shipment tracking for this WO (auto-poll every 2h, 8am–6pm ET)\n`; + text += `- **/trackTest** — one-off FedEx API lookup (not stored; for testing)\n`; + text += `\nFedEx tracking numbers in SC notes (e.g. \`FEDEX 876239420601\`) are registered automatically. Delivered packages notify the ops space.\n`; } else { text += `- **/woSummary ** — status of any work order\n`; text += `- **/woHistory ** — History of AV issues.\n`; text += `- **/woAttachments ** — download attachments for any work order\n`; text += `- **/avStatus ** — AV device status for any store\n`; text += `- **/woApprove ** — post proposal approval card (manual)\n`; + text += `- **/trackTest ** — one-off FedEx API lookup (not stored)\n`; } text += `\nType **/help** again in a 1:1 or group space for context-specific commands.`; diff --git a/src/commands/woTrackStatus.js b/src/commands/woTrackStatus.js new file mode 100644 index 0000000..21819dc --- /dev/null +++ b/src/commands/woTrackStatus.js @@ -0,0 +1,61 @@ +// src/commands/woTrackStatus.js +import botClient from '../integrations/webex/botClient.js'; +import db from '../db/mappings.js'; +import { logger } from '../utils/logger.js'; +import { resolveWoRoomContext } from '../utils/woRoomContext.js'; +import { + isFedExConfigured, + isShipmentTrackingEnabled, +} from '../integrations/fedex/client.js'; +import { + pollShipmentsForWorkOrder, + formatShipmentStatusMarkdown, +} from '../services/shipmentTrackingService.js'; +import webexService from '../services/webexService.js'; + +export async function handleWoTrackStatus(bot, trigger) { + logger('wo:trackStatus', 'HANDLER ENTERED'); + + if (!isShipmentTrackingEnabled() || !isFedExConfigured()) { + await bot.say('FedEx shipment tracking is not configured on this bot.'); + return; + } + + const roomId = trigger.message?.roomId; + const isGroup = trigger.message?.roomType === 'group'; + + const ctx = await resolveWoRoomContext({ db, botClient, roomId, isGroup }); + if (!ctx.ok) { + await bot.say(ctx.error); + return; + } + + await bot.say('Checking FedEx tracking status for this work order…'); + + const result = await pollShipmentsForWorkOrder({ + db, + webex: webexService, + workOrderId: ctx.workOrderId, + }); + + if (result.skipped) { + await bot.say('Tracking poll could not run.'); + return; + } + + if (result.checked === 0) { + await bot.say('No active FedEx shipments are registered for this work order.'); + return; + } + + const lines = (result.shipments || []).map((s) => formatShipmentStatusMarkdown(s)); + const summary = + `${result.checked} checked · ${result.updated} status update(s) posted · ${result.delivered} newly delivered`; + + await bot.say({ + markdown: + `**FedEx tracking refresh**\n\n` + + `${lines.join('\n')}\n\n` + + `_${summary}_`, + }); +} diff --git a/src/commands/woTrackTest.js b/src/commands/woTrackTest.js new file mode 100644 index 0000000..d9fb052 --- /dev/null +++ b/src/commands/woTrackTest.js @@ -0,0 +1,83 @@ +// src/commands/woTrackTest.js +// One-off FedEx API lookup — does not persist to shipment_tracking. + +import { logger } from '../utils/logger.js'; +import { isFedExConfigured, trackByNumbers } from '../integrations/fedex/client.js'; +import { parseFedExTrackingNumbers } from '../services/shipmentTrackingService.js'; + +const PLAIN_NUMBER_RE = /^\d{12,22}$/; + +function parseTrackingArgs(args = []) { + const raw = args.join(' ').trim(); + if (!raw) return []; + + const fromFedExPrefix = parseFedExTrackingNumbers(`FEDEX ${raw}`); + if (fromFedExPrefix.length) return fromFedExPrefix; + + const found = new Set(); + for (const part of raw.split(/[,\s]+/)) { + const num = part.trim().replace(/^fedex\s*/i, ''); + if (PLAIN_NUMBER_RE.test(num)) found.add(num); + } + return [...found]; +} + +function formatResult(result) { + const url = `https://www.fedex.com/fedextrack/?trknbr=${encodeURIComponent(result.trackingNumber)}`; + if (!result.ok) { + return `- **[${result.trackingNumber}](${url})** — API error: ${result.error}`; + } + + let line = `- **[${result.trackingNumber}](${url})**`; + if (result.statusDescription) line += ` — ${result.statusDescription}`; + else if (result.statusCode) line += ` — ${result.statusCode}`; + if (result.estimatedDelivery) { + line += ` · ETA ${new Date(result.estimatedDelivery).toLocaleDateString('en-US')}`; + } + if (result.deliveredAt) { + line += ` · Delivered ${new Date(result.deliveredAt).toLocaleDateString('en-US')}`; + } + return line; +} + +export async function handleWoTrackTest(bot, trigger) { + logger('wo:trackTest', 'HANDLER ENTERED'); + + if (!isFedExConfigured()) { + await bot.say('FedEx API is not configured (`FEDEX_CLIENT_ID`, `FEDEX_CLIENT_SECRET`, `FEDEX_API_BASE`).'); + return; + } + + const numbers = parseTrackingArgs(trigger.args); + if (!numbers.length) { + await bot.say( + 'Usage: `/trackTest 876239420601` or `/trackTest FEDEX 876239420601, 876239420612`\n\n' + + 'Looks up status via the FedEx API only — nothing is saved to ServChan.' + ); + return; + } + + if (numbers.length > 30) { + await bot.say('FedEx allows up to 30 tracking numbers per request. Please pass fewer numbers.'); + return; + } + + await bot.say(`Querying FedEx for **${numbers.length}** tracking number(s)…`); + + try { + const results = await trackByNumbers(numbers); + const lines = numbers.map((n) => formatResult(results.get(n) || { trackingNumber: n, ok: false, error: 'No response' })); + const apiBase = process.env.FEDEX_API_BASE || 'https://apis.fedex.com'; + + await bot.say({ + markdown: + `**FedEx API test** (\`${apiBase}\`)\n\n` + + `${lines.join('\n')}\n\n` + + `_Not stored — use SC notes or /trackStatus for ongoing WO tracking._`, + }); + logger('wo:trackTest', `Looked up ${numbers.length} number(s): ${numbers.join(', ')}`); + } catch (err) { + logger('wo:trackTest', `FedEx API error: ${err.message}`, 'error'); + await bot.say(`FedEx API request failed: ${err.message}`); + } +} diff --git a/src/config/secrets.js b/src/config/secrets.js index 3a120a7..aafeee5 100644 --- a/src/config/secrets.js +++ b/src/config/secrets.js @@ -66,6 +66,13 @@ export function loadSecrets() { optisign: { apiKey: optionalEnv('OPTISIGN_API_KEY'), }, + + fedex: { + clientId: optionalEnv('FEDEX_CLIENT_ID'), + clientSecret: optionalEnv('FEDEX_CLIENT_SECRET'), + accountNumber: optionalEnv('FEDEX_ACCOUNT_NUMBER'), + apiBase: optionalEnv('FEDEX_API_BASE') || 'https://apis.fedex.com', + }, }; } diff --git a/src/integrations/fedex/client.js b/src/integrations/fedex/client.js new file mode 100644 index 0000000..8dcc6b1 --- /dev/null +++ b/src/integrations/fedex/client.js @@ -0,0 +1,199 @@ +/** + * FedEx Track API client (OAuth 2.0 + track/v1/trackingnumbers). + */ + +import axios from 'axios'; +import { Mutex } from 'async-mutex'; +import qs from 'querystring'; +import { loadSecrets } from '../../config/secrets.js'; +import { logger } from '../../utils/logger.js'; + +const secrets = loadSecrets(); +const mutex = new Mutex(); + +let cachedToken = null; +let tokenExpiresAt = 0; + +function getApiBase() { + return (process.env.FEDEX_API_BASE || secrets.fedex?.apiBase || 'https://apis.fedex.com').replace(/\/$/, ''); +} + +export function isFedExConfigured() { + const clientId = process.env.FEDEX_CLIENT_ID || secrets.fedex?.clientId; + const clientSecret = process.env.FEDEX_CLIENT_SECRET || secrets.fedex?.clientSecret; + return Boolean(clientId && clientSecret); +} + +export function isShipmentTrackingEnabled() { + const raw = process.env.SHIPMENT_TRACKING_ENABLED; + if (raw == null || raw === '') return isFedExConfigured(); + return !/^(0|false|no|off)$/i.test(String(raw).trim()); +} + +async function getFedExToken(forceRefresh = false) { + const clientId = process.env.FEDEX_CLIENT_ID || secrets.fedex?.clientId; + const clientSecret = process.env.FEDEX_CLIENT_SECRET || secrets.fedex?.clientSecret; + + if (!clientId || !clientSecret) { + throw new Error('FedEx credentials not configured (FEDEX_CLIENT_ID, FEDEX_CLIENT_SECRET)'); + } + + const release = await mutex.acquire(); + try { + const now = Date.now(); + if (!forceRefresh && cachedToken && now < tokenExpiresAt) { + return cachedToken; + } + + const base = getApiBase(); + const response = await axios.post( + `${base}/oauth/token`, + qs.stringify({ + grant_type: 'client_credentials', + client_id: clientId, + client_secret: clientSecret, + }), + { + headers: { 'Content-Type': 'application/x-www-form-urlencoded' }, + timeout: 15000, + } + ); + + const { access_token, expires_in } = response.data; + cachedToken = access_token; + tokenExpiresAt = now + (Number(expires_in) || 3600) * 1000 - 300_000; + logger('fedex:oauth', 'New FedEx access token acquired'); + return access_token; + } catch (err) { + const msg = err.response?.data + ? JSON.stringify(err.response.data) + : err.message; + logger('fedex:oauth', `Token fetch failed: ${msg}`, 'error'); + throw new Error(`FedEx OAuth failed: ${msg}`); + } finally { + release(); + } +} + +function pickLatestStatus(trackResult) { + const detail = trackResult?.latestStatusDetail || {}; + const scanEvents = trackResult?.scanEvents || []; + const latestScan = scanEvents.length > 0 ? scanEvents[0] : null; + + const statusCode = detail.code || detail.derivedCode || latestScan?.eventType || ''; + const statusDescription = + detail.description || + detail.statusByLocale || + latestScan?.eventDescription || + latestScan?.derivedStatus || + ''; + + let estimatedDelivery = null; + const dateTimes = trackResult?.dateAndTimes || trackResult?.estimatedDeliveryTimeWindow || []; + if (Array.isArray(dateTimes)) { + const est = dateTimes.find((d) => + /ESTIMATED_DELIVERY|ANTICIPATED|TENDER/i.test(String(d.type || d.dateTimeType || '')) + ); + if (est?.dateTime) estimatedDelivery = est.dateTime; + } + if (!estimatedDelivery && trackResult?.standardTransitTimeWindow?.window?.ends) { + estimatedDelivery = trackResult.standardTransitTimeWindow.window.ends; + } + + let deliveredAt = null; + if (String(statusCode).toUpperCase() === 'DL' || /delivered/i.test(statusDescription)) { + const deliveredEvent = scanEvents.find((e) => + String(e.eventType || '').toUpperCase() === 'DL' || /delivered/i.test(e.eventDescription || '') + ); + deliveredAt = deliveredEvent?.date || detail.ancillaryDetails?.[0]?.value || new Date().toISOString(); + } + + return { + statusCode: String(statusCode || '').trim(), + statusDescription: String(statusDescription || '').trim(), + estimatedDelivery, + deliveredAt, + }; +} + +function normalizeTrackResult(trackingNumber, completeTrackResult) { + const trackResults = completeTrackResult?.trackResults || []; + const trackResult = trackResults[0] || null; + + if (!trackResult) { + const error = completeTrackResult?.error || completeTrackResult?.alerts?.[0]; + return { + trackingNumber, + ok: false, + error: error?.message || error?.code || 'No track results returned', + }; + } + + const status = pickLatestStatus(trackResult); + return { + trackingNumber, + ok: true, + ...status, + }; +} + +/** + * Track up to 30 FedEx tracking numbers in one API call. + * @returns {Promise>} trackingNumber → result + */ +export async function trackByNumbers(trackingNumbers = [], { _retried = false } = {}) { + const nums = [...new Set(trackingNumbers.map((n) => String(n).trim()).filter(Boolean))]; + if (!nums.length) return new Map(); + + const token = await getFedExToken(_retried); + const base = getApiBase(); + const accountNumber = process.env.FEDEX_ACCOUNT_NUMBER || secrets.fedex?.accountNumber || null; + + const body = { + includeDetailedScans: true, + trackingInfo: nums.map((trackingNumber) => ({ + trackingNumberInfo: { trackingNumber }, + })), + }; + + if (accountNumber) { + body.shipperAccountNumber = { value: String(accountNumber) }; + } + + try { + const response = await axios.post(`${base}/track/v1/trackingnumbers`, body, { + headers: { + Authorization: `Bearer ${token}`, + 'Content-Type': 'application/json', + 'X-locale': 'en_US', + }, + timeout: 30000, + }); + const results = new Map(); + const complete = response.data?.output?.completeTrackResults || []; + + for (const entry of complete) { + const num = entry?.trackingNumber || entry?.trackResults?.[0]?.trackingNumberInfo?.trackingNumber; + if (!num) continue; + results.set(String(num), normalizeTrackResult(String(num), entry)); + } + + for (const num of nums) { + if (!results.has(num)) { + results.set(num, { trackingNumber: num, ok: false, error: 'Not returned by FedEx API' }); + } + } + + return results; + } catch (err) { + if (err.response?.status === 401 && !_retried) { + await getFedExToken(true); + return trackByNumbers(nums, { _retried: true }); + } + const msg = err.response?.data ? JSON.stringify(err.response.data) : err.message; + logger('fedex:track', `Track request failed: ${msg}`, 'error'); + throw new Error(`FedEx track failed: ${msg}`); + } +} + +export default { isFedExConfigured, isShipmentTrackingEnabled, trackByNumbers }; diff --git a/src/server/adminAuth.js b/src/server/adminAuth.js index 46590c1..62ce95e 100644 --- a/src/server/adminAuth.js +++ b/src/server/adminAuth.js @@ -2,7 +2,7 @@ * src/server/adminAuth.js * * Small, dependency-free HTTP auth middleware for admin/report endpoints - * (/cleanup-test, /stale-workorders, etc.). + * (/cleanup-test, /stale-workorders, /track-backfill, etc.). * * Accepts either: * 1. `Authorization: Bearer ` header, or diff --git a/src/server/app.js b/src/server/app.js index fcdb82b..5e0046a 100644 --- a/src/server/app.js +++ b/src/server/app.js @@ -17,6 +17,8 @@ import cron from 'node-cron'; import { logger, ensureLogDir, getLogDir } from '../utils/logger.js'; import { verifyWebhook } from './webhookAuth.js'; import { requireAdmin } from './adminAuth.js'; +import { pollActiveShipments, backfillShipmentsFromAllMappings } from '../services/shipmentTrackingService.js'; +import { isShipmentTrackingEnabled } from '../integrations/fedex/client.js'; // ──────────────────────────────────────────────────────────────────────────── // HTML escaping — used everywhere we interpolate SC-sourced data into HTML @@ -249,6 +251,91 @@ export function createApp({ res.send(html); }); + // ──────────────────────────────────────────────────────────────────────── + // /track-backfill — protected admin endpoint + // + // Dry-run by default. Scans all mapped WOs for FedEx numbers in SC notes, + // registers them silently (no WO-room post), then polls FedEx once in live + // mode. Already-delivered packages are marked without ops-room notification. + // Requires ?dryRun=false&live=true to write. + // ──────────────────────────────────────────────────────────────────────── + app.get('/track-backfill', requireAdmin(), async (req, res) => { + const requestedLive = req.query.dryRun === 'false' && req.query.live === 'true'; + const dryRun = !requestedLive; + + logger('track-backfill', `Called dryRun=${dryRun} live=${requestedLive} by ${req.ip}`); + + const result = await backfillShipmentsFromAllMappings({ db, dryRun }); + + if (result.error) { + return res.status(503).send(`

Error

${esc(result.error)}
`); + } + + const { summary, results, mode } = result; + + let html = ` + + + + FedEx Track Backfill - ${esc(mode)} + + + +

FedEx Track Backfill - ${esc(mode)}

+ +

Summary

+
${esc(JSON.stringify(summary, null, 2))}
+ +

Work Orders With Tracking (${results.length})

+ + + + + + + + + + + + `; + + for (const r of results) { + html += ` + + + + + + + + `; + } + + html += ` + +
Work Order IDRoom IDFoundWould Register / RegisteredAlready TrackedError
${esc(r.workOrderId)}${esc(r.roomId)}${esc((r.trackingNumbers || []).join(', '))}${esc((r.registered || []).join(', '))}${esc((r.alreadyTracked || []).join(', '))}${esc(r.error || '')}
+ + `; + + res.send(html); + }); + // Stale Work Order Report (Phase 0) — protected admin endpoint if (getStaleWorkOrdersReport) { app.get('/stale-workorders', requireAdmin(), async (req, res) => { @@ -428,7 +515,7 @@ async function cleanupOldLogs({ maxAgeDays = 7 } = {}) { * Convenience helper to set up the daily log cleanup cron. * Can be called from the thin bootstrap. */ -export function setupCron({ db, runSpaceCleanup } = {}) { +export function setupCron({ db, runSpaceCleanup, webex } = {}) { cron.schedule('0 15 0,8,16 * * *', () => { cleanupOldLogs().catch(err => { logger('cron:cleanupOldLogs', `Unhandled: ${err.message}`, 'error'); @@ -446,6 +533,22 @@ export function setupCron({ db, runSpaceCleanup } = {}) { }); logger('cron', `Space cleanup consolidated digest scheduled (${reminderCron})`); } + + if (db && webex && isShipmentTrackingEnabled()) { + const trackingCron = process.env.SHIPMENT_TRACKING_CRON || '0 0 8,10,12,14,16,18 * * *'; + const trackingTz = process.env.SHIPMENT_TRACKING_TIMEZONE || 'America/New_York'; + cron.schedule( + trackingCron, + () => { + logger('cron:shipmentTracking', 'Starting FedEx shipment status poll'); + pollActiveShipments({ db, webex }).catch((err) => { + logger('cron:shipmentTracking', `Unhandled: ${err.message}`, 'error'); + }); + }, + { timezone: trackingTz } + ); + logger('cron', `FedEx shipment tracking poll scheduled (${trackingCron}, ${trackingTz})`); + } } export { cleanupOldLogs }; diff --git a/src/services/shipmentTrackingService.js b/src/services/shipmentTrackingService.js new file mode 100644 index 0000000..648a658 --- /dev/null +++ b/src/services/shipmentTrackingService.js @@ -0,0 +1,654 @@ +/** + * FedEx shipment tracking: parse SC notes, persist, poll, post to WO rooms. + */ + +import { logger } from '../utils/logger.js'; +import pLimit from 'p-limit'; +import { fetchWithRetry } from '../integrations/serviceChannel/client.js'; +import { + isFedExConfigured, + isShipmentTrackingEnabled, + trackByNumbers, +} from '../integrations/fedex/client.js'; +import webexService from './webexService.js'; +import { + getCompletedOperationsRoomId, + parseServChanRoomTitle, +} from '../utils/servchanRoomTitle.js'; + +const BATCH_SIZE = 30; +const BACKFILL_CONCURRENCY = 2; +const FEDEX_TRACKING_RE = /\bFEDEX\s+([\d,\s]+)/gi; +const FEDEX_NUMBER_RE = /^\d{12,22}$/; + +let bootLogged = false; + +export function parseFedExTrackingNumbers(noteText) { + if (!noteText) return []; + + const found = new Set(); + let match; + FEDEX_TRACKING_RE.lastIndex = 0; + + while ((match = FEDEX_TRACKING_RE.exec(String(noteText))) !== null) { + const chunk = match[1] || ''; + for (const part of chunk.split(/[,\s]+/)) { + const num = part.trim(); + if (FEDEX_NUMBER_RE.test(num)) found.add(num); + } + } + + return [...found]; +} + +function fedExTrackUrl(trackingNumber) { + return `https://www.fedex.com/fedextrack/?trknbr=${encodeURIComponent(trackingNumber)}`; +} + +function isDelivered(statusCode, statusDescription) { + if (String(statusCode).toUpperCase() === 'DL') return true; + return /delivered/i.test(String(statusDescription || '')); +} + +function formatStatusLabel(row) { + const parts = []; + if (row.statusDescription) parts.push(row.statusDescription); + else if (row.statusCode) parts.push(row.statusCode); + else parts.push('Status unknown'); + if (row.estimatedDelivery) { + parts.push(`ETA ${new Date(row.estimatedDelivery).toLocaleDateString('en-US')}`); + } + return parts.join(' · '); +} + +function formatStatusLabelFromResult(result) { + const parts = []; + if (result.statusDescription) parts.push(result.statusDescription); + else if (result.statusCode) parts.push(result.statusCode); + else parts.push('Status unknown'); + if (result.estimatedDelivery) { + parts.push(`ETA ${new Date(result.estimatedDelivery).toLocaleDateString('en-US')}`); + } + return parts.join(' · '); +} + +export function formatShipmentStatusMarkdown({ + trackingNumber, + statusDescription = null, + statusCode = null, + estimatedDelivery = null, + deliveredAt = null, + ok = true, + error = null, +}) { + const url = fedExTrackUrl(trackingNumber); + if (!ok || error) { + return `- **[${trackingNumber}](${url})** — ${error || 'Status unavailable'}`; + } + + let line = `- **[${trackingNumber}](${url})**`; + if (statusDescription) line += ` — ${statusDescription}`; + else if (statusCode) line += ` — ${statusCode}`; + if (estimatedDelivery) { + line += ` · ETA ${new Date(estimatedDelivery).toLocaleDateString('en-US')}`; + } + if (deliveredAt) { + line += ` · Delivered ${new Date(deliveredAt).toLocaleDateString('en-US')}`; + } + return line; +} + +function shipmentSnapshotFromTrackResult(row, result) { + if (!result.ok) { + return { + trackingNumber: row.trackingNumber, + ok: false, + error: result.error || 'Status unavailable', + }; + } + + const delivered = isDelivered(result.statusCode, result.statusDescription); + return { + trackingNumber: row.trackingNumber, + statusDescription: result.statusDescription, + statusCode: result.statusCode, + estimatedDelivery: result.estimatedDelivery, + deliveredAt: delivered ? (result.deliveredAt || new Date().toISOString()) : null, + ok: true, + }; +} + +function collectTrackingNumbersFromNotes(notes = []) { + const found = new Set(); + for (const note of notes) { + for (const num of parseFedExTrackingNumbers(note?.NoteData)) { + found.add(num); + } + } + return [...found]; +} + +function dbGetAllMappings(db) { + return new Promise((resolve, reject) => { + db.all('SELECT workOrderId, roomId FROM mappings', [], (err, rows) => { + if (err) reject(err); + else resolve(rows || []); + }); + }); +} + +async function fetchWorkOrderNotes(workOrderId) { + const notesRes = await fetchWithRetry(`/workorders/${workOrderId}/notes`, { timeout: 30000 }); + return notesRes.data?.Notes || []; +} + +async function applyBackfillTrackResult(db, row, result) { + const now = new Date().toISOString(); + + if (!result.ok) { + await dbUpdateShipment(db, row.trackingNumber, row.workOrderId, { + lastCheckedAt: now, + statusCode: row.statusCode, + statusDescription: row.statusDescription, + estimatedDelivery: row.estimatedDelivery, + deliveredAt: row.deliveredAt, + lastPostedStatus: row.lastPostedStatus, + }); + return { delivered: false, inTransit: false, error: result.error }; + } + + const statusLabel = formatStatusLabelFromResult(result); + const delivered = isDelivered(result.statusCode, result.statusDescription); + const deliveredAt = delivered ? (result.deliveredAt || now) : null; + + await dbUpdateShipment(db, row.trackingNumber, row.workOrderId, { + statusCode: result.statusCode, + statusDescription: result.statusDescription, + estimatedDelivery: result.estimatedDelivery, + deliveredAt, + lastCheckedAt: now, + lastPostedStatus: statusLabel, + }); + + if (delivered) { + await dbMarkOpsDeliveredNotified(db, row.workOrderId, row.trackingNumber); + } + + return { delivered, inTransit: !delivered, error: null }; +} + +function dbGetShipment(db, workOrderId, trackingNumber) { + return new Promise((resolve, reject) => { + db.get( + `SELECT * FROM shipment_tracking WHERE workOrderId = ? AND trackingNumber = ?`, + [workOrderId, trackingNumber], + (err, row) => (err ? reject(err) : resolve(row || null)) + ); + }); +} + +function dbInsertShipment(db, row) { + return new Promise((resolve, reject) => { + db.run( + `INSERT OR IGNORE INTO shipment_tracking + (workOrderId, roomId, carrier, trackingNumber, detectedAt, sourceNote) + VALUES (?, ?, 'FEDEX', ?, ?, ?)`, + [row.workOrderId, row.roomId, row.trackingNumber, row.detectedAt, row.sourceNote ?? null], + function onRun(err) { + if (err) reject(err); + else resolve(this.changes > 0); + } + ); + }); +} + +function dbUpdateShipment(db, trackingNumber, workOrderId, fields) { + return new Promise((resolve, reject) => { + db.run( + `UPDATE shipment_tracking SET + statusCode = ?, + statusDescription = ?, + estimatedDelivery = ?, + deliveredAt = ?, + lastCheckedAt = ?, + lastPostedStatus = ? + WHERE workOrderId = ? AND trackingNumber = ?`, + [ + fields.statusCode ?? null, + fields.statusDescription ?? null, + fields.estimatedDelivery ?? null, + fields.deliveredAt ?? null, + fields.lastCheckedAt, + fields.lastPostedStatus ?? null, + workOrderId, + trackingNumber, + ], + (err) => (err ? reject(err) : resolve()) + ); + }); +} + +function dbMarkOpsDeliveredNotified(db, workOrderId, trackingNumber) { + return new Promise((resolve, reject) => { + db.run( + `UPDATE shipment_tracking SET opsDeliveredNotifiedAt = ? WHERE workOrderId = ? AND trackingNumber = ?`, + [new Date().toISOString(), workOrderId, trackingNumber], + (err) => (err ? reject(err) : resolve()) + ); + }); +} + +async function resolveSpaceHeaderLine(roomId, workOrderId) { + try { + const botClientMod = await import('../integrations/webex/botClient.js'); + const room = await botClientMod.default.getRoom(roomId); + const parsed = parseServChanRoomTitle(room?.title); + return parsed?.headerLine || room?.title || `WO ${workOrderId}`; + } catch (err) { + logger('shipmentTracking', `Could not resolve room title for ${roomId}: ${err.message}`, 'warn'); + return `WO ${workOrderId}`; + } +} + +async function notifyOpsOnDelivery(webex, row) { + if (row.opsDeliveredNotifiedAt) return false; + + const opsRoomId = getCompletedOperationsRoomId(); + if (!opsRoomId) { + logger('shipmentTracking', 'COMPLETED_OPERATIONS_ROOM_ID not set — skipping ops delivery notification', 'warn'); + return false; + } + + const headerLine = await resolveSpaceHeaderLine(row.roomId, row.workOrderId); + const md = `${headerLine}\nFedEx ${row.trackingNumber} — successfully delivered`; + + const sender = getSender(webex); + try { + await sender.sendMarkdown(opsRoomId, md); + logger( + 'shipmentTracking', + `Posted delivery notification to ops for ${row.trackingNumber} (WO ${row.workOrderId})` + ); + return true; + } catch (err) { + logger('shipmentTracking', `Ops delivery post failed for ${row.trackingNumber}: ${err.message}`, 'warn'); + return false; + } +} + +function dbGetActiveShipments(db, workOrderId = null) { + return new Promise((resolve, reject) => { + const sql = workOrderId + ? `SELECT * FROM shipment_tracking WHERE deliveredAt IS NULL AND workOrderId = ? ORDER BY detectedAt` + : `SELECT * FROM shipment_tracking WHERE deliveredAt IS NULL ORDER BY detectedAt`; + const params = workOrderId ? [workOrderId] : []; + db.all(sql, params, (err, rows) => (err ? reject(err) : resolve(rows || []))); + }); +} + +function dbCountActiveForWo(db, workOrderId) { + return new Promise((resolve, reject) => { + db.get( + `SELECT COUNT(*) AS n FROM shipment_tracking WHERE workOrderId = ? AND deliveredAt IS NULL`, + [workOrderId], + (err, row) => (err ? reject(err) : resolve(row?.n ?? 0)) + ); + }); +} + +function getSender(webex) { + return webex && typeof webex.sendMarkdown === 'function' ? webex : webexService; +} + +export function logTrackingBootStatus() { + if (bootLogged) return; + bootLogged = true; + if (!isShipmentTrackingEnabled()) { + logger('shipmentTracking', 'Disabled (SHIPMENT_TRACKING_ENABLED or missing FedEx creds)', 'info'); + return; + } + if (!isFedExConfigured()) { + logger('shipmentTracking', 'Enabled but FedEx credentials missing — tracking inactive', 'warn'); + return; + } + logger('shipmentTracking', 'FedEx shipment tracking active', 'info'); +} + +export async function registerShipmentsFromNote({ + db, + webex, + workOrderId, + roomId, + woNumber, + noteText, +}) { + if (!isShipmentTrackingEnabled() || !isFedExConfigured()) return { registered: 0 }; + if (!db || !roomId || !workOrderId || !noteText) return { registered: 0 }; + + const numbers = parseFedExTrackingNumbers(noteText); + if (!numbers.length) return { registered: 0 }; + + const sender = getSender(webex); + const newlyRegistered = []; + + for (const trackingNumber of numbers) { + const existing = await dbGetShipment(db, workOrderId, trackingNumber); + if (existing) continue; + + const inserted = await dbInsertShipment(db, { + workOrderId, + roomId, + trackingNumber, + detectedAt: new Date().toISOString(), + sourceNote: noteText.substring(0, 500), + }); + + if (inserted) newlyRegistered.push(trackingNumber); + } + + if (!newlyRegistered.length) return { registered: 0 }; + + const lines = newlyRegistered.map( + (n) => `- [${n}](${fedExTrackUrl(n)})` + ); + const woLabel = woNumber ? `WO-${woNumber}` : `WO ${workOrderId}`; + const md = + `**FedEx shipment tracking registered** for ${woLabel}\n\n` + + `${lines.join('\n')}\n\n` + + `ServChan will post status updates every 2 hours (8am–6pm Eastern) while in transit.`; + + try { + await sender.sendMarkdown(roomId, md); + } catch (err) { + logger('shipmentTracking', `Initial post failed for WO ${workOrderId}: ${err.message}`, 'warn'); + } + + logger( + 'shipmentTracking', + `Registered ${newlyRegistered.length} tracking number(s) for WO ${workOrderId}: ${newlyRegistered.join(', ')}` + ); + + return { registered: newlyRegistered.length, trackingNumbers: newlyRegistered }; +} + +async function applyTrackResult(db, row, result, webex, woNumberByWoId) { + const now = new Date().toISOString(); + + if (!result.ok) { + logger( + 'shipmentTracking', + `Track failed for ${row.trackingNumber} (WO ${row.workOrderId}): ${result.error}`, + 'warn' + ); + await dbUpdateShipment(db, row.trackingNumber, row.workOrderId, { + lastCheckedAt: now, + statusCode: row.statusCode, + statusDescription: row.statusDescription, + estimatedDelivery: row.estimatedDelivery, + deliveredAt: row.deliveredAt, + lastPostedStatus: row.lastPostedStatus, + }); + return { updated: false, delivered: false }; + } + + const statusLabel = formatStatusLabel(result); + const delivered = isDelivered(result.statusCode, result.statusDescription); + const deliveredAt = delivered ? (result.deliveredAt || now) : null; + const statusChanged = statusLabel !== row.lastPostedStatus; + const newlyDelivered = delivered && !row.deliveredAt; + + await dbUpdateShipment(db, row.trackingNumber, row.workOrderId, { + statusCode: result.statusCode, + statusDescription: result.statusDescription, + estimatedDelivery: result.estimatedDelivery, + deliveredAt, + lastCheckedAt: now, + lastPostedStatus: statusChanged ? statusLabel : row.lastPostedStatus, + }); + + if (statusChanged) { + const sender = getSender(webex); + const woNum = woNumberByWoId?.get(row.workOrderId) || row.workOrderId; + const md = + `**FedEx update** · WO-${woNum}\n\n` + + `Tracking **[${row.trackingNumber}](${fedExTrackUrl(row.trackingNumber)})**: ${statusLabel}` + + (delivered ? '\n\nPackage delivered.' : ''); + + try { + await sender.sendMarkdown(row.roomId, md); + } catch (err) { + logger('shipmentTracking', `Update post failed for ${row.trackingNumber}: ${err.message}`, 'warn'); + } + } + + if (newlyDelivered) { + const opsNotified = await notifyOpsOnDelivery(webex, row); + if (opsNotified) { + await dbMarkOpsDeliveredNotified(db, row.workOrderId, row.trackingNumber); + } + } + + if (delivered) { + const remaining = await dbCountActiveForWo(db, row.workOrderId); + if (remaining === 0) { + const sender = getSender(webex); + const woNum = woNumberByWoId?.get(row.workOrderId) || row.workOrderId; + try { + await sender.sendMarkdown( + row.roomId, + `**All parts shipments delivered** for WO-${woNum}.` + ); + } catch (err) { + logger('shipmentTracking', `All-delivered post failed for WO ${row.workOrderId}: ${err.message}`, 'warn'); + } + } + } + + return { updated: statusChanged, delivered }; +} + +export async function pollActiveShipments({ db, webex, workOrderId = null } = {}) { + if (!isShipmentTrackingEnabled() || !isFedExConfigured()) { + return { checked: 0, updated: 0, delivered: 0, shipments: [], skipped: true }; + } + if (!db) return { checked: 0, updated: 0, delivered: 0, shipments: [], skipped: true }; + + const active = await dbGetActiveShipments(db, workOrderId); + if (!active.length) return { checked: 0, updated: 0, delivered: 0, shipments: [] }; + + let updated = 0; + let delivered = 0; + const shipments = []; + + for (let i = 0; i < active.length; i += BATCH_SIZE) { + const batch = active.slice(i, i + BATCH_SIZE); + const numbers = batch.map((r) => r.trackingNumber); + + let results; + try { + results = await trackByNumbers(numbers); + } catch (err) { + logger('shipmentTracking', `Batch track failed: ${err.message}`, 'error'); + for (const row of batch) { + shipments.push({ + trackingNumber: row.trackingNumber, + ok: false, + error: err.message, + }); + } + break; + } + + for (const row of batch) { + const result = results.get(row.trackingNumber) || { + ok: false, + error: 'Missing from batch response', + }; + const outcome = await applyTrackResult(db, row, result, webex, null); + if (outcome.updated) updated++; + if (outcome.delivered) delivered++; + shipments.push(shipmentSnapshotFromTrackResult(row, result)); + } + } + + logger( + 'shipmentTracking', + `Poll complete: checked=${active.length}, updated=${updated}, delivered=${delivered}` + + (workOrderId ? ` (WO ${workOrderId})` : '') + ); + + return { checked: active.length, updated, delivered, shipments }; +} + +export async function pollShipmentsForWorkOrder({ db, webex, workOrderId }) { + return pollActiveShipments({ db, webex, workOrderId }); +} + +/** + * One-shot backfill: scan all mapped WO spaces for FedEx numbers in SC notes, + * insert into shipment_tracking without WO-room registration posts, then (live + * mode) poll FedEx once. Option B: already-delivered packages get deliveredAt + + * opsDeliveredNotifiedAt without ops-room notification; lastPostedStatus is set + * so cron won't flood WO rooms. + */ +export async function backfillShipmentsFromAllMappings({ db, dryRun = true } = {}) { + const mode = dryRun ? 'DRY RUN' : 'LIVE'; + + if (!db) { + return { error: 'Database not available', dryRun, mode }; + } + if (!isShipmentTrackingEnabled() || !isFedExConfigured()) { + return { + error: 'Shipment tracking disabled or FedEx credentials missing', + dryRun, + mode, + }; + } + + const mappings = await dbGetAllMappings(db); + const limit = pLimit(BACKFILL_CONCURRENCY); + const results = []; + const newlyInsertedRows = []; + + const summary = { + mappingsScanned: mappings.length, + mappingsWithTracking: 0, + trackingNumbersFound: 0, + registered: 0, + alreadyTracked: 0, + fedExChecked: 0, + alreadyDelivered: 0, + inTransit: 0, + trackErrors: 0, + noteFetchErrors: 0, + }; + + await Promise.all( + mappings.map((mapping) => + limit(async () => { + const { workOrderId, roomId } = mapping; + const rowResult = { + workOrderId, + roomId, + trackingNumbers: [], + registered: [], + alreadyTracked: [], + }; + + try { + const notes = await fetchWorkOrderNotes(workOrderId); + const numbers = collectTrackingNumbersFromNotes(notes); + rowResult.trackingNumbers = numbers; + summary.trackingNumbersFound += numbers.length; + if (numbers.length) summary.mappingsWithTracking++; + + for (const trackingNumber of numbers) { + const existing = await dbGetShipment(db, workOrderId, trackingNumber); + if (existing) { + rowResult.alreadyTracked.push(trackingNumber); + summary.alreadyTracked++; + continue; + } + + if (dryRun) { + rowResult.registered.push(trackingNumber); + summary.registered++; + continue; + } + + const inserted = await dbInsertShipment(db, { + workOrderId, + roomId, + trackingNumber, + detectedAt: new Date().toISOString(), + sourceNote: 'backfill', + }); + + if (inserted) { + rowResult.registered.push(trackingNumber); + summary.registered++; + const row = await dbGetShipment(db, workOrderId, trackingNumber); + if (row) newlyInsertedRows.push(row); + } + } + } catch (err) { + rowResult.error = err.message; + summary.noteFetchErrors++; + logger( + 'shipmentTracking', + `Backfill note fetch failed for WO ${workOrderId}: ${err.message}`, + 'warn' + ); + } + + if (rowResult.trackingNumbers.length || rowResult.error) { + results.push(rowResult); + } + }) + ) + ); + + results.sort((a, b) => String(a.workOrderId).localeCompare(String(b.workOrderId))); + + if (!dryRun && newlyInsertedRows.length) { + for (let i = 0; i < newlyInsertedRows.length; i += BATCH_SIZE) { + const batch = newlyInsertedRows.slice(i, i + BATCH_SIZE); + const numbers = batch.map((r) => r.trackingNumber); + + let trackResults; + try { + trackResults = await trackByNumbers(numbers); + } catch (err) { + logger('shipmentTracking', `Backfill FedEx batch failed: ${err.message}`, 'error'); + summary.trackErrors += batch.length; + continue; + } + + for (const row of batch) { + summary.fedExChecked++; + const result = trackResults.get(row.trackingNumber) || { + ok: false, + error: 'Missing from batch response', + }; + const outcome = await applyBackfillTrackResult(db, row, result); + if (outcome.error) summary.trackErrors++; + else if (outcome.delivered) summary.alreadyDelivered++; + else if (outcome.inTransit) summary.inTransit++; + } + } + } + + logger('shipmentTracking', `Backfill complete dryRun=${dryRun} ${JSON.stringify(summary)}`); + + return { dryRun, mode, summary, results }; +} + +export default { + parseFedExTrackingNumbers, + logTrackingBootStatus, + registerShipmentsFromNote, + pollActiveShipments, + pollShipmentsForWorkOrder, + backfillShipmentsFromAllMappings, + formatShipmentStatusMarkdown, +}; diff --git a/src/services/webhookProcessor.js b/src/services/webhookProcessor.js index adea58e..8fb99d7 100644 --- a/src/services/webhookProcessor.js +++ b/src/services/webhookProcessor.js @@ -35,6 +35,7 @@ import { Mutex } from 'async-mutex'; import { logger } from '../utils/logger.js'; import approvalService from './approvalService.js'; import invoiceApprovalService from './invoiceApprovalService.js'; +import shipmentTrackingService from './shipmentTrackingService.js'; import attachmentService from './attachmentService.js'; export function createWebhookProcessor(options = {}) { @@ -241,6 +242,19 @@ export function createWebhookProcessor(options = {}) { logger('webhookProcessor:invoiceApproval', `Card remove error for ${workOrderId}: ${e.message}`, 'warn'); }); + if (p.EventType === 'WorkOrderNoteAdded' && noteDataForApproval) { + shipmentTrackingService.registerShipmentsFromNote({ + db, + webex, + roomId, + workOrderId, + woNumber: p.Object.Number, + noteText: noteDataForApproval, + }).catch((e) => { + logger('webhookProcessor:shipmentTracking', `Register error for ${workOrderId}: ${e.message}`, 'warn'); + }); + } + // NEW: Proposal approval card for WAITING FOR APPROVAL status (additive to text note) const ext = (p.Object?.Status?.Extended || '').toUpperCase(); if (ext.includes('WAITING FOR APPROVAL')) {