lib

Core libraries for Radroots
git clone https://radroots.dev/git/lib.git
Log | Files | Refs | README

engine_composition.rs (13047B)


      1 use std::sync::Arc;
      2 
      3 use futures_executor::block_on;
      4 use radroots_protocol::runtime::v1::{OPERATION_SCHEMA_VERSION, SyncCapabilityState, SyncHealth};
      5 use radroots_signing::{Error as SigningError, SignReceipt, SignRequest, Signer, SignerStatus};
      6 use radroots_storage::{event::SourceGeneration, memory::MemoryStorage, projection::ProjectionId};
      7 use radroots_sync::{
      8     Engine,
      9     policy::{Clock, DeadlinePolicy, Error, IdSource, OperationKind, SyncId, SyncStorage},
     10 };
     11 use radroots_transport::{
     12     DeliveryReceipt, DeliveryRequest, Error as TransportError, EventSink, EventSource, FetchPage,
     13     FetchRequest, SinkFailure, SinkStatus, SourceStatus,
     14     capability::{Availability, Maturity, SinkCapabilities, SourceCapabilities},
     15 };
     16 
     17 struct MockSource;
     18 struct MockSink;
     19 struct MockSigner;
     20 struct UnconfiguredSource;
     21 struct UnconfiguredSink;
     22 struct FixedClock;
     23 struct FixedIds;
     24 
     25 type TestDependencies = (
     26     Arc<dyn SyncStorage>,
     27     Arc<dyn Clock>,
     28     Arc<dyn IdSource>,
     29     DeadlinePolicy,
     30 );
     31 
     32 impl EventSource for MockSource {
     33     fn status(&self) -> radroots_transport::BoxFuture<'_, Result<SourceStatus, TransportError>> {
     34         Box::pin(async {
     35             Ok(SourceStatus::new(
     36                 radroots_transport::TransportId::NOSTR,
     37                 true,
     38                 Maturity::Stable,
     39                 Availability::Available,
     40                 SourceCapabilities::FETCH,
     41                 "ready",
     42             ))
     43         })
     44     }
     45 
     46     fn fetch(
     47         &self,
     48         _request: FetchRequest,
     49     ) -> radroots_transport::BoxFuture<'_, Result<FetchPage, TransportError>> {
     50         Box::pin(async { unreachable!("composition does not fetch") })
     51     }
     52 }
     53 
     54 impl EventSource for UnconfiguredSource {
     55     fn status(&self) -> radroots_transport::BoxFuture<'_, Result<SourceStatus, TransportError>> {
     56         Box::pin(async {
     57             Ok(SourceStatus::new(
     58                 radroots_transport::TransportId::NOSTR,
     59                 false,
     60                 Maturity::Preview,
     61                 Availability::Available,
     62                 SourceCapabilities::FETCH,
     63                 "not configured",
     64             ))
     65         })
     66     }
     67 
     68     fn fetch(
     69         &self,
     70         _request: FetchRequest,
     71     ) -> radroots_transport::BoxFuture<'_, Result<FetchPage, TransportError>> {
     72         Box::pin(async { unreachable!("status does not fetch") })
     73     }
     74 }
     75 
     76 impl EventSink for MockSink {
     77     fn status(&self) -> radroots_transport::BoxFuture<'_, Result<SinkStatus, TransportError>> {
     78         Box::pin(async {
     79             Ok(SinkStatus::new(
     80                 radroots_transport::TransportId::NOSTR,
     81                 true,
     82                 Maturity::Preview,
     83                 Availability::Degraded,
     84                 SinkCapabilities::DELIVER,
     85                 "degraded",
     86             ))
     87         })
     88     }
     89 
     90     fn deliver(
     91         &self,
     92         _request: DeliveryRequest,
     93     ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> {
     94         Box::pin(async { unreachable!("composition does not deliver") })
     95     }
     96 }
     97 
     98 impl EventSink for UnconfiguredSink {
     99     fn status(&self) -> radroots_transport::BoxFuture<'_, Result<SinkStatus, TransportError>> {
    100         Box::pin(async {
    101             Ok(SinkStatus::new(
    102                 radroots_transport::TransportId::NOSTR,
    103                 false,
    104                 Maturity::Preview,
    105                 Availability::Available,
    106                 SinkCapabilities::DELIVER,
    107                 "not configured",
    108             ))
    109         })
    110     }
    111 
    112     fn deliver(
    113         &self,
    114         _request: DeliveryRequest,
    115     ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> {
    116         Box::pin(async { unreachable!("status does not deliver") })
    117     }
    118 }
    119 
    120 impl Signer for MockSigner {
    121     fn status(
    122         &self,
    123     ) -> radroots_signing::signer::BoxFuture<'_, Result<SignerStatus, SigningError>> {
    124         Box::pin(async { Ok(SignerStatus::unavailable()) })
    125     }
    126 
    127     fn sign(
    128         &self,
    129         _request: SignRequest,
    130     ) -> radroots_signing::signer::BoxFuture<'_, Result<SignReceipt, SigningError>> {
    131         Box::pin(async { unreachable!("composition does not sign") })
    132     }
    133 }
    134 
    135 impl Clock for FixedClock {
    136     fn now_unix_ms(&self) -> Result<u64, Error> {
    137         Ok(1_700_000_000_000)
    138     }
    139 }
    140 
    141 impl IdSource for FixedIds {
    142     fn next_id(&self, operation: OperationKind) -> Result<SyncId, Error> {
    143         let byte = match operation {
    144             OperationKind::Pull => 1,
    145             OperationKind::Sign => 2,
    146             OperationKind::Deliver => 3,
    147             _ => 4,
    148         };
    149         SyncId::new([byte; 16])
    150     }
    151 }
    152 
    153 fn dependencies() -> TestDependencies {
    154     let generation = SourceGeneration::new([7; 32]).expect("generation");
    155     (
    156         Arc::new(MemoryStorage::new(generation)),
    157         Arc::new(FixedClock),
    158         Arc::new(FixedIds),
    159         DeadlinePolicy::new(10_000, 20_000, 30_000).expect("deadlines"),
    160     )
    161 }
    162 
    163 #[test]
    164 fn source_only_sink_only_and_full_compositions_are_explicit() {
    165     let (storage, clock, ids, deadlines) = dependencies();
    166     let source_only = Engine::builder(storage, clock, ids, deadlines)
    167         .source(Arc::new(MockSource))
    168         .build()
    169         .expect("source engine");
    170     assert!(source_only.source().is_some());
    171     assert!(source_only.sink().is_none());
    172     assert!(source_only.signer().is_none());
    173 
    174     let (storage, clock, ids, deadlines) = dependencies();
    175     let sink_only = Engine::builder(storage, clock, ids, deadlines)
    176         .sink(Arc::new(MockSink))
    177         .build()
    178         .expect("sink engine");
    179     assert!(sink_only.source().is_none());
    180     assert!(sink_only.sink().is_some());
    181     assert!(sink_only.signer().is_none());
    182 
    183     let (storage, clock, ids, deadlines) = dependencies();
    184     let full = Engine::builder(storage, clock, ids, deadlines)
    185         .source(Arc::new(MockSource))
    186         .sink(Arc::new(MockSink))
    187         .signer(Arc::new(MockSigner))
    188         .build()
    189         .expect("full engine");
    190     assert!(full.source().is_some());
    191     assert!(full.sink().is_some());
    192     assert!(full.signer().is_some());
    193     assert_eq!(
    194         full.clock().now_unix_ms().expect("clock"),
    195         1_700_000_000_000
    196     );
    197     assert_eq!(
    198         full.deadlines()
    199             .deadline_unix_ms(OperationKind::Deliver, 1_000)
    200             .expect("deadline"),
    201         31_000
    202     );
    203     assert_eq!(
    204         block_on(full.storage().storage_status())
    205             .expect("storage status")
    206             .shutdown(),
    207         radroots_storage::status::ShutdownState::Open
    208     );
    209     assert_ne!(
    210         full.ids()
    211             .next_id(OperationKind::Ingest)
    212             .expect("identity")
    213             .as_bytes(),
    214         &[0; 16]
    215     );
    216     assert!(format!("{full:?}").contains("Engine"));
    217     assert!(format!("{:?}", full.clone()).contains("signer: true"));
    218 }
    219 
    220 #[test]
    221 fn invalid_compositions_and_ambient_policy_inputs_fail_closed() {
    222     let (storage, clock, ids, deadlines) = dependencies();
    223     assert_eq!(
    224         Engine::builder(storage, clock, ids, deadlines)
    225             .build()
    226             .expect_err("missing transport"),
    227         Error::MissingTransportCapability
    228     );
    229 
    230     let (storage, clock, ids, deadlines) = dependencies();
    231     assert_eq!(
    232         Engine::builder(storage, clock, ids, deadlines)
    233             .signer(Arc::new(MockSigner))
    234             .build()
    235             .expect_err("signer without sink"),
    236         Error::SignerWithoutSink
    237     );
    238     assert_eq!(
    239         DeadlinePolicy::new(0, 1, 1),
    240         Err(Error::InvalidDeadlinePolicy)
    241     );
    242     assert_eq!(SyncId::new([0; 16]), Err(Error::InvalidSyncId));
    243     let id = SyncId::new([9; 16]).expect("sync identity");
    244     assert_eq!(id.as_bytes(), &[9; 16]);
    245     for invalid in [
    246         DeadlinePolicy::new(1, 0, 1),
    247         DeadlinePolicy::new(1, 1, 0),
    248         DeadlinePolicy::new(u64::MAX, 1, 1),
    249         DeadlinePolicy::new(1, u64::MAX, 1),
    250         DeadlinePolicy::new(1, 1, u64::MAX),
    251     ] {
    252         assert_eq!(invalid, Err(Error::InvalidDeadlinePolicy));
    253     }
    254     let deadlines = DeadlinePolicy::new(10, 20, 30).expect("deadlines");
    255     for (operation, expected) in [
    256         (OperationKind::Ingest, 10),
    257         (OperationKind::Projection, 10),
    258         (OperationKind::Pull, 10),
    259         (OperationKind::Sign, 20),
    260         (OperationKind::Deliver, 30),
    261     ] {
    262         assert_eq!(deadlines.timeout_ms(operation), expected);
    263     }
    264     assert_eq!(
    265         deadlines.deadline_unix_ms(OperationKind::Pull, 0),
    266         Err(Error::ClockUnavailable)
    267     );
    268     assert_eq!(
    269         deadlines.deadline_unix_ms(OperationKind::Pull, u64::MAX),
    270         Err(Error::DeadlineOverflow)
    271     );
    272     for error in [
    273         Error::InvalidSyncId,
    274         Error::InvalidDeadlinePolicy,
    275         Error::ClockUnavailable,
    276         Error::DeadlineOverflow,
    277         Error::MissingTransportCapability,
    278         Error::SignerWithoutSink,
    279         Error::VerificationFailed,
    280         Error::PolicyRejected,
    281         Error::StorageConflict,
    282         Error::StorageSpaceInsufficient,
    283         Error::StorageFailed,
    284         Error::InvalidIngestReceipt,
    285         Error::InvalidPullRequest,
    286         Error::MissingSource,
    287         Error::InvalidSourcePage,
    288         Error::InvalidProjectionRequest,
    289         Error::ReducerFailed,
    290         Error::InvalidReducerOutput,
    291         Error::InvalidPushRequest,
    292         Error::MissingSigner,
    293         Error::SignerFailed,
    294         Error::SignerDeadlineExceeded,
    295         Error::InvalidSignerOutput,
    296         Error::InvalidDeliveryRequest,
    297         Error::MissingSink,
    298         Error::InvalidStatusRequest,
    299     ] {
    300         assert!(!error.to_string().is_empty());
    301     }
    302 }
    303 
    304 #[test]
    305 fn status_aggregates_typed_capability_and_protocol_reports() {
    306     let (storage, clock, ids, deadlines) = dependencies();
    307     let full = Engine::builder(storage, clock, ids, deadlines)
    308         .source(Arc::new(MockSource))
    309         .sink(Arc::new(MockSink))
    310         .signer(Arc::new(MockSigner))
    311         .build()
    312         .expect("full engine");
    313     let projection = ProjectionId::parse("market-listings").expect("projection id");
    314     let status = block_on(full.status(std::slice::from_ref(&projection))).expect("sync status");
    315     assert_eq!(status.health(), SyncHealth::Degraded);
    316     assert_eq!(status.source().state(), SyncCapabilityState::Available);
    317     assert_eq!(status.sink().state(), SyncCapabilityState::Degraded);
    318     assert_eq!(status.signer().state(), SyncCapabilityState::Configured);
    319     assert_eq!(
    320         status.storage().shutdown(),
    321         radroots_storage::status::ShutdownState::Open
    322     );
    323     assert_eq!(
    324         status.events().health(),
    325         radroots_storage::status::EventStoreHealth::Available
    326     );
    327     assert_eq!(status.outbox().total(), Some(0));
    328     assert!(status.source().status().is_some());
    329     assert!(status.sink().status().is_some());
    330     assert!(status.signer().status().is_some());
    331     assert_eq!(status.projections()[0].projection_id(), &projection);
    332     assert!(status.projections()[0].status().is_none());
    333     let protocol = status.to_protocol();
    334     assert_eq!(protocol.schema_version, OPERATION_SCHEMA_VERSION);
    335     assert_eq!(protocol.health, SyncHealth::Degraded);
    336     assert_eq!(protocol.source, SyncCapabilityState::Available);
    337     assert_eq!(protocol.sink, SyncCapabilityState::Degraded);
    338     assert_eq!(protocol.signer, SyncCapabilityState::Configured);
    339     assert_eq!(protocol.projections.untracked, 1);
    340 
    341     let (storage, clock, ids, deadlines) = dependencies();
    342     let sink_only = Engine::builder(storage, clock, ids, deadlines)
    343         .sink(Arc::new(MockSink))
    344         .build()
    345         .expect("sink engine");
    346     let status = block_on(sink_only.status(&[])).expect("sink-only status");
    347     assert_eq!(status.source().state(), SyncCapabilityState::Unsupported);
    348     assert_eq!(status.signer().state(), SyncCapabilityState::Unsupported);
    349 
    350     let (storage, clock, ids, deadlines) = dependencies();
    351     let compiled = Engine::builder(storage, clock, ids, deadlines)
    352         .source(Arc::new(UnconfiguredSource))
    353         .build()
    354         .expect("unconfigured source engine");
    355     let status = block_on(compiled.status(&[])).expect("compiled status");
    356     assert_eq!(status.source().state(), SyncCapabilityState::Compiled);
    357     assert!(status.source().status().is_some());
    358     assert_eq!(status.sink().state(), SyncCapabilityState::Unsupported);
    359 
    360     let (storage, clock, ids, deadlines) = dependencies();
    361     let compiled_sink = Engine::builder(storage, clock, ids, deadlines)
    362         .sink(Arc::new(UnconfiguredSink))
    363         .build()
    364         .expect("unconfigured sink engine");
    365     let status = block_on(compiled_sink.status(&[])).expect("compiled sink status");
    366     assert_eq!(status.sink().state(), SyncCapabilityState::Compiled);
    367     assert!(status.sink().status().is_some());
    368 
    369     assert_eq!(
    370         block_on(full.status(&[projection.clone(), projection])),
    371         Err(Error::InvalidStatusRequest)
    372     );
    373     let too_many = vec![ProjectionId::parse("too-many-projections").expect("projection id"); 257];
    374     assert_eq!(
    375         block_on(full.status(&too_many)),
    376         Err(Error::InvalidStatusRequest)
    377     );
    378 }