Files
khadhroony-solana-project/docs/architecture/009-ACQUISITION_WORKERS_AND_JOBS.md
2026-08-14 13:05:50 +02:00

18 KiB
Raw Blame History

Acquisition, workers, jobs et pipelines spécialisés

Objet

Ce document constitue la sortie principale de 0.0.3-pre.007.

Il transforme les frontières durables D1D4 en modèle opérationnel et définit :

  • quatre pipelines spécialisés correspondant chacun à une frontière durable ;
  • quatre workers continus ;
  • le job de backfill historique ;
  • trois jobs de replay indépendants ;
  • la séparation worker/job lifecycle ;
  • le backlog durable par input + processor/version ;
  • les processing outcomes ;
  • claim/lease et concurrence multi-instance ;
  • at-least-once + idempotence plutôt qu'un faux exactly-once ;
  • reprise après crash ;
  • notification wake-up + periodic polling ;
  • hot reconfiguration du raw retriever ;
  • batching, backpressure et métriques opérationnelles minimales.

Les apps spécialisées, services autonomes, control plane, scenarios et frontière IPC sont détaillés dans 010-APPS_SERVICES_SCENARIOS_AND_CONTROL.md. Une application globale reste future ; un orchestrateur commun n'est pas une crate retenue actuellement.

Principe : une frontière de processing, une logique réutilisable

Un worker live et un job de replay ne doivent pas réimplémenter séparément la même transformation durable.

KSP introduit donc les premiers pipelines spécialisés justifiés par un besoin concret :

ksp-pipeline-raw-ingestion-lib
ksp-pipeline-core-processing-lib
ksp-pipeline-generic-materialization-lib
ksp-pipeline-domain-projection-lib

Il ne s'agit pas d'un retour vers un ksp-pipeline-lib monolithique.

Chaque pipeline correspond à une seule frontière durable.

transport model -> D1
D1 -> D2
D2 -> D3
D3 -> D4

Le pipeline porte la logique réutilisable de la frontière ; le worker/job porte le lifecycle et le scope d'exécution.

Frontière des pipelines

ksp-pipeline-raw-ingestion-lib

Mission :

homogeneous on-chain transport model
        |
        v
conversion D1
        |
        v
persistence D1
        |
        v
notification after commit

Le pipeline ne choisit pas et ne pilote pas la stratégie d'acquisition réseau.

Le worker live et le job de backfill lui fournissent les données déjà acquises.

Dépendances conceptuelles :

ksp-pipeline-raw-ingestion-lib
    -> ksp-onchain-transport-lib
    -> ksp-store-api
    -> ksp-core-lib
    -> ksp-logging-lib

Il ne dépend pas de ksp-store-lib.

Le backend concret est fourni par composition via les contrats de ksp-store-api.

ksp-pipeline-core-processing-lib

Mission :

D1 input
   |
   v
conversion vers ksp-program-api
   |
   v
decoder/registry fourni
   |
   v
Core output + processing outcome
   |
   v
persistence D2

Dépendances conceptuelles :

ksp-pipeline-core-processing-lib
    -> ksp-program-api
    -> ksp-store-api
    -> ksp-core-lib
    -> ksp-logging-lib

Il ne dépend pas de ksp-program-lib ni de ksp-store-lib.

L'implémentation officielle Program ou une extension externe est fournie par composition.

ksp-pipeline-generic-materialization-lib

Mission :

D2 input
   |
   v
GenericMaterializer fourni
   |
   v
generic output
   |
   v
D3 journal + processing outcome

Dépendances conceptuelles :

ksp-pipeline-generic-materialization-lib
    -> ksp-materializer-api
    -> ksp-store-api
    -> ksp-core-lib
    -> ksp-logging-lib

Il ne dépend pas de ksp-materializer-lib ni de ksp-store-lib.

ksp-pipeline-domain-projection-lib

Mission :

D3 input + contexte nécessaire
   |
   v
DomainProjector fourni
   |
   v
D4 projection + processing outcome

Dépendances conceptuelles :

ksp-pipeline-domain-projection-lib
    -> ksp-materializer-api
    -> ksp-store-api
    -> ksp-core-lib
    -> ksp-logging-lib

Il ne dépend pas de ksp-materializer-lib ni de ksp-store-lib.

Lorsqu'un projector est stateful, le pipeline/composant de composition charge via Store le contexte requis et le fournit au projector. Le projector reste indépendant du backend.

Worker API

ksp-worker-api reste une lifecycle API générique pour services continus.

Concepts candidats :

WorkerId
WorkerDescriptor
WorkerState
WorkerHealth
WorkerCapabilities

États conceptuels possibles :

Stopped
Starting
Running
Degraded
Stopping
Failed

Les noms Rust exacts ne sont pas figés.

