1mod provider_runtime;
4
5pub use provider_runtime::{
6 ProviderCancellation, ProviderCancellationToken, ProviderCompletionHandle,
7 ProviderEffectPublishReport, ProviderInvocation, ProviderObservationOutcome,
8 ProviderObservationStatus, ProviderRuntimeAuthorityHandle, ProviderRuntimeError,
9 ProviderRuntimeHandle, ProviderRuntimeState, ProviderWork, MAX_PROVIDER_EFFECT_CAPACITY,
10 MAX_PROVIDER_RUNTIMES,
11};
12
13use gate4agent_kernel::{
14 BackendIngress, BackendIngressOutcome, BackendSnapshot, KernelStep,
15 ToolAuthorityCommandOutcome, ToolRequestOutcome,
16};
17use gate4agent_tool_protocol::{
18 CapabilityCompletionBatch, CapabilityCompletionEnvelope, CapabilityOwner,
19 CapabilityProviderDescriptor, CapabilityRequestInput, CapabilityRequestSnapshot,
20 ConsumerBoundCapabilityRequest, ConsumerId, PolicyGrant, ProviderBindingId,
21 ProviderBoundCapabilityRequest, ToolActorId, ToolAuditEvent, ToolAuthorityCommand,
22 ToolAuthorityEnvelope, ToolAuthorityOutcome, ToolInstanceState, ToolProviderId,
23};
24use gate4agent_types::{
25 AgentInstanceId, CommandEnvelope, ControlEvent, ControlSnapshot, SessionGeneration,
26};
27use std::collections::{BTreeMap, BTreeSet};
28use std::sync::mpsc::{sync_channel, Receiver, SyncSender, TryRecvError, TrySendError};
29use std::sync::{Arc, Mutex, RwLock, Weak};
30use thiserror::Error;
31
32pub const MAX_CONTROL_SUBSCRIBERS: usize = 128;
33pub const MAX_TOOL_CLIENTS: usize = 128;
34pub const MAX_TOOL_SUBSCRIBERS_PER_CLIENT: usize = 32;
35pub const MAX_AUTHORITY_SUBSCRIBERS: usize = 32;
36pub const MAX_INGRESS_CAPACITY: usize = 4_096;
37pub const MAX_SUBSCRIPTION_CAPACITY: usize = 1_024;
38
39pub trait AgentPort: Send + Sync {
40 fn dispatch(&self, command: CommandEnvelope) -> Result<(), PortDispatchError>;
41 fn snapshot(&self) -> Arc<ControlSnapshot>;
42 fn subscribe(&self, capacity: usize) -> EventSubscription;
43}
44
45#[derive(Clone)]
46pub struct Gate4AgentHandle {
47 command_tx: GateIngressSender,
48 edge: Arc<EdgeState>,
49}
50
51pub struct KernelPort {
52 command_rx: Receiver<CommandEnvelope>,
53 edge: Arc<EdgeState>,
54}
55
56pub struct ControlPlaneKernelPort {
57 ingress_rx: Receiver<BackendIngress>,
58 provider_runtime: provider_runtime::ProviderRuntimePort,
59 edge: Arc<EdgeState>,
60 authority: Arc<AuthorityState>,
61}
62
63#[derive(Clone)]
64pub struct ToolAuthorityHandle {
65 edge: Arc<EdgeState>,
66 authority: Arc<AuthorityState>,
67}
68
69#[derive(Clone)]
70pub struct ToolClientHandle {
71 state: Arc<ClientState>,
72}
73
74pub struct EventSubscription {
75 receiver: Receiver<ControlEvent>,
76}
77
78pub struct ToolRequestOutcomeSubscription {
79 receiver: Receiver<ToolRequestOutcome>,
80}
81
82pub struct ToolCompletionSubscription {
83 receiver: Receiver<ToolCompletionDelivery>,
84}
85
86pub struct ToolAuthorityOutcomeSubscription {
87 receiver: Receiver<ToolAuthorityCommandOutcome>,
88}
89
90#[derive(Clone, Debug, Eq, PartialEq)]
91pub enum ToolCompletionDelivery {
92 SourceGap(ToolCompletionSourceGap),
93 Completion(CapabilityCompletionEnvelope),
94}
95
96#[derive(Clone, Copy, Debug, Eq, PartialEq)]
97pub struct ToolCompletionSourceGap {
98 pub dropped_since_last_drain: u64,
99 pub total_dropped: u64,
100 pub next_sequence: u64,
101 pub sequence_exhausted: bool,
102}
103
104#[derive(Clone, Debug, Eq, PartialEq)]
105pub struct ToolClientSnapshot {
106 pub backend_revision: u64,
107 pub logical_tick: u64,
108 pub tool_revision: u64,
109 pub current_tick: u64,
110 pub generations: Vec<(AgentInstanceId, SessionGeneration)>,
111 pub instance_states: Vec<(AgentInstanceId, ToolInstanceState)>,
112 pub providers: Vec<CapabilityProviderDescriptor>,
113 pub available_providers: Vec<ToolProviderId>,
114 pub grants: Vec<PolicyGrant>,
115 pub requests: Vec<CapabilityRequestSnapshot>,
116 pub audit_events: Vec<ToolAuditEvent>,
117 pub dropped_audit_events: u64,
118 pub next_completion_sequence: u64,
119 pub dropped_completions: u64,
120 pub completion_sequence_exhausted: bool,
121}
122
123#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
124pub struct PublishReport {
125 pub delivered: usize,
126 pub disconnected_slow: usize,
127 pub disconnected_closed: usize,
128}
129
130#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
131pub struct ControlPlanePublishReport {
132 pub control_events: PublishReport,
133 pub request_outcomes: PublishReport,
134 pub authority_outcomes: PublishReport,
135 pub completions: PublishReport,
136 pub provider_effects: ProviderEffectPublishReport,
137 pub closed_clients: usize,
138}
139
140#[derive(Clone, Copy, Debug, Eq, Error, PartialEq)]
141pub enum PortDispatchError {
142 #[error("gate4agent command ingress is full")]
143 Full,
144 #[error("gate4agent kernel ingress is disconnected")]
145 Disconnected,
146}
147
148#[derive(Clone, Debug, Eq, Error, PartialEq)]
149pub enum ToolClientDispatchError {
150 #[error("tool client is not active")]
151 Inactive,
152 #[error("tool provider '{provider_id}' has no active local runtime binding")]
153 ProviderUnavailable { provider_id: ToolProviderId },
154 #[error("gate4agent backend ingress is full")]
155 Full,
156 #[error("gate4agent backend ingress is disconnected")]
157 Disconnected,
158}
159
160#[derive(Clone, Copy, Debug, Eq, Error, PartialEq)]
161pub enum ToolAuthorityError {
162 #[error("tool client registry is full")]
163 ClientCapacityExceeded,
164 #[error("tool authority counter '{counter}' is exhausted")]
165 CounterExhausted { counter: &'static str },
166 #[error("tool client belongs to another authority")]
167 ForeignClient,
168 #[error("tool client is already closed")]
169 ClientClosed,
170 #[error("CloseClient must be dispatched through close_client")]
171 CloseRequiresMethod,
172 #[error("gate4agent backend ingress is full")]
173 Full,
174 #[error("gate4agent backend ingress is disconnected")]
175 Disconnected,
176}
177
178#[derive(Clone)]
179enum GateIngressSender {
180 Legacy(SyncSender<CommandEnvelope>),
181 Backend(SyncSender<BackendIngress>),
182}
183
184struct EdgeState {
185 publication: RwLock<EdgePublication>,
186 control_subscribers: Mutex<Vec<SyncSender<ControlEvent>>>,
187 completion_health: Mutex<CompletionPublishHealth>,
188}
189
190#[derive(Clone)]
191struct EdgePublication {
192 snapshot: Arc<BackendSnapshot>,
193 available_provider_bindings: BTreeMap<ToolProviderId, ProviderBindingId>,
194}
195
196#[derive(Clone, Copy, Default)]
197struct CompletionPublishHealth {
198 sequence_exhaustion_announced: bool,
199}
200
201struct AuthorityState {
202 ingress_tx: SyncSender<BackendIngress>,
203 inner: Mutex<AuthorityInner>,
204}
205
206struct AuthorityInner {
207 next_client_id: u64,
208 next_authority_sequence: u64,
209 clients: BTreeMap<ConsumerId, Arc<ClientState>>,
210 closing_by_sequence: BTreeMap<u64, ConsumerId>,
211 subscribers: Vec<SyncSender<ToolAuthorityCommandOutcome>>,
212}
213
214struct ClientState {
215 consumer_id: ConsumerId,
216 actor_id: ToolActorId,
217 authority: Weak<AuthorityState>,
218 edge: Arc<EdgeState>,
219 inner: Mutex<ClientInner>,
220}
221
222struct ClientInner {
223 lifecycle: ClientLifecycle,
224 request_subscribers: Vec<SyncSender<ToolRequestOutcome>>,
225 completion_subscribers: Vec<SyncSender<ToolCompletionDelivery>>,
226}
227
228#[derive(Clone)]
229enum ClientLifecycle {
230 Active,
231 ClosingSent { sequence: u64 },
232 ClosingRetry,
233 CloseSucceeded,
234 Closed,
235}
236
237pub fn bounded_port(command_capacity: usize) -> (Gate4AgentHandle, KernelPort) {
238 let (command_tx, command_rx) = sync_channel(bounded_ingress_capacity(command_capacity));
239 let edge = Arc::new(EdgeState::new());
240 (
241 Gate4AgentHandle {
242 command_tx: GateIngressSender::Legacy(command_tx),
243 edge: Arc::clone(&edge),
244 },
245 KernelPort { command_rx, edge },
246 )
247}
248
249pub fn bounded_control_plane(
250 ingress_capacity: usize,
251) -> (
252 Gate4AgentHandle,
253 ToolAuthorityHandle,
254 ControlPlaneKernelPort,
255) {
256 let (ingress_tx, ingress_rx) = sync_channel(bounded_ingress_capacity(ingress_capacity));
257 let (_provider_authority, provider_runtime) =
258 provider_runtime::provider_runtime(ingress_tx.clone());
259 let edge = Arc::new(EdgeState::new());
260 let authority = Arc::new(AuthorityState {
261 ingress_tx: ingress_tx.clone(),
262 inner: Mutex::new(AuthorityInner {
263 next_client_id: 1,
264 next_authority_sequence: 1,
265 clients: BTreeMap::new(),
266 closing_by_sequence: BTreeMap::new(),
267 subscribers: Vec::new(),
268 }),
269 });
270 (
271 Gate4AgentHandle {
272 command_tx: GateIngressSender::Backend(ingress_tx),
273 edge: Arc::clone(&edge),
274 },
275 ToolAuthorityHandle {
276 edge: Arc::clone(&edge),
277 authority: Arc::clone(&authority),
278 },
279 ControlPlaneKernelPort {
280 ingress_rx,
281 provider_runtime,
282 edge,
283 authority,
284 },
285 )
286}
287
288impl AgentPort for Gate4AgentHandle {
289 fn dispatch(&self, command: CommandEnvelope) -> Result<(), PortDispatchError> {
290 match &self.command_tx {
291 GateIngressSender::Legacy(sender) => sender.try_send(command).map_err(map_port_error),
292 GateIngressSender::Backend(sender) => sender
293 .try_send(BackendIngress::Control(command))
294 .map_err(map_port_error),
295 }
296 }
297
298 fn snapshot(&self) -> Arc<ControlSnapshot> {
299 Arc::clone(&self.edge.load_snapshot().control)
300 }
301
302 fn subscribe(&self, capacity: usize) -> EventSubscription {
303 let receiver = subscribe_bounded(
304 &self.edge.control_subscribers,
305 capacity,
306 MAX_CONTROL_SUBSCRIBERS,
307 );
308 EventSubscription { receiver }
309 }
310}
311
312impl Gate4AgentHandle {
313 pub fn dispatch(&self, command: CommandEnvelope) -> Result<(), PortDispatchError> {
314 AgentPort::dispatch(self, command)
315 }
316
317 pub fn snapshot(&self) -> Arc<ControlSnapshot> {
318 AgentPort::snapshot(self)
319 }
320
321 pub fn subscribe(&self, capacity: usize) -> EventSubscription {
322 AgentPort::subscribe(self, capacity)
323 }
324}
325
326impl KernelPort {
327 pub fn drain_commands(&self, limit: usize) -> Vec<CommandEnvelope> {
328 drain_receiver(&self.command_rx, limit)
329 }
330
331 pub fn publish_snapshot(&self, snapshot: ControlSnapshot) {
332 let current = self.edge.load_snapshot();
333 self.edge.publish_snapshot(BackendSnapshot {
334 revision: snapshot.revision,
335 logical_tick: current.logical_tick,
336 control: Arc::new(snapshot),
337 tools: Arc::clone(¤t.tools),
338 provider_runtime: current.provider_runtime.clone(),
339 });
340 }
341
342 pub fn publish_events(&self, events: impl IntoIterator<Item = ControlEvent>) -> PublishReport {
345 publish_many(&self.edge.control_subscribers, events)
346 }
347}
348
349impl ControlPlaneKernelPort {
350 pub fn drain_ingress(&self, limit: usize) -> Vec<BackendIngress> {
351 self.provider_runtime.flush_closing_bindings();
352 drain_receiver(&self.ingress_rx, limit)
353 }
354
355 pub fn provider_authority(&self) -> ProviderRuntimeAuthorityHandle {
356 self.provider_runtime.authority_handle()
357 }
358
359 pub fn publish_events(&self, events: impl IntoIterator<Item = ControlEvent>) -> PublishReport {
360 publish_many(&self.edge.control_subscribers, events)
361 }
362
363 pub fn publish_step(&self, step: &KernelStep) -> ControlPlanePublishReport {
372 let mut report = ControlPlanePublishReport::default();
373 for outcome in &step.ingress_outcomes {
374 if let BackendIngressOutcome::ToolProvider(outcome) = outcome {
375 self.provider_runtime.publish_outcome(outcome);
376 }
377 }
378 for effect in step.tool_effects.iter().cloned() {
379 report.provider_effects += self.provider_runtime.publish_effect(effect);
380 }
381 self.provider_runtime.finish_step();
382 self.provider_runtime
383 .reconcile_snapshot(&step.backend_snapshot.provider_runtime);
384 self.edge.publish_control_plane(
385 step.backend_snapshot.clone(),
386 self.provider_runtime.active_bindings(),
387 );
388
389 let mut completed_closes = Vec::new();
390 for outcome in &step.ingress_outcomes {
391 match outcome {
392 BackendIngressOutcome::Control(_) => {}
393 BackendIngressOutcome::ToolRequest(outcome) => {
394 report.request_outcomes += self.publish_request_outcome(outcome.clone());
395 }
396 BackendIngressOutcome::ToolAuthority(outcome) => {
397 let (published, completed_close) =
398 self.publish_authority_outcome(outcome.clone());
399 report.authority_outcomes += published;
400 if let Some(client) = completed_close {
401 completed_closes.push(client);
402 }
403 }
404 BackendIngressOutcome::ToolProvider(_) => {}
405 }
406 }
407
408 report.control_events = self.publish_events(step.events.iter().cloned());
409 report.completions = self.publish_completions(&step.tool_completions);
410 for client in completed_closes {
411 if self.finalize_client_close(&client) {
412 report.closed_clients += 1;
413 }
414 }
415 report
416 }
417
418 fn publish_request_outcome(&self, outcome: ToolRequestOutcome) -> PublishReport {
419 let Some(client) = self.client_for(
420 &outcome.request_key.consumer_id,
421 &outcome.request_key.actor_id,
422 ) else {
423 return PublishReport::default();
424 };
425 let mut inner = client
426 .inner
427 .lock()
428 .unwrap_or_else(|poisoned| poisoned.into_inner());
429 publish_one_locked(&mut inner.request_subscribers, outcome)
430 }
431
432 fn publish_authority_outcome(
433 &self,
434 outcome: ToolAuthorityCommandOutcome,
435 ) -> (PublishReport, Option<Arc<ClientState>>) {
436 let mut authority = self
437 .authority
438 .inner
439 .lock()
440 .unwrap_or_else(|poisoned| poisoned.into_inner());
441 let report = publish_one_locked(&mut authority.subscribers, outcome.clone());
442 let Some(consumer_id) = authority.closing_by_sequence.remove(&outcome.sequence) else {
443 return (report, None);
444 };
445 let client = authority.clients.get(&consumer_id).cloned();
446 if let Some(client) = &client {
447 let mut inner = client
448 .inner
449 .lock()
450 .unwrap_or_else(|poisoned| poisoned.into_inner());
451 if matches!(
452 outcome.result,
453 Ok(ToolAuthorityOutcome::ClientClosed { .. })
454 ) {
455 inner.lifecycle = ClientLifecycle::CloseSucceeded;
456 } else {
457 inner.lifecycle = ClientLifecycle::ClosingRetry;
458 return (report, None);
459 }
460 }
461 (report, client)
462 }
463
464 fn publish_completions(&self, batch: &CapabilityCompletionBatch) -> PublishReport {
465 let mut report = PublishReport::default();
466 for completion in &batch.completions {
467 let Some(client) = self.client_for(
468 &completion.request_key.consumer_id,
469 &completion.request_key.actor_id,
470 ) else {
471 continue;
472 };
473 let mut inner = client
474 .inner
475 .lock()
476 .unwrap_or_else(|poisoned| poisoned.into_inner());
477 report += publish_one_locked(
478 &mut inner.completion_subscribers,
479 ToolCompletionDelivery::Completion(completion.clone()),
480 );
481 }
482
483 let announce_gap = {
484 let mut health = self
485 .edge
486 .completion_health
487 .lock()
488 .unwrap_or_else(|poisoned| poisoned.into_inner());
489 let new_exhaustion = batch.sequence_exhausted && !health.sequence_exhaustion_announced;
490 health.sequence_exhaustion_announced = batch.sequence_exhausted;
491 batch.dropped_since_last_drain > 0 || new_exhaustion
492 };
493 if announce_gap {
494 let gap = ToolCompletionDelivery::SourceGap(ToolCompletionSourceGap {
495 dropped_since_last_drain: batch.dropped_since_last_drain,
496 total_dropped: batch.total_dropped,
497 next_sequence: batch.next_sequence,
498 sequence_exhausted: batch.sequence_exhausted,
499 });
500 let clients = self.all_clients();
501 for client in clients {
502 let mut inner = client
503 .inner
504 .lock()
505 .unwrap_or_else(|poisoned| poisoned.into_inner());
506 report += publish_one_locked(&mut inner.completion_subscribers, gap.clone());
507 }
508 }
509 report
510 }
511
512 fn client_for(
513 &self,
514 consumer_id: &ConsumerId,
515 actor_id: &ToolActorId,
516 ) -> Option<Arc<ClientState>> {
517 self.authority
518 .inner
519 .lock()
520 .unwrap_or_else(|poisoned| poisoned.into_inner())
521 .clients
522 .get(consumer_id)
523 .filter(|client| &client.actor_id == actor_id)
524 .cloned()
525 }
526
527 fn all_clients(&self) -> Vec<Arc<ClientState>> {
528 self.authority
529 .inner
530 .lock()
531 .unwrap_or_else(|poisoned| poisoned.into_inner())
532 .clients
533 .values()
534 .cloned()
535 .collect()
536 }
537
538 fn finalize_client_close(&self, client: &Arc<ClientState>) -> bool {
539 let mut authority = self
540 .authority
541 .inner
542 .lock()
543 .unwrap_or_else(|poisoned| poisoned.into_inner());
544 let is_registered = authority
545 .clients
546 .get(&client.consumer_id)
547 .is_some_and(|registered| Arc::ptr_eq(registered, client));
548 if !is_registered {
549 return false;
550 }
551 let mut inner = client
552 .inner
553 .lock()
554 .unwrap_or_else(|poisoned| poisoned.into_inner());
555 if !matches!(inner.lifecycle, ClientLifecycle::CloseSucceeded) {
556 return false;
557 }
558 inner.lifecycle = ClientLifecycle::Closed;
559 inner.request_subscribers.clear();
560 inner.completion_subscribers.clear();
561 authority.clients.remove(&client.consumer_id);
562 true
563 }
564}
565
566impl ToolAuthorityHandle {
567 pub fn bind_client(
568 &self,
569 actor_id: ToolActorId,
570 ) -> Result<ToolClientHandle, ToolAuthorityError> {
571 let mut authority = self
572 .authority
573 .inner
574 .lock()
575 .unwrap_or_else(|poisoned| poisoned.into_inner());
576 if authority.clients.len() >= MAX_TOOL_CLIENTS {
577 return Err(ToolAuthorityError::ClientCapacityExceeded);
578 }
579 let current = authority.next_client_id;
580 let next = current
581 .checked_add(1)
582 .ok_or(ToolAuthorityError::CounterExhausted {
583 counter: "consumer-id",
584 })?;
585 let consumer_id =
586 ConsumerId::new(format!("gate4agent-client-{current}")).map_err(|_| {
587 ToolAuthorityError::CounterExhausted {
588 counter: "consumer-id",
589 }
590 })?;
591 let state = Arc::new(ClientState {
592 consumer_id: consumer_id.clone(),
593 actor_id,
594 authority: Arc::downgrade(&self.authority),
595 edge: Arc::clone(&self.edge),
596 inner: Mutex::new(ClientInner {
597 lifecycle: ClientLifecycle::Active,
598 request_subscribers: Vec::new(),
599 completion_subscribers: Vec::new(),
600 }),
601 });
602 authority.next_client_id = next;
603 authority.clients.insert(consumer_id, Arc::clone(&state));
604 Ok(ToolClientHandle { state })
605 }
606
607 pub fn dispatch(&self, command: ToolAuthorityCommand) -> Result<u64, ToolAuthorityError> {
608 if matches!(command, ToolAuthorityCommand::CloseClient { .. }) {
609 return Err(ToolAuthorityError::CloseRequiresMethod);
610 }
611 let mut authority = self
612 .authority
613 .inner
614 .lock()
615 .unwrap_or_else(|poisoned| poisoned.into_inner());
616 let sequence = authority.next_authority_sequence;
617 let next = sequence
618 .checked_add(1)
619 .ok_or(ToolAuthorityError::CounterExhausted {
620 counter: "authority-sequence",
621 })?;
622 let envelope = ToolAuthorityEnvelope {
623 sequence,
624 command,
625 };
626 match self
627 .authority
628 .ingress_tx
629 .try_send(BackendIngress::ToolAuthority(envelope))
630 {
631 Ok(()) => {
632 authority.next_authority_sequence = next;
633 Ok(sequence)
634 }
635 Err(TrySendError::Full(_)) => Err(ToolAuthorityError::Full),
636 Err(TrySendError::Disconnected(_)) => Err(ToolAuthorityError::Disconnected),
637 }
638 }
639
640 pub fn close_client(&self, client: &ToolClientHandle) -> Result<u64, ToolAuthorityError> {
646 let Some(client_authority) = client.state.authority.upgrade() else {
647 return Err(ToolAuthorityError::ForeignClient);
648 };
649 if !Arc::ptr_eq(&client_authority, &self.authority) {
650 return Err(ToolAuthorityError::ForeignClient);
651 }
652
653 let mut authority = self
654 .authority
655 .inner
656 .lock()
657 .unwrap_or_else(|poisoned| poisoned.into_inner());
658 if !authority
659 .clients
660 .get(&client.state.consumer_id)
661 .is_some_and(|registered| Arc::ptr_eq(registered, &client.state))
662 {
663 return Err(ToolAuthorityError::ClientClosed);
664 }
665 let mut inner = client
666 .state
667 .inner
668 .lock()
669 .unwrap_or_else(|poisoned| poisoned.into_inner());
670 let sequence = match &inner.lifecycle {
671 ClientLifecycle::ClosingSent { sequence } => return Ok(*sequence),
672 ClientLifecycle::CloseSucceeded | ClientLifecycle::Closed => {
673 return Err(ToolAuthorityError::ClientClosed)
674 }
675 ClientLifecycle::Active | ClientLifecycle::ClosingRetry => {
676 let sequence = authority.next_authority_sequence;
677 sequence
678 .checked_add(1)
679 .ok_or(ToolAuthorityError::CounterExhausted {
680 counter: "authority-sequence",
681 })?;
682 inner.lifecycle = ClientLifecycle::ClosingRetry;
683 sequence
684 }
685 };
686 let next_sequence =
687 sequence
688 .checked_add(1)
689 .ok_or(ToolAuthorityError::CounterExhausted {
690 counter: "authority-sequence",
691 })?;
692 let envelope = ToolAuthorityEnvelope {
693 sequence,
694 command: ToolAuthorityCommand::CloseClient {
695 consumer_id: client.state.consumer_id.clone(),
696 actor_id: client.state.actor_id.clone(),
697 },
698 };
699 match self
700 .authority
701 .ingress_tx
702 .try_send(BackendIngress::ToolAuthority(envelope))
703 {
704 Ok(()) => {
705 authority.next_authority_sequence = next_sequence;
706 authority
707 .closing_by_sequence
708 .insert(sequence, client.state.consumer_id.clone());
709 inner.lifecycle = ClientLifecycle::ClosingSent { sequence };
710 Ok(sequence)
711 }
712 Err(TrySendError::Full(_)) => {
713 inner.lifecycle = ClientLifecycle::ClosingRetry;
714 Err(ToolAuthorityError::Full)
715 }
716 Err(TrySendError::Disconnected(_)) => Err(ToolAuthorityError::Disconnected),
717 }
718 }
719
720 pub fn snapshot(&self) -> Arc<BackendSnapshot> {
721 self.edge.load_snapshot()
722 }
723
724 pub fn subscribe_outcomes(&self, capacity: usize) -> ToolAuthorityOutcomeSubscription {
725 let (sender, receiver) = sync_channel(bounded_subscription_capacity(capacity));
726 let mut authority = self
727 .authority
728 .inner
729 .lock()
730 .unwrap_or_else(|poisoned| poisoned.into_inner());
731 if authority.subscribers.len() < MAX_AUTHORITY_SUBSCRIBERS {
732 authority.subscribers.push(sender);
733 }
734 ToolAuthorityOutcomeSubscription { receiver }
735 }
736}
737
738impl ToolClientHandle {
739 pub fn consumer_id(&self) -> &ConsumerId {
740 &self.state.consumer_id
741 }
742
743 pub fn actor_id(&self) -> &ToolActorId {
744 &self.state.actor_id
745 }
746
747 pub fn dispatch(
748 &self,
749 request: CapabilityRequestInput,
750 ) -> Result<gate4agent_tool_protocol::CapabilityRequestKey, ToolClientDispatchError> {
751 let mut inner = self
752 .state
753 .inner
754 .lock()
755 .unwrap_or_else(|poisoned| poisoned.into_inner());
756 if !matches!(inner.lifecycle, ClientLifecycle::Active) {
757 return Err(ToolClientDispatchError::Inactive);
758 }
759 let Some(authority) = self.state.authority.upgrade() else {
760 inner.lifecycle = ClientLifecycle::Closed;
761 return Err(ToolClientDispatchError::Disconnected);
762 };
763 let publication = self.state.edge.load_publication();
764 let provider_id = &request.provider_id;
765 let provider_is_registered = publication
766 .snapshot
767 .tools
768 .providers
769 .iter()
770 .any(|provider| &provider.id == provider_id);
771 let provider_binding_id = publication
772 .available_provider_bindings
773 .get(provider_id)
774 .copied();
775 if provider_is_registered && provider_binding_id.is_none() {
776 return Err(ToolClientDispatchError::ProviderUnavailable {
777 provider_id: provider_id.clone(),
778 });
779 }
780 let envelope = ConsumerBoundCapabilityRequest::new(
781 self.state.consumer_id.clone(),
782 self.state.actor_id.clone(),
783 request,
784 );
785 let request_key = envelope.key();
786 authority
787 .ingress_tx
788 .try_send(BackendIngress::ToolRequest(
789 ProviderBoundCapabilityRequest::new(provider_binding_id, envelope),
790 ))
791 .map_err(|error| match error {
792 TrySendError::Full(_) => ToolClientDispatchError::Full,
793 TrySendError::Disconnected(_) => ToolClientDispatchError::Disconnected,
794 })?;
795 Ok(request_key)
796 }
797
798 pub fn snapshot(&self) -> Arc<ToolClientSnapshot> {
799 let publication = self.state.edge.load_publication();
800 let snapshot = publication.snapshot;
801 let available_provider_bindings = publication.available_provider_bindings;
802 let tools = &snapshot.tools;
803 let consumer_id = &self.state.consumer_id;
804 let actor_id = &self.state.actor_id;
805 let grants = tools
806 .grants
807 .iter()
808 .filter(|grant| {
809 &grant.key.consumer_id == consumer_id && &grant.key.actor_id == actor_id
810 })
811 .cloned()
812 .collect::<Vec<_>>();
813 let requests = tools
814 .requests
815 .iter()
816 .filter(|request| {
817 &request.key.consumer_id == consumer_id && &request.key.actor_id == actor_id
818 })
819 .cloned()
820 .collect::<Vec<_>>();
821 let visible_instances = grants
822 .iter()
823 .map(|grant| grant.key.instance_id)
824 .chain(requests.iter().map(|request| request.instance_id))
825 .collect::<BTreeSet<_>>();
826 Arc::new(ToolClientSnapshot {
827 backend_revision: snapshot.revision,
828 logical_tick: snapshot.logical_tick,
829 tool_revision: tools.revision,
830 current_tick: tools.current_tick,
831 generations: tools
832 .generations
833 .iter()
834 .filter(|(instance_id, _)| visible_instances.contains(instance_id))
835 .copied()
836 .collect(),
837 instance_states: tools
838 .instance_states
839 .iter()
840 .filter(|(instance_id, _)| visible_instances.contains(instance_id))
841 .copied()
842 .collect(),
843 providers: tools
844 .providers
845 .iter()
846 .filter(|provider| match &provider.owner {
847 CapabilityOwner::Gate => true,
848 CapabilityOwner::Consumer(owner) => owner == consumer_id,
849 })
850 .cloned()
851 .collect(),
852 available_providers: available_provider_bindings
853 .iter()
854 .filter(|(provider_id, _)| {
855 tools.providers.iter().any(|provider| {
856 &provider.id == *provider_id
857 && match &provider.owner {
858 CapabilityOwner::Gate => true,
859 CapabilityOwner::Consumer(owner) => owner == consumer_id,
860 }
861 })
862 })
863 .map(|(provider_id, _)| provider_id.clone())
864 .collect(),
865 grants,
866 requests,
867 audit_events: tools
868 .audit_events
869 .iter()
870 .filter(|event| {
871 event.subject.as_ref().is_some_and(|subject| {
872 &subject.request_key.consumer_id == consumer_id
873 && &subject.request_key.actor_id == actor_id
874 })
875 })
876 .cloned()
877 .collect(),
878 dropped_audit_events: tools.dropped_audit_events,
879 next_completion_sequence: tools.next_completion_sequence,
880 dropped_completions: tools.dropped_completions,
881 completion_sequence_exhausted: tools.completion_sequence_exhausted,
882 })
883 }
884
885 pub fn subscribe_request_outcomes(&self, capacity: usize) -> ToolRequestOutcomeSubscription {
886 let (sender, receiver) = sync_channel(bounded_subscription_capacity(capacity));
887 let mut inner = self
888 .state
889 .inner
890 .lock()
891 .unwrap_or_else(|poisoned| poisoned.into_inner());
892 if !matches!(inner.lifecycle, ClientLifecycle::Closed)
893 && inner.request_subscribers.len() < MAX_TOOL_SUBSCRIBERS_PER_CLIENT
894 {
895 inner.request_subscribers.push(sender);
896 }
897 ToolRequestOutcomeSubscription { receiver }
898 }
899
900 pub fn subscribe_completions(&self, capacity: usize) -> ToolCompletionSubscription {
901 let (sender, receiver) = sync_channel(bounded_subscription_capacity(capacity));
902 let mut inner = self
903 .state
904 .inner
905 .lock()
906 .unwrap_or_else(|poisoned| poisoned.into_inner());
907 if !matches!(inner.lifecycle, ClientLifecycle::Closed)
908 && inner.completion_subscribers.len() < MAX_TOOL_SUBSCRIBERS_PER_CLIENT
909 {
910 inner.completion_subscribers.push(sender);
911 }
912 ToolCompletionSubscription { receiver }
913 }
914}
915
916impl EventSubscription {
917 pub fn try_recv(&self) -> Result<ControlEvent, TryRecvError> {
918 self.receiver.try_recv()
919 }
920}
921
922impl ToolRequestOutcomeSubscription {
923 pub fn try_recv(&self) -> Result<ToolRequestOutcome, TryRecvError> {
924 self.receiver.try_recv()
925 }
926}
927
928impl ToolCompletionSubscription {
929 pub fn try_recv(&self) -> Result<ToolCompletionDelivery, TryRecvError> {
930 self.receiver.try_recv()
931 }
932}
933
934impl ToolAuthorityOutcomeSubscription {
935 pub fn try_recv(&self) -> Result<ToolAuthorityCommandOutcome, TryRecvError> {
936 self.receiver.try_recv()
937 }
938}
939
940impl EdgeState {
941 fn new() -> Self {
942 Self {
943 publication: RwLock::new(EdgePublication {
944 snapshot: Arc::new(BackendSnapshot::default()),
945 available_provider_bindings: BTreeMap::new(),
946 }),
947 control_subscribers: Mutex::new(Vec::new()),
948 completion_health: Mutex::new(CompletionPublishHealth::default()),
949 }
950 }
951
952 fn load_snapshot(&self) -> Arc<BackendSnapshot> {
953 Arc::clone(
954 &self
955 .publication
956 .read()
957 .unwrap_or_else(|poisoned| poisoned.into_inner())
958 .snapshot,
959 )
960 }
961
962 fn load_publication(&self) -> EdgePublication {
963 self.publication
964 .read()
965 .unwrap_or_else(|poisoned| poisoned.into_inner())
966 .clone()
967 }
968
969 fn publish_snapshot(&self, snapshot: BackendSnapshot) {
970 self.publication
971 .write()
972 .unwrap_or_else(|poisoned| poisoned.into_inner())
973 .snapshot = Arc::new(snapshot);
974 }
975
976 fn publish_control_plane(
977 &self,
978 snapshot: BackendSnapshot,
979 available_provider_bindings: BTreeMap<ToolProviderId, ProviderBindingId>,
980 ) {
981 *self
982 .publication
983 .write()
984 .unwrap_or_else(|poisoned| poisoned.into_inner()) = EdgePublication {
985 snapshot: Arc::new(snapshot),
986 available_provider_bindings,
987 };
988 }
989}
990
991impl std::ops::AddAssign for PublishReport {
992 fn add_assign(&mut self, other: Self) {
993 self.delivered += other.delivered;
994 self.disconnected_slow += other.disconnected_slow;
995 self.disconnected_closed += other.disconnected_closed;
996 }
997}
998
999fn map_port_error<T>(error: TrySendError<T>) -> PortDispatchError {
1000 match error {
1001 TrySendError::Full(_) => PortDispatchError::Full,
1002 TrySendError::Disconnected(_) => PortDispatchError::Disconnected,
1003 }
1004}
1005
1006fn drain_receiver<T>(receiver: &Receiver<T>, limit: usize) -> Vec<T> {
1007 let mut items = Vec::new();
1008 for _ in 0..limit {
1009 match receiver.try_recv() {
1010 Ok(item) => items.push(item),
1011 Err(TryRecvError::Empty | TryRecvError::Disconnected) => break,
1012 }
1013 }
1014 items
1015}
1016
1017fn subscribe_bounded<T>(
1018 subscribers: &Mutex<Vec<SyncSender<T>>>,
1019 capacity: usize,
1020 max_subscribers: usize,
1021) -> Receiver<T> {
1022 let (sender, receiver) = sync_channel(bounded_subscription_capacity(capacity));
1023 let mut subscribers = subscribers
1024 .lock()
1025 .unwrap_or_else(|poisoned| poisoned.into_inner());
1026 if subscribers.len() < max_subscribers {
1027 subscribers.push(sender);
1028 }
1029 receiver
1030}
1031
1032fn bounded_ingress_capacity(requested: usize) -> usize {
1033 requested.clamp(1, MAX_INGRESS_CAPACITY)
1034}
1035
1036fn bounded_subscription_capacity(requested: usize) -> usize {
1037 requested.clamp(1, MAX_SUBSCRIPTION_CAPACITY)
1038}
1039
1040fn publish_many<T: Clone>(
1041 subscribers: &Mutex<Vec<SyncSender<T>>>,
1042 items: impl IntoIterator<Item = T>,
1043) -> PublishReport {
1044 let mut report = PublishReport::default();
1045 let mut subscribers = subscribers
1046 .lock()
1047 .unwrap_or_else(|poisoned| poisoned.into_inner());
1048 for item in items {
1049 report += publish_one_locked(&mut subscribers, item);
1050 }
1051 report
1052}
1053
1054fn publish_one_locked<T: Clone>(subscribers: &mut Vec<SyncSender<T>>, item: T) -> PublishReport {
1055 let mut report = PublishReport::default();
1056 let mut index = 0;
1057 while index < subscribers.len() {
1058 match subscribers[index].try_send(item.clone()) {
1059 Ok(()) => {
1060 report.delivered += 1;
1061 index += 1;
1062 }
1063 Err(TrySendError::Full(_)) => {
1064 subscribers.swap_remove(index);
1065 report.disconnected_slow += 1;
1066 }
1067 Err(TrySendError::Disconnected(_)) => {
1068 subscribers.swap_remove(index);
1069 report.disconnected_closed += 1;
1070 }
1071 }
1072 }
1073 report
1074}
1075
1076#[cfg(test)]
1077mod tests {
1078 use super::*;
1079 use gate4agent_types::{AgentId, CommandId, ControlCommand, ControlEventKind, TransportKind};
1080
1081 fn command(id: u64) -> CommandEnvelope {
1082 CommandEnvelope {
1083 id: CommandId(id),
1084 command: ControlCommand::Register {
1085 instance_id: AgentInstanceId(id),
1086 agent_id: AgentId::new("claude").unwrap(),
1087 transport: TransportKind::Pty,
1088 },
1089 }
1090 }
1091
1092 fn event(sequence: u64) -> ControlEvent {
1093 ControlEvent {
1094 sequence,
1095 command_id: None,
1096 instance_id: AgentInstanceId(1),
1097 generation: SessionGeneration(1),
1098 event: ControlEventKind::Registered,
1099 }
1100 }
1101
1102 #[test]
1103 fn legacy_command_ingress_remains_bounded_and_non_blocking() {
1104 let (handle, kernel) = bounded_port(1);
1105 handle.dispatch(command(1)).unwrap();
1106 assert_eq!(handle.dispatch(command(2)), Err(PortDispatchError::Full));
1107 assert_eq!(kernel.drain_commands(10), vec![command(1)]);
1108 }
1109
1110 #[test]
1111 fn legacy_snapshot_publication_replaces_current_control_truth() {
1112 let (handle, kernel) = bounded_port(1);
1113 let snapshot = ControlSnapshot {
1114 revision: 7,
1115 ..ControlSnapshot::default()
1116 };
1117 kernel.publish_snapshot(snapshot.clone());
1118 assert_eq!(*handle.snapshot(), snapshot);
1119 }
1120
1121 #[test]
1122 fn slow_control_subscriber_is_disconnected_without_silent_loss() {
1123 let (handle, kernel) = bounded_port(1);
1124 let subscription = handle.subscribe(1);
1125 assert_eq!(kernel.publish_events([event(1)]).delivered, 1);
1126
1127 let report = kernel.publish_events([event(2)]);
1128 assert_eq!(report.disconnected_slow, 1);
1129 assert_eq!(subscription.try_recv().unwrap(), event(1));
1130 assert_eq!(subscription.try_recv(), Err(TryRecvError::Disconnected));
1131 }
1132
1133 #[test]
1134 fn closed_control_subscriber_is_pruned_explicitly() {
1135 let (handle, kernel) = bounded_port(1);
1136 let subscription = handle.subscribe(1);
1137 drop(subscription);
1138
1139 let report = kernel.publish_events([event(1)]);
1140 assert_eq!(report.disconnected_closed, 1);
1141 assert_eq!(report.delivered, 0);
1142 }
1143}