Sincronizzazione tra un'app online e una offline con CloudEvents: l'implementazione in Python
Quarto articolo della serie: stessa architettura di sincronizzazione (evento CloudEvents 1.0, pattern outbox/inbox, push/pull/ack), questa volta in Python. Il server online usa FastAPI, il client offline è uno script con requests e un loop di sincronizzazione; entrambi persistono i dati con sqlite3 della libreria standard, senza bisogno di un ORM.
Riepilogo dell'architettura
L'evento viaggia come CloudEvents 1.0: specversion, id univoco (garantisce l'idempotenza), type, source, time, data. Tre operazioni sul server: push riceve il batch di eventi accumulati nell'outbox del client, pull restituisce gli eventi generati altrove dopo l'ultima sync, ack fa avanzare il cursore del client lato server una volta che ha applicato quanto ricevuto.
L'app online: persistenza
Il modulo sqlite3 della standard library basta: nessuna dipendenza extra per il livello dati.
"""db.py — persistenza dell'app online: event store append-only + cursore per client."""
import json
import sqlite3
from pathlib import Path
DB_PATH = Path(__file__).parent / "events.db"
_conn = sqlite3.connect(DB_PATH, check_same_thread=False)
_conn.execute("PRAGMA journal_mode=WAL")
_conn.execute(
"""
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'))
)
"""
)
_conn.execute(
"""
CREATE TABLE IF NOT EXISTS client_cursors (
client_id TEXT PRIMARY KEY,
last_seq INTEGER NOT NULL DEFAULT 0
)
"""
)
_conn.commit()
def insert_event_if_new(event: dict) -> bool:
"""Inserisce l'evento solo se il suo id non è già stato visto (idempotenza)."""
existing = _conn.execute("SELECT 1 FROM events WHERE id = ?", (event["id"],)).fetchone()
if existing:
return False
_conn.execute(
"""
INSERT INTO events (id, specversion, type, source, time, datacontenttype, data)
VALUES (?, ?, ?, ?, ?, ?, ?)
""",
(
event["id"],
event["specversion"],
event["type"],
event["source"],
event["time"],
event.get("datacontenttype", "application/json"),
json.dumps(event.get("data", {})),
),
)
_conn.commit()
return True
def events_since(seq: int, limit: int = 200) -> list[dict]:
rows = _conn.execute(
"""
SELECT seq, id, specversion, type, source, time, datacontenttype, data
FROM events WHERE seq > ? ORDER BY seq ASC LIMIT ?
""",
(seq, limit),
).fetchall()
return [
{
"_seq": row[0],
"specversion": row[2],
"id": row[1],
"type": row[3],
"source": row[4],
"time": row[5],
"datacontenttype": row[6],
"data": json.loads(row[7]),
}
for row in rows
]
def get_client_cursor(client_id: str) -> int:
row = _conn.execute(
"SELECT last_seq FROM client_cursors WHERE client_id = ?", (client_id,)
).fetchone()
return row[0] if row else 0
def set_client_cursor(client_id: str, seq: int) -> None:
_conn.execute(
"""
INSERT INTO client_cursors (client_id, last_seq) VALUES (?, ?)
ON CONFLICT(client_id) DO UPDATE SET last_seq = excluded.last_seq
""",
(client_id, seq),
)
_conn.commit()
L'app online: gli endpoint FastAPI
Per restare fedeli al comportamento visto negli articoli precedenti — accettare i singoli eventi validi anche quando altri nello stesso batch non lo sono — gli eventi in ingresso vengono letti come dizionari generici e validati a mano, invece di affidarsi a un modello Pydantic rigido per ciascun evento (che farebbe fallire l'intera richiesta al primo eventuale errore):
"""main.py — app online: server di sincronizzazione basato su CloudEvents 1.0."""
import os
from datetime import datetime, timezone
from fastapi import Depends, FastAPI, Header, HTTPException, Query
from pydantic import BaseModel
from db import events_since, get_client_cursor, insert_event_if_new, set_client_cursor
API_KEY = os.environ.get("API_KEY", "dev-secret")
app = FastAPI()
class PushRequest(BaseModel):
clientId: str
events: list[dict] = []
class AckRequest(BaseModel):
clientId: str
cursor: int
# Autenticazione minimale via API key statica (in produzione: JWT o mTLS)
def require_api_key(x_api_key: str = Header(default="")) -> None:
if x_api_key != API_KEY:
raise HTTPException(status_code=401, detail="unauthorized")
def is_valid_cloud_event(event: dict) -> bool:
return (
isinstance(event, dict)
and event.get("specversion") == "1.0"
and bool(event.get("id"))
and bool(event.get("type"))
and bool(event.get("source"))
and bool(event.get("time"))
)
# Il client offline invia gli eventi accumulati mentre non era connesso
@app.post("/api/sync/push")
def push(req: PushRequest, _: None = Depends(require_api_key)):
accepted = []
rejected = []
for event in req.events:
if not is_valid_cloud_event(event):
rejected.append({"id": event.get("id"), "reason": "evento non conforme a CloudEvents 1.0"})
continue
# inserted=False => id già visto in precedenza, scartato per idempotenza
inserted = insert_event_if_new(event)
accepted.append({"id": event["id"], "inserted": inserted})
return {
"accepted": accepted,
"rejected": rejected,
"serverTime": datetime.now(timezone.utc).isoformat(),
}
# Il client offline scarica gli eventi generati altrove dopo l'ultima sync
@app.get("/api/sync/pull")
def pull(clientId: str = Query(...), _: None = Depends(require_api_key)):
since = get_client_cursor(clientId)
events = events_since(since)
cursor = since
out = []
for ev in events:
cursor = ev["_seq"]
out.append({k: v for k, v in ev.items() if k != "_seq"})
return {"events": out, "cursor": cursor}
# Il client conferma fino a dove ha applicato gli eventi ricevuti
@app.post("/api/sync/ack")
def ack(req: AckRequest, _: None = Depends(require_api_key)):
set_client_cursor(req.clientId, req.cursor)
return {"ok": True}
fastapi>=0.115
uvicorn[standard]>=0.30
L'app offline: outbox e loop di sync
Stessa struttura a due tabelle vista negli articoli precedenti, qui con sqlite3 puro invece di un ORM:
"""db.py — persistenza locale dell'app offline: outbox ed eventi applicati."""
import json
import sqlite3
from pathlib import Path
DB_PATH = Path(__file__).parent / "offline.db"
_conn = sqlite3.connect(DB_PATH, check_same_thread=False)
_conn.execute("PRAGMA journal_mode=WAL")
_conn.execute(
"""
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
)
"""
)
_conn.execute(
"""
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'))
)
"""
)
_conn.commit()
def enqueue_event(event: dict) -> None:
"""Chiamata dal dominio offline ogni volta che avviene un cambiamento locale."""
_conn.execute(
"INSERT INTO outbox (id, type, source, time, data, sent) VALUES (?, ?, ?, ?, ?, 0)",
(event["id"], event["type"], event["source"], event["time"], json.dumps(event.get("data", {}))),
)
_conn.commit()
def get_unsent_events(limit: int = 200) -> list[dict]:
rows = _conn.execute(
"SELECT id, type, source, time, data FROM outbox WHERE sent = 0 ORDER BY rowid ASC LIMIT ?",
(limit,),
).fetchall()
return [
{
"specversion": "1.0",
"id": row[0],
"type": row[1],
"source": row[2],
"time": row[3],
"datacontenttype": "application/json",
"data": json.loads(row[4]),
}
for row in rows
]
def mark_sent(ids: list[str]) -> None:
_conn.executemany("UPDATE outbox SET sent = 1 WHERE id = ?", [(i,) for i in ids])
_conn.commit()
def apply_incoming_event(event: dict) -> None:
"""Applica un evento ricevuto dal server al dominio locale; idempotente sull'id."""
_conn.execute(
"INSERT OR IGNORE INTO applied_events (id, type, data) VALUES (?, ?, ?)",
(event["id"], event["type"], json.dumps(event.get("data", {}))),
)
_conn.commit()
# Qui va la logica reale: leggere event["type"]/event["data"] e aggiornare
# l'entità corrispondente nel modello di dominio offline.
"""client.py — app offline: prova a sincronizzarsi con il server a intervalli regolari.
Se la rete manca, la richiesta fallisce e si riprova al giro successivo: è questo
che rende l'app "offline-first" invece che "online-required".
"""
import time
import uuid
from datetime import datetime, timezone
import requests
from db import apply_incoming_event, enqueue_event, get_unsent_events, mark_sent
BASE_URL = "http://localhost:3000"
API_KEY = "dev-secret"
CLIENT_ID = "offline-client-01"
session = requests.Session()
session.headers.update({"x-api-key": API_KEY})
def try_sync() -> None:
# 1) PUSH: invia gli eventi locali accumulati mentre si era offline
pending = get_unsent_events()
if pending:
response = session.post(
f"{BASE_URL}/api/sync/push",
json={"clientId": CLIENT_ID, "events": pending},
timeout=10,
)
response.raise_for_status()
mark_sent([event["id"] for event in pending])
print(f"Inviati {len(pending)} eventi al server")
# 2) PULL: scarica gli eventi generati altrove dopo l'ultima sync
response = session.get(
f"{BASE_URL}/api/sync/pull", params={"clientId": CLIENT_ID}, timeout=10
)
response.raise_for_status()
payload = response.json()
for event in payload["events"]:
apply_incoming_event(event)
# 3) ACK: conferma al server fino a dove si è applicato
if payload["events"]:
session.post(
f"{BASE_URL}/api/sync/ack",
json={"clientId": CLIENT_ID, "cursor": payload["cursor"]},
timeout=10,
)
print(f"Applicati {len(payload['events'])} eventi ricevuti dal server")
def sync_loop() -> None:
try:
try_sync()
except requests.RequestException as exc:
# Rete assente o server irraggiungibile: comportamento atteso per un
# client offline-first. Si riprova semplicemente al prossimo giro.
print(f"Sync non riuscita (rete assente?): {exc}")
if __name__ == "__main__":
# Esempio: simula un cambiamento avvenuto nel dominio offline.
# Nella tua app reale questa chiamata va fatta subito dopo ogni
# scrittura locale rilevante.
enqueue_event(
{
"id": str(uuid.uuid4()),
"type": "com.gabrieleromanato.offlineapp.record.created",
"source": f"urn:client:{CLIENT_ID}",
"time": datetime.now(timezone.utc).isoformat(),
"data": {"recordId": 42, "note": "Creato mentre offline"},
}
)
while True:
sync_loop()
time.sleep(30)
requests>=2.32
Provarlo end-to-end
cd online-python
pip install -r requirements.txt
API_KEY=dev-secret uvicorn main:app --host 0.0.0.0 --port 3000
# in un altro terminale
cd offline-python
pip install -r requirements.txt
python3 client.py
Il client accoda l'evento di esempio, lo invia con push, e con la pull successiva lo riceve indietro applicandolo idempotentemente. Puoi ispezionare i due database con sqlite3 events.db "SELECT * FROM events": dopo un ciclo di sync l'evento generato offline compare sul server con lo stesso id, e il cursore del client in client_cursors avanza a 1.
Cosa manca per la produzione
Come negli articoli precedenti: autenticazione più solida di una API key statica (JWT con refresh o mTLS), retry con backoff esponenziale invece dell'intervallo fisso di 30 secondi, un filtro lato client sugli eventi con source uguale al proprio per non riapplicarsi da solo ciò che ha generato, validazione dello schema di data per ogni type (con Pydantic o JSON Schema), e una strategia esplicita di conflict resolution per le modifiche concorrenti allo stesso record. In produzione vale anche la pena eseguire uvicorn dietro un process manager (systemd, supervisord) con più worker se il volume di sync lo richiede.
Nel prossimo articolo vediamo la stessa architettura in PHP puro, senza framework, per capire cosa serve davvero sotto il cofano quando non c'è Laravel a fornire routing, ORM e client HTTP.