Sincronizzazione tra un'app online e una offline con CloudEvents: l'implementazione in Node.js

Sincronizzazione tra un'app online e una offline con CloudEvents: l'implementazione in Node.js

Secondo articolo della serie sulla sincronizzazione tra un'app online sempre raggiungibile e un'app offline che si collega solo saltuariamente. Nel primo articolo abbiamo visto l'architettura generale (evento CloudEvents, pattern outbox/inbox, push/pull/ack) implementata in Go con un semplice store su file JSON. Qui replichiamo lo stesso design in Node.js, questa volta con una persistenza vera basata su SQLite sia lato server che lato client.

Riepilogo dell'architettura

Il formato dell'evento resta CloudEvents 1.0: un envelope JSON con specversion, id (univoco, per l'idempotenza), type, source, time e data. Il protocollo di sync espone tre operazioni: push per inviare gli eventi accumulati nell'outbox locale, pull per scaricare quelli generati altrove da quando ci si è sincronizzati l'ultima volta, ack per confermare al server fino a dove sono stati applicati, così il cursore di quel client avanza.

L'app online: server Express con SQLite

Il modulo better-sqlite3 espone un'API sincrona che rende il codice più semplice da leggere rispetto a un driver asincrono, ed è più che sufficiente per il volume di scritture di un server di sync. Isoliamo la persistenza in db.js:

// db.js — persistenza dell'app online: event store append-only + cursore per client
const Database = require('better-sqlite3');
const path = require('path');

const db = new Database(path.join(__dirname, 'events.db'));
db.pragma('journal_mode = WAL');

db.exec(`
CREATE TABLE IF NOT EXISTS events (
  seq             INTEGER PRIMARY KEY AUTOINCREMENT,
  id              TEXT UNIQUE NOT NULL,
  specversion     TEXT NOT NULL,
  type            TEXT NOT NULL,
  source          TEXT NOT NULL,
  time            TEXT NOT NULL,
  datacontenttype TEXT,
  data            TEXT NOT NULL,
  received_at     TEXT NOT NULL DEFAULT (datetime('now'))
);

CREATE TABLE IF NOT EXISTS client_cursors (
  client_id TEXT PRIMARY KEY,
  last_seq  INTEGER NOT NULL DEFAULT 0
);
`);

const insertStmt = db.prepare(`
  INSERT INTO events (id, specversion, type, source, time, datacontenttype, data)
  VALUES (@id, @specversion, @type, @source, @time, @datacontenttype, @data)
`);
const existsStmt = db.prepare('SELECT 1 FROM events WHERE id = ?');
const sinceStmt = db.prepare(`
  SELECT seq, id, specversion, type, source, time, datacontenttype, data
  FROM events WHERE seq > ? ORDER BY seq ASC LIMIT ?
`);
const getCursorStmt = db.prepare('SELECT last_seq FROM client_cursors WHERE client_id = ?');
const setCursorStmt = db.prepare(`
  INSERT INTO client_cursors (client_id, last_seq) VALUES (?, ?)
  ON CONFLICT(client_id) DO UPDATE SET last_seq = excluded.last_seq
`);

// Inserisce l'evento solo se il suo id non è già stato visto (idempotenza)
function insertEventIfNew(event) {
  if (existsStmt.get(event.id)) return { inserted: false };

  insertStmt.run({
    id: event.id,
    specversion: event.specversion,
    type: event.type,
    source: event.source,
    time: event.time,
    datacontenttype: event.datacontenttype || 'application/json',
    data: JSON.stringify(event.data ?? {})
  });

  return { inserted: true };
}

function getEventsSince(seq, limit = 200) {
  const rows = sinceStmt.all(seq, limit);
  return rows.map((r) => ({
    specversion: r.specversion,
    id: r.id,
    type: r.type,
    source: r.source,
    time: r.time,
    datacontenttype: r.datacontenttype,
    data: JSON.parse(r.data),
    _seq: r.seq
  }));
}

function getClientCursor(clientId) {
  const row = getCursorStmt.get(clientId);
  return row ? row.last_seq : 0;
}

function setClientCursor(clientId, seq) {
  setCursorStmt.run(clientId, seq);
}

module.exports = { insertEventIfNew, getEventsSince, getClientCursor, setClientCursor };

E il server Express con i tre endpoint:

// server.js — app online: endpoint di sincronizzazione basati su CloudEvents 1.0
const express = require('express');
const {
  insertEventIfNew,
  getEventsSince,
  getClientCursor,
  setClientCursor
} = require('./db');

const app = express();
app.use(express.json({ limit: '5mb' }));

const API_KEY = process.env.API_KEY || 'dev-secret';

// Autenticazione minimale via API key statica.
// In produzione: JWT con refresh token, oppure mTLS se i client sono dispositivi fissi.
function requireApiKey(req, res, next) {
  if (req.header('x-api-key') !== API_KEY) {
    return res.status(401).json({ error: 'unauthorized' });
  }
  next();
}

