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 }