Les opérations communes attendues sont au minimum :

start
stop
status
health

Une capability telle que reconfigure n'est pas imposée à tous les workers. Un worker expose les capacités qu'il supporte réellement.

Le lifecycle Worker ne contient aucune progression terminale propre aux jobs.

ksp-worker-raw-retriever

Mission

Service continu/live ou quasi-live :

subscription/fetch live
        |
        v
homogeneous transport model
        |
        v
ksp-pipeline-raw-ingestion-lib
        |
        v
D1

Il ne décode pas, ne matérialise pas, ne backfill pas et ne rejoue pas les niveaux dérivés.

Hot reconfiguration

Le raw retriever doit pouvoir modifier à chaud au moins les dimensions réellement supportées par son transport, par exemple :

  • Program IDs suivis ;
  • accounts/adresses suivis ;
  • types de données/subscriptions ;
  • filtres d'acquisition ;
  • endpoints/providers lorsque le transport permet une transition sûre.

La configuration runtime doit distinguer :

desired configuration generation
effective/applied configuration generation

Une reconfiguration est considérée appliquée seulement lorsque les subscriptions/filtres correspondants sont réellement actifs.

En cas d'échec partiel, le status/health doit pouvoir refléter la divergence entre desired et effective.

La stratégie exacte de diff/subscription est propre au worker/transport, pas à ksp-worker-api.

ksp-worker-core-processor

Service continu :

D1 backlog
   |
   v
ksp-pipeline-core-processing-lib
   |
   v
D2

Il utilise :

  • ksp-worker-api pour son lifecycle ;
  • ksp-store-lib comme backend officiel de composition ;
  • ksp-program-lib ou d'autres implémentations compatibles ;
  • le pipeline Core commun.

Il peut être réveillé par une notification D1 mais reconstruit toujours son backlog depuis le Store.

ksp-worker-generic-materializer

Service continu :

D2 backlog
   |
   v
ksp-pipeline-generic-materialization-lib
   |
   v
D3

Il compose :

  • ksp-store-lib ;
  • ksp-materializer-lib et/ou materializers compatibles ;
  • le pipeline générique.

Le backlog est suivi par input + materializer identity/version.

ksp-worker-domain-projector

Service continu :

D3 backlog
   |
   v
ksp-pipeline-domain-projection-lib
   |
   v
D4

Le nom reste provisoire, mais la responsabilité D3 -> D4 est retenue.

Il compose Store, projectors officiels/externes et pipeline de projection.

Un projector peut nécessiter un contexte existant D4 ou canonique. Le worker/pipeline charge ce contexte depuis le Store et le fournit explicitement ; ksp-materializer-lib ne dépend toujours pas du Store.

Job API

ksp-job-api reste distinct de ksp-worker-api.

Concepts candidats :

JobId
JobDescriptor
JobState
JobProgress
JobResult
JobCapabilities

Le job est déclenché, borné et terminable.

Les capacités suivantes peuvent exister lorsqu'elles sont pertinentes :

cancel
pause
resume
checkpoint

Elles ne sont pas obligatoirement universelles.

Aucune ksp-job-control-lib n'est introduite.

ksp-job-backfill

Mission

Acquisition historique uniquement :

historical fetch/pagination
        |
        v
homogeneous transport model
        |
        v
ksp-pipeline-raw-ingestion-lib
        |
        v
D1

Il ne produit pas D2/D3/D4.

Live acquisition et backfill produisent donc le même contrat D1.

Scope et progression

Un backfill doit pouvoir enregistrer durablement selon sa stratégie :

  • network ;
  • provider/source ;
  • catégorie d'acquisition ;
  • filtre/range demandé ;
  • pagination/cursor ;
  • slot/signature ranges lorsque pertinents ;
  • candidats vus/traités ;
  • records D1 persistés ;
  • erreurs/retries ;
  • checkpoint de reprise.

Le checkpoint backfill décrit également la progression dans une source externe ; il ne se réduit donc pas à un processing outcome D1.

Jobs de replay

Trois jobs distincts sont retenus :

ksp-job-replay-core
ksp-job-replay-generic-materialization
ksp-job-replay-domain-projection

Ils réutilisent les mêmes pipelines que les workers continus.

ksp-job-replay-core

scope D1
    |
    v
ksp-pipeline-core-processing-lib
    |
    v
D2

Il compose l'implémentation Program cible et son identité/version.

ksp-job-replay-generic-materialization

scope D2
    |
    v
ksp-pipeline-generic-materialization-lib
    |
    v
D3

Il cible un ou plusieurs materializers/versions.

ksp-job-replay-domain-projection

scope D3
    |
    v
ksp-pipeline-domain-projection-lib
    |
    v
D4