function isValidCloudEvent(ev) {
  return (
    ev &&
    ev.specversion === '1.0' &&
    typeof ev.id === 'string' &&
    typeof ev.type === 'string' &&
    typeof ev.source === 'string' &&
    typeof ev.time === 'string'
  );
}

// Il client offline invia gli eventi accumulati mentre non era connesso
app.post('/api/sync/push', requireApiKey, (req, res) => {
  const { clientId, events } = req.body || {};
  if (!clientId || !Array.isArray(events)) {
    return res.status(400).json({ error: 'clientId ed events[] sono obbligatori' });
  }

  const accepted = [];
  const rejected = [];

  for (const ev of events) {
    if (!isValidCloudEvent(ev)) {
      rejected.push({ id: ev && ev.id, reason: 'evento non conforme a CloudEvents 1.0' });
      continue;
    }
    const { inserted } = insertEventIfNew(ev);
    // inserted=false => id già visto in precedenza, scartato per idempotenza
    accepted.push({ id: ev.id, inserted });
  }

  res.json({ accepted, rejected, serverTime: new Date().toISOString() });
});

// Il client offline scarica gli eventi generati altrove dopo l'ultima sync
app.get('/api/sync/pull', requireApiKey, (req, res) => {
  const { clientId } = req.query;
  if (!clientId) return res.status(400).json({ error: 'clientId obbligatorio' });

  const since = getClientCursor(clientId);
  const events = getEventsSince(since);
  const newCursor = events.length ? events[events.length - 1]._seq : since;

  res.json({
    events: events.map(({ _seq, ...ev }) => ev), // _seq è dettaglio interno, non fa parte del payload CloudEvents
    cursor: newCursor
  });
});

// Il client conferma fino a dove ha applicato gli eventi ricevuti
app.post('/api/sync/ack', requireApiKey, (req, res) => {
  const { clientId, cursor } = req.body || {};
  if (!clientId || typeof cursor !== 'number') {
    return res.status(400).json({ error: 'clientId e cursor sono obbligatori' });
  }
  setClientCursor(clientId, cursor);
  res.json({ ok: true });
});

const PORT = process.env.PORT || 3000;
app.listen(PORT, () => console.log(`Online sync server in ascolto su :${PORT}`));

L'app offline: outbox locale con SQLite e loop di sync

Lato client usiamo lo stesso modulo per un database SQLite locale, con una tabella outbox per gli eventi in uscita e una applied_events per restare idempotenti sugli eventi in ingresso:

// db.js — persistenza locale dell'app offline: outbox ed eventi applicati
const Database = require('better-sqlite3');
const path = require('path');

const db = new Database(path.join(__dirname, 'offline.db'));
db.pragma('journal_mode = WAL');

db.exec(`
CREATE TABLE IF NOT EXISTS outbox (
  id     TEXT PRIMARY KEY,
  type   TEXT NOT NULL,
  source TEXT NOT NULL,
  time   TEXT NOT NULL,
  data   TEXT NOT NULL,
  sent   INTEGER NOT NULL DEFAULT 0
);

CREATE TABLE IF NOT EXISTS applied_events (
  id         TEXT PRIMARY KEY,
  type       TEXT NOT NULL,
  data       TEXT NOT NULL,
  applied_at TEXT NOT NULL DEFAULT (datetime('now'))
);
`);

const enqueueStmt = db.prepare(`
  INSERT INTO outbox (id, type, source, time, data, sent) VALUES (@id, @type, @source, @time, @data, 0)
`);
const unsentStmt = db.prepare(`
  SELECT id, type, source, time, data FROM outbox WHERE sent = 0 ORDER BY rowid ASC LIMIT ?
`);
const markSentStmt = db.prepare('UPDATE outbox SET sent = 1 WHERE id = ?');
const applyStmt = db.prepare(`
  INSERT OR IGNORE INTO applied_events (id, type, data) VALUES (@id, @type, @data)
`);

// Chiamata dal dominio offline ogni volta che avviene un cambiamento locale
function enqueueEvent(ev) {
  enqueueStmt.run({
    id: ev.id,
    type: ev.type,
    source: ev.source,
    time: ev.time,
    data: JSON.stringify(ev.data ?? {})
  });
}

function getUnsentEvents(limit = 200) {
  return unsentStmt.all(limit).map((r) => ({
    specversion: '1.0',
    id: r.id,
    type: r.type,
    source: r.source,
    time: r.time,
    datacontenttype: 'application/json',
    data: JSON.parse(r.data)
  }));
}

function markSent(ids) {
  const tx = db.transaction((ids) => {
    for (const id of ids) markSentStmt.run(id);
  });
  tx(ids);
}

