1mod 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
86pub 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(¤t_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(¤t_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
1260pub 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 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 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 pub fn native_provider_shutdown_complete(&self) -> bool {
1449 self.provider_supervisors
1450 .values()
1451 .all(|supervisor| supervisor.state() == ProviderSupervisorState::Closed)
1452 }
1453
1454 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 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 pub fn tick_profile_snapshot(&self) -> tick_profile::TickProfileSnapshot {
1578 self.tick_profile.snapshot()
1579 }
1580
1581 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 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 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 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 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 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 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 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 for observation in shell.reclassify_foreground().await {
2565 if context.control_tx.send(observation).await.is_err() {
2566 return false;
2567 }
2568 }
2569
2570 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 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 ®istered,
3377 1,
3378 )
3379 .unwrap();
3380 let worktree_entries = catalog
3381 .catalog_for_workspace(
3382 &AgentId::new("claude").unwrap(),
3383 &worktree,
3384 ®istered,
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 ®istered,
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 ®istered,
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 #[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 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 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}