FireFeed Docs

Concurrency

Come funziona il locking della pipeline e i limiti tunable delle queue.

FireFeed bilancia tre tipi di concorrenza:

  1. Inter-project: pipeline di progetti diversi che girano in parallelo.
  2. Intra-project: il vincolo che impedisce due run sovrapposti dello stesso progetto.
  3. Intra-step: quanti job dello stesso tipo (import/rules/export) il worker processa in parallelo.

Lock intra-project: advisory lock Postgres

Il route handler POST /api/projects/{id}/pipeline/run (e l'alias run-now) avvolge il check-then-create dentro withProjectPipelineLock(prisma, projectId, ...) definito in packages/shared/src/services/pipeline-lock.service.ts.

return prisma.$transaction(async () => {
  const [{ acquired }] = await prisma.$queryRaw<...>`
    SELECT pg_try_advisory_xact_lock(${projectIdToLockKey(projectId)}) AS acquired`;
  if (!acquired) throw new PipelineLockBusy(projectId);
  return fn();
});

Caratteristiche:

  • Transaction-scoped: il lock viene rilasciato automaticamente al commit/rollback della transazione. Niente unlock manuale, niente deadlock se l'handler crasha.
  • Non-blocking: pg_try_advisory_xact_lock fa fail-fast (ritorna false), non aspetta. Due click consecutivi sul bottone "Run pipeline" producono uno 202 + un 409 immediato, non un timeout.
  • Per-project: la lock key deriva da SHA-256(projectId) letto come int64 signed. Progetti diversi mappano a chiavi diverse, quindi non si bloccano a vicenda.
  • Zero schema impact: no migration, no nuove tabelle, no row-level locks. Solo pg_try_advisory_xact_lock.

Concurrency intra-step: env-driven

I 4 worker BullMQ leggono la concurrency da env, con fallback ai valori storici:

Env varDefaultCodaSignificato
PIPELINE_CONCURRENCY1pipelineRoot orchestrator. Lascia a 1 (vedi sotto).
IMPORT_CONCURRENCY2importQuanti import processa in parallelo l'intera istanza worker
RULES_CONCURRENCY2rulesQuante rules-pass per-export in parallelo
EXPORT_CONCURRENCY4exportQuanti export writers in parallelo

Definito centralmente in packages/worker/src/queues/concurrency.ts. Modificare il valore in .env e riavviare il worker — nessuna ricompilazione.

Perché PIPELINE_CONCURRENCY=1 di default

Il job root orchestra l'intera pipeline (FETCH → PARSE → MERGE → RULES → EXPORT) per un singolo project. Tenerlo a 1 garantisce che ogni progetto venga processato come unità atomica osservabile (timeline UI coerente, log raggruppati per pipelineRunId).

Bumparlo a >1 non aumenta la velocità di un singolo run — il parallelismo intra-step è già regolato da IMPORT/RULES/EXPORT_CONCURRENCY. Serve solo se vuoi processare N progetti diversi in parallelo.

Quando portare a >1

Prerequisiti operativi (tutti necessari):

  1. Stress test verde: girare pnpm -F @firefeed/worker exec vitest run **/*.stress.test.ts con PIPELINE_CONCURRENCY=N e verificare no race su SchemaService operations (creazione rules_tmp_* simultanea da progetti diversi, drop temp orfane).
  2. DB pool sizing: Postgres connections totali = (PIPELINE × (IMPORT + RULES + EXPORT)) + web pool. Default è ~30; alzare connection_limit Prisma e max_connections Postgres prima di bumppare.
  3. Memoria worker: ogni pipeline carica chunk dati in memoria (default 2_000 righe × N rules-pass). Per N progetti paralleli moltiplica.
  4. Lock cross-project: il lock advisory è per projectId, quindi 2 progetti diversi non si bloccano. Verifica però che operazioni schema-level (es. CREATE SCHEMA company_<companyId>) siano serializzate dove serve.

Solo allora alza in produzione. Inizia da 2, monitora /metrics (vedi Observability), bumpa progressivamente.

Race window chiusa (WS3)

Prima:

Tab A: findFirst(PENDING|RUNNING) → null  ┐
Tab B: findFirst(PENDING|RUNNING) → null  ├─ entrambe fetch worker → 2 run
Tab A: fetch /internal/pipeline/run       │  per lo stesso project
Tab B: fetch /internal/pipeline/run       ┘

Dopo:

Tab A: pg_try_advisory_xact_lock(key) → true → findFirst → null → enqueue → 202
Tab B: pg_try_advisory_xact_lock(key) → false → 409 (immediato)

Test verifica: vedi pipeline-lock.int.test.ts quando creato (TODO Phase 1 follow-up).

Schema operations (cross-project)

SchemaService.createCompanySchema, createProjectTable, addColumn modificano DDL Postgres. Sono già transazionali per singola chiamata, ma due chiamate concorrenti su company diverse non si bloccano (intenzionale).

Su stessa company, il bottleneck DDL Postgres (lock implicito su pg_class) serializza naturalmente. Non c'è bisogno di lock applicativi qui.

In questa pagina