// Applica un evento ricevuto dal server al dominio locale.
// "INSERT OR IGNORE" garantisce idempotenza se lo stesso evento arriva due volte.
function applyIncomingEvent(ev) {
  applyStmt.run({ id: ev.id, type: ev.type, data: JSON.stringify(ev.data ?? {}) });
  // Qui va la logica reale: leggere ev.type/ev.data e aggiornare l'entità
  // corrispondente nel modello di dominio offline.
}

module.exports = { enqueueEvent, getUnsentEvents, markSent, applyIncomingEvent };

E il loop che prova a sincronizzarsi ogni 30 secondi, usando il fetch globale disponibile in Node.js dalla versione 18:

// client.js — app offline: prova a sincronizzarsi con il server a intervalli regolari.
// Se la rete manca, la fetch fallisce e si riprova al giro successivo: è questo
// che rende l'app "offline-first" invece che "online-required".
const crypto = require('crypto');
const { enqueueEvent, getUnsentEvents, markSent, applyIncomingEvent } = require('./db');

const BASE_URL = process.env.SYNC_SERVER_URL || 'http://localhost:3000';
const API_KEY = process.env.API_KEY || 'dev-secret';
const CLIENT_ID = process.env.CLIENT_ID || 'offline-client-01';

async function trySync() {
  // 1) PUSH: invia gli eventi locali accumulati mentre si era offline
  const pending = getUnsentEvents();
  if (pending.length > 0) {
    const pushRes = await fetch(`${BASE_URL}/api/sync/push`, {
      method: 'POST',
      headers: { 'Content-Type': 'application/json', 'x-api-key': API_KEY },
      body: JSON.stringify({ clientId: CLIENT_ID, events: pending })
    });
    if (!pushRes.ok) throw new Error(`push fallito: HTTP ${pushRes.status}`);

    markSent(pending.map((ev) => ev.id));
    console.log(`Inviati ${pending.length} eventi al server`);
  }

  // 2) PULL: scarica gli eventi generati altrove dopo l'ultima sync
  const pullRes = await fetch(`${BASE_URL}/api/sync/pull?clientId=${encodeURIComponent(CLIENT_ID)}`, {
    headers: { 'x-api-key': API_KEY }
  });
  if (!pullRes.ok) throw new Error(`pull fallito: HTTP ${pullRes.status}`);

  const { events, cursor } = await pullRes.json();
  for (const ev of events) applyIncomingEvent(ev);

  // 3) ACK: conferma al server fino a dove si è applicato, per avanzare il cursore
  if (events.length > 0) {
    await fetch(`${BASE_URL}/api/sync/ack`, {
      method: 'POST',
      headers: { 'Content-Type': 'application/json', 'x-api-key': API_KEY },
      body: JSON.stringify({ clientId: CLIENT_ID, cursor })
    });
    console.log(`Applicati ${events.length} eventi ricevuti dal server`);
  }
}

async function syncLoop() {
  try {
    await trySync();
  } catch (err) {
    // Rete assente o server irraggiungibile: comportamento atteso per un
    // client offline-first. Si riprova semplicemente al prossimo giro.
    console.warn(`Sync non riuscita (rete assente?): ${err.message}`);
  }
}

// Esempio: simula un cambiamento avvenuto nel dominio offline.
// Nella tua app reale questa chiamata va fatta subito dopo ogni
// scrittura locale rilevante.
enqueueEvent({
  id: crypto.randomUUID(),
  type: 'com.gabrieleromanato.offlineapp.record.created',
  source: `urn:client:${CLIENT_ID}`,
  time: new Date().toISOString(),
  data: { recordId: 42, note: 'Creato mentre offline' }
});

syncLoop();
setInterval(syncLoop, 30_000);

Provarlo end-to-end

cd online-node && npm install && npm start
# in un altro terminale
cd offline-node && npm install && npm start

Il client accoda subito un evento di esempio nel suo offline.db locale, lo invia al server con push, e con la pull successiva lo riceve indietro applicandolo in modo idempotente al proprio registro. Puoi verificarlo ispezionando i due database SQLite con qualunque client (anche sqlite3 events.db "SELECT * FROM events"): l'evento generato offline compare sul server con lo stesso id, e il cursore del client in client_cursors avanza a ogni sync.

Cosa manca per la produzione

Anche qui l'esempio è volutamente minimale. Prima di usarlo davvero servono: autenticazione più solida (JWT con refresh o mTLS) al posto della API key statica; retry con backoff esponenziale sul push invece di un intervallo fisso; un filtro lato client che escluda dalla pull gli eventi il cui source coincide con il client stesso, per non riapplicarsi da solo ciò che ha appena generato; validazione dello schema di data per ogni type (con JSON Schema, per esempio); e una strategia esplicita di conflict resolution se due client modificano lo stesso record mentre erano entrambi offline.

Nel prossimo articolo vediamo la stessa architettura in Laravel, con il server come applicazione web classica e il client come comando Artisan schedulabile.