Il cible un ou plusieurs projectors/versions.

Scope de replay

Selon la frontière, un replay peut être borné par :

  • network ;
  • slot/range ;
  • Program ID ;
  • type de donnée ;
  • processor/materializer/projector ;
  • domaine ;
  • autres filtres canoniques disponibles.

Les scopes exacts sont définis par les contrats Store réels et non par une query PostgreSQL fuite dans ksp-job-api.

Processing outcome durable

L'absence d'output ne signifie jamais automatiquement qu'un input n'a pas été traité.

Un traitement valide peut produire zéro output parce qu'il est :

  • non applicable ;
  • unsupported par cette version ;
  • volontairement sans output ;
  • reconnu mais ignoré selon le contrat du processor.

Chaque frontière dérivée doit donc conserver un outcome durable associé conceptuellement à :

input identity
processor identity
processor version
logical capability/output identity
attempt/provenance
outcome
timestamps

Outcomes conceptuels possibles :

Produced
NoOutput
NotApplicable
Unsupported
FailedDeterministic

Les noms exacts restent ouverts.

Les erreurs transitoires ne doivent pas nécessairement devenir immédiatement un outcome terminal.

Backlog par processor/version

La vérité du backlog n'est pas :

"aucune row de sortie"

La vérité est conceptuellement :

inputs applicables
MINUS
outcomes terminaux pour processor/version/capability cible

Ainsi :

D1 X + CoreProcessor v1 -> Unsupported

est traité pour v1.

Plus tard :

D1 X + CoreProcessor v2

constitue un nouveau travail si v2 a changé sa capability.

Pour D2 -> D3 et D3 -> D4, un même input peut être candidat pour plusieurs materializers/projectors. Le backlog existe donc séparément pour chaque identité/version/capability.

Claim / lease

Pour supporter plusieurs instances concurrentes et la reprise après crash, un candidat peut être temporairement possédé :

candidate
   |
   v
claim
   |
   v
lease(owner, expires_at)
   |
   v
process
   |
   v
complete

Si l'instance meurt :

lease expires
    |
    v
candidate claimable again

La sémantique de claim appartient au contrat Store.

L'implémentation PostgreSQL choisira plus tard la technique SQL appropriée.

Le claim n'est pas une preuve de completion.

Sémantique de livraison

KSP préfère :

at-least-once processing
+
idempotent durable writes
+
durable outcomes

à une promesse distribuée exactly-once.

Après crash, un même input peut être retraité.

L'idempotence garantit que :

same logical input
+ same processor/version
+ same logical output identity

ne crée pas un second fait logique incohérent.

Atomicité de l'unité de traitement

Pour une unité logique de processing, doivent être cohérents/atomiques du point de vue durable :

output(s)
processing outcome
notification request/publication contract

Un outcome Completed/Produced ne peut pas être durable si seule une partie des outputs obligatoires a été persistée.

L'unité atomique est généralement un input, sauf lorsqu'un projector définit explicitement un groupe indivisible.

La transaction backend concrète appartient à ksp-store-lib; le pipeline l'oriente via ksp-store-api.

Erreurs de processing

Trois familles conceptuelles au minimum :

Transient

Exemples :

  • timeout ;
  • backend temporairement indisponible ;
  • rate limit ;
  • ressource temporairement indisponible.

Traitement :

retry + backoff

avec limite/configuration.

Deterministic failure

Exemples :

  • donnée malformée connue ;
  • invariant impossible ;
  • erreur de transformation reproductible.

Un état durable permet d'éviter une boucle infinie pour la même version.

Not applicable / unsupported

Ce ne sont pas nécessairement des erreurs.

Ils constituent des outcomes terminaux pour la version/capability considérée et peuvent être réévalués avec une nouvelle version.

Reprise après crash

Un worker/job de processing doit pouvoir suivre la séquence :

restart
   |
   v
load configuration/scope
   |
   v
recover/ignore expired claims
   |
   v
query backlog from Store
   |
   v
resume processing

Aucune notification perdue pendant l'arrêt ne doit rendre la reprise impossible.

Cursor et checkpoint

Un cursor peut accélérer le scan :

last scanned slot
last durable id
last page

mais il ne prouve pas qu'un input a été traité.

Pour les workers de processing :

durable processing outcomes = vérité
cursor = optimisation

Pour ksp-job-backfill, le checkpoint possède aussi une sémantique de progression dans la source externe.

Replay normal et replay forcé

Deux intentions doivent être distinguées.

Compléter/reprendre

Traiter uniquement ce qui manque/échoue selon le scope et la version cible.

Force replay

Retraiter même si un résultat actuel existe.

Un force replay :

  • ne supprime pas silencieusement l'historique ;
  • conserve la provenance ;
  • produit ou supersède les résultats selon les règles du niveau ;
  • reste idempotent pour l'identité logique choisie.

