1mod assembly;
9mod composition;
10mod local;
11#[cfg(feature = "test-support")]
12pub mod test_support;
13mod transport;
14
15pub use assembly::{AssemblyError, AssemblyErrors};
16pub use composition::{Composition, CompositionBuilder, ImportTarget, RegisteredBox};
17pub use local::LocalBinding;
18pub use transport::{
19 TransportBinding, TransportExposure, TransportHandle, TransportJoinFuture, TransportRuntime,
20 TransportTaskTracker,
21};
22
23use std::collections::BTreeMap;
24use std::future::{Future, ready};
25use std::pin::Pin;
26use std::sync::{Arc, OnceLock};
27
28use boxology_contract::{
29 BoxId, CallContext, CapabilityId, Detail, ErasedCallError, ErasedCallTarget, ErasedTarget,
30 SlotValue, call_guarded,
31};
32
33type ErasedCallFuture<'a> =
34 Pin<Box<dyn Future<Output = Result<SlotValue, ErasedCallError>> + Send + 'a>>;
35
36pub trait RemoteImportTarget: ErasedCallTarget {
38 fn supports_capability(&self, capability: &CapabilityId) -> bool;
40}
41
42pub struct Imports {
47 handles: BTreeMap<BoxId, ImportHandle>,
48}
49
50impl Imports {
51 pub fn handle(&self, slot: &BoxId) -> Option<&ImportHandle> {
53 self.handles.get(slot)
54 }
55
56 pub(crate) fn cloned_handles(&self) -> BTreeMap<BoxId, ImportHandle> {
57 self.handles.clone()
58 }
59
60 pub(crate) fn new(imports: impl IntoIterator<Item = (BoxId, Vec<CapabilityId>)>) -> Self {
61 let handles = imports
62 .into_iter()
63 .map(|(slot, capabilities)| {
64 let handle = ImportHandle::new(slot.clone(), capabilities);
65 (slot, handle)
66 })
67 .collect();
68 Self { handles }
69 }
70}
71
72#[derive(Clone)]
77pub struct ImportHandle {
78 slot: BoxId,
79 capabilities: Arc<[CapabilityId]>,
80 target: Arc<OnceLock<Arc<dyn ErasedTarget>>>,
81}
82
83impl ImportHandle {
84 fn new(slot: BoxId, capabilities: Vec<CapabilityId>) -> Self {
85 Self {
86 slot,
87 capabilities: capabilities.into(),
88 target: Arc::new(OnceLock::new()),
89 }
90 }
91
92 pub fn slot_id(&self) -> &BoxId {
94 &self.slot
95 }
96
97 pub fn capabilities(&self) -> &[CapabilityId] {
99 &self.capabilities
100 }
101
102 pub fn call<'a>(
109 &'a self,
110 capability: &'a CapabilityId,
111 context: CallContext,
112 input: SlotValue,
113 ) -> ErasedCallFuture<'a> {
114 let Some(target) = self.target.get() else {
115 return Box::pin(ready(Err(ErasedCallError::Unavailable(Detail::new(
116 "unsealed_import",
117 )))));
118 };
119 if !self
120 .capabilities
121 .iter()
122 .any(|allowed| allowed == capability)
123 {
124 return Box::pin(ready(Err(ErasedCallError::ContractViolation(Detail::new(
125 "undeclared_import_capability",
126 )))));
127 }
128 if context
129 .deadline()
130 .is_some_and(|deadline| deadline.remaining().is_zero())
131 {
132 return Box::pin(ready(Err(ErasedCallError::Deadline)));
133 }
134 call_guarded(target.as_ref(), capability, context, input)
135 }
136
137 pub(crate) fn seal(&self, target: Arc<dyn ErasedTarget>) -> Result<(), Arc<dyn ErasedTarget>> {
138 self.target.set(target)
139 }
140}
141
142impl ErasedCallTarget for ImportHandle {
143 fn call<'a>(
144 &'a self,
145 capability: &'a CapabilityId,
146 context: CallContext,
147 input: SlotValue,
148 ) -> ErasedCallFuture<'a> {
149 ImportHandle::call(self, capability, context, input)
150 }
151}
152
153#[cfg(test)]
154mod tests {
155 use std::sync::atomic::{AtomicUsize, Ordering};
156 use std::sync::{Mutex, Weak};
157 use std::task::{Context, Poll, Waker};
158 use std::time::{Duration, Instant};
159
160 use boxology_contract::{
161 Caller, CancelToken, CapabilityDescriptor, CapabilityName, CapabilityShape,
162 ContractDescriptor, ContractRevision, ContractValue, Deadline, ExposureLevel, Idempotency,
163 ImplementationDescriptor, ImportDescriptor, TraceContext, TypeDescriptor,
164 };
165
166 use super::*;
167
168 enum Behavior {
169 Echo,
170 Error(ErasedCallError),
171 ConstructionPanic,
172 PollPanic,
173 }
174
175 struct Target {
176 calls: Arc<AtomicUsize>,
177 behavior: Behavior,
178 drops: Option<Arc<AtomicUsize>>,
179 }
180
181 struct RemoteTarget {
182 panic_on_call: bool,
183 panic_on_poll: bool,
184 }
185
186 impl ErasedCallTarget for RemoteTarget {
187 fn call<'a>(
188 &'a self,
189 _capability: &'a CapabilityId,
190 _context: CallContext,
191 input: SlotValue,
192 ) -> ErasedCallFuture<'a> {
193 assert!(!self.panic_on_call, "remote construction panic");
194 if self.panic_on_poll {
195 return Box::pin(std::future::poll_fn(|_| panic!("remote poll panic")));
196 }
197 Box::pin(ready(Ok(input)))
198 }
199 }
200
201 impl RemoteImportTarget for RemoteTarget {
202 fn supports_capability(&self, _capability: &CapabilityId) -> bool {
203 true
204 }
205 }
206
207 impl Drop for Target {
208 fn drop(&mut self) {
209 if let Some(drops) = &self.drops {
210 drops.fetch_add(1, Ordering::SeqCst);
211 }
212 }
213 }
214
215 impl ErasedTarget for Target {
216 fn call<'a>(
217 &'a self,
218 _capability: &'a CapabilityId,
219 _context: CallContext,
220 input: SlotValue,
221 ) -> ErasedCallFuture<'a> {
222 self.calls.fetch_add(1, Ordering::SeqCst);
223 match &self.behavior {
224 Behavior::Echo => Box::pin(ready(Ok(input))),
225 Behavior::Error(error) => Box::pin(ready(Err(error.clone()))),
226 Behavior::ConstructionPanic => panic!("construction panic"),
227 Behavior::PollPanic => Box::pin(std::future::poll_fn(|_| {
228 panic!("poll panic");
229 })),
230 }
231 }
232 }
233
234 fn box_id(value: &str) -> BoxId {
235 BoxId::new(value).unwrap()
236 }
237
238 fn capability(package: &str, name: &str) -> CapabilityId {
239 CapabilityId::new(box_id(package), CapabilityName::new(name).unwrap())
240 }
241
242 fn handle(names: &[&str]) -> ImportHandle {
243 let slot = box_id("service");
244 let capabilities = names
245 .iter()
246 .map(|name| capability("service", name))
247 .collect();
248 let imports = Imports::new([(slot.clone(), capabilities)]);
249 imports.handle(&slot).unwrap().clone()
250 }
251
252 fn context(deadline: Option<Deadline>) -> CallContext {
253 CallContext::new(
254 Caller::Anonymous,
255 deadline,
256 CancelToken::new(),
257 TraceContext::empty(),
258 None,
259 )
260 }
261
262 fn target(calls: &Arc<AtomicUsize>, behavior: Behavior) -> Arc<dyn ErasedTarget> {
263 Arc::new(Target {
264 calls: Arc::clone(calls),
265 behavior,
266 drops: None,
267 })
268 }
269
270 fn adapter(calls: &Arc<AtomicUsize>) -> Target {
271 Target {
272 calls: Arc::clone(calls),
273 behavior: Behavior::Echo,
274 drops: None,
275 }
276 }
277 fn implementation(
278 package: &str,
279 provided: &[&str],
280 imports: &[(&str, &[&str])],
281 ) -> ImplementationDescriptor {
282 let provided: Vec<_> = provided
283 .iter()
284 .map(|name| (*name, ExposureLevel::CodeOnly))
285 .collect();
286 implementation_at(package, &provided, imports)
287 }
288 fn implementation_at(
289 package: &str,
290 provided: &[(&str, ExposureLevel)],
291 imports: &[(&str, &[&str])],
292 ) -> ImplementationDescriptor {
293 let revision = ContractRevision::new("r1").unwrap();
294 let capabilities = provided.iter().map(|(name, maximum)| {
295 CapabilityDescriptor::new(
296 capability(package, name),
297 TypeDescriptor::bool(),
298 TypeDescriptor::bool(),
299 TypeDescriptor::bool(),
300 CapabilityShape::Unary,
301 *maximum,
302 Idempotency::None,
303 None,
304 )
305 });
306 let contract = Box::leak(Box::new(
307 ContractDescriptor::new(box_id(package), capabilities, revision.clone()).unwrap(),
308 ));
309 let imports = imports.iter().map(|(slot, names)| {
310 ImportDescriptor::new(
311 box_id(slot),
312 revision.clone(),
313 names.iter().map(|name| capability(slot, name)),
314 )
315 .unwrap()
316 });
317 ImplementationDescriptor::new(contract, imports).unwrap()
318 }
319 #[derive(Default)]
320 struct TransportProbe {
321 trace: Mutex<Vec<String>>,
322 runtimes: Mutex<Vec<Weak<TransportRuntime<()>>>>,
323 drops: AtomicUsize,
324 active_drops: Mutex<Vec<bool>>,
325 }
326 #[derive(Clone, Copy, PartialEq, Eq)]
327 enum ProbeFailure {
328 None,
329 Prepare,
330 Start,
331 }
332 struct ProbeBinding {
333 id: u8,
334 failure: ProbeFailure,
335 config: Arc<()>,
336 probe: Arc<TransportProbe>,
337 import: Option<ImportHandle>,
338 }
339 struct ProbeHandle {
340 id: u8,
341 runtime: Arc<TransportRuntime<()>>,
342 probe: Arc<TransportProbe>,
343 }
344 impl Drop for ProbeHandle {
345 fn drop(&mut self) {
346 self.probe.drops.fetch_add(1, Ordering::SeqCst);
347 self.probe
348 .active_drops
349 .lock()
350 .unwrap()
351 .push(self.runtime.is_active());
352 }
353 }
354 impl TransportHandle for ProbeHandle {
355 fn stop_intake(&self) {
356 self.record("stop");
357 }
358 fn cancel_tasks(&self) {
359 self.record("cancel");
360 }
361 fn abort_tasks(&self) {
362 self.record("abort");
363 }
364 fn join_tasks(self: Box<Self>) -> TransportJoinFuture {
365 Box::pin(ready(Ok(())))
366 }
367 }
368 impl ProbeHandle {
369 fn record(&self, phase: &str) {
370 self.probe
371 .trace
372 .lock()
373 .unwrap()
374 .push(format!("{phase}{}", self.id));
375 }
376 }
377 impl TransportBinding for ProbeBinding {
378 type Config = ();
379 type Handle = ProbeHandle;
380 fn config(&self) -> Arc<()> {
381 self.config.clone()
382 }
383 fn conform(
384 &self,
385 descriptor: &CapabilityDescriptor,
386 level: ExposureLevel,
387 ) -> Result<(), Detail> {
388 let name = descriptor.name().to_string();
389 let mut trace = self.probe.trace.lock().unwrap();
390 trace.push(format!("c{name}:{level:?}"));
391 drop(trace);
392 let rejected = self.id == 0
393 && ((name == "limited" && level == ExposureLevel::External)
394 || (name == "rejected" && level == ExposureLevel::Internal));
395 (!rejected)
396 .then_some(())
397 .ok_or_else(|| Detail::new("test_conformance"))
398 }
399 fn prepare(&self, descriptors: &[&'static CapabilityDescriptor]) -> Result<(), Detail> {
400 self.record('p', descriptors.iter().map(|item| item.id()));
401 match self.failure {
402 ProbeFailure::Prepare => Err(Detail::new("test_prepare")),
403 _ => Ok(()),
404 }
405 }
406 fn start(&self, runtime: TransportRuntime<()>) -> Result<ProbeHandle, Detail> {
407 let exposures = runtime.exposures();
408 let ids = exposures.iter().map(|item| item.descriptor().id());
409 self.record('s', ids);
410 assert!(std::ptr::eq(runtime.config(), self.config.as_ref()));
411 let mut active = Box::pin(runtime.wait_until_active());
412 assert!(!runtime.is_active() && matches!(poll_once(active.as_mut()), Poll::Pending));
413 drop(active);
414 if let Some(handle) = &self.import {
415 let capability = &handle.capabilities()[0];
416 let result = invoke(handle, capability, context(None), SlotValue::Null);
417 assert_eq!(
418 result,
419 Err(ErasedCallError::Unavailable(Detail::new("unsealed_import")))
420 );
421 }
422 if self.failure == ProbeFailure::Start {
423 return Err(Detail::new("test_start"));
424 }
425 let runtime = Arc::new(runtime);
426 let weak = Arc::downgrade(&runtime);
427 self.probe.runtimes.lock().unwrap().push(weak);
428 Ok(ProbeHandle {
429 id: self.id,
430 runtime,
431 probe: self.probe.clone(),
432 })
433 }
434 }
435 impl ProbeBinding {
436 fn new(id: u8, probe: &Arc<TransportProbe>, import: Option<ImportHandle>) -> Arc<Self> {
437 Self::with_failure(id, ProbeFailure::None, probe, import)
438 }
439 fn with_failure(
440 id: u8,
441 failure: ProbeFailure,
442 probe: &Arc<TransportProbe>,
443 import: Option<ImportHandle>,
444 ) -> Arc<Self> {
445 Arc::new(Self {
446 id,
447 failure,
448 config: Arc::new(()),
449 probe: probe.clone(),
450 import,
451 })
452 }
453 fn record<'a>(&self, phase: char, ids: impl Iterator<Item = &'a CapabilityId>) {
454 let ids = ids.map(ToString::to_string).collect::<Vec<_>>().join(",");
455 let mut trace = self.probe.trace.lock().unwrap();
456 trace.push(format!("{phase}{}:{ids}", self.id));
457 }
458 }
459 struct FailureSetup {
460 builder: CompositionBuilder,
461 handle: ImportHandle,
462 calls: Arc<AtomicUsize>,
463 target_drops: Arc<AtomicUsize>,
464 bindings: Vec<Weak<ProbeBinding>>,
465 configs: Vec<Weak<()>>,
466 }
467 fn failure_setup(
468 count: u8,
469 failure: ProbeFailure,
470 probe: &Arc<TransportProbe>,
471 ) -> FailureSetup {
472 let provider = box_id("provider");
473 let call = capability("provider", "call");
474 let calls = Arc::new(AtomicUsize::new(0));
475 let target_drops = Arc::new(AtomicUsize::new(0));
476 let mut captured = None;
477 let mut builder = CompositionBuilder::new();
478 builder.add_box(
479 implementation("consumer", &[], &[("provider", &["call"])]),
480 |imports| {
481 captured = Some(imports.handle(&provider).unwrap().clone());
482 adapter(&calls)
483 },
484 );
485 builder.add_box(implementation("provider", &["call"], &[]), |_| Target {
486 calls: calls.clone(),
487 behavior: Behavior::Echo,
488 drops: Some(target_drops.clone()),
489 });
490 builder.resolve_import(
491 box_id("consumer"),
492 provider.clone(),
493 ImportTarget::local(provider.clone()),
494 );
495 let handle = captured.unwrap();
496 let mut bindings = Vec::new();
497 let mut configs = Vec::new();
498 for id in 1..=count {
499 let failure = if id == count {
500 failure
501 } else {
502 ProbeFailure::None
503 };
504 let binding = ProbeBinding::with_failure(id, failure, probe, Some(handle.clone()));
505 bindings.push(Arc::downgrade(&binding));
506 configs.push(Arc::downgrade(&binding.config));
507 builder.expose(
508 provider.clone(),
509 call.clone(),
510 binding,
511 ExposureLevel::CodeOnly,
512 );
513 }
514 FailureSetup {
515 builder,
516 handle,
517 calls,
518 target_drops,
519 bindings,
520 configs,
521 }
522 }
523 fn lifecycle(probe: &TransportProbe) -> Vec<String> {
524 probe
525 .trace
526 .lock()
527 .unwrap()
528 .iter()
529 .filter(|event| !event.starts_with("ccall:"))
530 .cloned()
531 .collect()
532 }
533 fn poll_once<F: Future + ?Sized>(future: Pin<&mut F>) -> Poll<F::Output> {
534 future.poll(&mut Context::from_waker(Waker::noop()))
535 }
536 fn invoke(
537 handle: &ImportHandle,
538 capability: &CapabilityId,
539 context: CallContext,
540 input: SlotValue,
541 ) -> Result<SlotValue, ErasedCallError> {
542 let mut future = handle.call(capability, context, input);
543 loop {
544 match future
545 .as_mut()
546 .poll(&mut Context::from_waker(Waker::noop()))
547 {
548 Poll::Ready(output) => return output,
549 Poll::Pending => {}
550 }
551 }
552 }
553
554 #[test]
555 fn short_circuits_in_required_order_without_provider_invocation() {
556 let handle = handle(&["allowed"]);
557 let allowed = capability("service", "allowed");
558 let undeclared = capability("service", "undeclared");
559 let deadline = Deadline::at(Instant::now());
560 assert_eq!(deadline.remaining(), Duration::ZERO);
561 let calls = Arc::new(AtomicUsize::new(0));
562 let provider = target(&calls, Behavior::Echo);
563
564 assert_eq!(
565 invoke(
566 &handle,
567 &undeclared,
568 context(Some(deadline)),
569 SlotValue::Null,
570 ),
571 Err(ErasedCallError::Unavailable(Detail::new("unsealed_import")))
572 );
573 assert_eq!(calls.load(Ordering::SeqCst), 0);
574
575 assert!(handle.seal(provider).is_ok());
576 assert_eq!(
577 invoke(
578 &handle,
579 &undeclared,
580 context(Some(deadline)),
581 SlotValue::Null,
582 ),
583 Err(ErasedCallError::ContractViolation(Detail::new(
584 "undeclared_import_capability"
585 )))
586 );
587 assert_eq!(calls.load(Ordering::SeqCst), 0);
588 assert_eq!(
589 invoke(&handle, &allowed, context(Some(deadline)), SlotValue::Null,),
590 Err(ErasedCallError::Deadline)
591 );
592 assert_eq!(calls.load(Ordering::SeqCst), 0);
593 }
594
595 #[test]
596 fn sealed_calls_preserve_success_and_erased_errors() {
597 let capability = capability("service", "call");
598 let input = SlotValue::Value(ContractValue::string("payload"));
599 let success = handle(&["call"]);
600 let success_calls = Arc::new(AtomicUsize::new(0));
601 assert!(success.seal(target(&success_calls, Behavior::Echo)).is_ok());
602 let deadline = Deadline::at(Instant::now() + Duration::from_secs(3600));
603 assert_eq!(
604 invoke(
605 &success,
606 &capability,
607 context(Some(deadline)),
608 input.clone()
609 ),
610 Ok(input.clone())
611 );
612 assert_eq!(success_calls.load(Ordering::SeqCst), 1);
613
614 let expected = ErasedCallError::Domain {
615 error_tag: "ordinary".into(),
616 payload: input,
617 };
618 let failure = handle(&["call"]);
619 let failure_calls = Arc::new(AtomicUsize::new(0));
620 assert!(
621 failure
622 .seal(target(&failure_calls, Behavior::Error(expected.clone())))
623 .is_ok()
624 );
625 assert_eq!(
626 invoke(&failure, &capability, context(None), SlotValue::Null),
627 Err(expected)
628 );
629 assert_eq!(failure_calls.load(Ordering::SeqCst), 1);
630 }
631
632 #[test]
633 fn dispatch_panics_reach_internal_through_the_handle() {
634 let capability = capability("service", "call");
635 for behavior in [Behavior::ConstructionPanic, Behavior::PollPanic] {
636 let handle = handle(&["call"]);
637 let calls = Arc::new(AtomicUsize::new(0));
638 assert!(handle.seal(target(&calls, behavior)).is_ok());
639 let error = invoke(&handle, &capability, context(None), SlotValue::Null).unwrap_err();
640 let ErasedCallError::Internal(detail) = error else {
641 panic!("expected Internal, got {error:?}");
642 };
643 assert_eq!(detail.code(), "panic");
644 assert_eq!(calls.load(Ordering::SeqCst), 1);
645 }
646 }
647
648 #[test]
649 fn clones_observe_the_same_seal() {
650 let original = handle(&["call"]);
651 let clone = original.clone();
652 let calls = Arc::new(AtomicUsize::new(0));
653 assert!(original.seal(target(&calls, Behavior::Echo)).is_ok());
654
655 assert_eq!(
656 invoke(
657 &clone,
658 &capability("service", "call"),
659 context(None),
660 SlotValue::Null,
661 ),
662 Ok(SlotValue::Null)
663 );
664 assert_eq!(calls.load(Ordering::SeqCst), 1);
665 }
666 #[test]
667 fn composition_starts_grouped_transports_then_commits_all_traffic() {
668 let provider = box_id("provider");
669 let imported = capability("provider", "call");
670 let [b, c] = ["b", "c"].map(|name| capability("provider", name));
671 let selected = ImportTarget::local(provider.clone());
672 assert_eq!(selected, selected.clone());
673 let mut captured = None;
674 let mut factories = 0;
675 let consumer_calls = Arc::new(AtomicUsize::new(0));
676 let provider_calls = Arc::new(AtomicUsize::new(0));
677 let mut builder = CompositionBuilder::new();
678 builder.add_box(
679 implementation("consumer", &[], &[("provider", &["call"])]),
680 |imports| {
681 factories += 1;
682 captured = Some(imports.handle(&provider).unwrap().clone());
683 adapter(&consumer_calls)
684 },
685 );
686 assert_eq!(factories, 1);
687 builder.add_box(implementation("provider", &["call", "b", "c"], &[]), |_| {
688 factories += 1;
689 adapter(&provider_calls)
690 });
691 builder.resolve_import(box_id("consumer"), provider.clone(), selected);
692 let handle = captured.unwrap();
693 let probe = Arc::new(TransportProbe::default());
694 let first = ProbeBinding::new(1, &probe, Some(handle.clone()));
695 let second = ProbeBinding::new(2, &probe, Some(handle.clone()));
696 let binding_weaks = [Arc::downgrade(&first), Arc::downgrade(&second)];
697 let config_weaks = [
698 Arc::downgrade(&first.config),
699 Arc::downgrade(&second.config),
700 ];
701 assert!(!Arc::ptr_eq(&first, &second) && !Arc::ptr_eq(&first.config, &second.config));
702 let level = ExposureLevel::CodeOnly;
703 for (capability, binding) in [
704 (imported.clone(), first.clone()),
705 (b, second.clone()),
706 (c.clone(), first.clone()),
707 (c, first.clone()),
708 ] {
709 builder.expose(provider.clone(), capability, binding, level);
710 }
711 assert_eq!(builder.validate(), Ok(()));
712 assert_eq!(builder.validate(), Ok(()));
713 assert_eq!(
714 invoke(&handle, &imported, context(None), SlotValue::Null),
715 Err(ErasedCallError::Unavailable(Detail::new("unsealed_import")))
716 );
717 assert_eq!(provider_calls.load(Ordering::SeqCst), 0);
718 drop((first, second));
719 let _composition = builder.start().unwrap();
720 let trace = probe.trace.lock().unwrap();
721 assert_eq!(
722 trace[trace.len() - 4..].join("|"),
723 "p1:provider.call,provider.c,provider.c|p2:provider.b|s1:provider.call,provider.c,provider.c|s2:provider.b"
724 );
725 drop(trace);
726 let runtimes: Vec<_> = probe
727 .runtimes
728 .lock()
729 .unwrap()
730 .iter()
731 .map(Weak::upgrade)
732 .collect();
733 let runtimes: Vec<_> = runtimes.into_iter().map(Option::unwrap).collect();
734 assert!(TransportTaskTracker::ptr_eq(
735 runtimes[0].tracker(),
736 runtimes[1].tracker()
737 ));
738 assert!(runtimes.iter().all(|runtime| runtime.is_active()));
739 assert_eq!(
740 invoke(&handle, &imported, context(None), SlotValue::Null),
741 Ok(SlotValue::Null)
742 );
743 assert_eq!(provider_calls.load(Ordering::SeqCst), 1);
744 assert_eq!(consumer_calls.load(Ordering::SeqCst), 0);
745 assert_eq!(factories, 2);
746 assert!(binding_weaks.iter().all(|weak| weak.upgrade().is_some()));
747 assert!(config_weaks.iter().all(|weak| weak.upgrade().is_some()));
748 assert_eq!(probe.drops.load(Ordering::SeqCst), 0);
749 }
750
751 #[test]
752 fn composition_seals_remote_target_without_a_local_provider_and_contains_panics() {
753 let slot = box_id("remote");
754 let imported = capability("remote", "call");
755 let remote: Arc<dyn RemoteImportTarget> = Arc::new(RemoteTarget {
756 panic_on_call: false,
757 panic_on_poll: false,
758 });
759 let selected = ImportTarget::remote(remote.clone());
760 assert_eq!(selected, selected.clone());
761 assert_eq!(format!("{selected:?}"), "ImportTarget(Remote(<redacted>))");
762 assert_ne!(
763 selected,
764 ImportTarget::remote(Arc::new(RemoteTarget {
765 panic_on_call: false,
766 panic_on_poll: false,
767 }))
768 );
769 assert_ne!(selected, ImportTarget::local(slot.clone()));
770
771 for (target, expected) in [
772 (selected, Ok(SlotValue::Null)),
773 (
774 ImportTarget::remote(Arc::new(RemoteTarget {
775 panic_on_call: true,
776 panic_on_poll: false,
777 })),
778 Err(ErasedCallError::Internal(
779 Detail::new("panic").with_message("remote construction panic"),
780 )),
781 ),
782 (
783 ImportTarget::remote(Arc::new(RemoteTarget {
784 panic_on_call: false,
785 panic_on_poll: true,
786 })),
787 Err(ErasedCallError::Internal(
788 Detail::new("panic").with_message("remote poll panic"),
789 )),
790 ),
791 ] {
792 let mut captured = None;
793 let mut builder = CompositionBuilder::new();
794 builder.add_box(
795 implementation("consumer", &[], &[("remote", &["call"])]),
796 |imports| {
797 captured = Some(imports.handle(&slot).unwrap().clone());
798 adapter(&Arc::new(AtomicUsize::new(0)))
799 },
800 );
801 builder.resolve_import(box_id("consumer"), slot.clone(), target);
802 let handle = captured.unwrap();
803 assert_eq!(
804 invoke(&handle, &imported, context(None), SlotValue::Null),
805 Err(ErasedCallError::Unavailable(Detail::new("unsealed_import")))
806 );
807 let _composition = builder.start().unwrap();
808 assert_eq!(
809 invoke(&handle, &imported, context(None), SlotValue::Null),
810 expected
811 );
812 }
813 }
814
815 #[test]
816 fn prepare_failure_returns_without_starting_or_retaining_ownership() {
817 let probe = Arc::new(TransportProbe::default());
818 let FailureSetup {
819 builder,
820 handle,
821 calls,
822 target_drops,
823 bindings,
824 configs,
825 } = failure_setup(2, ProbeFailure::Prepare, &probe);
826 let error = builder.start().err().expect("prepare failure started");
827 assert_eq!(
828 lifecycle(&probe).join("|"),
829 "p1:provider.call|p2:provider.call"
830 );
831 assert_eq!(
832 error.errors(),
833 &[AssemblyError::TransportPrepareFailed {
834 detail: Detail::new("test_prepare"),
835 }]
836 );
837 assert_eq!(error.to_string(), "transport prepare failed: test_prepare");
838 assert_eq!(
839 invoke(
840 &handle,
841 &capability("provider", "call"),
842 context(None),
843 SlotValue::Null
844 ),
845 Err(ErasedCallError::Unavailable(Detail::new("unsealed_import")))
846 );
847 assert_eq!(calls.load(Ordering::SeqCst), 0);
848 assert!(probe.runtimes.lock().unwrap().is_empty());
849 assert_eq!(probe.drops.load(Ordering::SeqCst), 0);
850 assert!(probe.active_drops.lock().unwrap().is_empty());
851 assert!(bindings.iter().all(|weak| weak.upgrade().is_none()));
852 assert!(configs.iter().all(|weak| weak.upgrade().is_none()));
853 assert_eq!(target_drops.load(Ordering::SeqCst), 1);
854 }
855
856 #[test]
857 fn start_failure_rolls_back_three_handles_in_global_reverse_phases() {
858 let probe = Arc::new(TransportProbe::default());
859 let FailureSetup {
860 builder,
861 handle,
862 calls,
863 target_drops,
864 bindings,
865 configs,
866 } = failure_setup(4, ProbeFailure::Start, &probe);
867 let error = builder.start().err().expect("start failure committed");
868 let events = lifecycle(&probe);
869 assert_eq!(
870 events.join("|"),
871 "p1:provider.call|p2:provider.call|p3:provider.call|p4:provider.call|s1:provider.call|s2:provider.call|s3:provider.call|s4:provider.call|stop3|stop2|stop1|cancel3|cancel2|cancel1|abort3|abort2|abort1"
872 );
873 assert!(events[..4].iter().all(|event| event.starts_with('p')));
874 assert!(events[4..8].iter().all(|event| event.starts_with('s')));
875 assert_eq!(
876 error.errors(),
877 &[AssemblyError::TransportStartFailed {
878 detail: Detail::new("test_start"),
879 }]
880 );
881 assert_eq!(error.to_string(), "transport start failed: test_start");
882 assert_eq!(
883 invoke(
884 &handle,
885 &capability("provider", "call"),
886 context(None),
887 SlotValue::Null
888 ),
889 Err(ErasedCallError::Unavailable(Detail::new("unsealed_import")))
890 );
891 assert_eq!(calls.load(Ordering::SeqCst), 0);
892 let runtimes = probe.runtimes.lock().unwrap();
893 assert_eq!(runtimes.len(), 3);
894 assert!(runtimes.iter().all(|weak| weak.upgrade().is_none()));
895 drop(runtimes);
896 assert_eq!(probe.drops.load(Ordering::SeqCst), 3);
897 let active_drops = probe.active_drops.lock().unwrap();
898 assert_eq!(active_drops.len(), 3);
899 assert!(active_drops.iter().all(|active| !active));
900 drop(active_drops);
901 assert!(bindings.iter().all(|weak| weak.upgrade().is_none()));
902 assert!(configs.iter().all(|weak| weak.upgrade().is_none()));
903 assert_eq!(target_drops.load(Ordering::SeqCst), 1);
904 }
905
906 #[test]
907 fn composition_aggregates_every_failure_with_exact_precedence_and_no_seals() {
908 use ExposureLevel::{CodeOnly as C, External as E, Internal as I};
909 let c = box_id("consumer");
910 let valid = capability("consumer", "valid");
911 let limited = capability("consumer", "limited");
912 let rejected = capability("consumer", "rejected");
913 let calls = Arc::new(AtomicUsize::new(0));
914 let mut handles = Vec::new();
915 let mut factories = 0;
916 let mut builder = CompositionBuilder::new();
917 builder.add_box(
918 implementation_at(
919 "consumer",
920 &[("valid", E), ("limited", I), ("rejected", E)],
921 &[
922 ("missing", &["call"]),
923 ("partial", &["first", "second"]),
924 ("duplicate", &["call"]),
925 ("unknown", &["call"]),
926 ],
927 ),
928 |imports| {
929 factories += 1;
930 handles = imports.cloned_handles().into_values().collect();
931 adapter(&calls)
932 },
933 );
934 builder.add_box(implementation("consumer", &["only_second"], &[]), |_| {
935 factories += 1;
936 adapter(&calls)
937 });
938 builder.add_box(implementation("consumer", &[], &[]), |_| {
939 factories += 1;
940 adapter(&calls)
941 });
942 assert_eq!(factories, 3);
943 builder.add_box(implementation("partial", &[], &[]), |_| {
944 factories += 1;
945 adapter(&calls)
946 });
947 assert_eq!(factories, 4);
948 let mut resolve = |c, s, t| {
949 builder.resolve_import(box_id(c), box_id(s), ImportTarget::local(box_id(t)));
950 };
951 resolve("ghost", "slot", "partial");
952 resolve("ghost", "slot", "partial");
953 resolve("consumer", "undeclared", "partial");
954 resolve("consumer", "undeclared", "partial");
955 resolve("consumer", "duplicate", "partial");
956 resolve("consumer", "duplicate", "absent");
957 resolve("consumer", "unknown", "absent");
958 resolve("consumer", "partial", "partial");
959 let probe = Arc::new(TransportProbe::default());
960 let binding = ProbeBinding::new(0, &probe, None);
961 for (owner, capability, level) in [
962 (box_id("ghost"), rejected.clone(), E),
963 (c.clone(), capability("other", "rejected"), E),
964 (c.clone(), limited.clone(), E),
965 (c.clone(), rejected, I),
966 (c.clone(), capability("consumer", "only_second"), E),
967 (c.clone(), valid.clone(), E),
968 (c.clone(), limited, C),
969 (c.clone(), valid.clone(), I),
970 (c.clone(), valid, I),
971 ] {
972 builder.expose(owner, capability, binding.clone(), level);
973 }
974 use AssemblyError::*;
975 let unknown_consumer = |name| UnknownImportConsumer {
976 consumer: box_id(name),
977 };
978 let unknown_slot = || UnknownImportSlot {
979 consumer: c.clone(),
980 slot: box_id("undeclared"),
981 };
982 let missing_capability = |name| MissingImportedCapability {
983 consumer: c.clone(),
984 slot: box_id("partial"),
985 capability: capability("partial", name),
986 };
987 let expected = vec![
988 DuplicateBox { box_id: c.clone() },
989 DuplicateBox { box_id: c.clone() },
990 unknown_consumer("ghost"),
991 unknown_consumer("ghost"),
992 unknown_slot(),
993 unknown_slot(),
994 DuplicateImportResolution {
995 consumer: c.clone(),
996 slot: box_id("duplicate"),
997 },
998 UnknownImportTarget {
999 consumer: c.clone(),
1000 slot: box_id("unknown"),
1001 target: box_id("absent"),
1002 },
1003 MissingImportResolution {
1004 consumer: c.clone(),
1005 slot: box_id("missing"),
1006 },
1007 missing_capability("first"),
1008 missing_capability("second"),
1009 ];
1010 let validated = builder.validate().unwrap_err();
1011 assert_eq!(&validated.errors()[..expected.len()], expected);
1012 let display = validated.to_string();
1013 assert_eq!(
1014 display
1015 .lines()
1016 .skip(expected.len())
1017 .collect::<Vec<_>>()
1018 .join("|"),
1019 "unknown exposure provider: ghost|unknown exposed capability other.rejected for provider consumer|exposure external exceeds maximum internal for capability consumer.limited|transport conformance failed for capability consumer.rejected: test_conformance|unknown exposed capability consumer.only_second for provider consumer"
1020 );
1021 let conformed =
1022 "crejected:Internal|cvalid:External|climited:CodeOnly|cvalid:Internal|cvalid:Internal";
1023 assert_eq!(probe.trace.lock().unwrap().join("|"), conformed);
1024 assert_eq!(builder.validate().unwrap_err(), validated);
1025 assert_eq!(builder.validate().unwrap_err().to_string(), display);
1026 let started = builder.start().err().expect("invalid composition started");
1027 assert_eq!(started, validated);
1028 assert_eq!(started.to_string(), display);
1029 for handle in handles {
1030 let capability = handle.capabilities()[0].clone();
1031 assert_eq!(
1032 invoke(&handle, &capability, context(None), SlotValue::Null),
1033 Err(ErasedCallError::Unavailable(Detail::new("unsealed_import")))
1034 );
1035 }
1036 assert_eq!(calls.load(Ordering::SeqCst), 0);
1037 assert_eq!(
1038 probe.trace.lock().unwrap().join("|"),
1039 [conformed; 4].join("|")
1040 );
1041 }
1042
1043 #[test]
1044 fn public_carriers_and_returned_future_are_thread_safe() {
1045 fn assert_bounds<T: Send + Sync + 'static>() {}
1046 fn assert_send<T: Send>(value: T) -> T {
1047 value
1048 }
1049
1050 assert_bounds::<Imports>();
1051 assert_bounds::<ImportHandle>();
1052 assert_bounds::<CompositionBuilder>();
1053 assert_bounds::<ImportTarget>();
1054 assert_bounds::<Composition>();
1055 let service = box_id("service");
1056 let capabilities = vec![
1057 capability("service", "second"),
1058 capability("service", "first"),
1059 ];
1060 let imports = Imports::new([(service.clone(), capabilities.clone())]);
1061 let handle = imports.handle(&service).unwrap();
1062 assert_eq!(handle.slot_id(), &service);
1063 assert_eq!(handle.capabilities(), capabilities);
1064 assert!(imports.handle(&box_id("missing")).is_none());
1065
1066 let calls = Arc::new(AtomicUsize::new(0));
1067 assert!(handle.seal(target(&calls, Behavior::Echo)).is_ok());
1068 let capability = capability("service", "first");
1069 let carrier: &dyn ErasedCallTarget = handle;
1070 let mut future = assert_send(carrier.call(&capability, context(None), SlotValue::Null));
1071 assert_eq!(poll_once(future.as_mut()), Poll::Ready(Ok(SlotValue::Null)));
1072 assert_eq!(calls.load(Ordering::SeqCst), 1);
1073 }
1074}