[verified] fix(ordrestyring): stop raw duplicate growth
This commit is contained in:
@@ -0,0 +1,47 @@
|
||||
-- Apply explicitly to ordrestyring_local before enabling canonical sync.
|
||||
-- Additive only: never bootstrap from ambiguous legacy rows.
|
||||
CREATE TABLE IF NOT EXISTS ordrestyring_cases_current (
|
||||
case_number VARCHAR(191) COLLATE utf8mb4_bin NOT NULL PRIMARY KEY,
|
||||
source_account VARCHAR(128) COLLATE utf8mb4_bin NOT NULL,
|
||||
source_payload JSON NOT NULL,
|
||||
payload_hash CHAR(64) CHARACTER SET ascii COLLATE ascii_bin NOT NULL,
|
||||
source_created_at VARCHAR(128) NULL,
|
||||
source_updated_at VARCHAR(128) NULL,
|
||||
sync_run_id CHAR(36) CHARACTER SET ascii NOT NULL,
|
||||
mirrored_at DATETIME(6) NOT NULL
|
||||
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_bin;
|
||||
|
||||
CREATE TABLE IF NOT EXISTS ordrestyring_debtors_current (
|
||||
customer_number VARCHAR(191) COLLATE utf8mb4_bin NOT NULL PRIMARY KEY,
|
||||
source_account VARCHAR(128) COLLATE utf8mb4_bin NOT NULL,
|
||||
source_payload JSON NOT NULL,
|
||||
payload_hash CHAR(64) CHARACTER SET ascii COLLATE ascii_bin NOT NULL,
|
||||
source_created_at VARCHAR(128) NULL,
|
||||
source_updated_at VARCHAR(128) NULL,
|
||||
sync_run_id CHAR(36) CHARACTER SET ascii NOT NULL,
|
||||
mirrored_at DATETIME(6) NOT NULL
|
||||
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_bin;
|
||||
|
||||
CREATE TABLE IF NOT EXISTS ordrestyring_material_semantics (
|
||||
material_line_id VARCHAR(191) COLLATE utf8mb4_bin NOT NULL PRIMARY KEY,
|
||||
source_account VARCHAR(128) COLLATE utf8mb4_bin NOT NULL,
|
||||
source_payload JSON NOT NULL,
|
||||
payload_hash CHAR(64) CHARACTER SET ascii COLLATE ascii_bin NOT NULL,
|
||||
source_created_at VARCHAR(128) NULL,
|
||||
source_updated_at VARCHAR(128) NULL,
|
||||
sync_run_id CHAR(36) CHARACTER SET ascii NOT NULL,
|
||||
mirrored_at DATETIME(6) NOT NULL,
|
||||
semantic_class ENUM('normal_consumption', 'return', 'credit', 'correction', 'unclassified') NOT NULL DEFAULT 'unclassified',
|
||||
evidence_json JSON NOT NULL
|
||||
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_bin;
|
||||
|
||||
CREATE TABLE IF NOT EXISTS ordrestyring_projection_state (
|
||||
projection_kind VARCHAR(32) NOT NULL PRIMARY KEY,
|
||||
source_account VARCHAR(128) COLLATE utf8mb4_bin NULL,
|
||||
sync_run_id CHAR(36) CHARACTER SET ascii NULL,
|
||||
row_count INT UNSIGNED NOT NULL DEFAULT 0,
|
||||
mirrored_at DATETIME(6) NULL
|
||||
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_bin;
|
||||
|
||||
INSERT IGNORE INTO ordrestyring_projection_state (projection_kind)
|
||||
VALUES ('cases'), ('debtors'), ('materials');
|
||||
@@ -0,0 +1,42 @@
|
||||
process.env.ORDRESTYRING_API_TOKEN = 'test-token';
|
||||
process.env.DB_PASSWORD = 'test-password';
|
||||
process.env.ORDRESTYRING_SOURCE_ACCOUNT = 'fixture-account';
|
||||
jest.mock('axios', () => ({ get: jest.fn() }));
|
||||
jest.mock('mysql2/promise', () => ({ createConnection: jest.fn() }));
|
||||
jest.mock('../graphqlClient', () => ({}));
|
||||
jest.mock('../../utils/logger', () => ({ info: jest.fn(), error: jest.fn(), warn: jest.fn() }));
|
||||
const axios = require('axios');
|
||||
const mysql = require('mysql2/promise');
|
||||
const Service = require('../ordrestyringSyncService');
|
||||
const logger = require('../../utils/logger');
|
||||
|
||||
beforeEach(() => jest.clearAllMocks());
|
||||
test.each([
|
||||
['syncCases', 'case_number', 1934, 'ordrestyring_cases_current'],
|
||||
['syncDebtors', 'customer_number', 1564, 'ordrestyring_debtors_current']
|
||||
])('%s writes only canonical rows after the full source fetch', async (method, key, count, table) => {
|
||||
const db = { execute: jest.fn(async sql => /projection_state/.test(sql) && /SELECT/.test(sql)
|
||||
? [[{ source_account: null }]] : [[]]), beginTransaction: jest.fn(), commit: jest.fn(), rollback: jest.fn(),
|
||||
end: jest.fn().mockRejectedValue(new Error('private cleanup')) };
|
||||
axios.get.mockResolvedValue({ status: 200, data: Array.from({ length: count }, (_, i) => ({ [key]: String(i + 1) })) });
|
||||
mysql.createConnection.mockResolvedValue(db);
|
||||
await expect(new Service()[method]()).resolves.toEqual({ hasChanges: true, changes: count });
|
||||
const sql = db.execute.mock.calls.map(([query]) => query).join('\n');
|
||||
expect(sql).toContain(`INSERT INTO ${table}`);
|
||||
expect(sql).not.toMatch(/(?:FROM|INTO|UPDATE)\s+(?:cases|debtors)\b/);
|
||||
expect(axios.get.mock.invocationCallOrder[0]).toBeLessThan(mysql.createConnection.mock.invocationCallOrder[0]);
|
||||
expect(db.commit).toHaveBeenCalledTimes(1);
|
||||
expect(db.end).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
||||
test.each(['syncCases', 'syncDebtors'])('%s fails closed on malformed, empty, partial or failed responses without leaking errors', async method => {
|
||||
for (const response of [{ data: {} }, { data: [] }, { status: 206, data: [] },
|
||||
{ status: 200, headers: { link: '<private>; rel="next"' }, data: [] }]) {
|
||||
axios.get.mockResolvedValueOnce(response);
|
||||
await expect(new Service()[method]()).rejects.toThrow();
|
||||
}
|
||||
axios.get.mockRejectedValueOnce(new Error('private upstream payload/token'));
|
||||
await expect(new Service()[method]()).rejects.toThrow(/^Ordrestyring source fetch failed$/);
|
||||
expect(mysql.createConnection).not.toHaveBeenCalled();
|
||||
expect(JSON.stringify(logger.error.mock.calls)).not.toContain('private');
|
||||
});
|
||||
@@ -0,0 +1,167 @@
|
||||
const { projectCurrent, stablePayloadHash } = require('../ordrestyringCurrentProjection');
|
||||
const { classifyMaterial } = require('../ordrestyringMaterialSemantics');
|
||||
|
||||
// Transactional database double: committed state is distinct from in-flight state.
|
||||
function database() {
|
||||
let rows = new Map();
|
||||
let pending;
|
||||
let account = null;
|
||||
let pendingAccount;
|
||||
const db = {
|
||||
beginTransaction: jest.fn(async () => { pending = new Map(rows); pendingAccount = account; }),
|
||||
commit: jest.fn(async () => { rows = pending; account = pendingAccount; }),
|
||||
rollback: jest.fn(async () => { pending = null; }),
|
||||
execute: jest.fn(async (sql, args) => {
|
||||
if (/SELECT.*source_account.*FROM ordrestyring_projection_state/s.test(sql)) return [[{ source_account: account }]];
|
||||
if (/^SELECT/s.test(sql)) return [[...rows.values()]];
|
||||
if (/INSERT INTO ordrestyring_.*_current/.test(sql)) {
|
||||
pending.set(args[0], { identity: args[0], payload_hash: args[3], source_account: args[1] });
|
||||
}
|
||||
if (/UPDATE ordrestyring_projection_state/.test(sql)) pendingAccount = args[0];
|
||||
return [{ affectedRows: 1 }];
|
||||
}),
|
||||
rows: () => rows
|
||||
};
|
||||
return db;
|
||||
}
|
||||
const options = (connection, rows, extra = {}) => ({ connection, kind: 'cases', rows,
|
||||
sourceAccount: 'fixture-account', minimumRows: 1, ...extra });
|
||||
|
||||
test('stable normalized hashes ignore object key order but retain array order and source values', () => {
|
||||
expect(stablePayloadHash({ b: { y: 2, x: 1 }, a: [1, 2] }))
|
||||
.toBe(stablePayloadHash({ a: [1, 2], b: { x: 1, y: 2 } }));
|
||||
expect(stablePayloadHash({ a: [1, 2] })).not.toBe(stablePayloadHash({ a: [2, 1] }));
|
||||
expect(stablePayloadHash({ a: null })).not.toBe(stablePayloadHash({}));
|
||||
});
|
||||
|
||||
test.each([undefined, '', ' ', {}, [], true, 1.2, NaN])('rejects unsupported identity %p before writes', async identity => {
|
||||
const db = database();
|
||||
await expect(projectCurrent(options(db, [{ case_number: identity }]))).rejects.toThrow(/identity/);
|
||||
expect(db.beginTransaction).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
test.each(['cases', 'debtors'])('%s rejects normalized duplicates before writes', async kind => {
|
||||
const key = kind === 'cases' ? 'case_number' : 'customer_number';
|
||||
const db = database();
|
||||
await expect(projectCurrent(options(db, [{ [key]: ' 12 ' }, { [key]: 12 }], { kind })))
|
||||
.rejects.toThrow(/duplicate/i);
|
||||
expect(db.execute).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
test('second run is idempotent, preserves full payload, and never reads or writes legacy tables', async () => {
|
||||
const db = database();
|
||||
const row = { case_number: ' 12 ', contact: { name: 'private fixture' }, created_at: 123, updated_at: 456 };
|
||||
expect(await projectCurrent(options(db, [row]))).toEqual({ hasChanges: true, changes: 1 });
|
||||
expect(await projectCurrent(options(db, [{ ...row, case_number: 12 }]))).toEqual({ hasChanges: false, changes: 0 });
|
||||
expect(db.rows().size).toBe(1);
|
||||
const writes = db.execute.mock.calls.filter(([sql]) => /INSERT INTO/.test(sql));
|
||||
expect(JSON.parse(writes[0][1][2])).toEqual({ ...row, case_number: '12' });
|
||||
expect(writes[0][1]).toEqual(expect.arrayContaining(['123', '456']));
|
||||
expect(db.execute.mock.calls.map(([sql]) => sql).join('\n')).not.toMatch(/(?:FROM|INTO|UPDATE)\s+(?:cases|debtors)\b|\b(?:CREATE|ALTER|DROP)\b/i);
|
||||
});
|
||||
|
||||
test('changed payload changes hash and one business record', async () => {
|
||||
const db = database();
|
||||
await projectCurrent(options(db, [{ case_number: '12', status: 1 }]));
|
||||
const hash = db.rows().get('12').payload_hash;
|
||||
expect(await projectCurrent(options(db, [{ case_number: '12', status: 2 }]))).toEqual({ hasChanges: true, changes: 1 });
|
||||
expect(db.rows().get('12').payload_hash).not.toBe(hash);
|
||||
});
|
||||
|
||||
test('failed second write rolls back and sanitizes database errors', async () => {
|
||||
const db = database();
|
||||
await projectCurrent(options(db, [{ case_number: '12' }]));
|
||||
const execute = db.execute.getMockImplementation();
|
||||
db.execute.mockImplementation(async (sql, args) => {
|
||||
if (/INSERT INTO/.test(sql) && args[0] === '13') throw new Error('private SQL and customer data');
|
||||
return execute(sql, args);
|
||||
});
|
||||
await expect(projectCurrent(options(db, [{ case_number: '12', status: 2 }, { case_number: '13' }])))
|
||||
.rejects.toThrow(/^Ordrestyring projection write failed$/);
|
||||
expect(db.rollback).toHaveBeenCalledTimes(1);
|
||||
expect(db.rows().size).toBe(1);
|
||||
expect(db.commit).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
||||
test('rejects empty, bootstrap partial, and subsequent missing identities without committing', async () => {
|
||||
const db = database();
|
||||
await expect(projectCurrent(options(db, []))).rejects.toThrow(/empty/);
|
||||
await expect(projectCurrent(options(db, [{ case_number: '12' }], { minimumRows: 2 }))).rejects.toThrow(/partial/);
|
||||
await projectCurrent(options(db, [{ case_number: '12' }, { case_number: '13' }]));
|
||||
await expect(projectCurrent(options(db, [{ case_number: '12' }, { case_number: '14' }]))).rejects.toThrow(/partial/);
|
||||
expect(db.commit).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
||||
test('rejects source account changes even with disjoint identities', async () => {
|
||||
const db = database();
|
||||
await projectCurrent(options(db, [{ case_number: '12' }]));
|
||||
await expect(projectCurrent(options(db, [{ case_number: '99' }], { sourceAccount: 'another' })))
|
||||
.rejects.toThrow(/account/);
|
||||
});
|
||||
|
||||
test('missing schema fails explicitly without request-time DDL', async () => {
|
||||
const db = database();
|
||||
db.execute.mockRejectedValue({ code: 'ER_NO_SUCH_TABLE', message: 'private details' });
|
||||
await expect(projectCurrent(options(db, [{ case_number: '12' }]))).rejects.toThrow(/migration required/i);
|
||||
expect(db.rollback).toHaveBeenCalled();
|
||||
});
|
||||
|
||||
const contract = { reference: 'approved-test-contract-v1', units: ['m'], priceScale: 'major', vatState: 'exclusive' };
|
||||
const material = { quantity: 2, cost_price: '1.25', sales_price: 2, unit: 'm' };
|
||||
test('positive materials require explicit supported source/unit and monetary evidence', () => {
|
||||
expect(classifyMaterial(material).semanticClass).toBe('unclassified');
|
||||
expect(classifyMaterial(material, contract).semanticClass).toBe('normal_consumption');
|
||||
expect(classifyMaterial({ ...material, unit: null }, contract).semanticClass).toBe('unclassified');
|
||||
});
|
||||
test.each([-2, 0, null, undefined, '', '2oops', true])('quantity %p is excluded and preserved', quantity => {
|
||||
const result = classifyMaterial({ ...material, quantity }, contract);
|
||||
expect(result.semanticClass).toBe('unclassified');
|
||||
expect(result.payload.quantity).toBe(quantity);
|
||||
});
|
||||
test.each([-1, null, '', '1,25', '2oops', Infinity, true])('invalid or ambiguous price %p is excluded', price => {
|
||||
expect(classifyMaterial({ ...material, cost_price: price }, contract).semanticClass).toBe('unclassified');
|
||||
});
|
||||
|
||||
test('material projection persists signed source payload, explicit class and shared run provenance', async () => {
|
||||
const db = database();
|
||||
await projectCurrent(options(db, [{ material_line_id: '7', quantity: -2, cost_price: -3, unit: null }], { kind: 'materials' }));
|
||||
const write = db.execute.mock.calls.find(([sql]) => /INSERT INTO ordrestyring_material_semantics/.test(sql));
|
||||
expect(write[0]).toContain('semantic_class, evidence_json');
|
||||
expect(JSON.parse(write[1][2])).toEqual({ material_line_id: '7', quantity: -2, cost_price: -3, unit: null });
|
||||
expect(write[1][7]).toBe('unclassified');
|
||||
const state = db.execute.mock.calls.find(([sql]) => /UPDATE ordrestyring_projection_state/.test(sql));
|
||||
expect(state[1][1]).toBe(write[1][6]);
|
||||
expect(db.execute.mock.calls[0][0]).toContain('FOR UPDATE');
|
||||
});
|
||||
|
||||
test.each([
|
||||
{ sourceAccount: '' }, { sourceAccount: 'private account name' }, { kind: 'cases; DROP TABLE cases' },
|
||||
{ minimumRows: 0 }, { rows: {} }, { rows: [{ case_number: '12', contact: undefined }] },
|
||||
{ rows: [{ case_number: '12', updated_at: {} }] }
|
||||
])('invalid configuration/payload fails before opening a transaction: %p', async extra => {
|
||||
const db = database();
|
||||
await expect(projectCurrent(options(db, [{ case_number: '12' }], extra))).rejects.toThrow(/Ordrestyring projection/);
|
||||
expect(db.beginTransaction).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
test('missing seeded lock row is migration-required and rolls back', async () => {
|
||||
const db = database();
|
||||
db.execute.mockResolvedValue([[]]);
|
||||
await expect(projectCurrent(options(db, [{ case_number: '12' }]))).rejects.toThrow(/migration required/);
|
||||
expect(db.rollback).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
||||
test('migration defines additive, transactional canonical schemas and persistent locks', () => {
|
||||
const fs = require('fs');
|
||||
const path = require('path');
|
||||
const sql = fs.readFileSync(path.join(__dirname, '../../../migrations/20261004_ordrestyring_semantic_current.sql'), 'utf8');
|
||||
for (const [table, key] of [['ordrestyring_cases_current', 'case_number'],
|
||||
['ordrestyring_debtors_current', 'customer_number'], ['ordrestyring_material_semantics', 'material_line_id']]) {
|
||||
expect(sql).toContain(`CREATE TABLE IF NOT EXISTS ${table}`);
|
||||
expect(sql).toContain(`${key} VARCHAR(191) COLLATE utf8mb4_bin NOT NULL PRIMARY KEY`);
|
||||
}
|
||||
expect(sql.match(/ENGINE=InnoDB/g)).toHaveLength(4);
|
||||
expect(sql).toContain("DEFAULT 'unclassified'");
|
||||
expect(sql).toContain("VALUES ('cases'), ('debtors'), ('materials')");
|
||||
expect(sql).not.toMatch(/\b(?:DROP|TRUNCATE|DELETE|ALTER|AUTO_INCREMENT)\b/i);
|
||||
});
|
||||
@@ -138,7 +138,8 @@ describe('OrdrestyringSyncService.getSyncStatus', () => {
|
||||
}
|
||||
});
|
||||
|
||||
test.each(['syncCases', 'syncUsers', 'syncDebtors'])(
|
||||
// Canonical case/debtor cleanup is covered with nonempty source fixtures separately.
|
||||
test.each(['syncUsers'])(
|
||||
'%s contains connection-close failures after successful core work',
|
||||
async method => {
|
||||
axios.get.mockResolvedValueOnce({ data: [] });
|
||||
|
||||
@@ -0,0 +1,117 @@
|
||||
const { createHash, randomUUID } = require('crypto');
|
||||
const { classifyMaterial } = require('./ordrestyringMaterialSemantics');
|
||||
|
||||
const definitions = {
|
||||
cases: { table: 'ordrestyring_cases_current', key: 'case_number', minimum: 1934 },
|
||||
debtors: { table: 'ordrestyring_debtors_current', key: 'customer_number', minimum: 1564 },
|
||||
materials: { table: 'ordrestyring_material_semantics', key: 'material_line_id', minimum: 1 }
|
||||
};
|
||||
class ProjectionError extends Error {}
|
||||
const fail = message => { throw new ProjectionError(`Ordrestyring projection ${message}`); };
|
||||
|
||||
// JSON-only normalization rejects unsupported values instead of silently dropping
|
||||
// them. Sorting recursively makes hashes independent of API object key order.
|
||||
function stableJson(value) {
|
||||
if (value === null || typeof value === 'string' || typeof value === 'boolean') return JSON.stringify(value);
|
||||
if (typeof value === 'number' && Number.isFinite(value)) return JSON.stringify(value);
|
||||
if (Array.isArray(value)) return `[${value.map(stableJson).join(',')}]`;
|
||||
if (value && Object.getPrototypeOf(value) === Object.prototype) {
|
||||
return `{${Object.keys(value).sort().map(key => `${JSON.stringify(key)}:${stableJson(value[key])}`).join(',')}}`;
|
||||
}
|
||||
fail('invalid payload');
|
||||
}
|
||||
function stablePayloadHash(value) {
|
||||
return createHash('sha256').update(stableJson(value)).digest('hex');
|
||||
}
|
||||
function normalizeIdentity(value) {
|
||||
if (typeof value !== 'string' && !(typeof value === 'number' && Number.isSafeInteger(value) && value >= 0)) {
|
||||
fail('invalid identity');
|
||||
}
|
||||
const normalized = String(value).trim();
|
||||
if (!normalized || normalized.length > 191 || /[\x00-\x1f\x7f]/.test(normalized)) fail('invalid identity');
|
||||
return normalized;
|
||||
}
|
||||
function sourceTimestamp(value) {
|
||||
if (value === undefined || value === null) return null;
|
||||
if ((typeof value !== 'string' && !(typeof value === 'number' && Number.isSafeInteger(value)))
|
||||
|| String(value).length > 128) fail('invalid source timestamp');
|
||||
// Keep source representation and scale; do not invent date/unit semantics.
|
||||
return String(value);
|
||||
}
|
||||
function prepareCurrent({ kind, rows, sourceAccount, minimumRows, materialContract }) {
|
||||
const definition = definitions[kind];
|
||||
if (!definition) fail('invalid kind');
|
||||
if (typeof sourceAccount !== 'string' || !/^[a-zA-Z0-9._:-]{1,128}$/.test(sourceAccount)) fail('source account required');
|
||||
if (!Array.isArray(rows)) fail('invalid response');
|
||||
if (!rows.length) fail('empty response');
|
||||
const minimum = minimumRows === undefined ? definition.minimum : minimumRows;
|
||||
if (!Number.isSafeInteger(minimum) || minimum < 1) fail('invalid minimum');
|
||||
const keys = new Set();
|
||||
const prepared = rows.map(row => {
|
||||
if (!row || Object.getPrototypeOf(row) !== Object.prototype) fail('invalid identity');
|
||||
const identity = normalizeIdentity(row[definition.key]);
|
||||
if (keys.has(identity)) fail('duplicate identity');
|
||||
keys.add(identity);
|
||||
const payload = { ...row, [definition.key]: identity };
|
||||
const classification = kind === 'materials' ? classifyMaterial(payload, materialContract) : null;
|
||||
const payloadJson = stableJson(payload);
|
||||
return { identity, payloadJson,
|
||||
hash: stablePayloadHash(classification ? { payload, semanticClass: classification.semanticClass, evidence: classification.evidence } : payload),
|
||||
created: sourceTimestamp(row.created_at), updated: sourceTimestamp(row.updated_at), classification };
|
||||
});
|
||||
if (prepared.length < minimum) fail('partial response');
|
||||
return { definition, prepared, keys };
|
||||
}
|
||||
|
||||
async function projectCurrent(options) {
|
||||
const { connection, sourceAccount, kind } = options;
|
||||
const { definition, prepared, keys } = prepareCurrent(options);
|
||||
const { table, key } = definition;
|
||||
let started = false;
|
||||
try {
|
||||
await connection.beginTransaction();
|
||||
started = true;
|
||||
// Seeded by migration. Serializes writers even on the first empty mirror.
|
||||
const [state] = await connection.execute(
|
||||
'SELECT source_account FROM ordrestyring_projection_state WHERE projection_kind = ? FOR UPDATE', [kind]);
|
||||
if (state.length !== 1) fail('migration required');
|
||||
if (state[0].source_account !== null && state[0].source_account !== sourceAccount) fail('source account mismatch');
|
||||
const [current] = await connection.execute(`SELECT ${key} AS identity, payload_hash, source_account FROM ${table} FOR UPDATE`);
|
||||
if (current.some(row => row.source_account !== sourceAccount)) fail('source account mismatch');
|
||||
// Conservative policy: removals require an explicitly reviewed future workflow.
|
||||
// Count alone would miss a truncated response padded with newly created rows.
|
||||
if (current.some(row => !keys.has(row.identity))) fail('partial response');
|
||||
const hashes = new Map(current.map(row => [row.identity, row.payload_hash]));
|
||||
const runId = randomUUID();
|
||||
let changes = 0;
|
||||
for (const row of prepared) {
|
||||
if (hashes.get(row.identity) !== row.hash) changes += 1;
|
||||
const semanticColumns = row.classification ? ', semantic_class, evidence_json' : '';
|
||||
const semanticValues = row.classification ? ', ?, ?' : '';
|
||||
const semanticUpdates = row.classification
|
||||
? ', semantic_class = VALUES(semantic_class), evidence_json = VALUES(evidence_json)' : '';
|
||||
await connection.execute(`INSERT INTO ${table}
|
||||
(${key}, source_account, source_payload, payload_hash, source_created_at, source_updated_at, sync_run_id, mirrored_at${semanticColumns})
|
||||
VALUES (?, ?, ?, ?, ?, ?, ?, CURRENT_TIMESTAMP(6)${semanticValues})
|
||||
ON DUPLICATE KEY UPDATE source_payload = VALUES(source_payload), payload_hash = VALUES(payload_hash),
|
||||
source_created_at = VALUES(source_created_at), source_updated_at = VALUES(source_updated_at),
|
||||
sync_run_id = VALUES(sync_run_id), mirrored_at = VALUES(mirrored_at)${semanticUpdates}`,
|
||||
[row.identity, sourceAccount, row.payloadJson, row.hash, row.created, row.updated, runId,
|
||||
...(row.classification ? [row.classification.semanticClass, stableJson(row.classification.evidence)] : [])]);
|
||||
}
|
||||
await connection.execute(`UPDATE ordrestyring_projection_state
|
||||
SET source_account = ?, sync_run_id = ?, row_count = ?, mirrored_at = CURRENT_TIMESTAMP(6)
|
||||
WHERE projection_kind = ?`, [sourceAccount, runId, prepared.length, kind]);
|
||||
await connection.commit();
|
||||
started = false;
|
||||
return { hasChanges: changes > 0, changes };
|
||||
} catch (error) {
|
||||
if (started) {
|
||||
try { await connection.rollback(); } catch (_) { /* Never expose driver payloads. */ }
|
||||
}
|
||||
if (error instanceof ProjectionError) throw error;
|
||||
if (['ER_NO_SUCH_TABLE', 'ER_BAD_FIELD_ERROR'].includes(error?.code)) fail('migration required');
|
||||
fail('write failed');
|
||||
}
|
||||
}
|
||||
module.exports = { projectCurrent, prepareCurrent, stablePayloadHash };
|
||||
@@ -0,0 +1,30 @@
|
||||
// No supported contract is assumed for today's REST material feed. A caller must
|
||||
// supply a reviewed source contract; evidence must never come from row text/sign.
|
||||
function numeric(value) {
|
||||
if (typeof value === 'number') return Number.isFinite(value) ? value : null;
|
||||
if (typeof value !== 'string' || !/^-?\d+(?:\.\d+)?$/.test(value)) return null;
|
||||
const number = Number(value);
|
||||
return Number.isFinite(number) ? number : null;
|
||||
}
|
||||
|
||||
function classifyMaterial(payload, supportedContract = null) {
|
||||
const quantity = numeric(payload?.quantity);
|
||||
const cost = numeric(payload?.cost_price);
|
||||
const sales = numeric(payload?.sales_price);
|
||||
const supported = typeof supportedContract?.reference === 'string'
|
||||
&& supportedContract.reference.trim().length > 0
|
||||
&& Array.isArray(supportedContract.units)
|
||||
&& typeof payload?.unit === 'string' && payload.unit.trim().length > 0
|
||||
&& supportedContract.units.includes(payload.unit)
|
||||
&& ['major', 'minor'].includes(supportedContract.priceScale)
|
||||
&& ['inclusive', 'exclusive'].includes(supportedContract.vatState);
|
||||
return {
|
||||
semanticClass: supported && quantity !== null && quantity > 0
|
||||
&& cost !== null && cost >= 0 && sales !== null && sales >= 0
|
||||
? 'normal_consumption' : 'unclassified',
|
||||
payload,
|
||||
evidence: supported ? supportedContract : null
|
||||
};
|
||||
}
|
||||
|
||||
module.exports = { classifyMaterial };
|
||||
@@ -1,4 +1,5 @@
|
||||
const logger = require('../utils/logger');
|
||||
const { projectCurrent, prepareCurrent } = require('./ordrestyringCurrentProjection');
|
||||
const axios = require('axios');
|
||||
const mysql = require('mysql2/promise');
|
||||
const graphqlClient = require('./graphqlClient');
|
||||
@@ -235,57 +236,40 @@ class OrdrestyringSyncService {
|
||||
* Sync cases (sager) from API
|
||||
*/
|
||||
async syncCases() {
|
||||
let connection;
|
||||
return this.syncCanonicalCurrent('cases');
|
||||
}
|
||||
|
||||
async syncCanonicalCurrent(kind) {
|
||||
let response;
|
||||
try {
|
||||
const response = await axios.get(`${this.apiBaseUrl}/cases`, {
|
||||
auth: {
|
||||
username: this.apiToken,
|
||||
password: 'x'
|
||||
},
|
||||
headers: {
|
||||
'Accept': 'application/json'
|
||||
},
|
||||
response = await axios.get(`${this.apiBaseUrl}/${kind}`, {
|
||||
auth: { username: this.apiToken, password: 'x' },
|
||||
headers: { Accept: 'application/json' },
|
||||
timeout: 30000
|
||||
});
|
||||
|
||||
const apiCases = response.data;
|
||||
if (!Array.isArray(apiCases)) {
|
||||
throw new Error('Invalid API response for cases');
|
||||
} catch (_) {
|
||||
throw new Error('Ordrestyring source fetch failed');
|
||||
}
|
||||
// The evidenced REST contract is one complete array, not a paginated feed.
|
||||
// Fail closed if the server starts advertising pagination/partial content.
|
||||
const headers = response.headers || {};
|
||||
if (response.status !== 200 || headers['content-range']
|
||||
|| /rel\s*=\s*["']?next/i.test(headers.link || '')
|
||||
|| (headers['x-total-count'] !== undefined && Number(headers['x-total-count']) !== response.data?.length)) {
|
||||
throw new Error('Ordrestyring source response incomplete');
|
||||
}
|
||||
const options = { kind, rows: response.data, sourceAccount: process.env.ORDRESTYRING_SOURCE_ACCOUNT };
|
||||
prepareCurrent(options); // Full response validation before opening a database connection.
|
||||
let connection;
|
||||
try {
|
||||
try {
|
||||
connection = await mysql.createConnection(this.dbConfig);
|
||||
} catch (_) {
|
||||
throw new Error('Ordrestyring projection connection failed');
|
||||
}
|
||||
|
||||
connection = await mysql.createConnection(this.dbConfig);
|
||||
|
||||
// Get current cases from local database
|
||||
const [localCases] = await connection.execute(
|
||||
'SELECT id, case_number, creation_date, status FROM cases'
|
||||
);
|
||||
|
||||
// Create lookup maps
|
||||
const localCaseMap = new Map(localCases.map(c => [c.id, c]));
|
||||
let changes = 0;
|
||||
|
||||
// Process each API case
|
||||
for (const apiCase of apiCases) {
|
||||
const localCase = localCaseMap.get(apiCase.id);
|
||||
|
||||
if (!localCase) {
|
||||
// New case - insert
|
||||
await this.insertCase(connection, apiCase);
|
||||
changes++;
|
||||
} else if (this.caseNeedsUpdate(localCase, apiCase)) {
|
||||
// Existing case - update if changed
|
||||
await this.updateCase(connection, apiCase);
|
||||
changes++;
|
||||
}
|
||||
}
|
||||
|
||||
return { hasChanges: changes > 0, changes };
|
||||
|
||||
} catch (error) {
|
||||
logger.error('Error syncing cases:', error);
|
||||
throw error;
|
||||
return await projectCurrent({ ...options, connection });
|
||||
} finally {
|
||||
await this.closeConnectionSafely(connection, 'Cases sync');
|
||||
await this.closeConnectionSafely(connection, 'Canonical sync');
|
||||
}
|
||||
}
|
||||
|
||||
@@ -623,60 +607,7 @@ class OrdrestyringSyncService {
|
||||
* Sync debtors (kunder) from API
|
||||
*/
|
||||
async syncDebtors() {
|
||||
let connection;
|
||||
try {
|
||||
const response = await axios.get(`${this.apiBaseUrl}/debtors`, {
|
||||
auth: {
|
||||
username: this.apiToken,
|
||||
password: 'x'
|
||||
},
|
||||
headers: {
|
||||
'Accept': 'application/json'
|
||||
},
|
||||
timeout: 30000
|
||||
});
|
||||
|
||||
const apiDebtors = response.data;
|
||||
if (!Array.isArray(apiDebtors)) {
|
||||
return { hasChanges: false, changes: 0 };
|
||||
}
|
||||
|
||||
connection = await mysql.createConnection(this.dbConfig);
|
||||
|
||||
const [localDebtors] = await connection.execute(
|
||||
'SELECT customer_number, customer_name, customer_telephone FROM debtors'
|
||||
);
|
||||
|
||||
const localDebtorMap = new Map(localDebtors.map(d => [d.customer_number, d]));
|
||||
let changes = 0;
|
||||
|
||||
for (const apiDebtor of apiDebtors) {
|
||||
const localDebtor = localDebtorMap.get(apiDebtor.customer_number);
|
||||
|
||||
if (!localDebtor) {
|
||||
await this.insertDebtor(connection, apiDebtor);
|
||||
changes++;
|
||||
} else if (this.debtorNeedsUpdate(localDebtor, apiDebtor)) {
|
||||
await this.updateDebtor(connection, apiDebtor);
|
||||
changes++;
|
||||
}
|
||||
}
|
||||
|
||||
return { hasChanges: changes > 0, changes };
|
||||
|
||||
} catch (error) {
|
||||
logger.error('Error syncing debtors:', error);
|
||||
throw error;
|
||||
} finally {
|
||||
await this.closeConnectionSafely(connection, 'Debtors sync');
|
||||
}
|
||||
}
|
||||
|
||||
// Helper methods for checking if updates are needed
|
||||
caseNeedsUpdate(local, api) {
|
||||
return local.status !== api.status ||
|
||||
local.case_number !== api.case_number ||
|
||||
Math.abs(local.creation_date - api.creation_date) > 1;
|
||||
return this.syncCanonicalCurrent('debtors');
|
||||
}
|
||||
|
||||
userNeedsUpdate(local, api) {
|
||||
@@ -697,64 +628,6 @@ class OrdrestyringSyncService {
|
||||
local.case_id !== api.case_id;
|
||||
}
|
||||
|
||||
debtorNeedsUpdate(local, api) {
|
||||
return local.customer_name !== api.customer_name ||
|
||||
local.customer_telephone !== api.customer_telephone;
|
||||
}
|
||||
|
||||
// Database insert/update methods
|
||||
async insertCase(connection, caseData) {
|
||||
const query = `
|
||||
INSERT INTO cases (
|
||||
id, case_number, yourref, description, creation_date, status,
|
||||
main_technician, additional_technicians, contact, delivery_address,
|
||||
remarks, case_type, customer_number
|
||||
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
||||
ON DUPLICATE KEY UPDATE
|
||||
case_number = VALUES(case_number),
|
||||
status = VALUES(status),
|
||||
description = VALUES(description)
|
||||
`;
|
||||
|
||||
// Helper function to convert values properly
|
||||
const convertValue = (val) => {
|
||||
if (val === null || val === undefined) return null;
|
||||
if (typeof val === 'object') return JSON.stringify(val);
|
||||
return val;
|
||||
};
|
||||
|
||||
const values = [
|
||||
convertValue(caseData.id),
|
||||
convertValue(caseData.case_number),
|
||||
convertValue(caseData.yourref),
|
||||
convertValue(caseData.description),
|
||||
convertValue(caseData.creation_date),
|
||||
convertValue(caseData.status),
|
||||
convertValue(caseData.main_technician),
|
||||
convertValue(caseData.additional_technicians),
|
||||
convertValue(caseData.contact),
|
||||
convertValue(caseData.delivery_address),
|
||||
convertValue(caseData.remarks),
|
||||
convertValue(caseData.case_type),
|
||||
convertValue(caseData.customer_number)
|
||||
];
|
||||
|
||||
try {
|
||||
await connection.execute(query, values);
|
||||
} catch (err) {
|
||||
logger.error('Error inserting case:', {
|
||||
error: err.message,
|
||||
caseId: caseData.id,
|
||||
caseNumber: caseData.case_number
|
||||
});
|
||||
throw err;
|
||||
}
|
||||
}
|
||||
|
||||
async updateCase(connection, caseData) {
|
||||
await this.insertCase(connection, caseData); // Uses ON DUPLICATE KEY UPDATE
|
||||
}
|
||||
|
||||
async insertUser(connection, userData) {
|
||||
const query = `
|
||||
INSERT INTO users (
|
||||
@@ -866,41 +739,7 @@ class OrdrestyringSyncService {
|
||||
await this.insertMaterial(connection, materialData); // Uses ON DUPLICATE KEY UPDATE
|
||||
}
|
||||
|
||||
async insertDebtor(connection, debtorData) {
|
||||
const query = `
|
||||
INSERT INTO debtors (
|
||||
customer_number, customer_name, customer_address, customer_telephone,
|
||||
customer_email, created_at, updated_at
|
||||
) VALUES (?, ?, ?, ?, ?, ?, ?)
|
||||
ON DUPLICATE KEY UPDATE
|
||||
customer_name = VALUES(customer_name),
|
||||
customer_address = VALUES(customer_address),
|
||||
customer_telephone = VALUES(customer_telephone),
|
||||
customer_email = VALUES(customer_email),
|
||||
updated_at = VALUES(updated_at)
|
||||
`;
|
||||
|
||||
// Helper function to convert values properly
|
||||
const convertValue = (val) => {
|
||||
if (val === null || val === undefined) return null;
|
||||
if (typeof val === 'object') return JSON.stringify(val);
|
||||
return val;
|
||||
};
|
||||
|
||||
await connection.execute(query, [
|
||||
convertValue(debtorData.customer_number),
|
||||
convertValue(debtorData.customer_name),
|
||||
convertValue(debtorData.customer_address),
|
||||
convertValue(debtorData.customer_telephone),
|
||||
convertValue(debtorData.customer_email),
|
||||
convertValue(debtorData.created_at),
|
||||
convertValue(debtorData.updated_at)
|
||||
]);
|
||||
}
|
||||
|
||||
async updateDebtor(connection, debtorData) {
|
||||
await this.insertDebtor(connection, debtorData); // Uses ON DUPLICATE KEY UPDATE
|
||||
}
|
||||
|
||||
parseOrdrestyringDate(value) {
|
||||
if (!value && value !== 0) {
|
||||
|
||||
Reference in New Issue
Block a user