1use std::collections::HashMap;
2use std::str::FromStr;
3use std::sync::Arc;
4
5use anyhow::{Context, Result, anyhow};
6use async_trait::async_trait;
7use greentic_session::{SessionData, SessionKey as StoreSessionKey};
8use greentic_types::{
9 EnvId, FlowId, GreenticError, PackId, ReplyScope, SessionCursor as TypesSessionCursor,
10 TenantCtx, TenantId, UserId,
11};
12
13use crate::telemetry::attr_keys;
14use rand::{RngExt, rng};
15use serde::{Deserialize, Serialize};
16use serde_json::{Value, json};
17use sha2::{Digest, Sha256};
18
19use super::api::{RunFlowRequest, RunnerApi};
20use super::builder::{Runner, RunnerBuilder};
21use super::error::{GResult, RunnerError};
22use super::glue::{FnSecretsHost, FnTelemetryHost};
23use super::host::{HostBundle, SecretsHost, SessionHost, StateHost};
24use super::policy::Policy;
25use super::registry::{Adapter, AdapterCall, AdapterRegistry};
26use super::shims::{InMemorySessionHost, InMemoryStateHost};
27use super::state_machine::{FlowDefinition, FlowStep, PAYLOAD_FROM_LAST_INPUT};
28
29use crate::config::{HostConfig, SecretsPolicy};
30use crate::pack::FlowDescriptor;
31use crate::runner::engine::{FlowContext, FlowEngine, FlowSnapshot, FlowStatus, FlowWait};
32use crate::runner::mocks::MockLayer;
33use crate::secrets::{DynSecretsManager, read_secret_blocking};
34use crate::storage::session::DynSessionStore;
35use crate::trace::audit_sink::AuditSink;
36use crate::trace::{PackTraceInfo, TraceContext, TraceMode, TraceRecorder};
37
38const DEFAULT_ENV: &str = "local";
39const PACK_FLOW_ADAPTER: &str = "pack_flow";
40
41#[derive(Clone)]
42pub struct FlowResumeStore {
43 store: DynSessionStore,
44}
45
46impl FlowResumeStore {
47 pub fn new(store: DynSessionStore) -> Self {
48 Self { store }
49 }
50
51 pub fn fetch(&self, envelope: &IngressEnvelope) -> GResult<Option<FlowSnapshot>> {
52 let (mut ctx, user, _, scope) = build_store_ctx(envelope)?;
53 ctx = ctx.with_user(Some(user.clone()));
54
55 let mut scopes = vec![scope.clone()];
56 if scope.correlation.is_some() {
57 let mut base = scope.clone();
58 base.correlation = None;
59 scopes.push(base);
60 }
61
62 for lookup in scopes {
63 if let Some(key) = self
64 .store
65 .find_wait_by_scope(&ctx, &user, &lookup)
66 .map_err(map_store_error)?
67 {
68 let Some(data) = self.store.get_session(&key).map_err(map_store_error)? else {
69 continue;
70 };
71 let record: FlowResumeRecord =
72 serde_json::from_str(&data.context_json).map_err(|err| {
73 RunnerError::Session {
74 reason: format!("failed to decode flow resume snapshot: {err}"),
75 }
76 })?;
77 if let Some(pack_id) = envelope.pack_id.as_deref()
78 && record.snapshot.pack_id != pack_id
79 {
80 return Err(RunnerError::Session {
81 reason: format!(
82 "resume pack mismatch: expected {pack_id}, found {}",
83 record.snapshot.pack_id
84 ),
85 });
86 }
87 return Ok(Some(record.snapshot));
88 }
89 }
90
91 Ok(None)
92 }
93
94 pub fn save(&self, envelope: &IngressEnvelope, wait: &FlowWait) -> GResult<ReplyScope> {
95 let (ctx, user, hint, scope) = build_store_ctx(envelope)?;
96 let record = FlowResumeRecord {
97 snapshot: wait.snapshot.clone(),
98 reason: wait.reason.clone(),
99 };
100 let data = record_to_session_data(&record, ctx.clone(), &user, &hint)?;
101 let mut reply_scope = scope.clone();
102 if reply_scope.correlation.is_none() {
103 reply_scope.correlation = Some(generate_correlation_id());
104 }
105 let mut store_scope = scope;
106 store_scope.correlation = None;
107 let session_key = StoreSessionKey::new(format!("{hint}::{}", store_scope.scope_hash()));
108 self.store
109 .register_wait(&ctx, &user, &store_scope, &session_key, data, None)
110 .map_err(map_store_error)?;
111 Ok(reply_scope)
112 }
113
114 pub fn clear(&self, envelope: &IngressEnvelope) -> GResult<()> {
115 let (ctx, user, _, scope) = build_store_ctx(envelope)?;
116 let mut scopes = vec![scope.clone()];
117 if scope.correlation.is_some() {
118 let mut base = scope;
119 base.correlation = None;
120 scopes.push(base);
121 }
122 for lookup in scopes {
123 self.store
124 .clear_wait(&ctx, &user, &lookup)
125 .map_err(map_store_error)?;
126 }
127 Ok(())
128 }
129
130 pub(crate) fn contact_identity(envelope: &IngressEnvelope) -> GResult<(TenantCtx, UserId)> {
137 let (ctx, user, _, _) = build_store_ctx(envelope)?;
138 Ok((ctx, user))
139 }
140}
141
142#[derive(Serialize, Deserialize)]
143struct FlowResumeRecord {
144 snapshot: FlowSnapshot,
145 #[serde(default)]
146 reason: Option<String>,
147}
148
149fn build_store_ctx(envelope: &IngressEnvelope) -> GResult<(TenantCtx, UserId, String, ReplyScope)> {
150 let base_hint = envelope
151 .session_hint
152 .clone()
153 .unwrap_or_else(|| envelope.canonical_session_hint());
154 let hint = if let Some(pack_id) = envelope.pack_id.as_deref() {
155 format!("{base_hint}::pack={pack_id}")
156 } else {
157 base_hint.clone()
158 };
159 let user = derive_user_id(&hint)?;
160 let scope = envelope
161 .reply_scope
162 .clone()
163 .ok_or_else(|| RunnerError::Session {
164 reason: "Cannot suspend: reply_scope missing; provider plugin must supply ReplyScope"
165 .to_string(),
166 })?;
167 let mut ctx = envelope.tenant_ctx();
168 ctx = ctx.with_session(hint.clone());
169 ctx = ctx.with_user(Some(user.clone()));
170 Ok((ctx, user, hint, scope))
171}
172
173fn record_to_session_data(
174 record: &FlowResumeRecord,
175 ctx: TenantCtx,
176 user: &UserId,
177 session_hint: &str,
178) -> GResult<SessionData> {
179 let flow = FlowId::from_str(record.snapshot.flow_id.as_str()).map_err(map_store_error)?;
180 let pack = PackId::from_str(record.snapshot.pack_id.as_str()).map_err(map_store_error)?;
181 let mut cursor = TypesSessionCursor::new(record.snapshot.next_node.clone());
182 if let Some(reason) = record.reason.clone() {
183 cursor = cursor.with_wait_reason(reason);
184 }
185 let context_json = serde_json::to_string(record).map_err(|err| RunnerError::Session {
186 reason: format!("failed to encode flow resume snapshot: {err}"),
187 })?;
188 let ctx = ctx
189 .with_user(Some(user.clone()))
190 .with_session(session_hint.to_string())
191 .with_flow(record.snapshot.flow_id.clone());
192 Ok(SessionData {
193 tenant_ctx: ctx,
194 flow_id: flow,
195 pack_id: Some(pack),
196 cursor,
197 context_json,
198 })
199}
200
201fn derive_user_id(hint: &str) -> GResult<UserId> {
202 let digest = Sha256::digest(hint.as_bytes());
203 let slug = format!("sess{}", hex::encode(&digest[..8]));
204 UserId::from_str(&slug).map_err(map_store_error)
205}
206
207fn map_store_error(err: GreenticError) -> RunnerError {
208 RunnerError::Session {
209 reason: err.to_string(),
210 }
211}
212
213fn generate_correlation_id() -> String {
214 let mut bytes = [0u8; 16];
215 rng().fill(&mut bytes);
216 hex::encode(bytes)
217}
218
219#[cfg(test)]
220mod tests {
221 use super::*;
222 use crate::runner::engine::ExecutionState;
223 use crate::storage::session::new_session_store;
224 use serde_json::json;
225
226 fn sample_envelope() -> IngressEnvelope {
227 IngressEnvelope {
228 tenant: "demo".into(),
229 env: Some("local".into()),
230 pack_id: Some("pack.demo".into()),
231 flow_id: "flow.main".into(),
232 flow_type: None,
233 action: Some("messaging".into()),
234 session_hint: Some("demo:provider:chan:conv:user".into()),
235 provider: Some("provider".into()),
236 messaging_endpoint_id: None,
237 channel: Some("chan".into()),
238 conversation: Some("conv".into()),
239 user: Some("user".into()),
240 entry_node: None,
241 activity_id: Some("act-1".into()),
242 timestamp: None,
243 payload: json!({ "text": "hi" }),
244 metadata: None,
245 reply_scope: Some(ReplyScope {
246 conversation: "conv".into(),
247 thread: None,
248 reply_to: None,
249 correlation: None,
250 }),
251 }
252 }
253
254 fn sample_wait() -> FlowWait {
255 let state: ExecutionState = serde_json::from_value(json!({
256 "input": { "text": "hi" },
257 "nodes": {},
258 "egress": []
259 }))
260 .expect("state");
261 FlowWait {
262 reason: Some("await-user".into()),
263 snapshot: FlowSnapshot {
264 pack_id: "pack.demo".into(),
265 flow_id: "flow.main".into(),
266 next_flow: None,
267 next_node: "node-2".into(),
268 awaiting_submit: false,
269 state,
270 },
271 }
272 }
273
274 #[test]
275 fn derive_user_id_is_stable() {
276 let hint = "some-tenant::session-key";
277 let a = derive_user_id(hint).unwrap();
278 let b = derive_user_id(hint).unwrap();
279 assert_eq!(a, b);
280 assert!(a.as_str().starts_with("sess"));
281 }
282
283 #[test]
284 fn resume_store_roundtrip() -> GResult<()> {
285 let store = FlowResumeStore::new(new_session_store());
286 let envelope = sample_envelope();
287 assert!(store.fetch(&envelope)?.is_none());
288
289 let wait = sample_wait();
290 let _ = store.save(&envelope, &wait)?;
291 let snapshot = store.fetch(&envelope)?.expect("snapshot missing");
292 assert_eq!(snapshot.flow_id, wait.snapshot.flow_id);
293 assert_eq!(snapshot.next_node, wait.snapshot.next_node);
294
295 store.clear(&envelope)?;
296 assert!(store.fetch(&envelope)?.is_none());
297 Ok(())
298 }
299
300 #[test]
301 fn resume_store_overwrites_existing() -> GResult<()> {
302 let store = FlowResumeStore::new(new_session_store());
303 let envelope = sample_envelope();
304 let mut wait = sample_wait();
305 let _ = store.save(&envelope, &wait)?;
306
307 wait.snapshot.next_node = "node-3".into();
308 wait.reason = Some("retry".into());
309 let _ = store.save(&envelope, &wait)?;
310
311 let snapshot = store.fetch(&envelope)?.expect("snapshot missing");
312 assert_eq!(snapshot.next_node, "node-3");
313 store.clear(&envelope)?;
314 Ok(())
315 }
316
317 #[test]
318 fn resume_store_uses_snapshot_even_if_envelope_flow_differs() -> GResult<()> {
319 let store = FlowResumeStore::new(new_session_store());
320 let envelope = sample_envelope();
321 let wait = sample_wait();
322 let _ = store.save(&envelope, &wait)?;
323
324 let mut redirected = envelope.clone();
325 redirected.flow_id = "flow.other".into();
326 let snapshot = store.fetch(&redirected)?.expect("snapshot missing");
327 assert_eq!(snapshot.flow_id, wait.snapshot.flow_id);
328
329 store.clear(&envelope)?;
330 Ok(())
331 }
332
333 #[test]
334 fn canonicalize_populates_defaults() {
335 let envelope = IngressEnvelope {
336 tenant: "demo".into(),
337 env: None,
338 pack_id: None,
339 flow_id: "flow.main".into(),
340 flow_type: None,
341 action: None,
342 session_hint: None,
343 provider: None,
344 messaging_endpoint_id: None,
345 channel: None,
346 conversation: None,
347 user: None,
348 entry_node: None,
349 activity_id: Some("activity-1".into()),
350 timestamp: None,
351 payload: json!({}),
352 metadata: None,
353 reply_scope: None,
354 }
355 .canonicalize();
356
357 assert_eq!(envelope.provider.as_deref(), Some("provider"));
358 assert_eq!(envelope.channel.as_deref(), Some("flow.main"));
359 assert_eq!(envelope.conversation.as_deref(), Some("flow.main"));
360 assert_eq!(envelope.user.as_deref(), Some("activity-1"));
361 assert!(envelope.session_hint.is_some());
362 }
363
364 #[test]
365 fn canonical_session_hint_is_pure_structured_form() {
366 let mut envelope = sample_envelope();
370 envelope.session_hint = None;
371 assert_eq!(
372 envelope.canonical_session_hint(),
373 "demo:provider:chan:conv:user"
374 );
375 envelope.messaging_endpoint_id = Some("teams-legal".into());
376 assert_eq!(
377 envelope.canonical_session_hint(),
378 "demo:provider:chan:conv:user"
379 );
380 }
381
382 #[test]
383 fn canonicalize_unchanged_when_endpoint_id_none() {
384 let mut derived = sample_envelope();
388 derived.session_hint = None;
389 let derived = derived.canonicalize();
390 assert_eq!(
391 derived.session_hint.as_deref(),
392 Some("demo:provider:chan:conv:user")
393 );
394
395 let explicit = sample_envelope().canonicalize();
396 assert_eq!(
397 explicit.session_hint.as_deref(),
398 Some("demo:provider:chan:conv:user")
399 );
400 }
401
402 #[test]
403 fn canonicalize_namespaces_derived_hint_when_endpoint_id_set() {
404 let mut envelope = sample_envelope();
405 envelope.session_hint = None;
406 envelope.messaging_endpoint_id = Some("teams-legal".into());
407 let envelope = envelope.canonicalize();
408 assert_eq!(
409 envelope.session_hint.as_deref(),
410 Some("ep=teams-legal::demo:provider:chan:conv:user")
411 );
412 }
413
414 #[test]
415 fn canonicalize_partitions_explicit_hints_across_endpoints() {
416 let raw = "shared-session-key";
421
422 let mut a = sample_envelope();
423 a.session_hint = Some(raw.into());
424 a.messaging_endpoint_id = Some("teams-legal".into());
425 let a = a.canonicalize();
426
427 let mut b = sample_envelope();
428 b.session_hint = Some(raw.into());
429 b.messaging_endpoint_id = Some("teams-accounting".into());
430 let b = b.canonicalize();
431
432 assert_ne!(a.session_hint, b.session_hint);
433 assert_eq!(
434 a.session_hint.as_deref(),
435 Some("ep=teams-legal::shared-session-key")
436 );
437 assert_eq!(
438 b.session_hint.as_deref(),
439 Some("ep=teams-accounting::shared-session-key")
440 );
441 }
442
443 #[test]
444 fn canonicalize_is_idempotent_for_endpoint_prefix() {
445 let mut envelope = sample_envelope();
447 envelope.session_hint = Some("raw-key".into());
448 envelope.messaging_endpoint_id = Some("teams-legal".into());
449 let once = envelope.canonicalize();
450 let twice = once.clone().canonicalize();
451 assert_eq!(once.session_hint, twice.session_hint);
452 assert_eq!(
453 twice.session_hint.as_deref(),
454 Some("ep=teams-legal::raw-key")
455 );
456 }
457
458 #[test]
459 fn canonicalize_drops_invalid_endpoint_id_to_none() {
460 let cases = [
465 ("empty", ""),
466 ("colon embedded", "teams:legal"), ("space", "teams legal"),
468 ("control char", "teams\nlegal"),
469 ("oversized", &"a".repeat(129)),
470 ];
471 for (label, bad) in cases {
472 let mut envelope = sample_envelope();
473 envelope.session_hint = Some("raw-key".into());
474 envelope.messaging_endpoint_id = Some(bad.into());
475 let canon = envelope.canonicalize();
476 assert!(
477 canon.messaging_endpoint_id.is_none(),
478 "{label}: invalid eid {bad:?} must drop to None"
479 );
480 assert_eq!(
482 canon.session_hint.as_deref(),
483 Some("raw-key"),
484 "{label}: dropped eid must leave hint un-namespaced"
485 );
486 }
487 }
488
489 #[test]
490 fn canonicalize_preserves_valid_endpoint_id_forms() {
491 for valid in ["teams-legal", "01HA1ABCDE", "teams_legal.v2"] {
494 let mut envelope = sample_envelope();
495 envelope.session_hint = Some("raw-key".into());
496 envelope.messaging_endpoint_id = Some(valid.into());
497 let canon = envelope.canonicalize();
498 assert_eq!(
499 canon.messaging_endpoint_id.as_deref(),
500 Some(valid),
501 "valid eid {valid:?} must be preserved"
502 );
503 assert_eq!(
504 canon.session_hint.as_deref(),
505 Some(format!("ep={valid}::raw-key").as_str()),
506 "valid eid {valid:?} must still prefix the hint"
507 );
508 }
509 }
510
511 #[test]
512 fn tenant_ctx_stamps_messaging_endpoint_id() {
513 let mut envelope = sample_envelope();
514 envelope.messaging_endpoint_id = Some("teams-legal".into());
515 let ctx = envelope.tenant_ctx();
516 assert_eq!(
517 ctx.attributes.get(attr_keys::MESSAGING_ENDPOINT_ID),
518 Some(&"teams-legal".to_string())
519 );
520 }
521
522 #[test]
523 fn tenant_ctx_omits_messaging_endpoint_id_when_unset() {
524 let envelope = sample_envelope();
525 let ctx = envelope.tenant_ctx();
526 assert!(
527 !ctx.attributes
528 .contains_key(attr_keys::MESSAGING_ENDPOINT_ID)
529 );
530 }
531}
532
533pub struct StateMachineRuntime {
534 runner: Runner,
535}
536
537impl StateMachineRuntime {
538 pub fn new(flows: Vec<FlowDefinition>) -> GResult<Self> {
540 let secrets = Arc::new(FnSecretsHost::new(|name| {
541 Err(RunnerError::Secrets {
542 reason: format!("secret {name} unavailable (noop host)"),
543 })
544 }));
545 let telemetry = Arc::new(FnTelemetryHost::new(|_, _| Ok(())));
546 let session = Arc::new(InMemorySessionHost::new());
547 let state = Arc::new(InMemoryStateHost::new());
548 let host = HostBundle::new(secrets, telemetry, session, state);
549
550 let adapters = AdapterRegistry::default();
551 let policy = Policy::default();
552
553 let mut builder = RunnerBuilder::new()
554 .with_host(host)
555 .with_adapters(adapters)
556 .with_policy(policy);
557 for flow in flows {
558 builder = builder.with_flow(flow);
559 }
560 let runner = builder.build()?;
561 Ok(Self { runner })
562 }
563
564 #[allow(clippy::too_many_arguments)]
566 pub fn from_flow_engine(
567 config: Arc<HostConfig>,
568 engine: Arc<FlowEngine>,
569 pack_trace: HashMap<String, PackTraceInfo>,
570 session_host: Arc<dyn SessionHost>,
571 session_store: DynSessionStore,
572 state_host: Arc<dyn StateHost>,
573 secrets_manager: DynSecretsManager,
574 mocks: Option<Arc<MockLayer>>,
575 audit_nats_client: Option<async_nats::Client>,
576 ) -> Result<Self> {
577 let policy = Arc::new(config.secrets_policy.clone());
578 let tenant_ctx = config.tenant_ctx();
579 let secrets = Arc::new(PolicySecretsHost::new(policy, secrets_manager, tenant_ctx));
580 let telemetry = Arc::new(FnTelemetryHost::new(|span, fields| {
581 tracing::debug!(?span, ?fields, "telemetry emit");
582 Ok(())
583 }));
584 let host = HostBundle::new(secrets, telemetry, session_host, state_host);
585 let resume_store = FlowResumeStore::new(session_store);
586
587 let mut adapters = AdapterRegistry::default();
588 adapters.register(
589 PACK_FLOW_ADAPTER,
590 Box::new(PackFlowAdapter::new(
591 Arc::clone(&config),
592 Arc::clone(&engine),
593 pack_trace,
594 resume_store,
595 mocks,
596 audit_nats_client,
597 )),
598 );
599
600 let flows = build_flow_definitions(engine.flows());
601 let mut builder = RunnerBuilder::new()
602 .with_host(host)
603 .with_adapters(adapters)
604 .with_policy(Policy::default());
605 for flow in flows {
606 builder = builder.with_flow(flow);
607 }
608 let runner = builder
609 .build()
610 .map_err(|err| anyhow!("state machine init failed: {err}"))?;
611 Ok(Self { runner })
612 }
613
614 pub async fn handle(&self, envelope: IngressEnvelope) -> Result<Value> {
616 let tenant_ctx = envelope.tenant_ctx();
617 let session_hint = envelope
618 .session_hint
619 .clone()
620 .unwrap_or_else(|| envelope.canonical_session_hint());
621 let pack_id = envelope.pack_id.clone().ok_or_else(|| {
622 anyhow!("pack_id missing; ingress must specify pack_id for multi-pack flows")
623 })?;
624 let input =
625 serde_json::to_value(&envelope).context("failed to serialise ingress envelope")?;
626 let request = RunFlowRequest {
627 tenant: tenant_ctx,
628 pack_id,
629 flow_id: envelope.flow_id.clone(),
630 input,
631 session_hint: Some(session_hint),
632 };
633 let result: super::api::RunFlowResult = self
634 .runner
635 .run_flow(request)
636 .await
637 .map_err(|err| anyhow!("flow execution failed: {err}"))?;
638 let outcome = result.outcome;
639 Ok(outcome.get("response").cloned().unwrap_or(outcome))
640 }
641}
642
643struct PolicySecretsHost {
644 policy: Arc<SecretsPolicy>,
645 manager: DynSecretsManager,
646 tenant_ctx: TenantCtx,
647}
648
649impl PolicySecretsHost {
650 fn new(policy: Arc<SecretsPolicy>, manager: DynSecretsManager, tenant_ctx: TenantCtx) -> Self {
651 Self {
652 policy,
653 manager,
654 tenant_ctx,
655 }
656 }
657}
658
659const POLICY_SECRETS_PACK_ID: &str = "_runner";
660
661#[async_trait]
662impl SecretsHost for PolicySecretsHost {
663 async fn get(&self, name: &str) -> GResult<String> {
664 if !self.policy.is_allowed(name) {
665 return Err(RunnerError::Secrets {
666 reason: format!("secret {name} denied by policy"),
667 });
668 }
669 let bytes = read_secret_blocking(
670 &self.manager,
671 &self.tenant_ctx,
672 POLICY_SECRETS_PACK_ID,
673 name,
674 )
675 .map_err(|err| RunnerError::Secrets {
676 reason: format!("secret {name} unavailable: {err}"),
677 })?;
678 String::from_utf8(bytes).map_err(|err| RunnerError::Secrets {
679 reason: format!("secret {name} not valid UTF-8: {err}"),
680 })
681 }
682}
683
684fn build_flow_definitions(flows: &[FlowDescriptor]) -> Vec<FlowDefinition> {
685 flows
686 .iter()
687 .map(|descriptor| {
688 FlowDefinition::new(
689 super::api::FlowSummary {
690 pack_id: descriptor.pack_id.clone(),
691 id: descriptor.id.clone(),
692 name: descriptor
693 .description
694 .clone()
695 .unwrap_or_else(|| descriptor.id.clone()),
696 version: descriptor.version.clone(),
697 description: descriptor.description.clone(),
698 },
699 serde_json::json!({
700 "type": "object"
701 }),
702 vec![FlowStep::Adapter(AdapterCall {
703 adapter: PACK_FLOW_ADAPTER.into(),
704 operation: descriptor.id.clone(),
705 payload: Value::String(PAYLOAD_FROM_LAST_INPUT.into()),
706 })],
707 )
708 })
709 .collect()
710}
711
712struct PackFlowAdapter {
713 tenant: String,
714 config: Arc<HostConfig>,
715 engine: Arc<FlowEngine>,
716 pack_trace: HashMap<String, PackTraceInfo>,
717 resume: FlowResumeStore,
718 mocks: Option<Arc<MockLayer>>,
719 audit_nats_client: Option<async_nats::Client>,
723}
724
725impl PackFlowAdapter {
726 fn new(
727 config: Arc<HostConfig>,
728 engine: Arc<FlowEngine>,
729 pack_trace: HashMap<String, PackTraceInfo>,
730 resume: FlowResumeStore,
731 mocks: Option<Arc<MockLayer>>,
732 audit_nats_client: Option<async_nats::Client>,
733 ) -> Self {
734 Self {
735 tenant: config.tenant.clone(),
736 config,
737 engine,
738 pack_trace,
739 resume,
740 mocks,
741 audit_nats_client,
742 }
743 }
744}
745
746#[async_trait::async_trait]
747impl Adapter for PackFlowAdapter {
748 async fn call(&self, call: &AdapterCall) -> GResult<Value> {
749 let envelope: IngressEnvelope =
750 serde_json::from_value(call.payload.clone()).map_err(|err| {
751 RunnerError::AdapterCall {
752 reason: format!("invalid ingress payload: {err}"),
753 }
754 })?;
755 let envelope = envelope.canonicalize();
756 let flow_id = call.operation.clone();
757 let action_owned = envelope.action.clone();
758 let session_owned = envelope
759 .session_hint
760 .clone()
761 .unwrap_or_else(|| envelope.canonical_session_hint());
762 let provider_owned = envelope.provider.clone();
763 let payload = envelope.payload.clone();
764 let retry_config = self.config.retry_config().into();
765 let resume_snapshot = self.resume.fetch(&envelope)?;
766 let resume_flow_id = resume_snapshot
767 .as_ref()
768 .and_then(|snapshot| snapshot.next_flow.clone())
769 .or_else(|| {
770 resume_snapshot
771 .as_ref()
772 .map(|snapshot| snapshot.flow_id.clone())
773 });
774 let effective_flow_id = resume_flow_id.clone().unwrap_or_else(|| flow_id.clone());
775 let effective_pack_id = if let Some(snapshot) = resume_snapshot.as_ref() {
776 snapshot.pack_id.clone()
777 } else if let Some(pack_id) = envelope.pack_id.as_deref() {
778 let found = self
779 .engine
780 .flow_by_key(pack_id, effective_flow_id.as_str())
781 .is_some();
782 if !found {
783 return Err(RunnerError::AdapterCall {
784 reason: format!(
785 "flow {} not registered for pack {pack_id}",
786 effective_flow_id
787 ),
788 });
789 }
790 pack_id.to_string()
791 } else if let Some(flow) = self.engine.flow_by_id(effective_flow_id.as_str()) {
792 flow.pack_id.clone()
793 } else {
794 return Err(RunnerError::AdapterCall {
795 reason: format!(
796 "flow {} is ambiguous; pack_id is required",
797 effective_flow_id
798 ),
799 });
800 };
801
802 let trace_config = self.config.trace.clone();
803 let flow_version = self
804 .engine
805 .flow_by_key(effective_pack_id.as_str(), effective_flow_id.as_str())
806 .map(|desc| desc.version.clone())
807 .unwrap_or_else(|| "unknown".to_string());
808 let pack_trace = self
809 .pack_trace
810 .get(effective_pack_id.as_str())
811 .cloned()
812 .unwrap_or_else(|| PackTraceInfo {
813 pack_ref: effective_pack_id.clone(),
814 resolved_digest: None,
815 });
816 let trace_ctx = TraceContext {
817 pack_ref: pack_trace.pack_ref,
818 resolved_digest: pack_trace.resolved_digest,
819 flow_id: effective_flow_id.clone(),
820 flow_version,
821 };
822 let trace = if trace_config.mode == TraceMode::Off {
823 None
824 } else {
825 let sink = self.audit_nats_client.clone().map(AuditSink::new);
833 let audit_tenant = sink.is_some().then(|| envelope.tenant_ctx());
834 Some(TraceRecorder::new_with_audit(
835 trace_config,
836 trace_ctx,
837 sink,
838 audit_tenant,
839 ))
840 };
841
842 let mocks = self.mocks.as_deref();
843 let ctx = FlowContext {
844 tenant: &self.tenant,
845 pack_id: effective_pack_id.as_str(),
846 flow_id: effective_flow_id.as_str(),
847 node_id: None,
848 tool: None,
849 action: action_owned.as_deref(),
850 session_id: Some(session_owned.as_str()),
851 provider_id: provider_owned.as_deref(),
852 reply_scope: envelope.reply_scope.as_ref(),
856 retry_config,
857 attempt: 1,
858 observer: trace
859 .as_ref()
860 .map(|recorder| recorder as &dyn crate::runner::engine::ExecutionObserver),
861 mocks,
862 };
863
864 let execution = if let Some(snapshot) = resume_snapshot {
865 let resume_pack_id = snapshot.pack_id.clone();
866 let resume_flow_id = snapshot
867 .next_flow
868 .clone()
869 .unwrap_or_else(|| snapshot.flow_id.clone());
870 let resume_ctx = FlowContext {
871 pack_id: resume_pack_id.as_str(),
872 flow_id: resume_flow_id.as_str(),
873 ..ctx
874 };
875 self.engine.resume(resume_ctx, snapshot, payload).await
876 } else if let Some(entry) = envelope.entry_node.as_deref() {
877 self.engine.execute_from(ctx, payload, entry).await
878 } else {
879 self.engine.execute(ctx, payload).await
880 };
881 let execution = match execution {
882 Ok(execution) => {
883 if let Some(recorder) = trace.as_ref()
884 && let Err(err) = recorder.flush_success()
885 {
886 tracing::warn!(error = %err, "failed to write trace");
887 }
888 execution
889 }
890 Err(err) => {
891 if let Some(recorder) = trace.as_ref()
892 && let Err(write_err) = recorder.flush_error(err.as_ref())
893 {
894 tracing::warn!(error = %write_err, "failed to write trace");
895 }
896 return Err(RunnerError::AdapterCall {
897 reason: err.to_string(),
898 });
899 }
900 };
901
902 match execution.status {
903 FlowStatus::Completed => {
904 self.resume.clear(&envelope)?;
905 Ok(execution.output)
906 }
907 FlowStatus::Waiting(wait) => {
908 let reply_scope = self.resume.save(&envelope, &wait)?;
909 Ok(json!({
910 "status": "pending",
911 "reason": wait.reason,
912 "resume": wait.snapshot,
913 "reply_scope": reply_scope,
914 "response": execution.output,
915 }))
916 }
917 }
918 }
919}
920
921#[derive(Clone, Debug, Serialize, Deserialize)]
922pub struct IngressEnvelope {
923 pub tenant: String,
924 #[serde(default, skip_serializing_if = "Option::is_none")]
925 pub env: Option<String>,
926 #[serde(default, skip_serializing_if = "Option::is_none")]
927 pub pack_id: Option<String>,
928 pub flow_id: String,
929 #[serde(default, skip_serializing_if = "Option::is_none")]
930 pub flow_type: Option<String>,
931 #[serde(default, skip_serializing_if = "Option::is_none")]
932 pub action: Option<String>,
933 #[serde(default, skip_serializing_if = "Option::is_none")]
934 pub session_hint: Option<String>,
935 #[serde(default, skip_serializing_if = "Option::is_none")]
936 pub provider: Option<String>,
937 #[serde(default, skip_serializing_if = "Option::is_none")]
942 pub messaging_endpoint_id: Option<String>,
943 #[serde(default, skip_serializing_if = "Option::is_none")]
944 pub channel: Option<String>,
945 #[serde(default, skip_serializing_if = "Option::is_none")]
946 pub conversation: Option<String>,
947 #[serde(default, skip_serializing_if = "Option::is_none")]
948 pub user: Option<String>,
949 #[serde(default, skip_serializing_if = "Option::is_none")]
956 pub entry_node: Option<String>,
957 #[serde(default, skip_serializing_if = "Option::is_none")]
958 pub activity_id: Option<String>,
959 #[serde(default, skip_serializing_if = "Option::is_none")]
960 pub timestamp: Option<String>,
961 #[serde(default)]
962 pub payload: Value,
963 #[serde(default, skip_serializing_if = "Option::is_none")]
964 pub metadata: Option<Value>,
965 #[serde(default, skip_serializing_if = "Option::is_none")]
966 pub reply_scope: Option<ReplyScope>,
967}
968
969fn endpoint_id_is_valid(raw: &str) -> bool {
985 if raw.is_empty() || raw.len() > 128 {
986 return false;
987 }
988 raw.bytes()
989 .all(|b| b.is_ascii_alphanumeric() || matches!(b, b'-' | b'_' | b'.'))
990}
991
992impl IngressEnvelope {
993 pub fn canonicalize(mut self) -> Self {
994 if self.provider.is_none() {
995 self.provider = Some("provider".into());
996 }
997 if self.channel.is_none() {
998 self.channel = Some(self.flow_id.clone());
999 }
1000 if self.conversation.is_none() {
1001 self.conversation = self.channel.clone();
1002 }
1003 if self.user.is_none() {
1004 if let Some(ref hint) = self.session_hint {
1005 self.user = Some(hint.clone());
1006 } else if let Some(ref activity) = self.activity_id {
1007 self.user = Some(activity.clone());
1008 } else {
1009 self.user = Some("user".into());
1010 }
1011 }
1012 if let Some(eid) = self.messaging_endpoint_id.as_deref()
1024 && !endpoint_id_is_valid(eid)
1025 {
1026 tracing::warn!(
1027 tenant = %self.tenant,
1028 messaging_endpoint_id = %eid,
1029 "M1.4: invalid messaging_endpoint_id dropped at runner-host canonicalize \
1030 — request will run unscoped (legacy session bucket). \
1031 Producer bug or non-HTTP caller bypassing the boundary validator."
1032 );
1033 self.messaging_endpoint_id = None;
1034 }
1035 let base = self
1039 .session_hint
1040 .clone()
1041 .unwrap_or_else(|| self.canonical_session_hint());
1042 self.session_hint = Some(match &self.messaging_endpoint_id {
1043 Some(eid) => {
1044 let prefix = format!("ep={eid}::");
1045 if base.starts_with(&prefix) {
1046 base
1047 } else {
1048 format!("{prefix}{base}")
1049 }
1050 }
1051 None => base,
1052 });
1053 if self.reply_scope.is_none()
1054 && let Some(conversation) = self.conversation.clone()
1055 {
1056 self.reply_scope = Some(ReplyScope {
1057 conversation,
1058 thread: None,
1059 reply_to: None,
1060 correlation: None,
1061 });
1062 }
1063 self
1064 }
1065
1066 pub fn canonical_session_hint(&self) -> String {
1067 format!(
1072 "{}:{}:{}:{}:{}",
1073 self.tenant,
1074 self.provider.as_deref().unwrap_or("provider"),
1075 self.channel.as_deref().unwrap_or("channel"),
1076 self.conversation.as_deref().unwrap_or("conversation"),
1077 self.user.as_deref().unwrap_or("user")
1078 )
1079 }
1080
1081 pub fn tenant_ctx(&self) -> TenantCtx {
1082 let env_raw = self.env.clone().unwrap_or_else(|| DEFAULT_ENV.into());
1083 let env = EnvId::from_str(env_raw.as_str())
1084 .unwrap_or_else(|_| EnvId::from_str(DEFAULT_ENV).expect("default env must be valid"));
1085 let tenant_id = TenantId::from_str(self.tenant.as_str()).unwrap_or_else(|_| {
1086 TenantId::from_str("tenant.default").expect("tenant fallback must be valid")
1087 });
1088 let mut ctx = TenantCtx::new(env, tenant_id).with_flow(self.flow_id.clone());
1089 if let Some(provider) = &self.provider {
1090 ctx = ctx.with_provider(provider.clone());
1091 }
1092 if let Some(session) = &self.session_hint {
1093 ctx = ctx.with_session(session.clone());
1094 }
1095 if let Some(eid) = &self.messaging_endpoint_id {
1099 ctx.attributes
1100 .insert(attr_keys::MESSAGING_ENDPOINT_ID.to_string(), eid.clone());
1101 }
1102 ctx
1103 }
1104}