Skip to main content

gate4agent_handle/
lib.rs

1//! Bounded in-process ports for the Gate4Agent backend control plane.
2
3mod 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(&current.tools),
338            provider_runtime: current.provider_runtime.clone(),
339        });
340    }
341
342    /// Slow subscribers are disconnected instead of silently losing ordered
343    /// control events. Their next receive observes channel disconnection.
344    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    /// Publishes one already-reduced kernel step in a fixed edge order.
364    ///
365    /// Provider outcomes and effects first reconcile the private executor
366    /// boundary. The immutable combined snapshot is then replaced before
367    /// request/authority outcomes, control events, exact scoped completions,
368    /// and their trailing source-gap marker. A successfully closed client's
369    /// subscriptions are disconnected only after its `ClientClosed`
370    /// completions have been offered.
371    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    /// Deactivates every clone before attempting to enqueue `CloseClient`.
641    ///
642    /// A full ingress leaves the client inactive and retryable without
643    /// reserving the global authority sequence or blocking unrelated authority
644    /// commands. Once enqueued, repeated calls return the one queued sequence.
645    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}