mirror of
https://github.com/R0m1k3/Apiflow-gateway.git
synced 2026-10-11 17:27:08 +02:00
feat: Implement a new Dockerized PostgreSQL data synchronization service with cron scheduling.
This commit is contained in:
1 parent
420da6a409
commit
62b7093172
10 files changed
+606
No files matched your search
@@ -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
|
||||
@@ -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
|
||||
|
||||
@@ -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"]
|
||||
+41
@@ -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 };
|
||||
@@ -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));
|
||||
});
|
||||
}
|
||||
@@ -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"
|
||||
}
|
||||
}
|
||||
+281
@@ -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 = `<!DOCTYPE html>
|
||||
<html lang="fr">
|
||||
<head>
|
||||
<meta charset="UTF-8">
|
||||
<meta name="viewport" content="width=device-width,initial-scale=1">
|
||||
<title>Sync Hostinger</title>
|
||||
<style>
|
||||
*{box-sizing:border-box;margin:0;padding:0}
|
||||
body{font-family:system-ui,sans-serif;background:#0f172a;color:#e2e8f0;min-height:100vh}
|
||||
header{background:#1e293b;padding:1rem 2rem;display:flex;align-items:center;gap:1rem;border-bottom:1px solid #334155}
|
||||
header h1{font-size:1.2rem;font-weight:700;color:#f8fafc}
|
||||
.subtitle{font-size:.75rem;color:#64748b}
|
||||
.badge{font-size:.75rem;padding:.25rem .6rem;border-radius:9999px;font-weight:600}
|
||||
.badge.idle{background:#1e3a5f;color:#93c5fd}
|
||||
.badge.running{background:#713f12;color:#fde68a;animation:pulse 1s infinite}
|
||||
@keyframes pulse{0%,100%{opacity:1}50%{opacity:.5}}
|
||||
main{padding:1.5rem 2rem;max-width:1100px;margin:0 auto}
|
||||
.actions{display:flex;align-items:center;gap:1rem;flex-wrap:wrap;margin-bottom:1.5rem}
|
||||
button{padding:.55rem 1.3rem;border:none;border-radius:.5rem;cursor:pointer;font-weight:600;font-size:.85rem;transition:.15s}
|
||||
#btnSync{background:#2563eb;color:#fff}
|
||||
#btnSync:hover:not(:disabled){background:#1d4ed8}
|
||||
#btnSync:disabled{opacity:.4;cursor:not-allowed}
|
||||
#btnRefresh{background:#334155;color:#cbd5e1}
|
||||
#btnRefresh:hover{background:#475569}
|
||||
#nextSync{font-size:.8rem;color:#64748b;margin-left:auto}
|
||||
#msg{font-size:.82rem;padding:.35rem .75rem;border-radius:.4rem;display:none}
|
||||
#msg.ok{background:#14532d;color:#86efac;display:inline-block}
|
||||
#msg.err{background:#450a0a;color:#fca5a5;display:inline-block}
|
||||
#progressBox{background:#1e293b;border:1px solid #334155;border-radius:.75rem;padding:1rem 1.25rem;margin-bottom:1.5rem;display:none}
|
||||
#progressBox.visible{display:block}
|
||||
.progress-header{display:flex;justify-content:space-between;margin-bottom:.6rem;font-size:.85rem}
|
||||
.progress-bar-bg{background:#0f172a;border-radius:9999px;height:12px;overflow:hidden}
|
||||
.progress-bar-fill{height:100%;background:linear-gradient(90deg,#059669,#2563eb);border-radius:9999px;transition:width .5s ease}
|
||||
.progress-detail{margin-top:.5rem;font-size:.78rem;color:#64748b}
|
||||
.grid{display:grid;grid-template-columns:1fr 1fr;gap:1.5rem}
|
||||
@media(max-width:700px){.grid{grid-template-columns:1fr}}
|
||||
.card{background:#1e293b;border:1px solid #334155;border-radius:.75rem;overflow:hidden}
|
||||
.card-header{padding:.65rem 1rem;background:#0f172a;border-bottom:1px solid #334155;font-size:.75rem;font-weight:600;text-transform:uppercase;letter-spacing:.05em;color:#94a3b8}
|
||||
table{width:100%;border-collapse:collapse;font-size:.85rem}
|
||||
th{padding:.55rem 1rem;text-align:left;color:#64748b;font-weight:500;border-bottom:1px solid #334155}
|
||||
td{padding:.55rem 1rem;border-bottom:1px solid #1e293b}
|
||||
tr:last-child td{border-bottom:none}
|
||||
.ok{color:#4ade80}.err{color:#f87171}
|
||||
.log-box{height:320px;overflow-y:auto;font-family:monospace;font-size:.76rem;padding:.75rem 1rem;background:#0a0f1e;color:#a5f3fc;line-height:1.7;user-select:text;cursor:text}
|
||||
.log-line{white-space:pre-wrap;word-break:break-all}
|
||||
.log-line.error{color:#f87171}
|
||||
.ts{color:#475569;margin-right:.5rem}
|
||||
#btnCopyLogs{background:#334155;color:#cbd5e1;font-size:.75rem;padding:.35rem .8rem}
|
||||
#btnCopyLogs:hover{background:#475569}
|
||||
#lastRefresh{font-size:.75rem;color:#475569}
|
||||
</style>
|
||||
</head>
|
||||
<body>
|
||||
<header>
|
||||
<h1>⚡ Sync Hostinger</h1>
|
||||
<span class="subtitle">PG Debian → PG Hostinger</span>
|
||||
<span id="statusBadge" class="badge idle">Idle</span>
|
||||
<span id="lastRefresh" style="margin-left:auto"></span>
|
||||
</header>
|
||||
<main>
|
||||
<div class="actions">
|
||||
<button id="btnSync" onclick="runSync()">▶ Lancer sync maintenant</button>
|
||||
<button id="btnRefresh" onclick="refresh()">↺ Rafraîchir</button>
|
||||
<span id="msg"></span>
|
||||
<span id="nextSync"></span>
|
||||
</div>
|
||||
|
||||
<div id="progressBox">
|
||||
<div class="progress-header">
|
||||
<span id="progressLabel">Synchronisation en cours…</span>
|
||||
<span id="progressPct" style="font-weight:700;color:#93c5fd">0%</span>
|
||||
</div>
|
||||
<div class="progress-bar-bg">
|
||||
<div class="progress-bar-fill" id="progressBar" style="width:0%"></div>
|
||||
</div>
|
||||
<div class="progress-detail" id="progressDetail"></div>
|
||||
</div>
|
||||
|
||||
<div class="grid">
|
||||
<div class="card" style="grid-column:1/-1">
|
||||
<div class="card-header">Statut des tables</div>
|
||||
<table>
|
||||
<thead><tr><th>Table</th><th>Dernière sync</th><th>Lignes</th><th>Statut</th></tr></thead>
|
||||
<tbody id="tableBody"><tr><td colspan="4" style="color:#475569;padding:1rem">Chargement…</td></tr></tbody>
|
||||
</table>
|
||||
</div>
|
||||
<div class="card" style="grid-column:1/-1">
|
||||
<div class="card-header" style="display:flex;justify-content:space-between;align-items:center">
|
||||
<span>Logs</span>
|
||||
<button id="btnCopyLogs" onclick="copyLogs()">Copier les logs</button>
|
||||
</div>
|
||||
<div class="log-box" id="logBox">Chargement…</div>
|
||||
</div>
|
||||
</div>
|
||||
</main>
|
||||
<script>
|
||||
function fmtDate(iso){
|
||||
if(!iso) return '—';
|
||||
return new Date(iso).toLocaleString('fr-FR');
|
||||
}
|
||||
function renderProgress(progress, running){
|
||||
const box = document.getElementById('progressBox');
|
||||
if(!running || !progress){ box.classList.remove('visible'); return; }
|
||||
box.classList.add('visible');
|
||||
document.getElementById('progressBar').style.width = progress.pct + '%';
|
||||
document.getElementById('progressPct').textContent = progress.pct + '%';
|
||||
document.getElementById('progressLabel').textContent =
|
||||
'Sync en cours — table ' + progress.index + '/' + progress.total;
|
||||
const elapsed = Math.round((Date.now() - new Date(progress.startedAt)) / 1000);
|
||||
document.getElementById('progressDetail').textContent =
|
||||
'En cours : ' + progress.current + ' | Durée : ' + elapsed + 's';
|
||||
}
|
||||
function renderStatus(data){
|
||||
const badge = document.getElementById('statusBadge');
|
||||
badge.textContent = data.running ? 'En cours…' : 'Idle';
|
||||
badge.className = 'badge ' + (data.running ? 'running' : 'idle');
|
||||
document.getElementById('btnSync').disabled = data.running;
|
||||
renderProgress(data.progress, data.running);
|
||||
if(!data.tables.length){
|
||||
document.getElementById('tableBody').innerHTML =
|
||||
'<tr><td colspan="4" style="color:#475569;padding:1rem">Aucune sync — cliquez sur Lancer sync maintenant</td></tr>';
|
||||
return;
|
||||
}
|
||||
document.getElementById('tableBody').innerHTML = data.tables.map(t => \`
|
||||
<tr>
|
||||
<td><strong>\${t.table_name}</strong></td>
|
||||
<td>\${fmtDate(t.last_sync)}</td>
|
||||
<td>\${t.rows_synced ?? '—'}</td>
|
||||
<td class="\${t.status==='ok'?'ok':'err'}">\${t.status ?? '—'}\${t.error_msg?' — '+t.error_msg:''}</td>
|
||||
</tr>
|
||||
\`).join('');
|
||||
}
|
||||
function renderLogs(logs){
|
||||
const box = document.getElementById('logBox');
|
||||
if(!logs.length){ box.textContent = 'Aucun log.'; return; }
|
||||
const atBottom = box.scrollTop + box.clientHeight >= box.scrollHeight - 10;
|
||||
box.innerHTML = logs.map(l => \`
|
||||
<div class="log-line \${l.level}">
|
||||
<span class="ts">\${new Date(l.ts).toLocaleTimeString('fr-FR')}</span>\${l.msg}
|
||||
</div>
|
||||
\`).join('');
|
||||
if(atBottom) box.scrollTop = box.scrollHeight;
|
||||
}
|
||||
async function refresh(){
|
||||
try{
|
||||
const [s, logs, next] = await Promise.all([
|
||||
fetch('/api/status').then(r=>r.json()),
|
||||
fetch('/api/logs?limit=150').then(r=>r.json()),
|
||||
fetch('/api/next').then(r=>r.json()),
|
||||
]);
|
||||
renderStatus(s);
|
||||
renderLogs(logs);
|
||||
document.getElementById('nextSync').textContent = 'Prochaine sync auto : ' + fmtDate(next.next);
|
||||
document.getElementById('lastRefresh').textContent = 'Mis à jour ' + new Date().toLocaleTimeString('fr-FR');
|
||||
} catch(e){ console.error(e); }
|
||||
}
|
||||
function copyLogs(){
|
||||
const box = document.getElementById('logBox');
|
||||
const text = Array.from(box.querySelectorAll('.log-line')).map(l => l.textContent).join('\\n');
|
||||
const btn = document.getElementById('btnCopyLogs');
|
||||
function onSuccess(){ btn.textContent = 'Copié !'; setTimeout(() => { btn.textContent = 'Copier les logs'; }, 2000); }
|
||||
if(navigator.clipboard && window.isSecureContext){
|
||||
navigator.clipboard.writeText(text).then(onSuccess).catch(fallback);
|
||||
} else { fallback(); }
|
||||
function fallback(){
|
||||
const ta = document.createElement('textarea');
|
||||
ta.value = text;
|
||||
ta.style.cssText = 'position:fixed;top:-9999px;left:-9999px;opacity:0';
|
||||
document.body.appendChild(ta); ta.focus(); ta.select();
|
||||
try{ document.execCommand('copy'); onSuccess(); } catch(e){ alert('Sélectionnez les logs et Ctrl+C'); }
|
||||
document.body.removeChild(ta);
|
||||
}
|
||||
}
|
||||
async function runSync(){
|
||||
const msg = document.getElementById('msg');
|
||||
msg.style.display = 'none';
|
||||
const r = await fetch('/api/sync/run',{method:'POST'});
|
||||
const data = await r.json();
|
||||
msg.textContent = r.ok ? '✓ ' + data.message : '✗ ' + data.error;
|
||||
msg.className = r.ok ? 'ok' : 'err';
|
||||
setTimeout(()=>{ msg.style.display='none'; }, 5000);
|
||||
setTimeout(refresh, 800);
|
||||
}
|
||||
let timer;
|
||||
async function loop(){
|
||||
await refresh();
|
||||
const s = await fetch('/api/status').then(r=>r.json()).catch(()=>({}));
|
||||
timer = setTimeout(loop, s.running ? 3000 : 15000);
|
||||
}
|
||||
loop();
|
||||
</script>
|
||||
</body>
|
||||
</html>`;
|
||||
@@ -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) ==="
|
||||
@@ -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 };
|
||||
@@ -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 };
|
||||
Reference in new issue
Block a user