18 KiB
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 D1–D4 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-apipour son lifecycle ;ksp-store-libcomme backend officiel de composition ;ksp-program-libou 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-libet/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
DomainProjectorstateful ; - 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.