361 lines
17 KiB
Markdown
361 lines
17 KiB
Markdown
<!-- file: crates/ksp-worker-raw-transaction-ingest-lib/USAGE.md -->
|
|
<!-- version: 5 -->
|
|
|
|
# Utilisation de ksp-worker-raw-transaction-ingest-lib
|
|
|
|
Cette page décrit la façade publique de `ksp-worker-raw-transaction-ingest-lib`. Le caller possède la composition du runtime Tokio, du `Store` et des ressources Transport ; le Worker ne lit pas Config, ne construit pas un backend physique et ne lit pas de secret depuis l'environnement.
|
|
|
|
## Construire les settings
|
|
|
|
Les identités sont validées par leurs couches propriétaires :
|
|
|
|
```rust
|
|
fn worker_settings() -> ksp_core_lib::Result<ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestSettings> {
|
|
let network = match ksp_store_lib::RawNetworkId::new("mainnet") {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
|
};
|
|
let worker_id = match ksp_worker_api::WorkerId::new("raw-ingest-mainnet-0001") {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
|
};
|
|
return std::result::Result::Ok(ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestSettings::with_defaults(network, worker_id));
|
|
}
|
|
```
|
|
|
|
Le Worker kind est fixe :
|
|
|
|
```rust
|
|
assert_eq!(ksp_worker_raw_transaction_ingest_lib::RAW_TRANSACTION_INGEST_WORKER_KIND_CODE, "raw_transaction_ingest");
|
|
```
|
|
|
|
Pour des limites explicites :
|
|
|
|
```rust
|
|
let settings = ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestSettings::new(
|
|
network,
|
|
worker_id,
|
|
512,
|
|
16,
|
|
std::time::Duration::from_secs(10),
|
|
);
|
|
```
|
|
|
|
Bornes publiques :
|
|
|
|
```text
|
|
admission_queue_capacity 1 ..= 65_536 défaut 256
|
|
persistence_concurrency 1 ..= 64 défaut 8
|
|
shutdown_drain_timeout 100 ms ..= 30 s défaut 5 s
|
|
```
|
|
|
|
Une valeur hors borne retourne `worker_raw_transaction_ingest.settings_invalid` avec uniquement le nom stable du champ invalide.
|
|
|
|
## Choisir le mode de démarrage
|
|
|
|
### Fondation sans source productive
|
|
|
|
`RawTransactionIngestWorker::start(settings, store)` démarre la fondation runtime sans source Transport. Ce mode reste utile aux tests/compositions qui veulent uniquement le lifecycle, les snapshots et le contrat de shutdown.
|
|
|
|
```rust
|
|
let handle = match ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestWorker::start(settings, store) {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
|
};
|
|
```
|
|
|
|
### Source Yellowstone productive
|
|
|
|
Pour l'ingestion live, le caller compose d'abord les ressources via les crates propriétaires, puis construit :
|
|
|
|
```rust
|
|
fn runtime_resources(
|
|
yellowstone_channel: ksp_onchain_transport_lib::YellowstoneGrpcChannel,
|
|
subscribe_request: ksp_onchain_transport_lib::YellowstoneSubscribeRequest,
|
|
http_pool: ksp_onchain_transport_lib::HttpTransportPool,
|
|
hydration_role: ksp_onchain_transport_lib::HttpRoleName,
|
|
) -> ksp_core_lib::Result<ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestRuntimeResources> {
|
|
let source = match ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestYellowstoneSource::new(
|
|
yellowstone_channel,
|
|
subscribe_request,
|
|
http_pool,
|
|
hydration_role,
|
|
) {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
|
};
|
|
return std::result::Result::Ok(ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestRuntimeResources::new(source));
|
|
}
|
|
```
|
|
|
|
Puis :
|
|
|
|
```rust
|
|
let handle = match ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestWorker::start_with_runtime_resources(
|
|
settings,
|
|
store,
|
|
runtime_resources,
|
|
) {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
|
};
|
|
```
|
|
|
|
Les deux entrées exigent un runtime Tokio courant et un `Store` portant exactement le même `RawNetworkId` que les settings. Le démarrage avec ressources exige également que chaque source composée cible ce même réseau. Une exécution productive accepte une source unique Yellowstone, Standard Logs, Standard Block ou Helius Transaction ; une collection de plusieurs sources est validable à la composition mais reste rejetée avant spawn tant que la supervision simultanée n'est pas disponible.
|
|
|
|
### Source Standard Logs productive
|
|
|
|
Une source standard Solana WS se compose ainsi :
|
|
|
|
```rust
|
|
fn standard_logs_runtime_resources(
|
|
ws_endpoint: ksp_onchain_transport_lib::WsEndpointSettings,
|
|
filter: ksp_onchain_transport_lib::SolanaLogsSubscribeFilter,
|
|
commitment: ksp_onchain_transport_lib::SolanaCommitment,
|
|
http_pool: ksp_onchain_transport_lib::HttpTransportPool,
|
|
hydration_role: ksp_onchain_transport_lib::HttpRoleName,
|
|
) -> ksp_core_lib::Result<ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestRuntimeResources> {
|
|
let source = match ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestStandardLogsSource::new(
|
|
ws_endpoint,
|
|
filter,
|
|
commitment,
|
|
http_pool,
|
|
hydration_role,
|
|
) {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
|
};
|
|
return std::result::Result::Ok(ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestRuntimeResources::from_standard_logs_source(source));
|
|
}
|
|
```
|
|
|
|
`ws_endpoint` doit être un endpoint Transport valide de kind `solana_standard`. Le commitment doit être explicitement `Confirmed` ou `Finalized`. Le pool HTTP doit exposer `getTransaction` via le rôle indiqué sur le même réseau. `All`, `AllWithVotes` et `Mentions(pubkey)` sont acceptés par le contrat Transport ; la valeur du filtre reste privée dans le Worker.
|
|
|
|
### Source Standard Block productive
|
|
|
|
Une source `blockSubscribe` standard se compose sans pool HTTP :
|
|
|
|
```rust
|
|
fn standard_block_runtime_resources(
|
|
ws_endpoint: ksp_onchain_transport_lib::WsEndpointSettings,
|
|
filter: ksp_onchain_transport_lib::SolanaBlockSubscribeFilter,
|
|
commitment: ksp_onchain_transport_lib::SolanaCommitment,
|
|
) -> ksp_core_lib::Result<ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestRuntimeResources> {
|
|
let source = match ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestStandardBlockSource::new(ws_endpoint, filter, commitment) {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
|
};
|
|
return std::result::Result::Ok(ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestRuntimeResources::from_standard_block_source(source));
|
|
}
|
|
```
|
|
|
|
Le Worker demande `Base64`, `Full`, `maxSupportedTransactionVersion = 1` et `showRewards = false`. Une transaction n'est RAW-direct que si sa version est explicitement `Legacy`, `0` ou `1`. Une version omise/nulle ou supérieure, `block: null`, une erreur de bloc, un champ transactions absent/nul ou une transaction non Base64 provoque une faute sûre ; ces cas ne sont jamais assimilés à une progression vide.
|
|
|
|
### Source Helius Transaction productive
|
|
|
|
Une source Helius LaserStream `transactionSubscribe` se compose avec le même pattern caller-owned :
|
|
|
|
```rust
|
|
fn helius_transaction_runtime_resources(
|
|
ws_endpoint: ksp_onchain_transport_lib::WsEndpointSettings,
|
|
filter: ksp_onchain_transport_lib::HeliusTransactionSubscribeFilter,
|
|
commitment: ksp_onchain_transport_lib::SolanaCommitment,
|
|
http_pool: ksp_onchain_transport_lib::HttpTransportPool,
|
|
hydration_role: ksp_onchain_transport_lib::HttpRoleName,
|
|
) -> ksp_core_lib::Result<ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestRuntimeResources> {
|
|
let source = match ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestHeliusTransactionSource::new(
|
|
ws_endpoint,
|
|
filter,
|
|
commitment,
|
|
http_pool,
|
|
hydration_role,
|
|
) {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
|
};
|
|
return std::result::Result::Ok(ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestRuntimeResources::from_helius_transaction_source(source));
|
|
}
|
|
```
|
|
|
|
`ws_endpoint` doit être un endpoint Transport de kind `helius_laserstream`. Le caller supérieur résout éventuellement `KSP_SECRET_HELIUS_API_KEY` via Config avant de construire l'endpoint ; le Worker ne lit jamais l'environnement ni Config. Le commitment est limité à `Confirmed`/`Finalized` et la route HTTP doit supporter `getTransaction` sur le même réseau.
|
|
|
|
Le Worker demande la forme Helius `Full` avec `Base64`, `showRewards = false` et `maxSupportedTransactionVersion = 1`, mais ne fait pas confiance au nested payload pour construire directement le Common RAW. Il conserve seulement signature/slot/index et hydrate par `getTransaction observed`. Une notification d'une autre forme est fail-closed.
|
|
|
|
## Préparer la source Yellowstone
|
|
|
|
La `YellowstoneSubscribeRequest` doit :
|
|
|
|
- être valide selon Transport ;
|
|
- contenir au moins une famille transaction-bearing admise par le Worker ;
|
|
- utiliser explicitement `Confirmed` ou `Finalized` ;
|
|
- rester compatible avec le réseau du `YellowstoneGrpcChannel`.
|
|
|
|
Le `HttpTransportPool` doit posséder au moins une route compatible avec le rôle d'hydration et `getTransaction` sur le même réseau. La construction de `RawTransactionIngestYellowstoneSource` vérifie ces invariants sans ouvrir la connexion réseau.
|
|
|
|
Le caller ne passe pas de signature, `program_id`, plage de slots ou limite historique au Worker. Ces paramètres appartiennent à un Job Backfill, pas au service continu.
|
|
|
|
## Comprendre la pipeline live
|
|
|
|
Les familles productives sont traitées ainsi :
|
|
|
|
```text
|
|
Yellowstone Transaction -> signal -> HTTP getTransaction -> Common RAW -> admission
|
|
Yellowstone TransactionStatus -> signal -> HTTP getTransaction -> Common RAW -> admission
|
|
Yellowstone Block -> un signal par transaction -> HTTP getTransaction -> Common RAW -> admission
|
|
Standard WS logsSubscribe -> context.slot + signature -> HTTP getTransaction -> Common RAW -> admission
|
|
Standard WS blockSubscribe -> Full/Base64 Legacy|V0|V1 -> Common RAW direct par transaction -> admission
|
|
Helius transactionSubscribe -> Full envelope -> signature/slot/index -> HTTP getTransaction -> Common RAW -> admission
|
|
Yellowstone BlockMeta -> continuity-only
|
|
Yellowstone Slot -> continuity-only
|
|
Yellowstone Account/Ping/Pong/Entry -> sans RAW Transaction dans cette verticale
|
|
```
|
|
|
|
Les signaux de même `(network, signature, commitment)` sont coalescés avant l'hydration HTTP. Le Worker ne possède pas une boucle de retry HTTP : reroutage/retry/backoff restent dans `ksp-onchain-transport-lib`.
|
|
|
|
## Observer le snapshot concret
|
|
|
|
```rust
|
|
let source = handle.snapshot_source();
|
|
let current = source.current();
|
|
|
|
let sequence = current.worker_snapshot().sequence();
|
|
let state = current.worker_snapshot().state();
|
|
let queue_depth = current.admission_queue_depth();
|
|
let in_flight = current.in_flight_persistence();
|
|
let hydration_pending = current.hydration_pending();
|
|
let frontier = current.processing_frontier_slot();
|
|
let oldest_pending = current.oldest_pending_slot();
|
|
let source_state = current.source_state();
|
|
```
|
|
|
|
Pour attendre une valeur plus récente :
|
|
|
|
```rust
|
|
let observed = source.current().worker_snapshot().sequence();
|
|
let newer = source.wait_for_change(observed).await;
|
|
assert!(newer.worker_snapshot().sequence().is_after(observed));
|
|
```
|
|
|
|
Le flux est latest-value : les transitions intermédiaires peuvent être coalescées. Ne pas l'utiliser comme journal exhaustif.
|
|
|
|
## Utiliser la projection Worker API
|
|
|
|
La même source implémente `ksp_worker_api::WorkerSnapshotSource` :
|
|
|
|
```rust
|
|
let source = handle.worker_snapshot_source();
|
|
let current = ksp_worker_api::WorkerSnapshotSource::current(&source);
|
|
let observed = current.sequence();
|
|
let newer = ksp_worker_api::WorkerSnapshotSource::wait_for_change(&source, observed).await;
|
|
assert!(newer.sequence().is_after(observed));
|
|
```
|
|
|
|
La projection commune contient seulement les dimensions génériques Worker. Les compteurs et frontiers spécifiques restent sur `RawTransactionIngestSnapshot`.
|
|
|
|
## Lire les compteurs
|
|
|
|
```rust
|
|
let snapshot = handle.snapshot_source().current();
|
|
|
|
let admitted = snapshot.admitted_total();
|
|
let canonicalized = snapshot.canonicalized_total();
|
|
let persisted = snapshot.persisted_total();
|
|
let inserted = snapshot.entity_inserted_total();
|
|
let existing = snapshot.entity_already_present_total();
|
|
let purged = snapshot.entity_skipped_purged_total();
|
|
let observations = snapshot.observation_inserted_total();
|
|
let conflicts = snapshot.content_conflict_total();
|
|
let store_failures = snapshot.store_failure_total();
|
|
let source_failures = snapshot.source_failure_total();
|
|
let backpressure = snapshot.backpressure_wait_total();
|
|
let reconnects = snapshot.source_reconnect_total();
|
|
let replay_attempts = snapshot.source_replay_attempt_total();
|
|
let proven_gaps = snapshot.source_continuity_gap_total();
|
|
```
|
|
|
|
`admission_queue_depth()`, `in_flight_persistence()` et `hydration_pending()` sont des gauges latest-value. Les compteurs cumulés ne wrapent jamais silencieusement ; l'épuisement est terminal avec `worker_raw_transaction_ingest.counter_exhausted`.
|
|
|
|
## Interpréter la processing frontier
|
|
|
|
`processing_frontier_slot()` est la plus haute slot de travail source réellement observé qui n'est pas bloquée par un pending plus ancien connu. `oldest_pending_slot()` expose ce plus ancien pending lorsqu'il existe.
|
|
|
|
Cette frontier est strictement run-local :
|
|
|
|
```text
|
|
elle ne prouve pas que toutes les transactions blockchain d'une slot ont été observées
|
|
elle ne prouve pas la persistence durable des ingress déjà envoyés à l'admission
|
|
elle n'est pas persistée entre deux runs
|
|
elle n'est pas un checkpoint Backfill
|
|
```
|
|
|
|
Un `Missing` HTTP règle le signal du point de vue source-processing sans créer de RAW. Un ingress `Available` n'est réglé qu'après envoi réussi vers l'admission centrale. Une hydration abandonnée au shutdown est retirée des pending sans faire avancer artificiellement la frontier.
|
|
|
|
## Interpréter reconnect et replay
|
|
|
|
`source_state()` peut retourner `Active`, `Reconnecting`, `Closing`, `Closed` ou `Failed` après démarrage de la source productive.
|
|
|
|
Les compteurs ont des sémantiques distinctes :
|
|
|
|
```text
|
|
source_reconnect_total reconnects automatiques réussis observés
|
|
source_replay_attempt_total tentatives de reconnect portant une demande de replay
|
|
source_continuity_gap_total gaps de rétention prouvés par Transport
|
|
```
|
|
|
|
Une tentative de replay n'est pas une preuve de continuité. Le Worker ne choisit pas `from_slot` et ne traite pas directement `SubscribeReplayInfo` ; ces mécanismes appartiennent à Transport.
|
|
|
|
Si `source_continuity_gap_total` augmente, le Worker fault avec `worker_raw_transaction_ingest.source_failed`. Il ne déclenche pas automatiquement `ksp-job-backfill-lib`.
|
|
|
|
## Demander un stop et attendre le terminal
|
|
|
|
```rust
|
|
let accepted = handle.request_stop();
|
|
let terminal = handle.wait_terminal().await;
|
|
|
|
match terminal {
|
|
std::result::Result::Ok(state) => {
|
|
assert!(state.is_terminal());
|
|
}
|
|
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
|
}
|
|
|
|
if accepted {
|
|
// La première demande de stop a été remise au runtime vivant.
|
|
}
|
|
```
|
|
|
|
`request_stop()` est idempotent. Le terminal n'est publié qu'après le drain borné et la récupération des tâches possédées.
|
|
|
|
## Interpréter les faults
|
|
|
|
```text
|
|
settings_invalid settings techniques hors contrat
|
|
runtime_invalid invariant runtime/lifecycle impossible
|
|
store_failed erreur Store non-conflict classifiée
|
|
content_conflict contenu canonique incompatible avec l'identité durable
|
|
counter_exhausted compteur/séquence monotone arrivé à sa borne
|
|
source_failed source/replay/hydration terminée par une erreur classifiée
|
|
drain_timeout drain de shutdown hors deadline
|
|
```
|
|
|
|
Un consumer doit traiter l'`ErrorCode` comme contrat stable et ne pas dépendre d'un texte backend/provider arbitraire.
|
|
|
|
## Composition supérieure
|
|
|
|
Le pattern attendu est :
|
|
|
|
```text
|
|
Config / application / service owner
|
|
-> résout endpoints, credentials et rôles
|
|
-> construit une source Transport Yellowstone, Standard Logs, Standard Block ou Helius Transaction
|
|
-> construit HttpTransportPool + hydration role seulement pour les sources hydratées
|
|
-> construit le Store
|
|
-> construit RawTransactionIngestRuntimeResources
|
|
-> construit RawTransactionIngestSettings
|
|
-> démarre RawTransactionIngestWorker::start_with_runtime_resources
|
|
-> observe RawTransactionIngestSnapshotSource
|
|
-> demande stop lorsque nécessaire
|
|
```
|
|
|
|
Le Worker ne reçoit pas de requête historique métier et ne dépend pas de Config. Une campagne `signature/program_id/plage/limite` appartient à `ksp-job-backfill-lib`.
|