Les règles SQL exactes seront fixées avec les premiers modèles D2/D3/D4.

Notifications : mécanisme de référence

Le mécanisme initial de référence retenu pour PostgreSQL est :

LISTEN / NOTIFY

sans en faire une garantie de livraison.

Le modèle attendu est :

persist outputs/outcome
        |
        v
commit
        |
        v
NOTIFY / wake-up

L'implémentation exacte peut exploiter le comportement transactionnel PostgreSQL de notification.

Chaque consumer combine :

notification wake-up
+
periodic backlog polling

Le polling périodique constitue le filet de sécurité.

Le contrat public reste dans ksp-store-api et ne dépend pas de PostgreSQL.

Batching et concurrence

Un worker/job de processing peut suivre :

claim bounded batch
    |
    v
process with bounded concurrency
    |
    v
persist outcomes
    |
    v
repeat

Paramètres potentiels :

  • batch size ;
  • max concurrency ;
  • poll interval ;
  • retry/backoff ;
  • lease duration.

Ces paramètres ne sont pas nécessairement des champs universels de ksp-worker-api ou ksp-job-api.

Ils appartiennent au composant qui sait les interpréter.

Backpressure

Un backlog croissant n'est pas automatiquement une erreur :

D1 production rate > D2 processing rate

Le système doit pouvoir observer au minimum conceptuellement :

  • backlog count ;
  • âge du plus ancien input pending ;
  • processing rate ;
  • retry/failure rate ;
  • état degraded éventuel.

Le raw retriever n'est pas automatiquement ralenti par un backlog downstream. D1 sert précisément de tampon durable permettant de préserver l'acquisition.

Une stratégie de backpressure globale pourra être décidée plus tard par un manager/orchestrateur ou une configuration produit.

Logging

Tous les pipelines, workers et jobs runtime utilisent ksp-logging-lib.

Les logs doivent couvrir selon pertinence :

  • lifecycle ;
  • claims/releases ;
  • batch start/end ;
  • retries/backoff ;
  • outcomes unsupported/not applicable ;
  • failures ;
  • checkpoints ;
  • replay scopes ;
  • hot reconfiguration ;
  • divergences desired/effective ;
  • notifications/polling.

Les secrets et données sensibles ne sont jamais loggés.

Graphe de composition

Raw live

ksp-worker-raw-retriever
    -> ksp-worker-api
    -> ksp-pipeline-raw-ingestion-lib
    -> ksp-onchain-transport-lib
    -> ksp-store-lib
    -> ksp-config-lib
    -> ksp-logging-lib

Backfill

ksp-job-backfill
    -> ksp-job-api
    -> ksp-pipeline-raw-ingestion-lib
    -> ksp-onchain-transport-lib
    -> ksp-store-lib
    -> ksp-config-lib
    -> ksp-logging-lib

Core live/replay

ksp-worker-core-processor
ksp-job-replay-core
    -> lifecycle API correspondant
    -> ksp-pipeline-core-processing-lib
    -> ksp-program-lib ou implémentation compatible
    -> ksp-store-lib
    -> ksp-config-lib
    -> ksp-logging-lib

Generic materialization live/replay

ksp-worker-generic-materializer
ksp-job-replay-generic-materialization
    -> lifecycle API correspondant
    -> ksp-pipeline-generic-materialization-lib
    -> ksp-materializer-lib ou implémentation compatible
    -> ksp-store-lib
    -> ksp-config-lib
    -> ksp-logging-lib

Domain projection live/replay

ksp-worker-domain-projector
ksp-job-replay-domain-projection
    -> lifecycle API correspondant
    -> ksp-pipeline-domain-projection-lib
    -> ksp-materializer-lib ou implémentation compatible
    -> ksp-store-lib
    -> ksp-config-lib
    -> ksp-logging-lib

Questions laissées ouvertes

Les premières implémentations doivent encore fixer :

  • types Rust précis Worker/Job state, health et progress ;
  • modèle SQL exact de claim/lease ;
  • durée/renouvellement d'une lease ;
  • granularité des transactions par frontière ;
  • nombre de retries et politique de backoff ;
  • conflict handling entre plusieurs processor implementations ;
  • type exact des processing outcomes ;
  • contexte déclarable d'un DomainProjector stateful ;
  • stratégie de graceful shutdown d'un batch en cours ;
  • politique précise de pause/resume des jobs ;
  • métriques/export telemetry au-delà des logs et health.

La couche supérieure est désormais cadrée dans 010-APPS_SERVICES_SCENARIOS_AND_CONTROL.md. Le mécanisme IPC exact et l'éventuel orchestrateur restent à décider uniquement sur besoin concret.