1#![deny(missing_docs)]
11
12use std::collections::{BTreeMap, VecDeque};
13use std::error::Error;
14use std::fmt::{self, Display, Formatter};
15
16use wip_http::http::{Request, Response};
17use wip_http::{
18 ClientResponseError, DecodedResponse, Endpoint, Limits, decode_call_operation_response,
19 decode_fetch_interface_response, decode_observe_response, encode_call_operation_request,
20 encode_fetch_interface_request, encode_observe_request,
21};
22use wip_protocol::{
23 CallOperationRequest, CallOperationResponse, FetchInterfaceRequest, InterfaceDescriptor,
24 InterfaceReference, InterfaceTarget, Object, ObjectObservation as ProtocolObjectObservation,
25 ObserveRequest, OperationDeclaration, ProtocolError, ProtocolErrorCode, Target, Value,
26};
27
28#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
30pub struct RequestId(u64);
31
32impl RequestId {
33 #[must_use]
35 pub const fn get(self) -> u64 {
36 self.0
37 }
38}
39
40#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)]
45pub struct SecurityContext(String);
46
47impl SecurityContext {
48 #[must_use]
50 pub fn new(identity: impl Into<String>) -> Self {
51 Self(identity.into())
52 }
53
54 #[must_use]
56 pub fn as_str(&self) -> &str {
57 &self.0
58 }
59}
60
61#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)]
63pub struct SessionId {
64 endpoint: String,
65 security_context: SecurityContext,
66}
67
68impl SessionId {
69 #[must_use]
71 pub fn endpoint(&self) -> &str {
72 &self.endpoint
73 }
74
75 #[must_use]
77 pub fn security_context(&self) -> &SecurityContext {
78 &self.security_context
79 }
80}
81
82#[derive(Debug, Clone, Copy, PartialEq, Eq)]
84pub struct ClientLimits {
85 pub sessions: usize,
87 pub objects_per_session: usize,
89 pub interfaces_per_session: usize,
91 pub in_flight: usize,
93 pub history: usize,
95}
96
97impl ClientLimits {
98 pub fn new(
100 sessions: usize,
101 objects_per_session: usize,
102 interfaces_per_session: usize,
103 in_flight: usize,
104 history: usize,
105 ) -> Result<Self, ConfigurationError> {
106 if sessions == 0
107 || objects_per_session == 0
108 || interfaces_per_session == 0
109 || in_flight == 0
110 || history == 0
111 {
112 return Err(ConfigurationError);
113 }
114 Ok(Self {
115 sessions,
116 objects_per_session,
117 interfaces_per_session,
118 in_flight,
119 history,
120 })
121 }
122}
123
124#[derive(Debug, Clone, Copy, PartialEq, Eq)]
126pub struct ConfigurationError;
127
128impl Display for ConfigurationError {
129 fn fmt(&self, formatter: &mut Formatter<'_>) -> fmt::Result {
130 formatter.write_str("client limits must be nonzero")
131 }
132}
133
134impl Error for ConfigurationError {}
135
136#[derive(Debug, Clone, Copy, PartialEq, Eq)]
138pub enum ClientErrorKind {
139 InvalidResponse,
141 UnsupportedDescriptorFormat,
143 Transport,
145 ProtocolFailure,
147 InvalidRequest,
149 MissingObservation,
151 Capacity,
153 StaleResponse,
155}
156
157#[derive(Debug, Clone, PartialEq, Eq)]
159pub struct ClientError {
160 pub kind: ClientErrorKind,
162 pub detail: String,
164}
165
166impl ClientError {
167 fn new(kind: ClientErrorKind, detail: impl Into<String>) -> Self {
168 Self {
169 kind,
170 detail: detail.into(),
171 }
172 }
173}
174
175impl Display for ClientError {
176 fn fmt(&self, formatter: &mut Formatter<'_>) -> fmt::Result {
177 write!(
178 formatter,
179 "client failure ({:?}): {}",
180 self.kind, self.detail
181 )
182 }
183}
184
185impl Error for ClientError {}
186
187#[derive(Debug, Clone, PartialEq, Eq)]
189pub enum ObservationState {
190 Loading(RequestId),
192 Fresh,
194 Stale,
196 Error(ClientError),
198}
199
200#[derive(Debug, Clone, PartialEq, Eq)]
202pub struct ObjectObservation {
203 pub path: String,
205 pub object: Option<Object>,
207 pub validator: Option<Vec<u8>>,
209 pub state: ObservationState,
211}
212
213#[derive(Debug, Clone, PartialEq, Eq)]
215pub struct TreeObservation {
216 pub path: String,
218 pub children: Vec<String>,
220 pub revision: RequestId,
222 pub state: ObservationState,
224}
225
226#[derive(Debug, Clone, PartialEq, Eq)]
233pub struct InterfaceObservation {
234 pub reference: InterfaceReference,
236 pub descriptor: Option<InterfaceDescriptor>,
238 pub validator: Option<Vec<u8>>,
241 pub scope_ref: Option<String>,
243 pub state: ObservationState,
245}
246
247#[derive(Debug, Clone, PartialEq)]
249pub struct CallContext {
250 session: SessionId,
251 object_key: u64,
252 object: ObjectObservation,
253 interface: InterfaceObservation,
254 interface_revision: RequestId,
255 descriptor: InterfaceDescriptor,
256 operation: OperationDeclaration,
257 request: CallOperationRequest,
258}
259
260impl CallContext {
261 #[must_use]
263 pub fn session(&self) -> &SessionId {
264 &self.session
265 }
266
267 #[must_use]
269 pub fn object(&self) -> &ObjectObservation {
270 &self.object
271 }
272
273 #[must_use]
275 pub fn interface(&self) -> &InterfaceObservation {
276 &self.interface
277 }
278
279 #[must_use]
281 pub fn descriptor(&self) -> &InterfaceDescriptor {
282 &self.descriptor
283 }
284
285 #[must_use]
287 pub fn operation(&self) -> &OperationDeclaration {
288 &self.operation
289 }
290
291 #[must_use]
293 pub fn request(&self) -> &CallOperationRequest {
294 &self.request
295 }
296}
297
298#[derive(Debug, Clone, Copy, PartialEq, Eq)]
300pub enum OutcomeUnknownReason {
301 Protocol,
303 Timeout,
305 Disconnect,
307 ResponseDecode,
309}
310
311#[derive(Debug, Clone, PartialEq)]
313pub enum CallOutcome {
314 Success(CallOperationResponse),
316 ProtocolFailure(ProtocolError),
318 Unknown {
320 reason: OutcomeUnknownReason,
322 failure: Option<ClientError>,
324 },
325 NotDispatched(ClientError),
327}
328
329#[derive(Debug, Clone, PartialEq)]
331pub struct CallRecord {
332 pub request_id: RequestId,
334 pub context: CallContext,
336 pub outcome: CallOutcome,
338}
339
340#[derive(Debug, Clone, Copy, PartialEq, Eq)]
342pub enum DispatchedTransportFailure {
343 Timeout,
345 Disconnect,
347}
348
349#[derive(Debug, Clone, Copy, PartialEq, Eq)]
351pub enum DiagnosticKind {
352 HttpStatusMismatch,
354 StaleResponse,
356 ClientFailure,
358}
359
360#[derive(Debug, Clone, PartialEq, Eq)]
362pub struct Diagnostic {
363 pub request_id: RequestId,
365 pub kind: DiagnosticKind,
367 pub detail: String,
369}
370
371#[derive(Debug)]
373pub struct PreparedRequest {
374 pub id: RequestId,
376 pub request: Request<Vec<u8>>,
378}
379
380#[derive(Debug, Clone, PartialEq)]
382pub enum Completion {
383 Object(ObjectObservation),
385 Interface(InterfaceObservation),
387 Call(Box<CallRecord>),
389 ProtocolFailure(ProtocolError),
391 StaleResponseRejected(RequestId),
393}
394
395#[derive(Debug, Clone)]
396struct EntryRecord {
397 object_key: Option<u64>,
398 revision: RequestId,
399 state: ObservationState,
400 confirmed_missing: bool,
403}
404
405#[derive(Debug, Clone)]
406struct ObjectRecord {
407 object: Object,
408 revision: RequestId,
409 call_owners: usize,
410 touched: u64,
411}
412
413#[derive(Debug, Clone)]
414struct Session {
415 entries: BTreeMap<String, EntryRecord>,
416 objects: BTreeMap<u64, ObjectRecord>,
417 refs: BTreeMap<String, u64>,
418 trees: BTreeMap<String, TreeObservation>,
419 interfaces: BTreeMap<InterfaceReference, InterfaceRecord>,
420 calls: VecDeque<CallRecord>,
421 diagnostics: VecDeque<Diagnostic>,
422 next_object_key: u64,
423 clock: u64,
424}
425
426#[derive(Debug, Clone)]
427struct InterfaceRecord {
428 observation: InterfaceObservation,
429 revision: RequestId,
431 accepted_revision: RequestId,
434 scope_binding: Option<String>,
438 binding_revision: RequestId,
439 touched: u64,
440}
441
442impl InterfaceRecord {
443 fn retire(&mut self, revision: RequestId) {
444 self.revision = self.revision.max(revision);
445 self.observation.descriptor = None;
446 self.observation.validator = None;
447 self.observation.state = ObservationState::Stale;
448 }
449}
450
451impl Session {
452 fn new() -> Self {
453 Self {
454 entries: BTreeMap::new(),
455 objects: BTreeMap::new(),
456 refs: BTreeMap::new(),
457 trees: BTreeMap::new(),
458 interfaces: BTreeMap::new(),
459 calls: VecDeque::new(),
460 diagnostics: VecDeque::new(),
461 next_object_key: 1,
462 clock: 0,
463 }
464 }
465
466 fn tick(&mut self) -> u64 {
467 self.clock = self.clock.saturating_add(1);
468 self.clock
469 }
470
471 fn object_observation(&self, path: &str) -> Option<ObjectObservation> {
472 let entry = self.entries.get(path)?;
473 let object = entry
474 .object_key
475 .and_then(|key| self.objects.get(&key))
476 .map(|record| record.object.clone());
477 let validator = object.as_ref().and_then(|value| value.validator.clone());
478 Some(ObjectObservation {
479 path: path.to_owned(),
480 object,
481 validator,
482 state: entry.state.clone(),
483 })
484 }
485
486 fn mark_object_stale(&mut self, key: u64, revision: RequestId) {
487 if let Some(record) = self.objects.get_mut(&key) {
488 record.revision = record.revision.max(revision);
489 }
490 for entry in self.entries.values_mut() {
491 if entry.object_key == Some(key) {
492 entry.state = ObservationState::Stale;
493 }
494 }
495 }
496
497 fn pin_call_object(&mut self, key: u64) {
498 let record = self
499 .objects
500 .get_mut(&key)
501 .expect("call context object remains observed");
502 record.call_owners = record.call_owners.saturating_add(1);
503 }
504
505 fn release_call_object(&mut self, key: u64) {
506 let Some(record) = self.objects.get_mut(&key) else {
507 return;
508 };
509 record.call_owners = record.call_owners.saturating_sub(1);
510 let remove = record.call_owners == 0
511 && !self
512 .entries
513 .values()
514 .any(|entry| entry.object_key == Some(key));
515 if remove {
516 remove_object(self, key);
517 }
518 }
519
520 fn push_diagnostic(&mut self, value: Diagnostic, maximum: usize) {
521 push_bounded(&mut self.diagnostics, value, maximum);
522 }
523
524 fn push_call(&mut self, value: CallRecord, maximum: usize) {
525 push_bounded(&mut self.calls, value, maximum);
526 }
527}
528
529#[derive(Debug)]
530enum Pending {
531 Observe {
532 session: SessionId,
533 request: ObserveRequest,
534 },
535 Interface {
536 session: SessionId,
537 request: FetchInterfaceRequest,
538 },
539 Call {
540 context: Box<CallContext>,
541 dispatched: bool,
542 },
543}
544
545impl Pending {
546 fn session(&self) -> &SessionId {
547 match self {
548 Self::Observe { session, .. } | Self::Interface { session, .. } => session,
549 Self::Call { context, .. } => context.session(),
550 }
551 }
552}
553
554#[derive(Debug)]
556pub struct Client {
557 limits: ClientLimits,
558 wire_limits: Limits,
559 sessions: BTreeMap<SessionId, Session>,
560 pending: BTreeMap<RequestId, Pending>,
561 next_request_id: u64,
562}
563
564impl Client {
565 #[must_use]
567 pub fn new(limits: ClientLimits, wire_limits: Limits) -> Self {
568 Self {
569 limits,
570 wire_limits,
571 sessions: BTreeMap::new(),
572 pending: BTreeMap::new(),
573 next_request_id: 1,
574 }
575 }
576
577 pub fn open_session(
579 &mut self,
580 endpoint: &str,
581 security_context: SecurityContext,
582 ) -> Result<SessionId, ClientError> {
583 let endpoint = Endpoint::parse(endpoint).map_err(|error| {
584 ClientError::new(ClientErrorKind::InvalidRequest, error.to_string())
585 })?;
586 let id = SessionId {
587 endpoint: endpoint.to_string(),
588 security_context,
589 };
590 if !self.sessions.contains_key(&id) {
591 if self.sessions.len() >= self.limits.sessions {
592 return Err(ClientError::new(
593 ClientErrorKind::Capacity,
594 "session capacity reached",
595 ));
596 }
597 self.sessions.insert(id.clone(), Session::new());
598 }
599 Ok(id)
600 }
601
602 #[must_use]
604 pub fn object(&self, session: &SessionId, path: &str) -> Option<ObjectObservation> {
605 self.sessions.get(session)?.object_observation(path)
606 }
607
608 #[must_use]
610 pub fn known_space(&self, session: &SessionId) -> Option<Vec<ObjectObservation>> {
611 let state = self.sessions.get(session)?;
612 Some(
613 state
614 .entries
615 .keys()
616 .filter_map(|path| state.object_observation(path))
617 .filter(|observation| observation.object.is_some())
618 .collect(),
619 )
620 }
621
622 #[must_use]
624 pub fn tree(&self, session: &SessionId, path: &str) -> Option<&TreeObservation> {
625 self.sessions.get(session)?.trees.get(path)
626 }
627
628 #[must_use]
631 pub fn interface(
632 &self,
633 session: &SessionId,
634 reference: &InterfaceReference,
635 ) -> Option<&InterfaceObservation> {
636 self.sessions
637 .get(session)?
638 .interfaces
639 .get(reference)
640 .map(|record| &record.observation)
641 }
642
643 #[must_use]
645 pub fn call_history(&self, session: &SessionId) -> Option<&VecDeque<CallRecord>> {
646 Some(&self.sessions.get(session)?.calls)
647 }
648
649 #[must_use]
651 pub fn diagnostics(&self, session: &SessionId) -> Option<&VecDeque<Diagnostic>> {
652 Some(&self.sessions.get(session)?.diagnostics)
653 }
654
655 #[must_use]
657 pub fn in_flight_len(&self) -> usize {
658 self.pending.len()
659 }
660
661 #[must_use]
663 pub fn pending_call_context(&self, id: RequestId) -> Option<&CallContext> {
664 match self.pending.get(&id)? {
665 Pending::Call { context, .. } => Some(context),
666 _ => None,
667 }
668 }
669
670 pub fn ensure_observed(
678 &mut self,
679 session: &SessionId,
680 path: impl Into<String>,
681 depth: u32,
682 ) -> Result<Option<PreparedRequest>, ClientError> {
683 let path = path.into();
684 wip_protocol::validate_path(&path).map_err(|error| {
685 ClientError::new(ClientErrorKind::InvalidRequest, error.to_string())
686 })?;
687 let coverage = observation_coverage(self.session(session)?, &path, depth);
688 if coverage != ObservationCoverage::Missing || self.pending.len() >= self.limits.in_flight {
689 return Ok(None);
690 }
691 self.refresh_observed(session, path, depth).map(Some)
692 }
693
694 pub fn ensure_children_observed(
697 &mut self,
698 session: &SessionId,
699 path: impl Into<String>,
700 ) -> Result<Option<PreparedRequest>, ClientError> {
701 self.ensure_observed(session, path, 1)
702 }
703
704 pub fn refresh_object(
707 &mut self,
708 session: &SessionId,
709 path: impl Into<String>,
710 ) -> Result<PreparedRequest, ClientError> {
711 self.refresh_observed(session, path, 0)
712 }
713
714 pub fn refresh_children(
716 &mut self,
717 session: &SessionId,
718 path: impl Into<String>,
719 ) -> Result<PreparedRequest, ClientError> {
720 self.refresh_observed(session, path, 1)
721 }
722
723 pub fn refresh_observed(
728 &mut self,
729 session: &SessionId,
730 path: impl Into<String>,
731 depth: u32,
732 ) -> Result<PreparedRequest, ClientError> {
733 self.ensure_capacity()?;
734 let request = ObserveRequest {
735 path: path.into(),
736 depth,
737 };
738 let endpoint = session_endpoint(session)?;
739 let encoded =
740 encode_observe_request(&endpoint, &request, self.wire_limits).map_err(|error| {
741 ClientError::new(ClientErrorKind::InvalidRequest, error.to_string())
742 })?;
743 let id = self.allocate_request_id();
744 let maximum = self.limits.objects_per_session;
745 let state = self.session_mut(session)?;
746 ensure_entry_slot(state, &request.path, maximum)?;
747 if request.depth > 0 {
748 ensure_tree_slot(state, &request.path, maximum)?;
749 }
750 let entry = state
751 .entries
752 .entry(request.path.clone())
753 .or_insert(EntryRecord {
754 object_key: None,
755 revision: id,
756 state: ObservationState::Loading(id),
757 confirmed_missing: false,
758 });
759 entry.revision = id;
760 entry.state = ObservationState::Loading(id);
761 if request.depth > 0 {
762 let tree = state
763 .trees
764 .entry(request.path.clone())
765 .or_insert(TreeObservation {
766 path: request.path.clone(),
767 children: Vec::new(),
768 revision: id,
769 state: ObservationState::Loading(id),
770 });
771 tree.revision = id;
772 tree.state = ObservationState::Loading(id);
773 } else if let Some(tree) = state.trees.get_mut(&request.path)
774 && matches!(tree.state, ObservationState::Loading(_))
775 {
776 tree.state = ObservationState::Stale;
777 }
778 self.pending.insert(
779 id,
780 Pending::Observe {
781 session: session.clone(),
782 request,
783 },
784 );
785 Ok(PreparedRequest {
786 id,
787 request: encoded,
788 })
789 }
790
791 pub fn ensure_interface(
800 &mut self,
801 session: &SessionId,
802 reference: InterfaceReference,
803 ) -> Result<Option<PreparedRequest>, ClientError> {
804 let state = self.session(session)?;
805 if state.interfaces.contains_key(&reference)
806 || state.interfaces.len() >= self.limits.interfaces_per_session
807 || self.pending.len() >= self.limits.in_flight
808 {
809 return Ok(None);
810 }
811 self.prepare_interface(session, reference).map(Some)
812 }
813
814 pub fn prepare_interface(
819 &mut self,
820 session: &SessionId,
821 reference: InterfaceReference,
822 ) -> Result<PreparedRequest, ClientError> {
823 self.ensure_capacity()?;
824 let request = FetchInterfaceRequest {
825 interface: reference,
826 };
827 let endpoint = session_endpoint(session)?;
828 let encoded = encode_fetch_interface_request(&endpoint, &request, self.wire_limits)
829 .map_err(|error| {
830 ClientError::new(ClientErrorKind::InvalidRequest, error.to_string())
831 })?;
832 let id = self.allocate_request_id();
833 let maximum = self.limits.interfaces_per_session;
834 let state = self.session_mut(session)?;
835 ensure_interface_slot(state, &request.interface, maximum)?;
836 let (scope_binding, binding_revision) =
837 latest_scope_binding(state, &request.interface.scope);
838 let touched = state.tick();
839 let record = state
840 .interfaces
841 .entry(request.interface.clone())
842 .or_insert(InterfaceRecord {
843 observation: InterfaceObservation {
844 reference: request.interface.clone(),
845 scope_ref: None,
846 descriptor: None,
847 validator: None,
848 state: ObservationState::Loading(id),
849 },
850 revision: id,
851 accepted_revision: RequestId(0),
852 scope_binding,
853 binding_revision,
854 touched,
855 });
856 record.revision = id;
857 record.observation.state = ObservationState::Loading(id);
858 record.touched = touched;
859 self.pending.insert(
860 id,
861 Pending::Interface {
862 session: session.clone(),
863 request,
864 },
865 );
866 Ok(PreparedRequest {
867 id,
868 request: encoded,
869 })
870 }
871
872 pub fn prepare_call(
880 &mut self,
881 session: &SessionId,
882 path: &str,
883 interface: &InterfaceReference,
884 operation: impl Into<String>,
885 arguments: BTreeMap<String, Value>,
886 ) -> Result<PreparedRequest, ClientError> {
887 self.ensure_capacity()?;
888 let operation = operation.into();
889 let context = {
890 let state = self.session(session)?;
891 let object_observation = state.object_observation(path).ok_or_else(|| {
892 ClientError::new(
893 ClientErrorKind::MissingObservation,
894 "object is not observed",
895 )
896 })?;
897 if object_observation.state != ObservationState::Fresh {
898 return Err(ClientError::new(
899 ClientErrorKind::MissingObservation,
900 "object observation is not fresh",
901 ));
902 }
903 let object = object_observation
904 .object
905 .as_ref()
906 .expect("fresh object exists");
907 if !object.interfaces.iter().any(|member| member == interface) {
908 return Err(ClientError::new(
909 ClientErrorKind::InvalidRequest,
910 "interface is not a member of the observed object",
911 ));
912 }
913 let interface_observation = state
914 .interfaces
915 .get(interface)
916 .map(|record| record.observation.clone())
917 .ok_or_else(|| {
918 ClientError::new(
919 ClientErrorKind::MissingObservation,
920 "interface is not observed",
921 )
922 })?;
923 if interface_observation.state != ObservationState::Fresh {
924 return Err(ClientError::new(
925 ClientErrorKind::MissingObservation,
926 "interface observation is not fresh",
927 ));
928 }
929 let interface_revision = state.interfaces[interface].revision;
930 let descriptor = interface_observation
931 .descriptor
932 .clone()
933 .expect("fresh interface exists");
934 descriptor
935 .validate_arguments(&operation, &arguments)
936 .map_err(|error| {
937 ClientError::new(ClientErrorKind::InvalidRequest, error.to_string())
938 })?;
939 let operation_declaration = descriptor
940 .operations
941 .iter()
942 .find(|declaration| declaration.name == operation)
943 .expect("validated operation exists")
944 .clone();
945 let request = CallOperationRequest {
946 target: Target {
947 path: path.to_owned(),
948 validator: object.validator.clone(),
949 },
950 interface: InterfaceTarget {
951 reference: interface.clone(),
952 scope_ref: interface_observation.scope_ref.clone(),
953 validator: interface_observation.validator.clone(),
954 },
955 operation,
956 arguments,
957 };
958 let object_key = state.entries[path]
959 .object_key
960 .expect("fresh entry has object key");
961 CallContext {
962 session: session.clone(),
963 object_key,
964 object: object_observation,
965 interface: interface_observation,
966 interface_revision,
967 descriptor,
968 operation: operation_declaration,
969 request,
970 }
971 };
972 let endpoint = session_endpoint(session)?;
973 let encoded = encode_call_operation_request(
974 &endpoint,
975 &context.request,
976 &context.descriptor,
977 self.wire_limits,
978 )
979 .map_err(|error| ClientError::new(ClientErrorKind::InvalidRequest, error.to_string()))?;
980 let id = self.allocate_request_id();
981 self.session_mut(&context.session)?
982 .pin_call_object(context.object_key);
983 self.pending.insert(
984 id,
985 Pending::Call {
986 context: Box::new(context),
987 dispatched: false,
988 },
989 );
990 Ok(PreparedRequest {
991 id,
992 request: encoded,
993 })
994 }
995
996 pub fn mark_dispatched(&mut self, id: RequestId) -> Result<(), ClientError> {
998 match self.pending.get_mut(&id) {
999 Some(Pending::Call { dispatched, .. }) => {
1000 *dispatched = true;
1001 Ok(())
1002 }
1003 Some(_) => Err(ClientError::new(
1004 ClientErrorKind::InvalidRequest,
1005 "only operation calls have a dispatch boundary",
1006 )),
1007 None => Err(unknown_request()),
1008 }
1009 }
1010
1011 pub fn complete(
1013 &mut self,
1014 id: RequestId,
1015 response: Response<Vec<u8>>,
1016 ) -> Result<Completion, ClientError> {
1017 let pending = self.pending.remove(&id).ok_or_else(unknown_request)?;
1018 match pending {
1019 Pending::Observe { session, request } => {
1020 self.complete_observe(id, &session, request, response)
1021 }
1022 Pending::Interface { session, request } => {
1023 self.complete_interface(id, &session, request, response)
1024 }
1025 Pending::Call { context, .. } => {
1026 let session = context.session.clone();
1027 let object_key = context.object_key;
1028 let result = self.complete_call(id, *context, response);
1029 self.session_mut(&session)?.release_call_object(object_key);
1030 result
1031 }
1032 }
1033 }
1034
1035 pub fn fail_transport(
1037 &mut self,
1038 id: RequestId,
1039 detail: impl Into<String>,
1040 ) -> Result<Completion, ClientError> {
1041 let detail = detail.into();
1042 let pending = self.pending.remove(&id).ok_or_else(unknown_request)?;
1043 match pending {
1044 Pending::Call {
1045 context,
1046 dispatched: true,
1047 } => {
1048 let session = context.session.clone();
1049 let object_key = context.object_key;
1050 let result = self.finish_unknown_call(
1051 id,
1052 *context,
1053 OutcomeUnknownReason::Disconnect,
1054 ClientError::new(ClientErrorKind::Transport, detail),
1055 );
1056 self.session_mut(&session)?.release_call_object(object_key);
1057 result
1058 }
1059 Pending::Call {
1060 context,
1061 dispatched: false,
1062 } => {
1063 let session = context.session.clone();
1064 let object_key = context.object_key;
1065 let record = CallRecord {
1066 request_id: id,
1067 context: *context,
1068 outcome: CallOutcome::NotDispatched(ClientError::new(
1069 ClientErrorKind::Transport,
1070 detail,
1071 )),
1072 };
1073 let maximum = self.limits.history;
1074 let state = self.session_mut(&session)?;
1075 state.push_call(record.clone(), maximum);
1076 state.release_call_object(object_key);
1077 Ok(Completion::Call(Box::new(record)))
1078 }
1079 other => {
1080 let session = other.session().clone();
1081 let is_current = match &other {
1082 Pending::Observe { request, .. } => {
1083 self.entry_request_is_current(&session, &request.path, id)
1084 }
1085 Pending::Interface { request, .. } => {
1086 self.interface_request_is_current(&session, &request.interface, id)
1087 }
1088 Pending::Call { .. } => unreachable!("calls handled above"),
1089 };
1090 if !is_current {
1091 return self.reject_stale(&session, id);
1092 }
1093 let error = ClientError::new(ClientErrorKind::Transport, detail);
1094 self.mark_retrieval_error(id, other, error.clone())?;
1095 self.push_failure_diagnostic(&session, id, &error)?;
1096 Err(error)
1097 }
1098 }
1099 }
1100
1101 pub fn fail_dispatched_call(
1103 &mut self,
1104 id: RequestId,
1105 failure: DispatchedTransportFailure,
1106 detail: impl Into<String>,
1107 ) -> Result<Completion, ClientError> {
1108 match self.pending.get(&id) {
1109 Some(Pending::Call {
1110 dispatched: true, ..
1111 }) => {}
1112 Some(_) => {
1113 return Err(ClientError::new(
1114 ClientErrorKind::InvalidRequest,
1115 "request was not a dispatched operation call",
1116 ));
1117 }
1118 None => return Err(unknown_request()),
1119 }
1120 let Some(Pending::Call { context, .. }) = self.pending.remove(&id) else {
1121 unreachable!("validated dispatched call remains pending")
1122 };
1123 let reason = match failure {
1124 DispatchedTransportFailure::Timeout => OutcomeUnknownReason::Timeout,
1125 DispatchedTransportFailure::Disconnect => OutcomeUnknownReason::Disconnect,
1126 };
1127 let session = context.session.clone();
1128 let object_key = context.object_key;
1129 let result = self.finish_unknown_call(
1130 id,
1131 *context,
1132 reason,
1133 ClientError::new(ClientErrorKind::Transport, detail),
1134 );
1135 self.session_mut(&session)?.release_call_object(object_key);
1136 result
1137 }
1138
1139 fn complete_observe(
1140 &mut self,
1141 id: RequestId,
1142 session: &SessionId,
1143 request: ObserveRequest,
1144 response: Response<Vec<u8>>,
1145 ) -> Result<Completion, ClientError> {
1146 if !self.entry_request_is_current(session, &request.path, id) {
1147 return self.reject_stale(session, id);
1148 }
1149 match decode_observe_response(&request, &response, self.wire_limits) {
1150 Ok(DecodedResponse::Success(value)) => {
1151 let maximum = self.limits.objects_per_session;
1152 let mut candidate = self.session(session)?.clone();
1153 if let Err(error) =
1154 apply_observation(&mut candidate, &request.path, value, maximum, id)
1155 {
1156 set_observe_subject_state(
1157 self.session_mut(session)?,
1158 &request,
1159 id,
1160 ObservationState::Error(error.clone()),
1161 );
1162 self.push_failure_diagnostic(session, id, &error)?;
1163 return Err(error);
1164 }
1165 let observation = candidate
1166 .object_observation(&request.path)
1167 .expect("inserted observation exists");
1168 self.sessions.insert(session.clone(), candidate);
1169 Ok(Completion::Object(observation))
1170 }
1171 Ok(DecodedResponse::ProtocolFailure {
1172 error,
1173 status_mismatch,
1174 }) => {
1175 self.record_status_mismatch(session, id, status_mismatch)?;
1176 let observation_state = if error.code == ProtocolErrorCode::NotFound {
1177 ObservationState::Stale
1178 } else {
1179 ObservationState::Error(ClientError::new(
1180 ClientErrorKind::ProtocolFailure,
1181 error.to_string(),
1182 ))
1183 };
1184 if error.code == ProtocolErrorCode::NotFound {
1185 let state = self.session_mut(session)?;
1186 retire_scope_subtree(state, &request.path, id);
1187 mark_subtree_stale(state, &request.path, id);
1188 }
1189 set_observe_subject_state(
1190 self.session_mut(session)?,
1191 &request,
1192 id,
1193 observation_state,
1194 );
1195 Ok(Completion::ProtocolFailure(error))
1196 }
1197 Err(error) => {
1198 self.retrieval_decode_error(id, session, RetrievalSubject::Observe(request), error)
1199 }
1200 }
1201 }
1202
1203 fn complete_interface(
1204 &mut self,
1205 id: RequestId,
1206 session: &SessionId,
1207 request: FetchInterfaceRequest,
1208 response: Response<Vec<u8>>,
1209 ) -> Result<Completion, ClientError> {
1210 if !self.interface_request_is_current(session, &request.interface, id) {
1211 return self.reject_stale(session, id);
1212 }
1213 match decode_fetch_interface_response(&request, &response, self.wire_limits) {
1214 Ok(DecodedResponse::Success(value)) => {
1215 let state = self.session_mut(session)?;
1216 let (scope_binding, binding_revision) =
1219 latest_scope_binding(state, &request.interface.scope);
1220 if binding_revision > id
1221 && value.scope_ref.is_some()
1222 && value.scope_ref != scope_binding
1223 {
1224 state
1225 .interfaces
1226 .get_mut(&request.interface)
1227 .expect("interface exists")
1228 .retire(binding_revision);
1229 return self.reject_stale(session, id);
1230 }
1231 let existing_record = &state.interfaces[&request.interface];
1232 let existing = &existing_record.observation;
1233 let lost_wire_binding = existing.scope_ref.is_some() && value.scope_ref.is_none();
1234 let current_binding = if lost_wire_binding || value.scope_ref.is_some() {
1238 value.scope_ref.clone()
1239 } else {
1240 scope_binding
1241 };
1242 let same_lifetime = existing.scope_ref == value.scope_ref
1243 || (current_binding.is_some()
1244 && existing_record.scope_binding == current_binding);
1245 if let Some(descriptor) = &existing.descriptor
1246 && same_lifetime
1247 && ((existing.validator.is_none()
1248 && (value.validator.is_some() || descriptor != &value.descriptor))
1249 || (existing.validator.is_some()
1250 && existing.validator == value.validator
1251 && descriptor != &value.descriptor))
1252 {
1253 let error = ClientError::new(
1254 ClientErrorKind::InvalidResponse,
1255 "interface representation changed for an immutable observation identity",
1256 );
1257 state
1258 .interfaces
1259 .get_mut(&request.interface)
1260 .expect("interface exists")
1261 .observation
1262 .state = ObservationState::Error(error.clone());
1263 self.push_failure_diagnostic(session, id, &error)?;
1264 return Err(error);
1265 }
1266 update_scope_binding(
1267 state,
1268 &request.interface.scope,
1269 current_binding.clone(),
1270 id,
1271 false,
1272 Some(&request.interface),
1273 );
1274 let touched = state.tick();
1275 let observation = InterfaceObservation {
1276 reference: value.interface,
1277 scope_ref: value.scope_ref,
1278 descriptor: Some(value.descriptor),
1279 validator: value.validator,
1280 state: ObservationState::Fresh,
1281 };
1282 let record = state
1283 .interfaces
1284 .get_mut(&request.interface)
1285 .expect("interface exists");
1286 record.observation = observation.clone();
1287 record.accepted_revision = id;
1288 record.scope_binding = current_binding;
1289 record.binding_revision = id;
1290 record.touched = touched;
1291 Ok(Completion::Interface(observation))
1292 }
1293 Ok(DecodedResponse::ProtocolFailure {
1294 error,
1295 status_mismatch,
1296 }) => {
1297 self.record_status_mismatch(session, id, status_mismatch)?;
1298 self.session_mut(session)?
1299 .interfaces
1300 .get_mut(&request.interface)
1301 .expect("interface exists")
1302 .observation
1303 .state = ObservationState::Error(ClientError::new(
1304 ClientErrorKind::ProtocolFailure,
1305 error.to_string(),
1306 ));
1307 Ok(Completion::ProtocolFailure(error))
1308 }
1309 Err(error) => self.retrieval_decode_error(
1310 id,
1311 session,
1312 RetrievalSubject::Interface(request.interface),
1313 error,
1314 ),
1315 }
1316 }
1317
1318 fn complete_call(
1319 &mut self,
1320 id: RequestId,
1321 context: CallContext,
1322 response: Response<Vec<u8>>,
1323 ) -> Result<Completion, ClientError> {
1324 match decode_call_operation_response(
1325 &context.request,
1326 &context.descriptor,
1327 &response,
1328 self.wire_limits,
1329 ) {
1330 Ok(DecodedResponse::Success(value)) => {
1331 let session_id = context.session.clone();
1332 let maximum = self.limits.history;
1333 let state = self.session_mut(&session_id)?;
1334 if let Some(object) = state.objects.get_mut(&context.object_key)
1335 && object.revision <= id
1336 {
1337 if let Some(validator) = value.validator.clone() {
1338 object.object.validator = Some(validator);
1339 }
1340 object.revision = id;
1341 }
1342 let record = CallRecord {
1343 request_id: id,
1344 context,
1345 outcome: CallOutcome::Success(value),
1346 };
1347 state.push_call(record.clone(), maximum);
1348 Ok(Completion::Call(Box::new(record)))
1349 }
1350 Ok(DecodedResponse::ProtocolFailure {
1351 error,
1352 status_mismatch,
1353 }) => {
1354 let session_id = context.session.clone();
1355 self.record_status_mismatch(&session_id, id, status_mismatch)?;
1356 if error.code == ProtocolErrorCode::OperationOutcomeUnknown {
1357 let maximum = self.limits.history;
1358 let state = self.session_mut(&session_id)?;
1359 state.mark_object_stale(context.object_key, id);
1360 let record = CallRecord {
1361 request_id: id,
1362 context,
1363 outcome: CallOutcome::Unknown {
1364 reason: OutcomeUnknownReason::Protocol,
1365 failure: None,
1366 },
1367 };
1368 state.push_call(record.clone(), maximum);
1369 return Ok(Completion::Call(Box::new(record)));
1370 }
1371 let maximum = self.limits.history;
1372 let state = self.session_mut(&session_id)?;
1373 match error.code {
1374 ProtocolErrorCode::ValidatorRequired
1375 | ProtocolErrorCode::ValidatorMismatch
1376 | ProtocolErrorCode::NotFound => {
1377 state.mark_object_stale(context.object_key, id);
1378 }
1379 ProtocolErrorCode::InterfaceMismatch => {
1380 state.mark_object_stale(context.object_key, id);
1381 if let Some(interface) =
1382 state.interfaces.get_mut(&context.interface.reference)
1383 && interface.revision == context.interface_revision
1384 {
1385 interface.revision = interface.revision.max(id);
1389 interface.observation.state = ObservationState::Stale;
1390 }
1391 }
1392 ProtocolErrorCode::InterfaceValidatorRequired
1393 | ProtocolErrorCode::InterfaceValidatorMismatch => {
1394 if let Some(interface) = state
1395 .interfaces
1396 .get_mut(&context.request.interface.reference)
1397 && interface.revision == context.interface_revision
1398 {
1399 interface.observation.state = ObservationState::Stale;
1400 }
1401 }
1402 _ => {}
1403 }
1404 let record = CallRecord {
1405 request_id: id,
1406 context,
1407 outcome: CallOutcome::ProtocolFailure(error),
1408 };
1409 state.push_call(record.clone(), maximum);
1410 Ok(Completion::Call(Box::new(record)))
1411 }
1412 Err(error) => {
1413 let failure = classify_response_error(error);
1414 self.finish_unknown_call(id, context, OutcomeUnknownReason::ResponseDecode, failure)
1415 }
1416 }
1417 }
1418
1419 fn finish_unknown_call(
1420 &mut self,
1421 id: RequestId,
1422 context: CallContext,
1423 reason: OutcomeUnknownReason,
1424 failure: ClientError,
1425 ) -> Result<Completion, ClientError> {
1426 let session_id = context.session.clone();
1427 let maximum = self.limits.history;
1428 let state = self.session_mut(&session_id)?;
1429 state.mark_object_stale(context.object_key, id);
1430 state.push_diagnostic(
1431 Diagnostic {
1432 request_id: id,
1433 kind: DiagnosticKind::ClientFailure,
1434 detail: failure.to_string(),
1435 },
1436 maximum,
1437 );
1438 let record = CallRecord {
1439 request_id: id,
1440 context,
1441 outcome: CallOutcome::Unknown {
1442 reason,
1443 failure: Some(failure),
1444 },
1445 };
1446 state.push_call(record.clone(), maximum);
1447 Ok(Completion::Call(Box::new(record)))
1448 }
1449
1450 fn retrieval_decode_error(
1451 &mut self,
1452 id: RequestId,
1453 session: &SessionId,
1454 subject: RetrievalSubject,
1455 error: ClientResponseError,
1456 ) -> Result<Completion, ClientError> {
1457 let error = classify_response_error(error);
1458 match subject {
1459 RetrievalSubject::Observe(request) => {
1460 set_observe_subject_state(
1461 self.session_mut(session)?,
1462 &request,
1463 id,
1464 ObservationState::Error(error.clone()),
1465 );
1466 }
1467 RetrievalSubject::Interface(reference) => {
1468 self.session_mut(session)?
1469 .interfaces
1470 .get_mut(&reference)
1471 .expect("interface exists")
1472 .observation
1473 .state = ObservationState::Error(error.clone());
1474 }
1475 }
1476 self.push_failure_diagnostic(session, id, &error)?;
1477 Err(error)
1478 }
1479
1480 fn mark_retrieval_error(
1481 &mut self,
1482 id: RequestId,
1483 pending: Pending,
1484 error: ClientError,
1485 ) -> Result<(), ClientError> {
1486 match pending {
1487 Pending::Observe { session, request } => {
1488 set_observe_subject_state(
1489 self.session_mut(&session)?,
1490 &request,
1491 id,
1492 ObservationState::Error(error),
1493 );
1494 }
1495 Pending::Interface { session, request } => {
1496 self.session_mut(&session)?
1497 .interfaces
1498 .get_mut(&request.interface)
1499 .expect("interface exists")
1500 .observation
1501 .state = ObservationState::Error(error);
1502 }
1503 Pending::Call { .. } => unreachable!("handled above"),
1504 }
1505 Ok(())
1506 }
1507
1508 fn record_status_mismatch(
1509 &mut self,
1510 session: &SessionId,
1511 id: RequestId,
1512 mismatch: Option<wip_http::StatusMismatch>,
1513 ) -> Result<(), ClientError> {
1514 if let Some(mismatch) = mismatch {
1515 let maximum = self.limits.history;
1516 self.session_mut(session)?.push_diagnostic(
1517 Diagnostic {
1518 request_id: id,
1519 kind: DiagnosticKind::HttpStatusMismatch,
1520 detail: format!(
1521 "received HTTP {}; canonical status is {}",
1522 mismatch.actual, mismatch.expected
1523 ),
1524 },
1525 maximum,
1526 );
1527 }
1528 Ok(())
1529 }
1530
1531 fn push_failure_diagnostic(
1532 &mut self,
1533 session: &SessionId,
1534 id: RequestId,
1535 error: &ClientError,
1536 ) -> Result<(), ClientError> {
1537 let maximum = self.limits.history;
1538 self.session_mut(session)?.push_diagnostic(
1539 Diagnostic {
1540 request_id: id,
1541 kind: DiagnosticKind::ClientFailure,
1542 detail: error.to_string(),
1543 },
1544 maximum,
1545 );
1546 Ok(())
1547 }
1548
1549 fn reject_stale(
1550 &mut self,
1551 session: &SessionId,
1552 id: RequestId,
1553 ) -> Result<Completion, ClientError> {
1554 let maximum = self.limits.history;
1555 self.session_mut(session)?.push_diagnostic(
1556 Diagnostic {
1557 request_id: id,
1558 kind: DiagnosticKind::StaleResponse,
1559 detail: "response superseded by a newer request".into(),
1560 },
1561 maximum,
1562 );
1563 Ok(Completion::StaleResponseRejected(id))
1564 }
1565
1566 fn entry_request_is_current(&self, session: &SessionId, path: &str, id: RequestId) -> bool {
1567 matches!(
1568 self.sessions
1569 .get(session)
1570 .and_then(|state| state.entries.get(path))
1571 .map(|entry| &entry.state),
1572 Some(ObservationState::Loading(current)) if *current == id
1573 )
1574 }
1575
1576 fn interface_request_is_current(
1577 &self,
1578 session: &SessionId,
1579 reference: &InterfaceReference,
1580 id: RequestId,
1581 ) -> bool {
1582 matches!(
1583 self.sessions
1584 .get(session)
1585 .and_then(|state| state.interfaces.get(reference))
1586 .map(|interface| &interface.observation.state),
1587 Some(ObservationState::Loading(current)) if *current == id
1588 )
1589 }
1590
1591 fn ensure_capacity(&self) -> Result<(), ClientError> {
1592 if self.pending.len() >= self.limits.in_flight {
1593 Err(ClientError::new(
1594 ClientErrorKind::Capacity,
1595 "in-flight request capacity reached",
1596 ))
1597 } else {
1598 Ok(())
1599 }
1600 }
1601
1602 fn allocate_request_id(&mut self) -> RequestId {
1603 let id = RequestId(self.next_request_id);
1604 self.next_request_id = self.next_request_id.saturating_add(1);
1605 id
1606 }
1607
1608 fn session(&self, id: &SessionId) -> Result<&Session, ClientError> {
1609 self.sessions.get(id).ok_or_else(|| {
1610 ClientError::new(ClientErrorKind::InvalidRequest, "unknown client session")
1611 })
1612 }
1613
1614 fn session_mut(&mut self, id: &SessionId) -> Result<&mut Session, ClientError> {
1615 self.sessions.get_mut(id).ok_or_else(|| {
1616 ClientError::new(ClientErrorKind::InvalidRequest, "unknown client session")
1617 })
1618 }
1619}
1620
1621#[derive(Debug)]
1622enum RetrievalSubject {
1623 Observe(ObserveRequest),
1624 Interface(InterfaceReference),
1625}
1626
1627#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1628enum ObservationCoverage {
1629 Missing,
1630 Loading,
1631 Covered,
1632 Blocked,
1633}
1634
1635fn observation_coverage(state: &Session, path: &str, depth: u32) -> ObservationCoverage {
1636 let Some(entry) = state.entries.get(path) else {
1637 return ObservationCoverage::Missing;
1638 };
1639 match &entry.state {
1640 ObservationState::Loading(_) => return ObservationCoverage::Loading,
1641 ObservationState::Stale | ObservationState::Error(_) => {
1642 return ObservationCoverage::Blocked;
1643 }
1644 ObservationState::Fresh => {}
1645 }
1646 if entry
1647 .object_key
1648 .is_none_or(|key| !state.objects.contains_key(&key))
1649 {
1650 return ObservationCoverage::Missing;
1651 }
1652 if depth == 0 {
1653 return ObservationCoverage::Covered;
1654 }
1655
1656 let Some(tree) = state.trees.get(path) else {
1657 return ObservationCoverage::Missing;
1658 };
1659 match &tree.state {
1660 ObservationState::Loading(_) => return ObservationCoverage::Loading,
1661 ObservationState::Stale => return ObservationCoverage::Missing,
1662 ObservationState::Error(_) => return ObservationCoverage::Blocked,
1663 ObservationState::Fresh => {}
1664 }
1665
1666 let mut saw_loading = false;
1667 let mut saw_missing = false;
1668 for child in &tree.children {
1669 match observation_coverage(state, child, depth - 1) {
1670 ObservationCoverage::Blocked => return ObservationCoverage::Blocked,
1671 ObservationCoverage::Loading => saw_loading = true,
1672 ObservationCoverage::Missing => saw_missing = true,
1673 ObservationCoverage::Covered => {}
1674 }
1675 }
1676 if saw_loading {
1677 ObservationCoverage::Loading
1678 } else if saw_missing {
1679 ObservationCoverage::Missing
1680 } else {
1681 ObservationCoverage::Covered
1682 }
1683}
1684
1685fn session_endpoint(session: &SessionId) -> Result<Endpoint, ClientError> {
1686 Endpoint::parse(&session.endpoint)
1687 .map_err(|error| ClientError::new(ClientErrorKind::InvalidRequest, error.to_string()))
1688}
1689
1690fn unknown_request() -> ClientError {
1691 ClientError::new(ClientErrorKind::InvalidRequest, "unknown request identity")
1692}
1693
1694fn classify_response_error(error: ClientResponseError) -> ClientError {
1695 match error {
1696 ClientResponseError::InvalidResponse { .. } => {
1697 ClientError::new(ClientErrorKind::InvalidResponse, error.to_string())
1698 }
1699 ClientResponseError::UnsupportedDescriptorFormat { .. } => ClientError::new(
1700 ClientErrorKind::UnsupportedDescriptorFormat,
1701 error.to_string(),
1702 ),
1703 ClientResponseError::TransportBinding(_) => {
1704 ClientError::new(ClientErrorKind::Transport, error.to_string())
1705 }
1706 }
1707}
1708
1709fn set_observe_subject_state(
1710 state: &mut Session,
1711 request: &ObserveRequest,
1712 revision: RequestId,
1713 observation_state: ObservationState,
1714) {
1715 if request.depth > 0
1716 && let Some(tree) = state.trees.get_mut(&request.path)
1717 && tree.revision == revision
1718 {
1719 tree.state = observation_state.clone();
1720 }
1721 if let Some(entry) = state.entries.get_mut(&request.path)
1722 && entry.revision == revision
1723 {
1724 entry.state = observation_state;
1725 }
1726}
1727
1728fn ensure_interface_slot(
1729 state: &mut Session,
1730 reference: &InterfaceReference,
1731 maximum: usize,
1732) -> Result<(), ClientError> {
1733 if state.interfaces.contains_key(reference) || state.interfaces.len() < maximum {
1734 return Ok(());
1735 }
1736 let candidate = state
1737 .interfaces
1738 .iter()
1739 .filter(|(_, record)| !matches!(record.observation.state, ObservationState::Loading(_)))
1740 .min_by_key(|(_, record)| record.touched)
1741 .map(|(reference, _)| reference.clone())
1742 .ok_or_else(|| {
1743 ClientError::new(
1744 ClientErrorKind::Capacity,
1745 "all interface cache entries are loading",
1746 )
1747 })?;
1748 state.interfaces.remove(&candidate);
1749 Ok(())
1750}
1751
1752fn ensure_entry_slot(state: &mut Session, path: &str, maximum: usize) -> Result<(), ClientError> {
1753 if state.entries.contains_key(path) || state.entries.len() < maximum {
1754 return Ok(());
1755 }
1756 let candidate = state
1757 .entries
1758 .iter()
1759 .find(|(path, entry)| {
1760 !matches!(entry.state, ObservationState::Loading(_))
1761 && !entry.object_key.is_some_and(|key| {
1762 state
1763 .objects
1764 .get(&key)
1765 .is_some_and(|record| record.call_owners > 0)
1766 })
1767 && !state
1768 .trees
1769 .get(*path)
1770 .is_some_and(|tree| matches!(tree.state, ObservationState::Loading(_)))
1771 })
1772 .map(|(path, _)| path.clone())
1773 .ok_or_else(|| {
1774 ClientError::new(
1775 ClientErrorKind::Capacity,
1776 "all entry observations are owned by in-flight work",
1777 )
1778 })?;
1779 remove_entry(state, &candidate);
1780 Ok(())
1781}
1782
1783fn ensure_tree_slot(state: &mut Session, path: &str, maximum: usize) -> Result<(), ClientError> {
1784 if state.trees.contains_key(path) || state.trees.len() < maximum {
1785 return Ok(());
1786 }
1787 let candidate = state
1788 .trees
1789 .iter()
1790 .find(|(_, tree)| !matches!(tree.state, ObservationState::Loading(_)))
1791 .map(|(path, _)| path.clone())
1792 .ok_or_else(|| {
1793 ClientError::new(
1794 ClientErrorKind::Capacity,
1795 "all tree observations are loading",
1796 )
1797 })?;
1798 state.trees.remove(&candidate);
1799 Ok(())
1800}
1801
1802fn remove_entry(state: &mut Session, path: &str) {
1803 let object_key = state
1804 .entries
1805 .remove(path)
1806 .and_then(|entry| entry.object_key);
1807 state.trees.remove(path);
1808 for tree in state.trees.values_mut() {
1809 tree.children.retain(|child| child != path);
1810 }
1811 if let Some(key) = object_key
1812 && !state
1813 .entries
1814 .values()
1815 .any(|entry| entry.object_key == Some(key))
1816 {
1817 remove_object(state, key);
1818 }
1819}
1820
1821fn object_ownership_revision(state: &Session, key: u64) -> Option<RequestId> {
1822 state
1823 .entries
1824 .values()
1825 .filter(|entry| entry.object_key == Some(key))
1826 .map(|entry| entry.revision)
1827 .chain(state.objects.get(&key).map(|record| record.revision))
1828 .max()
1829}
1830
1831fn validate_object_consistency(state: &Session, object: &Object) -> Result<(), ClientError> {
1832 let Some(reference) = &object.r#ref else {
1833 return Ok(());
1834 };
1835 let Some(existing) = state
1836 .refs
1837 .get(reference)
1838 .and_then(|key| state.objects.get(key))
1839 else {
1840 return Ok(());
1841 };
1842 if existing.object.validator.is_some()
1843 && existing.object.validator == object.validator
1844 && existing.object != *object
1845 {
1846 return Err(ClientError::new(
1847 ClientErrorKind::InvalidResponse,
1848 "object representation changed for the same ref and validator",
1849 ));
1850 }
1851 Ok(())
1852}
1853
1854fn insert_object(
1855 state: &mut Session,
1856 path: &str,
1857 object: Object,
1858 maximum: usize,
1859 revision: RequestId,
1860) -> Result<u64, ClientError> {
1861 validate_object_consistency(state, &object)?;
1862 ensure_entry_slot(state, path, maximum)?;
1863 let existing_path_entry = state.entries.get(path).cloned();
1864 let existing_path_key = existing_path_entry
1865 .as_ref()
1866 .and_then(|entry| entry.object_key);
1867 let previous_scope_ref = existing_path_entry
1868 .as_ref()
1869 .filter(|entry| !entry.confirmed_missing)
1870 .and_then(|entry| entry.object_key)
1871 .and_then(|key| state.objects.get(&key))
1872 .and_then(|record| record.object.r#ref.clone());
1873 let replaceable_key = existing_path_key.filter(|key| {
1874 state
1875 .objects
1876 .get(key)
1877 .is_none_or(|record| record.call_owners == 0)
1878 && !state
1879 .entries
1880 .iter()
1881 .any(|(entry_path, entry)| entry_path != path && entry.object_key == Some(*key))
1882 });
1883 let key = if let Some(reference) = &object.r#ref {
1884 if let Some(key) = state.refs.get(reference).copied() {
1885 key
1886 } else {
1887 make_object_slot(state, maximum, replaceable_key)?
1888 }
1889 } else if let Some(key) = existing_path_key
1890 && state
1891 .objects
1892 .get(&key)
1893 .is_some_and(|record| record.object.r#ref.is_none())
1894 {
1895 key
1896 } else {
1897 make_object_slot(state, maximum, replaceable_key)?
1898 };
1899
1900 if let Some(old_key) = existing_path_key
1901 && old_key != key
1902 && !state
1903 .entries
1904 .iter()
1905 .any(|(entry_path, entry)| entry_path != path && entry.object_key == Some(old_key))
1906 {
1907 remove_object(state, old_key);
1908 }
1909 if let Some(reference) = &object.r#ref {
1910 state.refs.insert(reference.clone(), key);
1911 }
1912 let touched = state.tick();
1913 let ownership_revision = object_ownership_revision(state, key).unwrap_or(revision);
1914 let suppressed = ownership_revision > revision;
1915 let call_owners = state
1916 .objects
1917 .get(&key)
1918 .map_or(0, |record| record.call_owners);
1919 match state.objects.get_mut(&key) {
1920 Some(existing) if suppressed => {
1921 existing.revision = ownership_revision;
1922 existing.touched = touched;
1923 }
1924 _ => {
1925 state.objects.insert(
1926 key,
1927 ObjectRecord {
1928 object,
1929 revision,
1930 call_owners,
1931 touched,
1932 },
1933 );
1934 }
1935 }
1936 let entry = if suppressed {
1937 match existing_path_entry {
1938 Some(mut entry) if entry.object_key == Some(key) => {
1939 entry.revision = ownership_revision;
1940 if matches!(entry.state, ObservationState::Loading(current) if current == revision)
1941 {
1942 entry.state = ObservationState::Fresh;
1943 }
1944 entry
1945 }
1946 _ => EntryRecord {
1947 object_key: Some(key),
1948 revision: ownership_revision,
1949 state: ObservationState::Stale,
1950 confirmed_missing: false,
1951 },
1952 }
1953 } else {
1954 EntryRecord {
1955 object_key: Some(key),
1956 revision,
1957 state: ObservationState::Fresh,
1958 confirmed_missing: false,
1959 }
1960 };
1961 state.entries.insert(path.to_owned(), entry);
1962 if !suppressed {
1963 let scope_ref = state.objects[&key].object.r#ref.clone();
1964 let changed = previous_scope_ref.is_some() && previous_scope_ref != scope_ref;
1965 update_scope_binding(state, path, scope_ref, revision, changed, None);
1966 }
1967 Ok(key)
1968}
1969
1970fn make_object_slot(
1971 state: &mut Session,
1972 maximum: usize,
1973 replaceable: Option<u64>,
1974) -> Result<u64, ClientError> {
1975 if state.objects.len() >= maximum {
1976 let candidate = replaceable.or_else(|| {
1977 state
1978 .objects
1979 .iter()
1980 .filter(|(key, record)| {
1981 record.call_owners == 0
1982 && !state.entries.iter().any(|(path, entry)| {
1983 entry.object_key == Some(**key)
1984 && (matches!(entry.state, ObservationState::Loading(_))
1985 || state.trees.get(path).is_some_and(|tree| {
1986 matches!(tree.state, ObservationState::Loading(_))
1987 }))
1988 })
1989 })
1990 .min_by_key(|(_, record)| record.touched)
1991 .map(|(key, _)| *key)
1992 });
1993 let candidate = candidate.ok_or_else(|| {
1994 ClientError::new(
1995 ClientErrorKind::Capacity,
1996 "all object cache entries are owned by in-flight work",
1997 )
1998 })?;
1999 remove_object(state, candidate);
2000 }
2001 let key = state.next_object_key;
2002 state.next_object_key = state.next_object_key.saturating_add(1);
2003 Ok(key)
2004}
2005
2006fn remove_object(state: &mut Session, key: u64) {
2007 if state
2008 .objects
2009 .get(&key)
2010 .is_some_and(|record| record.call_owners > 0)
2011 {
2012 return;
2013 }
2014 if let Some(record) = state.objects.remove(&key)
2015 && let Some(reference) = record.object.r#ref
2016 && state.refs.get(&reference) == Some(&key)
2017 {
2018 state.refs.remove(&reference);
2019 }
2020 let removed_paths: Vec<String> = state
2021 .entries
2022 .iter()
2023 .filter(|(_, entry)| entry.object_key == Some(key))
2024 .map(|(path, _)| path.clone())
2025 .collect();
2026 state
2027 .entries
2028 .retain(|_, entry| entry.object_key != Some(key));
2029 let removable_paths: Vec<&String> = removed_paths
2030 .iter()
2031 .filter(|path| {
2032 !state
2033 .trees
2034 .get(*path)
2035 .is_some_and(|tree| matches!(tree.state, ObservationState::Loading(_)))
2036 })
2037 .collect();
2038 for path in &removable_paths {
2039 state.trees.remove(*path);
2040 }
2041 for tree in state.trees.values_mut() {
2042 tree.children.retain(|path| {
2043 !removable_paths
2044 .iter()
2045 .any(|removed| removed.as_str() == path)
2046 });
2047 }
2048}
2049
2050fn reserve_observation_metadata(
2051 state: &mut Session,
2052 paths: &[String],
2053 observed_children_paths: &[String],
2054 maximum: usize,
2055) -> Result<(), ClientError> {
2056 while state.entries.len()
2057 + paths
2058 .iter()
2059 .filter(|path| !state.entries.contains_key(*path))
2060 .count()
2061 > maximum
2062 {
2063 let candidate = state
2064 .entries
2065 .iter()
2066 .find(|(path, entry)| {
2067 !paths.contains(path)
2068 && !matches!(entry.state, ObservationState::Loading(_))
2069 && !entry.object_key.is_some_and(|key| {
2070 state
2071 .objects
2072 .get(&key)
2073 .is_some_and(|record| record.call_owners > 0)
2074 })
2075 && !state
2076 .trees
2077 .get(*path)
2078 .is_some_and(|tree| matches!(tree.state, ObservationState::Loading(_)))
2079 })
2080 .map(|(path, _)| path.clone())
2081 .ok_or_else(|| {
2082 ClientError::new(
2083 ClientErrorKind::Capacity,
2084 "tree response cannot evict entries owned by in-flight work",
2085 )
2086 })?;
2087 remove_entry(state, &candidate);
2088 }
2089 while state.trees.len()
2090 + observed_children_paths
2091 .iter()
2092 .filter(|path| !state.trees.contains_key(*path))
2093 .count()
2094 > maximum
2095 {
2096 let candidate = state
2097 .trees
2098 .iter()
2099 .find(|(path, tree)| {
2100 !observed_children_paths.contains(path)
2101 && !matches!(tree.state, ObservationState::Loading(_))
2102 })
2103 .map(|(path, _)| path.clone())
2104 .ok_or_else(|| {
2105 ClientError::new(
2106 ClientErrorKind::Capacity,
2107 "tree response cannot evict loading tree observations",
2108 )
2109 })?;
2110 state.trees.remove(&candidate);
2111 }
2112 Ok(())
2113}
2114
2115fn latest_scope_binding(state: &Session, scope: &str) -> (Option<String>, RequestId) {
2119 let object_binding = state
2120 .entries
2121 .get(scope)
2122 .filter(|entry| !entry.confirmed_missing)
2123 .and_then(|entry| entry.object_key)
2124 .and_then(|key| state.objects.get(&key))
2125 .map(|record| (record.object.r#ref.clone(), record.revision));
2126 state
2127 .interfaces
2128 .iter()
2129 .filter(|(reference, _)| reference.scope == scope)
2130 .map(|(_, record)| (record.scope_binding.clone(), record.binding_revision))
2131 .chain(object_binding)
2132 .max_by_key(|(_, revision)| *revision)
2133 .unwrap_or((None, RequestId(0)))
2134}
2135
2136fn update_scope_binding(
2137 state: &mut Session,
2138 scope: &str,
2139 scope_ref: Option<String>,
2140 revision: RequestId,
2141 confirmed_change: bool,
2142 except: Option<&InterfaceReference>,
2143) {
2144 for (reference, interface) in &mut state.interfaces {
2145 if reference.scope != scope
2146 || except == Some(reference)
2147 || interface.accepted_revision > revision
2148 || interface.binding_revision > revision
2149 {
2150 continue;
2151 }
2152 if confirmed_change
2153 || (interface.scope_binding.is_some() && interface.scope_binding != scope_ref)
2154 {
2155 interface.retire(revision);
2156 }
2157 interface.scope_binding = scope_ref.clone();
2160 interface.binding_revision = revision;
2161 }
2162}
2163
2164fn is_subtree_path(root: &str, path: &str) -> bool {
2165 root == "/"
2166 || path == root
2167 || path
2168 .strip_prefix(root)
2169 .is_some_and(|suffix| suffix.starts_with('/'))
2170}
2171
2172fn retire_scope_subtree(state: &mut Session, root: &str, revision: RequestId) {
2175 for (path, entry) in &mut state.entries {
2176 if is_subtree_path(root, path) && entry.revision <= revision {
2177 entry.confirmed_missing = true;
2178 }
2179 }
2180 state.interfaces.retain(|reference, record| {
2181 !is_subtree_path(root, &reference.scope) || record.accepted_revision > revision
2182 });
2183}
2184
2185fn mark_subtree_stale(state: &mut Session, root: &str, request_id: RequestId) {
2186 for (path, entry) in &mut state.entries {
2187 if is_subtree_path(root, path) && entry.revision <= request_id {
2188 entry.revision = request_id;
2189 entry.state = ObservationState::Stale;
2190 }
2191 }
2192 for (path, tree) in &mut state.trees {
2193 if is_subtree_path(root, path) && tree.revision <= request_id {
2194 tree.revision = request_id;
2195 tree.state = ObservationState::Stale;
2196 }
2197 }
2198}
2199
2200fn apply_observation(
2201 state: &mut Session,
2202 root_path: &str,
2203 response: ProtocolObjectObservation,
2204 maximum: usize,
2205 request_id: RequestId,
2206) -> Result<(), ClientError> {
2207 let mut nodes = Vec::new();
2208 flatten_observation(root_path, &response, &mut nodes);
2209 if nodes.len() > maximum {
2210 return Err(ClientError::new(
2211 ClientErrorKind::Capacity,
2212 "observation response exceeds the local observation bound",
2213 ));
2214 }
2215 let paths: Vec<String> = nodes.iter().map(|(path, _, _)| path.clone()).collect();
2216 let observed_children_paths: Vec<String> = nodes
2217 .iter()
2218 .filter_map(|(path, _, children)| children.as_ref().map(|_| path.clone()))
2219 .collect();
2220 reserve_observation_metadata(state, &paths, &observed_children_paths, maximum)?;
2221
2222 for (path, _, children) in &nodes {
2223 if children.is_none()
2224 && let Some(tree) = state.trees.get_mut(path)
2225 && tree.revision < request_id
2226 && matches!(tree.state, ObservationState::Loading(_))
2227 {
2228 tree.state = ObservationState::Stale;
2229 }
2230 }
2231
2232 let mut omitted = Vec::new();
2233 for (path, _, children) in &nodes {
2234 let Some(children) = children else {
2235 continue;
2236 };
2237 if let Some(previous) = state
2238 .trees
2239 .get(path)
2240 .filter(|previous| previous.revision <= request_id)
2241 {
2242 omitted.extend(
2243 previous
2244 .children
2245 .iter()
2246 .filter(|child| !children.contains(child))
2247 .cloned(),
2248 );
2249 }
2250 }
2251 for path in omitted {
2252 mark_subtree_stale(state, &path, request_id);
2253 }
2254
2255 for (path, object, _) in &nodes {
2256 validate_object_consistency(state, object)?;
2257 let newer_owner = state
2258 .entries
2259 .get(path)
2260 .is_some_and(|entry| entry.revision > request_id);
2261 if !newer_owner {
2262 insert_object(state, path, object.clone(), maximum, request_id)?;
2263 }
2264 }
2265 for (path, _, children) in nodes {
2266 let Some(children) = children else {
2267 continue;
2268 };
2269 let newer_owner = state
2270 .trees
2271 .get(&path)
2272 .is_some_and(|tree| tree.revision > request_id);
2273 if newer_owner {
2274 continue;
2275 }
2276 state.trees.insert(
2277 path.clone(),
2278 TreeObservation {
2279 path,
2280 children,
2281 revision: request_id,
2282 state: ObservationState::Fresh,
2283 },
2284 );
2285 }
2286 Ok(())
2287}
2288
2289fn flatten_observation(
2290 path: &str,
2291 observation: &ProtocolObjectObservation,
2292 output: &mut Vec<(String, Object, Option<Vec<String>>)>,
2293) {
2294 let children = observation.children.as_ref().map(|children| {
2295 children
2296 .iter()
2297 .map(|child| child_path(path, &child.object.name))
2298 .collect::<Vec<_>>()
2299 });
2300 output.push((
2301 path.to_owned(),
2302 observation.object.clone(),
2303 children.clone(),
2304 ));
2305 if let (Some(observations), Some(paths)) = (&observation.children, children) {
2306 for (child, child_path) in observations.iter().zip(paths) {
2307 flatten_observation(&child_path, child, output);
2308 }
2309 }
2310}
2311
2312fn child_path(parent: &str, name: &str) -> String {
2313 if parent == "/" {
2314 format!("/{name}")
2315 } else {
2316 format!("{parent}/{name}")
2317 }
2318}
2319
2320fn push_bounded<T>(history: &mut VecDeque<T>, value: T, maximum: usize) {
2321 if history.len() == maximum {
2322 history.pop_front();
2323 }
2324 history.push_back(value);
2325}