diff --git a/Cargo.toml b/Cargo.toml index d096272..010f346 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,12 +1,12 @@ # file: Cargo.toml -# version: 616 +# version: 617 [workspace] resolver = "3" members = ["crates/ksp-app-backfill-desk", "crates/ksp-app-config-desk", "crates/ksp-app-raw-transaction-ingest-desk", "crates/ksp-app-solprices-desk", "crates/ksp-app-store-desk", "crates/ksp-app-wallet-desk", "crates/ksp-config-lib", "crates/ksp-core-lib", "crates/ksp-interface-lib", "crates/ksp-job-api", "crates/ksp-job-backfill-lib", "crates/ksp-logging-lib", "crates/ksp-offchain-transport-lib", "crates/ksp-onchain-transport-lib", "crates/ksp-program-api", "crates/ksp-raw-transaction-lib", "crates/ksp-store-api", "crates/ksp-store-lib", "crates/ksp-store-postgres-lib", "crates/ksp-wallet-lib", "crates/ksp-worker-api", "crates/ksp-worker-raw-transaction-ingest-lib"] [workspace.package] -version = "0.3.15-pre.9.fix.2" +version = "0.3.15-pre.9.fix.3" edition = "2024" license = "MIT" repository = "https://git.sasedev.com/Sasedev/khadhroony-solana-project" diff --git a/crates/ksp-worker-raw-transaction-ingest-lib/src/runtime.rs b/crates/ksp-worker-raw-transaction-ingest-lib/src/runtime.rs index a98e751..f46ea33 100644 --- a/crates/ksp-worker-raw-transaction-ingest-lib/src/runtime.rs +++ b/crates/ksp-worker-raw-transaction-ingest-lib/src/runtime.rs @@ -1,5 +1,5 @@ // file: crates/ksp-worker-raw-transaction-ingest-lib/src/runtime.rs -// version: 18 +// version: 19 type PersistencePort = std::sync::Arc; type PersistenceTasks = tokio::task::JoinSet>; @@ -196,7 +196,8 @@ async fn drain_admission_and_persistence( snapshots: &mut crate::RawTransactionIngestSnapshotPublisher, ) -> std::option::Option { let mut fault = std::option::Option::None; - admission.close(); + // Keep admission open while cooperative sources observe the stop signal and drop their senders. + // Closing the receiver here races with nested source-stop propagation. loop { while persistence.len() >= settings.persistence_concurrency() { let joined = persistence.join_next().await; diff --git a/crates/ksp-worker-raw-transaction-ingest-lib/tests/release_completeness.rs b/crates/ksp-worker-raw-transaction-ingest-lib/tests/release_completeness.rs index 604c9f9..2abf9f8 100644 --- a/crates/ksp-worker-raw-transaction-ingest-lib/tests/release_completeness.rs +++ b/crates/ksp-worker-raw-transaction-ingest-lib/tests/release_completeness.rs @@ -1,5 +1,5 @@ // file: crates/ksp-worker-raw-transaction-ingest-lib/tests/release_completeness.rs -// version: 33 +// version: 34 //! Release-completeness canaries through the `v0.3.15-pre.004` WebSocket capability-enforcement tranche. @@ -36,6 +36,7 @@ fn pre_010_production_module_inventory_is_exact() -> std::io::Result<()> { names, std::vec![ "admission.rs", + "constants.rs", "continuity.rs", "error.rs", "identity.rs", diff --git a/crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime.rs b/crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime.rs index bf0a81e..699fd86 100644 --- a/crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime.rs +++ b/crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime.rs @@ -1,5 +1,5 @@ // file: crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime.rs -// version: 11 +// version: 12 struct ActiveTaskGuard { active: std::sync::Arc, @@ -994,3 +994,58 @@ async fn v0_3_13_pre_011_drain_timeout_prevents_late_persistence_completion_afte assert_eq!(retained.in_flight_persistence(), 0); return; } + +#[tokio::test(flavor = "current_thread")] +async fn v0_3_15_pre_009_fix_003_shutdown_keeps_admission_open_until_cooperative_source_observes_stop() { + let settings = match settings("mainnet") { + std::option::Option::Some(value) => value, + std::option::Option::None => return, + }; + let (started_sender, started_receiver) = tokio::sync::oneshot::channel::<()>(); + let (release_sender, release_receiver) = tokio::sync::oneshot::channel::<()>(); + let handle = match super::start_foundation_with_source_spawner( + settings, + tokio::runtime::Handle::current(), + std::option::Option::None, + move |children, mut stop_receiver, admission_sender| { + let _abort_handle = children.spawn(async move { + let _started = started_sender.send(()); + if release_receiver.await.is_err() { + return std::result::Result::Err(crate::runtime_error("test.release_closed")); + } + if admission_sender.is_closed() { + return std::result::Result::Err(crate::runtime_error("test.admission_closed_before_source_stop")); + } + if !*stop_receiver.borrow() { + let changed = stop_receiver.changed().await; + if changed.is_err() || !*stop_receiver.borrow() { + return std::result::Result::Err(crate::runtime_error("test.stop_not_observed")); + } + } + return std::result::Result::Ok(()); + }); + }, + ) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return, + }; + if started_receiver.await.is_err() { + return; + } + assert!(handle.request_stop()); + for _ in 0..64 { + if handle.snapshot_source().current().worker_snapshot().state() == ksp_worker_api::WorkerState::Stopping { + break; + } + tokio::task::yield_now().await; + } + assert_eq!(handle.snapshot_source().current().worker_snapshot().state(), ksp_worker_api::WorkerState::Stopping); + let _released = release_sender.send(()); + let terminal = match handle.wait_terminal().await { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return, + }; + assert_eq!(terminal, ksp_worker_api::WorkerState::Stopped); + assert_eq!(handle.snapshot_source().current().source_failure_total(), 0); + return; +} diff --git a/deltas/0.3.15/pre.009-fix.003.md b/deltas/0.3.15/pre.009-fix.003.md new file mode 100644 index 0000000..10f518f --- /dev/null +++ b/deltas/0.3.15/pre.009-fix.003.md @@ -0,0 +1,244 @@ + + + +# Delta `0.3.15-pre.009-fix.003` + +## Base requise + +```text +workspace state : 0.3.15-pre.009-fix.002 +Cargo version : 0.3.15-pre.9.fix.2 +base delta : deltas/0.3.15/pre.009-fix.002.md +base delta zip : ksp-general-0.3.15-pre.009-fix.002.zip +base zip SHA-256: aaf199604be3351eb2befbcf4903fe6c933895d71512a5fa6ab6a7476e31baa3 +``` + +`pre.009-fix.002` est un delta d'échange. L'état de base requis est donc le workspace `fix.001` après application exacte de ce delta, et non l'extraction isolée de l'archive delta dans un répertoire vide. + +## Objectif + +Fermer le dernier échec de gate statique de `fix.002` et corriger la race de shutdown désormais qualifiée de manière identique sur Mainnet par Yellowstone et HTTP Block Polling : + +```text +worker_raw_transaction_ingest.runtime_invalid +condition=source.admission_closed +transport_domain=none +transport_code=none +``` + +Le correctif doit conserver le drain borné, ne pas masquer les vraies fautes Transport et ne pas introduire de traitement spécifique par route lorsque la cause appartient au superviseur Worker commun. + +## Résultat du gate `fix.002` pris en compte + +Le gate opérateur confirme : + +```text +Rust rule audit PASS +Markdown table audit PASS (318 tables / 210 fichiers) +cargo check --workspace PASS +Clippy strict workspace PASS +Worker unit tests PASS : 158 / 158 +workspace_logging PASS +Worker all-targets/all-features FAIL : 1 canari release_completeness +workspace all-targets/all-features FAIL : même canari Worker +cargo tauri dev lancé et application opérationnelle +``` + +L'échec statique est déterministe : + +```text +pre_010_production_module_inventory_is_exact +réel : admission.rs, constants.rs, continuity.rs, ... +attendu : admission.rs, continuity.rs, ... +``` + +`constants.rs`, ajouté volontairement par `fix.002` pour satisfaire le contrat de tracing workspace, doit appartenir à cet inventaire exact. + +## Qualification live autoritaire + +Mainnet Yellowstone : + +```text +yellowstone-hydrated +route productive via Yellowstone Block + HTTP getBlock +Stop demandé +~5 s plus tard : runtime_invalid +condition=source.admission_closed +transport_domain=none +transport_code=none +terminal=Faulted +``` + +Mainnet HTTP polling : + +```text +http-block-polling +route productive via getSlot/getBlocks/getBlock +Stop demandé +~10 ms plus tard : runtime_invalid +condition=source.admission_closed +transport_domain=none +transport_code=none +terminal=Faulted +``` + +Le Yellowstone Devnet observé dans le même live est distinct : OrbitFlare retourne `onchain_transport.grpc_status` avec refus de permission lors de `SubscribeOpen`, avant Stop. Ce correctif ne reclasse pas cette faute provider/auth. + +## Cause racine + +Le superviseur Worker externe suit actuellement l'ordre suivant après Stop : + +```text +source_stop_sender = true +begin_stopping() +drain_owned_work() + -> drain_admission_and_persistence() + -> admission.close() +``` + +Le child `run_live_sources` reçoit le Stop externe puis doit encore le relayer à ses sources internes. Si l'admission est fermée avant ce relais, un dernier `send()` interne peut observer : + +```text +admission receiver fermé +source-local stop_receiver encore false +``` + +La source conclut alors à tort `source.admission_closed` et le shutdown coopératif devient `Faulted`. + +## Corrections + +- le receiver d'admission n'est plus fermé au début du drain coopératif ; +- l'ordre de shutdown continue à signaler le Stop aux sources avant le drain ; +- le drain continue à consommer la file jusqu'à destruction naturelle des derniers senders après observation du Stop ; +- `admission.close()` reste présent dans le chemin de `shutdown_drain_timeout`, juste avant abort/join forcé des tâches non coopératives ; +- un test unitaire force une source coopérative à retarder l'observation du Stop jusqu'à l'état `Stopping` et exige que l'admission soit encore ouverte, puis que le Worker termine `Stopped` avec `source_failure_total == 0` ; +- le canari `pre_010_production_module_inventory_is_exact` inclut désormais `constants.rs`. + +## Non-claims + +Ce correctif ne transforme pas les fautes Transport en succès et ne modifie pas : + +```text +Devnet Yellowstone / OrbitFlare permission denied +standard-logs-hydrated / rate-limit ou timeout provider +standard-block-direct / rpc_application_error +helius-transaction-hydrated / rpc_application_error +``` + +Ces fautes restent séparées du shutdown Mainnet qualifié ici. + +## Version + +```text +header racine : 616 -> 617 +workspace.package.version : 0.3.15-pre.9.fix.2 -> 0.3.15-pre.9.fix.3 +livraison : 0.3.15-pre.009-fix.003 +``` + +## Fichiers ajoutés + +```text +deltas/0.3.15/pre.009-fix.003.md +``` + +## Fichiers modifiés + +```text +Cargo.toml +crates/ksp-worker-raw-transaction-ingest-lib/src/runtime.rs +crates/ksp-worker-raw-transaction-ingest-lib/tests/release_completeness.rs +crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime.rs +docs/plans/036-V0_3_15_RAW_TRANSACTION_INGEST_DESK_PLAN.md +docs/validation/032-V0_3_15_RAW_TRANSACTION_INGEST_DESK.md +``` + +## Fichiers supprimés + +```text +aucun +``` + +## Validations exécutées dans l'environnement d'assemblage + +```text +python3 scripts/audit_rust_workspace_rules.py + General Rust rule audit: clean + Rust export completeness audit: 0 candidate(s) + KSP workspace Rust rule audit: clean + +python3 scripts/audit_markdown_tables.py README.md RULES.md ROADMAP.md CHANGELOG.md docs prompts crates deltas/0.3.15 + Markdown table audit: clean (318 tables / 211 fichiers) + +canari inventaire production Worker + src/constants.rs présent dans l'inventaire exact attendu + +archive delta + entrées : 7 + unzip -t : PASS + +reconstruction depuis l'état fix.002 + reconstructed files : 1938 + work files : 1938 + only reconstructed : 0 + only work : 0 + hash mismatches : 0 + résultat : byte-exact +``` + +## Validations non exécutées dans l'environnement d'assemblage + +`cargo` et `rustfmt` ne sont pas disponibles dans cet environnement. Les commandes suivantes ne sont donc pas déclarées PASS pour `fix.003` : + +```text +cargo fmt --all +cargo fmt --all -- --check +cargo check --workspace +cargo clippy --workspace --all-targets --all-features -- -D warnings +cargo test -p ksp-worker-raw-transaction-ingest-lib --all-targets --all-features +cargo test -p ksp-core-lib --test workspace_logging +cargo test --workspace --all-targets --all-features +cargo tauri dev +``` + +## Gate opérateur requis + +```bash +cargo fmt --all +cargo fmt --all -- --check +python3 scripts/audit_rust_workspace_rules.py +python3 scripts/audit_markdown_tables.py README.md RULES.md ROADMAP.md CHANGELOG.md docs prompts crates deltas/0.3.15 +cargo check --workspace +cargo clippy --workspace --all-targets --all-features -- -D warnings +cargo test -p ksp-worker-raw-transaction-ingest-lib --all-targets --all-features +cargo test -p ksp-core-lib --test workspace_logging +cargo test --workspace --all-targets --all-features +(cd crates/ksp-app-raw-transaction-ingest-desk && cargo tauri dev) +``` + +Pour le live Mainnet, tester au minimum : + +```text +yellowstone-hydrated -> Stop -> attendu Stopped +http-block-polling -> Stop -> attendu Stopped +``` + +Une nouvelle ligne `RAW transaction ingest live source reached terminal failure` sur l'une de ces routes doit être copiée intégralement avec `condition`, `transport_domain` et `transport_code` ; elle ne doit pas être reclassée sans qualification. + +## Décisions + +```text +corriger la race au niveau du drain Worker commun, pas par source +laisser l'admission ouverte pendant la propagation coopérative du Stop +réserver admission.close au cutoff forcé du shutdown timeout +conserver shutdown_drain_timeout à 10 s +ajouter constants.rs à l'inventaire de production exact +``` + +## Questions ouvertes + +```text +validation live Mainnet Yellowstone -> Stopped après fix.003 +validation live Mainnet HTTP Block Polling -> Stopped après fix.003 +qualification/provider strategy des rpc_application_error Standard Block et Helius +stratégie de débit Standard Logs face au rate-limit du RPC public +``` diff --git a/docs/plans/036-V0_3_15_RAW_TRANSACTION_INGEST_DESK_PLAN.md b/docs/plans/036-V0_3_15_RAW_TRANSACTION_INGEST_DESK_PLAN.md index c9e2cb0..302cc6b 100644 --- a/docs/plans/036-V0_3_15_RAW_TRANSACTION_INGEST_DESK_PLAN.md +++ b/docs/plans/036-V0_3_15_RAW_TRANSACTION_INGEST_DESK_PLAN.md @@ -1,5 +1,5 @@ - + # Plan v0.3.15 — Raw Transaction Ingest Desk @@ -904,3 +904,21 @@ Le gate opérateur de `fix.001` invalide la conclusion selon laquelle le défaut Deux régressions statiques de `fix.001` doivent être fermées avant toute nouvelle hypothèse live : le test des defaults Worker doit attendre 10 s comme la constante publique, et la cible `tracing` explicite doit être possédée par `src/constants.rs` puis consommée via la façade crate. La trace sûre de faute source est étendue avec le contexte interne `condition` déjà borné par `runtime_error`; aucune valeur provider, endpoint, signature, payload ou secret n'est ajoutée. Le prochain gate live doit relever cette `condition` pour Yellowstone et HTTP polling avant toute correction fonctionnelle supplémentaire. Les faults `standard-logs-hydrated` (`timeout` après rate-limit), `standard-block-direct` (`rpc_application_error`) et `helius-transaction-hydrated` (`rpc_application_error`) sont des fautes antérieures au Stop et restent séparées du diagnostic shutdown. + +### `pre.009-fix.003` — race Stop / admission drain + +Le gate opérateur de `fix.002` ferme les deux régressions statiques de `fix.001` : les 158 tests unitaires Worker passent et le canari workspace de cible tracing passe. Un unique test d’intégration Worker reste rouge parce que `src/constants.rs`, ajouté par `fix.002`, n’a pas été ajouté à l’inventaire exact de `tests/release_completeness.rs`. + +Les nouvelles traces live qualifient désormais la cause Mainnet des deux faults de shutdown : + +```text +yellowstone-hydrated : Stop -> runtime_invalid / condition=source.admission_closed +http-block-polling : Stop -> runtime_invalid / condition=source.admission_closed +``` + +La cause est la même pour les deux routes et se situe dans le drain Worker commun. Après `request_stop`, le superviseur externe signale le Stop au child `run_live_sources`, puis `drain_admission_and_persistence` ferme immédiatement le receiver d’admission. Le child doit encore relayer ce Stop à ses sources internes ; pendant cette fenêtre, un dernier `send()` peut donc observer un receiver fermé avant que son propre `stop_receiver` ne soit passé à `true`, et la fermeture coopérative devient artificiellement `source.admission_closed`. + +`fix.003` garde le receiver d’admission ouvert pendant le drain coopératif. Les sources ont déjà reçu l’ordre de Stop ; le drain continue à consommer les éléments déjà produits jusqu’à ce que les derniers senders soient naturellement détruits. La fermeture forcée de l’admission reste réservée au chemin de timeout borné, qui abort/join les tâches non coopératives. Un test unitaire reproduit la fenêtre de propagation retardée et exige un terminal `Stopped` sans `source_failure_total`. + +Le fault Yellowstone Devnet observé dans le même gate reste distinct : OrbitFlare répond `onchain_transport.grpc_status` avec refus de permission dès l’ouverture du Subscribe, avant toute demande de Stop. Il n’est pas reclassé par ce correctif. + diff --git a/docs/validation/032-V0_3_15_RAW_TRANSACTION_INGEST_DESK.md b/docs/validation/032-V0_3_15_RAW_TRANSACTION_INGEST_DESK.md index 0562d3a..c663c6e 100644 --- a/docs/validation/032-V0_3_15_RAW_TRANSACTION_INGEST_DESK.md +++ b/docs/validation/032-V0_3_15_RAW_TRANSACTION_INGEST_DESK.md @@ -1,5 +1,5 @@ - + # Validation v0.3.15 — Raw Transaction Ingest Desk @@ -591,3 +591,59 @@ helius-transaction-hydrated ``` Conclusion : le passage du drain par défaut à 10 s reste conservé mais n'est pas une preuve de correction Yellowstone. `fix.002` répare les deux gates statiques et expose dans la trace source le contexte interne statique `condition` des erreurs `runtime_invalid`. Le prochain live doit fournir cette `condition` exacte pour Yellowstone et HTTP polling ; aucune conversion générique `runtime_invalid -> Stopped` n'est autorisée. + +### `pre.009-fix.003` — résultat du gate opérateur `fix.002` + +Le gate opérateur de `0.3.15-pre.009-fix.002` confirme : + +```text +audit Rust workspace PASS +audit Markdown PASS (318 tables / 210 fichiers) +cargo check --workspace PASS +Clippy workspace strict -D warnings PASS +tests unitaires Worker PASS : 158 / 158 +workspace_logging PASS +tests Worker complets FAIL : 1 canari release_completeness +tests workspace FAIL : même canari Worker +cargo tauri dev lancé et application opérationnelle +``` + +L’échec statique restant est exact : + +```text +pre_010_production_module_inventory_is_exact + réel : admission.rs, constants.rs, continuity.rs, ... + attendu : admission.rs, continuity.rs, ... +``` + +`constants.rs` est une production source volontaire de `fix.002`; le canari doit donc l’inclure. + +Le live apporte surtout la qualification attendue : + +```text +Mainnet / yellowstone-hydrated + route productive via Yellowstone Block + HTTP getBlock + Stop demandé + ~5 s plus tard : worker_raw_transaction_ingest.runtime_invalid + condition=source.admission_closed + transport_domain=none + transport_code=none + +Mainnet / http-block-polling + route productive via getSlot/getBlocks/getBlock + Stop demandé + ~10 ms plus tard : worker_raw_transaction_ingest.runtime_invalid + condition=source.admission_closed + transport_domain=none + transport_code=none + +Devnet / yellowstone-hydrated + fault avant Stop lors de SubscribeOpen + source_failed / onchain_transport.grpc_status + OrbitFlare refuse la permission +``` + +La cause Mainnet est une race interne commune : le superviseur Worker ferme le receiver d’admission au début du drain alors que `run_live_sources` n’a pas nécessairement encore relayé le Stop à sa source interne. Un dernier `send()` voit alors l’admission fermée alors que son `stop_receiver` local vaut encore `false`, ce qui fabrique `source.admission_closed`. + +Validation attendue après `fix.003` : le drain coopératif ne ferme plus l’admission au début ; il consomme la file jusqu’à destruction naturelle des senders après propagation du Stop. Le hard-close reste dans le chemin `shutdown_drain_timeout`. Un test unitaire force cette fenêtre de propagation et exige `Stopped` sans source failure. +