1#![recursion_limit = "256"]
2#![allow(clippy::result_large_err, clippy::cloned_ref_to_slice_refs)]
3pub use harn_clock as clock;
23
24mod runtime_stack;
25pub use runtime_stack::RUNTIME_STACK_SIZE;
26
27pub mod a2a;
28pub mod actor_chain;
29pub mod agent_events;
30pub(crate) mod agent_session_journal;
31pub mod agent_sessions;
32pub mod agent_transcript_budget;
33pub mod atomic_io;
34pub mod autonomy;
35#[cfg(feature = "cloud-aws")]
36pub(crate) mod aws_sigv4;
37#[cfg(not(feature = "cloud-aws"))]
38#[path = "aws_sigv4_disabled.rs"]
39pub(crate) mod aws_sigv4;
40pub mod boundary;
41pub mod bridge;
42pub use bridge::{inject_leading_authorities, inject_leading_authority};
43pub mod builtin_profile;
44pub mod bytecode_cache;
45pub mod call_budget;
46pub mod canonical_json;
47pub mod channel_guardrails;
48pub mod channels;
49pub mod checkpoint;
50mod chunk;
51mod compiler;
52pub mod composition;
53pub mod conditional_replace;
54pub mod config;
55pub mod connectors;
56pub mod context_manifest;
57pub mod corrections;
58pub mod coverage;
59pub(crate) mod durable_rate_limit;
60pub mod duration_parse;
61pub mod egress;
62pub mod environment_registry;
63pub mod event_log;
64pub mod events;
65pub mod external_agent;
66pub mod flow;
67pub mod harness;
68pub mod harness_auth;
69pub(crate) mod harness_crypto;
70pub mod harness_net;
71pub mod harness_system;
72pub mod harness_tenant;
73pub mod host_attachments;
74mod http;
75pub mod jsonrpc;
76pub(crate) mod limits;
77pub mod linked_program;
78pub mod llm;
79pub mod llm_config;
80pub mod mcp;
81pub mod mcp_allowlist;
82pub mod mcp_auth;
83pub mod mcp_bulk_auth;
84pub mod mcp_card;
85pub mod mcp_client_roots;
86pub mod mcp_elicit;
87pub mod mcp_host;
88pub mod mcp_identity;
89pub mod mcp_input;
90pub mod mcp_json_discovery;
91pub mod mcp_oauth;
92pub mod mcp_presets;
93pub mod mcp_progress;
94pub mod mcp_protocol;
95pub mod mcp_registry;
96pub mod mcp_sampling;
97pub mod mcp_server;
98pub mod metadata;
99pub mod module_artifact;
100pub mod module_source;
101pub mod observability;
102pub mod op_interrupt;
103pub mod orchestration;
104mod persistent_state;
105pub mod personas;
106pub mod portable;
107mod prepared_module;
108pub mod process_sandbox;
109pub mod profile;
110pub mod provenance;
111pub mod provider_catalog;
112pub mod receipts;
113pub mod record_filter;
114pub mod redact;
115pub mod run_events;
116pub mod runtime_context;
117pub(crate) mod runtime_guards;
118pub mod runtime_limits;
119pub mod runtime_paths;
120pub(crate) mod runtime_sqlite;
121pub mod schema;
122pub(crate) mod secret_patterns;
123pub mod secrets;
124pub mod security;
125pub mod session_bundle;
126pub mod session_timeline;
127pub mod sessions;
128pub(crate) mod shared_state;
129pub mod shells;
130pub mod skills;
131pub mod stdlib;
132pub mod stdlib_modules;
133pub mod step_runtime;
134pub mod store;
135pub(crate) mod synchronization;
136pub mod tenant;
137pub(crate) mod term;
138pub(crate) mod test_env;
139pub mod testbench;
140pub mod text;
141pub mod text_diff;
142pub mod tool_annotations;
143pub mod tool_call_cancellations;
144pub mod tool_surface;
145pub mod tracing;
146pub mod triggers;
147pub mod trust_graph;
148pub(crate) mod url_encoding;
149pub mod user_dirs;
150
151pub fn initialize_runtime_assets() {
159 secret_patterns::initialize_default_secret_patterns();
160}
161
162pub fn initialize_runtime() -> Result<(), environment_registry::EnvironmentValidationError> {
165 environment_registry::validate_startup_environment()?;
166 initialize_runtime_assets();
167 Ok(())
168}
169
170pub mod clock_mock {
175 pub(crate) use crate::triggers::test_util::clock::scope_capability_clock;
176 pub use crate::triggers::test_util::clock::{
177 active_clock, active_mock_clock, advance, clear_overrides, install_override, instant_now,
178 is_mocked, now_ms, now_utc, sleep, ClockInstant, ClockOverrideGuard, MockClock,
179 };
180
181 pub mod leak_audit {
185 pub use crate::triggers::test_util::clock_leak::{
186 drain, enter_scope, install_scope, instant_now, reset, snapshot, wall_now, ClockLeak,
187 ClockLeakScope, ClockLeakScopeGuard,
188 };
189 }
190}
191
192pub(crate) mod text_index;
193pub mod typecheck;
194pub mod value;
195pub mod verification;
196pub mod visible_text;
197mod vm;
198pub(crate) mod wait_for_graph;
199pub mod waitpoints;
200pub mod windows_path;
201pub mod workspace_anchor;
202pub mod workspace_path;
203
204pub use persistent_state::{register_persistent_state_builtins_at_root, PersistentStateRoot};
205pub use prepared_module::{PreparedModuleCache, PreparedModuleCacheStats};
206
207pub use actor_chain::{
208 ActorChain, ActorChainEntry, ActorChainError, Principal, ScopeAttenuationMode,
209 ScopeAttenuationPolicy, ScopeAttenuationViolation,
210};
211pub use call_budget::{
212 charge_mcp_call, charge_pg_query, install_mcp_call_budget, install_pg_query_budget,
213 McpCallBudgetGuard, PgQueryBudgetGuard,
214};
215pub use checkpoint::register_checkpoint_builtins;
216pub use chunk::*;
217pub use compiler::*;
218pub use connectors::{
219 active_connector_client, active_metrics_registry, clear_active_connector_clients,
220 clear_active_metrics_registry, connector_export_denied_builtin_reason,
221 connector_export_denied_harness_method_reason, connector_export_effect_class,
222 cron::{CatchupMode, CronConnector},
223 default_connector_export_policy,
224 harn_module::{
225 load_contract as load_harn_connector_contract, HarnConnector, HarnConnectorContract,
226 },
227 hmac::{verify_hmac_signed, SIGNATURE_VERIFY_AUDIT_TOPIC},
228 install_active_connector_clients, install_active_metrics_registry,
229 postprocess_normalized_event, ActivationHandle, ClientError, Connector, ConnectorClient,
230 ConnectorCtx, ConnectorError, ConnectorExportEffectClass, ConnectorHttpResponse,
231 ConnectorMetricsSnapshot, ConnectorNormalizeResult, ConnectorRegistry, GenericWebhookConnector,
232 HarnConnectorEffectPolicies, MetricsRegistry, PostNormalizeOutcome, ProviderPayloadSchema,
233 RateLimitConfig, RateLimiterFactory, RawInbound, StreamConnector, TriggerBinding, TriggerKind,
234 TriggerRegistry, WebhookSignatureVariant,
235};
236pub use corrections::{
237 append_correction_record, apply_corrections_to_policy, correction_query_filters_from_json,
238 correction_record_from_json, policy_with_corrections, query_correction_records,
239 CorrectionQueryFilters, CorrectionRecord, CorrectionScope, CORRECTIONS_TOPIC,
240 CORRECTION_EVENT_KIND, CORRECTION_SCHEMA_V0,
241};
242pub use harn_kernel::BuiltinId;
243pub use harness::{
244 DenyEvent, Harness, HarnessAgent, HarnessCall, HarnessChannels, HarnessClock, HarnessEnv,
245 HarnessFs, HarnessKind, HarnessLlm, HarnessMemory, HarnessNet, HarnessObs, HarnessPostgres,
246 HarnessProcess, HarnessRandom, HarnessSecrets, HarnessSqlite, HarnessStdio, HarnessSystem,
247 HarnessTenant, HarnessTerm, HarnessTesting, MockHarnessBuilder, VmHarness,
248};
249pub use harness_auth::{
250 current_auth_principal, enter_auth_principal, AuthPrincipal, AuthPrincipalScopeGuard,
251 MISSING_PRINCIPAL_MESSAGE,
252};
253pub use harness_net::{
254 bypass_enabled as net_policy_bypass_enabled, NetMatcher, NetPolicy, NetPolicyAudit,
255 NetPolicyDecision, NetPolicyDefault, NetPolicyRule, OnViolation, HARN_NET_POLICY_BYPASS_ENV,
256 NET_POLICY_AUDIT_TOPIC,
257};
258pub use harness_tenant::{
259 current_tenant_id, enter_tenant, TenantScopeGuard, MISSING_TENANT_MESSAGE,
260};
261pub use http::{register_http_builtins, reset_http_state};
262pub use llm::register_llm_builtins;
263pub use llm::trigger_predicate::TriggerPredicateBudget;
264pub use llm::{
265 current_agent_session_id, install_llm_cost_budget, install_llm_token_budget,
266 peek_llm_cost_budget, peek_llm_token_budget, register_session_end_hook, set_llm_cost_budget,
267 set_llm_token_budget, LlmBudgetGuard, LlmTokenBudgetGuard, SessionEndHookRegistration,
268};
269pub use mcp::{connect_mcp_server_from_json, connect_mcp_server_from_spec, register_mcp_builtins};
270pub use mcp_allowlist::{
271 build_catalog as build_mcp_catalog, catalog_for_request as mcp_catalog_for_request,
272 AdvertisedItem as McpAdvertisedItem, CatalogRequest as McpCatalogRequest, McpAllowlist,
273 McpAllowlistItem, McpCatalog, McpCatalogItem, McpCatalogServer, McpItemKind,
274 MCP_ALLOWLIST_SCHEMA_VERSION,
275};
276pub use mcp_card::{fetch_server_card, load_server_card_from_path, CardError};
277pub use mcp_host::{
278 cache_stats as mcp_host_cache_stats, set_allowlist as set_mcp_host_allowlist,
279 AllowlistDecision as McpHostAllowlistDecision, AllowlistGuard as McpHostAllowlistGuard,
280 BreakerState as McpHostBreakerState, CacheStats as McpHostCacheStats, McpHostStatus,
281 SpawnOptions as McpHostSpawnOptions, SupervisionPolicy as McpHostSupervisionPolicy,
282};
283pub use mcp_registry::{
284 active_handle as mcp_active_handle, ensure_active as mcp_ensure_active,
285 get_registration as mcp_get_registration, install_active as mcp_install_active,
286 is_registered as mcp_is_registered, register_servers as mcp_register_servers,
287 release as mcp_release, reset as mcp_reset_registry, snapshot_status as mcp_snapshot_status,
288 sweep_expired as mcp_sweep_expired, RegisteredMcpServer, RegistryStatus,
289};
290pub use mcp_server::{
291 take_mcp_serve_metadata, take_mcp_serve_prompts, take_mcp_serve_registry,
292 take_mcp_serve_resource_templates, take_mcp_serve_resources, tool_registry_to_mcp_tools,
293 McpServer, McpServerMetadata,
294};
295pub use metadata::register_metadata_builtins;
296pub use observability::audit::{audit_events as audit_obs_events, AuditFinding, AuditFindingKind};
297pub use observability::execution_scope::{
298 current_execution_scope, enter_execution_scope, mint_execution_scope, ExecutionScopeGuard,
299};
300pub use observability::request_id::{current_request_id, enter_request_id, RequestIdScopeGuard};
301pub use orchestration::{
302 benchmark_adapted_replay_pair, benchmark_replay_trace, build_replay_benchmark_report,
303 OpenCodeJsonlAdapter, ReplayBenchmarkCloudIngest, ReplayBenchmarkError,
304 ReplayBenchmarkFixtureReceipt, ReplayBenchmarkFixtureReport, ReplayBenchmarkMetrics,
305 ReplayBenchmarkReport, ReplayBenchmarkSuiteIdentity, ReplayBenchmarkSummary,
306 ReplayCategoryMetric, ReplayDebuggingProxyMetrics, ReplayRuntimeCostMetrics,
307 ReplayTraceAdapter, OPENCODE_JSONL_ADAPTER_ID, OPENCODE_JSONL_ADAPTER_SCHEMA_VERSION,
308 REPLAY_BENCHMARK_CLOUD_INGEST_KIND, REPLAY_BENCHMARK_REPORT_SCHEMA_VERSION,
309};
310pub use orchestration::{
311 canonicalize_run, first_divergence, run_replay_oracle_trace, ReplayAllowlistRule,
312 ReplayDivergence, ReplayExpectation, ReplayOracleError, ReplayOracleReport, ReplayOracleTrace,
313 ReplayTraceRun, ReplayTraceRunCounts, REPLAY_TRACE_SCHEMA_VERSION,
314};
315pub use orchestration::{
316 install_handoff_routes, snapshot_handoff_routes, HandoffRouteConfig,
317 HandoffRouteDecisionRecord, HandoffRouteTargetConfig,
318};
319pub use personas::{
320 disable_persona, fire_schedule as fire_persona_schedule, fire_trigger as fire_persona_trigger,
321 format_ms as format_persona_ms, now_ms as persona_now_ms, parse_rfc3339_ms as parse_persona_ms,
322 pause_persona, persona_status, record_persona_spend, register_persona_supervision_sink,
323 register_persona_value_sink, report_repair_worker_status, restore_persona_checkpoint,
324 resume_persona, PersonaAssignmentStatus, PersonaBudgetPolicy, PersonaBudgetStatus,
325 PersonaCheckpointAction, PersonaCheckpointRestoreOutcome, PersonaCheckpointRestoreRequest,
326 PersonaCheckpointResume, PersonaCheckpointUpdate, PersonaHandoffInboxItem, PersonaLease,
327 PersonaLifecycleState, PersonaQueuePositionUpdate, PersonaQueuedWork, PersonaReceiptUpdate,
328 PersonaRepairWorkerLifecycle, PersonaRepairWorkerStatusUpdate, PersonaRunCost,
329 PersonaRunReceipt, PersonaRuntimeBinding, PersonaStatus, PersonaSupervisionEvent,
330 PersonaSupervisionSink, PersonaSupervisionSinkRegistration, PersonaTriggerEnvelope,
331 PersonaValueEvent, PersonaValueEventKind, PersonaValueReceipt, PersonaValueSink,
332 PersonaValueSinkRegistration, StageDecl, StageExit, PERSONA_RUNTIME_TOPIC,
333};
334pub use provenance::{
335 build_signed_receipt, load_or_generate_agent_signing_key, verify_receipt, ProvenanceReceipt,
336 ReceiptBuildOptions, ReceiptVerificationReport,
337};
338pub use receipts::{
339 Receipt, ReceiptSink, ReceiptStatus, ReceiptValidationError, RedactingReceiptSink,
340 RedactionClass, RECEIPT_SCHEMA_ID, RECEIPT_SCHEMA_JSON, RECEIPT_SCHEMA_VERSION,
341};
342pub use record_filter::{normalize_record_filter_expression, CompiledRecordFilter};
343pub use runtime_limits::{
344 RuntimeLimitDescription, RuntimeLimitEntry, RuntimeLimits, RuntimeLimitsReport,
345 RUNTIME_LIMIT_DESCRIPTIONS,
346};
347pub use schema::json_to_vm_value;
348pub use sessions::{
349 CreateSession, ExpireSession, Session, SessionAttributes, SessionError, SessionStore,
350 TouchSession, SESSIONS_TOPIC,
351};
352pub use stdlib::fs::ignore_policy;
358pub use stdlib::hitl::{
359 append_hitl_response, ApprovalRequest, HitlHostResponse, HITL_APPROVALS_TOPIC,
360 HITL_DUAL_CONTROL_TOPIC, HITL_ESCALATIONS_TOPIC, HITL_QUESTIONS_TOPIC,
361};
362pub use stdlib::host::turn_cache as host_turn_cache;
364pub use stdlib::host::{
365 clear_host_call_bridge, dispatch_host_operation, host_call_ready, set_host_call_bridge,
366 HostCallBridge, HostCallDispatchFuture,
367};
368pub use stdlib::http_response::{
369 parse_envelope as parse_http_envelope, HttpEnvelope, HttpHeaderValue, WsUpgradeSpec,
370 HTTP_RESPONSE_TAG_KEY, HTTP_RESPONSE_TAG_VERSION,
371};
372#[cfg(feature = "postgres")]
373pub use stdlib::install_shared_pool_registry;
374pub use stdlib::io::{
375 reserve_stdio_for_current_thread, set_stdout_passthrough, take_stderr_buffer,
376 StdioReservationGuard,
377};
378pub use stdlib::long_running::cancel_handle as cancel_long_running_handle;
379pub use stdlib::observability::install_default_backend as install_obs_default_backend;
380pub use stdlib::secret_scan::{
381 append_secret_scan_audit, audit_secret_scan_active, scan_content as secret_scan_content,
382 SecretFinding, SECRET_SCAN_AUDIT_TOPIC,
383};
384pub use stdlib::template::{
385 lookup_prompt_consumers, lookup_prompt_span, prompt_render_indices, record_prompt_render_index,
386 PromptSourceSpan, PromptSpanKind,
387};
388pub use stdlib::waitpoint::{
389 process_waitpoint_resume_event, service_waitpoints_once, WAITPOINT_RESUME_TOPIC,
390};
391pub use stdlib::workflow_messages::{
392 workflow_pause_for_base, workflow_publish_query_for_base, workflow_query_for_base,
393 workflow_respond_update_for_base, workflow_resume_for_base, workflow_signal_for_base,
394 workflow_update_for_base, WorkflowMailboxState,
395};
396pub use stdlib::{
397 register_agent_stdlib, register_core_stdlib, register_io_stdlib, register_vm_stdlib,
398};
399pub use store::register_store_builtins;
400pub use tenant::{
401 tenant_event_topic_prefix, tenant_secret_namespace, tenant_topic, validate_tenant_id, ApiKeyId,
402 TenantApiKeyRecord, TenantBudget, TenantEventLog, TenantRecord, TenantRegistrySnapshot,
403 TenantResolutionError, TenantScope, TenantSecretProvider, TenantStatus, TenantStore,
404 TENANT_EVENT_TOPIC_PREFIX, TENANT_REGISTRY_DIR, TENANT_REGISTRY_FILE,
405 TENANT_SECRET_NAMESPACE_PREFIX,
406};
407pub use triggers::{
408 append_dispatch_cancel_request, begin_in_flight, binding_autonomy_budget_would_exceed,
409 binding_budget_would_exceed, binding_version_as_of, classify_trigger_dlq_error,
410 clear_dispatcher_state, clear_orchestrator_budget, clear_trigger_registry, drain,
411 dynamic_deregister, dynamic_register, expected_predicate_cost_usd_micros, finish_in_flight,
412 install_manifest_triggers, install_orchestrator_budget, micros_to_usd,
413 note_autonomous_decision, note_orchestrator_budget_cost, orchestrator_budget_would_exceed,
414 parse_flow_control_duration, pause, pin_trigger_binding, provider_metadata,
415 record_predicate_cost_sample, redact_headers, register_provider_schemas,
416 registered_provider_metadata, registered_provider_schema_names, reset_binding_budget_windows,
417 reset_provider_catalog, resolve_live_or_as_of, resolve_live_trigger_binding,
418 resolve_trigger_binding_as_of, resume, run_trigger_harness_fixture, scheduler_in_flight_by_key,
419 scheduler_ready_stats_by_key, snapshot_dispatcher_stats, snapshot_orchestrator_budget,
420 snapshot_trigger_bindings, unpin_trigger_binding, usd_to_micros, worker_claims_topic_name,
421 worker_job_topic_name, worker_response_topic_name, ClaimedWorkerJob, DispatchCancelRequest,
422 DispatchError, DispatchOutcome, DispatchStatus, Dispatcher, DispatcherDrainReport,
423 DispatcherStatsSnapshot, ExtensionProviderPayload, FairnessKey, HeaderRedactionPolicy,
424 InboxIndex, OrchestratorBudgetConfig, OrchestratorBudgetSnapshot, ProviderCatalog,
425 ProviderCatalogError, ProviderId, ProviderMetadata, ProviderOutboundMethod, ProviderPayload,
426 ProviderRuntimeMetadata, ProviderSchema, ProviderSecretRequirement, ReadyKeyStats,
427 RecordedTriggerBinding, RetryPolicy, SchedulableJob, SchedulerKeyStat, SchedulerPolicy,
428 SchedulerSnapshot, SchedulerState, SchedulerStrategy, SignatureStatus,
429 SignatureVerificationMetadata, StreamEventPayload, TenantId, TraceId, TriggerBatchConfig,
430 TriggerBindingSnapshot, TriggerBindingSource, TriggerBindingSpec,
431 TriggerBudgetExhaustionStrategy, TriggerConcurrencyConfig, TriggerDebounceConfig,
432 TriggerDispatchOutcome, TriggerEvent, TriggerEventId, TriggerExpressionSpec,
433 TriggerFlowControlConfig, TriggerHandlerSpec, TriggerHarnessResult, TriggerId,
434 TriggerMetricsSnapshot, TriggerPredicateSpec, TriggerPriorityOrderConfig,
435 TriggerRateLimitConfig, TriggerRegistryError, TriggerRetryConfig, TriggerSingletonConfig,
436 TriggerState, TriggerThrottleConfig, WorkerQueue, WorkerQueueClaimHandle,
437 WorkerQueueEnqueueReceipt, WorkerQueueInspectSnapshot, WorkerQueueJob, WorkerQueueJobState,
438 WorkerQueuePriority, WorkerQueueResponseRecord, WorkerQueueState, WorkerQueueSummary,
439 DEFAULT_INBOX_RETENTION_DAYS, DEFAULT_STARVATION_AGE_MS, TRIGGERS_LIFECYCLE_TOPIC,
440 TRIGGER_ATTEMPTS_TOPIC, TRIGGER_CANCEL_REQUESTS_TOPIC, TRIGGER_DLQ_TOPIC,
441 TRIGGER_INBOX_CLAIMS_TOPIC, TRIGGER_INBOX_ENVELOPES_TOPIC, TRIGGER_INBOX_LEGACY_TOPIC,
442 TRIGGER_INBOX_OBSERVABILITY_TOPIC, TRIGGER_OPERATION_AUDIT_TOPIC, TRIGGER_OUTBOX_TOPIC,
443 TRIGGER_TEST_FIXTURES, WORKER_QUEUE_CATALOG_TOPIC,
444};
445pub use trust_graph::{
446 append_active_scope_attenuation_alert, append_active_trust_record,
447 append_scope_attenuation_alert, append_trust_record, export_trust_chain,
448 group_trust_records_by_trace, policy_for_agent, policy_for_autonomy_tier,
449 query_trust_graph_records, query_trust_records, resolve_agent_autonomy_tier,
450 summarize_trust_records, topic_for_agent, trust_score_for, verify_trust_chain, AutonomyTier,
451 TrustAgentSummary, TrustChainExport, TrustChainExportMetadata, TrustChainExportProducer,
452 TrustChainReport, TrustGraphRecord, TrustOutcome, TrustQueryFilters, TrustRecord,
453 TrustRecordActionKind, TrustScore, TrustTraceGroup, METADATA_KEY_ACTOR_CHAIN,
454 METADATA_KEY_ACTOR_CHAIN_ALERT, METADATA_KEY_EFFECTS_GRANT, METADATA_KEY_EFFECTS_USED,
455 METADATA_KEY_PARENT_RECORD_ID, OPENTRUSTGRAPH_ACCEPTED_SCHEMAS, OPENTRUSTGRAPH_CHAIN_SCHEMA_V0,
456 OPENTRUSTGRAPH_SCHEMA_V0, OPENTRUSTGRAPH_SCHEMA_V0_1, TRUST_ACTION_RELEASE,
457 TRUST_GRAPH_GLOBAL_TOPIC, TRUST_GRAPH_LEGACY_GLOBAL_TOPIC, TRUST_GRAPH_LEGACY_TOPIC_PREFIX,
458 TRUST_GRAPH_RECORDS_TOPIC, TRUST_GRAPH_TOPIC_PREFIX,
459};
460pub use value::*;
461pub use vm::*;
462
463#[cfg(feature = "vm-bench-internals")]
464#[doc(hidden)]
465pub mod bench_internals;
466
467pub fn compile_source(source: &str) -> Result<Chunk, String> {
472 let program = harn_parser::check_source_strict(source).map_err(|e| e.to_string())?;
473 Compiler::new().compile(&program).map_err(|e| e.to_string())
474}
475
476pub fn compile_source_named(source: &str, pipeline_name: &str) -> Result<Chunk, String> {
481 let program = harn_parser::check_source_strict(source).map_err(|e| e.to_string())?;
482 let has_pipeline = program.iter().any(|sn| {
483 let (_, inner) = harn_parser::peel_attributes(sn);
484 matches!(&inner.node, harn_parser::Node::Pipeline { name, .. } if name == pipeline_name)
485 });
486 if !has_pipeline {
487 return Err(format!("no pipeline named `{pipeline_name}` in source"));
488 }
489 Compiler::new()
490 .compile_named(&program, pipeline_name)
491 .map_err(|e| e.to_string())
492}
493
494pub struct SchemaAliasResolver {
502 compiler: compiler::Compiler,
503}
504
505impl SchemaAliasResolver {
506 pub fn empty() -> Self {
509 Self {
510 compiler: compiler::Compiler::new(),
511 }
512 }
513
514 pub fn from_program(program: &[harn_parser::SNode]) -> Self {
517 let mut compiler = compiler::Compiler::new();
518 compiler.collect_type_aliases(program);
519 Self { compiler }
520 }
521
522 pub fn json_schema_for_type_expr(
525 &self,
526 type_expr: &harn_parser::TypeExpr,
527 ) -> Option<serde_json::Value> {
528 let expanded = self.compiler.expand_alias(type_expr);
529 let schema = compiler::Compiler::type_expr_to_schema_value(&expanded)?;
530 let json_schema = schema::schema_to_json_schema_value(&schema).ok()?;
531 Some(llm::vm_value_to_json(&json_schema))
532 }
533
534 pub fn json_schema_for_typed_params(
537 &self,
538 params: &[harn_parser::TypedParam],
539 ) -> serde_json::Value {
540 let mut properties = serde_json::Map::new();
541 let mut required = Vec::new();
542
543 for param in params {
544 let param_schema = param
545 .type_expr
546 .as_ref()
547 .and_then(|type_expr| self.json_schema_for_type_expr(type_expr))
548 .unwrap_or_else(|| serde_json::json!({}));
549 if param.default_value.is_none() {
550 required.push(serde_json::Value::String(param.name.clone()));
551 }
552 properties.insert(param.name.clone(), param_schema);
553 }
554
555 let mut schema = serde_json::Map::new();
556 schema.insert(
557 "type".to_string(),
558 serde_json::Value::String("object".to_string()),
559 );
560 schema.insert(
561 "properties".to_string(),
562 serde_json::Value::Object(properties),
563 );
564 if !required.is_empty() {
565 schema.insert("required".to_string(), serde_json::Value::Array(required));
566 }
567 serde_json::Value::Object(schema)
568 }
569}
570
571pub fn json_schema_for_type_expr(type_expr: &harn_parser::TypeExpr) -> Option<serde_json::Value> {
575 SchemaAliasResolver::empty().json_schema_for_type_expr(type_expr)
576}
577
578pub fn json_schema_for_typed_params(params: &[harn_parser::TypedParam]) -> serde_json::Value {
579 SchemaAliasResolver::empty().json_schema_for_typed_params(params)
580}
581
582#[cfg(test)]
583mod schema_alias_resolver_tests {
584 use super::*;
585
586 fn fn_params_schema(src: &str) -> serde_json::Value {
587 let program = harn_parser::parse_source(src).expect("parse test source");
588 let resolver = SchemaAliasResolver::from_program(&program);
589 for node in &program {
590 let (_, inner) = harn_parser::peel_attributes(node);
591 if let harn_parser::Node::FnDecl { params, .. } = &inner.node {
592 return resolver.json_schema_for_typed_params(params);
593 }
594 }
595 panic!("no fn decl in test source");
596 }
597
598 #[test]
599 fn named_shape_alias_projects_like_inline_shape() {
600 let inline = fn_params_schema("pub fn f(p: {kind: string, path: string}) {}");
601 let aliased =
602 fn_params_schema("type Src = {kind: string, path: string}\npub fn f(p: Src) {}");
603 assert_eq!(
604 aliased, inline,
605 "a named shape alias must project the same inputSchema as its inline shape",
606 );
607 assert_ne!(
608 aliased["properties"]["p"],
609 serde_json::json!({}),
610 "the alias parameter must not erase to an empty schema",
611 );
612 }
613
614 #[test]
615 fn literal_union_alias_projects_json_enum() {
616 let schema = fn_params_schema("type Kind = \"local\" | \"ssh\"\npub fn f(p: Kind) {}");
617 let p = &schema["properties"]["p"];
618 assert_eq!(p["type"], "string");
619 assert_eq!(p["enum"], serde_json::json!(["local", "ssh"]));
620 }
621
622 #[test]
623 fn unknown_named_type_still_erases_to_empty() {
624 let schema = fn_params_schema("pub fn f(p: Unknown) {}");
626 assert_eq!(schema["properties"]["p"], serde_json::json!({}));
627 }
628}
629
630fn reset_llm_state_for_thread_reset() {
631 llm::reset_llm_state();
632 #[cfg(test)]
633 reset_thread_local_state_test_hooks::before_llm_global_reset();
634 llm::reset_rate_limit_registry();
637 llm_config::clear_user_overrides();
638 llm_config::clear_runtime_provider_endpoint_overrides();
639}
640
641#[cfg(test)]
642mod reset_thread_local_state_test_hooks {
643 use std::sync::{Arc, Mutex, OnceLock};
644
645 type Hook = Arc<dyn Fn() + Send + Sync + 'static>;
646
647 static BEFORE_LLM_GLOBAL_RESET: OnceLock<Mutex<Option<Hook>>> = OnceLock::new();
648
649 fn before_llm_global_reset_hook() -> &'static Mutex<Option<Hook>> {
650 BEFORE_LLM_GLOBAL_RESET.get_or_init(|| Mutex::new(None))
651 }
652
653 pub(crate) struct HookGuard;
654
655 impl Drop for HookGuard {
656 fn drop(&mut self) {
657 let mut hook = before_llm_global_reset_hook()
658 .lock()
659 .unwrap_or_else(std::sync::PoisonError::into_inner);
660 *hook = None;
661 }
662 }
663
664 pub(crate) fn install_before_llm_global_reset(hook: Hook) -> HookGuard {
665 let mut slot = before_llm_global_reset_hook()
666 .lock()
667 .unwrap_or_else(std::sync::PoisonError::into_inner);
668 *slot = Some(hook);
669 HookGuard
670 }
671
672 pub(crate) fn before_llm_global_reset() {
673 let hook = before_llm_global_reset_hook()
674 .lock()
675 .unwrap_or_else(std::sync::PoisonError::into_inner)
676 .clone();
677 if let Some(hook) = hook {
678 hook();
679 }
680 }
681}
682
683pub fn reset_thread_local_state() {
685 #[cfg(test)]
686 {
687 let _guard = llm::env_guard();
692 reset_llm_state_for_thread_reset();
693 }
694 #[cfg(not(test))]
695 reset_llm_state_for_thread_reset();
696
697 http::reset_http_state();
698 channels::reset_channel_state();
699 event_log::reset_active_event_log();
700 egress::clear_explicit_egress_policy_requirement_for_host();
701 egress::clear_ssrf_guard_requirement_for_host();
702 stdlib::reset_stdlib_state();
703 connectors::clear_active_connector_clients();
704 orchestration::clear_runtime_hooks();
705 orchestration::clear_file_edit_queue();
706 orchestration::clear_execution_policy_stacks();
707 orchestration::clear_command_policies();
708 orchestration::clear_pipeline_on_finish();
709 orchestration::reset_lifecycle_receipt_registry();
710 orchestration::agent_inbox::reset();
711 tool_call_cancellations::reset_registry();
712 redact::clear_policy_stack();
713 security::reset_thread_state();
714 triggers::clear_dispatcher_state();
715 triggers::clear_trigger_registry();
716 events::reset_event_sinks();
717 tracing::set_tracing_enabled(false);
718 tracing::reset_tracing();
719 agent_events::reset_all_sinks();
731 agent_sessions::reset_session_store();
732 mcp_registry::reset();
733 mcp_host::reset_for_tests();
734 call_budget::reset_call_budget_state();
735 clock_mock::leak_audit::reset();
736}
737
738#[cfg(test)]
739mod reset_leak_tests {
740 use super::*;
746 use crate::value::VmValue;
747
748 #[test]
749 fn reset_drains_pending_file_edit_notifications() {
750 orchestration::queue_file_edited("stale.harn", serde_json::json!({"operation": "write"}));
751
752 reset_thread_local_state();
753
754 assert!(
755 orchestration::drain_file_edits().is_empty(),
756 "a later VM run must not receive file edits queued by the previous run"
757 );
758 }
759
760 #[test]
769 fn builtin_profile_recording_ends_with_its_run_not_with_a_global_reset() {
770 let _lock = builtin_profile::test_lock()
771 .lock()
772 .unwrap_or_else(std::sync::PoisonError::into_inner);
773 let recording = builtin_profile::enable();
774 builtin_profile::record("run_shell", std::time::Duration::from_millis(5));
775 assert!(builtin_profile::is_enabled());
776 assert!(!builtin_profile::snapshot().is_empty());
777
778 reset_thread_local_state();
779
780 assert!(
781 builtin_profile::is_enabled(),
782 "an unrelated global reset must not disarm an in-flight profiled run"
783 );
784 assert!(
785 !builtin_profile::snapshot().is_empty(),
786 "an unrelated global reset must not drop totals the run still owns"
787 );
788
789 drop(recording);
790
791 assert!(
792 !builtin_profile::is_enabled(),
793 "a profiled run must not leave the recorder on for the next one"
794 );
795 assert!(
796 builtin_profile::snapshot().is_empty(),
797 "builtin totals must be empty once the run's guard drops"
798 );
799 }
800
801 #[test]
806 fn reset_drains_session_changed_paths() {
807 let session = "sess-leak";
808 agent_sessions::open_or_create(Some(session.to_string()));
809 agent_sessions::record_session_changed_path(session, "/tmp/written-by-a-dead-run.txt");
810 assert!(!agent_sessions::session_changed_paths(session).is_empty());
811 reset_thread_local_state();
812 assert!(
813 agent_sessions::session_changed_paths(session).is_empty(),
814 "a receipt must not inherit an abandoned session's writes"
815 );
816 }
817
818 #[test]
819 fn reset_drains_agent_inbox() {
820 orchestration::agent_inbox::reset();
821 orchestration::agent_inbox::push("sess-2660", "note", "leak", "test");
822 assert!(orchestration::agent_inbox::session_count() > 0);
823 reset_thread_local_state();
824 assert_eq!(
825 orchestration::agent_inbox::session_count(),
826 0,
827 "agent_inbox must be empty after reset"
828 );
829 }
830
831 #[test]
832 fn reset_drains_tool_call_cancellation_registry() {
833 tool_call_cancellations::reset_registry();
834 let registered = tool_call_cancellations::register("sess-2660", "call-1", "tool");
837 if let Some((_handle, guard)) = registered {
838 std::mem::forget(guard);
839 }
840 assert!(tool_call_cancellations::registry_len() > 0);
841 reset_thread_local_state();
842 assert_eq!(
843 tool_call_cancellations::registry_len(),
844 0,
845 "tool-call cancellation registry must be empty after reset"
846 );
847 }
848
849 #[test]
850 fn reset_drains_routing_policy_registry() {
851 llm::routing::clear_policy_registry();
852 let mut config: crate::value::DictMap = crate::value::DictMap::new();
853 config.insert(
854 crate::value::intern_key("chain"),
855 VmValue::List(std::sync::Arc::new(vec![VmValue::String(
856 arcstr::ArcStr::from("mock:mock"),
857 )])),
858 );
859 llm::routing::build_routing_policy(&config).expect("intern a routing policy");
860 assert!(llm::routing::policy_registry_len() > 0);
861 reset_thread_local_state();
862 assert_eq!(
863 llm::routing::policy_registry_len(),
864 0,
865 "routing policy registry must be empty after reset"
866 );
867 }
868
869 #[test]
870 fn reset_holds_llm_env_guard_while_wiping_llm_globals() {
871 let observed = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
872 let observed_hook = std::sync::Arc::clone(&observed);
873 let _hook = reset_thread_local_state_test_hooks::install_before_llm_global_reset(
874 std::sync::Arc::new(move || {
875 assert!(
876 matches!(
877 llm::env_lock().try_lock(),
878 Err(std::sync::TryLockError::WouldBlock)
879 ),
880 "reset_thread_local_state must hold env_guard before wiping LLM globals"
881 );
882 observed_hook.store(true, std::sync::atomic::Ordering::SeqCst);
883 }),
884 );
885
886 reset_thread_local_state();
887 assert!(
888 observed.load(std::sync::atomic::Ordering::SeqCst),
889 "LLM global reset hook should have run"
890 );
891 }
892}