394 lines
7.7 KiB
Markdown
394 lines
7.7 KiB
Markdown
<!-- file: docs/architecture/009-ACQUISITION_WORKERS_AND_JOBS.md -->
|
|
<!-- version: 4 -->
|
|
|
|
# Acquisition, workers, jobs et pipelines spécialisés
|
|
|
|
## Objet
|
|
|
|
Ce document définit le lifecycle opérationnel autour des couches :
|
|
|
|
```text
|
|
RAW -> CORE -> DECODE -> SPECIALIZED
|
|
```
|
|
|
|
Il conserve la séparation stricte entre :
|
|
|
|
- **pipeline/processor** : logique réutilisable d'une transformation ;
|
|
- **worker** : service continu/live ;
|
|
- **job** : traitement borné/historique/replay ;
|
|
- **application** : contrôle/visualisation, sans duplication de la logique métier.
|
|
|
|
## Règle de progression
|
|
|
|
RAW et CORE sont les deux premières couches horizontales. Elles ne nécessitent aucun decoder Program.
|
|
|
|
Pour chacune, KSP peut terminer successivement :
|
|
|
|
```text
|
|
persistence
|
|
-> backfill/replay
|
|
-> worker live si nécessaire
|
|
-> application de contrôle/inspection si utile
|
|
```
|
|
|
|
À partir de DECODE, les processors/jobs/workers/scenarios sont introduits **avec le groupe Program concerné**, en vertical slice, au lieu de créer à l'avance une grande flotte générique de workers de décodage/materialisation sans programme réel.
|
|
|
|
## RAW
|
|
|
|
### Pipeline RAW ingestion
|
|
|
|
Concept :
|
|
|
|
```text
|
|
homogeneous transport model
|
|
|
|
|
v
|
|
conversion RAW DTO
|
|
|
|
|
v
|
|
ksp-store-api
|
|
|
|
|
v
|
|
persistence D1 RAW
|
|
|
|
|
v
|
|
notification after commit
|
|
```
|
|
|
|
Une crate spécialisée `ksp-pipeline-raw-ingestion-lib` peut être introduite lorsque la réutilisation worker + job le justifie réellement.
|
|
|
|
Elle ne choisit pas le provider réseau et ne pilote pas le range historique.
|
|
|
|
### `ksp-job-backfill`
|
|
|
|
Le premier backfill historique appartient à la couche RAW :
|
|
|
|
```text
|
|
range/cursor historique
|
|
|
|
|
v
|
|
ksp-onchain-transport-lib
|
|
|
|
|
v
|
|
RAW ingestion
|
|
|
|
|
v
|
|
D1 RAW
|
|
```
|
|
|
|
Le job :
|
|
|
|
- utilise `ksp-job-api` pour son lifecycle ;
|
|
- gère scope/range/pagination/checkpoint ;
|
|
- n'effectue aucun décodage Program ;
|
|
- n'écrit pas directement des faits CORE/DECODE/SPECIALIZED.
|
|
|
|
### Worker RAW live
|
|
|
|
Le worker live futur :
|
|
|
|
```text
|
|
subscriptions/fetch live
|
|
|
|
|
v
|
|
ksp-onchain-transport-lib
|
|
|
|
|
v
|
|
RAW ingestion
|
|
|
|
|
v
|
|
D1 RAW
|
|
```
|
|
|
|
Il ne décode pas et ne matérialise pas.
|
|
|
|
Sa hot reconfiguration peut concerner selon le transport :
|
|
|
|
- endpoints/providers ;
|
|
- rôles ;
|
|
- accounts/program IDs suivis ;
|
|
- types de subscriptions ;
|
|
- filtres ;
|
|
- limites/concurrence.
|
|
|
|
La configuration desired/effective reste distinguée lorsqu'une reconfiguration est asynchrone.
|
|
|
|
## CORE
|
|
|
|
### Pipeline RAW -> CORE
|
|
|
|
La normalisation CORE est générique Solana :
|
|
|
|
```text
|
|
D1 RAW
|
|
|
|
|
v
|
|
Solana generic normalizer
|
|
|
|
|
v
|
|
D2 CORE
|
|
```
|
|
|
|
Dépendances interdites :
|
|
|
|
```text
|
|
CORE normalizer -X-> ksp-program-api
|
|
CORE normalizer -X-> ksp-program-lib
|
|
CORE normalizer -X-> ksp-materializer-api
|
|
```
|
|
|
|
`ksp-interface-lib` peut fournir les wires Solana génériques nécessaires à la structure blockchain.
|
|
|
|
### Job replay CORE
|
|
|
|
Un job borné peut rejouer :
|
|
|
|
```text
|
|
RAW persisted range
|
|
|
|
|
v
|
|
CORE normalizer
|
|
|
|
|
v
|
|
D2 CORE
|
|
```
|
|
|
|
sans redemander les données au réseau.
|
|
|
|
### Worker CORE
|
|
|
|
Un worker CORE continu peut consommer le backlog RAW nouvellement persisté et produire CORE.
|
|
|
|
Le Store reste source de vérité du backlog ; les notifications ne sont qu'un wake-up.
|
|
|
|
## DECODE et SPECIALIZED
|
|
|
|
### Introduction par groupe fonctionnel
|
|
|
|
KSP ne crée pas d'abord un unique « worker decoder de tout Solana » puis tous les materializers plusieurs séries plus tard.
|
|
|
|
Chaque groupe prioritaire introduit les capacités nécessaires :
|
|
|
|
```text
|
|
CORE inputs du groupe
|
|
|
|
|
v
|
|
Program decoder
|
|
|
|
|
v
|
|
decoded facts
|
|
|
|
|
v
|
|
generic materialization / DECODE persistence
|
|
|
|
|
v
|
|
SPECIALIZED projection si utile
|
|
```
|
|
|
|
Puis le même groupe avance vers préparation d'exécution, policy, execution et scénarios Devnet.
|
|
|
|
Les jobs de replay et workers live correspondants réutilisent les mêmes processors du groupe.
|
|
|
|
### Ordre fonctionnel
|
|
|
|
Direction actuelle :
|
|
|
|
```text
|
|
Solana Core Programs
|
|
-> SPL token/trading
|
|
-> token metadata
|
|
-> Anchor
|
|
-> Meteora
|
|
-> Raydium
|
|
-> Pump
|
|
-> Orca
|
|
-> Market Desk V1
|
|
-> Jupiter / OKX routing
|
|
-> Market Desk V2
|
|
-> trading-adjacent
|
|
-> general decoding
|
|
```
|
|
|
|
Un satellite protocolaire reste avec son groupe : Meteora vaults avec Meteora, Pump fee avec Pump, etc.
|
|
|
|
## Worker API
|
|
|
|
`ksp-worker-api` reste une lifecycle API générique pour services continus.
|
|
|
|
Concepts candidats :
|
|
|
|
```text
|
|
WorkerId
|
|
WorkerDescriptor
|
|
WorkerState
|
|
WorkerHealth
|
|
WorkerCapabilities
|
|
```
|
|
|
|
Opérations minimales candidates :
|
|
|
|
```text
|
|
start
|
|
stop
|
|
status
|
|
health
|
|
```
|
|
|
|
Une capability comme `reconfigure` n'est pas imposée à tous les workers.
|
|
|
|
## Job API
|
|
|
|
`ksp-job-api` reste distinct de Worker API.
|
|
|
|
Concepts candidats :
|
|
|
|
```text
|
|
JobId
|
|
JobDescriptor
|
|
JobState
|
|
JobProgress
|
|
JobOutcome
|
|
JobCapabilities
|
|
```
|
|
|
|
Un job est borné/terminable et peut exposer selon besoin :
|
|
|
|
```text
|
|
start
|
|
pause
|
|
resume
|
|
cancel
|
|
status
|
|
progress
|
|
```
|
|
|
|
Les types exacts sont décidés à `0.3.3` avec le premier vrai backfill.
|
|
|
|
Aucune `ksp-job-control-lib` n'est créée sans duplication concrète.
|
|
|
|
## Backlog, claim et idempotence
|
|
|
|
Pour les processors asynchrones :
|
|
|
|
```text
|
|
notification wake-up
|
|
+
|
|
periodic polling
|
|
|
|
|
v
|
|
Store backlog query
|
|
|
|
|
v
|
|
claim/lease bounded batch
|
|
|
|
|
v
|
|
processor
|
|
|
|
|
v
|
|
output + durable outcome
|
|
```
|
|
|
|
Les notifications ne remplacent jamais le Store.
|
|
|
|
La sémantique cible reste at-least-once avec idempotence durable, plutôt qu'un faux exactly-once.
|
|
|
|
## Processing outcomes
|
|
|
|
Les états exacts seront définis avec le premier processor durable, mais doivent distinguer au minimum les familles conceptuelles :
|
|
|
|
- success ;
|
|
- transient failure ;
|
|
- deterministic failure ;
|
|
- not applicable ;
|
|
- unsupported ;
|
|
- superseded/replayed lorsque pertinent.
|
|
|
|
## Reprise après crash
|
|
|
|
Un worker/job doit reconstruire son état depuis :
|
|
|
|
- inputs persistés ;
|
|
- claims/leases ;
|
|
- outcomes ;
|
|
- cursors/checkpoints ;
|
|
- versions de processor.
|
|
|
|
La mémoire du processus ne constitue jamais l'unique source de reprise.
|
|
|
|
## Logging
|
|
|
|
Tous les pipelines/workers/jobs runtime utilisent `ksp-logging-lib`.
|
|
|
|
Les événements utiles comprennent notamment :
|
|
|
|
- start/stop ;
|
|
- batch/range ;
|
|
- progression ;
|
|
- retries ;
|
|
- claim/lease ;
|
|
- rate-limit/backpressure ;
|
|
- processing outcome ;
|
|
- checkpoint ;
|
|
- reconfiguration desired/effective ;
|
|
- erreurs redacted.
|
|
|
|
## Dépendances de composition
|
|
|
|
### RAW backfill
|
|
|
|
```text
|
|
ksp-job-backfill
|
|
-> ksp-job-api
|
|
-> ksp-onchain-transport-lib
|
|
-> ksp-store-api / backend injecté
|
|
-> ksp-config-lib # orchestration/config, pas ownership transport
|
|
-> ksp-logging-lib
|
|
```
|
|
|
|
### RAW worker
|
|
|
|
```text
|
|
ksp-worker-raw-retriever
|
|
-> ksp-worker-api
|
|
-> ksp-onchain-transport-lib
|
|
-> ksp-store-api / backend injecté
|
|
-> ksp-config-lib
|
|
-> ksp-logging-lib
|
|
```
|
|
|
|
### CORE replay/worker
|
|
|
|
```text
|
|
CORE processor
|
|
-> ksp-interface-lib si wires génériques nécessaires
|
|
-> ksp-store-api
|
|
-> ksp-core-lib
|
|
-> ksp-logging-lib
|
|
```
|
|
|
|
Pas de Program API.
|
|
|
|
### Groupes DECODE/SPECIALIZED
|
|
|
|
Le composant de composition du groupe peut utiliser :
|
|
|
|
```text
|
|
ksp-program-api
|
|
ksp-program-lib ou extension compatible
|
|
ksp-materializer-api
|
|
ksp-materializer-lib ou implementation compatible
|
|
ksp-store-api
|
|
```
|
|
|
|
selon les capacités réellement introduites.
|
|
|
|
## Questions laissées ouvertes
|
|
|
|
- nom final de la crate pipeline RAW si la réutilisation justifie une crate dédiée ;
|
|
- nom final du worker RAW ;
|
|
- contrat exact de `ksp-job-api` ;
|
|
- modèle de claim/lease PostgreSQL ;
|
|
- taille de batch et stratégie backpressure ;
|
|
- découpage des workers DECODE/SPECIALIZED par groupe lorsque les premiers groupes existent ;
|
|
- mécanisme IPC des applications de contrôle futures.
|