Sincronizzazione tra un'app online e una offline con CloudEvents: l'implementazione in Laravel

Sincronizzazione tra un'app online e una offline con CloudEvents: l'implementazione in Laravel

Terzo articolo della serie sulla sincronizzazione tra un'app online sempre raggiungibile e un'app offline che si collega solo saltuariamente (nei precedenti l'abbiamo vista in Go e in Node.js). Qui la stessa architettura — evento CloudEvents, pattern outbox/inbox, push/pull/ack — prende la forma di due applicazioni Laravel: una web application classica per il server online, e un comando Artisan di lunga durata per il client offline.

Riepilogo dell'architettura

Ogni cambiamento viaggia come evento CloudEvents 1.0: specversion, id univoco (per l'idempotenza), type, source, time, data. Il server espone tre operazioni: push per ricevere il batch di eventi accumulati nell'outbox del client, pull perché il client scarichi ciò che è successo altrove dall'ultima sync, ack perché il client confermi fino a dove ha applicato gli eventi, facendo avanzare il proprio cursore lato server.

L'app online: migrazioni e modelli

Due tabelle: events, l'event store append-only, e client_cursors, un cursore per ogni client che si sincronizza. Notiamo che il campo CloudEvents id è mappato su event_id nella tabella, per non entrare in conflitto con la chiave primaria seq usata per l'ordinamento.

<?php

use Illuminate\Database\Migrations\Migration;
use Illuminate\Database\Schema\Blueprint;
use Illuminate\Support\Facades\Schema;

return new class extends Migration
{
    public function up(): void
    {
        Schema::create('events', function (Blueprint $table) {
            $table->id('seq');
            $table->uuid('event_id')->unique();
            $table->string('specversion');
            $table->string('type');
            $table->string('source');
            $table->string('time');
            $table->string('datacontenttype')->nullable();
            $table->json('data');
            $table->timestamp('received_at')->useCurrent();
        });
    }

    public function down(): void
    {
        Schema::dropIfExists('events');
    }
};
<?php

use Illuminate\Database\Migrations\Migration;
use Illuminate\Database\Schema\Blueprint;
use Illuminate\Support\Facades\Schema;

return new class extends Migration
{
    public function up(): void
    {
        Schema::create('client_cursors', function (Blueprint $table) {
            $table->string('client_id')->primary();
            $table->unsignedBigInteger('last_seq')->default(0);
        });
    }

    public function down(): void
    {
        Schema::dropIfExists('client_cursors');
    }
};

I modelli Eloquent corrispondenti, con il cast automatico di data da/verso JSON:

<?php

namespace App\Models;

use Illuminate\Database\Eloquent\Model;

class SyncEvent extends Model
{
    protected $table = 'events';

    protected $primaryKey = 'seq';

    public $timestamps = false;

    protected $fillable = [
        'event_id', 'specversion', 'type', 'source', 'time', 'datacontenttype', 'data',
    ];

    protected $casts = [
        'data' => 'array',
    ];

    // Rappresenta il record nel formato CloudEvents 1.0 atteso dal client
    public function toCloudEvent(): array
    {
        return [
            'specversion' => $this->specversion,
            'id' => $this->event_id,
            'type' => $this->type,
            'source' => $this->source,
            'time' => $this->time,
            'datacontenttype' => $this->datacontenttype ?? 'application/json',
            'data' => $this->data,
        ];
    }
}
<?php

namespace App\Models;

use Illuminate\Database\Eloquent\Model;

class ClientCursor extends Model
{
    protected $table = 'client_cursors';

    protected $primaryKey = 'client_id';

    public $incrementing = false;

    protected $keyType = 'string';

    public $timestamps = false;

    protected $fillable = ['client_id', 'last_seq'];
}

L'app online: middleware, rotte e controller

Un middleware verifica l'API key prima di lasciar passare le richieste verso gli endpoint di sync (in produzione andrebbe sostituito con un guard basato su token/JWT o con mTLS):

<?php

namespace App\Http\Middleware;

use Closure;
use Illuminate\Http\Request;
use Symfony\Component\HttpFoundation\Response;

// Autenticazione minimale via API key statica (in produzione: JWT o mTLS)
class EnsureApiKey
{
    public function handle(Request $request, Closure $next): Response
    {
        if ($request->header('X-Api-Key') !== config('services.sync.api_key')) {
            return response()->json(['error' => 'unauthorized'], 401);
        }

        return $next($request);
    }
}

Registriamo l'alias in bootstrap/app.php (Laravel 11+) e aggiungiamo la chiave in config/services.php:

// Estratto di bootstrap/app.php — registrazione dell'alias del middleware
use App\Http\Middleware\EnsureApiKey;
use Illuminate\Foundation\Configuration\Middleware;

return Application::configure(basePath: dirname(__DIR__))
    ->withRouting(
        web: __DIR__.'/../routes/web.php',
        api: __DIR__.'/../routes/api.php',
        commands: __DIR__.'/../routes/console.php',
        health: '/up',
    )
    ->withMiddleware(function (Middleware $middleware) {
        $middleware->alias([
            'api.key' => EnsureApiKey::class,
        ]);
    })
    ->create();
<?php

// Estratto di config/services.php
return [

    // ... altre voci di configurazione già presenti ...

    'sync' => [
        'api_key' => env('SYNC_API_KEY', 'dev-secret'),
    ],

];
<?php

use App\Http\Controllers\SyncController;
use Illuminate\Support\Facades\Route;

Route::middleware('api.key')->prefix('sync')->group(function () {
    Route::post('/push', [SyncController::class, 'push']);
    Route::get('/pull', [SyncController::class, 'pull']);
    Route::post('/ack', [SyncController::class, 'ack']);
});

E il controller con i tre metodi:

<?php

namespace App\Http\Controllers;

use App\Models\ClientCursor;
use App\Models\SyncEvent;
use Illuminate\Http\Request;

class SyncController extends Controller
{
    // Il client offline invia gli eventi accumulati mentre non era connesso
    public function push(Request $request)
    {
        $clientId = $request->input('clientId');
        $events = $request->input('events', []);

        if (! $clientId || ! is_array($events)) {
            return response()->json(['error' => 'clientId ed events[] sono obbligatori'], 400);
        }

        $accepted = [];
        $rejected = [];

        foreach ($events as $event) {
            if (! $this->isValidCloudEvent($event)) {
                $rejected[] = ['id' => $event['id'] ?? null, 'reason' => 'evento non conforme a CloudEvents 1.0'];
                continue;
            }

            // Se event_id esiste già è un duplicato (retry di rete): non va reinserito
            $existing = SyncEvent::where('event_id', $event['id'])->exists();

            if (! $existing) {
                SyncEvent::create([
                    'event_id' => $event['id'],
                    'specversion' => $event['specversion'],
                    'type' => $event['type'],
                    'source' => $event['source'],
                    'time' => $event['time'],
                    'datacontenttype' => $event['datacontenttype'] ?? 'application/json',
                    'data' => $event['data'] ?? [],
                ]);
            }

            $accepted[] = ['id' => $event['id'], 'inserted' => ! $existing];
        }

        return response()->json([
            'accepted' => $accepted,
            'rejected' => $rejected,
            'serverTime' => now()->toIso8601String(),
        ]);
    }

    // Il client offline scarica gli eventi generati altrove dopo l'ultima sync
    public function pull(Request $request)
    {
        $clientId = $request->query('clientId');
        if (! $clientId) {
            return response()->json(['error' => 'clientId obbligatorio'], 400);
        }

        $since = ClientCursor::find($clientId)?->last_seq ?? 0;

        $events = SyncEvent::where('seq', '>', $since)
            ->orderBy('seq')
            ->limit(200)
            ->get();

        $cursor = $events->isNotEmpty() ? $events->last()->seq : $since;

        return response()->json([
            'events' => $events->map->toCloudEvent(),
            'cursor' => $cursor,
        ]);
    }

    // Il client conferma fino a dove ha applicato gli eventi ricevuti
    public function ack(Request $request)
    {
        $clientId = $request->input('clientId');
        $cursor = $request->input('cursor');

        if (! $clientId || ! is_numeric($cursor)) {
            return response()->json(['error' => 'clientId e cursor sono obbligatori'], 400);
        }

        ClientCursor::updateOrCreate(['client_id' => $clientId], ['last_seq' => $cursor]);

        return response()->json(['ok' => true]);
    }

    private function isValidCloudEvent(mixed $event): bool
    {
        return is_array($event)
            && ($event['specversion'] ?? null) === '1.0'
            && ! empty($event['id'])
            && ! empty($event['type'])
            && ! empty($event['source'])
            && ! empty($event['time']);
    }
}

L'app offline: outbox, inbox e comando Artisan

Sul lato offline le tabelle sono speculari: outbox per gli eventi da inviare, applied_events per quelli già applicati (idempotenza in ingresso).

<?php

use Illuminate\Database\Migrations\Migration;
use Illuminate\Database\Schema\Blueprint;
use Illuminate\Support\Facades\Schema;

return new class extends Migration
{
    public function up(): void
    {
        Schema::create('outbox', function (Blueprint $table) {
            $table->uuid('event_id')->primary();
            $table->string('type');
            $table->string('source');
            $table->string('time');
            $table->json('data');
            $table->boolean('sent')->default(false);
        });
    }

    public function down(): void
    {
        Schema::dropIfExists('outbox');
    }
};
<?php

use Illuminate\Database\Migrations\Migration;
use Illuminate\Database\Schema\Blueprint;
use Illuminate\Support\Facades\Schema;

return new class extends Migration
{
    public function up(): void
    {
        Schema::create('applied_events', function (Blueprint $table) {
            $table->uuid('event_id')->primary();
            $table->string('type');
            $table->json('data');
            $table->timestamp('applied_at')->useCurrent();
        });
    }

    public function down(): void
    {
        Schema::dropIfExists('applied_events');
    }
};
<?php

namespace App\Models;

use Illuminate\Database\Eloquent\Model;

class OutboxEvent extends Model
{
    protected $table = 'outbox';

    protected $primaryKey = 'event_id';

    public $incrementing = false;

    protected $keyType = 'string';

    public $timestamps = false;

    protected $fillable = ['event_id', 'type', 'source', 'time', 'data', 'sent'];

    protected $casts = [
        'data' => 'array',
        'sent' => 'boolean',
    ];

    public function toCloudEvent(): array
    {
        return [
            'specversion' => '1.0',
            'id' => $this->event_id,
            'type' => $this->type,
            'source' => $this->source,
            'time' => $this->time,
            'datacontenttype' => 'application/json',
            'data' => $this->data,
        ];
    }
}
<?php

namespace App\Models;

use Illuminate\Database\Eloquent\Model;

class AppliedEvent extends Model
{
    protected $table = 'applied_events';

    protected $primaryKey = 'event_id';

    public $incrementing = false;

    protected $keyType = 'string';

    public $timestamps = false;

    protected $fillable = ['event_id', 'type', 'data'];

    protected $casts = ['data' => 'array'];
}

Il cuore del client è un comando Artisan che, per impostazione predefinita, gira in loop: prova a sincronizzarsi, dorme 30 secondi, riprova. L'opzione --once esegue un solo ciclo, comoda se preferisci schedularlo tu (via Supervisor o cron) invece di tenerlo sempre in esecuzione:

<?php

namespace App\Console\Commands;

use App\Models\AppliedEvent;
use App\Models\OutboxEvent;
use Illuminate\Console\Command;
use Illuminate\Http\Client\ConnectionException;
use Illuminate\Http\Client\RequestException;
use Illuminate\Support\Facades\Http;

class SyncRun extends Command
{
    protected $signature = 'sync:run {--once : Esegue un solo ciclo invece del loop infinito}';

    protected $description = "Sincronizza l'outbox locale con il server online (CloudEvents)";

    private string $clientId;

    private string $baseUrl;

    private string $apiKey;

    public function handle(): int
    {
        $this->clientId = config('services.sync.client_id', 'offline-client-01');
        $this->baseUrl = config('services.sync.server_url', 'http://localhost:8000');
        $this->apiKey = config('services.sync.api_key', 'dev-secret');

        do {
            $this->trySync();

            if (! $this->option('once')) {
                sleep(30);
            }
        } while (! $this->option('once'));

        return self::SUCCESS;
    }

    private function trySync(): void
    {
        try {
            // 1) PUSH: invia gli eventi locali accumulati mentre si era offline
            $pending = OutboxEvent::where('sent', false)->limit(200)->get();

            if ($pending->isNotEmpty()) {
                Http::withHeaders(['X-Api-Key' => $this->apiKey])
                    ->timeout(10)
                    ->post("{$this->baseUrl}/api/sync/push", [
                        'clientId' => $this->clientId,
                        'events' => $pending->map->toCloudEvent()->all(),
                    ])
                    ->throw();

                OutboxEvent::whereIn('event_id', $pending->pluck('event_id'))->update(['sent' => true]);
                $this->info("Inviati {$pending->count()} eventi al server");
            }

            // 2) PULL: scarica gli eventi generati altrove dopo l'ultima sync
            $response = Http::withHeaders(['X-Api-Key' => $this->apiKey])
                ->timeout(10)
                ->get("{$this->baseUrl}/api/sync/pull", ['clientId' => $this->clientId])
                ->throw()
                ->json();

            foreach ($response['events'] as $event) {
                // updateOrCreate è idempotente: un evento con lo stesso id non viene duplicato
                AppliedEvent::updateOrCreate(
                    ['event_id' => $event['id']],
                    ['type' => $event['type'], 'data' => $event['data']]
                );

                // Qui va la logica reale: leggere $event['type']/$event['data'] e
                // aggiornare l'entità corrispondente nel modello di dominio offline.
            }

            // 3) ACK: conferma al server fino a dove si è applicato
            if (! empty($response['events'])) {
                Http::withHeaders(['X-Api-Key' => $this->apiKey])
                    ->timeout(10)
                    ->post("{$this->baseUrl}/api/sync/ack", [
                        'clientId' => $this->clientId,
                        'cursor' => $response['cursor'],
                    ])
                    ->throw();

                $this->info('Applicati '.count($response['events']).' eventi ricevuti dal server');
            }
        } catch (ConnectionException|RequestException $e) {
            // Rete assente o server irraggiungibile: comportamento atteso per
            // un client offline-first. Si riprova al ciclo successivo.
            $this->warn("Sync non riuscita (rete assente?): {$e->getMessage()}");
        }
    }
}

Il comando è auto-scoperto da Laravel (nessuna registrazione manuale necessaria dalla 11 in poi), quindi basta aggiungere le stesse chiavi di configurazione in config/services.php:

<?php

// Estratto di config/services.php
return [

    // ... altre voci di configurazione già presenti ...

    'sync' => [
        'client_id' => env('SYNC_CLIENT_ID', 'offline-client-01'),
        'server_url' => env('SYNC_SERVER_URL', 'http://localhost:8000'),
        'api_key' => env('SYNC_API_KEY', 'dev-secret'),
    ],

];

Provarlo

# app online
composer create-project laravel/laravel online-laravel
# copia i file dei paragrafi precedenti nelle posizioni indicate, poi:
php artisan migrate
php artisan serve

# app offline, in un altro progetto
composer create-project laravel/laravel offline-laravel
php artisan migrate
php artisan sync:run

Tutto il codice di questo articolo segue le convenzioni di Laravel 11+ (routing e middleware registrati in bootstrap/app.php, comandi auto-scoperti) ed è stato validato a livello di sintassi PHP; non essendoci qui un ambiente con Composer per installare il framework, provalo end-to-end nel tuo ambiente di sviluppo con i comandi sopra.

Cosa manca per la produzione

Anche in questa versione mancano gli stessi accorgimenti visti negli articoli precedenti: autenticazione più robusta di una API key statica (Sanctum con token per client, o mTLS), retry con backoff invece del sleep(30) fisso, 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, e una strategia di conflict resolution esplicita per le modifiche concorrenti allo stesso record. In un'app Laravel reale, il comando sync:run --once schedulato ogni minuto da Supervisor è spesso più robusto del loop infinito, perché un crash del processo non lascia il client senza sync finché qualcuno non se ne accorge.

Nel prossimo articolo vediamo la stessa architettura in Python, con FastAPI lato server e uno script client basato su requests.