From 62b709317236cbc6b8de4b6e25c99b45e59e7c40 Mon Sep 17 00:00:00 2001 From: Michael SCHAL Date: Fri, 20 Mar 2026 07:31:55 +0100 Subject: [PATCH] feat: Implement a new Dockerized PostgreSQL data synchronization service with cron scheduling. --- .env.example | 6 + docker-compose.yml | 26 +++++ sync/Dockerfile | 12 ++ sync/db.js | 41 +++++++ sync/index.js | 84 ++++++++++++++ sync/package.json | 13 +++ sync/server.js | 281 +++++++++++++++++++++++++++++++++++++++++++++ sync/sync.sh | 20 ++++ sync/tables.js | 34 ++++++ sync/utils.js | 89 ++++++++++++++ 10 files changed, 606 insertions(+) create mode 100644 sync/Dockerfile create mode 100644 sync/db.js create mode 100644 sync/index.js create mode 100644 sync/package.json create mode 100644 sync/server.js create mode 100644 sync/sync.sh create mode 100644 sync/tables.js create mode 100644 sync/utils.js diff --git a/.env.example b/.env.example index 32b87e2..f0e3e37 100644 --- a/.env.example +++ b/.env.example @@ -3,3 +3,9 @@ # Token partagé entre frps (Hostinger) et frpc (magasin) # Doit être identique dans frps.toml ET frpc.toml côté magasin FRP_TOKEN=CHANGE_ME_FRP_SECRET_TOKEN + +# Credentials PostgreSQL Debian (source du sync) +# Même valeurs que dans foirfouille-api/.env.docker +SRC_DB=CHANGE_ME_DB_NAME +SRC_USER=CHANGE_ME_DB_USER +SRC_PASS=CHANGE_ME_DB_PASSWORD diff --git a/docker-compose.yml b/docker-compose.yml index 8539d43..547501d 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -14,6 +14,32 @@ services: - nginx_network restart: unless-stopped + # Sync PG Debian → PG Hostinger + # Source : frps:15432 (tunnel frp vers postgres Debian) + # Destination : postgres_db (container existant dans nginx_network) + sync: + build: ./sync + container_name: ff_sync_hostinger + environment: + SRC_HOST: frps + SRC_PORT: 15432 + SRC_DB: ${SRC_DB} + SRC_USER: ${SRC_USER} + SRC_PASS: ${SRC_PASS} + DST_HOST: postgres_db + DST_PORT: 5432 + DST_DB: post + DST_USER: post + DST_PASS: post25 + ports: + - "3002:3002" # dashboard sync accessible localement + networks: + - default # pour atteindre frps:15432 + - nginx_network # pour atteindre postgres_db + depends_on: + - frps + restart: unless-stopped + networks: nginx_network: external: true diff --git a/sync/Dockerfile b/sync/Dockerfile new file mode 100644 index 0000000..d6d24fd --- /dev/null +++ b/sync/Dockerfile @@ -0,0 +1,12 @@ +FROM node:20-alpine + +WORKDIR /app + +COPY package.json ./ +RUN npm install --omit=dev + +COPY . . + +EXPOSE 3002 + +CMD ["node", "index.js"] diff --git a/sync/db.js b/sync/db.js new file mode 100644 index 0000000..1d6f13c --- /dev/null +++ b/sync/db.js @@ -0,0 +1,41 @@ +'use strict'; +const { Pool } = require('pg'); + +let srcPool = null; +let dstPool = null; + +function getSrc() { + if (!srcPool) { + srcPool = new Pool({ + host: process.env.SRC_HOST, + port: parseInt(process.env.SRC_PORT) || 5432, + database: process.env.SRC_DB, + user: process.env.SRC_USER, + password: process.env.SRC_PASS, + max: 3, + connectionTimeoutMillis: 10000, + }); + } + return srcPool; +} + +function getDst() { + if (!dstPool) { + dstPool = new Pool({ + host: process.env.DST_HOST, + port: parseInt(process.env.DST_PORT) || 5432, + database: process.env.DST_DB, + user: process.env.DST_USER, + password: process.env.DST_PASS, + max: 5, + }); + } + return dstPool; +} + +async function closeAll() { + if (srcPool) { await srcPool.end(); srcPool = null; } + if (dstPool) { await dstPool.end(); dstPool = null; } +} + +module.exports = { getSrc, getDst, closeAll }; diff --git a/sync/index.js b/sync/index.js new file mode 100644 index 0000000..1d70abd --- /dev/null +++ b/sync/index.js @@ -0,0 +1,84 @@ +'use strict'; +const cron = require('node-cron'); +const { getSrc, getDst, closeAll } = require('./db'); +const { batchUpsert, fullRefresh, logSync } = require('./utils'); +const { TABLES } = require('./tables'); +const { + startDashboard, registerSyncFn, + setSyncRunning, isSyncRunning, + setProgress, clearProgress, +} = require('./server'); + +const forceMode = process.argv.includes('--force'); + +async function syncTable(src, dst, table, pk) { + const { rows } = await src.query(`SELECT * FROM ${table}`); + if (!rows.length) { + await logSync(dst, table, 0, 'ok'); + console.log(`[${table}] 0 lignes (table vide)`); + return 0; + } + + let count; + if (pk.length === 0) { + count = await fullRefresh(dst, table, rows); + } else { + count = await batchUpsert(dst, table, rows, pk); + } + await logSync(dst, table, count, 'ok'); + console.log(`[${table}] ${count} lignes`); + return count; +} + +async function syncAll(force = false) { + if (isSyncRunning()) { + console.log('Sync déjà en cours — ignorée.'); + return; + } + setSyncRunning(true); + clearProgress(); + const start = Date.now(); + console.log(`\n=== SYNC START ${new Date().toISOString()} ===`); + + const src = getSrc(); + const dst = getDst(); + + try { + for (let i = 0; i < TABLES.length; i++) { + const { name, pk } = TABLES[i]; + setProgress(name, i + 1); + console.log(`[${i + 1}/${TABLES.length}] Sync ${name}…`); + try { + await syncTable(src, dst, name, pk); + } catch (err) { + await logSync(dst, name, 0, 'error', err.message).catch(() => {}); + console.error(`[${name}] ERREUR: ${err.message}`); + } + } + } finally { + setSyncRunning(false); + clearProgress(); + } + + const elapsed = ((Date.now() - start) / 1000).toFixed(1); + console.log(`=== SYNC DONE en ${elapsed}s ===\n`); +} + +registerSyncFn(syncAll); + +if (forceMode) { + console.log('Mode force : sync complète en cours...'); + syncAll(true) + .then(() => closeAll()) + .then(() => process.exit(0)) + .catch(err => { + console.error('Erreur fatale sync:', err.message); + process.exit(1); + }); +} else { + startDashboard(); + console.log('Service sync Hostinger démarré — cron: 30 3 * * * (3h30 chaque nuit)'); + cron.schedule('30 3 * * *', () => { + syncAll(false).catch(err => console.error('Erreur cron sync:', err.message)); + }); +} diff --git a/sync/package.json b/sync/package.json new file mode 100644 index 0000000..4f7957b --- /dev/null +++ b/sync/package.json @@ -0,0 +1,13 @@ +{ + "name": "ff-sync-hostinger", + "version": "1.0.0", + "description": "Sync nightly PG Debian → PG Hostinger", + "main": "index.js", + "license": "ISC", + "type": "commonjs", + "dependencies": { + "express": "^5.1.0", + "node-cron": "^3.0.3", + "pg": "^8.13.3" + } +} diff --git a/sync/server.js b/sync/server.js new file mode 100644 index 0000000..0278215 --- /dev/null +++ b/sync/server.js @@ -0,0 +1,281 @@ +'use strict'; +const express = require('express'); +const { getDst } = require('./db'); +const { TABLES } = require('./tables'); + +const app = express(); +const PORT = process.env.SYNC_DASHBOARD_PORT || 3002; + +const TABLES_ORDER = TABLES.map(t => t.name); +const TABLES_TOTAL = TABLES_ORDER.length; + +const state = { + running: false, + progress: null, +}; + +const logBuffer = []; +function pushLog(msg, level = 'info') { + logBuffer.push({ ts: new Date().toISOString(), msg, level }); + if (logBuffer.length > 300) logBuffer.shift(); +} + +const _log = console.log.bind(console); +const _err = console.error.bind(console); +console.log = (...a) => { pushLog(a.join(' '), 'info'); _log(...a); }; +console.error = (...a) => { pushLog(a.join(' '), 'error'); _err(...a); }; + +function setProgress(tableName, index) { + state.progress = { + current: tableName, + index, + total: TABLES_TOTAL, + pct: Math.round((index / TABLES_TOTAL) * 100), + startedAt: state.progress?.startedAt || new Date().toISOString(), + }; +} +function clearProgress() { state.progress = null; } + +let syncAllFn = null; +function registerSyncFn(fn) { syncAllFn = fn; } +function setSyncRunning(val) { state.running = val; } +function isSyncRunning() { return state.running; } + +app.use(express.json()); + +app.get('/api/status', async (req, res) => { + try { + const dst = getDst(); + const { rows } = await dst.query( + `SELECT table_name, last_sync, rows_synced, status, error_msg + FROM sync_log ORDER BY table_name` + ); + res.json({ running: state.running, progress: state.progress, tables: rows }); + } catch (err) { + res.status(500).json({ error: err.message }); + } +}); + +app.get('/api/logs', (req, res) => { + const limit = parseInt(req.query.limit) || 150; + res.json(logBuffer.slice(-limit)); +}); + +app.get('/api/next', (req, res) => { + const now = new Date(); + const next = new Date(now); + next.setHours(3, 30, 0, 0); + if (next <= now) next.setDate(next.getDate() + 1); + res.json({ next: next.toISOString() }); +}); + +app.post('/api/sync/run', async (req, res) => { + if (state.running) return res.status(409).json({ error: 'Sync déjà en cours' }); + if (!syncAllFn) return res.status(500).json({ error: 'syncAll non enregistré' }); + res.json({ message: 'Sync lancée' }); + syncAllFn(true).catch(err => console.error('Erreur sync manuelle:', err.message)); +}); + +app.get('/', (req, res) => res.send(HTML)); + +function startDashboard() { + app.listen(PORT, '0.0.0.0', () => { + console.log(`Dashboard sync Hostinger : http://0.0.0.0:${PORT}`); + }); +} + +module.exports = { startDashboard, registerSyncFn, setSyncRunning, isSyncRunning, setProgress, clearProgress }; + +const HTML = ` + + + + +Sync Hostinger + + + +
+

