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

905 lines
18 KiB
Markdown
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
<!-- file: docs/architecture/009-ACQUISITION_WORKERS_AND_JOBS.md -->
<!-- version: 2 -->
# 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 :
```text
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**.
```text
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 :
```text
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 :
```text
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 :
```text
D1 input
|
v
conversion vers ksp-program-api
|
v
decoder/registry fourni
|
v
Core output + processing outcome
|
v
persistence D2
```
Dépendances conceptuelles :
```text
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 :
```text
D2 input
|
v
GenericMaterializer fourni
|
v
generic output
|
v
D3 journal + processing outcome
```
Dépendances conceptuelles :
```text
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 :
```text
D3 input + contexte nécessaire
|
v
DomainProjector fourni
|
v
D4 projection + processing outcome
```
Dépendances conceptuelles :
```text
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 :
```text
WorkerId
WorkerDescriptor
WorkerState
WorkerHealth
WorkerCapabilities
```
États conceptuels possibles :
```text
Stopped
Starting
Running
Degraded
Stopping
Failed
```
Les noms Rust exacts ne sont pas figés.
Les opérations communes attendues sont au minimum :
```text
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 :
```text
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 :
```text
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 :
```text
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 :
```text
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 :
```text
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 :
```text
JobId
JobDescriptor
JobState
JobProgress
JobResult
JobCapabilities
```
Le job est déclenché, borné et terminable.
Les capacités suivantes peuvent exister lorsqu'elles sont pertinentes :
```text
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 :
```text
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 :
```text
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`
```text
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`
```text
scope D2
|
v
ksp-pipeline-generic-materialization-lib
|
v
D3
```
Il cible un ou plusieurs materializers/versions.
## `ksp-job-replay-domain-projection`
```text
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 à :
```text
input identity
processor identity
processor version
logical capability/output identity
attempt/provenance
outcome
timestamps
```
Outcomes conceptuels possibles :
```text
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 :
```text
"aucune row de sortie"
```
La vérité est conceptuellement :
```text
inputs applicables
MINUS
outcomes terminaux pour processor/version/capability cible
```
Ainsi :
```text
D1 X + CoreProcessor v1 -> Unsupported
```
est traité pour v1.
Plus tard :
```text
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é :
```text
candidate
|
v
claim
|
v
lease(owner, expires_at)
|
v
process
|
v
complete
```
Si l'instance meurt :
```text
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 :
```text
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 :
```text
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 :
```text
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 :
```text
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 :
```text
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 :
```text
last scanned slot
last durable id
last page
```
mais il ne prouve pas qu'un input a été traité.
Pour les workers de processing :
```text
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 :
```text
LISTEN / NOTIFY
```
sans en faire une garantie de livraison.
Le modèle attendu est :
```text
persist outputs/outcome
|
v
commit
|
v
NOTIFY / wake-up
```
L'implémentation exacte peut exploiter le comportement transactionnel PostgreSQL de notification.
Chaque consumer combine :
```text
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 :
```text
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 :
```text
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
```text
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
```text
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
```text
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
```text
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
```text
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.