Files
Apiflow/sync/utils.js
T
MichaelandClaude Sonnet 4.6 d565ed9f10 Initial commit : API Node.js + PostgreSQL + sync service Docker
- API Express/pg connectée à PostgreSQL (converti depuis mssql)
- 8 routes : articles, fournisseurs, commandes, stock, mouvements, performance, schema, sync
- Service sync nightly (mssql → pg, cron 3h00)
- Docker Compose : postgres + api + sync + frpc (tunnel sortant)
- Schéma PostgreSQL 17 tables dans postgres/init.sql

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-03-18 10:10:47 +01:00

92 lines
2.8 KiB
JavaScript

/**
* Upsert en masse dans PostgreSQL (chunks de 500 lignes).
* @param {import('pg').Pool} pg
* @param {string} table - nom de la table PostgreSQL (lowercase)
* @param {object[]} rows - lignes avec les valeurs à insérer
* @param {string[]} pk - colonnes formant la clé primaire
* @param {string[]} cols - toutes les colonnes à upsert
* @returns {number} total de lignes traitées
*/
async function batchUpsert(pg, table, rows, pk, cols) {
if (!rows.length) return 0;
const CHUNK = 500;
let total = 0;
for (let i = 0; i < rows.length; i += CHUNK) {
const chunk = rows.slice(i, i + CHUNK);
const values = [];
const placeholders = chunk.map((row, ri) => {
const ph = cols.map((col, ci) => {
values.push(row[col] ?? null);
return `$${ri * cols.length + ci + 1}`;
});
return `(${ph.join(', ')})`;
});
const updateCols = cols.filter(c => !pk.includes(c));
const updateSet = updateCols.map(c => `${c} = EXCLUDED.${c}`).join(', ');
const sql = `
INSERT INTO ${table} (${cols.join(', ')})
VALUES ${placeholders.join(', ')}
ON CONFLICT (${pk.join(', ')}) DO UPDATE SET ${updateSet}
`;
await pg.query(sql, values);
total += chunk.length;
}
return total;
}
/**
* Full refresh : TRUNCATE puis INSERT en masse.
*/
async function fullRefresh(pg, table, rows, cols) {
if (!rows.length) {
await pg.query(`TRUNCATE TABLE ${table}`);
return 0;
}
const CHUNK = 500;
await pg.query(`TRUNCATE TABLE ${table}`);
let total = 0;
for (let i = 0; i < rows.length; i += CHUNK) {
const chunk = rows.slice(i, i + CHUNK);
const values = [];
const placeholders = chunk.map((row, ri) => {
const ph = cols.map((col, ci) => {
values.push(row[col] ?? null);
return `$${ri * cols.length + ci + 1}`;
});
return `(${ph.join(', ')})`;
});
const sql = `INSERT INTO ${table} (${cols.join(', ')}) VALUES ${placeholders.join(', ')}`;
await pg.query(sql, values);
total += chunk.length;
}
return total;
}
/**
* Lit la date du dernier sync réussi pour une table.
*/
async function getLastSync(pg, tableName) {
const res = await pg.query(
`SELECT last_sync FROM sync_log WHERE table_name = $1 AND status = 'ok'`,
[tableName]
);
return res.rows[0]?.last_sync || null;
}
/**
* Enregistre le résultat d'un sync dans sync_log.
*/
async function logSync(pg, tableName, rowsSynced, status, errorMsg = null) {
await pg.query(`
INSERT INTO sync_log (table_name, last_sync, rows_synced, status, error_msg)
VALUES ($1, NOW(), $2, $3, $4)
ON CONFLICT (table_name) DO UPDATE SET
last_sync = NOW(),
rows_synced = $2,
status = $3,
error_msg = $4
`, [tableName, rowsSynced, status, errorMsg]);
}
module.exports = { batchUpsert, fullRefresh, getLastSync, logSync };