Skip to main content

harn_vm/
lib.rs

1#![recursion_limit = "256"]
2#![allow(clippy::result_large_err, clippy::cloned_ref_to_slice_refs)]
3//! # harn-vm
4//!
5//! The Harn compiler, virtual machine, standard library, provider/LLM layer,
6//! orchestration runtime, and host bridge.
7//!
8//! ## Stability
9//!
10//! This crate is consumed both by the in-tree surfaces (`harn-cli`,
11//! `harn-serve`, the LSP and DAP) and by external embedders. The intended
12//! embedding entry points are `Vm`, `Harness`, `compile_source`, and the
13//! `llm`, `orchestration`, `agent_events`, `agent_sessions`, `config`, and
14//! `security` modules. Other public items exist primarily for in-workspace use
15//! and may change between minor releases; anything marked `#[doc(hidden)]` is
16//! an implementation detail with no stability guarantee. The crate follows the
17//! workspace version and is pre-1.0, so the public surface may still evolve.
18
19/// Re-export of the unified clock substrate so downstream crates (CLI,
20/// orchestrator, and cloud runtimes) can depend on a single canonical `Clock`
21/// trait without each adding `harn-clock` as a direct dependency.
22pub 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
151/// Initialize process-wide assets whose construction should happen before an
152/// embedding host enters an async request or VM execution stack.
153///
154/// New embedding hosts should call [`initialize_runtime`] instead so startup
155/// also validates the Harn-owned environment namespace. This asset-only
156/// operation remains for compatibility, and VM construction retains it as a
157/// fallback for embedders without an explicit bootstrap phase.
158pub fn initialize_runtime_assets() {
159    secret_patterns::initialize_default_secret_patterns();
160}
161
162/// Validate the Harn-owned environment namespace and initialize process-wide
163/// runtime assets through the same bootstrap boundary used by the CLI.
164pub fn initialize_runtime() -> Result<(), environment_registry::EnvironmentValidationError> {
165    environment_registry::validate_startup_environment()?;
166    initialize_runtime_assets();
167    Ok(())
168}
169
170/// Crate-wide deterministic clock mock used by stdlib time builtins, the
171/// trigger dispatcher, the cron scheduler, and Rust-side tests. Re-exports
172/// the long-lived implementation under `triggers::test_util::clock` so all
173/// callers go through one source of truth.
174pub 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    /// Runtime audit for capabilities that observe real wall-clock or
182    /// monotonic time while a testbench mock is installed. See the module
183    /// docs for the full design.
184    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};
352/// The single owner of ignore policy for every Harn filesystem walk.
353///
354/// Re-exported so embedders that enumerate files on behalf of Harn scripts
355/// (today: the `harn-hostlib` deterministic-tool builtins) skip exactly the
356/// same paths the in-VM builtins do.
357pub 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};
362/// Per-turn memo for turn-stable host reads. See [`stdlib::host::turn_cache`].
363pub 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
467/// Lex, parse, type-check, and compile source to bytecode in one call.
468/// Bails on the first type error. For callers that need diagnostics
469/// rather than early exit, use `harn_parser::check_source` directly
470/// and then call `Compiler::new().compile(&program)`.
471pub 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
476/// Same as [`compile_source`] but compiles a specific named pipeline as
477/// the program entry point instead of the default-pipeline-or-first
478/// selection rule. Returns a runtime error when no pipeline with
479/// `pipeline_name` exists in the source.
480pub 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
494/// Lowers Harn `TypeExpr`s to JSON Schema with `type`-alias EXPANSION, built once
495/// from a parsed program's alias declarations. Without expansion, a tool parameter
496/// typed as a user alias (`p: EvalSource`, `p: FunnelStage`) erases to an empty
497/// `{}` schema because the low-level lowering only recognizes built-in type names —
498/// the exporter must first resolve the alias to its underlying shape/union (a
499/// literal-union alias then round-trips as a JSON `enum`). Cycle-safe via the same
500/// `expand_alias` guard the compiler and typechecker share.
501pub struct SchemaAliasResolver {
502    compiler: compiler::Compiler,
503}
504
505impl SchemaAliasResolver {
506    /// A resolver with no aliases in scope — lowering is identical to the raw
507    /// (unexpanded) form, so `Named(alias)` still lowers to `{}` when unknown.
508    pub fn empty() -> Self {
509        Self {
510            compiler: compiler::Compiler::new(),
511        }
512    }
513
514    /// Collect every `type` alias declared in `program`, so named references in
515    /// tool signatures resolve to their bodies.
516    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    /// JSON Schema for one `TypeExpr`, expanding any named alias first. `None`
523    /// when the (expanded) type has no JSON-Schema form (function types, ...).
524    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    /// JSON Schema `object` for a parameter list (a served tool's `inputSchema`),
535    /// expanding aliases per parameter.
536    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
571/// Raw lowering with no program aliases in scope. Prefer
572/// [`SchemaAliasResolver::from_program`] when serving a module so named-alias
573/// parameters resolve instead of erasing to `{}`.
574pub 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        // No alias declared: unchanged behavior — an unknown named type lowers to {}.
625        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    // This full wipe is necessary between Harn programs to clear durable
635    // cooldowns that would otherwise stall a later run under a paused clock.
636    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
683/// Reset all thread-local state that can leak between test runs.
684pub fn reset_thread_local_state() {
685    #[cfg(test)]
686    {
687        // `reset_thread_local_state` is also used by in-process unit tests. It
688        // clears process-global LLM config/rate-limit state, so share the same
689        // lock used by LLM env tests; otherwise a sibling reset can erase a
690        // parked rate-limit test's registry while the test still owns a permit.
691        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    // `builtin_profile` is deliberately NOT reset here. Its recorder is
720    // process-global (`static ENABLED` / `static TOTALS`), and this function
721    // runs from ~150 test setups and from production entry points like
722    // `execute_conformance_source` and the orchestrator lifecycle. Every one
723    // of those calls disarmed the recorder that a concurrently running
724    // profiled run had just enabled, so `harn run --profile` reported
725    // `vm/residual 100%` and named nothing. `builtin_profile::enable()`
726    // already discards the previous run's totals, so the profiling entry
727    // point owns the lifecycle without help from here. Same reasoning as
728    // `llm::rate_limit::reset_runtime_rate_limit_overrides` and the
729    // `long_running::reset_state` exclusion in `stdlib::reset_stdlib_state`.
730    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    //! Regression coverage for harn#2660: process-/thread-global
741    //! registries that accumulated one entry per test because they were
742    //! never drained by `reset_thread_local_state`. Each case populates a
743    //! registry through its real entry point, runs the reset, and asserts
744    //! the registry is empty again.
745    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    /// The recorder is enabled per RUN but lives for the PROCESS, so an
761    /// embedder that runs one script with `--profile` and the next without it
762    /// would keep paying for bookkeeping nobody reads and fold the second
763    /// run's builtins into the first run's totals. Enablement therefore ends
764    /// with the run that asked for it — the guard `enable()` returns — and NOT
765    /// in `reset_thread_local_state`, which fires from ~150 test setups and
766    /// from production entry points that know nothing about an in-flight
767    /// profiled run.
768    #[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    /// The changed-path map is the authoritative source for a sub-agent's
802    /// `files_written` receipt, is process-global, and is drained only at
803    /// teardown — which a session that errors never reaches. A later session
804    /// reusing the id would report writes it never made.
805    #[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        // Leak the guard so the entry survives until the reset runs —
835        // this mirrors a dispatch abandoned mid-flight.
836        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}