# 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 { 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 10 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 { 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. L'agrégat accepte 1 à 32 sources logiques et les lance simultanément ; les doublons d'identité et les mélanges de réseaux sont refusés avant spawn. Il n'existe pas de source primaire, standby ou fallback implicite : toute source configurée fait partie du run. ### Composer plusieurs sources simultanées Une fois un premier `RawTransactionIngestRuntimeResources` construit, le caller ajoute les autres sources avec les méthodes `try_push_*` correspondant à leur capability. Exemple conceptuel : ```rust let mut resources = ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestRuntimeResources::new(yellowstone_source); if let std::result::Result::Err(error) = resources.try_push_standard_logs_source(standard_logs_source) { return std::result::Result::Err(error); } if let std::result::Result::Err(error) = resources.try_push_http_block_polling_source(http_polling_source) { return std::result::Result::Err(error); } ``` La validation est transactionnelle à chaque ajout : limite globale 32, réseau unique et `source_key` logique unique. Au démarrage, toutes les sources présentes sont supervisées ensemble. Une défaillance de source ne permet la continuation des siblings que lorsqu'une plage de perte sûre est disponible, que cette perte est entièrement réconciliée et que les sources encore actives prouvent tout le `TargetCoverage` futur. Sinon le Worker arrête et joint les autres sources. Aucune équivalence n'est déduite du seul provider, protocole ou nom de filtre. ### 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 { 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` déclarant explicitement `WsSubscriptionKind::Logs`. Un endpoint legacy dont les capabilities sont absentes est refusé fail-closed. 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 { 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)); } ``` `ws_endpoint` doit déclarer explicitement `WsSubscriptionKind::Block`; l'absence de déclaration est refusée avant toute connexion. 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 { 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` déclarant explicitement `WsSubscriptionKind::HeliusTransaction`. L'état capability undeclared est refusé avant toute connexion. 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. ### Source HTTP Block Polling productive Une source HTTP live peut être composée sans WebSocket ni gRPC : ```rust fn http_block_polling_runtime_resources( http_pool: ksp_onchain_transport_lib::HttpTransportPool, polling_role: ksp_onchain_transport_lib::HttpRoleName, commitment: ksp_onchain_transport_lib::SolanaCommitment, ) -> ksp_core_lib::Result { let source = match ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestHttpBlockPollingSource::new( http_pool, polling_role, 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_http_block_polling_source(source)); } ``` Le rôle HTTP doit disposer, sur un même réseau, de `getSlot`, `getBlocksWithLimit` et `getBlock`. Le commitment est limité à `Confirmed`/`Finalized`. Par défaut, le Worker interroge toutes les secondes et borne la découverte à 128 blocs par cycle. `new_with_limits` permet de choisir une cadence entre 100 ms et 30 s et une limite entre 1 et 1024 blocs par cycle ; les limites de débit physiques restent celles de Transport. Au démarrage du run, la première valeur `getSlot(commitment)` devient la borne inférieure stricte du poller. Il ne demande jamais de slot antérieur. `getBlocksWithLimit` détermine les slots réellement matérialisables, puis `getBlock observed` produit directement le Common RAW en Full/Base64 pour Legacy/V0/V1. Un slot listé dont `getBlock` retourne `null` reste la tête de reprise du prochain cycle et n'est jamais transformé en progression silencieuse. ## 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 sur le même réseau. Les familles Transaction/TransactionStatus exigent `getTransaction`; un abonnement Yellowstone Block pur exige `getBlock`. 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 -> trigger slot -> HTTP getBlock Full/Base64 -> Common RAW par transaction -> 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 HTTP live block polling -> getSlot -> getBlocksWithLimit -> getBlock observed -> Common RAW direct Legacy|V0|V1 -> 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 globalement avant l'hydration HTTP. Si plusieurs sources produisent ensuite la même transaction canonique, le Worker conserve une seule entité RAW et enregistre séparément les observations déterministes propres à chaque source. Pour Yellowstone Block, les slots gRPC sont coordonnés dans une borne pending distincte et plusieurs `getBlock` peuvent progresser jusqu'au quota in-flight de la source ; si l'aval sature, Transport attend de la capacité sur sa queue bornée et propage la backpressure au stream gRPC au lieu de dropper l'update ou de produire un overflow local. Le Worker ne possède pas une boucle de retry HTTP : reroutage/retry/backoff restent dans `ksp-onchain-transport-lib`. Les sources qui nécessitent `getTransaction` partagent des quotas bornés et déterministes. Pour une composition valide : ```text nombre de sources avec hydration <= admission_queue_capacity nombre de sources avec hydration <= persistence_concurrency somme de leurs pending <= admission_queue_capacity somme de leurs tâches d'hydration actives <= persistence_concurrency ``` Ces contraintes garantissent au moins une part à chaque source reference-bearing sans introduire de scheduler pondéré. Si les settings ne permettent pas cette répartition, le démarrage échoue avant spawn avec `runtime_invalid`; le caller doit augmenter la capacité concernée ou réduire le nombre de sources nécessitant une hydration. Pour une même identité `(network, signature)`, le Worker ne choisit ni majorité ni provider préféré et n'écrase pas un contenu divergent ; le Store reste l'autorité durable finale. En `0.3.15`, un canonique complet peut accepter comme observation supplémentaire un entrant identique hors `meta.logMessages` lorsque celui-ci contient exactement un marqueur `Log truncated` après un préfixe identique. Le canonique n'est pas réécrit. Le sens tronqué -> complet et toute autre divergence non prouvée restent des `content_conflict` terminaux ; leur stockage en variantes et leur résolution sont réservés à `0.3.16`. ## 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_total = current.source_total(); let source_active = current.source_active(); let source_reconnecting = current.source_reconnecting(); let source_failed = current.source_failed(); 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(); let gaps = snapshot.gaps(); let open_gaps = snapshot.open_gap_count(); let repairing_gaps = snapshot.repairing_gap_count(); let repaired_gaps = snapshot.repaired_gap_total(); let unresolved_gaps = snapshot.unresolved_gap_total(); let oldest_gap = snapshot.oldest_open_gap_start_slot(); ``` `admission_queue_depth()`, `in_flight_persistence()`, `hydration_pending()`, `source_total()`, `source_active()`, `source_reconnecting()` et `source_failed()` sont des gauges latest-value. Les compteurs cumulés et les agrégats multi-source ne wrapent ni ne saturent 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` comme état agrégé source-neutral. Les gauges `source_total()`, `source_active()`, `source_reconnecting()` et `source_failed()` permettent d'interpréter une composition multi-source sans exposer provider, endpoint, filtre ou `source_key`. Lorsque la policy de continuité est active, la health commune est plus stricte qu'une simple lecture des états source. Pendant `Running`, reconnect en cours, gap ouvert, continuity frontier différente de la processing frontier ou `TargetCoverage` futur non couvert donnent `Unhealthy`. Toutes les sources attendues `Active` avec continuité réconciliée permettent `Healthy`. Une source `Failed` ne donne `Degraded` que si sa perte est entièrement réconciliée et si les sources restantes couvrent encore tout le `TargetCoverage` futur. Une transition Worker `Faulted` reste `Unhealthy`. 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. Un gap est projeté via `gaps()` avec un état `Pending`, `Repairing`, `Repaired` ou `Unresolved`, une raison source-neutral et éventuellement la dernière méthode de réparation. Les mécanismes publics décrits par `RawTransactionIngestRepairMethod` sont `Replay`, `RedundantCoverage`, `HttpScan`, `BlockFetch` et `TransactionHydration`. Une perte source ou un gap de rétention ne déclenche jamais `ksp-job-backfill-lib`. Le Worker continue seulement lorsqu'il peut prouver la réconciliation passée et la coverage future dans les bornes du run courant ; sinon il fault. ## 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. Le défaut de drain est de 10 s afin de laisser une fermeture Yellowstone bornée à 5 s se terminer sans collision de deadline. Si la deadline de drain expire, toutes les tâches source/persistence encore possédées sont abortées puis jointes avant publication terminale ; une persistence libérée après ce terminal ne peut donc pas produire une complétion tardive. Une faute source déjà observée reste prioritaire face à un stop concurrent, sauf si le drain lui-même expire et devient le terminal `drain_timeout`. ## 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, Helius Transaction ou HTTP Block Polling -> construit HttpTransportPool + rôle HTTP adapté à la source -> 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`.