⚡ Sync Hostinger

+ PG Debian → PG Hostinger + Idle + +
+
+
+ + + + +
+ +
+
+ Synchronisation en cours… + 0% +
+
+
+
+
+
+ +
+
+
Statut des tables
+ + + +
TableDernière syncLignesStatut
Chargement…
+
+
+
+ Logs + +
+
Chargement…
+
+
+
+ + +`; diff --git a/sync/sync.sh b/sync/sync.sh new file mode 100644 index 0000000..6eecfea --- /dev/null +++ b/sync/sync.sh @@ -0,0 +1,20 @@ +#!/bin/sh +# Sync PostgreSQL Debian → PostgreSQL Hostinger +# Source : frps:15432 (tunnel frp vers postgres Debian) +# Destination : postgres_db (container existant sur Hostinger) + +set -e + +echo "=== SYNC START $(date) ===" + +PGPASSWORD=$SRC_PASS pg_dump \ + -h "$SRC_HOST" -p "$SRC_PORT" \ + -U "$SRC_USER" -d "$SRC_DB" \ + --data-only --if-exists --clean \ + --no-acl --no-owner \ +| PGPASSWORD=$DST_PASS psql \ + -h "$DST_HOST" -p "$DST_PORT" \ + -U "$DST_USER" -d "$DST_DB" \ + --set ON_ERROR_STOP=0 + +echo "=== SYNC DONE $(date) ===" diff --git a/sync/tables.js b/sync/tables.js new file mode 100644 index 0000000..44f6854 --- /dev/null +++ b/sync/tables.js @@ -0,0 +1,34 @@ +'use strict'; +// Définition des tables à synchroniser PG Debian → PG Hostinger +// pk: [] = full refresh (TRUNCATE + INSERT), sinon ON CONFLICT DO UPDATE + +const TABLES = [ + // Référentiel + { name: 'nomenclature', pk: ['no_id'] }, + { name: 'gammes', pk: ['no_id'] }, + { name: 'saisons', pk: ['no_id'] }, + // Articles + { name: 'articles', pk: ['no_id'] }, + { name: 'article_infosup', pk: ['artnoid'] }, + { name: 'art_gtin', pk: ['idarticle', 'gtin'] }, + { name: 'art_gamme_saison', pk: ['artnoid', 'idgamme', 'idsaison'] }, + // Fournisseurs + { name: 'fouadr1', pk: ['code', 'sit_code'] }, + { name: 'artfou1', pk: ['no_id'] }, + { name: 'artfou2', pk: ['idartfou1'] }, + // Stock & prix + { name: 'cube_pa', pk: ['artnoid'] }, + { name: 'cube_pv', pk: ['artnoid', 'site'] }, + { name: 'cube_stock', pk: ['artnoid', 'site'] }, + // Mouvements (sans PK → full refresh) + { name: 'mvtart', pk: [] }, + { name: 'mvtreg', pk: [] }, + // Ranking + { name: 'ranking', pk: ['gencod', 'site'] }, + // Commandes + { name: 'cdefou_reception', pk: ['no_id'] }, + { name: 'cdefou_vivant', pk: ['cdefou_ligne_com_no_id', 'artfou1_no_id'] }, + { name: 'cdefou_receplig', pk: ['no_id'] }, +]; + +module.exports = { TABLES }; diff --git a/sync/utils.js b/sync/utils.js new file mode 100644 index 0000000..8f457c9 --- /dev/null +++ b/sync/utils.js @@ -0,0 +1,89 @@ +'use strict'; + +/** + * Upsert en masse dans PG destination (chunks de 500 lignes). + * Les colonnes sont déduites dynamiquement des lignes source. + */ +async function batchUpsert(pg, table, rows, pk) { + if (!rows.length) return 0; + const cols = Object.keys(rows[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.length + ? updateCols.map(c => `${c} = EXCLUDED.${c}`).join(', ') + : `${pk[0]} = EXCLUDED.${pk[0]}`; // fallback si toutes les cols sont PK + + 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) { + await pg.query(`TRUNCATE TABLE ${table} CASCADE`); + if (!rows.length) return 0; + + const cols = Object.keys(rows[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(', ')})`; + }); + await pg.query( + `INSERT INTO ${table} (${cols.join(', ')}) VALUES ${placeholders.join(', ')}`, + values + ); + total += chunk.length; + } + return total; +} + +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]); +} + +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; +} + +module.exports = { batchUpsert, fullRefresh, logSync, getLastSync };