399 lines
11 KiB
Markdown
399 lines
11 KiB
Markdown
<!-- file: docs/architecture/009-ACQUISITION_WORKERS_AND_JOBS.md -->
|
|
<!-- version: 10 -->
|
|
|
|
# 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-lib`
|
|
|
|
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 identité, lifecycle, annulation abstraite et observation latest-value ;
|
|
- expose `LatestAddress`, `BeforeAddress`, `AfterAddress` et `ExplicitSignatures` avec bornes explicites de pages/candidats/concurrence ;
|
|
- porte explicitement le réseau logique du Store dans le scope et dans l'identité `(network, signature)` de chaque transaction candidate ;
|
|
- exclut rôle HTTP, provider, endpoint et protocole du fingerprint sémantique et de l'identité transactionnelle ;
|
|
- hydrate uniquement via la voie observée `getTransaction`, afin de conserver la provenance du provider/endpoint réellement gagnant ;
|
|
- produit un RAW v1 canonique déterministe puis persiste transaction + observation atomiquement via `ksp-store-lib` en mode normal ;
|
|
- respecte les tombstones `Purged`, distingue missing/conflit/idempotence et ne pré-lit pas le Store avant hydratation ;
|
|
- limite les hydrations concurrentes, avance seulement une frontier contiguë durable et retourne un checkpoint opaque caller-owned ;
|
|
- arrête coopérativement les nouvelles admissions lors d'une annulation et draine une persistence Store déjà soumise ;
|
|
- publie des snapshots latest-value sûrs sans payload RAW ni secrets/URLs Transport ;
|
|
- n'effectue aucun décodage Program et n'écrit aucun fait 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 et volontairement runtime-neutral.
|
|
|
|
Contrats communs actuels :
|
|
|
|
```text
|
|
JobId
|
|
JobKindCode
|
|
JobState
|
|
JobCompletion
|
|
JobLifecycle
|
|
JobCancellationToken
|
|
JobNotificationSequence
|
|
JobNotification<Snapshot>
|
|
JobSnapshotSource
|
|
```
|
|
|
|
Le lifecycle commun est borné aux transitions explicitement validées entre `Created`, `Running`, `Cancelling` et les états terminaux `Completed(Complete|Partial)`, `Cancelled`, `Failed`. Il ne définit ni `pause`, ni `resume`, ni scheduler, ni runtime d'exécution générique.
|
|
|
|
`JobSnapshotSource` suit une sémantique latest-value : un listener lit une valeur complète courante puis peut attendre une séquence plus récente ; les valeurs intermédiaires peuvent être coalescées. Le snapshot métier reste possédé par le job concret.
|
|
|
|
L'annulation commune est une intention coopérative. Le job concret décide quelles attentes peuvent être interrompues et quelles opérations engagées doivent être drainées.
|
|
|
|
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 taille du batch, la priorité et la stratégie de sélection appartiennent au worker/job/executor. Store fournit les primitives de query/cursor/claim nécessaires et peut refléter une contrainte physique du backend, mais `ksp-store-lib` n'invente pas un plafond métier global inférieur à la capacité réellement disponible.
|
|
|
|
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 qui promet une reprise après crash doit reconstruire son état depuis des données durables : inputs persistés, claims/leases/outcomes lorsqu'ils existent, cursors/checkpoints et versions de processor.
|
|
|
|
Le premier backfill RAW retourne un `BackfillCheckpoint` caller-owned lié au `JobId` et au fingerprint de scope. La crate ne persiste pas ce checkpoint elle-même : tant qu'un caller ne l'enregistre pas durablement, il s'agit d'une primitive de reprise contrôlée, pas d'une promesse crash-safe automatique.
|
|
|
|
La mémoire du processus ne constitue jamais l'unique source d'une garantie de reprise durable.
|
|
|
|
## 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-lib
|
|
-> ksp-job-api
|
|
-> ksp-core-lib
|
|
-> ksp-logging-lib
|
|
-> ksp-onchain-transport-lib
|
|
-> ksp-store-lib # façade Store ; default-features=false côté Job
|
|
-> futures-util/tokio # runtime privé de Backfill
|
|
-> serde_json/sha2 # RAW v1 canonique + digest
|
|
```
|
|
|
|
### RAW worker
|
|
|
|
```text
|
|
ksp-worker-raw-retriever
|
|
-> ksp-worker-api
|
|
-> ksp-onchain-transport-lib
|
|
-> ksp-interface-lib # seulement si un fait passif partagé aide la composition live
|
|
-> ksp-store-lib # façade Store ; backend sélectionné par feature + Config
|
|
-> ksp-config-lib
|
|
-> ksp-logging-lib
|
|
```
|
|
|
|
Les événements Interface peuvent servir de signal provider-neutral à la composition live, mais ne constituent jamais le backlog durable. Après crash ou perte d'un événement, la reprise s'appuie sur Store et sur les primitives de replay/hydratation appropriées.
|
|
|
|
### 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 ;
|
|
- modèle de claim/lease PostgreSQL pour les futurs processors continus ;
|
|
- 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.
|