1use super::errors::PackageError;
2use super::*;
3
4pub(crate) fn manifest_capabilities(
5 manifest: &Manifest,
6) -> Option<&harn_vm::llm::capabilities::CapabilitiesFile> {
7 manifest.capabilities.as_ref()
8}
9
10pub(crate) fn is_empty_capabilities(file: &harn_vm::llm::capabilities::CapabilitiesFile) -> bool {
11 file.provider.is_empty() && file.provider_family.is_empty()
12}
13
14pub fn validate_runtime_manifest_extensions(anchor: &Path) -> Result<(), PackageError> {
15 let Some((manifest, _manifest_dir)) = load_nearest_manifest(anchor).into_result()? else {
16 return Ok(());
17 };
18 validate_handoff_routes(&manifest.handoff_routes, &manifest)?;
19 validate_contributions(&manifest)
20}
21
22pub fn try_load_runtime_extensions(anchor: &Path) -> Result<RuntimeExtensions, PackageError> {
25 ensure_dependencies_materialized(anchor)?;
26 let Some((root_manifest, manifest_dir)) = load_nearest_manifest(anchor).into_result()? else {
27 return Ok(RuntimeExtensions::default());
28 };
29
30 let mut llm = harn_vm::llm_config::ProvidersConfig::default();
31 let mut capabilities = harn_vm::llm::capabilities::CapabilitiesFile::default();
32 let mut hooks = Vec::new();
33 let mut triggers = Vec::new();
34
35 llm.merge_from(&root_manifest.llm);
36 if let Some(file) = manifest_capabilities(&root_manifest) {
37 merge_capability_overrides(&mut capabilities, file);
38 }
39 hooks.extend(resolved_hooks_from_manifest(&root_manifest, &manifest_dir));
40 triggers.extend(resolved_triggers_from_manifest(
41 &root_manifest,
42 &manifest_dir,
43 ));
44 let handoff_routes = root_manifest.handoff_routes.clone();
45 validate_handoff_routes(&handoff_routes, &root_manifest)?;
46 let mut provider_connectors =
47 resolved_provider_connectors_from_manifest(&root_manifest, &manifest_dir);
48 let package_snapshot =
49 dependency_package_snapshot(&root_manifest, &manifest_dir)?.map(Arc::new);
50 if let Some(snapshot) = package_snapshot.as_ref() {
51 provider_connectors.extend(installed_package_provider_connectors(
52 snapshot,
53 snapshot.packages_root(),
54 )?);
55 }
56 provider_connectors = dedupe_provider_connectors(provider_connectors);
57 let root_manifest_path = manifest_dir.join(MANIFEST);
58 let runtime_personas = resolve_runtime_personas(
59 root_manifest.clone(),
60 root_manifest_path.clone(),
61 manifest_dir.clone(),
62 package_snapshot,
63 )?;
64 triggers.extend(installed_persona_trigger_configs(&runtime_personas)?);
65
66 Ok(RuntimeExtensions {
67 root_manifest_path: Some(root_manifest_path),
68 root_manifest_dir: Some(manifest_dir),
69 root_manifest: Some(root_manifest),
70 runtime_personas,
71 llm: (!llm.is_empty()).then_some(llm),
72 capabilities: (!is_empty_capabilities(&capabilities)).then_some(capabilities),
73 hooks,
74 triggers,
75 handoff_routes,
76 provider_connectors,
77 })
78}
79
80pub fn try_load_runtime_extensions_from_manifest(
85 manifest_path: &Path,
86) -> Result<Option<RuntimeExtensions>, PackageError> {
87 let manifest_path = if manifest_path.is_dir() {
88 manifest_path.join(MANIFEST)
89 } else {
90 manifest_path.to_path_buf()
91 };
92 if manifest_path.extension().and_then(|value| value.to_str()) == Some("harn") {
93 return Ok(None);
94 }
95 if manifest_path.file_name() != Some(OsStr::new(MANIFEST)) {
96 return Ok(None);
97 }
98 if read_manifest_from_path(&manifest_path).is_err() {
99 return Ok(None);
100 }
101 try_load_runtime_extensions(&manifest_path).map(Some)
102}
103
104fn installed_package_provider_connectors(
105 snapshot: &harn_modules::package_snapshot::PackageSnapshot,
106 packages_dir: &Path,
107) -> Result<Vec<ResolvedProviderConnectorConfig>, PackageError> {
108 let lock = LockFile::load(snapshot.lock_path())?.ok_or_else(|| {
109 PackageError::Lockfile(format!(
110 "published package generation is missing {}",
111 snapshot.lock_path().display()
112 ))
113 })?;
114 let mut providers = Vec::new();
115 for entry in &lock.packages {
116 validate_package_alias(&entry.name)?;
117 let package_dir = packages_dir.join(&entry.name);
118 if package_dir.is_dir() {
119 if let Some(manifest) = read_package_manifest_from_dir(&package_dir)? {
120 providers.extend(resolved_provider_connectors_from_manifest(
121 &manifest,
122 &package_dir,
123 ));
124 }
125 continue;
126 }
127
128 let package_file = packages_dir.join(format!("{}.harn", entry.name));
129 if package_file.is_file() {
130 continue;
131 }
132
133 return Err(PackageError::Manifest(format!(
134 "installed package {} is missing under {}; run `harn install`",
135 entry.name,
136 packages_dir.display()
137 )));
138 }
139 Ok(providers)
140}
141
142fn dedupe_provider_connectors(
143 providers: Vec<ResolvedProviderConnectorConfig>,
144) -> Vec<ResolvedProviderConnectorConfig> {
145 let mut seen = std::collections::BTreeSet::new();
146 let mut out = Vec::new();
147 for provider in providers {
148 if seen.insert(provider.id.as_str().to_string()) {
149 out.push(provider);
150 }
151 }
152 out
153}
154
155pub fn load_runtime_extensions(anchor: &Path) -> RuntimeExtensions {
156 match try_load_runtime_extensions(anchor) {
157 Ok(extensions) => extensions,
158 Err(error) => {
159 eprintln!("error: {error}");
160 process::exit(1);
161 }
162 }
163}
164
165pub fn install_runtime_extensions(extensions: &RuntimeExtensions) {
167 harn_vm::llm_config::set_user_overrides(extensions.llm.clone());
168 harn_vm::llm::capabilities::set_user_overrides(extensions.capabilities.clone());
169 install_manifest_handoff_routes(extensions);
170 install_orchestrator_budget(extensions);
171}
172
173pub fn install_manifest_handoff_routes(extensions: &RuntimeExtensions) {
174 harn_vm::install_handoff_routes(extensions.handoff_routes.clone());
175}
176
177pub fn install_orchestrator_budget(extensions: &RuntimeExtensions) {
178 let budget = extensions
179 .root_manifest
180 .as_ref()
181 .map(|manifest| harn_vm::OrchestratorBudgetConfig {
182 daily_cost_usd: manifest.orchestrator.budget.daily_cost_usd,
183 hourly_cost_usd: manifest.orchestrator.budget.hourly_cost_usd,
184 })
185 .unwrap_or_default();
186 harn_vm::install_orchestrator_budget(budget);
187}
188
189pub async fn install_manifest_hooks(
190 vm: &mut harn_vm::Vm,
191 extensions: &RuntimeExtensions,
192) -> Result<(), PackageError> {
193 install_manifest_hooks_with_mode(vm, extensions, false).await
194}
195
196pub async fn install_manifest_hooks_with_mode(
204 vm: &mut harn_vm::Vm,
205 extensions: &RuntimeExtensions,
206 lazy: bool,
207) -> Result<(), PackageError> {
208 harn_vm::orchestration::clear_runtime_hooks();
209 let mut loaded_exports: HashMap<ManifestModuleCacheKey, ManifestModuleExports> = HashMap::new();
210 let mut module_signatures: HashMap<PathBuf, Vec<CachedModuleCallableSignatures>> =
211 HashMap::new();
212 for hook in &extensions.hooks {
213 let Some((module_name, function_name)) = hook.handler.rsplit_once("::") else {
214 return Err(format!(
215 "invalid hook handler '{}': expected <module>::<function>",
216 hook.handler
217 )
218 .into());
219 };
220 let module_path = crate::package::manifest_module_source_path(
221 &hook.manifest_dir,
222 hook.package_name.as_deref(),
223 &hook.exports,
224 Some(module_name),
225 )?;
226 let signatures =
227 cached_module_callable_signatures(&mut module_signatures, &module_path, None)?;
228 if signatures
229 .get(function_name)
230 .is_none_or(|signature| !signature.is_pub)
231 {
232 return Err(format!(
233 "hook handler '{function_name}' is not exported by module '{module_name}'"
234 )
235 .into());
236 }
237 if lazy {
238 harn_vm::orchestration::register_vm_hook_lazy(
239 hook.event,
240 hook.pattern.clone(),
241 hook.handler.clone(),
242 harn_vm::LazyVmCallable::new(module_path, function_name),
243 );
244 continue;
245 }
246 let cache_key = (
247 hook.manifest_dir.clone(),
248 hook.package_name.clone(),
249 Some(module_name.to_string()),
250 );
251 if !loaded_exports.contains_key(&cache_key) {
252 let exports = resolve_manifest_exports(
253 vm,
254 &hook.manifest_dir,
255 hook.package_name.as_deref(),
256 &hook.exports,
257 Some(module_name),
258 )
259 .await?;
260 loaded_exports.insert(cache_key.clone(), exports);
261 }
262 let exports = loaded_exports
263 .get(&cache_key)
264 .expect("manifest hook exports cached");
265 let Some(closure) = exports.get(function_name) else {
266 return Err(format!(
267 "hook handler '{function_name}' is not exported by module '{module_name}'"
268 )
269 .into());
270 };
271 harn_vm::orchestration::register_vm_hook(
272 hook.event,
273 hook.pattern.clone(),
274 hook.handler.clone(),
275 closure.clone(),
276 );
277 }
278 Ok(())
279}
280
281pub async fn collect_manifest_triggers(
282 vm: &mut harn_vm::Vm,
283 extensions: &RuntimeExtensions,
284) -> Result<Vec<CollectedManifestTrigger>, PackageError> {
285 collect_manifest_triggers_with_mode(vm, extensions, false).await
286}
287
288async fn collect_manifest_triggers_with_mode(
289 vm: &mut harn_vm::Vm,
290 extensions: &RuntimeExtensions,
291 lazy_vm_callables: bool,
292) -> Result<Vec<CollectedManifestTrigger>, PackageError> {
293 let _provider_schema_guard = lock_manifest_provider_schemas().await;
294 let provider_schemas = build_manifest_provider_schemas(extensions).await?;
295 let provider_catalog = manifest_provider_catalog(provider_schemas.clone())?;
296 validate_orchestrator_budget(extensions.root_manifest.as_ref())?;
297 validate_static_trigger_configs(&extensions.triggers, &provider_catalog)?;
298 let mut loaded_exports: HashMap<ManifestModuleCacheKey, ManifestModuleExports> = HashMap::new();
299 let mut module_signatures: HashMap<PathBuf, Vec<CachedModuleCallableSignatures>> =
300 HashMap::new();
301 let mut validated = Vec::with_capacity(extensions.triggers.len());
302 for trigger in &extensions.triggers {
303 validated.push(validate_trigger_callable_declarations(
304 trigger,
305 &mut module_signatures,
306 )?);
307 }
308 let mut collected = Vec::new();
309
310 for (trigger, declarations) in extensions.triggers.iter().zip(validated) {
311 let mut effective_config = trigger.clone();
312 let collected_handler = match declarations.handler {
313 TriggerHandlerUri::Local(reference) => {
314 let module_path = declarations
315 .local_handler_path
316 .expect("validated local trigger handler has a source path");
317 let callable = collect_manifest_vm_callable(
318 vm,
319 &mut loaded_exports,
320 trigger,
321 &reference,
322 &module_path,
323 lazy_vm_callables,
324 "handler",
325 )
326 .await?;
327 CollectedTriggerHandler::Local {
328 reference,
329 callable,
330 }
331 }
332 TriggerHandlerUri::A2a {
333 target,
334 allow_cleartext,
335 } => CollectedTriggerHandler::A2a {
336 target,
337 allow_cleartext,
338 },
339 TriggerHandlerUri::Worker { queue } => CollectedTriggerHandler::Worker { queue },
340 TriggerHandlerUri::Persona { name } => {
341 let (binding, callable, autonomy_ceiling) =
342 persona_runtime_handler_for_trigger(extensions, trigger, &name)?;
343 effective_config.autonomy_tier =
344 effective_config.autonomy_tier.min(autonomy_ceiling);
345 CollectedTriggerHandler::Persona { binding, callable }
346 }
347 TriggerHandlerUri::EvalPack { target } => {
348 let manifest = eval_pack_manifest_for_handler(trigger, &target)?;
349 let ledger_options = eval_pack_ledger_options_for_handler(trigger)?;
350 CollectedTriggerHandler::EvalPack {
351 target,
352 manifest: Box::new(manifest),
353 ledger_options,
354 }
355 }
356 };
357
358 let collected_when = if let Some((reference, source_path)) = declarations.when {
359 let callable = collect_manifest_vm_callable(
360 vm,
361 &mut loaded_exports,
362 trigger,
363 &reference,
364 &source_path,
365 lazy_vm_callables,
366 "when predicate",
367 )
368 .await?;
369
370 Some(CollectedTriggerPredicate {
371 reference,
372 callable,
373 })
374 } else {
375 None
376 };
377
378 let flow_control = collect_trigger_flow_control(vm, trigger).await?;
379
380 collected.push(CollectedManifestTrigger {
381 config: effective_config,
382 handler: collected_handler,
383 when: collected_when,
384 flow_control,
385 });
386 }
387
388 register_manifest_provider_schemas(provider_schemas)?;
389 Ok(collected)
390}
391
392struct ValidatedTriggerCallableDeclarations {
393 handler: TriggerHandlerUri,
394 local_handler_path: Option<PathBuf>,
395 when: Option<(TriggerFunctionRef, PathBuf)>,
396}
397
398struct CachedModuleCallableSignatures {
399 execution_guard: Option<Arc<harn_modules::package_execution::PackageExecutionGuard>>,
400 signatures: BTreeMap<String, ModuleCallableSignature>,
401}
402
403fn validate_trigger_callable_declarations(
404 trigger: &ResolvedTriggerConfig,
405 module_signatures: &mut HashMap<PathBuf, Vec<CachedModuleCallableSignatures>>,
406) -> Result<ValidatedTriggerCallableDeclarations, PackageError> {
407 let handler = parse_trigger_handler_uri(trigger)?;
408 let local_handler_path = if let TriggerHandlerUri::Local(reference) = &handler {
409 let module_path = trigger_function_source_path(trigger, reference)?;
410 let signatures = cached_module_callable_signatures(
411 module_signatures,
412 &module_path,
413 trigger.execution_guard.as_ref(),
414 )
415 .map_err(|error| trigger_error(trigger, error))?;
416 if signatures
417 .get(&reference.function_name)
418 .is_none_or(|signature| !signature.is_pub)
419 {
420 return Err(trigger_error(
421 trigger,
422 format!(
423 "handler '{}' is not exported by the resolved module",
424 reference.raw
425 ),
426 ));
427 }
428 Some(module_path)
429 } else {
430 None
431 };
432 let when = if let Some(when_raw) = &trigger.when {
433 let reference = parse_local_trigger_ref(when_raw, "when", trigger)?;
434 let source_path = trigger_function_source_path(trigger, &reference)?;
435 let signatures = cached_module_callable_signatures(
436 module_signatures,
437 &source_path,
438 trigger.execution_guard.as_ref(),
439 )
440 .map_err(|error| trigger_error(trigger, error))?;
441 let Some(signature) = signatures.get(&reference.function_name) else {
442 return Err(trigger_error(
443 trigger,
444 format!(
445 "when predicate '{}' must resolve to a function declaration",
446 reference.raw
447 ),
448 ));
449 };
450 if !signature.is_pub {
451 return Err(trigger_error(
452 trigger,
453 format!(
454 "when predicate '{}' is not exported by the resolved module",
455 reference.raw
456 ),
457 ));
458 }
459 if signature.params.len() != 1
460 || signature.params[0]
461 .as_ref()
462 .is_none_or(|param| !is_trigger_event_type(param))
463 {
464 return Err(trigger_error(
465 trigger,
466 format!(
467 "when predicate '{}' must have signature fn(TriggerEvent) -> bool",
468 reference.raw
469 ),
470 ));
471 }
472 if signature
473 .return_type
474 .as_ref()
475 .is_none_or(|return_type| !is_predicate_return_type(return_type))
476 {
477 return Err(trigger_error(
478 trigger,
479 format!(
480 "when predicate '{}' must have signature fn(TriggerEvent) -> bool or Result<bool, _>",
481 reference.raw
482 ),
483 ));
484 }
485 Some((reference, source_path))
486 } else {
487 None
488 };
489 Ok(ValidatedTriggerCallableDeclarations {
490 handler,
491 local_handler_path,
492 when,
493 })
494}
495
496fn trigger_function_source_path(
497 trigger: &ResolvedTriggerConfig,
498 reference: &TriggerFunctionRef,
499) -> Result<PathBuf, PackageError> {
500 manifest_module_source_path(
501 &trigger.manifest_dir,
502 trigger.package_name.as_deref(),
503 &trigger.exports,
504 reference.module_name.as_deref(),
505 )
506 .map_err(|error| trigger_error(trigger, error))
507}
508
509async fn collect_manifest_vm_callable(
510 vm: &mut harn_vm::Vm,
511 loaded_exports: &mut HashMap<ManifestModuleCacheKey, ManifestModuleExports>,
512 trigger: &ResolvedTriggerConfig,
513 reference: &TriggerFunctionRef,
514 module_path: &Path,
515 lazy: bool,
516 role: &str,
517) -> Result<harn_vm::VmCallable, PackageError> {
518 let mut deferred =
519 harn_vm::LazyVmCallable::new(module_path.to_path_buf(), reference.function_name.clone());
520 if let Some(guard) = &trigger.execution_guard {
521 deferred = deferred.with_package_execution_guard(Arc::clone(guard));
522 }
523 if lazy {
524 return Ok(harn_vm::VmCallable::Lazy(deferred));
525 }
526 if trigger.execution_guard.is_some() {
527 let closure = vm
528 .resolve_callable(&harn_vm::VmCallable::Lazy(deferred))
529 .await
530 .map_err(|error| trigger_error(trigger, error.to_string()))?;
531 return Ok(harn_vm::VmCallable::Eager(closure));
532 }
533
534 let cache_key = (
535 trigger.manifest_dir.clone(),
536 trigger.package_name.clone(),
537 reference.module_name.clone(),
538 );
539 if !loaded_exports.contains_key(&cache_key) {
540 let exports = resolve_manifest_exports(
541 vm,
542 &trigger.manifest_dir,
543 trigger.package_name.as_deref(),
544 &trigger.exports,
545 reference.module_name.as_deref(),
546 )
547 .await
548 .map_err(|error| trigger_error(trigger, error))?;
549 loaded_exports.insert(cache_key.clone(), exports);
550 }
551 let exports = loaded_exports
552 .get(&cache_key)
553 .expect("manifest trigger exports cached");
554 let closure = exports.get(&reference.function_name).ok_or_else(|| {
555 trigger_error(
556 trigger,
557 format!(
558 "{role} '{}' is not exported by the resolved module",
559 reference.raw
560 ),
561 )
562 })?;
563 Ok(harn_vm::VmCallable::Eager(closure.clone()))
564}
565
566fn cached_module_callable_signatures<'a>(
567 cache: &'a mut HashMap<PathBuf, Vec<CachedModuleCallableSignatures>>,
568 source_path: &Path,
569 execution_guard: Option<&Arc<harn_modules::package_execution::PackageExecutionGuard>>,
570) -> Result<&'a BTreeMap<String, ModuleCallableSignature>, PackageError> {
571 let entries = cache.entry(source_path.to_path_buf()).or_default();
572 if let Some(index) = entries
573 .iter()
574 .position(|entry| entry.execution_guard.as_ref() == execution_guard)
575 {
576 return Ok(&entries[index].signatures);
577 }
578 let signatures = if let Some(guard) = execution_guard {
579 load_guarded_module_callable_signatures(source_path, guard)?
580 } else {
581 load_module_callable_signatures(source_path)?
582 };
583 entries.push(CachedModuleCallableSignatures {
584 execution_guard: execution_guard.cloned(),
585 signatures,
586 });
587 Ok(&entries
588 .last()
589 .expect("signature cache entry inserted")
590 .signatures)
591}
592
593pub(crate) async fn collect_trigger_flow_control(
594 vm: &mut harn_vm::Vm,
595 trigger: &ResolvedTriggerConfig,
596) -> Result<harn_vm::TriggerFlowControlConfig, PackageError> {
597 let mut flow = harn_vm::TriggerFlowControlConfig::default();
598
599 let concurrency = if let Some(spec) = &trigger.concurrency {
600 Some(spec.clone())
601 } else if let Some(max) = trigger.budget.max_concurrent {
602 eprintln!(
603 "warning: {} uses deprecated budget.max_concurrent; prefer concurrency = {{ max = {} }}",
604 manifest_trigger_location(trigger),
605 max
606 );
607 Some(TriggerConcurrencyManifestSpec { key: None, max })
608 } else {
609 None
610 };
611 if let Some(spec) = concurrency {
612 flow.concurrency = Some(harn_vm::TriggerConcurrencyConfig {
613 key: compile_optional_trigger_expression(
614 vm,
615 trigger,
616 "concurrency.key",
617 spec.key.as_deref(),
618 )
619 .await?,
620 max: spec.max,
621 });
622 }
623
624 if let Some(spec) = &trigger.throttle {
625 flow.throttle = Some(harn_vm::TriggerThrottleConfig {
626 key: compile_optional_trigger_expression(
627 vm,
628 trigger,
629 "throttle.key",
630 spec.key.as_deref(),
631 )
632 .await?,
633 period: harn_vm::parse_flow_control_duration(&spec.period)
634 .map_err(|error| trigger_error(trigger, format!("throttle.period {error}")))?,
635 max: spec.max,
636 });
637 }
638
639 if let Some(spec) = &trigger.rate_limit {
640 flow.rate_limit = Some(harn_vm::TriggerRateLimitConfig {
641 key: compile_optional_trigger_expression(
642 vm,
643 trigger,
644 "rate_limit.key",
645 spec.key.as_deref(),
646 )
647 .await?,
648 period: harn_vm::parse_flow_control_duration(&spec.period)
649 .map_err(|error| trigger_error(trigger, format!("rate_limit.period {error}")))?,
650 max: spec.max,
651 });
652 }
653
654 if let Some(spec) = &trigger.debounce {
655 flow.debounce = Some(harn_vm::TriggerDebounceConfig {
656 key: compile_trigger_expression(vm, trigger, "debounce.key", &spec.key).await?,
657 period: harn_vm::parse_flow_control_duration(&spec.period)
658 .map_err(|error| trigger_error(trigger, format!("debounce.period {error}")))?,
659 });
660 }
661
662 if let Some(spec) = &trigger.singleton {
663 flow.singleton = Some(harn_vm::TriggerSingletonConfig {
664 key: compile_optional_trigger_expression(
665 vm,
666 trigger,
667 "singleton.key",
668 spec.key.as_deref(),
669 )
670 .await?,
671 });
672 }
673
674 if let Some(spec) = &trigger.batch {
675 flow.batch = Some(harn_vm::TriggerBatchConfig {
676 key: compile_optional_trigger_expression(vm, trigger, "batch.key", spec.key.as_deref())
677 .await?,
678 size: spec.size,
679 timeout: harn_vm::parse_flow_control_duration(&spec.timeout)
680 .map_err(|error| trigger_error(trigger, format!("batch.timeout {error}")))?,
681 });
682 }
683
684 if let Some(spec) = &trigger.priority_flow {
685 flow.priority = Some(harn_vm::TriggerPriorityOrderConfig {
686 key: compile_trigger_expression(vm, trigger, "priority.key", &spec.key).await?,
687 order: spec.order.clone(),
688 });
689 }
690
691 Ok(flow)
692}
693
694fn eval_pack_manifest_for_handler(
695 trigger: &ResolvedTriggerConfig,
696 target: &str,
697) -> Result<harn_vm::orchestration::EvalPackManifest, PackageError> {
698 if eval_pack_target_is_path(target) {
699 let path = resolve_eval_pack_target_path(&trigger.manifest_dir, target);
700 return harn_vm::orchestration::load_eval_pack_manifest(&path).map_err(|error| {
701 trigger_error(
702 trigger,
703 format!(
704 "handler eval_pack://{target} failed to load eval pack {}: {error}",
705 path.display()
706 ),
707 )
708 });
709 }
710
711 let paths = load_package_eval_pack_paths(Some(&trigger.manifest_path))
712 .map_err(|error| trigger_error(trigger, error))?;
713 let mut matches = Vec::new();
714 for path in paths {
715 let manifest = harn_vm::orchestration::load_eval_pack_manifest(&path).map_err(|error| {
716 trigger_error(
717 trigger,
718 format!(
719 "failed to load package eval pack {}: {error}",
720 path.display()
721 ),
722 )
723 })?;
724 let file_stem = path.file_stem().and_then(|stem| stem.to_str());
725 if manifest.id == target
726 || manifest.name.as_deref() == Some(target)
727 || file_stem == Some(target)
728 {
729 matches.push((path, manifest));
730 }
731 }
732
733 match matches.len() {
734 0 => Err(trigger_error(
735 trigger,
736 format!(
737 "handler eval_pack://{target} did not match any package eval pack by id, name, or file stem",
738 ),
739 )),
740 1 => Ok(matches.remove(0).1),
741 _ => Err(trigger_error(
742 trigger,
743 format!("handler eval_pack://{target} matched multiple package eval packs"),
744 )),
745 }
746}
747
748fn eval_pack_target_is_path(target: &str) -> bool {
749 target.contains('/')
750 || target.contains('\\')
751 || target.ends_with(".toml")
752 || target.ends_with(".json")
753}
754
755fn resolve_eval_pack_target_path(manifest_dir: &Path, target: &str) -> PathBuf {
756 let path = PathBuf::from(target);
757 if path.is_absolute() {
758 path
759 } else {
760 manifest_dir.join(path)
761 }
762}
763
764fn eval_pack_ledger_options_for_handler(
765 trigger: &ResolvedTriggerConfig,
766) -> Result<Option<serde_json::Value>, PackageError> {
767 let value = trigger
768 .kind_specific
769 .get("eval_options")
770 .or_else(|| trigger.kind_specific.get("ledger"));
771 value
772 .map(|value| {
773 serde_json::to_value(value).map_err(|error| {
774 trigger_error(trigger, format!("invalid eval ledger options: {error}"))
775 })
776 })
777 .transpose()
778}
779
780pub(crate) async fn compile_optional_trigger_expression(
781 vm: &mut harn_vm::Vm,
782 trigger: &ResolvedTriggerConfig,
783 field_name: &str,
784 expr: Option<&str>,
785) -> Result<Option<harn_vm::TriggerExpressionSpec>, PackageError> {
786 match expr {
787 Some(expr) => compile_trigger_expression(vm, trigger, field_name, expr)
788 .await
789 .map(Some),
790 None => Ok(None),
791 }
792}
793
794pub(crate) async fn compile_trigger_expression(
795 vm: &mut harn_vm::Vm,
796 trigger: &ResolvedTriggerConfig,
797 field_name: &str,
798 expr: &str,
799) -> Result<harn_vm::TriggerExpressionSpec, PackageError> {
800 let synthetic = PathBuf::from(format!(
801 "<trigger-expr>/{}/{:04}-{}.harn",
802 harn_vm::event_log::sanitize_topic_component(&trigger.id),
803 trigger.table_index,
804 harn_vm::event_log::sanitize_topic_component(field_name),
805 ));
806 let source = format!(
807 "import \"std/triggers\"\n\npub fn __trigger_expr(event: TriggerEvent) -> any {{\n return {expr}\n}}\n"
808 );
809 let exports = vm
810 .load_module_exports_from_source(synthetic, &source)
811 .await
812 .map_err(|error| {
813 trigger_error(
814 trigger,
815 format!("{field_name} '{expr}' is invalid Harn expression: {error}"),
816 )
817 })?;
818 let closure = exports.get("__trigger_expr").ok_or_else(|| {
819 trigger_error(
820 trigger,
821 format!("{field_name} '{expr}' did not compile into an exported closure"),
822 )
823 })?;
824 Ok(harn_vm::TriggerExpressionSpec {
825 raw: expr.to_string(),
826 callable: harn_vm::VmCallable::Eager(closure.clone()),
827 })
828}
829
830pub(crate) fn trigger_kind_label(kind: TriggerKind) -> &'static str {
831 match kind {
832 TriggerKind::Webhook => "webhook",
833 TriggerKind::Cron => "cron",
834 TriggerKind::Poll => "poll",
835 TriggerKind::Stream => "stream",
836 TriggerKind::Predicate => "predicate",
837 TriggerKind::A2aPush => "a2a-push",
838 }
839}
840
841pub(crate) fn worker_queue_priority(
842 priority: TriggerDispatchPriority,
843) -> harn_vm::WorkerQueuePriority {
844 match priority {
845 TriggerDispatchPriority::High => harn_vm::WorkerQueuePriority::High,
846 TriggerDispatchPriority::Normal => harn_vm::WorkerQueuePriority::Normal,
847 TriggerDispatchPriority::Low => harn_vm::WorkerQueuePriority::Low,
848 }
849}
850
851pub fn manifest_trigger_binding_spec(
852 trigger: CollectedManifestTrigger,
853) -> harn_vm::TriggerBindingSpec {
854 let flow_control = trigger.flow_control.clone();
855 let config = trigger.config;
856 let (handler, handler_descriptor) = match trigger.handler {
857 CollectedTriggerHandler::Local {
858 reference,
859 callable,
860 } => (
861 harn_vm::TriggerHandlerSpec::Local {
862 raw: reference.raw.clone(),
863 callable,
864 },
865 serde_json::json!({
866 "kind": "local",
867 "raw": reference.raw,
868 }),
869 ),
870 CollectedTriggerHandler::A2a {
871 target,
872 allow_cleartext,
873 } => (
874 harn_vm::TriggerHandlerSpec::A2a {
875 target: target.clone(),
876 allow_cleartext,
877 },
878 serde_json::json!({
879 "kind": "a2a",
880 "target": target,
881 "allow_cleartext": allow_cleartext,
882 }),
883 ),
884 CollectedTriggerHandler::Worker { queue } => (
885 harn_vm::TriggerHandlerSpec::Worker {
886 queue: queue.clone(),
887 },
888 serde_json::json!({
889 "kind": "worker",
890 "queue": queue,
891 }),
892 ),
893 CollectedTriggerHandler::Persona { binding, callable } => (
894 harn_vm::TriggerHandlerSpec::Persona {
895 binding: binding.clone(),
896 callable,
897 },
898 serde_json::json!({
899 "kind": "persona",
900 "name": binding.name,
901 "entry_workflow": binding.entry_workflow,
902 }),
903 ),
904 CollectedTriggerHandler::EvalPack {
905 target,
906 manifest,
907 ledger_options,
908 } => {
909 let pack_id = manifest.id.clone();
910 let harness_config_fingerprint =
911 harn_vm::orchestration::eval_pack_harness_config_fingerprint(manifest.as_ref())
912 .ok();
913 (
914 harn_vm::TriggerHandlerSpec::EvalPack {
915 target: target.clone(),
916 manifest,
917 ledger_options: ledger_options.clone(),
918 },
919 serde_json::json!({
920 "kind": "eval_pack",
921 "target": target,
922 "pack_id": pack_id,
923 "harness_config_fingerprint": harness_config_fingerprint,
924 "ledger_options": ledger_options,
925 }),
926 )
927 }
928 };
929
930 let when_raw = trigger
931 .when
932 .as_ref()
933 .map(|predicate| predicate.reference.raw.clone());
934 let when = trigger.when.map(|predicate| harn_vm::TriggerPredicateSpec {
935 raw: predicate.reference.raw,
936 callable: predicate.callable,
937 });
938 let mut when_budget = config
939 .when_budget
940 .as_ref()
941 .map(|budget| {
942 Ok::<harn_vm::TriggerPredicateBudget, String>(harn_vm::TriggerPredicateBudget {
943 max_cost_usd: budget.max_cost_usd,
944 tokens_max: budget.tokens_max,
945 timeout_ms: budget
946 .timeout
947 .as_deref()
948 .map(parse_duration_millis)
949 .transpose()?,
950 })
951 })
952 .transpose()
953 .unwrap_or_default();
954 if config.budget.max_cost_usd.is_some() || config.budget.max_tokens.is_some() {
955 let budget = when_budget.get_or_insert_with(harn_vm::TriggerPredicateBudget::default);
956 if budget.max_cost_usd.is_none() {
957 budget.max_cost_usd = config.budget.max_cost_usd;
958 }
959 if budget.tokens_max.is_none() {
960 budget.tokens_max = config.budget.max_tokens;
961 }
962 }
963 let id = config.id.clone();
964 let kind = trigger_kind_label(config.kind).to_string();
965 let provider = config.provider.clone();
966 let autonomy_tier = config.autonomy_tier;
967 let match_events = config.match_.events.clone();
968 let dedupe_key = config.dedupe_key.clone();
969 let retry = harn_vm::TriggerRetryConfig::new(
970 config.retry.max,
971 match config.retry.backoff {
972 TriggerRetryBackoff::Immediate => harn_vm::RetryPolicy::Linear { delay_ms: 0 },
973 TriggerRetryBackoff::Svix => harn_vm::RetryPolicy::Svix,
974 },
975 );
976 let filter = config.filter.clone();
977 let dedupe_retention_days = config.retry.retention_days;
978 let daily_cost_usd = config.budget.daily_cost_usd;
979 let hourly_cost_usd = config.budget.hourly_cost_usd;
980 let max_autonomous_decisions_per_hour = config.budget.max_autonomous_decisions_per_hour;
981 let max_autonomous_decisions_per_day = config.budget.max_autonomous_decisions_per_day;
982 let on_budget_exhausted = config.budget.on_budget_exhausted;
983 let max_concurrent = flow_control.concurrency.as_ref().map(|config| config.max);
984 let manifest_path = Some(config.manifest_path.clone());
985 let package_name = config.package_name.clone();
986
987 let fingerprint = serde_json::to_string(&serde_json::json!({
988 "id": &id,
989 "kind": &kind,
990 "provider": provider.as_str(),
991 "autonomy_tier": autonomy_tier,
992 "match": config.match_,
993 "when": when_raw,
994 "when_budget": config.when_budget,
995 "handler": handler_descriptor,
996 "dedupe_key": &dedupe_key,
997 "retry": config.retry,
998 "dispatch_priority": config.dispatch_priority,
999 "budget": config.budget,
1000 "flow_control": {
1001 "concurrency": config.concurrency,
1002 "throttle": config.throttle,
1003 "rate_limit": config.rate_limit,
1004 "debounce": config.debounce,
1005 "singleton": config.singleton,
1006 "batch": config.batch,
1007 "priority": config.priority_flow,
1008 },
1009 "window": config.window,
1010 "secrets": config.secrets,
1011 "filter": &filter,
1012 "kind_specific": config.kind_specific,
1013 "manifest_path": &manifest_path,
1014 "package_name": &package_name,
1015 }))
1016 .unwrap_or_else(|_| format!("{}:{}:{}", id, kind, provider.as_str()));
1017
1018 harn_vm::TriggerBindingSpec {
1019 id,
1020 source: harn_vm::TriggerBindingSource::Manifest,
1021 kind,
1022 provider,
1023 autonomy_tier,
1024 handler,
1025 dispatch_priority: worker_queue_priority(config.dispatch_priority),
1026 when,
1027 when_budget,
1028 retry,
1029 match_events,
1030 dedupe_key,
1031 filter,
1032 dedupe_retention_days,
1033 daily_cost_usd,
1034 hourly_cost_usd,
1035 max_autonomous_decisions_per_hour,
1036 max_autonomous_decisions_per_day,
1037 on_budget_exhausted,
1038 max_concurrent,
1039 flow_control,
1040 aggregation: None,
1041 manifest_path,
1042 package_name,
1043 definition_fingerprint: fingerprint,
1044 }
1045}
1046
1047pub async fn install_manifest_triggers(
1048 vm: &mut harn_vm::Vm,
1049 extensions: &RuntimeExtensions,
1050) -> Result<(), PackageError> {
1051 install_manifest_triggers_with_mode(vm, extensions, false).await
1052}
1053
1054pub async fn install_manifest_triggers_with_mode(
1059 vm: &mut harn_vm::Vm,
1060 extensions: &RuntimeExtensions,
1061 lazy_vm_callables: bool,
1062) -> Result<(), PackageError> {
1063 install_orchestrator_budget(extensions);
1064 let collected = collect_manifest_triggers_with_mode(vm, extensions, lazy_vm_callables).await?;
1065 let mut bindings: Vec<_> = collected
1066 .iter()
1067 .cloned()
1068 .map(manifest_trigger_binding_spec)
1069 .collect();
1070 bindings.extend(collect_persona_trigger_binding_specs(extensions)?);
1071 harn_vm::install_manifest_triggers(bindings)
1072 .await
1073 .map_err(|error| PackageError::Extensions(error.to_string()))
1074}
1075
1076pub async fn install_collected_manifest_triggers(
1077 collected: &[CollectedManifestTrigger],
1078) -> Result<(), PackageError> {
1079 let bindings = collected
1080 .iter()
1081 .cloned()
1082 .map(manifest_trigger_binding_spec)
1083 .collect();
1084 harn_vm::install_manifest_triggers(bindings)
1085 .await
1086 .map_err(|error| PackageError::Extensions(error.to_string()))
1087}
1088
1089pub fn load_personas_from_manifest_path(
1090 manifest_path: &Path,
1091) -> Result<ResolvedPersonaManifest, Vec<PersonaValidationError>> {
1092 let manifest_path = if manifest_path.is_dir() {
1093 manifest_path.join(MANIFEST)
1094 } else {
1095 manifest_path.to_path_buf()
1096 };
1097 let manifest_dir = manifest_path
1098 .parent()
1099 .map(Path::to_path_buf)
1100 .unwrap_or_else(|| PathBuf::from("."));
1101 if manifest_path.extension().and_then(|ext| ext.to_str()) == Some("harn") {
1102 return match harn_modules::personas::parse_persona_source_file(&manifest_path) {
1103 Ok(document) if !document.personas.is_empty() => {
1104 validate_and_resolve_standalone_personas(
1105 document.personas,
1106 manifest_path,
1107 manifest_dir,
1108 )
1109 }
1110 Ok(_) => Err(vec![PersonaValidationError {
1111 manifest_path: manifest_path.clone(),
1112 field_path: "persona".to_string(),
1113 message: "no @persona declarations found".to_string(),
1114 }]),
1115 Err(message) => Err(vec![PersonaValidationError {
1116 manifest_path: manifest_path.clone(),
1117 field_path: "persona".to_string(),
1118 message,
1119 }]),
1120 };
1121 }
1122 let manifest = match read_manifest_from_path(&manifest_path) {
1123 Ok(manifest) => manifest,
1124 Err(message) => {
1125 if let Ok(document) =
1126 harn_modules::personas::parse_persona_manifest_file(&manifest_path)
1127 {
1128 if !document.personas.is_empty() {
1129 return validate_and_resolve_standalone_personas(
1130 document.personas,
1131 manifest_path,
1132 manifest_dir,
1133 );
1134 }
1135 }
1136 return Err(vec![PersonaValidationError {
1137 manifest_path: manifest_path.clone(),
1138 field_path: "harn.toml".to_string(),
1139 message: message.to_string(),
1140 }]);
1141 }
1142 };
1143 if manifest.personas.is_empty() {
1144 if let Ok(document) = harn_modules::personas::parse_persona_manifest_file(&manifest_path) {
1145 if !document.personas.is_empty() {
1146 return validate_and_resolve_standalone_personas(
1147 document.personas,
1148 manifest_path,
1149 manifest_dir,
1150 );
1151 }
1152 }
1153 }
1154 validate_and_resolve_personas(manifest, manifest_path, manifest_dir)
1155}
1156
1157pub(crate) fn load_personas_from_verified_package_manifest(
1158 manifest_path: &Path,
1159 source: &str,
1160) -> Result<ResolvedPersonaManifest, Vec<PersonaValidationError>> {
1161 let manifest_path = manifest_path.to_path_buf();
1162 let manifest_dir = manifest_path
1163 .parent()
1164 .map(Path::to_path_buf)
1165 .unwrap_or_else(|| PathBuf::from("."));
1166 let manifest = toml::from_str::<Manifest>(source).map_err(|error| {
1167 vec![PersonaValidationError {
1168 manifest_path: manifest_path.clone(),
1169 field_path: "harn.toml".to_string(),
1170 message: format!("failed to parse {}: {error}", manifest_path.display()),
1171 }]
1172 })?;
1173 validate_and_resolve_personas(manifest, manifest_path, manifest_dir)
1174}
1175
1176fn validate_and_resolve_standalone_personas(
1177 personas: Vec<PersonaManifestEntry>,
1178 manifest_path: PathBuf,
1179 manifest_dir: PathBuf,
1180) -> Result<ResolvedPersonaManifest, Vec<PersonaValidationError>> {
1181 let known_names = personas
1182 .iter()
1183 .filter_map(|persona| persona.name.as_ref())
1184 .filter(|name| !name.trim().is_empty())
1185 .cloned()
1186 .collect();
1187 let context = harn_modules::personas::PersonaValidationContext {
1188 known_capabilities: harn_modules::personas::default_persona_capabilities(),
1189 known_tools: BTreeSet::new(),
1190 known_names,
1191 };
1192 harn_modules::personas::validate_persona_manifests(&manifest_path, &personas, &context)?;
1193 Ok(ResolvedPersonaManifest {
1194 manifest_path,
1195 manifest_dir,
1196 personas,
1197 })
1198}
1199
1200pub fn load_personas_config(
1201 anchor: Option<&Path>,
1202) -> Result<Option<ResolvedPersonaManifest>, Vec<PersonaValidationError>> {
1203 let anchor = anchor
1204 .map(Path::to_path_buf)
1205 .unwrap_or_else(|| std::env::current_dir().unwrap_or_else(|_| PathBuf::from(".")));
1206 let Some((manifest, dir)) = nearest_manifest_or_warn(&anchor) else {
1207 return Ok(None);
1208 };
1209 let manifest_path = dir.join(MANIFEST);
1210 validate_and_resolve_personas(manifest, manifest_path, dir).map(Some)
1211}
1212
1213pub(crate) fn validate_and_resolve_personas(
1214 manifest: Manifest,
1215 manifest_path: PathBuf,
1216 manifest_dir: PathBuf,
1217) -> Result<ResolvedPersonaManifest, Vec<PersonaValidationError>> {
1218 let known_capabilities = known_persona_capabilities(&manifest, &manifest_dir);
1219 let known_tools = known_persona_tools(&manifest);
1220 let known_names: BTreeSet<String> = manifest
1221 .personas
1222 .iter()
1223 .filter_map(|persona| persona.name.as_ref())
1224 .filter(|name| !name.trim().is_empty())
1225 .cloned()
1226 .collect();
1227 let context = harn_modules::personas::PersonaValidationContext {
1228 known_capabilities,
1229 known_tools,
1230 known_names,
1231 };
1232 if let Err(errors) = harn_modules::personas::validate_persona_manifests(
1233 &manifest_path,
1234 &manifest.personas,
1235 &context,
1236 ) {
1237 Err(errors)
1238 } else {
1239 let mut personas = manifest.personas;
1240 attach_entry_workflow_steps(&mut personas, &manifest_dir);
1241 Ok(ResolvedPersonaManifest {
1242 manifest_path,
1243 manifest_dir,
1244 personas,
1245 })
1246 }
1247}
1248
1249fn attach_entry_workflow_steps(personas: &mut [PersonaManifestEntry], manifest_dir: &Path) {
1250 for persona in personas {
1251 if !persona.steps.is_empty() {
1252 continue;
1253 }
1254 let Some(entry_workflow) = persona.entry_workflow.as_deref() else {
1255 continue;
1256 };
1257 let Some((path, entry_name)) = entry_workflow.split_once('#') else {
1258 continue;
1259 };
1260 if !path.ends_with(".harn") {
1261 continue;
1262 }
1263 let source_path = manifest_dir.join(path);
1264 let Ok(document) = harn_modules::personas::parse_persona_source_file(&source_path) else {
1265 continue;
1266 };
1267 let entry_name = entry_name.trim();
1268 if let Some(source_persona) = document.personas.iter().find(|candidate| {
1269 candidate.entry_workflow.as_deref() == Some(entry_name)
1270 || candidate.name.as_deref() == persona.name.as_deref()
1271 }) {
1272 persona.steps.clone_from(&source_persona.steps);
1273 }
1274 }
1275}
1276
1277pub(crate) fn known_persona_capabilities(
1278 manifest: &Manifest,
1279 manifest_dir: &Path,
1280) -> BTreeSet<String> {
1281 let mut capabilities = BTreeSet::new();
1282 for (capability, operations) in default_persona_capability_map() {
1283 for operation in operations {
1284 capabilities.insert(format!("{capability}.{operation}"));
1285 }
1286 }
1287 for (capability, operations) in &manifest.check.host_capabilities {
1288 for operation in operations {
1289 capabilities.insert(format!("{capability}.{operation}"));
1290 }
1291 }
1292 if let Some(path) = manifest.check.host_capabilities_path.as_deref() {
1293 let path = PathBuf::from(path);
1294 let path = if path.is_absolute() {
1295 path
1296 } else {
1297 manifest_dir.join(path)
1298 };
1299 if let Ok(content) = fs::read_to_string(path) {
1300 let parsed_json = serde_json::from_str::<serde_json::Value>(&content).ok();
1301 let parsed_toml = toml::from_str::<toml::Value>(&content)
1302 .ok()
1303 .and_then(|value| serde_json::to_value(value).ok());
1304 if let Some(value) = parsed_json.or(parsed_toml) {
1305 collect_persona_capabilities_from_json(&value, &mut capabilities);
1306 }
1307 }
1308 }
1309 capabilities
1310}
1311
1312pub(crate) fn collect_persona_capabilities_from_json(
1313 value: &serde_json::Value,
1314 out: &mut BTreeSet<String>,
1315) {
1316 let root = value.get("capabilities").unwrap_or(value);
1317 let Some(capabilities) = root.as_object() else {
1318 return;
1319 };
1320 for (capability, entry) in capabilities {
1321 if let Some(list) = entry.as_array() {
1322 for item in list {
1323 if let Some(operation) = item.as_str() {
1324 out.insert(format!("{capability}.{operation}"));
1325 }
1326 }
1327 } else if let Some(obj) = entry.as_object() {
1328 if let Some(list) = obj
1329 .get("operations")
1330 .or_else(|| obj.get("ops"))
1331 .and_then(|v| v.as_array())
1332 {
1333 for item in list {
1334 if let Some(operation) = item.as_str() {
1335 out.insert(format!("{capability}.{operation}"));
1336 }
1337 }
1338 } else {
1339 for (operation, enabled) in obj {
1340 if enabled.as_bool().unwrap_or(true) {
1341 out.insert(format!("{capability}.{operation}"));
1342 }
1343 }
1344 }
1345 }
1346 }
1347}
1348
1349pub(crate) fn default_persona_capability_map() -> BTreeMap<&'static str, Vec<&'static str>> {
1350 harn_modules::personas::default_persona_capability_map()
1351}
1352
1353pub(crate) fn known_persona_tools(manifest: &Manifest) -> BTreeSet<String> {
1354 let mut tools = BTreeSet::from([
1355 "a2a".to_string(),
1356 "acp".to_string(),
1357 "ci".to_string(),
1358 "filesystem".to_string(),
1359 "github".to_string(),
1360 "linear".to_string(),
1361 "mcp".to_string(),
1362 "notion".to_string(),
1363 "pagerduty".to_string(),
1364 "shell".to_string(),
1365 "slack".to_string(),
1366 ]);
1367 for server in &manifest.mcp {
1368 tools.insert(server.name.clone());
1369 }
1370 for provider in &manifest.providers {
1371 tools.insert(provider.id.as_str().to_string());
1372 }
1373 for trigger in &manifest.triggers {
1374 if let Some(provider) = trigger.provider.as_ref() {
1375 tools.insert(provider.as_str().to_string());
1376 }
1377 for source in &trigger.sources {
1378 tools.insert(source.provider.as_str().to_string());
1379 }
1380 }
1381 tools
1382}
1383
1384#[cfg(test)]
1385#[path = "extensions_tests.rs"]
1386mod tests;
1387
1388#[cfg(test)]
1389#[path = "extensions_lazy_tests.rs"]
1390mod lazy_tests;
1391
1392#[cfg(test)]
1393#[path = "extensions_provider_tests.rs"]
1394mod provider_tests;
1395
1396#[cfg(test)]
1397#[path = "persona_runtime_tests.rs"]
1398mod persona_tests;