Skip to main content

gate4agent_runtime_native/
lib.rs

1//! Tick-driven native runtime for embedding gate4agent in an owning app core.
2
3mod launch_profiles;
4pub mod shell_efficiency;
5pub mod tick_profile;
6mod vendor_contract;
7
8pub use launch_profiles::{
9    NativeChildEnvironmentResolveError, NativeChildEnvironmentResolver, NativeLaunchProfile,
10    NativeMcpServerLaunchOverlay, NativeInstanceLaunchOverlay, NativeLaunchEnvironmentOverlay,
11    NativeLaunchProfileControl,
12    NativeLaunchProfileDescriptor, NativeLaunchProfileError, NativeLaunchProfileId,
13    ZAI_GLM_ANTHROPIC_BASE_URL, ZAI_GLM_CLAUDE_OPTIONAL_ENV_KEYS,
14    ZAI_GLM_CLAUDE_OWNED_ENV_KEYS, ZAI_GLM_CLAUDE_PROFILE, ZAI_GLM_CLAUDE_PROFILE_ID,
15    ZAI_GLM_CLAUDE_PROFILE_REVISION, ZAI_GLM_CLAUDE_REQUIRED_ENV_KEYS,
16};
17pub use vendor_contract::{
18    probe_installed_vendor_version, resolve_vendor_contract, VendorCapabilitySet,
19    VendorCapabilityUnavailableReason, VendorCapabilityVerdict, VendorCliFamily,
20    VendorContractResolution, VendorFallbackReason, VendorLauncherIdentity, VendorPlatform,
21    VendorRuntimeMode, VendorVersionProbeCache, VendorVersionProbeResult,
22    VendorVersionProbeStatus, VendorVersionStatus,
23    CLAUDE_WINDOWS_X86_64_2_1_223_CONTRACT_ID,
24    CLAUDE_WINDOWS_X86_64_2_1_224_CONTRACT_ID,
25};
26
27use gate4agent_catalog::{builtin_registry, AgentRegistry, EnvMutation, McpServerSpec};
28use gate4agent_handle::{
29    bounded_control_plane, ControlPlaneKernelPort, Gate4AgentHandle,
30    ProviderRuntimeError, PublishReport, ToolAuthorityHandle,
31};
32use gate4agent_kernel::{CommandOutcome, Gate4AgentKernel};
33use gate4agent_provider_ports::{
34    discover_history, load_history_session, prepare_resume, HistoryCandidate,
35    HistoryDiscoveryRequest, HistoryLoadRequest, PreparedResume, ResumeAuthority,
36    ResumeAuthorityDecision, ResumeOutcome, ResumeRequest,
37    HISTORY_DISCOVERY_LIMIT_MAX as PROVIDER_HISTORY_DISCOVERY_LIMIT_MAX,
38};
39use gate4agent_shell_capabilities::NativeCapabilityProbeAuthority;
40use gate4agent_shell_history::NativeHistoryAuthority;
41pub use gate4agent_shell_history::{orca_home_roots, NativeHistoryConfig, NativeHistoryRoot};
42pub use gate4agent_adapters::{HistorySourceLayout, OneShotSessionPersistence};
43pub use gate4agent_shell_native::{
44    NativeProviderExecutor, NativeProviderExit, NativeProviderOperation,
45    NativeProviderOperationError, NativeProviderResultPoll, PhysicalExitAck,
46    ProviderOperationKey, ProviderStopCause, ProviderSupervisorFault,
47    ProviderSupervisorFaultKind, ProviderSupervisorSnapshot, ProviderSupervisorState,
48};
49use gate4agent_shell_native::{
50    NativeEffectShell, ProviderSupervisor, ProviderSupervisorBuildError,
51    MAX_PROVIDER_SUPERVISOR_EVENTS,
52};
53use shell_efficiency::ShellEfficiencyProfile;
54use gate4agent_tool_engine::{
55    CapabilityOwner, CapabilityProviderDescriptor, ProviderBindingId, ToolEngineError,
56    ToolProviderId,
57};
58use gate4agent_types::{
59    AgentId, AgentInstanceId, CapabilityProbeFailure, ControlEffect, ControlObservation,
60    EffectEnvelope, HistoryCandidateSummary, HistoryMessageRecord, HistoryMessageRole,
61    HistorySessionRecord, NativeSessionCatalogEntry, NativeSessionCatalogScope,
62    NativeSessionCatalogWindow, NativeSessionExternalGroup, NativeSessionExternalGroupKind,
63    NativeSessionPreview,
64    NativeSessionPreviewMessage, ObservationEnvelope, ProviderSessionIdentity,
65    ProviderRuntimeCapability, ProviderRuntimePolicy,
66    PipeProtocol, ResumeAuthorityTarget, ResumeLaunchRequest, SessionGeneration,
67    TransportKind,
68};
69use std::collections::{BTreeMap, HashMap, VecDeque};
70use std::convert::Infallible;
71use std::ffi::OsString;
72use std::path::{Component, Path, PathBuf};
73use std::sync::atomic::{AtomicUsize, Ordering};
74use std::sync::{Arc, Mutex};
75use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
76use thiserror::Error;
77use tokio::sync::mpsc::{self, error::TrySendError, Receiver, Sender};
78use tokio::time::MissedTickBehavior;
79
80type TerminalFrameKey = (AgentInstanceId, SessionGeneration);
81const PROVIDER_EVENT_DRAIN_QUANTUM: usize = 64;
82const MAX_PROVIDER_SHUTDOWN_TIMEOUT_MS: u64 = 86_400_000;
83const NATIVE_SESSION_INDEX_REFRESH_INTERVAL: Duration = Duration::from_secs(30);
84const NATIVE_SESSION_RECENT_WINDOW_MS: u64 = 7 * 24 * 60 * 60 * 1_000;
85
86/// Independent read-only authority for projecting provider-native history into
87/// bounded session metadata. It does not create kernel or provider sessions.
88pub struct NativeSessionCatalogAuthority {
89    catalog: AgentRegistry,
90    history_config: NativeHistoryConfig,
91    authorities: HashMap<AgentId, NativeHistoryAuthority>,
92    indexes: HashMap<AgentId, NativeProviderSessionIndex>,
93}
94
95struct NativeProviderSessionIndex {
96    refreshed_at: Instant,
97    revision: u64,
98    recent_cutoff_unix_ms: u64,
99    sessions: Vec<IndexedNativeSession>,
100}
101
102#[derive(Clone, Eq, PartialEq)]
103struct IndexedNativeSession {
104    candidate: HistoryCandidate,
105    selection_id: String,
106    modified_at_unix_ms: Option<u64>,
107    cwd_key: Option<NativePathKey>,
108    session_id: String,
109    title: Option<String>,
110    model: Option<String>,
111    project_label: String,
112    provider_session: ProviderSessionIdentity,
113}
114
115#[derive(Clone, Debug, Eq, PartialEq)]
116pub struct ScopedNativeSessionCatalogEntry {
117    pub metadata: NativeSessionCatalogEntry,
118    pub external_group: Option<NativeSessionExternalGroup>,
119    pub provider_session: ProviderSessionIdentity,
120}
121
122#[derive(Clone, Debug, Eq, PartialEq)]
123pub struct ScopedNativeSessionCatalogPage {
124    pub revision: u64,
125    pub cutoff_unix_ms: u64,
126    pub window: NativeSessionCatalogWindow,
127    pub entries: Vec<ScopedNativeSessionCatalogEntry>,
128    pub recent_total_count: u32,
129    pub older_total_count: u32,
130    pub next_after_selection_id: Option<String>,
131    pub remaining_count: u32,
132}
133
134#[derive(Clone, Debug, Eq, PartialEq)]
135pub struct NativeSessionSelectionResolution {
136    pub identity: ProviderSessionIdentity,
137    pub external_group: Option<NativeSessionExternalGroup>,
138}
139
140#[derive(Clone, Debug, Eq, PartialEq)]
141pub struct NativeSessionCatalogPage {
142    pub revision: u64,
143    pub cutoff_unix_ms: u64,
144    pub window: NativeSessionCatalogWindow,
145    pub entries: Vec<NativeSessionCatalogEntry>,
146    pub recent_total_count: u32,
147    pub older_total_count: u32,
148    pub next_after_selection_id: Option<String>,
149    pub remaining_count: u32,
150}
151
152impl NativeSessionCatalogAuthority {
153    pub fn new(config: NativeHistoryConfig) -> Self {
154        Self {
155            catalog: builtin_registry().clone(),
156            history_config: config,
157            authorities: HashMap::new(),
158            indexes: HashMap::new(),
159        }
160    }
161
162    pub fn catalog(
163        &mut self,
164        provider: &AgentId,
165        canonical_workspace: &Path,
166        limit: u16,
167    ) -> Result<Vec<NativeSessionCatalogEntry>, NativeSessionCatalogError> {
168        self.catalog_for_workspace(
169            provider,
170            canonical_workspace,
171            &[canonical_workspace.to_path_buf()],
172            limit,
173        )
174    }
175
176    pub fn catalog_for_workspace(
177        &mut self,
178        provider: &AgentId,
179        canonical_workspace: &Path,
180        registered_workspace_roots: &[PathBuf],
181        limit: u16,
182    ) -> Result<Vec<NativeSessionCatalogEntry>, NativeSessionCatalogError> {
183        self.catalog_initial_for_workspace(
184            provider,
185            canonical_workspace,
186            registered_workspace_roots,
187            limit,
188        )
189        .map(|page| page.entries)
190    }
191
192    pub fn catalog_initial_for_workspace(
193        &mut self,
194        provider: &AgentId,
195        canonical_workspace: &Path,
196        registered_workspace_roots: &[PathBuf],
197        limit: u16,
198    ) -> Result<NativeSessionCatalogPage, NativeSessionCatalogError> {
199        self.catalog_initial_for_scope(
200            provider,
201            NativeSessionCatalogScope::Workspace,
202            Some(canonical_workspace),
203            registered_workspace_roots,
204            limit,
205        )
206        .map(unscoped_native_catalog_page)
207    }
208
209    pub fn catalog_initial_for_scope(
210        &mut self,
211        provider: &AgentId,
212        scope: NativeSessionCatalogScope,
213        canonical_workspace: Option<&Path>,
214        registered_workspace_roots: &[PathBuf],
215        limit: u16,
216    ) -> Result<ScopedNativeSessionCatalogPage, NativeSessionCatalogError> {
217        let ownership = NativeSessionOwnership::new(
218            scope,
219            canonical_workspace,
220            registered_workspace_roots,
221        )
222        .map_err(|_| NativeSessionCatalogError::WorkspaceUnavailable)?;
223        self.ensure_provider_index(provider)
224            .map_err(map_catalog_index_error)?;
225        let (index_revision, cutoff_unix_ms) = self
226            .indexes
227            .get(provider)
228            .map(|index| (index.revision, index.recent_cutoff_unix_ms))
229            .ok_or(NativeSessionCatalogError::CatalogUnavailable)?;
230        let revision = scoped_native_catalog_snapshot_revision(
231            index_revision,
232            cutoff_unix_ms,
233            ownership.fingerprint(),
234        );
235        self.catalog_page_for_scope(
236            provider,
237            scope,
238            canonical_workspace,
239            registered_workspace_roots,
240            NativeSessionCatalogWindow::Recent,
241            revision,
242            cutoff_unix_ms,
243            None,
244            limit,
245        )
246    }
247
248    pub fn catalog_page_for_workspace(
249        &mut self,
250        provider: &AgentId,
251        canonical_workspace: &Path,
252        registered_workspace_roots: &[PathBuf],
253        window: NativeSessionCatalogWindow,
254        revision: u64,
255        cutoff_unix_ms: u64,
256        after_selection_id: Option<&str>,
257        limit: u16,
258    ) -> Result<NativeSessionCatalogPage, NativeSessionCatalogError> {
259        self.catalog_page_for_scope(
260            provider,
261            NativeSessionCatalogScope::Workspace,
262            Some(canonical_workspace),
263            registered_workspace_roots,
264            window,
265            revision,
266            cutoff_unix_ms,
267            after_selection_id,
268            limit,
269        )
270        .map(unscoped_native_catalog_page)
271    }
272
273    pub fn catalog_page_for_scope(
274        &mut self,
275        provider: &AgentId,
276        scope: NativeSessionCatalogScope,
277        canonical_workspace: Option<&Path>,
278        registered_workspace_roots: &[PathBuf],
279        window: NativeSessionCatalogWindow,
280        revision: u64,
281        cutoff_unix_ms: u64,
282        after_selection_id: Option<&str>,
283        limit: u16,
284    ) -> Result<ScopedNativeSessionCatalogPage, NativeSessionCatalogError> {
285        if !(1..=gate4agent_types::NATIVE_SESSION_CATALOG_LIMIT_MAX).contains(&limit) {
286            return Err(NativeSessionCatalogError::InvalidLimit);
287        }
288        let ownership = NativeSessionOwnership::new(
289            scope,
290            canonical_workspace,
291            registered_workspace_roots,
292        )
293            .map_err(|_| NativeSessionCatalogError::WorkspaceUnavailable)?;
294        self.ensure_provider_index(provider)
295            .map_err(map_catalog_index_error)?;
296        let index = self
297            .indexes
298            .get(provider)
299            .expect("ensured native provider index must exist");
300        if scoped_native_catalog_snapshot_revision(
301            index.revision,
302            cutoff_unix_ms,
303            ownership.fingerprint(),
304        ) != revision
305        {
306            return Err(NativeSessionCatalogError::StaleCatalog);
307        }
308        let external_group_ids = ownership.external_group_ids(&index.sessions);
309        let mut seen_sessions = std::collections::HashSet::new();
310        let mut recent = Vec::new();
311        let mut older = Vec::new();
312        for indexed in &index.sessions {
313            if !ownership.owns(&indexed.cwd_key)
314                || !seen_sessions.insert(indexed.session_id.clone())
315            {
316                continue;
317            }
318            let entry = NativeSessionCatalogEntry {
319                selection_id: indexed.selection_id.clone(),
320                session_id: indexed.session_id.clone(),
321                title: indexed.title.clone(),
322                modified_at_unix_ms: indexed.modified_at_unix_ms,
323                model: indexed.model.clone(),
324                message_count: 0,
325                completed_turn_count: None,
326            };
327            let external_group = ownership.external_group(indexed, &external_group_ids);
328            if entry.validate().is_ok()
329                && external_group
330                    .as_ref()
331                    .map_or(true, |group| group.validate().is_ok())
332            {
333                let entry = ScopedNativeSessionCatalogEntry {
334                    metadata: entry,
335                    external_group,
336                    provider_session: indexed.provider_session.clone(),
337                };
338                match native_catalog_window_for_modified(
339                    indexed.modified_at_unix_ms,
340                    cutoff_unix_ms,
341                ) {
342                    NativeSessionCatalogWindow::Recent => recent.push(entry),
343                    NativeSessionCatalogWindow::Older => older.push(entry),
344                }
345            }
346        }
347        let recent_total_count = u32::try_from(recent.len()).unwrap_or(u32::MAX);
348        let older_total_count = u32::try_from(older.len()).unwrap_or(u32::MAX);
349        let selected = match window {
350            NativeSessionCatalogWindow::Recent => recent,
351            NativeSessionCatalogWindow::Older => older,
352        };
353        let start = match after_selection_id {
354            Some(cursor) => selected
355                .iter()
356                .position(|entry| entry.metadata.selection_id == cursor)
357                .map(|index| index + 1)
358                .ok_or(NativeSessionCatalogError::StaleCatalog)?,
359            None => 0,
360        };
361        let end = start.saturating_add(usize::from(limit)).min(selected.len());
362        let entries = selected[start..end].to_vec();
363        let next_after_selection_id = (end < selected.len())
364            .then(|| entries.last().map(|entry| entry.metadata.selection_id.clone()))
365            .flatten();
366        let remaining_count = u32::try_from(selected.len().saturating_sub(end))
367            .unwrap_or(u32::MAX);
368        Ok(ScopedNativeSessionCatalogPage {
369            revision,
370            cutoff_unix_ms,
371            window,
372            entries,
373            recent_total_count,
374            older_total_count,
375            next_after_selection_id,
376            remaining_count,
377        })
378    }
379
380    fn ensure_provider_index(
381        &mut self,
382        provider: &AgentId,
383    ) -> Result<(), NativeSessionIndexError> {
384        if self.indexes.get(provider).is_some_and(|index| {
385            index.refreshed_at.elapsed() < NATIVE_SESSION_INDEX_REFRESH_INTERVAL
386        }) {
387            return Ok(());
388        }
389        let spec = self
390            .catalog
391            .get(provider)
392            .cloned()
393            .ok_or(NativeSessionIndexError::UnsupportedProvider)?;
394        let request = HistoryDiscoveryRequest::from_spec(
395            &spec,
396            None,
397            PROVIDER_HISTORY_DISCOVERY_LIMIT_MAX,
398        )
399        .map_err(|_| NativeSessionIndexError::UnsupportedProvider)?;
400        let authority = self
401            .authorities
402            .entry(provider.clone())
403            .or_insert_with(|| NativeHistoryAuthority::new(self.history_config.clone()));
404        let locators = authority.discover_locator_index(&request)
405            .map_err(|_| NativeSessionIndexError::Unavailable)?;
406        if !authority.take_discovery_issues().is_empty() {
407            return Err(NativeSessionIndexError::Unavailable);
408        }
409        let mut sessions = Vec::with_capacity(locators.len());
410        for locator in locators {
411            let indexed_candidate = locator.candidate().clone();
412            let candidate = locator.candidate();
413            let selection_id = candidate.id().as_str().to_owned();
414            let modified_at_unix_ms = candidate.modified_at_unix_ms();
415            let cwd_key = locator.cwd().and_then(NativePathKey::from_declared_cwd);
416            let project_label = if cwd_key.is_some() {
417                locator
418                    .cwd()
419                    .map(native_external_group_label)
420                    .unwrap_or_else(|| "Global".to_owned())
421            } else {
422                "Global".to_owned()
423            };
424            let load = HistoryLoadRequest::new(&request, indexed_candidate.clone())
425                .map_err(|_| NativeSessionIndexError::Unavailable)?;
426            let provider_session = authority
427                .resume_provider_session(&load, locator.session_id().to_owned())
428                .map_err(|_| NativeSessionIndexError::Unavailable)?;
429            sessions.push(IndexedNativeSession {
430                candidate: indexed_candidate,
431                selection_id,
432                modified_at_unix_ms,
433                cwd_key,
434                session_id: locator.session_id().to_owned(),
435                title: locator.title().map(str::to_owned),
436                model: locator.model().map(str::to_owned),
437                project_label,
438                provider_session,
439            });
440        }
441        sessions.sort_by(|left, right| {
442            right
443                .modified_at_unix_ms
444                .cmp(&left.modified_at_unix_ms)
445                .then_with(|| left.session_id.cmp(&right.session_id))
446                .then_with(|| left.selection_id.cmp(&right.selection_id))
447        });
448        if let Some(current) = self.indexes.get_mut(provider) {
449            if current.sessions == sessions {
450                current.refreshed_at = Instant::now();
451                return Ok(());
452            }
453        }
454        let revision = self
455            .indexes
456            .get(provider)
457            .map(|index| index.revision.saturating_add(1))
458            .unwrap_or(1);
459        let recent_cutoff_unix_ms = SystemTime::now()
460            .duration_since(UNIX_EPOCH)
461            .ok()
462            .and_then(|duration| u64::try_from(duration.as_millis()).ok())
463            .map(|now| now.saturating_sub(NATIVE_SESSION_RECENT_WINDOW_MS))
464            .ok_or(NativeSessionIndexError::Unavailable)?;
465        self.indexes.insert(
466            provider.clone(),
467            NativeProviderSessionIndex {
468                refreshed_at: Instant::now(),
469                revision,
470                recent_cutoff_unix_ms,
471                sessions,
472            },
473        );
474        Ok(())
475    }
476
477    pub fn preview(
478        &mut self,
479        provider: &AgentId,
480        canonical_workspace: &Path,
481        selection_id: &str,
482        message_limit: u16,
483    ) -> Result<NativeSessionPreview, NativeSessionPreviewError> {
484        self.preview_for_workspace(
485            provider,
486            canonical_workspace,
487            &[canonical_workspace.to_path_buf()],
488            selection_id,
489            message_limit,
490        )
491    }
492
493    pub fn preview_for_workspace(
494        &mut self,
495        provider: &AgentId,
496        canonical_workspace: &Path,
497        registered_workspace_roots: &[PathBuf],
498        selection_id: &str,
499        message_limit: u16,
500    ) -> Result<NativeSessionPreview, NativeSessionPreviewError> {
501        self.preview_selected(
502            provider,
503            NativeSessionCatalogScope::Workspace,
504            Some(canonical_workspace),
505            registered_workspace_roots,
506            None,
507            NativeSessionPreviewSelector::Candidate(selection_id),
508            message_limit,
509        )
510    }
511
512    pub fn preview_for_scope(
513        &mut self,
514        provider: &AgentId,
515        scope: NativeSessionCatalogScope,
516        canonical_workspace: Option<&Path>,
517        registered_workspace_roots: &[PathBuf],
518        catalog_revision: u64,
519        recent_cutoff_unix_ms: u64,
520        selection_id: &str,
521        message_limit: u16,
522    ) -> Result<NativeSessionPreview, NativeSessionPreviewError> {
523        self.preview_selected(
524            provider,
525            scope,
526            canonical_workspace,
527            registered_workspace_roots,
528            Some((catalog_revision, recent_cutoff_unix_ms)),
529            NativeSessionPreviewSelector::Candidate(selection_id),
530            message_limit,
531        )
532    }
533
534    pub fn preview_session_id(
535        &mut self,
536        provider: &AgentId,
537        canonical_workspace: &Path,
538        session_id: &str,
539        message_limit: u16,
540    ) -> Result<NativeSessionPreview, NativeSessionPreviewError> {
541        self.preview_session_id_for_workspace(
542            provider,
543            canonical_workspace,
544            &[canonical_workspace.to_path_buf()],
545            session_id,
546            message_limit,
547        )
548    }
549
550    pub fn preview_session_id_for_workspace(
551        &mut self,
552        provider: &AgentId,
553        canonical_workspace: &Path,
554        registered_workspace_roots: &[PathBuf],
555        session_id: &str,
556        message_limit: u16,
557    ) -> Result<NativeSessionPreview, NativeSessionPreviewError> {
558        self.preview_selected(
559            provider,
560            NativeSessionCatalogScope::Workspace,
561            Some(canonical_workspace),
562            registered_workspace_roots,
563            None,
564            NativeSessionPreviewSelector::Session(session_id),
565            message_limit,
566        )
567    }
568
569    fn preview_selected(
570        &mut self,
571        provider: &AgentId,
572        scope: NativeSessionCatalogScope,
573        canonical_workspace: Option<&Path>,
574        registered_workspace_roots: &[PathBuf],
575        catalog_snapshot: Option<(u64, u64)>,
576        selector: NativeSessionPreviewSelector<'_>,
577        message_limit: u16,
578    ) -> Result<NativeSessionPreview, NativeSessionPreviewError> {
579        if !(1..=gate4agent_types::NATIVE_SESSION_PREVIEW_MESSAGE_LIMIT_MAX)
580            .contains(&message_limit)
581        {
582            return Err(NativeSessionPreviewError::InvalidLimit);
583        }
584        let ownership = NativeSessionOwnership::new(
585            scope,
586            canonical_workspace,
587            registered_workspace_roots,
588        )
589            .map_err(|_| NativeSessionPreviewError::WorkspaceUnavailable)?;
590        let cached_candidate = match selector {
591            NativeSessionPreviewSelector::Candidate(selection_id) => self
592                .indexes
593                .get(provider)
594                .is_some_and(|index| {
595                    index.sessions.iter().any(|session| {
596                        session.selection_id == selection_id
597                            && ownership.owns(&session.cwd_key)
598                    })
599                }),
600            NativeSessionPreviewSelector::Session(_) => false,
601        };
602        if !cached_candidate {
603            self.ensure_provider_index(provider)
604                .map_err(|error| match error {
605                    NativeSessionIndexError::UnsupportedProvider => {
606                        NativeSessionPreviewError::UnsupportedProvider
607                    }
608                    NativeSessionIndexError::Unavailable => {
609                        NativeSessionPreviewError::PreviewUnavailable
610                    }
611                })?;
612        }
613        let index = self
614            .indexes
615            .get(provider)
616            .ok_or(NativeSessionPreviewError::SessionNotFound)?;
617        if let Some((catalog_revision, recent_cutoff_unix_ms)) = catalog_snapshot {
618            if index.recent_cutoff_unix_ms != recent_cutoff_unix_ms
619                || scoped_native_catalog_snapshot_revision(
620                    index.revision,
621                    recent_cutoff_unix_ms,
622                    ownership.fingerprint(),
623                ) != catalog_revision
624            {
625                return Err(NativeSessionPreviewError::StaleCatalog);
626            }
627        }
628        let mut matched = None;
629        for indexed in &index.sessions {
630            if matches!(selector, NativeSessionPreviewSelector::Candidate(selection_id)
631                if indexed.selection_id != selection_id)
632                || matches!(selector, NativeSessionPreviewSelector::Session(session_id)
633                    if indexed.session_id != session_id)
634                || !ownership.owns(&indexed.cwd_key)
635            {
636                continue;
637            }
638            if matched.is_some() {
639                return Err(NativeSessionPreviewError::AmbiguousSession);
640            }
641            matched = Some((
642                indexed.candidate.clone(),
643                indexed.session_id.clone(),
644                indexed.cwd_key.clone(),
645                indexed.modified_at_unix_ms,
646            ));
647        }
648        let (candidate, expected_session_id, expected_cwd, modified_at_unix_ms) =
649            matched.ok_or(NativeSessionPreviewError::SessionNotFound)?;
650        let spec = self
651            .catalog
652            .get(provider)
653            .cloned()
654            .ok_or(NativeSessionPreviewError::UnsupportedProvider)?;
655        let request = HistoryDiscoveryRequest::from_spec(
656            &spec,
657            None,
658            PROVIDER_HISTORY_DISCOVERY_LIMIT_MAX,
659        )
660        .map_err(|_| NativeSessionPreviewError::UnsupportedProvider)?;
661        let load = HistoryLoadRequest::new(&request, candidate)
662            .map_err(|_| NativeSessionPreviewError::PreviewUnavailable)?;
663        let authority = self
664            .authorities
665            .get_mut(provider)
666            .ok_or(NativeSessionPreviewError::PreviewUnavailable)?;
667        let preview_load = authority.load_preview_session(&load)
668            .map_err(|_| NativeSessionPreviewError::PreviewUnavailable)?;
669        let source_truncated = preview_load.source_truncated();
670        let session = preview_load.session();
671        let current_cwd = session
672            .cwd
673            .as_deref()
674            .and_then(NativePathKey::from_declared_cwd);
675        if session.session_id != expected_session_id
676            || current_cwd != expected_cwd
677            || !ownership.owns(&current_cwd)
678        {
679            return Err(NativeSessionPreviewError::PreviewUnavailable);
680        }
681        project_native_session_preview(
682            session,
683            modified_at_unix_ms,
684            source_truncated,
685            message_limit,
686        )
687    }
688
689    pub fn resolve_selection_for_scope(
690        &mut self,
691        provider: &AgentId,
692        scope: NativeSessionCatalogScope,
693        canonical_workspace: Option<&Path>,
694        registered_workspace_roots: &[PathBuf],
695        catalog_revision: u64,
696        recent_cutoff_unix_ms: u64,
697        selection_id: &str,
698    ) -> Result<NativeSessionSelectionResolution, NativeSessionPreviewError> {
699        gate4agent_types::validate_candidate_id(selection_id)
700            .map_err(|_| NativeSessionPreviewError::SessionNotFound)?;
701        let ownership = NativeSessionOwnership::new(
702            scope,
703            canonical_workspace,
704            registered_workspace_roots,
705        )
706        .map_err(|_| NativeSessionPreviewError::WorkspaceUnavailable)?;
707        self.ensure_provider_index(provider)
708            .map_err(|error| match error {
709                NativeSessionIndexError::UnsupportedProvider => {
710                    NativeSessionPreviewError::UnsupportedProvider
711                }
712                NativeSessionIndexError::Unavailable => {
713                    NativeSessionPreviewError::PreviewUnavailable
714                }
715            })?;
716        let index = self
717            .indexes
718            .get(provider)
719            .ok_or(NativeSessionPreviewError::SessionNotFound)?;
720        if index.recent_cutoff_unix_ms != recent_cutoff_unix_ms
721            || scoped_native_catalog_snapshot_revision(
722                index.revision,
723                recent_cutoff_unix_ms,
724                ownership.fingerprint(),
725            ) != catalog_revision
726        {
727            return Err(NativeSessionPreviewError::StaleCatalog);
728        }
729        let external_group_ids = ownership.external_group_ids(&index.sessions);
730        let mut matched = None;
731        for indexed in &index.sessions {
732            if indexed.selection_id != selection_id || !ownership.owns(&indexed.cwd_key) {
733                continue;
734            }
735            if matched.is_some() {
736                return Err(NativeSessionPreviewError::AmbiguousSession);
737            }
738            matched = Some((
739                indexed.candidate.clone(),
740                indexed.session_id.clone(),
741                indexed.cwd_key.clone(),
742                ownership.external_group(indexed, &external_group_ids),
743            ));
744        }
745        let (candidate, expected_session_id, expected_cwd, external_group) =
746            matched.ok_or(NativeSessionPreviewError::SessionNotFound)?;
747        let spec = self
748            .catalog
749            .get(provider)
750            .cloned()
751            .ok_or(NativeSessionPreviewError::UnsupportedProvider)?;
752        let request = HistoryDiscoveryRequest::from_spec(
753            &spec,
754            None,
755            PROVIDER_HISTORY_DISCOVERY_LIMIT_MAX,
756        )
757        .map_err(|_| NativeSessionPreviewError::UnsupportedProvider)?;
758        let load = HistoryLoadRequest::new(&request, candidate)
759            .map_err(|_| NativeSessionPreviewError::PreviewUnavailable)?;
760        let authority = self
761            .authorities
762            .get_mut(provider)
763            .ok_or(NativeSessionPreviewError::PreviewUnavailable)?;
764        let session = authority
765            .load_preview_session(&load)
766            .map_err(|_| NativeSessionPreviewError::PreviewUnavailable)?
767            .session();
768        let current_cwd = session
769            .cwd
770            .as_deref()
771            .and_then(NativePathKey::from_declared_cwd);
772        if session.session_id != expected_session_id
773            || current_cwd != expected_cwd
774            || !ownership.owns(&current_cwd)
775        {
776            return Err(NativeSessionPreviewError::PreviewUnavailable);
777        }
778        let identity = authority
779            .resume_provider_session(&load, session.session_id)
780            .map_err(|_| NativeSessionPreviewError::PreviewUnavailable)?;
781        identity
782            .validate()
783            .map_err(|_| NativeSessionPreviewError::PreviewUnavailable)?;
784        if external_group
785            .as_ref()
786            .is_some_and(|group| group.validate().is_err())
787        {
788            return Err(NativeSessionPreviewError::PreviewUnavailable);
789        }
790        Ok(NativeSessionSelectionResolution {
791            identity,
792            external_group,
793        })
794    }
795}
796
797struct NativeSessionOwnership {
798    scope: NativeSessionCatalogScope,
799    selected: Option<NativePathKey>,
800    registered: Vec<NativePathKey>,
801}
802
803impl NativeSessionOwnership {
804    fn new(
805        scope: NativeSessionCatalogScope,
806        selected: Option<&Path>,
807        registered: &[PathBuf],
808    ) -> Result<Self, ()> {
809        let mut registered = registered
810            .iter()
811            .filter(|root| root.is_absolute())
812            .map(std::fs::canonicalize)
813            .collect::<Result<Vec<_>, _>>()
814            .map_err(|_| ())?
815            .into_iter()
816            .map(|root| NativePathKey::from_canonical_root(&root).ok_or(()))
817            .collect::<Result<Vec<_>, _>>()?;
818        let selected = match (scope, selected) {
819            (NativeSessionCatalogScope::Workspace, Some(selected)) => {
820                let selected = std::fs::canonicalize(selected).map_err(|_| ())?;
821                if !selected.is_dir() {
822                    return Err(());
823                }
824                let selected = NativePathKey::from_canonical_root(&selected).ok_or(())?;
825                if !registered.iter().any(|root| root == &selected) {
826                    registered.push(selected.clone());
827                }
828                Some(selected)
829            }
830            (NativeSessionCatalogScope::Unregistered, None) => None,
831            _ => return Err(()),
832        };
833        registered.sort();
834        registered.dedup();
835        Ok(Self {
836            scope,
837            selected,
838            registered,
839        })
840    }
841
842    fn owns(&self, cwd: &Option<NativePathKey>) -> bool {
843        let Some(cwd) = cwd.as_ref() else {
844            return matches!(self.scope, NativeSessionCatalogScope::Unregistered);
845        };
846        let owner = self.registered
847            .iter()
848            .filter(|root| cwd.starts_with(root))
849            .max_by_key(|root| root.components.len());
850        match self.scope {
851            NativeSessionCatalogScope::Workspace => owner == self.selected.as_ref(),
852            NativeSessionCatalogScope::Unregistered => owner.is_none(),
853        }
854    }
855
856    fn external_group(
857        &self,
858        indexed: &IndexedNativeSession,
859        group_ids: &BTreeMap<NativeExternalGroupKey, String>,
860    ) -> Option<NativeSessionExternalGroup> {
861        matches!(self.scope, NativeSessionCatalogScope::Unregistered).then(|| {
862            let key = NativeExternalGroupKey::from_cwd(indexed.cwd_key.as_ref());
863            NativeSessionExternalGroup {
864                group_id: group_ids
865                    .get(&key)
866                    .cloned()
867                    .expect("owned unregistered session must have an opaque group token"),
868                kind: match key {
869                    NativeExternalGroupKey::Project(_) => NativeSessionExternalGroupKind::Project,
870                    NativeExternalGroupKey::Global => NativeSessionExternalGroupKind::Global,
871                },
872                display_name: indexed.project_label.clone(),
873            }
874        })
875    }
876
877    fn external_group_ids(
878        &self,
879        sessions: &[IndexedNativeSession],
880    ) -> BTreeMap<NativeExternalGroupKey, String> {
881        if !matches!(self.scope, NativeSessionCatalogScope::Unregistered) {
882            return BTreeMap::new();
883        }
884        let mut paths = sessions
885            .iter()
886            .filter(|session| self.owns(&session.cwd_key))
887            .map(|session| NativeExternalGroupKey::from_cwd(session.cwd_key.as_ref()))
888            .collect::<Vec<_>>();
889        paths.sort();
890        paths.dedup();
891        paths
892            .into_iter()
893            .enumerate()
894            .map(|(index, path)| (path, format!("external-{:04}", index + 1)))
895            .collect()
896    }
897
898    fn fingerprint(&self) -> u64 {
899        let mut hash = 0xcbf2_9ce4_8422_2325u64;
900        hash = native_hash_bytes(
901            hash,
902            &[match self.scope {
903                NativeSessionCatalogScope::Workspace => 1,
904                NativeSessionCatalogScope::Unregistered => 2,
905            }],
906        );
907        if let Some(selected) = self.selected.as_ref() {
908            hash = native_hash_path(hash, selected);
909        }
910        for root in &self.registered {
911            hash = native_hash_bytes(hash, &[0xff]);
912            hash = native_hash_path(hash, root);
913        }
914        hash.max(1)
915    }
916}
917
918#[derive(Clone, Debug, Eq, Ord, PartialEq, PartialOrd)]
919enum NativeExternalGroupKey {
920    Global,
921    Project(NativePathKey),
922}
923
924impl NativeExternalGroupKey {
925    fn from_cwd(cwd: Option<&NativePathKey>) -> Self {
926        match cwd {
927            Some(cwd) => Self::Project(cwd.clone()),
928            None => Self::Global,
929        }
930    }
931}
932
933#[derive(Clone, Debug, Eq, Ord, PartialEq, PartialOrd)]
934struct NativePathKey {
935    components: Vec<String>,
936}
937
938impl NativePathKey {
939    fn from_canonical_root(path: &Path) -> Option<Self> {
940        Self::from_absolute_path(path)
941    }
942
943    fn from_declared_cwd(value: &str) -> Option<Self> {
944        let path = Path::new(value);
945        if !path.is_absolute()
946            || path
947                .components()
948                .any(|component| matches!(component, Component::ParentDir))
949        {
950            return None;
951        }
952        if let Ok(canonical) = std::fs::canonicalize(path) {
953            return Self::from_absolute_path(&canonical);
954        }
955        let mut existing = path;
956        let mut missing = Vec::<OsString>::new();
957        while !existing.exists() {
958            missing.push(existing.file_name()?.to_os_string());
959            existing = existing.parent()?;
960        }
961        let mut resolved = std::fs::canonicalize(existing).ok()?;
962        for component in missing.into_iter().rev() {
963            resolved.push(component);
964        }
965        Self::from_absolute_path(&resolved)
966    }
967
968    fn from_absolute_path(path: &Path) -> Option<Self> {
969        if !path.is_absolute() {
970            return None;
971        }
972        let mut components = Vec::new();
973        for component in path.components() {
974            match component {
975                Component::Prefix(prefix) => {
976                    components.push(normalize_native_path_component(
977                        &prefix.as_os_str().to_string_lossy(),
978                    ));
979                }
980                Component::RootDir => {}
981                Component::CurDir => {}
982                Component::ParentDir => {
983                    if components.len() <= 1 {
984                        return None;
985                    }
986                    components.pop();
987                }
988                Component::Normal(value) => components.push(normalize_native_path_component(
989                    &value.to_string_lossy(),
990                )),
991            }
992        }
993        (!components.is_empty()).then_some(Self { components })
994    }
995
996    fn starts_with(&self, root: &Self) -> bool {
997        self.components.starts_with(&root.components)
998    }
999}
1000
1001fn native_external_group_label(cwd: &str) -> String {
1002    let raw = Path::new(cwd)
1003        .file_name()
1004        .map(|name| name.to_string_lossy().into_owned())
1005        .unwrap_or_else(|| "project".to_owned());
1006    let sanitized = raw
1007        .chars()
1008        .filter(|character| !character.is_control())
1009        .collect::<String>();
1010    let sanitized = sanitized.trim();
1011    let sanitized = if sanitized.is_empty() { "project" } else { sanitized };
1012    let mut end = sanitized
1013        .len()
1014        .min(gate4agent_types::NATIVE_SESSION_EXTERNAL_GROUP_LABEL_MAX_BYTES);
1015    while !sanitized.is_char_boundary(end) {
1016        end -= 1;
1017    }
1018    sanitized[..end].to_owned()
1019}
1020
1021fn native_hash_path(mut hash: u64, path: &NativePathKey) -> u64 {
1022    for component in &path.components {
1023        hash = native_hash_bytes(hash, &(component.len() as u64).to_le_bytes());
1024        hash = native_hash_bytes(hash, component.as_bytes());
1025    }
1026    hash
1027}
1028
1029fn native_hash_bytes(mut hash: u64, bytes: &[u8]) -> u64 {
1030    for byte in bytes {
1031        hash ^= u64::from(*byte);
1032        hash = hash.wrapping_mul(0x0000_0100_0000_01b3);
1033    }
1034    hash
1035}
1036
1037#[cfg(windows)]
1038fn normalize_native_path_component(value: &str) -> String {
1039    if let Some(unc) = value.strip_prefix(r"\\?\UNC\") {
1040        format!(r"\\{unc}").to_lowercase()
1041    } else {
1042        value.trim_start_matches(r"\\?\").to_lowercase()
1043    }
1044}
1045
1046#[cfg(not(windows))]
1047fn normalize_native_path_component(value: &str) -> String {
1048    value.to_owned()
1049}
1050
1051#[derive(Clone, Copy)]
1052enum NativeSessionPreviewSelector<'a> {
1053    Candidate(&'a str),
1054    Session(&'a str),
1055}
1056
1057fn normalize_preview_text(text: &str) -> String {
1058    let normalized = text.replace("\r\n", "\n").replace('\r', "\n");
1059    let mut end = normalized.len().min(gate4agent_types::NATIVE_SESSION_PREVIEW_TEXT_MAX_BYTES);
1060    while !normalized.is_char_boundary(end) {
1061        end -= 1;
1062    }
1063    normalized[..end]
1064        .chars()
1065        .filter(|character| {
1066            !character.is_control() || matches!(character, '\n' | '\t')
1067        })
1068        .collect()
1069}
1070
1071fn project_native_session_preview(
1072    session: gate4agent_adapters::HistorySession,
1073    modified_at_unix_ms: Option<u64>,
1074    source_truncated: bool,
1075    message_limit: u16,
1076) -> Result<NativeSessionPreview, NativeSessionPreviewError> {
1077    let skip = session
1078        .messages
1079        .len()
1080        .saturating_sub(usize::from(message_limit));
1081    let messages: Vec<NativeSessionPreviewMessage> = session
1082        .messages
1083        .into_iter()
1084        .skip(skip)
1085        .map(|message| NativeSessionPreviewMessage {
1086            role: match message.role {
1087                gate4agent_adapters::HistoryRole::User => HistoryMessageRole::User,
1088                gate4agent_adapters::HistoryRole::Assistant => HistoryMessageRole::Assistant,
1089            },
1090            text: normalize_preview_text(&message.text),
1091        })
1092        .collect();
1093    let truncated = source_truncated || session.message_count > messages.len() as u64;
1094    let preview = NativeSessionPreview {
1095        session_id: session.session_id,
1096        title: session.title,
1097        modified_at_unix_ms,
1098        model: session.model,
1099        message_count: session.message_count,
1100        message_count_exact: !source_truncated,
1101        completed_turn_count: session.completed_turn_count,
1102        total_tokens: session
1103            .total_tokens_observed
1104            .then_some(session.total_tokens),
1105        truncated,
1106        messages,
1107    };
1108    preview
1109        .validate()
1110        .map_err(|_| NativeSessionPreviewError::PreviewUnavailable)?;
1111    Ok(preview)
1112}
1113
1114fn native_catalog_window_for_modified(
1115    modified_at_unix_ms: Option<u64>,
1116    recent_cutoff: u64,
1117) -> NativeSessionCatalogWindow {
1118    if modified_at_unix_ms.is_some_and(|modified| modified >= recent_cutoff) {
1119        NativeSessionCatalogWindow::Recent
1120    } else {
1121        NativeSessionCatalogWindow::Older
1122    }
1123}
1124
1125fn native_catalog_snapshot_revision(index_revision: u64, cutoff_unix_ms: u64) -> u64 {
1126    let mut hash = 0xcbf2_9ce4_8422_2325u64;
1127    for byte in index_revision
1128        .to_le_bytes()
1129        .into_iter()
1130        .chain(cutoff_unix_ms.to_le_bytes())
1131    {
1132        hash ^= u64::from(byte);
1133        hash = hash.wrapping_mul(0x0000_0100_0000_01b3);
1134    }
1135    hash.max(1)
1136}
1137
1138fn scoped_native_catalog_snapshot_revision(
1139    index_revision: u64,
1140    cutoff_unix_ms: u64,
1141    ownership_fingerprint: u64,
1142) -> u64 {
1143    native_hash_bytes(
1144        native_catalog_snapshot_revision(index_revision, cutoff_unix_ms),
1145        &ownership_fingerprint.to_le_bytes(),
1146    )
1147    .max(1)
1148}
1149
1150fn unscoped_native_catalog_page(
1151    page: ScopedNativeSessionCatalogPage,
1152) -> NativeSessionCatalogPage {
1153    NativeSessionCatalogPage {
1154        revision: page.revision,
1155        cutoff_unix_ms: page.cutoff_unix_ms,
1156        window: page.window,
1157        entries: page
1158            .entries
1159            .into_iter()
1160            .map(|entry| entry.metadata)
1161            .collect(),
1162        recent_total_count: page.recent_total_count,
1163        older_total_count: page.older_total_count,
1164        next_after_selection_id: page.next_after_selection_id,
1165        remaining_count: page.remaining_count,
1166    }
1167}
1168
1169#[derive(Clone, Copy, Debug, Eq, PartialEq)]
1170enum NativeSessionIndexError {
1171    UnsupportedProvider,
1172    Unavailable,
1173}
1174
1175fn map_catalog_index_error(error: NativeSessionIndexError) -> NativeSessionCatalogError {
1176    match error {
1177        NativeSessionIndexError::UnsupportedProvider => {
1178            NativeSessionCatalogError::UnsupportedProvider
1179        }
1180        NativeSessionIndexError::Unavailable => NativeSessionCatalogError::CatalogUnavailable,
1181    }
1182}
1183
1184#[derive(Clone, Copy, Debug, Error, Eq, PartialEq)]
1185pub enum NativeSessionCatalogError {
1186    #[error("native session catalog limit is outside the supported bounded range")]
1187    InvalidLimit,
1188    #[error("native session catalog provider is unsupported")]
1189    UnsupportedProvider,
1190    #[error("native session catalog workspace is unavailable")]
1191    WorkspaceUnavailable,
1192    #[error("native session catalog is unavailable")]
1193    CatalogUnavailable,
1194    #[error("native session catalog revision or cursor is stale")]
1195    StaleCatalog,
1196}
1197
1198#[derive(Clone, Copy, Debug, Error, Eq, PartialEq)]
1199pub enum NativeSessionPreviewError {
1200    #[error("native session preview limit is outside the supported bounded range")]
1201    InvalidLimit,
1202    #[error("native session preview provider is unsupported")]
1203    UnsupportedProvider,
1204    #[error("native session preview workspace is unavailable")]
1205    WorkspaceUnavailable,
1206    #[error("native session preview is unavailable")]
1207    PreviewUnavailable,
1208    #[error("native session preview was not found")]
1209    SessionNotFound,
1210    #[error("native session catalog revision is stale")]
1211    StaleCatalog,
1212    #[error("native session preview is ambiguous")]
1213    AmbiguousSession,
1214}
1215
1216fn drain_queue<T>(queue: &mut VecDeque<T>, limit: usize) -> Vec<T> {
1217    let count = limit.min(queue.len());
1218    queue.drain(..count).collect()
1219}
1220
1221#[derive(Clone, Copy, Debug, Eq, PartialEq)]
1222pub struct NativeRuntimeConfig {
1223    pub command_capacity: usize,
1224    pub max_commands_per_tick: usize,
1225    pub effect_capacity_per_session: usize,
1226    pub observation_capacity: usize,
1227    pub max_observations_per_tick: usize,
1228    pub provider_stop_grace_ms: u64,
1229    pub provider_shutdown_timeout_ms: u64,
1230    pub worker_poll_interval_ms: u64,
1231    pub worker_idle_timeout_ms: u64,
1232}
1233
1234impl Default for NativeRuntimeConfig {
1235    fn default() -> Self {
1236        Self {
1237            command_capacity: 256,
1238            max_commands_per_tick: 64,
1239            effect_capacity_per_session: 16,
1240            observation_capacity: 1_024,
1241            max_observations_per_tick: 256,
1242            provider_stop_grace_ms: 5_000,
1243            provider_shutdown_timeout_ms: 15_000,
1244            worker_poll_interval_ms: 20,
1245            worker_idle_timeout_ms: 60_000,
1246        }
1247    }
1248}
1249
1250#[derive(Clone, Debug, Eq, PartialEq)]
1251pub struct NativeRuntimeTick {
1252    pub command_outcomes: Vec<CommandOutcome>,
1253    pub effects_dispatched: usize,
1254    pub observations_applied: usize,
1255    pub terminal_frames_collected: usize,
1256    pub snapshot_revision: u64,
1257    pub publish_report: PublishReport,
1258}
1259
1260/// Owns the kernel tick and all native effect workers. Product apps retain
1261/// only the bounded handle returned by [`NativeRuntime::new`].
1262pub struct NativeRuntime {
1263    config: NativeRuntimeConfig,
1264    kernel: Gate4AgentKernel,
1265    tool_authority: ToolAuthorityHandle,
1266    tool_providers: BTreeMap<ToolProviderId, CapabilityProviderDescriptor>,
1267    provider_supervisors: BTreeMap<ToolProviderId, ProviderSupervisor>,
1268    port: ControlPlaneKernelPort,
1269    provider_exit_acks: VecDeque<PhysicalExitAck>,
1270    provider_faults: VecDeque<ProviderSupervisorFault>,
1271    provider_ack_cursor: Option<ToolProviderId>,
1272    provider_fault_cursor: Option<ToolProviderId>,
1273    effects: NativeEffectDispatcher,
1274    /// Per-phase timing distributions for [`Self::tick`] -- see
1275    /// `tick_profile`'s own doc comment for why this is a fixed-window
1276    /// ring, not an average, and why it stays on unconditionally.
1277    tick_profile: tick_profile::TickPhaseProfiler,
1278}
1279
1280impl NativeRuntime {
1281    pub fn new(catalog: AgentRegistry, config: NativeRuntimeConfig) -> (Gate4AgentHandle, Self) {
1282        Self::new_with_optional_history(catalog, config, None)
1283    }
1284
1285    pub fn new_with_history(
1286        catalog: AgentRegistry,
1287        config: NativeRuntimeConfig,
1288        history: NativeHistoryConfig,
1289    ) -> (Gate4AgentHandle, Self) {
1290        Self::new_with_optional_history(catalog, config, Some(history))
1291    }
1292
1293    pub fn new_with_tool_providers(
1294        catalog: AgentRegistry,
1295        config: NativeRuntimeConfig,
1296        providers: impl IntoIterator<Item = CapabilityProviderDescriptor>,
1297    ) -> Result<(Gate4AgentHandle, Self), ToolEngineError> {
1298        let providers = providers.into_iter().collect::<Vec<_>>();
1299        let kernel = Gate4AgentKernel::with_tool_providers(catalog.clone(), providers.clone())?;
1300        Ok(Self::new_with_kernel_and_optional_history(
1301            catalog, config, None, kernel, providers,
1302        ))
1303    }
1304
1305    pub fn new_with_history_and_tool_providers(
1306        catalog: AgentRegistry,
1307        config: NativeRuntimeConfig,
1308        history: NativeHistoryConfig,
1309        providers: impl IntoIterator<Item = CapabilityProviderDescriptor>,
1310    ) -> Result<(Gate4AgentHandle, Self), ToolEngineError> {
1311        let providers = providers.into_iter().collect::<Vec<_>>();
1312        let kernel = Gate4AgentKernel::with_tool_providers(catalog.clone(), providers.clone())?;
1313        Ok(Self::new_with_kernel_and_optional_history(
1314            catalog,
1315            config,
1316            Some(history),
1317            kernel,
1318            providers,
1319        ))
1320    }
1321
1322    fn new_with_optional_history(
1323        catalog: AgentRegistry,
1324        config: NativeRuntimeConfig,
1325        history: Option<NativeHistoryConfig>,
1326    ) -> (Gate4AgentHandle, Self) {
1327        let kernel = Gate4AgentKernel::new(catalog.clone());
1328        Self::new_with_kernel_and_optional_history(catalog, config, history, kernel, Vec::new())
1329    }
1330
1331    fn new_with_kernel_and_optional_history(
1332        catalog: AgentRegistry,
1333        config: NativeRuntimeConfig,
1334        history: Option<NativeHistoryConfig>,
1335        kernel: Gate4AgentKernel,
1336        providers: Vec<CapabilityProviderDescriptor>,
1337    ) -> (Gate4AgentHandle, Self) {
1338        let (handle, tool_authority, port) = bounded_control_plane(config.command_capacity);
1339        let runtime = Self {
1340            config,
1341            kernel,
1342            tool_authority,
1343            tool_providers: providers
1344                .into_iter()
1345                .map(|descriptor| (descriptor.id.clone(), descriptor))
1346                .collect(),
1347            provider_supervisors: BTreeMap::new(),
1348            port,
1349            provider_exit_acks: VecDeque::new(),
1350            provider_faults: VecDeque::new(),
1351            provider_ack_cursor: None,
1352            provider_fault_cursor: None,
1353            effects: NativeEffectDispatcher::new(catalog, config, history),
1354            tick_profile: tick_profile::TickPhaseProfiler::default(),
1355        };
1356        (handle, runtime)
1357    }
1358
1359    pub fn tool_authority(&self) -> ToolAuthorityHandle {
1360        self.tool_authority.clone()
1361    }
1362
1363    pub fn install_native_provider(
1364        &mut self,
1365        provider_id: &ToolProviderId,
1366        work_capacity: usize,
1367        executor: Box<dyn NativeProviderExecutor>,
1368    ) -> Result<ProviderBindingId, NativeProviderControlError> {
1369        self.collect_provider_supervisor_events();
1370        if let Some(existing) = self.provider_supervisors.get(provider_id) {
1371            let snapshot = existing.snapshot();
1372            if snapshot.state != ProviderSupervisorState::Closed {
1373                return Err(NativeProviderControlError::AlreadyInstalled {
1374                    state: snapshot.state,
1375                });
1376            }
1377            if snapshot.buffered_exit_acks != 0 || snapshot.buffered_faults != 0 {
1378                return Err(NativeProviderControlError::PendingSupervisorEvents);
1379            }
1380        }
1381        self.provider_supervisors.remove(provider_id);
1382
1383        let descriptor = self
1384            .tool_providers
1385            .get(provider_id)
1386            .cloned()
1387            .ok_or_else(|| NativeProviderControlError::UnknownProvider {
1388                provider_id: provider_id.clone(),
1389            })?;
1390        if !matches!(&descriptor.owner, CapabilityOwner::Gate) {
1391            return Err(NativeProviderControlError::UnsupportedOwner {
1392                provider_id: provider_id.clone(),
1393            });
1394        }
1395        let runtime = self
1396            .port
1397            .provider_authority()
1398            .bind_provider(provider_id.clone(), work_capacity)?;
1399        let binding_id = runtime.binding_id();
1400        let supervisor = ProviderSupervisor::new_with_stop_grace(
1401            descriptor,
1402            runtime,
1403            executor,
1404            Duration::from_millis(self.config.provider_stop_grace_ms.max(1)),
1405        )
1406        .map_err(NativeProviderControlError::Build)?;
1407        self.provider_supervisors
1408            .insert(provider_id.clone(), supervisor);
1409        Ok(binding_id)
1410    }
1411
1412    pub fn retire_native_provider(
1413        &mut self,
1414        provider_id: &ToolProviderId,
1415    ) -> Result<(), NativeProviderControlError> {
1416        let supervisor = self
1417            .provider_supervisors
1418            .get_mut(provider_id)
1419            .ok_or_else(|| NativeProviderControlError::NotInstalled {
1420                provider_id: provider_id.clone(),
1421            })?;
1422        supervisor.begin_retirement()?;
1423        Ok(())
1424    }
1425
1426    /// Begins non-blocking retirement for every installed native provider.
1427    ///
1428    /// The owner must keep calling [`NativeRuntime::tick`] until
1429    /// [`NativeRuntime::native_provider_shutdown_complete`] returns `true`
1430    /// before dropping the runtime. Dropping the runtime does not acknowledge
1431    /// physical provider exit or complete coordinated shutdown.
1432    pub fn retire_all_native_providers(&mut self) -> Result<(), NativeProviderControlError> {
1433        let mut first_error = None;
1434        for supervisor in self.provider_supervisors.values_mut() {
1435            if let Err(error) = supervisor.begin_retirement() {
1436                if first_error.is_none() {
1437                    first_error = Some(error);
1438                }
1439            }
1440        }
1441        match first_error {
1442            Some(error) => Err(error.into()),
1443            None => Ok(()),
1444        }
1445    }
1446
1447    /// Returns `true` only after every installed provider supervisor is closed.
1448    pub fn native_provider_shutdown_complete(&self) -> bool {
1449        self.provider_supervisors
1450            .values()
1451            .all(|supervisor| supervisor.state() == ProviderSupervisorState::Closed)
1452    }
1453
1454    /// Drives the canonical runtime path until every native provider has
1455    /// detached after physical teardown, or returns the still-owned snapshots
1456    /// at the configured shutdown deadline. Lifecycle events remain buffered
1457    /// for explicit operator drain after this method returns.
1458    pub async fn shutdown_native_providers(
1459        &mut self,
1460    ) -> Result<(), NativeProviderShutdownError> {
1461        let retirement_error = self.retire_all_native_providers().err();
1462        let shutdown_timeout_ms = self
1463            .config
1464            .provider_shutdown_timeout_ms
1465            .clamp(1, MAX_PROVIDER_SHUTDOWN_TIMEOUT_MS);
1466        let deadline = Instant::now() + Duration::from_millis(shutdown_timeout_ms);
1467        while !self.native_provider_shutdown_complete() {
1468            self.tick().await;
1469            if self.native_provider_shutdown_complete() {
1470                break;
1471            }
1472            let now = Instant::now();
1473            if now >= deadline {
1474                return Err(NativeProviderShutdownError::TimedOut {
1475                    pending: self
1476                        .provider_supervisors
1477                        .values()
1478                        .filter(|supervisor| {
1479                            supervisor.state() != ProviderSupervisorState::Closed
1480                        })
1481                        .map(ProviderSupervisor::snapshot)
1482                        .collect(),
1483                });
1484            }
1485            let poll_interval =
1486                Duration::from_millis(self.config.worker_poll_interval_ms.max(1));
1487            tokio::time::sleep(poll_interval.min(deadline.duration_since(now))).await;
1488        }
1489        match retirement_error {
1490            Some(error) => Err(NativeProviderShutdownError::Control(error)),
1491            None => Ok(()),
1492        }
1493    }
1494
1495    pub fn native_provider_snapshot(
1496        &self,
1497        provider_id: &ToolProviderId,
1498    ) -> Option<ProviderSupervisorSnapshot> {
1499        self.provider_supervisors
1500            .get(provider_id)
1501            .map(ProviderSupervisor::snapshot)
1502    }
1503
1504    pub fn drain_provider_exit_acks(&mut self, limit: usize) -> Vec<PhysicalExitAck> {
1505        self.collect_provider_supervisor_events();
1506        drain_queue(&mut self.provider_exit_acks, limit)
1507    }
1508
1509    pub fn drain_provider_faults(&mut self, limit: usize) -> Vec<ProviderSupervisorFault> {
1510        self.collect_provider_supervisor_events();
1511        drain_queue(&mut self.provider_faults, limit)
1512    }
1513
1514    pub fn history_enabled(&self) -> bool {
1515        self.effects.history_config.is_some()
1516    }
1517
1518    /// Run one non-blocking host tick. Effects are dispatched in session order
1519    /// to per-instance workers; their observations enter later ticks.
1520    pub async fn tick(&mut self) -> NativeRuntimeTick {
1521        let drain_observations_start = Instant::now();
1522        let (observations, terminal_frames_collected) = self
1523            .effects
1524            .drain_observations(self.config.max_observations_per_tick.max(1));
1525        let observations_applied = observations.len();
1526        self.tick_profile
1527            .record_drain_observations(drain_observations_start.elapsed());
1528
1529        let drain_ingress_start = Instant::now();
1530        let ingress = self
1531            .port
1532            .drain_ingress(self.config.max_commands_per_tick.max(1));
1533        self.tick_profile
1534            .record_drain_ingress(drain_ingress_start.elapsed());
1535
1536        let step_control_plane_start = Instant::now();
1537        let step = self.kernel.step_control_plane(ingress, observations);
1538        let effects_dispatched = step.effects.len();
1539        self.tick_profile
1540            .record_step_control_plane(step_control_plane_start.elapsed());
1541
1542        let dispatch_effects_start = Instant::now();
1543        for effect in step.effects.iter().cloned() {
1544            self.effects.dispatch(effect);
1545        }
1546        self.tick_profile
1547            .record_dispatch_effects(dispatch_effects_start.elapsed());
1548
1549        let snapshot_revision = step.snapshot.revision;
1550        let publish_step_start = Instant::now();
1551        let publish_report = self.port.publish_step(&step).control_events;
1552        self.tick_profile
1553            .record_publish_step(publish_step_start.elapsed());
1554
1555        let provider_supervisors_start = Instant::now();
1556        for supervisor in self.provider_supervisors.values_mut() {
1557            supervisor.tick();
1558        }
1559        self.collect_provider_supervisor_events();
1560        self.tick_profile
1561            .record_provider_supervisors(provider_supervisors_start.elapsed());
1562
1563        NativeRuntimeTick {
1564            command_outcomes: step.command_outcomes,
1565            effects_dispatched,
1566            observations_applied,
1567            terminal_frames_collected,
1568            snapshot_revision,
1569            publish_report,
1570        }
1571    }
1572
1573    /// Snapshot of the last [`tick_profile::SAMPLE_WINDOW`] ticks' own
1574    /// per-phase timing distributions -- the number `/metrics` reads to
1575    /// name which piece of an otherwise-idle [`Self::tick`] is spending the
1576    /// CPU.
1577    pub fn tick_profile_snapshot(&self) -> tick_profile::TickProfileSnapshot {
1578        self.tick_profile.snapshot()
1579    }
1580
1581    /// Snapshot of every shell-efficiency series folded in from the
1582    /// per-instance worker loops -- see `shell_efficiency`'s own doc
1583    /// comment. Locked only for this read, mirroring `tick_profile_snapshot`.
1584    pub fn shell_efficiency_snapshot(&self) -> shell_efficiency::ShellEfficiencyProfileSnapshot {
1585        self.effects
1586            .shell_efficiency
1587            .lock()
1588            .unwrap_or_else(|poisoned| poisoned.into_inner())
1589            .snapshot()
1590    }
1591
1592    /// The live profile, for a consumer that wants to decide WHEN to pay
1593    /// for a snapshot.
1594    ///
1595    /// `snapshot()` sorts every ring it reads to produce percentiles, and
1596    /// both this module and `tick_profile` state that the sort belongs on
1597    /// the read rather than on the tick. Handing out the `Arc` lets `GET
1598    /// /metrics` pay it per request instead of the node's drive loop paying
1599    /// it ninety times a second to keep a copy nobody has asked for warm.
1600    pub fn shell_efficiency_profile(&self) -> Arc<Mutex<ShellEfficiencyProfile>> {
1601        Arc::clone(&self.effects.shell_efficiency)
1602    }
1603
1604    fn collect_provider_supervisor_events(&mut self) {
1605        self.collect_provider_exit_acks();
1606        self.collect_provider_faults();
1607    }
1608
1609    fn collect_provider_exit_acks(&mut self) {
1610        let remaining =
1611            MAX_PROVIDER_SUPERVISOR_EVENTS.saturating_sub(self.provider_exit_acks.len());
1612        if remaining == 0 {
1613            return;
1614        }
1615        let provider_ids = provider_ids_after(
1616            &self.provider_supervisors,
1617            self.provider_ack_cursor.as_ref(),
1618        );
1619        let mut remaining = remaining;
1620        for provider_id in provider_ids {
1621            let limit = remaining.min(PROVIDER_EVENT_DRAIN_QUANTUM);
1622            let drained = self
1623                .provider_supervisors
1624                .get_mut(&provider_id)
1625                .map_or_else(Vec::new, |supervisor| supervisor.drain_exit_acks(limit));
1626            if !drained.is_empty() {
1627                remaining -= drained.len();
1628                self.provider_exit_acks.extend(drained);
1629                self.provider_ack_cursor = Some(provider_id);
1630            }
1631            if remaining == 0 {
1632                break;
1633            }
1634        }
1635    }
1636
1637    fn collect_provider_faults(&mut self) {
1638        let remaining =
1639            MAX_PROVIDER_SUPERVISOR_EVENTS.saturating_sub(self.provider_faults.len());
1640        if remaining == 0 {
1641            return;
1642        }
1643        let provider_ids = provider_ids_after(
1644            &self.provider_supervisors,
1645            self.provider_fault_cursor.as_ref(),
1646        );
1647        let mut remaining = remaining;
1648        for provider_id in provider_ids {
1649            let limit = remaining.min(PROVIDER_EVENT_DRAIN_QUANTUM);
1650            let drained = self
1651                .provider_supervisors
1652                .get_mut(&provider_id)
1653                .map_or_else(Vec::new, |supervisor| supervisor.drain_faults(limit));
1654            if !drained.is_empty() {
1655                remaining -= drained.len();
1656                self.provider_faults.extend(drained);
1657                self.provider_fault_cursor = Some(provider_id);
1658            }
1659            if remaining == 0 {
1660                break;
1661            }
1662        }
1663    }
1664
1665    pub fn active_native_sessions(&self) -> usize {
1666        self.effects.active_sessions.load(Ordering::Acquire)
1667    }
1668
1669    /// Installs or replaces a bounded host-only profile for future spawns.
1670    pub fn upsert_native_launch_profile(
1671        &mut self,
1672        profile: NativeLaunchProfile,
1673    ) -> Result<(), NativeLaunchProfileError> {
1674        self.effects.launch_profiles.upsert(profile)
1675    }
1676
1677    /// Removes an unselected host-only profile.
1678    pub fn remove_native_launch_profile(
1679        &mut self,
1680        profile_id: &NativeLaunchProfileId,
1681    ) -> Result<bool, NativeLaunchProfileError> {
1682        self.effects.launch_profiles.remove(profile_id)
1683    }
1684
1685    /// Selects a profile for future spawns of one exact instance.
1686    ///
1687    /// Spawn dispatch is the linearization point: a selection change does not
1688    /// alter a child that was already dispatched or started.
1689    pub fn select_native_launch_profile(
1690        &mut self,
1691        instance_id: AgentInstanceId,
1692        profile_id: NativeLaunchProfileId,
1693    ) -> Result<(), NativeLaunchProfileError> {
1694        self.effects
1695            .launch_profiles
1696            .select_native_launch_profile(instance_id, profile_id)
1697    }
1698
1699    /// Clears one instance selection for future spawns only.
1700    pub fn clear_native_launch_profile_selection(
1701        &mut self,
1702        instance_id: AgentInstanceId,
1703    ) -> bool {
1704        self.effects
1705            .launch_profiles
1706            .clear_native_launch_profile_selection(instance_id)
1707    }
1708
1709    /// Returns a clonable selector for future native launch-profile spawns.
1710    pub fn native_launch_profile_control(&self) -> NativeLaunchProfileControl {
1711        self.effects.launch_profiles.clone()
1712    }
1713
1714}
1715
1716fn provider_ids_after(
1717    supervisors: &BTreeMap<ToolProviderId, ProviderSupervisor>,
1718    cursor: Option<&ToolProviderId>,
1719) -> Vec<ToolProviderId> {
1720    let mut provider_ids = supervisors.keys().cloned().collect::<Vec<_>>();
1721    let Some(cursor) = cursor else {
1722        return provider_ids;
1723    };
1724    let start = provider_ids
1725        .iter()
1726        .position(|provider_id| provider_id > cursor)
1727        .unwrap_or(0);
1728    provider_ids.rotate_left(start);
1729    provider_ids
1730}
1731
1732#[derive(Debug, Error)]
1733pub enum NativeProviderControlError {
1734    #[error("tool provider '{provider_id}' is not registered in this native runtime")]
1735    UnknownProvider { provider_id: ToolProviderId },
1736    #[error("tool provider '{provider_id}' is not owned by gate4agent")]
1737    UnsupportedOwner { provider_id: ToolProviderId },
1738    #[error("native provider supervisor is already installed in state {state:?}")]
1739    AlreadyInstalled { state: ProviderSupervisorState },
1740    #[error("native provider supervisor is not installed for '{provider_id}'")]
1741    NotInstalled { provider_id: ToolProviderId },
1742    #[error("retired provider supervisor still has undelivered lifecycle events")]
1743    PendingSupervisorEvents,
1744    #[error("native provider runtime failed: {0}")]
1745    Runtime(#[from] ProviderRuntimeError),
1746    #[error("native provider supervisor build failed: {0:?}")]
1747    Build(ProviderSupervisorBuildError),
1748}
1749
1750#[derive(Debug, Error)]
1751pub enum NativeProviderShutdownError {
1752    #[error("native provider retirement reported an error: {0}")]
1753    Control(NativeProviderControlError),
1754    #[error("native provider shutdown timed out with physical owners retained")]
1755    TimedOut { pending: Vec<ProviderSupervisorSnapshot> },
1756}
1757
1758struct EffectWorker {
1759    sender: Sender<NativeEffectRequest>,
1760}
1761
1762struct NativeSpawnOverlay {
1763    environment: Vec<EnvMutation>,
1764    extra_args: Vec<OsString>,
1765    one_shot_session_persistence: OneShotSessionPersistence,
1766    mcp_server: Option<McpServerSpec>,
1767}
1768
1769impl Default for NativeSpawnOverlay {
1770    fn default() -> Self {
1771        Self {
1772            environment: Vec::new(),
1773            extra_args: Vec::new(),
1774            one_shot_session_persistence: OneShotSessionPersistence::Ephemeral,
1775            mcp_server: None,
1776        }
1777    }
1778}
1779
1780struct NativeEffectRequest {
1781    effect: EffectEnvelope,
1782    spawn_env: Vec<EnvMutation>,
1783    spawn_extra_args: Vec<OsString>,
1784    one_shot_session_persistence: OneShotSessionPersistence,
1785    spawn_mcp_server: Option<McpServerSpec>,
1786}
1787
1788#[derive(Clone)]
1789struct NativeWorkerContext {
1790    control_tx: Sender<ObservationEnvelope>,
1791    terminal_frames: Arc<Mutex<BTreeMap<TerminalFrameKey, ObservationEnvelope>>>,
1792    active_sessions: Arc<AtomicUsize>,
1793    /// Shared with every other instance's worker and with
1794    /// `NativeEffectDispatcher` -- see `shell_efficiency`'s own doc comment
1795    /// for why folding here, once per worker-loop iteration, is the only
1796    /// place a shell-native fact can become part of a distribution.
1797    shell_efficiency: Arc<Mutex<ShellEfficiencyProfile>>,
1798    poll_interval: Duration,
1799    idle_timeout: Duration,
1800}
1801
1802struct NativeEffectDispatcher {
1803    catalog: AgentRegistry,
1804    config: NativeRuntimeConfig,
1805    launch_profiles: NativeLaunchProfileControl,
1806    workers: HashMap<AgentInstanceId, EffectWorker>,
1807    authority_worker: Option<EffectWorker>,
1808    capability_worker: Option<EffectWorker>,
1809    history_config: Option<NativeHistoryConfig>,
1810    control_tx: Sender<ObservationEnvelope>,
1811    control_rx: Receiver<ObservationEnvelope>,
1812    pending_failures: VecDeque<ObservationEnvelope>,
1813    terminal_frames: Arc<Mutex<BTreeMap<TerminalFrameKey, ObservationEnvelope>>>,
1814    active_sessions: Arc<AtomicUsize>,
1815    shell_efficiency: Arc<Mutex<ShellEfficiencyProfile>>,
1816}
1817
1818impl NativeEffectDispatcher {
1819    fn new(
1820        catalog: AgentRegistry,
1821        config: NativeRuntimeConfig,
1822        history_config: Option<NativeHistoryConfig>,
1823    ) -> Self {
1824        let (control_tx, control_rx) = mpsc::channel(config.observation_capacity.max(1));
1825        Self {
1826            catalog,
1827            config,
1828            launch_profiles: NativeLaunchProfileControl::new(),
1829            workers: HashMap::new(),
1830            authority_worker: None,
1831            capability_worker: None,
1832            history_config,
1833            control_tx,
1834            control_rx,
1835            pending_failures: VecDeque::new(),
1836            terminal_frames: Arc::new(Mutex::new(BTreeMap::new())),
1837            active_sessions: Arc::new(AtomicUsize::new(0)),
1838            shell_efficiency: Arc::new(Mutex::new(ShellEfficiencyProfile::default())),
1839        }
1840    }
1841
1842    fn dispatch(&mut self, effect: EffectEnvelope) {
1843        if matches!(effect.effect, ControlEffect::ProbeCapabilities { .. }) {
1844            self.dispatch_capability(effect);
1845            return;
1846        }
1847        if matches!(
1848            effect.effect,
1849            ControlEffect::DiscoverHistory { .. }
1850                | ControlEffect::LoadHistory { .. }
1851                | ControlEffect::AuthorizeResume { .. }
1852        ) {
1853            self.dispatch_authority(effect);
1854            return;
1855        }
1856        self.workers.retain(|_, worker| !worker.sender.is_closed());
1857        if let Err(message) = validate_effect_runtime_policy(&effect.effect) {
1858            self.pending_failures
1859                .push_back(effect_failure(effect, message));
1860            return;
1861        }
1862        let instance_id = effect.instance_id;
1863        let spawn_overlay = match self.compose_spawn_overlay(&effect) {
1864            Ok(overlay) => overlay,
1865            Err(message) => {
1866                self.pending_failures
1867                    .push_back(effect_failure(effect, message));
1868                return;
1869            }
1870        };
1871        let mut pending = NativeEffectRequest {
1872            effect,
1873            spawn_env: spawn_overlay.environment,
1874            spawn_extra_args: spawn_overlay.extra_args,
1875            one_shot_session_persistence: spawn_overlay.one_shot_session_persistence,
1876            spawn_mcp_server: spawn_overlay.mcp_server,
1877        };
1878        for _ in 0..2 {
1879            let sender = self.worker_sender(instance_id);
1880            match sender.try_send(pending) {
1881                Ok(()) => return,
1882                Err(TrySendError::Closed(request)) => {
1883                    self.workers.remove(&instance_id);
1884                    pending = request;
1885                }
1886                Err(TrySendError::Full(request)) => {
1887                    self.pending_failures.push_back(effect_failure(
1888                        request.effect,
1889                        "native session effect queue is full".to_owned(),
1890                    ));
1891                    return;
1892                }
1893            }
1894        }
1895        self.pending_failures.push_back(effect_failure(
1896            pending.effect,
1897            "native session effect worker is unavailable".to_owned(),
1898        ));
1899    }
1900
1901    fn dispatch_capability(&mut self, effect: EffectEnvelope) {
1902        let mut pending = NativeEffectRequest {
1903            effect,
1904            spawn_env: Vec::new(),
1905            spawn_extra_args: Vec::new(),
1906            one_shot_session_persistence: OneShotSessionPersistence::Ephemeral,
1907            spawn_mcp_server: None,
1908        };
1909        for _ in 0..2 {
1910            let sender = self.capability_sender();
1911            match sender.try_send(pending) {
1912                Ok(()) => return,
1913                Err(TrySendError::Closed(request)) => {
1914                    self.capability_worker = None;
1915                    pending = request;
1916                }
1917                Err(TrySendError::Full(request)) => {
1918                    self.pending_failures.push_back(effect_failure(
1919                        request.effect,
1920                        "native capability effect queue is full".to_owned(),
1921                    ));
1922                    return;
1923                }
1924            }
1925        }
1926        self.pending_failures.push_back(effect_failure(
1927            pending.effect,
1928            "native capability effect worker is unavailable".to_owned(),
1929        ));
1930    }
1931
1932    fn capability_sender(&mut self) -> Sender<NativeEffectRequest> {
1933        if let Some(worker) = &self.capability_worker {
1934            if !worker.sender.is_closed() {
1935                return worker.sender.clone();
1936            }
1937        }
1938        let (sender, receiver) = mpsc::channel(self.config.effect_capacity_per_session.max(1));
1939        tokio::spawn(run_capability_worker(
1940            self.catalog.clone(),
1941            receiver,
1942            self.control_tx.clone(),
1943        ));
1944        self.capability_worker = Some(EffectWorker {
1945            sender: sender.clone(),
1946        });
1947        sender
1948    }
1949
1950    fn dispatch_authority(&mut self, effect: EffectEnvelope) {
1951        if self.history_config.is_none()
1952            && matches!(
1953                effect.effect,
1954                ControlEffect::DiscoverHistory { .. } | ControlEffect::LoadHistory { .. }
1955            )
1956        {
1957            self.pending_failures.push_back(effect_failure(
1958                effect,
1959                "native history authority is not configured".to_owned(),
1960            ));
1961            return;
1962        }
1963        let mut pending = NativeEffectRequest {
1964            effect,
1965            spawn_env: Vec::new(),
1966            spawn_extra_args: Vec::new(),
1967            one_shot_session_persistence: OneShotSessionPersistence::Ephemeral,
1968            spawn_mcp_server: None,
1969        };
1970        for _ in 0..2 {
1971            let sender = self.authority_sender();
1972            match sender.try_send(pending) {
1973                Ok(()) => return,
1974                Err(TrySendError::Closed(request)) => {
1975                    self.authority_worker = None;
1976                    pending = request;
1977                }
1978                Err(TrySendError::Full(request)) => {
1979                    self.pending_failures.push_back(effect_failure(
1980                        request.effect,
1981                        "native authority effect queue is full".to_owned(),
1982                    ));
1983                    return;
1984                }
1985            }
1986        }
1987        self.pending_failures.push_back(effect_failure(
1988            pending.effect,
1989            "native authority effect worker is unavailable".to_owned(),
1990        ));
1991    }
1992
1993    fn authority_sender(&mut self) -> Sender<NativeEffectRequest> {
1994        if let Some(worker) = &self.authority_worker {
1995            if !worker.sender.is_closed() {
1996                return worker.sender.clone();
1997            }
1998        }
1999        let (sender, receiver) = mpsc::channel(self.config.effect_capacity_per_session.max(1));
2000        tokio::spawn(run_authority_worker(
2001            self.catalog.clone(),
2002            self.history_config.clone(),
2003            receiver,
2004            self.control_tx.clone(),
2005        ));
2006        self.authority_worker = Some(EffectWorker {
2007            sender: sender.clone(),
2008        });
2009        sender
2010    }
2011
2012    fn compose_spawn_overlay(
2013        &self,
2014        effect: &EffectEnvelope,
2015    ) -> Result<NativeSpawnOverlay, String> {
2016        self.profile_spawn_overlay(effect)
2017    }
2018
2019    fn profile_spawn_overlay(&self, effect: &EffectEnvelope) -> Result<NativeSpawnOverlay, String> {
2020        let (agent_id, transport) = match &effect.effect {
2021            ControlEffect::Spawn {
2022                agent_id,
2023                transport,
2024                ..
2025            } => (agent_id, *transport),
2026            ControlEffect::SpawnResume {
2027                agent_id,
2028                transport,
2029                ..
2030            } => (agent_id, *transport),
2031            _ => return Ok(NativeSpawnOverlay::default()),
2032        };
2033        let pipe_binding_is_exact_one_shot = transport != TransportKind::Pipe
2034            || self.catalog.get(agent_id).is_some_and(|spec| {
2035                spec.capabilities.transports.pipe.as_ref().is_some_and(|pipe| {
2036                    pipe.protocol == PipeProtocol::OneShotText
2037                        && spec.capabilities.adapters.one_shot.as_ref() == Some(&pipe.adapter)
2038                })
2039            });
2040        self.launch_profiles
2041            .resolve_launch_overlay(
2042                effect.instance_id,
2043                agent_id,
2044                transport,
2045                pipe_binding_is_exact_one_shot,
2046            )
2047            .map(|overlay| NativeSpawnOverlay {
2048                environment: overlay.environment,
2049                extra_args: overlay.extra_args,
2050                one_shot_session_persistence: overlay.one_shot_session_persistence,
2051                mcp_server: overlay.mcp_server,
2052            })
2053            .map_err(|error| error.to_string())
2054    }
2055
2056    fn worker_sender(&mut self, instance_id: AgentInstanceId) -> Sender<NativeEffectRequest> {
2057        if let Some(worker) = self.workers.get(&instance_id) {
2058            return worker.sender.clone();
2059        }
2060        let (sender, receiver) = mpsc::channel(self.config.effect_capacity_per_session.max(1));
2061        tokio::spawn(run_effect_worker(
2062            self.catalog.clone(),
2063            receiver,
2064            NativeWorkerContext {
2065                control_tx: self.control_tx.clone(),
2066                terminal_frames: Arc::clone(&self.terminal_frames),
2067                active_sessions: Arc::clone(&self.active_sessions),
2068                shell_efficiency: Arc::clone(&self.shell_efficiency),
2069                poll_interval: Duration::from_millis(self.config.worker_poll_interval_ms.max(1)),
2070                idle_timeout: Duration::from_millis(self.config.worker_idle_timeout_ms.max(1)),
2071            },
2072        ));
2073        self.workers.insert(
2074            instance_id,
2075            EffectWorker {
2076                sender: sender.clone(),
2077            },
2078        );
2079        sender
2080    }
2081
2082    fn drain_observations(&mut self, limit: usize) -> (Vec<ObservationEnvelope>, usize) {
2083        self.workers.retain(|_, worker| !worker.sender.is_closed());
2084        let mut observations = Vec::new();
2085        while observations.len() < limit {
2086            if let Some(observation) = self.pending_failures.pop_front() {
2087                observations.push(observation);
2088                continue;
2089            }
2090            match self.control_rx.try_recv() {
2091                Ok(observation) => observations.push(observation),
2092                Err(mpsc::error::TryRecvError::Empty | mpsc::error::TryRecvError::Disconnected) => {
2093                    break;
2094                }
2095            }
2096        }
2097
2098        let remaining = limit.saturating_sub(observations.len());
2099        let mut frames = self
2100            .terminal_frames
2101            .lock()
2102            .unwrap_or_else(|poisoned| poisoned.into_inner());
2103        let keys: Vec<_> = frames.keys().copied().take(remaining).collect();
2104        let terminal_frames_collected = keys.len();
2105        for key in keys {
2106            if let Some(frame) = frames.remove(&key) {
2107                observations.push(frame);
2108            }
2109        }
2110        (observations, terminal_frames_collected)
2111    }
2112}
2113
2114const HISTORY_INSTANCE_DISCOVERIES_MAX: usize = 1_024;
2115
2116struct InstanceHistoryDiscovery {
2117    generation: SessionGeneration,
2118    agent_id: AgentId,
2119    request: HistoryDiscoveryRequest,
2120    candidates: HashMap<String, HistoryCandidate>,
2121}
2122
2123struct NativeAuthorityWorkerState {
2124    catalog: AgentRegistry,
2125    authority: Option<NativeHistoryAuthority>,
2126    discoveries: HashMap<AgentInstanceId, InstanceHistoryDiscovery>,
2127    discovery_order: VecDeque<AgentInstanceId>,
2128}
2129
2130impl NativeAuthorityWorkerState {
2131    fn new(catalog: AgentRegistry, config: Option<NativeHistoryConfig>) -> Self {
2132        Self {
2133            catalog,
2134            authority: config.map(NativeHistoryAuthority::new),
2135            discoveries: HashMap::new(),
2136            discovery_order: VecDeque::new(),
2137        }
2138    }
2139
2140    fn execute(&mut self, envelope: EffectEnvelope) -> ObservationEnvelope {
2141        let EffectEnvelope {
2142            operation_id,
2143            instance_id,
2144            generation,
2145            effect,
2146        } = envelope;
2147        let observation = match effect {
2148            ControlEffect::DiscoverHistory { agent_id, query } => {
2149                self.discover(instance_id, generation, agent_id, query)
2150            }
2151            ControlEffect::LoadHistory {
2152                agent_id,
2153                candidate_id,
2154            } => self.load(instance_id, generation, agent_id, candidate_id),
2155            ControlEffect::AuthorizeResume {
2156                agent_id,
2157                target,
2158                request,
2159            } => self.authorize_resume(instance_id, generation, agent_id, target, request),
2160            _ => history_failure("native authority worker received an invalid effect".to_owned()),
2161        };
2162        ObservationEnvelope {
2163            operation_id: Some(operation_id),
2164            instance_id,
2165            generation,
2166            observation,
2167        }
2168    }
2169
2170    fn discover(
2171        &mut self,
2172        instance_id: AgentInstanceId,
2173        generation: SessionGeneration,
2174        agent_id: AgentId,
2175        query: gate4agent_types::HistoryQuery,
2176    ) -> ControlObservation {
2177        let Some(spec) = self.catalog.get(&agent_id).cloned() else {
2178            return history_failure(format!("agent '{agent_id}' is absent from native catalog"));
2179        };
2180        let request =
2181            match HistoryDiscoveryRequest::from_spec(&spec, query.working_directory, query.limit) {
2182                Ok(request) => request,
2183                Err(error) => return history_failure(error.to_string()),
2184            };
2185        let Some(authority) = self.authority.as_mut() else {
2186            return history_failure("native history authority is not configured".to_owned());
2187        };
2188        let candidates = match discover_history(authority, &request) {
2189            Ok(candidates) => candidates,
2190            Err(error) => return history_failure(error.to_string()),
2191        };
2192        let summaries = candidates
2193            .iter()
2194            .map(|candidate| HistoryCandidateSummary {
2195                id: candidate.id().as_str().to_owned(),
2196                session_id_hint: candidate.session_id_hint().to_owned(),
2197                modified_at_unix_ms: candidate.modified_at_unix_ms(),
2198            })
2199            .collect::<Vec<_>>();
2200        if summaries
2201            .iter()
2202            .any(|candidate| candidate.validate().is_err())
2203        {
2204            return history_failure(
2205                "native history authority returned an invalid candidate".to_owned(),
2206            );
2207        }
2208        let candidates = candidates
2209            .into_iter()
2210            .map(|candidate| (candidate.id().as_str().to_owned(), candidate))
2211            .collect();
2212        self.discoveries.insert(
2213            instance_id,
2214            InstanceHistoryDiscovery {
2215                generation,
2216                agent_id,
2217                request,
2218                candidates,
2219            },
2220        );
2221        self.touch_discovery(instance_id);
2222        ControlObservation::HistoryDiscovered {
2223            candidates: summaries,
2224        }
2225    }
2226
2227    fn load(
2228        &mut self,
2229        instance_id: AgentInstanceId,
2230        generation: SessionGeneration,
2231        agent_id: AgentId,
2232        candidate_id: String,
2233    ) -> ControlObservation {
2234        let Some(discovery) = self.discoveries.get(&instance_id) else {
2235            return history_failure("history candidate discovery is expired".to_owned());
2236        };
2237        if discovery.generation != generation || discovery.agent_id != agent_id {
2238            return history_failure("history candidate generation is stale".to_owned());
2239        }
2240        let Some(candidate) = discovery.candidates.get(&candidate_id).cloned() else {
2241            return history_failure("history candidate is expired or unknown".to_owned());
2242        };
2243        let request = discovery.request.clone();
2244        self.touch_discovery(instance_id);
2245        let load = match HistoryLoadRequest::new(&request, candidate) {
2246            Ok(load) => load,
2247            Err(error) => return history_failure(error.to_string()),
2248        };
2249        let Some(authority) = self.authority.as_mut() else {
2250            return history_failure("native history authority is not configured".to_owned());
2251        };
2252        match load_history_session(authority, &load) {
2253            Ok(session) => {
2254                let session = history_session_record(session);
2255                if let Err(error) = session.validate() {
2256                    history_failure(error.to_string())
2257                } else {
2258                    ControlObservation::HistoryLoaded { session }
2259                }
2260            }
2261            Err(error) => history_failure(error.to_string()),
2262        }
2263    }
2264
2265    fn authorize_resume(
2266        &mut self,
2267        instance_id: AgentInstanceId,
2268        generation: SessionGeneration,
2269        agent_id: AgentId,
2270        target: ResumeAuthorityTarget,
2271        request: ResumeLaunchRequest,
2272    ) -> ControlObservation {
2273        if let Err(error) = request.validate() {
2274            return resume_failure(error.to_string());
2275        }
2276        let Some(spec) = self.catalog.get(&agent_id).cloned() else {
2277            return resume_failure(format!("agent '{agent_id}' is absent from native catalog"));
2278        };
2279
2280        let provider_session = match target {
2281            ResumeAuthorityTarget::ProviderSession { identity } => identity,
2282            ResumeAuthorityTarget::HistoryCandidate { candidate_id } => {
2283                let Some(discovery) = self.discoveries.get(&instance_id) else {
2284                    return resume_failure("history candidate discovery is expired".to_owned());
2285                };
2286                if discovery.generation != generation || discovery.agent_id != agent_id {
2287                    return resume_failure("history candidate generation is stale".to_owned());
2288                }
2289                let Some(candidate) = discovery.candidates.get(&candidate_id).cloned() else {
2290                    return resume_failure("history candidate is expired or unknown".to_owned());
2291                };
2292                let discovery_request = discovery.request.clone();
2293                self.touch_discovery(instance_id);
2294                let load = match HistoryLoadRequest::new(&discovery_request, candidate) {
2295                    Ok(load) => load,
2296                    Err(error) => return resume_failure(error.to_string()),
2297                };
2298                let Some(authority) = self.authority.as_mut() else {
2299                    return resume_failure("native history authority is not configured".to_owned());
2300                };
2301                let session = match load_history_session(authority, &load) {
2302                    Ok(session) => session,
2303                    Err(error) => return resume_failure(error.to_string()),
2304                };
2305                match authority.resume_provider_session(&load, session.session_id) {
2306                    Ok(identity) => identity,
2307                    Err(error) => return resume_failure(error.to_string()),
2308                }
2309            }
2310        };
2311        let resume_request = match ResumeRequest::from_provider_session(
2312            &spec,
2313            provider_session,
2314            Some(request.working_directory.clone()),
2315        ) {
2316            Ok(request) => request,
2317            Err(error) => return resume_failure(error.to_string()),
2318        };
2319        match prepare_resume(&mut ExplicitResumeAuthority, resume_request) {
2320            Ok(ResumeOutcome::Authorized(prepared)) => {
2321                if prepared.working_directory() != Some(request.working_directory.as_str()) {
2322                    return resume_failure(
2323                        "resume authority changed the requested working directory".to_owned(),
2324                    );
2325                }
2326                ControlObservation::ResumeAuthorized {
2327                    provider_session: prepared.provider_session().clone(),
2328                }
2329            }
2330            Ok(ResumeOutcome::Denied { reason }) => ControlObservation::ResumeDenied { reason },
2331            Err(error) => resume_failure(error.to_string()),
2332        }
2333    }
2334
2335    fn touch_discovery(&mut self, instance_id: AgentInstanceId) {
2336        self.discovery_order
2337            .retain(|candidate| *candidate != instance_id);
2338        self.discovery_order.push_back(instance_id);
2339        while self.discoveries.len() > HISTORY_INSTANCE_DISCOVERIES_MAX {
2340            if let Some(expired) = self.discovery_order.pop_front() {
2341                self.discoveries.remove(&expired);
2342            }
2343        }
2344    }
2345}
2346
2347async fn run_authority_worker(
2348    catalog: AgentRegistry,
2349    config: Option<NativeHistoryConfig>,
2350    mut effects: Receiver<NativeEffectRequest>,
2351    control_tx: Sender<ObservationEnvelope>,
2352) {
2353    let state = Arc::new(Mutex::new(NativeAuthorityWorkerState::new(catalog, config)));
2354    while let Some(request) = effects.recv().await {
2355        let fallback = request.effect.clone();
2356        let worker_state = Arc::clone(&state);
2357        let completion = match tokio::task::spawn_blocking(move || {
2358            worker_state
2359                .lock()
2360                .unwrap_or_else(|poisoned| poisoned.into_inner())
2361                .execute(request.effect)
2362        })
2363        .await
2364        {
2365            Ok(completion) => completion,
2366            Err(_) => effect_failure(fallback, "native authority worker task failed".to_owned()),
2367        };
2368        if control_tx.send(completion).await.is_err() {
2369            break;
2370        }
2371    }
2372}
2373
2374async fn run_capability_worker(
2375    catalog: AgentRegistry,
2376    mut effects: Receiver<NativeEffectRequest>,
2377    control_tx: Sender<ObservationEnvelope>,
2378) {
2379    let mut authority = NativeCapabilityProbeAuthority::default();
2380    while let Some(request) = effects.recv().await {
2381        let EffectEnvelope {
2382            operation_id,
2383            instance_id,
2384            generation,
2385            effect,
2386        } = request.effect;
2387        let observation = if let ControlEffect::ProbeCapabilities { agent_id, request } = effect {
2388            if request.validate().is_err() {
2389                ControlObservation::CapabilityProbeFailed {
2390                    failure: CapabilityProbeFailure::AuthorityRejected,
2391                }
2392            } else if let Some(spec) = catalog.get(&agent_id) {
2393                match authority.probe(spec, &request.working_directory).await {
2394                    Ok(session_option_models) => ControlObservation::CapabilitiesProbed {
2395                        session_option_models,
2396                    },
2397                    Err(failure) => ControlObservation::CapabilityProbeFailed { failure },
2398                }
2399            } else {
2400                ControlObservation::CapabilityProbeFailed {
2401                    failure: CapabilityProbeFailure::AuthorityRejected,
2402                }
2403            }
2404        } else {
2405            ControlObservation::CapabilityProbeFailed {
2406                failure: CapabilityProbeFailure::AuthorityRejected,
2407            }
2408        };
2409        if control_tx
2410            .send(ObservationEnvelope {
2411                operation_id: Some(operation_id),
2412                instance_id,
2413                generation,
2414                observation,
2415            })
2416            .await
2417            .is_err()
2418        {
2419            break;
2420        }
2421    }
2422}
2423
2424fn history_failure(message: String) -> ControlObservation {
2425    ControlObservation::HistoryFailed { message }
2426}
2427
2428fn resume_failure(message: String) -> ControlObservation {
2429    ControlObservation::ResumeFailed { message }
2430}
2431
2432struct ExplicitResumeAuthority;
2433
2434impl ResumeAuthority for ExplicitResumeAuthority {
2435    type Error = Infallible;
2436
2437    fn authorize(
2438        &mut self,
2439        _prepared: &PreparedResume,
2440    ) -> Result<ResumeAuthorityDecision, Self::Error> {
2441        Ok(ResumeAuthorityDecision::Authorized)
2442    }
2443}
2444
2445fn history_session_record(session: gate4agent_adapters::HistorySession) -> HistorySessionRecord {
2446    HistorySessionRecord {
2447        session_id: session.session_id,
2448        title: session.title,
2449        cwd: session.cwd,
2450        model: session.model,
2451        message_count: session.message_count,
2452        completed_turn_count: session.completed_turn_count,
2453        total_tokens: session.total_tokens,
2454        messages: session
2455            .messages
2456            .into_iter()
2457            .map(|message| HistoryMessageRecord {
2458                role: match message.role {
2459                    gate4agent_adapters::HistoryRole::User => HistoryMessageRole::User,
2460                    gate4agent_adapters::HistoryRole::Assistant => HistoryMessageRole::Assistant,
2461                },
2462                text: message.text,
2463            })
2464            .collect(),
2465    }
2466}
2467
2468async fn run_effect_worker(
2469    catalog: AgentRegistry,
2470    mut effects: Receiver<NativeEffectRequest>,
2471    context: NativeWorkerContext,
2472) {
2473    let mut shell = NativeEffectShell::new(catalog);
2474    let mut interval = tokio::time::interval(context.poll_interval);
2475    interval.set_missed_tick_behavior(MissedTickBehavior::Skip);
2476    let mut last_effect = Instant::now();
2477
2478    loop {
2479        tokio::select! {
2480            request = effects.recv() => {
2481                let Some(request) = request else {
2482                    break;
2483                };
2484                last_effect = Instant::now();
2485                let before = shell.active_session_count();
2486                let completion = shell
2487                    .execute_with_launch_context(
2488                        request.effect,
2489                        request.spawn_env,
2490                        request.spawn_extra_args,
2491                        request.one_shot_session_persistence,
2492                        request.spawn_mcp_server,
2493                    )
2494                    .await;
2495                update_active_count(&context.active_sessions, before, shell.active_session_count());
2496                if closes_terminal_session(&completion.observation) {
2497                    context
2498                        .terminal_frames
2499                        .lock()
2500                        .unwrap_or_else(|poisoned| poisoned.into_inner())
2501                        .remove(&(completion.instance_id, completion.generation));
2502                }
2503                if context.control_tx.send(completion).await.is_err() {
2504                    break;
2505                }
2506                if !publish_shell_observations(&mut shell, &context).await {
2507                    break;
2508                }
2509            }
2510            _ = interval.tick() => {
2511                if !publish_shell_observations(&mut shell, &context).await {
2512                    break;
2513                }
2514                if shell.active_session_count() == 0
2515                    && last_effect.elapsed() >= context.idle_timeout
2516                {
2517                    break;
2518                }
2519            }
2520        }
2521    }
2522
2523    let remaining = shell.active_session_count();
2524    if remaining > 0 {
2525        context
2526            .active_sessions
2527            .fetch_sub(remaining, Ordering::AcqRel);
2528    }
2529}
2530
2531async fn publish_shell_observations(
2532    shell: &mut NativeEffectShell,
2533    context: &NativeWorkerContext,
2534) -> bool {
2535    for observation in shell.collect_provider_events() {
2536        if context.control_tx.send(observation).await.is_err() {
2537            return false;
2538        }
2539    }
2540
2541    for observation in shell.collect_terminal_frames() {
2542        if matches!(
2543            &observation.observation,
2544            ControlObservation::TerminalFrame { .. }
2545        ) {
2546            let key = (observation.instance_id, observation.generation);
2547            context
2548                .terminal_frames
2549                .lock()
2550                .unwrap_or_else(|poisoned| poisoned.into_inner())
2551                .insert(key, observation);
2552        } else if context.control_tx.send(observation).await.is_err() {
2553            return false;
2554        }
2555    }
2556
2557    // Foreground reclassification runs after the text pass above so a
2558    // session `collect_terminal_frames` just re-armed (text went from
2559    // `Ready` to a gate/failure) is picked up in this same tick instead of
2560    // waiting a full `FOREGROUND_RECLASSIFY_INTERVAL`. Kept out of the text
2561    // loop itself: this performs a real OS process-tree walk per due
2562    // session, which the text pass's own sequence gate exists specifically
2563    // to avoid paying on every changed frame.
2564    for observation in shell.reclassify_foreground().await {
2565        if context.control_tx.send(observation).await.is_err() {
2566            return false;
2567        }
2568    }
2569
2570    // Drain this iteration's shell-side efficiency facts (both calls above
2571    // recorded into them) and fold them into the shared profile
2572    // `NativeRuntime::shell_efficiency_snapshot` reads -- locked only for
2573    // this fold, dropped before `collect_exits`'s own `.await` below.
2574    context
2575        .shell_efficiency
2576        .lock()
2577        .unwrap_or_else(|poisoned| poisoned.into_inner())
2578        .fold(&shell.take_efficiency_facts());
2579
2580    let before = shell.active_session_count();
2581    for observation in shell.collect_exits().await {
2582        context
2583            .terminal_frames
2584            .lock()
2585            .unwrap_or_else(|poisoned| poisoned.into_inner())
2586            .remove(&(observation.instance_id, observation.generation));
2587        if context.control_tx.send(observation).await.is_err() {
2588            return false;
2589        }
2590    }
2591    update_active_count(
2592        &context.active_sessions,
2593        before,
2594        shell.active_session_count(),
2595    );
2596    true
2597}
2598
2599fn closes_terminal_session(observation: &ControlObservation) -> bool {
2600    matches!(
2601        observation,
2602        ControlObservation::StopCompleted { .. }
2603            | ControlObservation::StopFailed { .. }
2604            | ControlObservation::ProcessExited { .. }
2605    )
2606}
2607
2608fn update_active_count(counter: &AtomicUsize, before: usize, after: usize) {
2609    if after > before {
2610        counter.fetch_add(after - before, Ordering::AcqRel);
2611    } else if before > after {
2612        counter.fetch_sub(before - after, Ordering::AcqRel);
2613    }
2614}
2615
2616fn validate_effect_runtime_policy(effect: &ControlEffect) -> Result<(), String> {
2617    let (transport, policy, has_initial_prompt, is_resume) = match effect {
2618        ControlEffect::Spawn {
2619            transport,
2620            runtime_policy,
2621            request,
2622            ..
2623        } => (*transport, *runtime_policy, request.initial_prompt.is_some(), false),
2624        ControlEffect::SpawnResume {
2625            transport,
2626            runtime_policy,
2627            request,
2628            ..
2629        } => (*transport, *runtime_policy, request.initial_prompt.is_some(), true),
2630        ControlEffect::Stop { .. }
2631        | ControlEffect::WriteInput { .. }
2632        | ControlEffect::SubmitPrompt { .. }
2633        | ControlEffect::Interrupt
2634        | ControlEffect::Resize { .. }
2635        | ControlEffect::ObserveForeground
2636        | ControlEffect::ProbeCapabilities { .. }
2637        | ControlEffect::DiscoverHistory { .. }
2638        | ControlEffect::LoadHistory { .. }
2639        | ControlEffect::AuthorizeResume { .. }
2640        | ControlEffect::ResolveInteraction { .. }
2641        | ControlEffect::SetSessionMode { .. }
2642        | ControlEffect::SetSessionConfigOption { .. }
2643        | ControlEffect::SetSessionModel { .. } => return Ok(()),
2644    };
2645    policy
2646        .validate()
2647        .map_err(|error| format!("provider runtime policy is invalid: {error}"))?;
2648    // ACP speaks a structured protocol over stdio, not a PTY -- none of the
2649    // capabilities below (raw-pty lifecycle, semantic readiness, structured
2650    // prompt, provider session identity, semantic resume) describe anything
2651    // that exists for it; they all gate inferring provider state from
2652    // terminal text. The transport-support gate for ACP already lives in
2653    // the kernel (`spec.capabilities.transports.acp.is_some()`), so this
2654    // PTY-semantic policy simply does not apply here.
2655    if transport != TransportKind::Acp {
2656        require_runtime_capability(policy, ProviderRuntimeCapability::RawPtyLifecycle)?;
2657        if transport != TransportKind::Pty {
2658            require_runtime_capability(policy, ProviderRuntimeCapability::SemanticReadiness)?;
2659        }
2660        if has_initial_prompt {
2661            require_runtime_capability(policy, ProviderRuntimeCapability::SemanticReadiness)?;
2662            require_runtime_capability(policy, ProviderRuntimeCapability::StructuredPrompt)?;
2663        }
2664        if is_resume && has_initial_prompt {
2665            require_runtime_capability(policy, ProviderRuntimeCapability::ProviderSessionIdentity)?;
2666            require_runtime_capability(policy, ProviderRuntimeCapability::SemanticResume)?;
2667        }
2668    }
2669    Ok(())
2670}
2671
2672fn require_runtime_capability(
2673    policy: ProviderRuntimePolicy,
2674    capability: ProviderRuntimeCapability,
2675) -> Result<(), String> {
2676    if policy.admits(capability) {
2677        Ok(())
2678    } else {
2679        Err(format!(
2680            "provider runtime capability {capability:?} is not admitted"
2681        ))
2682    }
2683}
2684
2685fn effect_failure(effect: EffectEnvelope, message: String) -> ObservationEnvelope {
2686    let observation = match effect.effect {
2687        ControlEffect::Spawn { .. } | ControlEffect::SpawnResume { .. } => {
2688            ControlObservation::SpawnFailed { message }
2689        }
2690        ControlEffect::Stop { .. } => ControlObservation::StopFailed { message },
2691        ControlEffect::WriteInput { .. } => ControlObservation::InputFailed { message },
2692        ControlEffect::SubmitPrompt { .. } | ControlEffect::Interrupt => {
2693            ControlObservation::InputFailed { message }
2694        }
2695        ControlEffect::ResolveInteraction { target, .. } => {
2696            ControlObservation::InteractionResolutionFailed {
2697                interaction_id: target.interaction_id,
2698                message,
2699            }
2700        }
2701        ControlEffect::SetSessionMode { .. } => ControlObservation::SessionModeSetFailed { message },
2702        ControlEffect::SetSessionConfigOption { .. } => {
2703            ControlObservation::SessionConfigOptionSetFailed { message }
2704        }
2705        ControlEffect::SetSessionModel { .. } => ControlObservation::SessionModelSetFailed { message },
2706        ControlEffect::Resize { .. } => ControlObservation::ResizeFailed { message },
2707        ControlEffect::ObserveForeground => ControlObservation::ForegroundFailed { message },
2708        ControlEffect::ProbeCapabilities { .. } => ControlObservation::CapabilityProbeFailed {
2709            failure: CapabilityProbeFailure::ExecutorUnavailable,
2710        },
2711        ControlEffect::DiscoverHistory { .. } | ControlEffect::LoadHistory { .. } => {
2712            ControlObservation::HistoryFailed { message }
2713        }
2714        ControlEffect::AuthorizeResume { .. } => ControlObservation::ResumeFailed { message },
2715    };
2716    ObservationEnvelope {
2717        operation_id: Some(effect.operation_id),
2718        instance_id: effect.instance_id,
2719        generation: effect.generation,
2720        observation,
2721    }
2722}
2723
2724#[cfg(test)]
2725mod tests {
2726    use super::*;
2727    use gate4agent_types::{ApprovalLevel, OperationId, StartRequest, TerminalSize};
2728
2729    #[test]
2730    fn history_record_preserves_completed_turn_count() {
2731        let record = history_session_record(gate4agent_adapters::HistorySession {
2732            session_id: "session-1".to_owned(),
2733            title: None,
2734            cwd: None,
2735            model: None,
2736            message_count: 5,
2737            completed_turn_count: Some(2),
2738            total_tokens: 8,
2739            total_tokens_observed: true,
2740            messages: Vec::new(),
2741        });
2742
2743        assert_eq!(record.completed_turn_count, Some(2));
2744    }
2745
2746    #[test]
2747    fn runtime_preview_maps_known_and_unknown_token_totals_exactly() {
2748        let session = |total_tokens_observed, total_tokens| {
2749            gate4agent_adapters::HistorySession {
2750                session_id: "session-1".to_owned(),
2751                title: None,
2752                cwd: None,
2753                model: None,
2754                message_count: 0,
2755                completed_turn_count: None,
2756                total_tokens,
2757                total_tokens_observed,
2758                messages: Vec::new(),
2759            }
2760        };
2761
2762        let unknown = project_native_session_preview(session(false, 0), None, false, 1).unwrap();
2763        let observed_zero =
2764            project_native_session_preview(session(true, 0), None, false, 1).unwrap();
2765        let observed_nonzero =
2766            project_native_session_preview(session(true, 42), None, false, 1).unwrap();
2767
2768        assert_eq!(unknown.total_tokens, None);
2769        assert_eq!(observed_zero.total_tokens, Some(0));
2770        assert_eq!(observed_nonzero.total_tokens, Some(42));
2771    }
2772
2773    #[test]
2774    fn native_catalog_high_cardinality_outside_workspace_does_not_hide_match() {
2775        let unique = std::time::SystemTime::now()
2776            .duration_since(std::time::UNIX_EPOCH)
2777            .unwrap()
2778            .as_nanos();
2779        let root = std::env::temp_dir().join(format!(
2780            "gate4agent-native-catalog-prelimit-{}-{unique}",
2781            std::process::id(),
2782        ));
2783        let workspace = root.join("workspace");
2784        let outside = root.join("outside");
2785        let history = root.join("central-history");
2786        std::fs::create_dir_all(&workspace).unwrap();
2787        std::fs::create_dir_all(&outside).unwrap();
2788        std::fs::create_dir_all(&history).unwrap();
2789
2790        let transcript = |session_id: &str, cwd: &Path| {
2791            let cwd = cwd
2792                .to_str()
2793                .unwrap()
2794                .replace('\\', "\\\\")
2795                .replace('"', "\\\"");
2796            format!(
2797                "{{\"type\":\"user\",\"sessionId\":\"{session_id}\",\"cwd\":\"{cwd}\",\"message\":{{\"content\":\"question\"}}}}\n\
2798                 {{\"type\":\"assistant\",\"sessionId\":\"{session_id}\",\"cwd\":\"{cwd}\",\"message\":{{\"content\":\"answer\"}}}}"
2799            )
2800        };
2801        std::fs::write(
2802            history.join("zz-workspace-session.jsonl"),
2803            transcript("workspace-session", &workspace),
2804        )
2805        .unwrap();
2806        for index in 0..1_100 {
2807            let session_id = format!("outside-session-{index:04}");
2808            std::fs::write(
2809                history.join(format!("aa-outside-session-{index:04}.jsonl")),
2810                transcript(&session_id, &outside),
2811            )
2812            .unwrap();
2813        }
2814
2815        let config = NativeHistoryConfig::new(vec![
2816            NativeHistoryRoot::new(
2817                gate4agent_types::AdapterId::new("claude-code").unwrap(),
2818                HistorySourceLayout::SingleNdjson,
2819                &history,
2820            )
2821            .unwrap(),
2822        ])
2823        .unwrap();
2824        let mut catalog = NativeSessionCatalogAuthority::new(config);
2825        let entries = catalog
2826            .catalog(&AgentId::new("claude").unwrap(), &workspace, 1)
2827            .unwrap();
2828
2829        assert_eq!(entries.len(), 1);
2830        assert_eq!(entries[0].session_id, "workspace-session");
2831        std::fs::remove_dir_all(root).unwrap();
2832    }
2833
2834    #[test]
2835    fn native_catalog_attributes_deleted_historical_cwd_lexically() {
2836        let unique = std::time::SystemTime::now()
2837            .duration_since(std::time::UNIX_EPOCH)
2838            .unwrap()
2839            .as_nanos();
2840        let root = std::env::temp_dir().join(format!(
2841            "gate4agent-native-catalog-deleted-cwd-{}-{unique}",
2842            std::process::id(),
2843        ));
2844        let workspace = root.join("workspace");
2845        let deleted_cwd = workspace.join("removed-subdirectory");
2846        let history = root.join("central-history");
2847        std::fs::create_dir_all(&workspace).unwrap();
2848        std::fs::create_dir_all(&history).unwrap();
2849        let cwd = deleted_cwd
2850            .to_str()
2851            .unwrap()
2852            .replace('\\', "\\\\")
2853            .replace('"', "\\\"");
2854        let workspace_key = NativePathKey::from_canonical_root(
2855            &std::fs::canonicalize(&workspace).unwrap(),
2856        )
2857        .unwrap();
2858        let deleted_key = NativePathKey::from_declared_cwd(
2859            deleted_cwd.to_str().unwrap(),
2860        )
2861        .unwrap();
2862        assert!(
2863            deleted_key.starts_with(&workspace_key),
2864            "declared cwd {deleted_key:?} must remain beneath workspace {workspace_key:?}",
2865        );
2866        std::fs::write(
2867            history.join("deleted-cwd-session.jsonl"),
2868            format!(
2869                "{{\"type\":\"user\",\"sessionId\":\"deleted-cwd-session\",\"cwd\":\"{cwd}\",\"message\":{{\"content\":\"question\"}}}}\n\
2870                 {{\"type\":\"assistant\",\"sessionId\":\"deleted-cwd-session\",\"cwd\":\"{cwd}\",\"message\":{{\"content\":\"answer\"}}}}"
2871            ),
2872        )
2873        .unwrap();
2874
2875        let config = NativeHistoryConfig::new(vec![
2876            NativeHistoryRoot::new(
2877                gate4agent_types::AdapterId::new("claude-code").unwrap(),
2878                HistorySourceLayout::SingleNdjson,
2879                &history,
2880            )
2881            .unwrap(),
2882        ])
2883        .unwrap();
2884        let mut catalog = NativeSessionCatalogAuthority::new(config);
2885        let entries = catalog
2886            .catalog(&AgentId::new("claude").unwrap(), &workspace, 10)
2887            .unwrap();
2888
2889        assert_eq!(entries.len(), 1);
2890        assert_eq!(entries[0].session_id, "deleted-cwd-session");
2891        let preview = catalog
2892            .preview(
2893                &AgentId::new("claude").unwrap(),
2894                &workspace,
2895                &entries[0].selection_id,
2896                10,
2897            )
2898            .unwrap();
2899        assert_eq!(preview.session_id, "deleted-cwd-session");
2900        std::fs::remove_dir_all(root).unwrap();
2901    }
2902
2903    #[test]
2904    fn native_catalog_keeps_provider_candidate_indexes_isolated() {
2905        let unique = std::time::SystemTime::now()
2906            .duration_since(std::time::UNIX_EPOCH)
2907            .unwrap()
2908            .as_nanos();
2909        let root = std::env::temp_dir().join(format!(
2910            "gate4agent-native-catalog-provider-isolation-{}-{unique}",
2911            std::process::id(),
2912        ));
2913        let workspace = root.join("workspace");
2914        let claude_history = root.join("claude-history");
2915        let second_history = root.join("codex-history");
2916        std::fs::create_dir_all(&workspace).unwrap();
2917        std::fs::create_dir_all(&claude_history).unwrap();
2918        std::fs::create_dir_all(&second_history).unwrap();
2919        let cwd = workspace
2920            .to_str()
2921            .unwrap()
2922            .replace('\\', "\\\\")
2923            .replace('"', "\\\"");
2924        std::fs::write(
2925            claude_history.join("claude-session.jsonl"),
2926            format!(
2927                "{{\"type\":\"user\",\"sessionId\":\"claude-session\",\"cwd\":\"{cwd}\",\"message\":{{\"content\":\"question\"}}}}\n\
2928                 {{\"type\":\"assistant\",\"sessionId\":\"claude-session\",\"cwd\":\"{cwd}\",\"message\":{{\"content\":\"answer\"}}}}"
2929            ),
2930        )
2931        .unwrap();
2932        std::fs::write(
2933            second_history.join("codex-session.jsonl"),
2934            format!(
2935                "{{\"type\":\"turn_context\",\"payload\":{{\"cwd\":\"{cwd}\"}}}}\n\
2936                 {{\"type\":\"event_msg\",\"payload\":{{\"type\":\"user_message\",\"message\":\"question\"}}}}\n\
2937                 {{\"type\":\"event_msg\",\"payload\":{{\"type\":\"agent_message\",\"message\":\"answer\"}}}}"
2938            ),
2939        )
2940        .unwrap();
2941        let limits = gate4agent_shell_history::NativeHistoryLimits {
2942            max_candidates: 1,
2943            ..gate4agent_shell_history::NativeHistoryLimits::default()
2944        };
2945        let config = NativeHistoryConfig::with_limits(
2946            vec![
2947                NativeHistoryRoot::new(
2948                    gate4agent_types::AdapterId::new("claude-code").unwrap(),
2949                    HistorySourceLayout::SingleNdjson,
2950                    &claude_history,
2951                )
2952                .unwrap(),
2953                NativeHistoryRoot::new(
2954                    gate4agent_types::AdapterId::new("codex").unwrap(),
2955                    HistorySourceLayout::NdjsonWithOptionalIndex,
2956                    &second_history,
2957                )
2958                .unwrap(),
2959            ],
2960            limits,
2961        )
2962        .unwrap();
2963        let mut catalog = NativeSessionCatalogAuthority {
2964            catalog: builtin_registry().clone(),
2965            history_config: config,
2966            authorities: HashMap::new(),
2967            indexes: HashMap::new(),
2968        };
2969        let claude = catalog
2970            .catalog(&AgentId::new("claude").unwrap(), &workspace, 1)
2971            .unwrap();
2972        let second = catalog
2973            .catalog(&AgentId::new("codex").unwrap(), &workspace, 1)
2974            .unwrap();
2975
2976        assert_eq!(claude.len(), 1);
2977        assert_eq!(second.len(), 1);
2978        let preview = catalog
2979            .preview(
2980                &AgentId::new("claude").unwrap(),
2981                &workspace,
2982                &claude[0].selection_id,
2983                10,
2984            )
2985            .unwrap();
2986        assert_eq!(preview.session_id, "claude-session");
2987        std::fs::remove_dir_all(root).unwrap();
2988    }
2989
2990    #[test]
2991    fn native_catalog_rejects_incomplete_provider_scan() {
2992        let unique = std::time::SystemTime::now()
2993            .duration_since(std::time::UNIX_EPOCH)
2994            .unwrap()
2995            .as_nanos();
2996        let root = std::env::temp_dir().join(format!(
2997            "gate4agent-native-catalog-incomplete-{}-{unique}",
2998            std::process::id(),
2999        ));
3000        let workspace = root.join("workspace");
3001        let history = root.join("history");
3002        std::fs::create_dir_all(&workspace).unwrap();
3003        std::fs::create_dir_all(&history).unwrap();
3004        let cwd = workspace
3005            .to_str()
3006            .unwrap()
3007            .replace('\\', "\\\\")
3008            .replace('"', "\\\"");
3009        for index in 0..2 {
3010            std::fs::write(
3011                history.join(format!("session-{index}.jsonl")),
3012                format!(
3013                    "{{\"type\":\"user\",\"sessionId\":\"session-{index}\",\"cwd\":\"{cwd}\",\"message\":{{\"content\":\"question\"}}}}"
3014                ),
3015            )
3016            .unwrap();
3017        }
3018        let limits = gate4agent_shell_history::NativeHistoryLimits {
3019            max_walk_entries: 1,
3020            ..gate4agent_shell_history::NativeHistoryLimits::default()
3021        };
3022        let config = NativeHistoryConfig::with_limits(
3023            vec![NativeHistoryRoot::new(
3024                gate4agent_types::AdapterId::new("claude-code").unwrap(),
3025                HistorySourceLayout::SingleNdjson,
3026                &history,
3027            )
3028            .unwrap()],
3029            limits,
3030        )
3031        .unwrap();
3032        let mut catalog = NativeSessionCatalogAuthority::new(config);
3033
3034        assert_eq!(
3035            catalog.catalog(&AgentId::new("claude").unwrap(), &workspace, 10),
3036            Err(NativeSessionCatalogError::CatalogUnavailable),
3037        );
3038        std::fs::remove_dir_all(root).unwrap();
3039    }
3040
3041    #[test]
3042    fn native_catalog_initial_window_excludes_older_until_explicit_page() {
3043        let unique = std::time::SystemTime::now()
3044            .duration_since(std::time::UNIX_EPOCH)
3045            .unwrap()
3046            .as_nanos();
3047        let workspace = std::env::temp_dir().join(format!(
3048            "gate4agent-native-catalog-window-{}-{unique}",
3049            std::process::id(),
3050        ));
3051        std::fs::create_dir_all(&workspace).unwrap();
3052        let provider = AgentId::new("claude").unwrap();
3053        let mut catalog = NativeSessionCatalogAuthority::new(
3054            NativeHistoryConfig::new(Vec::new()).unwrap(),
3055        );
3056        let spec = catalog.catalog.get(&provider).unwrap().clone();
3057        let request = HistoryDiscoveryRequest::from_spec(&spec, None, 2).unwrap();
3058        let cwd_key = NativePathKey::from_canonical_root(
3059            &std::fs::canonicalize(&workspace).unwrap(),
3060        )
3061        .unwrap();
3062        let now = SystemTime::now()
3063            .duration_since(UNIX_EPOCH)
3064            .unwrap()
3065            .as_millis() as u64;
3066        let indexed = |id: &str, modified_at_unix_ms| IndexedNativeSession {
3067            candidate: HistoryCandidate::new(
3068                &request,
3069                gate4agent_provider_ports::HistoryCandidateId::new(id).unwrap(),
3070                id,
3071                Some(modified_at_unix_ms),
3072            )
3073            .unwrap(),
3074            selection_id: id.to_owned(),
3075            modified_at_unix_ms: Some(modified_at_unix_ms),
3076            cwd_key: Some(cwd_key.clone()),
3077            session_id: id.to_owned(),
3078            title: Some(id.to_owned()),
3079            model: None,
3080            project_label: "workspace".to_owned(),
3081            provider_session: ProviderSessionIdentity {
3082                key: gate4agent_types::ProviderSessionKey::SessionId,
3083                id: id.to_owned(),
3084                transcript_path: None,
3085            },
3086        };
3087        catalog.indexes.insert(
3088            provider.clone(),
3089            NativeProviderSessionIndex {
3090                refreshed_at: Instant::now(),
3091                revision: 1,
3092                recent_cutoff_unix_ms: now.saturating_sub(NATIVE_SESSION_RECENT_WINDOW_MS),
3093                sessions: vec![
3094                    indexed("recent-session", now),
3095                    indexed("recent-session-2", now.saturating_sub(1)),
3096                    indexed(
3097                        "older-session",
3098                        now.saturating_sub(NATIVE_SESSION_RECENT_WINDOW_MS + 1),
3099                    ),
3100                ],
3101            },
3102        );
3103
3104        let recent = catalog
3105            .catalog_initial_for_workspace(
3106                &provider,
3107                &workspace,
3108                &[workspace.clone()],
3109                1,
3110            )
3111            .unwrap();
3112        let older = catalog
3113            .catalog_page_for_workspace(
3114                &provider,
3115                &workspace,
3116                &[workspace.clone()],
3117                NativeSessionCatalogWindow::Older,
3118                recent.revision,
3119                recent.cutoff_unix_ms,
3120                None,
3121                64,
3122            )
3123            .unwrap();
3124
3125        assert_eq!(recent.entries.len(), 1);
3126        assert_eq!(recent.entries[0].session_id, "recent-session");
3127        assert_eq!(recent.recent_total_count, 2);
3128        assert_eq!(recent.older_total_count, 1);
3129        assert_eq!(recent.next_after_selection_id.as_deref(), Some("recent-session"));
3130        let next = catalog
3131            .catalog_page_for_workspace(
3132                &provider,
3133                &workspace,
3134                &[workspace.clone()],
3135                NativeSessionCatalogWindow::Recent,
3136                recent.revision,
3137                recent.cutoff_unix_ms,
3138                recent.next_after_selection_id.as_deref(),
3139                1,
3140            )
3141            .unwrap();
3142        assert_eq!(next.entries[0].session_id, "recent-session-2");
3143        assert_eq!(next.next_after_selection_id, None);
3144        assert_eq!(older.entries.len(), 1);
3145        assert_eq!(older.entries[0].session_id, "older-session");
3146        assert_eq!(
3147            catalog.catalog_page_for_workspace(
3148                &provider,
3149                &workspace,
3150                &[workspace.clone()],
3151                NativeSessionCatalogWindow::Recent,
3152                recent.revision + 1,
3153                recent.cutoff_unix_ms,
3154                None,
3155                1,
3156            ),
3157            Err(NativeSessionCatalogError::StaleCatalog),
3158        );
3159        assert_eq!(
3160            catalog.catalog_page_for_workspace(
3161                &provider,
3162                &workspace,
3163                &[workspace.clone()],
3164                NativeSessionCatalogWindow::Recent,
3165                recent.revision,
3166                recent.cutoff_unix_ms + 1,
3167                None,
3168                1,
3169            ),
3170            Err(NativeSessionCatalogError::StaleCatalog),
3171        );
3172        std::fs::remove_dir_all(workspace).unwrap();
3173    }
3174
3175    #[test]
3176    fn unregistered_catalog_uses_revision_bound_opaque_groups_for_same_basename() {
3177        let unique = std::time::SystemTime::now()
3178            .duration_since(std::time::UNIX_EPOCH)
3179            .unwrap()
3180            .as_nanos();
3181        let root = std::env::temp_dir().join(format!(
3182            "gate4agent-native-external-groups-{}-{unique}",
3183            std::process::id(),
3184        ));
3185        let first = root.join("private-a").join("shared");
3186        let second = root.join("private-b").join("shared");
3187        std::fs::create_dir_all(&first).unwrap();
3188        std::fs::create_dir_all(&second).unwrap();
3189        let provider = AgentId::new("claude").unwrap();
3190        let mut catalog = NativeSessionCatalogAuthority::new(
3191            NativeHistoryConfig::new(Vec::new()).unwrap(),
3192        );
3193        let spec = catalog.catalog.get(&provider).unwrap().clone();
3194        let request = HistoryDiscoveryRequest::from_spec(&spec, None, 2).unwrap();
3195        let now = SystemTime::now()
3196            .duration_since(UNIX_EPOCH)
3197            .unwrap()
3198            .as_millis() as u64;
3199        let indexed = |id: &str, cwd: &Path| IndexedNativeSession {
3200            candidate: HistoryCandidate::new(
3201                &request,
3202                gate4agent_provider_ports::HistoryCandidateId::new(id).unwrap(),
3203                id,
3204                Some(now),
3205            )
3206            .unwrap(),
3207            selection_id: id.to_owned(),
3208            modified_at_unix_ms: Some(now),
3209            cwd_key: Some(
3210                NativePathKey::from_declared_cwd(cwd.to_str().unwrap()).unwrap(),
3211            ),
3212            session_id: id.to_owned(),
3213            title: None,
3214            model: None,
3215            project_label: "shared".to_owned(),
3216            provider_session: ProviderSessionIdentity {
3217                key: gate4agent_types::ProviderSessionKey::SessionId,
3218                id: id.to_owned(),
3219                transcript_path: None,
3220            },
3221        };
3222        catalog.indexes.insert(
3223            provider.clone(),
3224            NativeProviderSessionIndex {
3225                refreshed_at: Instant::now(),
3226                revision: 1,
3227                recent_cutoff_unix_ms: now.saturating_sub(NATIVE_SESSION_RECENT_WINDOW_MS),
3228                sessions: vec![indexed("session-a", &first), indexed("session-b", &second)],
3229            },
3230        );
3231
3232        let page = catalog
3233            .catalog_initial_for_scope(
3234                &provider,
3235                NativeSessionCatalogScope::Unregistered,
3236                None,
3237                &[],
3238                64,
3239            )
3240            .unwrap();
3241        let groups = page
3242            .entries
3243            .iter()
3244            .map(|entry| entry.external_group.as_ref().unwrap())
3245            .collect::<Vec<_>>();
3246        assert_eq!(groups.len(), 2);
3247        assert_eq!(groups[0].display_name, "shared");
3248        assert_eq!(groups[1].display_name, "shared");
3249        assert_ne!(groups[0].group_id, groups[1].group_id);
3250        assert!(groups.iter().all(|group| group.group_id.starts_with("external-")));
3251        assert!(groups.iter().all(|group| !group.group_id.contains("private")));
3252        std::fs::remove_dir_all(root).unwrap();
3253    }
3254
3255    #[test]
3256    fn native_catalog_revision_changes_only_when_locator_metadata_changes() {
3257        let unique = std::time::SystemTime::now()
3258            .duration_since(std::time::UNIX_EPOCH)
3259            .unwrap()
3260            .as_nanos();
3261        let root = std::env::temp_dir().join(format!(
3262            "gate4agent-native-catalog-revision-{}-{unique}",
3263            std::process::id(),
3264        ));
3265        let workspace = root.join("workspace");
3266        let history = root.join("history");
3267        std::fs::create_dir_all(&workspace).unwrap();
3268        std::fs::create_dir_all(&history).unwrap();
3269        let cwd = workspace.to_str().unwrap().replace('\\', "\\\\");
3270        let transcript = |title: &str| {
3271            format!(
3272                "{{\"type\":\"user\",\"sessionId\":\"revision-session\",\"cwd\":\"{cwd}\",\"message\":{{\"content\":\"{title}\"}}}}"
3273            )
3274        };
3275        let path = history.join("revision-session.jsonl");
3276        std::fs::write(&path, transcript("first title")).unwrap();
3277        let config = NativeHistoryConfig::new(vec![
3278            NativeHistoryRoot::new(
3279                gate4agent_types::AdapterId::new("claude-code").unwrap(),
3280                HistorySourceLayout::SingleNdjson,
3281                &history,
3282            )
3283            .unwrap(),
3284        ])
3285        .unwrap();
3286        let provider = AgentId::new("claude").unwrap();
3287        let mut catalog = NativeSessionCatalogAuthority::new(config);
3288
3289        let first = catalog
3290            .catalog_initial_for_workspace(&provider, &workspace, &[workspace.clone()], 64)
3291            .unwrap();
3292        catalog.indexes.get_mut(&provider).unwrap().refreshed_at =
3293            Instant::now() - NATIVE_SESSION_INDEX_REFRESH_INTERVAL;
3294        let unchanged = catalog
3295            .catalog_initial_for_workspace(&provider, &workspace, &[workspace.clone()], 64)
3296            .unwrap();
3297        assert_eq!(unchanged.revision, first.revision);
3298
3299        std::fs::write(&path, transcript("changed title")).unwrap();
3300        catalog.indexes.get_mut(&provider).unwrap().refreshed_at =
3301            Instant::now() - NATIVE_SESSION_INDEX_REFRESH_INTERVAL;
3302        let changed = catalog
3303            .catalog_initial_for_workspace(&provider, &workspace, &[workspace.clone()], 64)
3304            .unwrap();
3305        assert_ne!(changed.revision, first.revision);
3306        std::fs::remove_dir_all(root).unwrap();
3307    }
3308
3309    #[test]
3310    fn native_catalog_high_cardinality_nested_worktree_uses_deepest_owner() {
3311        let unique = std::time::SystemTime::now()
3312            .duration_since(std::time::UNIX_EPOCH)
3313            .unwrap()
3314            .as_nanos();
3315        let root = std::env::temp_dir().join(format!(
3316            "gate4agent-native-catalog-worktree-owner-{}-{unique}",
3317            std::process::id(),
3318        ));
3319        let parent = root.join("workspace");
3320        let worktree = parent.join("linked-worktree");
3321        let history = root.join("central-history");
3322        std::fs::create_dir_all(&worktree).unwrap();
3323        std::fs::create_dir_all(&history).unwrap();
3324        let cwd = worktree
3325            .to_str()
3326            .unwrap()
3327            .replace('\\', "\\\\")
3328            .replace('"', "\\\"");
3329        std::fs::write(
3330            history.join("worktree-session.jsonl"),
3331            format!(
3332                "{{\"type\":\"user\",\"sessionId\":\"worktree-session\",\"cwd\":\"{cwd}\",\"message\":{{\"content\":\"question\"}}}}\n\
3333                 {{\"type\":\"assistant\",\"sessionId\":\"worktree-session\",\"cwd\":\"{cwd}\",\"message\":{{\"content\":\"answer\"}}}}"
3334            ),
3335        )
3336        .unwrap();
3337        let parent_cwd = parent
3338            .to_str()
3339            .unwrap()
3340            .replace('\\', "\\\\")
3341            .replace('"', "\\\"");
3342        std::fs::write(
3343            history.join("parent-session.jsonl"),
3344            format!(
3345                "{{\"type\":\"user\",\"sessionId\":\"parent-session\",\"cwd\":\"{parent_cwd}\",\"message\":{{\"content\":\"parent question\"}}}}\n\
3346                 {{\"type\":\"assistant\",\"sessionId\":\"parent-session\",\"cwd\":\"{parent_cwd}\",\"message\":{{\"content\":\"parent answer\"}}}}"
3347            ),
3348        )
3349        .unwrap();
3350        for index in 0..1_100 {
3351            std::fs::write(
3352                history.join(format!("child-session-{index:04}.jsonl")),
3353                format!(
3354                    "{{\"type\":\"user\",\"sessionId\":\"child-session-{index:04}\",\"cwd\":\"{cwd}\",\"message\":{{\"content\":\"child question\"}}}}\n\
3355                     {{\"type\":\"assistant\",\"sessionId\":\"child-session-{index:04}\",\"cwd\":\"{cwd}\",\"message\":{{\"content\":\"child answer\"}}}}"
3356                ),
3357            )
3358            .unwrap();
3359        }
3360
3361        let config = NativeHistoryConfig::new(vec![
3362            NativeHistoryRoot::new(
3363                gate4agent_types::AdapterId::new("claude-code").unwrap(),
3364                HistorySourceLayout::SingleNdjson,
3365                &history,
3366            )
3367            .unwrap(),
3368        ])
3369        .unwrap();
3370        let registered = vec![parent.clone(), worktree.clone()];
3371        let mut catalog = NativeSessionCatalogAuthority::new(config);
3372        let parent_entries = catalog
3373            .catalog_for_workspace(
3374                &AgentId::new("claude").unwrap(),
3375                &parent,
3376                &registered,
3377                1,
3378            )
3379            .unwrap();
3380        let worktree_entries = catalog
3381            .catalog_for_workspace(
3382                &AgentId::new("claude").unwrap(),
3383                &worktree,
3384                &registered,
3385                10,
3386            )
3387            .unwrap();
3388
3389        assert_eq!(parent_entries.len(), 1);
3390        assert_eq!(parent_entries[0].session_id, "parent-session");
3391        assert_eq!(worktree_entries.len(), 10);
3392        assert_eq!(
3393            catalog.preview_for_workspace(
3394                &AgentId::new("claude").unwrap(),
3395                &parent,
3396                &registered,
3397                &worktree_entries[0].selection_id,
3398                10,
3399            ),
3400            Err(NativeSessionPreviewError::SessionNotFound),
3401        );
3402        let preview = catalog
3403            .preview_for_workspace(
3404                &AgentId::new("claude").unwrap(),
3405                &worktree,
3406                &registered,
3407                &worktree_entries[0].selection_id,
3408                10,
3409            )
3410            .unwrap();
3411        assert_eq!(preview.session_id, worktree_entries[0].session_id);
3412        std::fs::remove_dir_all(root).unwrap();
3413    }
3414
3415    #[test]
3416    fn runtime_dispatcher_rejects_raw_semantic_prompt_before_worker_dispatch() {
3417        let mut dispatcher = NativeEffectDispatcher::new(
3418            AgentRegistry::new([]).unwrap(),
3419            NativeRuntimeConfig::default(),
3420            None,
3421        );
3422        dispatcher.dispatch(EffectEnvelope {
3423            operation_id: OperationId(1),
3424            instance_id: AgentInstanceId(1),
3425            generation: SessionGeneration(1),
3426            effect: ControlEffect::Spawn {
3427                agent_id: AgentId::new("unknown-provider").unwrap(),
3428                transport: TransportKind::Pty,
3429                runtime_policy: ProviderRuntimePolicy::raw_pty(),
3430                request: StartRequest {
3431                    working_directory: ".".to_owned(),
3432                    terminal_size: TerminalSize { rows: 24, columns: 80 },
3433                    initial_prompt: Some("must-not-run".to_owned()),
3434                    session_options: None,
3435                    approval_level: ApprovalLevel::default(),
3436                },
3437            },
3438        });
3439
3440        let (observations, terminal_frames) = dispatcher.drain_observations(1);
3441        assert_eq!(terminal_frames, 0);
3442        assert!(dispatcher.workers.is_empty());
3443        assert!(matches!(
3444            observations.as_slice(),
3445            [ObservationEnvelope {
3446                observation: ControlObservation::SpawnFailed { message },
3447                ..
3448            }] if message.contains("SemanticReadiness")
3449        ));
3450    }
3451
3452    /// Repro for the standing observation that `gate4agent-node` burns CPU
3453    /// with zero live PTY sessions: drives `tick()` repeatedly on a runtime
3454    /// with no spawned sessions and no installed native providers -- the
3455    /// same "idle" shape as the live stand -- and prints which phase the
3456    /// time actually landed in. Run with `--nocapture` to read the numbers;
3457    /// the assertions only check that every phase got a full window of
3458    /// samples, since the absolute timings are hardware-dependent and the
3459    /// point of this test is the printed breakdown, not a threshold.
3460    #[tokio::test]
3461    async fn tick_profile_snapshot_records_every_phase_on_an_idle_runtime() {
3462        let catalog = builtin_registry().clone();
3463        let (_, mut runtime) = NativeRuntime::new(catalog, NativeRuntimeConfig::default());
3464        const IDLE_TICKS: usize = 64;
3465        for _ in 0..IDLE_TICKS {
3466            runtime.tick().await;
3467        }
3468        let snapshot = runtime.tick_profile_snapshot();
3469        println!(
3470            "idle tick phases (us, p50/p95/max/n over {IDLE_TICKS} ticks): \
3471             drain_observations={:?} drain_ingress={:?} step_control_plane={:?} \
3472             dispatch_effects={:?} publish_step={:?} provider_supervisors={:?}",
3473            snapshot.drain_observations_us,
3474            snapshot.drain_ingress_us,
3475            snapshot.step_control_plane_us,
3476            snapshot.dispatch_effects_us,
3477            snapshot.publish_step_us,
3478            snapshot.provider_supervisors_us,
3479        );
3480        for distribution in [
3481            snapshot.drain_observations_us,
3482            snapshot.drain_ingress_us,
3483            snapshot.step_control_plane_us,
3484            snapshot.dispatch_effects_us,
3485            snapshot.publish_step_us,
3486            snapshot.provider_supervisors_us,
3487        ] {
3488            assert_eq!(distribution.count, IDLE_TICKS);
3489        }
3490    }
3491
3492    #[test]
3493    fn runtime_policy_admits_raw_native_resume_without_prompt_only() {
3494        let raw = ProviderRuntimePolicy::raw_pty();
3495        let raw_resume = ControlEffect::SpawnResume {
3496            agent_id: AgentId::new("unknown-provider").unwrap(),
3497            transport: TransportKind::Pty,
3498            provider_session: gate4agent_types::ProviderSessionIdentity {
3499                key: gate4agent_types::ProviderSessionKey::SessionId,
3500                id: "session-1".to_owned(),
3501                transcript_path: None,
3502            },
3503            runtime_policy: raw,
3504            request: ResumeLaunchRequest {
3505                working_directory: ".".to_owned(),
3506                terminal_size: gate4agent_types::TerminalSize { rows: 24, columns: 80 },
3507                initial_prompt: None,
3508            },
3509        };
3510
3511        assert!(validate_effect_runtime_policy(&raw_resume).is_ok());
3512
3513        let resume_with_prompt = ControlEffect::SpawnResume {
3514            agent_id: AgentId::new("unknown-provider").unwrap(),
3515            transport: TransportKind::Pty,
3516            provider_session: gate4agent_types::ProviderSessionIdentity {
3517                key: gate4agent_types::ProviderSessionKey::SessionId,
3518                id: "session-1".to_owned(),
3519                transcript_path: None,
3520            },
3521            runtime_policy: raw,
3522            request: ResumeLaunchRequest {
3523                working_directory: ".".to_owned(),
3524                terminal_size: gate4agent_types::TerminalSize { rows: 24, columns: 80 },
3525                initial_prompt: Some("must-not-run".to_owned()),
3526            },
3527        };
3528        assert!(validate_effect_runtime_policy(&resume_with_prompt)
3529            .unwrap_err()
3530            .contains("SemanticReadiness"));
3531    }
3532
3533    #[test]
3534    fn acp_transport_bypasses_the_pty_semantic_policy_gate_pty_still_enforces_it() {
3535        // No raw PTY, no verified terminal semantics -- exactly grok's real
3536        // policy (ACP-only, no PTY transport at all). `TransportKind::Acp`
3537        // must not care: it has no terminal to infer state from.
3538        let no_pty_capabilities_at_all = ProviderRuntimePolicy::new(
3539            false, false, false, false, false, false,
3540        )
3541        .expect("an all-false policy is internally valid");
3542        let acp_spawn = ControlEffect::SpawnResume {
3543            agent_id: AgentId::new("grok").unwrap(),
3544            transport: TransportKind::Acp,
3545            provider_session: gate4agent_types::ProviderSessionIdentity {
3546                key: gate4agent_types::ProviderSessionKey::SessionId,
3547                id: "session-1".to_owned(),
3548                transcript_path: None,
3549            },
3550            runtime_policy: no_pty_capabilities_at_all,
3551            request: ResumeLaunchRequest {
3552                working_directory: ".".to_owned(),
3553                terminal_size: gate4agent_types::TerminalSize { rows: 24, columns: 80 },
3554                initial_prompt: None,
3555            },
3556        };
3557        assert!(validate_effect_runtime_policy(&acp_spawn).is_ok());
3558
3559        // The same all-false policy is still correctly refused for Pty --
3560        // this fix narrows the gate to skip Acp specifically, it does not
3561        // weaken it for the transport that does need it.
3562        let pty_spawn = ControlEffect::SpawnResume {
3563            agent_id: AgentId::new("grok").unwrap(),
3564            transport: TransportKind::Pty,
3565            provider_session: gate4agent_types::ProviderSessionIdentity {
3566                key: gate4agent_types::ProviderSessionKey::SessionId,
3567                id: "session-1".to_owned(),
3568                transcript_path: None,
3569            },
3570            runtime_policy: no_pty_capabilities_at_all,
3571            request: ResumeLaunchRequest {
3572                working_directory: ".".to_owned(),
3573                terminal_size: gate4agent_types::TerminalSize { rows: 24, columns: 80 },
3574                initial_prompt: None,
3575            },
3576        };
3577        assert!(validate_effect_runtime_policy(&pty_spawn)
3578            .unwrap_err()
3579            .contains("RawPtyLifecycle"));
3580    }
3581}