1use std::time::Duration;
20
21use crate::{Error, Result};
22use zenkey::origin::{HostId, ServiceOrigin};
23use zenkey::qos::QosProfile;
24use zenkey::{Declared, Fanout, ProcedureKind};
25use zenoh::Session;
26
27use crate::bus::query::FleetAnswer;
28use crate::model::registry::SliceSet;
29use crate::report::{
30 CallAnswer, CallError, CallOutcome, CallReport, ConcurrentLane, HlcReference, TRACE_CHAIN_RULE,
31 TRACE_EXCLUDED, TraceReport,
32};
33
34pub struct Publication {
36 publisher: zenoh::pubsub::Publisher<'static>,
37 encoding: Option<String>,
38}
39
40impl std::fmt::Debug for Publication {
41 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
42 f.debug_struct("Publication")
43 .field("key", &self.publisher.key_expr().as_str())
44 .finish_non_exhaustive()
45 }
46}
47
48pub async fn declare_publication(
55 session: &Session,
56 key: &str,
57 qos: QosProfile,
58 encoding: Option<&str>,
59) -> Result<Publication> {
60 let publisher = session
61 .declare_publisher(key.to_string())
62 .reliability(qos.reliability())
63 .congestion_control(qos.congestion_control())
64 .priority(qos.priority())
65 .express(qos.express())
66 .await
67 .map_err(|e| Error::bus("declare publisher", key, e))?;
68 Ok(Publication {
69 publisher,
70 encoding: encoding.map(str::to_string),
71 })
72}
73
74impl Publication {
75 pub async fn send(&self, payload: Vec<u8>, attachment: Option<Vec<u8>>) -> Result<()> {
80 self.send_stamped(payload, attachment, None).await
81 }
82
83 pub async fn send_stamped(
89 &self,
90 payload: Vec<u8>,
91 attachment: Option<Vec<u8>>,
92 timestamp: Option<zenoh::time::Timestamp>,
93 ) -> Result<()> {
94 let put = self.publisher.put(payload);
95 let put = match &self.encoding {
96 Some(e) => put.encoding(e.as_str()),
97 None => put,
98 };
99 let put = match attachment {
100 Some(a) => put.attachment(a),
101 None => put,
102 };
103 let put = match timestamp {
104 Some(ts) => put.timestamp(ts),
105 None => put,
106 };
107 put.await
108 .map_err(|e| Error::bus("put", self.publisher.key_expr().as_str(), e))
109 }
110
111 pub async fn retire(&self) -> Result<()> {
117 self.publisher
118 .delete()
119 .await
120 .map_err(|e| Error::bus("delete", self.publisher.key_expr().as_str(), e))
121 }
122
123 pub async fn undeclare(self) -> Result<()> {
125 self.publisher
126 .undeclare()
127 .await
128 .map_err(|e| Error::bus("undeclare publisher", "", e))
129 }
130
131 pub async fn matching_status(&self) -> Result<bool> {
138 self.publisher
139 .matching_status()
140 .await
141 .map(|s| s.matching())
142 .map_err(|e| Error::bus("matching status", "", e))
143 }
144
145 pub async fn matching_events(&self) -> Result<MatchingEvents> {
148 let listener = self
149 .publisher
150 .matching_listener()
151 .await
152 .map_err(|e| Error::bus("matching listener", "", e))?;
153 Ok(MatchingEvents { listener })
154 }
155}
156
157#[derive(Debug, Clone, PartialEq, Eq)]
160pub enum RetireClass {
161 State {
163 registered: bool,
167 ttl_s: Option<i64>,
170 },
171 NonState { class: String },
174 Unclassified { reason: String },
176}
177
178pub fn check_retire(
188 base: &str,
189 key: &str,
190 slices: Option<&SliceSet>,
191 force: bool,
192) -> Result<RetireClass> {
193 if key.contains('*') || key.contains('$') {
194 return Err(Error::unaskable(
195 key,
196 "is a wildcard — a tombstone is addressed to one concrete key; a \
197 wildcard delete is not an operator act, it is a blast radius \
198 (RFC 04 §1.2, v1.12). Not overridable.",
199 ));
200 }
201 let facts = crate::model::facts::describe_key(base, key, slices).facts;
202 use crate::model::facts::{ClassKind, KeyShape, Registration};
203 match &facts.shape {
204 KeyShape::V1(v) if v.class_kind == ClassKind::State => {
205 let (registered, ttl_s) = match &facts.registration {
206 Registration::Registered(s) => (true, s.ttl_s),
207 _ => (false, None),
208 };
209 Ok(RetireClass::State { registered, ttl_s })
210 }
211 KeyShape::V1(v) if matches!(v.class_kind, ClassKind::Telemetry | ClassKind::Events) => {
212 if force {
213 return Ok(RetireClass::NonState {
214 class: v.class.clone(),
215 });
216 }
217 Err(Error::unaskable(
218 key,
219 format!(
220 "is {}-shaped — RFC 04 §1: a delete there is meaningless and \
221 MUST NOT be sent by the class's publisher. Retiring it anyway \
222 is an operator cleanup (RFC 04 §1.2, v1.12) — pass --i-know \
223 to mean it.",
224 v.class
225 ),
226 ))
227 }
228 KeyShape::V1(v) => {
229 if force {
230 return Ok(RetireClass::NonState {
231 class: v.class.clone(),
232 });
233 }
234 Err(Error::unaskable(
235 key,
236 format!(
237 "sits on the {} plane — a plane key answers GETs or carries \
238 frames; a tombstone there is at most a storage purge \
239 (RFC 04 §1.2, v1.12) — pass --i-know to mean it.",
240 v.class
241 ),
242 ))
243 }
244 KeyShape::NotUnderBase | KeyShape::Unparsed { .. } => {
245 let reason = match &facts.shape {
246 KeyShape::Unparsed { reason } => reason.clone(),
247 _ => format!("not under base {base:?}"),
248 };
249 if force {
250 return Ok(RetireClass::Unclassified { reason });
251 }
252 Err(Error::unaskable(
253 key,
254 format!(
255 "cannot be classified under base {base:?} ({reason}) — 'not \
256 asked' is not 'state' (RFC 09 §5.1 O4); pass --i-know to \
257 retire an unclassified key."
258 ),
259 ))
260 }
261 }
262}
263
264pub struct MatchingEvents {
270 listener: zenoh::matching::MatchingListener<
271 zenoh::handlers::FifoChannelHandler<zenoh::matching::MatchingStatus>,
272 >,
273}
274
275impl MatchingEvents {
276 pub(crate) async fn for_querier(querier: &zenoh::query::Querier<'_>) -> Result<Self> {
277 let listener = querier
278 .matching_listener()
279 .await
280 .map_err(|e| Error::bus("matching listener", "", e))?;
281 Ok(MatchingEvents { listener })
282 }
283
284 pub async fn recv(&self) -> Option<bool> {
287 self.listener.recv_async().await.ok().map(|s| s.matching())
288 }
289
290 pub fn stream(&self) -> impl futures_core::Stream<Item = bool> + '_ {
296 futures_util::StreamExt::map(self.listener.stream(), |s| s.matching())
297 }
298}
299
300#[derive(Debug, Clone, PartialEq, Eq)]
304pub enum CallTarget {
305 Host(HostId),
307 Fleet,
310 Service(ServiceOrigin),
312}
313
314impl CallTarget {
315 pub fn parse(s: &str) -> Result<CallTarget> {
319 if s == "*" {
320 return Ok(CallTarget::Fleet);
321 }
322 if s.starts_with('@') {
323 return Ok(CallTarget::Service(ServiceOrigin::new(s)?));
324 }
325 HostId::parse(s).map(CallTarget::Host).map_err(|e| {
326 Error::unaskable(
327 "origin",
328 format!("{e} — a hostname is not an origin; resolve it first (RFC 06 §6)"),
329 )
330 })
331 }
332}
333
334fn attachment_value(bytes: &[u8]) -> serde_json::Value {
340 if let Ok(v) = serde_json::from_slice::<serde_json::Value>(bytes) {
341 v
342 } else if let Ok(s) = std::str::from_utf8(bytes) {
343 serde_json::Value::String(s.to_string())
344 } else {
345 serde_json::Value::String(format!("<{} bytes>", bytes.len()))
346 }
347}
348
349pub struct CallSpec<'a> {
356 pub target: &'a CallTarget,
357 pub producer: &'a str,
358 pub procedure: &'a str,
360 pub params: &'a [String],
362 pub body: Option<Vec<u8>>,
364 pub attachment: Option<Vec<u8>>,
367 pub timeout: Duration,
368 pub slices: Option<&'a SliceSet>,
371}
372
373pub async fn call(fleet: &crate::Fleet<'_>, spec: CallSpec<'_>) -> Result<CallReport> {
389 let (key, timeout, answers) = call_answers(fleet, spec).await?;
390 Ok(project_call(key, timeout, &answers))
391}
392
393async fn call_answers(
398 fleet: &crate::Fleet<'_>,
399 spec: CallSpec<'_>,
400) -> Result<(String, Duration, Vec<FleetAnswer>)> {
401 let CallSpec {
402 target,
403 producer,
404 procedure,
405 params,
406 body,
407 attachment,
408 timeout,
409 slices,
410 } = spec;
411 if matches!(target, CallTarget::Fleet)
412 && let Some(slices) = slices
413 && let Some(slice) = slices.get(producer)
414 && let Some(proc_decl) = slice.procedures.iter().find(|p| p.path == procedure)
415 {
416 let forbidden = match proc_decl.fanout.as_ref().and_then(Declared::known) {
422 Some(Fanout::Forbidden) => true,
423 Some(Fanout::Allowed) => false,
424 None => matches!(
429 proc_decl.kind.as_ref().and_then(Declared::known),
430 Some(ProcedureKind::Write)
431 ),
432 };
433 if forbidden {
434 let declared = match proc_decl.fanout.as_ref() {
441 Some(f) if f.is(&Fanout::Forbidden) => {
442 "declares fanout = \"forbidden\"".to_string()
443 }
444 Some(f) => format!(
445 "declares fanout = {:?}, a token this build does not know — \
446 RFC 08 §2 defaults a write to forbidden and an unreadable \
447 spelling is not a licence",
448 f.token()
449 ),
450 None => "is a write with no declared fanout, which defaults to forbidden \
451 (RFC 08 §2)"
452 .to_string(),
453 };
454 return Err(Error::unaskable(
455 format!("procedure {producer}/{procedure}"),
456 format!(
457 "{declared} — a fleet (`*`) call to it is refused \
458 (RFC 05 §2.1); name one origin"
459 ),
460 ));
461 }
462 }
463
464 let segments: Vec<&str> = procedure.split('/').collect();
465 let relative = match target {
466 CallTarget::Host(id) => {
467 let origin = zenkey::origin::RemoteOrigin::from_host(id.clone());
468 zenkey::selector::rpc_at(&origin, producer, &segments).to_string()
469 }
470 CallTarget::Fleet => zenkey::selector::fleet_rpc(producer, &segments).to_string(),
471 CallTarget::Service(origin) => zenkey::selector::service_rpc(origin, &segments).to_string(),
472 };
473 let mut key = fleet.wire(relative);
474 if !params.is_empty() {
475 key.push('?');
476 key.push_str(¶ms.join(";"));
477 }
478
479 let answers = crate::bus::query::fleet_get(
480 fleet,
481 &key,
482 &crate::bus::query::GetOpts::new(timeout)
483 .payload(body)
484 .attachment(attachment),
485 )
486 .await?;
487 Ok((key, timeout, answers))
488}
489
490fn project_call(key: String, timeout: Duration, answers: &[FleetAnswer]) -> CallReport {
492 CallReport {
493 key,
494 timeout_s: timeout.as_secs_f64(),
497 answers: answers
498 .iter()
499 .map(|a| {
500 let (att, att_bytes) = match &a.attachment {
504 Some(z) => {
505 let bytes = z.to_bytes();
506 (Some(attachment_value(&bytes)), Some(bytes.len()))
507 }
508 None => (None, None),
509 };
510 let outcome = match &a.answer {
511 crate::bus::query::Answer::Value(bytes) => {
512 let bytes = bytes.to_bytes();
513 match serde_json::from_slice::<serde_json::Value>(&bytes) {
514 Ok(v) => CallOutcome::Ok {
515 value: Some(v),
516 text: None,
517 },
518 Err(_) => CallOutcome::Ok {
519 value: None,
520 text: Some(String::from_utf8_lossy(&bytes).to_string()),
521 },
522 }
523 }
524 crate::bus::query::Answer::Error { name, message } => {
525 CallOutcome::Err(CallError {
526 name: name.clone(),
527 message: message.clone(),
528 })
529 }
530 };
531 CallAnswer {
532 origin: a.origin.clone(),
533 outcome,
534 attachment: att,
535 attachment_bytes: att_bytes,
536 }
537 })
538 .collect(),
539 }
540}
541
542#[derive(Debug, Clone, Copy, PartialEq, Eq)]
544pub struct TraceSpec {
545 pub window: Duration,
550}
551
552const ORIGIN_CAPACITY: usize = 4096;
556const FLEET_CAPACITY: usize = 8192;
559
560pub async fn call_traced(
585 fleet: &crate::Fleet<'_>,
586 spec: CallSpec<'_>,
587 trace: TraceSpec,
588) -> Result<TraceReport> {
589 use crate::bus::monitor::{FleetEvent, Monitor, MonitorSpec, StreamItem};
590 use crate::model::examples::Examples;
591 use crate::model::facts::describe_key;
592 use crate::model::timeline::TimelineRow;
593 use crate::model::trace::{TraceTarget, idiom_of, trace_row};
594 use zenkey::selector::{Scope, all_under};
595
596 let (scope, producer) = match spec.target {
597 CallTarget::Fleet => {
598 return Err(Error::unaskable(
599 "--trace",
600 "a trace attributes what it sees to one origin, and a fleet (`*`) call \
601 has none to attribute to — name one origin",
602 ));
603 }
604 CallTarget::Host(id) => (
605 Scope::origin(&zenkey::origin::RemoteOrigin::from_host(id.clone())),
606 Some(spec.producer.to_string()),
607 ),
608 CallTarget::Service(o) => (Scope::origin(o), None),
609 };
610 let origin = scope.chunk().to_string();
611 let slices = spec.slices;
612 let idiom = idiom_of(
613 slices
614 .and_then(|s| s.get(spec.producer))
615 .and_then(|s| s.procedures.iter().find(|p| p.path == spec.procedure)),
616 );
617 let target = TraceTarget {
618 origin: origin.clone(),
619 producer,
620 chain_chunk: spec
621 .procedure
622 .split('/')
623 .next()
624 .unwrap_or_default()
625 .to_string(),
626 registry_loaded: slices.is_some(),
627 };
628 let base = fleet.base();
629
630 let origin_scope = fleet.wire(all_under(scope));
634 let fleet_scope = fleet.wire(all_under(Scope::fleet()));
635 let session = fleet.session();
636 let origin_monitor = Monitor::start(
637 session,
638 MonitorSpec {
639 capacity: ORIGIN_CAPACITY,
640 ..MonitorSpec::default()
641 },
642 )
643 .await?;
644 let mut origin_events = origin_monitor.events();
645 origin_monitor.watch(&origin_scope).await?;
646 let fleet_monitor = Monitor::start(
647 session,
648 MonitorSpec {
649 capacity: FLEET_CAPACITY,
650 ..MonitorSpec::default()
651 },
652 )
653 .await?;
654 let mut fleet_events = fleet_monitor.events();
655 fleet_monitor.watch(&fleet_scope).await?;
656
657 let t0 = std::time::Instant::now();
659 let t0_unix_s = std::time::SystemTime::now()
660 .duration_since(std::time::UNIX_EPOCH)
661 .map(|d| d.as_secs_f64())
662 .unwrap_or(0.0);
663 let (key, timeout, answers) = call_answers(fleet, spec).await?;
664 let call_returned_ms = t0.elapsed().as_secs_f64() * 1_000.0;
665 let reply_hlc = answers.iter().find_map(|a| a.timestamp);
666 let call = project_call(key, timeout, &answers);
667
668 let mut attributed = Vec::new();
671 let mut same_origin = Vec::new();
672 let mut pending_attributed = 0u64;
673 let mut pending_same_origin = 0u64;
674 let mut dropped = 0u64;
675 let mut concurrent_samples = 0u64;
676 let mut concurrent_dropped = 0u64;
677 let mut concurrent_keys = std::collections::HashSet::new();
678 let mut concurrent_examples = Examples::new(crate::judge::common::EXPANSION_CAP);
679 let reply_ntp64 = reply_hlc.map(|t| t.get_time().as_u64());
680 let deadline = tokio::time::sleep(trace.window);
681 tokio::pin!(deadline);
682 let mut origin_open = true;
683 let mut fleet_open = true;
684 while origin_open || fleet_open {
685 tokio::select! {
686 () = &mut deadline => break,
687 item = origin_events.recv(), if origin_open => match item {
688 None => origin_open = false,
689 Some(StreamItem::Dropped(n)) => {
690 dropped += n;
691 pending_attributed += n;
692 pending_same_origin += n;
693 }
694 Some(StreamItem::Event(FleetEvent::Sample(view))) => {
695 let desc = describe_key(base, &view.key, slices);
696 let Some(relation) = target.relation_of(&desc) else {
697 continue;
701 };
702 let row = TimelineRow::from_view(&view, t0, base);
703 let (lane, pending) = match relation {
704 crate::report::TraceRelation::DeclaredChain =>
705 (&mut attributed, &mut pending_attributed),
706 _ => (&mut same_origin, &mut pending_same_origin),
707 };
708 let break_before = (*pending > 0).then_some(*pending);
709 *pending = 0;
710 lane.push(trace_row(&row, relation, reply_ntp64, break_before));
711 }
712 Some(StreamItem::Event(_)) => {}
713 },
714 item = fleet_events.recv(), if fleet_open => match item {
715 None => fleet_open = false,
716 Some(StreamItem::Dropped(n)) => concurrent_dropped += n,
717 Some(StreamItem::Event(FleetEvent::Sample(view))) => {
718 let desc = describe_key(base, &view.key, None);
719 if target.relation_of(&desc).is_some() {
720 continue;
723 }
724 concurrent_samples += 1;
725 if concurrent_keys.insert(view.key.clone()) {
726 concurrent_examples.push_with(|| view.key.clone());
727 }
728 }
729 Some(StreamItem::Event(_)) => {}
730 },
731 }
732 }
733 let keys_evicted = origin_monitor.core().keys_evicted();
734 origin_monitor.stop();
735 fleet_monitor.stop();
736
737 Ok(TraceReport {
738 call,
739 scopes: vec![origin_scope, fleet_scope],
740 excluded: TRACE_EXCLUDED,
741 window_s: trace.window.as_secs_f64(),
742 subscribed_before_call: true,
743 t0_unix_s,
744 call_returned_ms,
745 hlc_reference: if reply_hlc.is_some() {
746 HlcReference::Reply
747 } else {
748 HlcReference::None
749 },
750 reply_hlc: reply_hlc.map(|t| t.to_string()),
751 chain_rule: TRACE_CHAIN_RULE,
752 registry_loaded: slices.is_some(),
753 idiom,
754 attributed,
755 same_origin,
756 concurrent: ConcurrentLane {
757 samples: concurrent_samples,
758 keys: concurrent_keys.len() as u64,
759 examples: concurrent_examples.into_vec(),
760 dropped: concurrent_dropped,
761 },
762 dropped,
763 keys_evicted,
764 })
765}
766
767#[cfg(test)]
768mod tests {
769 use super::*;
770 use zenkey::slice::{ProcedureDecl, RegistrySlice, SubjectDecl};
771
772 fn slice_with_state_subject() -> SliceSet {
773 let mut health = SubjectDecl::new("health", zenkey::Class::State);
774 health.type_name = "Health".into();
775 health.ttl_s = Some(900);
776 let mut slice = RegistrySlice::new("1.0", "t", "sysinfo");
777 slice.subjects = vec![health];
778 SliceSet::from_slices(vec![slice])
779 }
780
781 #[test]
784 fn a_wildcard_retire_is_refused_unconditionally() {
785 for force in [false, true] {
786 let err = check_retire("", "v1/h-3fa9c2d41b7e/state/sysinfo/**", None, force)
787 .unwrap_err()
788 .to_string();
789 assert!(err.contains("blast radius"), "{err}");
790 }
791 }
792
793 #[test]
794 fn a_state_key_retires_without_a_registry() {
795 let got = check_retire("", "v1/h-3fa9c2d41b7e/state/sysinfo/health", None, false).unwrap();
798 assert_eq!(
799 got,
800 RetireClass::State {
801 registered: false,
802 ttl_s: None
803 }
804 );
805 let slices = slice_with_state_subject();
807 let got = check_retire(
808 "",
809 "v1/h-3fa9c2d41b7e/state/sysinfo/health",
810 Some(&slices),
811 false,
812 )
813 .unwrap();
814 assert_eq!(
815 got,
816 RetireClass::State {
817 registered: true,
818 ttl_s: Some(900)
819 }
820 );
821 }
822
823 #[test]
824 fn a_telemetry_retire_needs_i_know_and_cites_the_rfc() {
825 let key = "v1/h-3fa9c2d41b7e/telemetry/sysinfo/cpu/usage";
826 let err = check_retire("", key, None, false).unwrap_err().to_string();
827 assert!(err.contains("MUST NOT"), "{err}");
828 assert!(err.contains("v1.12"), "{err}");
829 assert!(err.contains("--i-know"), "{err}");
830 assert_eq!(
831 check_retire("", key, None, true).unwrap(),
832 RetireClass::NonState {
833 class: "telemetry".to_string()
834 }
835 );
836 }
837
838 #[test]
839 fn a_plane_retire_needs_i_know_too() {
840 let key = "v1/h-3fa9c2d41b7e/@rpc/sysinfo/introspect";
841 let err = check_retire("", key, None, false).unwrap_err().to_string();
842 assert!(err.contains("plane"), "{err}");
843 assert!(matches!(
844 check_retire("", key, None, true).unwrap(),
845 RetireClass::NonState { class } if class == "@rpc"
846 ));
847 }
848
849 #[test]
850 fn an_unclassified_retire_needs_i_know_and_names_o4() {
851 let err = check_retire("", "some/foreign/key", None, false)
853 .unwrap_err()
854 .to_string();
855 assert!(err.contains("O4"), "{err}");
856 assert!(matches!(
857 check_retire("", "some/foreign/key", None, true).unwrap(),
858 RetireClass::Unclassified { .. }
859 ));
860 let err = check_retire("acme", "other/v1/h-3fa9c2d41b7e/state/x/y", None, false)
862 .unwrap_err()
863 .to_string();
864 assert!(err.contains("cannot be classified"), "{err}");
865 }
866
867 fn slice_with_proc(kind: &str, fanout: Option<&str>) -> SliceSet {
868 let mut trigger = ProcedureDecl::new("capture/trigger");
869 trigger.kind = Some(Declared::parse(kind));
870 trigger.reply = Some("Ack".into());
871 trigger.fanout = fanout.map(Declared::parse);
872 trigger.idempotent = Some(false);
873 let mut slice = RegistrySlice::new("1.0", "t", "netring");
874 slice.procedures = vec![trigger];
875 SliceSet::from_slices(vec![slice])
876 }
877
878 #[test]
879 fn call_targets_parse_and_validate() {
880 assert_eq!(CallTarget::parse("*").unwrap(), CallTarget::Fleet);
881 assert!(matches!(
882 CallTarget::parse("@catalog").unwrap(),
883 CallTarget::Service(_)
884 ));
885 assert!(matches!(
886 CallTarget::parse("h-3fa9c2d41b7e").unwrap(),
887 CallTarget::Host(_)
888 ));
889 let err = CallTarget::parse("toolbx").unwrap_err().to_string();
891 assert!(err.contains("RFC 06 §6"), "{err}");
892 }
893
894 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
897 async fn fleet_calls_to_forbidden_fanout_are_refused() {
898 let session = crate::bus::session::open(&[], &[], false).await.unwrap();
899 let slices = slice_with_proc("write", Some("forbidden"));
900 let err = call(
901 &crate::Fleet::new(&session, ""),
902 CallSpec {
903 target: &CallTarget::Fleet,
904 producer: "netring",
905 procedure: "capture/trigger",
906 params: &[],
907 body: None,
908 attachment: None,
909 timeout: Duration::from_millis(100),
910 slices: Some(&slices),
911 },
912 )
913 .await
914 .unwrap_err()
915 .to_string();
916 assert!(err.contains("fanout"), "{err}");
917 assert!(err.contains("RFC 05 §2.1"), "{err}");
918
919 let err = call(
923 &crate::Fleet::new(&session, ""),
924 CallSpec {
925 target: &CallTarget::Fleet,
926 producer: "netring",
927 procedure: "capture/trigger",
928 params: &[],
929 body: None,
930 attachment: None,
931 timeout: Duration::from_millis(100),
932 slices: Some(&slice_with_proc("write", None)),
933 },
934 )
935 .await
936 .unwrap_err()
937 .to_string();
938 assert!(err.contains("defaults to forbidden"), "{err}");
939 assert!(err.contains("RFC 08 §2"), "{err}");
940 assert!(err.contains("RFC 05 §2.1"), "{err}");
941
942 let err = call(
947 &crate::Fleet::new(&session, ""),
948 CallSpec {
949 target: &CallTarget::Fleet,
950 producer: "netring",
951 procedure: "capture/trigger",
952 params: &[],
953 body: None,
954 attachment: None,
955 timeout: Duration::from_millis(100),
956 slices: Some(&slice_with_proc("write", Some("per-iface"))),
957 },
958 )
959 .await
960 .unwrap_err()
961 .to_string();
962 assert!(err.contains("per-iface"), "{err}");
963 assert!(err.contains("does not know"), "{err}");
964 assert!(
965 !err.contains("declares fanout = \"forbidden\""),
966 "the slice declared no such thing: {err}"
967 );
968 assert!(err.contains("RFC 05 §2.1"), "{err}");
969
970 for slices in [
974 slice_with_proc("write", Some("allowed")),
975 slice_with_proc("read", None),
976 ] {
977 let report = call(
978 &crate::Fleet::new(&session, ""),
979 CallSpec {
980 target: &CallTarget::Fleet,
981 producer: "netring",
982 procedure: "capture/trigger",
983 params: &[],
984 body: None,
985 attachment: None,
986 timeout: Duration::from_millis(100),
987 slices: Some(&slices),
988 },
989 )
990 .await
991 .unwrap();
992 assert_eq!(report.exit_code(), 2, "silence stays exit 2");
993 }
994 }
995}