CI - Test & Build / Lint & Type Check (push) Canceled after 0s
CI - Test & Build / Backend Unit Tests (push) Canceled after 0s
CI - Test & Build / Frontend Build (push) Canceled after 0s
CI - Test & Build / Security Scan (push) Canceled after 0s
CI - Test & Build / E2E Tests (Playwright) (push) Canceled after 0s
CI - Test & Build / CI Summary (push) Canceled after 0s
* [verified] fix(ordrestyring): initialize database before preflight * [verified] fix(ordrestyring): preserve promise pool preflight adapter * test(ordrestyring): support pool-only service adapters --------- Co-authored-by: alexpolo1 <[email protected]>
219 lines
13 KiB
JavaScript
219 lines
13 KiB
JavaScript
const crypto = require('crypto');
|
|
const { normalizeOrdrestyringMoney } = require('./ordrestyringMoneyService');
|
|
const { classifyMaterial } = require('./ordrestyringMaterialSemantics');
|
|
|
|
/** @typedef {'legacy'|'compare'|'canonical'} ReadMode */
|
|
/** @typedef {'features'|'planning'|'history'|'suggestions'|'realism'} ReadDomain */
|
|
/** @typedef {{case_number: string, customer_number: string, description: string, total_hours: number, hour_entries: number}} CaseReadRow */
|
|
/** @typedef {{reference: string, units: string[], priceScale: 'major'|'minor', vatState: 'inclusive'|'exclusive'}} MaterialEvidence */
|
|
/** @typedef {{case_number: string, quantity: number, unit: string, unitPrice: number, totalPrice: number, vatState: 'inclusive'|'exclusive', evidence: MaterialEvidence}} MaterialReadRow */
|
|
/** @typedef {Map<string, MaterialReadRow[]> & {excludedCounts: Map<string, number>, negativeCounts: Map<string, number>, unclassifiedCounts: Map<string, number>}} MaterialReadResult */
|
|
/** @typedef {{legacyCount: number, canonicalCount: number, legacyHash: string, canonicalHash: string}} ComparisonEvent */
|
|
/** @typedef {{env?: Object<string, string>, collector?: (event: ComparisonEvent) => void|Promise<void>, logger?: {info: Function}}} ReadOptions */
|
|
// ORDRESTYRING_READ_MODE defaults to canonical. Explicit rollback: legacy.
|
|
// Per-domain overrides use ORDRESTYRING_READ_MODE_<DOMAIN>, e.g. _PLANNING.
|
|
const DOMAINS = ['features', 'planning', 'history', 'suggestions', 'realism'];
|
|
const MODES = ['legacy', 'compare', 'canonical'];
|
|
const fields = (names) => names.map(name => `NULLIF(JSON_UNQUOTE(JSON_EXTRACT(source_payload, '$.${name}')), 'null') AS ${name}`).join(', ');
|
|
const casesRelation = `(SELECT case_number, case_number AS id, ${fields(['customer_number', 'description', 'remarks', 'creation_date', 'created_at', 'updated_at', 'status', 'case_type', 'main_technician', 'additional_technicians', 'contact', 'delivery_address', 'yourref', 'work_done', 'offer_number'])} FROM ordrestyring_local.ordrestyring_cases_current)`;
|
|
const debtorsRelation = `(SELECT customer_number, ${fields(['customer_name', 'customer_address', 'customer_telephone', 'customer_email'])} FROM ordrestyring_local.ordrestyring_debtors_current)`;
|
|
const seconds = name => `(CASE WHEN ${name} > 100000000000 THEN ${name} / 1000 ELSE ${name} END)`;
|
|
const timestampEpoch = name => `(CASE
|
|
WHEN ${name} REGEXP '^[0-9]+([.][0-9]+)?$' THEN
|
|
CASE WHEN CAST(${name} AS DECIMAL(20,3)) > 100000000000
|
|
THEN CAST(${name} AS DECIMAL(20,3)) / 1000
|
|
ELSE CAST(${name} AS DECIMAL(20,3)) END
|
|
ELSE UNIX_TIMESTAMP(${name}) END)`;
|
|
const hoursRelation = `(SELECT id, emp_id, new_case_number, remark, hour_type, ${seconds('start_time')} AS start_time, ${seconds('stop_time')} AS stop_time FROM ordrestyring_local.hours h WHERE exported IS NOT NULL AND exported <> -1 AND EXISTS (SELECT 1 FROM ordrestyring_local.ordrestyring_hour_mirror_ids mirror WHERE mirror.hour_id = h.id) AND start_time > 0 AND ${seconds('stop_time')} > ${seconds('start_time')})`;
|
|
// Planning also represents local allocations whose end is not known yet.
|
|
const planningHoursRelation = `(SELECT id, emp_id, new_case_number, remark, hour_type, ${seconds('start_time')} AS start_time, ${seconds('stop_time')} AS stop_time FROM ordrestyring_local.hours WHERE start_time > 0 AND (stop_time = 0 OR ${seconds('stop_time')} > ${seconds('start_time')}))`;
|
|
const parse = value => typeof value === 'string' ? JSON.parse(value) : value;
|
|
|
|
class OrdrestyringReadModel {
|
|
/** @param {{execute?: Function, query?: Function}} connection @param {ReadOptions} options */
|
|
constructor(connection, { env = process.env, collector, logger } = {}) {
|
|
this.connection = connection;
|
|
this.collector = collector || (logger ? event => logger.info('Ordrestyring read comparison', event) : () => {});
|
|
this.modes = {};
|
|
const defaultMode = env.ORDRESTYRING_READ_MODE ?? 'canonical';
|
|
if (!MODES.includes(defaultMode)) throw new Error('Invalid Ordrestyring read mode');
|
|
for (const domain of DOMAINS) {
|
|
const mode = env[`ORDRESTYRING_READ_MODE_${domain.toUpperCase()}`] ?? defaultMode;
|
|
if (!MODES.includes(mode)) throw new Error('Invalid Ordrestyring read mode');
|
|
this.modes[domain] = mode;
|
|
}
|
|
this.hashKey = crypto.randomBytes(32);
|
|
}
|
|
|
|
/** @param {ReadDomain} domain @returns {ReadMode} */
|
|
mode(domain) {
|
|
if (!DOMAINS.includes(domain)) throw new Error('Invalid Ordrestyring read domain');
|
|
return this.modes[domain];
|
|
}
|
|
|
|
/**
|
|
* @template T
|
|
* @param {ReadDomain} domain
|
|
* @param {() => Promise<T>} legacy
|
|
* @param {() => Promise<T>} canonical
|
|
* @returns {Promise<T>}
|
|
*/
|
|
async read(domain, legacy, canonical) {
|
|
const mode = this.mode(domain);
|
|
if (mode === 'canonical') return canonical();
|
|
const result = await legacy();
|
|
if (mode === 'compare') {
|
|
let current;
|
|
try { current = await canonical(); } catch (_) { current = null; }
|
|
const summarize = value => {
|
|
const normalized = value instanceof Map ? [...value.entries()] : value;
|
|
return { count: value === null ? -1 : (value instanceof Map ? value.size : Array.isArray(value) ? value.length : 1),
|
|
hash: crypto.createHmac('sha256', this.hashKey).update(JSON.stringify(normalized) ?? 'null').digest('hex') };
|
|
};
|
|
const old = summarize(result);
|
|
const next = summarize(current);
|
|
// Never send payloads, identifiers, SQL, or exception messages to telemetry.
|
|
try { await this.collector({ legacyCount: old.count, canonicalCount: next.count, legacyHash: old.hash, canonicalHash: next.hash }); } catch (_) { /* Observation cannot change legacy output. */ }
|
|
}
|
|
return result;
|
|
}
|
|
|
|
async rows(sql, params = []) {
|
|
if (this.connection && Object.prototype.hasOwnProperty.call(this.connection, 'pool')
|
|
&& typeof this.connection.query === 'function' && typeof this.connection.execute !== 'function') {
|
|
return this.connection.query(sql, params);
|
|
}
|
|
const connection = typeof this.connection?.execute === 'function'
|
|
? this.connection
|
|
: this.connection?.pool || this.connection;
|
|
if (!connection) throw new Error('Ordrestyring canonical database unavailable');
|
|
if (connection.execute) {
|
|
const [rows] = await connection.execute(sql, params);
|
|
return rows;
|
|
}
|
|
return connection.query(sql, params);
|
|
}
|
|
|
|
// JSON_UNQUOTE returns Unix timestamps as strings; expose dates consistently to consumers.
|
|
async caseRows(sql, params) {
|
|
const rows = await this.rows(sql, params);
|
|
return rows.map(row => {
|
|
const normalized = { ...row };
|
|
for (const field of ['creation_date', 'created_at', 'updated_at', 'latest_activity_at']) {
|
|
const value = row[field];
|
|
if (typeof value !== 'number' && !(typeof value === 'string' && /^\d+(?:\.\d+)?$/.test(value.trim()))) continue;
|
|
const numeric = Number(value);
|
|
const date = new Date(numeric > 100000000000 ? numeric : numeric * 1000);
|
|
normalized[field] = Number.isFinite(date.getTime()) ? date.toISOString() : null;
|
|
}
|
|
return normalized;
|
|
});
|
|
}
|
|
|
|
/** @param {string[]} numbers @returns {Promise<CaseReadRow[]>} */
|
|
async cases(numbers) {
|
|
if (!numbers.length) return [];
|
|
return this.caseRows(`SELECT c.* FROM ${casesRelation} c WHERE c.case_number IN (${numbers.map(() => '?').join(',')})`, numbers);
|
|
}
|
|
|
|
async caseDetail(number) {
|
|
return this.caseRows(`SELECT c.*, d.customer_name, d.customer_address,
|
|
d.customer_telephone, d.customer_email, cs.text AS status_text
|
|
FROM ${casesRelation} c
|
|
LEFT JOIN ${debtorsRelation} d ON d.customer_number = c.customer_number
|
|
LEFT JOIN ordrestyring_local.case_statuses cs ON cs.id = c.status
|
|
WHERE c.case_number = ?`, [number]);
|
|
}
|
|
|
|
/** @param {string[]} numbers @returns {Promise<Map<string, string>>} */
|
|
async debtors(numbers) {
|
|
if (!numbers.length) return new Map();
|
|
const rows = await this.rows(`SELECT d.* FROM ${debtorsRelation} d WHERE d.customer_number IN (${numbers.map(() => '?').join(',')})`, numbers);
|
|
return new Map(rows.map(row => [row.customer_number, row.customer_name || '']));
|
|
}
|
|
|
|
/** Unique current relations and preaggregated hours prevent crowd-out before LIMIT. @returns {Promise<CaseReadRow[]>} */
|
|
async candidates({ tokens = [], customerNumber, customerName, limit = 40, offset = 0, regex } = {}) {
|
|
const clauses = [];
|
|
const params = [];
|
|
if (customerNumber) { clauses.push('c.customer_number = ?'); params.push(customerNumber); }
|
|
if (customerName) { clauses.push('LOWER(d.customer_name) LIKE ?'); params.push(`%${customerName.toLowerCase()}%`); }
|
|
for (const token of tokens.slice(0, 6)) {
|
|
clauses.push("LOWER(CONCAT_WS(' ', c.case_number, c.description, c.remarks, c.work_done)) LIKE ?");
|
|
params.push(`%${token.toLowerCase()}%`);
|
|
}
|
|
if (regex) { clauses.push("LOWER(CONCAT_WS(' ', c.description, c.remarks)) REGEXP ?"); params.push(regex); }
|
|
const latestActivity = `GREATEST(COALESCE(${timestampEpoch('c.updated_at')}, 0),
|
|
COALESCE(${timestampEpoch('c.creation_date')}, 0), COALESCE(${timestampEpoch('c.created_at')}, 0))`;
|
|
return this.caseRows(`SELECT c.*, d.customer_name, d.customer_address, h.total_hours, h.hour_entries,
|
|
${latestActivity} AS latest_activity_at
|
|
FROM ${casesRelation} c LEFT JOIN ${debtorsRelation} d ON d.customer_number = c.customer_number
|
|
JOIN (SELECT new_case_number, SUM((stop_time - start_time) / 3600) AS total_hours, COUNT(*) AS hour_entries
|
|
FROM ${hoursRelation} verified_hours GROUP BY new_case_number) h ON h.new_case_number = c.case_number
|
|
WHERE ${clauses.length ? `(${clauses.join(' OR ')})` : '1=1'}
|
|
ORDER BY latest_activity_at DESC, c.case_number ASC LIMIT ? OFFSET ?`, [...params, Math.max(1, Math.min(1000, Math.floor(limit))), Math.max(0, Math.min(1000, Math.floor(offset)))]);
|
|
}
|
|
|
|
async hours(numbers) {
|
|
if (!numbers.length) return new Map();
|
|
const rows = await this.rows(`SELECT *, new_case_number AS case_number FROM ${hoursRelation} h WHERE new_case_number IN (${numbers.map(() => '?').join(',')})`, numbers);
|
|
const grouped = new Map();
|
|
for (const row of rows) {
|
|
if (!(Number(row.start_time) > 0 && Number(row.stop_time) > Number(row.start_time))) continue;
|
|
const key = row.case_number;
|
|
if (!grouped.has(key)) grouped.set(key, []);
|
|
grouped.get(key).push({ ...row, start_at: new Date(row.start_time * 1000), stop_at: new Date(row.stop_time * 1000) });
|
|
}
|
|
return grouped;
|
|
}
|
|
|
|
/** @param {string[]} numbers @returns {Promise<MaterialReadResult>} */
|
|
async materials(numbers) {
|
|
const grouped = new Map();
|
|
grouped.excludedCounts = new Map();
|
|
grouped.negativeCounts = new Map();
|
|
grouped.unclassifiedCounts = new Map();
|
|
if (!numbers.length) return grouped;
|
|
const rows = await this.rows(`SELECT material_line_id, source_payload, semantic_class, evidence_json
|
|
FROM ordrestyring_local.ordrestyring_material_semantics
|
|
WHERE JSON_UNQUOTE(JSON_EXTRACT(source_payload, '$.case_number')) IN (${numbers.map(() => '?').join(',')})`, numbers);
|
|
for (const row of rows) {
|
|
const payload = parse(row.source_payload);
|
|
const evidence = parse(row.evidence_json);
|
|
const key = String(payload.case_number);
|
|
const eligible = row.semantic_class === 'normal_consumption'
|
|
&& classifyMaterial(payload, evidence).semanticClass === 'normal_consumption';
|
|
if (Number(payload.quantity) < 0) grouped.negativeCounts.set(key, (grouped.negativeCounts.get(key) || 0) + 1);
|
|
if (row.semantic_class === 'unclassified' || (row.semantic_class === 'normal_consumption' && !eligible)) {
|
|
grouped.unclassifiedCounts.set(key, (grouped.unclassifiedCounts.get(key) || 0) + 1);
|
|
}
|
|
if (!eligible) {
|
|
grouped.excludedCounts.set(key, (grouped.excludedCounts.get(key) || 0) + 1);
|
|
continue;
|
|
}
|
|
const unitPrice = Number(payload.sales_price) / (evidence.priceScale === 'minor' ? 100 : 1);
|
|
if (!grouped.has(key)) grouped.set(key, []);
|
|
const unitPriceMoney = normalizeOrdrestyringMoney(unitPrice * 100);
|
|
grouped.get(key).push({ ...payload, quantity: Number(payload.quantity), unitPrice,
|
|
unitPriceMoney, currency: unitPriceMoney.currency, minorUnit: unitPriceMoney.minorUnit,
|
|
vatState: evidence.vatState, semanticClass: row.semantic_class, evidence,
|
|
totalPrice: Number(payload.quantity) * unitPrice });
|
|
}
|
|
return grouped;
|
|
}
|
|
|
|
async materialSummary(number) {
|
|
const materials = await this.materials([number]);
|
|
const positive = materials.get(number) || [];
|
|
// A common ex-VAT value is only available because eligibility requires VAT evidence.
|
|
const value = positive.reduce((sum, row) => sum + row.totalPrice / (row.vatState === 'inclusive' ? 1.25 : 1), 0);
|
|
return [{ materials_count: positive.length, total_value: value, positive_consumption_count: positive.length,
|
|
positive_consumption_value: value,
|
|
negative_material_count: materials.negativeCounts.get(number) || 0,
|
|
unclassified_material_count: materials.unclassifiedCounts.get(number) || 0,
|
|
excluded_material_count: materials.excludedCounts.get(number) || 0 }];
|
|
}
|
|
}
|
|
|
|
module.exports = { OrdrestyringReadModel, casesRelation, debtorsRelation, hoursRelation, planningHoursRelation };
|