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 activity_id: Some("act-1".into()),
241 timestamp: None,
242 payload: json!({ "text": "hi" }),
243 metadata: None,
244 reply_scope: Some(ReplyScope {
245 conversation: "conv".into(),
246 thread: None,
247 reply_to: None,
248 correlation: None,
249 }),
250 }
251 }
252
253 fn sample_wait() -> FlowWait {
254 let state: ExecutionState = serde_json::from_value(json!({
255 "input": { "text": "hi" },
256 "nodes": {},
257 "egress": []
258 }))
259 .expect("state");
260 FlowWait {
261 reason: Some("await-user".into()),
262 snapshot: FlowSnapshot {
263 pack_id: "pack.demo".into(),
264 flow_id: "flow.main".into(),
265 next_flow: None,
266 next_node: "node-2".into(),
267 state,
268 },
269 }
270 }
271
272 #[test]
273 fn derive_user_id_is_stable() {
274 let hint = "some-tenant::session-key";
275 let a = derive_user_id(hint).unwrap();
276 let b = derive_user_id(hint).unwrap();
277 assert_eq!(a, b);
278 assert!(a.as_str().starts_with("sess"));
279 }
280
281 #[test]
282 fn resume_store_roundtrip() -> GResult<()> {
283 let store = FlowResumeStore::new(new_session_store());
284 let envelope = sample_envelope();
285 assert!(store.fetch(&envelope)?.is_none());
286
287 let wait = sample_wait();
288 let _ = store.save(&envelope, &wait)?;
289 let snapshot = store.fetch(&envelope)?.expect("snapshot missing");
290 assert_eq!(snapshot.flow_id, wait.snapshot.flow_id);
291 assert_eq!(snapshot.next_node, wait.snapshot.next_node);
292
293 store.clear(&envelope)?;
294 assert!(store.fetch(&envelope)?.is_none());
295 Ok(())
296 }
297
298 #[test]
299 fn resume_store_overwrites_existing() -> GResult<()> {
300 let store = FlowResumeStore::new(new_session_store());
301 let envelope = sample_envelope();
302 let mut wait = sample_wait();
303 let _ = store.save(&envelope, &wait)?;
304
305 wait.snapshot.next_node = "node-3".into();
306 wait.reason = Some("retry".into());
307 let _ = store.save(&envelope, &wait)?;
308
309 let snapshot = store.fetch(&envelope)?.expect("snapshot missing");
310 assert_eq!(snapshot.next_node, "node-3");
311 store.clear(&envelope)?;
312 Ok(())
313 }
314
315 #[test]
316 fn resume_store_uses_snapshot_even_if_envelope_flow_differs() -> GResult<()> {
317 let store = FlowResumeStore::new(new_session_store());
318 let envelope = sample_envelope();
319 let wait = sample_wait();
320 let _ = store.save(&envelope, &wait)?;
321
322 let mut redirected = envelope.clone();
323 redirected.flow_id = "flow.other".into();
324 let snapshot = store.fetch(&redirected)?.expect("snapshot missing");
325 assert_eq!(snapshot.flow_id, wait.snapshot.flow_id);
326
327 store.clear(&envelope)?;
328 Ok(())
329 }
330
331 #[test]
332 fn canonicalize_populates_defaults() {
333 let envelope = IngressEnvelope {
334 tenant: "demo".into(),
335 env: None,
336 pack_id: None,
337 flow_id: "flow.main".into(),
338 flow_type: None,
339 action: None,
340 session_hint: None,
341 provider: None,
342 messaging_endpoint_id: None,
343 channel: None,
344 conversation: None,
345 user: None,
346 activity_id: Some("activity-1".into()),
347 timestamp: None,
348 payload: json!({}),
349 metadata: None,
350 reply_scope: None,
351 }
352 .canonicalize();
353
354 assert_eq!(envelope.provider.as_deref(), Some("provider"));
355 assert_eq!(envelope.channel.as_deref(), Some("flow.main"));
356 assert_eq!(envelope.conversation.as_deref(), Some("flow.main"));
357 assert_eq!(envelope.user.as_deref(), Some("activity-1"));
358 assert!(envelope.session_hint.is_some());
359 }
360
361 #[test]
362 fn canonical_session_hint_is_pure_structured_form() {
363 let mut envelope = sample_envelope();
367 envelope.session_hint = None;
368 assert_eq!(
369 envelope.canonical_session_hint(),
370 "demo:provider:chan:conv:user"
371 );
372 envelope.messaging_endpoint_id = Some("teams-legal".into());
373 assert_eq!(
374 envelope.canonical_session_hint(),
375 "demo:provider:chan:conv:user"
376 );
377 }
378
379 #[test]
380 fn canonicalize_unchanged_when_endpoint_id_none() {
381 let mut derived = sample_envelope();
385 derived.session_hint = None;
386 let derived = derived.canonicalize();
387 assert_eq!(
388 derived.session_hint.as_deref(),
389 Some("demo:provider:chan:conv:user")
390 );
391
392 let explicit = sample_envelope().canonicalize();
393 assert_eq!(
394 explicit.session_hint.as_deref(),
395 Some("demo:provider:chan:conv:user")
396 );
397 }
398
399 #[test]
400 fn canonicalize_namespaces_derived_hint_when_endpoint_id_set() {
401 let mut envelope = sample_envelope();
402 envelope.session_hint = None;
403 envelope.messaging_endpoint_id = Some("teams-legal".into());
404 let envelope = envelope.canonicalize();
405 assert_eq!(
406 envelope.session_hint.as_deref(),
407 Some("ep=teams-legal::demo:provider:chan:conv:user")
408 );
409 }
410
411 #[test]
412 fn canonicalize_partitions_explicit_hints_across_endpoints() {
413 let raw = "shared-session-key";
418
419 let mut a = sample_envelope();
420 a.session_hint = Some(raw.into());
421 a.messaging_endpoint_id = Some("teams-legal".into());
422 let a = a.canonicalize();
423
424 let mut b = sample_envelope();
425 b.session_hint = Some(raw.into());
426 b.messaging_endpoint_id = Some("teams-accounting".into());
427 let b = b.canonicalize();
428
429 assert_ne!(a.session_hint, b.session_hint);
430 assert_eq!(
431 a.session_hint.as_deref(),
432 Some("ep=teams-legal::shared-session-key")
433 );
434 assert_eq!(
435 b.session_hint.as_deref(),
436 Some("ep=teams-accounting::shared-session-key")
437 );
438 }
439
440 #[test]
441 fn canonicalize_is_idempotent_for_endpoint_prefix() {
442 let mut envelope = sample_envelope();
444 envelope.session_hint = Some("raw-key".into());
445 envelope.messaging_endpoint_id = Some("teams-legal".into());
446 let once = envelope.canonicalize();
447 let twice = once.clone().canonicalize();
448 assert_eq!(once.session_hint, twice.session_hint);
449 assert_eq!(
450 twice.session_hint.as_deref(),
451 Some("ep=teams-legal::raw-key")
452 );
453 }
454
455 #[test]
456 fn canonicalize_drops_invalid_endpoint_id_to_none() {
457 let cases = [
462 ("empty", ""),
463 ("colon embedded", "teams:legal"), ("space", "teams legal"),
465 ("control char", "teams\nlegal"),
466 ("oversized", &"a".repeat(129)),
467 ];
468 for (label, bad) in cases {
469 let mut envelope = sample_envelope();
470 envelope.session_hint = Some("raw-key".into());
471 envelope.messaging_endpoint_id = Some(bad.into());
472 let canon = envelope.canonicalize();
473 assert!(
474 canon.messaging_endpoint_id.is_none(),
475 "{label}: invalid eid {bad:?} must drop to None"
476 );
477 assert_eq!(
479 canon.session_hint.as_deref(),
480 Some("raw-key"),
481 "{label}: dropped eid must leave hint un-namespaced"
482 );
483 }
484 }
485
486 #[test]
487 fn canonicalize_preserves_valid_endpoint_id_forms() {
488 for valid in ["teams-legal", "01HA1ABCDE", "teams_legal.v2"] {
491 let mut envelope = sample_envelope();
492 envelope.session_hint = Some("raw-key".into());
493 envelope.messaging_endpoint_id = Some(valid.into());
494 let canon = envelope.canonicalize();
495 assert_eq!(
496 canon.messaging_endpoint_id.as_deref(),
497 Some(valid),
498 "valid eid {valid:?} must be preserved"
499 );
500 assert_eq!(
501 canon.session_hint.as_deref(),
502 Some(format!("ep={valid}::raw-key").as_str()),
503 "valid eid {valid:?} must still prefix the hint"
504 );
505 }
506 }
507
508 #[test]
509 fn tenant_ctx_stamps_messaging_endpoint_id() {
510 let mut envelope = sample_envelope();
511 envelope.messaging_endpoint_id = Some("teams-legal".into());
512 let ctx = envelope.tenant_ctx();
513 assert_eq!(
514 ctx.attributes.get(attr_keys::MESSAGING_ENDPOINT_ID),
515 Some(&"teams-legal".to_string())
516 );
517 }
518
519 #[test]
520 fn tenant_ctx_omits_messaging_endpoint_id_when_unset() {
521 let envelope = sample_envelope();
522 let ctx = envelope.tenant_ctx();
523 assert!(
524 !ctx.attributes
525 .contains_key(attr_keys::MESSAGING_ENDPOINT_ID)
526 );
527 }
528}
529
530pub struct StateMachineRuntime {
531 runner: Runner,
532}
533
534impl StateMachineRuntime {
535 pub fn new(flows: Vec<FlowDefinition>) -> GResult<Self> {
537 let secrets = Arc::new(FnSecretsHost::new(|name| {
538 Err(RunnerError::Secrets {
539 reason: format!("secret {name} unavailable (noop host)"),
540 })
541 }));
542 let telemetry = Arc::new(FnTelemetryHost::new(|_, _| Ok(())));
543 let session = Arc::new(InMemorySessionHost::new());
544 let state = Arc::new(InMemoryStateHost::new());
545 let host = HostBundle::new(secrets, telemetry, session, state);
546
547 let adapters = AdapterRegistry::default();
548 let policy = Policy::default();
549
550 let mut builder = RunnerBuilder::new()
551 .with_host(host)
552 .with_adapters(adapters)
553 .with_policy(policy);
554 for flow in flows {
555 builder = builder.with_flow(flow);
556 }
557 let runner = builder.build()?;
558 Ok(Self { runner })
559 }
560
561 #[allow(clippy::too_many_arguments)]
563 pub fn from_flow_engine(
564 config: Arc<HostConfig>,
565 engine: Arc<FlowEngine>,
566 pack_trace: HashMap<String, PackTraceInfo>,
567 session_host: Arc<dyn SessionHost>,
568 session_store: DynSessionStore,
569 state_host: Arc<dyn StateHost>,
570 secrets_manager: DynSecretsManager,
571 mocks: Option<Arc<MockLayer>>,
572 audit_nats_client: Option<async_nats::Client>,
573 ) -> Result<Self> {
574 let policy = Arc::new(config.secrets_policy.clone());
575 let tenant_ctx = config.tenant_ctx();
576 let secrets = Arc::new(PolicySecretsHost::new(policy, secrets_manager, tenant_ctx));
577 let telemetry = Arc::new(FnTelemetryHost::new(|span, fields| {
578 tracing::debug!(?span, ?fields, "telemetry emit");
579 Ok(())
580 }));
581 let host = HostBundle::new(secrets, telemetry, session_host, state_host);
582 let resume_store = FlowResumeStore::new(session_store);
583
584 let mut adapters = AdapterRegistry::default();
585 adapters.register(
586 PACK_FLOW_ADAPTER,
587 Box::new(PackFlowAdapter::new(
588 Arc::clone(&config),
589 Arc::clone(&engine),
590 pack_trace,
591 resume_store,
592 mocks,
593 audit_nats_client,
594 )),
595 );
596
597 let flows = build_flow_definitions(engine.flows());
598 let mut builder = RunnerBuilder::new()
599 .with_host(host)
600 .with_adapters(adapters)
601 .with_policy(Policy::default());
602 for flow in flows {
603 builder = builder.with_flow(flow);
604 }
605 let runner = builder
606 .build()
607 .map_err(|err| anyhow!("state machine init failed: {err}"))?;
608 Ok(Self { runner })
609 }
610
611 pub async fn handle(&self, envelope: IngressEnvelope) -> Result<Value> {
613 let tenant_ctx = envelope.tenant_ctx();
614 let session_hint = envelope
615 .session_hint
616 .clone()
617 .unwrap_or_else(|| envelope.canonical_session_hint());
618 let pack_id = envelope.pack_id.clone().ok_or_else(|| {
619 anyhow!("pack_id missing; ingress must specify pack_id for multi-pack flows")
620 })?;
621 let input =
622 serde_json::to_value(&envelope).context("failed to serialise ingress envelope")?;
623 let request = RunFlowRequest {
624 tenant: tenant_ctx,
625 pack_id,
626 flow_id: envelope.flow_id.clone(),
627 input,
628 session_hint: Some(session_hint),
629 };
630 let result: super::api::RunFlowResult = self
631 .runner
632 .run_flow(request)
633 .await
634 .map_err(|err| anyhow!("flow execution failed: {err}"))?;
635 let outcome = result.outcome;
636 Ok(outcome.get("response").cloned().unwrap_or(outcome))
637 }
638}
639
640struct PolicySecretsHost {
641 policy: Arc<SecretsPolicy>,
642 manager: DynSecretsManager,
643 tenant_ctx: TenantCtx,
644}
645
646impl PolicySecretsHost {
647 fn new(policy: Arc<SecretsPolicy>, manager: DynSecretsManager, tenant_ctx: TenantCtx) -> Self {
648 Self {
649 policy,
650 manager,
651 tenant_ctx,
652 }
653 }
654}
655
656const POLICY_SECRETS_PACK_ID: &str = "_runner";
657
658#[async_trait]
659impl SecretsHost for PolicySecretsHost {
660 async fn get(&self, name: &str) -> GResult<String> {
661 if !self.policy.is_allowed(name) {
662 return Err(RunnerError::Secrets {
663 reason: format!("secret {name} denied by policy"),
664 });
665 }
666 let bytes = read_secret_blocking(
667 &self.manager,
668 &self.tenant_ctx,
669 POLICY_SECRETS_PACK_ID,
670 name,
671 )
672 .map_err(|err| RunnerError::Secrets {
673 reason: format!("secret {name} unavailable: {err}"),
674 })?;
675 String::from_utf8(bytes).map_err(|err| RunnerError::Secrets {
676 reason: format!("secret {name} not valid UTF-8: {err}"),
677 })
678 }
679}
680
681fn build_flow_definitions(flows: &[FlowDescriptor]) -> Vec<FlowDefinition> {
682 flows
683 .iter()
684 .map(|descriptor| {
685 FlowDefinition::new(
686 super::api::FlowSummary {
687 pack_id: descriptor.pack_id.clone(),
688 id: descriptor.id.clone(),
689 name: descriptor
690 .description
691 .clone()
692 .unwrap_or_else(|| descriptor.id.clone()),
693 version: descriptor.version.clone(),
694 description: descriptor.description.clone(),
695 },
696 serde_json::json!({
697 "type": "object"
698 }),
699 vec![FlowStep::Adapter(AdapterCall {
700 adapter: PACK_FLOW_ADAPTER.into(),
701 operation: descriptor.id.clone(),
702 payload: Value::String(PAYLOAD_FROM_LAST_INPUT.into()),
703 })],
704 )
705 })
706 .collect()
707}
708
709struct PackFlowAdapter {
710 tenant: String,
711 config: Arc<HostConfig>,
712 engine: Arc<FlowEngine>,
713 pack_trace: HashMap<String, PackTraceInfo>,
714 resume: FlowResumeStore,
715 mocks: Option<Arc<MockLayer>>,
716 audit_nats_client: Option<async_nats::Client>,
720}
721
722impl PackFlowAdapter {
723 fn new(
724 config: Arc<HostConfig>,
725 engine: Arc<FlowEngine>,
726 pack_trace: HashMap<String, PackTraceInfo>,
727 resume: FlowResumeStore,
728 mocks: Option<Arc<MockLayer>>,
729 audit_nats_client: Option<async_nats::Client>,
730 ) -> Self {
731 Self {
732 tenant: config.tenant.clone(),
733 config,
734 engine,
735 pack_trace,
736 resume,
737 mocks,
738 audit_nats_client,
739 }
740 }
741}
742
743#[async_trait::async_trait]
744impl Adapter for PackFlowAdapter {
745 async fn call(&self, call: &AdapterCall) -> GResult<Value> {
746 let envelope: IngressEnvelope =
747 serde_json::from_value(call.payload.clone()).map_err(|err| {
748 RunnerError::AdapterCall {
749 reason: format!("invalid ingress payload: {err}"),
750 }
751 })?;
752 let envelope = envelope.canonicalize();
753 let flow_id = call.operation.clone();
754 let action_owned = envelope.action.clone();
755 let session_owned = envelope
756 .session_hint
757 .clone()
758 .unwrap_or_else(|| envelope.canonical_session_hint());
759 let provider_owned = envelope.provider.clone();
760 let payload = envelope.payload.clone();
761 let retry_config = self.config.retry_config().into();
762 let resume_snapshot = self.resume.fetch(&envelope)?;
763 let resume_flow_id = resume_snapshot
764 .as_ref()
765 .and_then(|snapshot| snapshot.next_flow.clone())
766 .or_else(|| {
767 resume_snapshot
768 .as_ref()
769 .map(|snapshot| snapshot.flow_id.clone())
770 });
771 let effective_flow_id = resume_flow_id.clone().unwrap_or_else(|| flow_id.clone());
772 let effective_pack_id = if let Some(snapshot) = resume_snapshot.as_ref() {
773 snapshot.pack_id.clone()
774 } else if let Some(pack_id) = envelope.pack_id.as_deref() {
775 let found = self
776 .engine
777 .flow_by_key(pack_id, effective_flow_id.as_str())
778 .is_some();
779 if !found {
780 return Err(RunnerError::AdapterCall {
781 reason: format!(
782 "flow {} not registered for pack {pack_id}",
783 effective_flow_id
784 ),
785 });
786 }
787 pack_id.to_string()
788 } else if let Some(flow) = self.engine.flow_by_id(effective_flow_id.as_str()) {
789 flow.pack_id.clone()
790 } else {
791 return Err(RunnerError::AdapterCall {
792 reason: format!(
793 "flow {} is ambiguous; pack_id is required",
794 effective_flow_id
795 ),
796 });
797 };
798
799 let trace_config = self.config.trace.clone();
800 let flow_version = self
801 .engine
802 .flow_by_key(effective_pack_id.as_str(), effective_flow_id.as_str())
803 .map(|desc| desc.version.clone())
804 .unwrap_or_else(|| "unknown".to_string());
805 let pack_trace = self
806 .pack_trace
807 .get(effective_pack_id.as_str())
808 .cloned()
809 .unwrap_or_else(|| PackTraceInfo {
810 pack_ref: effective_pack_id.clone(),
811 resolved_digest: None,
812 });
813 let trace_ctx = TraceContext {
814 pack_ref: pack_trace.pack_ref,
815 resolved_digest: pack_trace.resolved_digest,
816 flow_id: effective_flow_id.clone(),
817 flow_version,
818 };
819 let trace = if trace_config.mode == TraceMode::Off {
820 None
821 } else {
822 let sink = self.audit_nats_client.clone().map(AuditSink::new);
830 let audit_tenant = sink.is_some().then(|| envelope.tenant_ctx());
831 Some(TraceRecorder::new_with_audit(
832 trace_config,
833 trace_ctx,
834 sink,
835 audit_tenant,
836 ))
837 };
838
839 let mocks = self.mocks.as_deref();
840 let ctx = FlowContext {
841 tenant: &self.tenant,
842 pack_id: effective_pack_id.as_str(),
843 flow_id: effective_flow_id.as_str(),
844 node_id: None,
845 tool: None,
846 action: action_owned.as_deref(),
847 session_id: Some(session_owned.as_str()),
848 provider_id: provider_owned.as_deref(),
849 reply_scope: envelope.reply_scope.as_ref(),
853 retry_config,
854 attempt: 1,
855 observer: trace
856 .as_ref()
857 .map(|recorder| recorder as &dyn crate::runner::engine::ExecutionObserver),
858 mocks,
859 };
860
861 let execution = if let Some(snapshot) = resume_snapshot {
862 let resume_pack_id = snapshot.pack_id.clone();
863 let resume_flow_id = snapshot
864 .next_flow
865 .clone()
866 .unwrap_or_else(|| snapshot.flow_id.clone());
867 let resume_ctx = FlowContext {
868 pack_id: resume_pack_id.as_str(),
869 flow_id: resume_flow_id.as_str(),
870 ..ctx
871 };
872 self.engine.resume(resume_ctx, snapshot, payload).await
873 } else {
874 self.engine.execute(ctx, payload).await
875 };
876 let execution = match execution {
877 Ok(execution) => {
878 if let Some(recorder) = trace.as_ref()
879 && let Err(err) = recorder.flush_success()
880 {
881 tracing::warn!(error = %err, "failed to write trace");
882 }
883 execution
884 }
885 Err(err) => {
886 if let Some(recorder) = trace.as_ref()
887 && let Err(write_err) = recorder.flush_error(err.as_ref())
888 {
889 tracing::warn!(error = %write_err, "failed to write trace");
890 }
891 return Err(RunnerError::AdapterCall {
892 reason: err.to_string(),
893 });
894 }
895 };
896
897 match execution.status {
898 FlowStatus::Completed => {
899 self.resume.clear(&envelope)?;
900 Ok(execution.output)
901 }
902 FlowStatus::Waiting(wait) => {
903 let reply_scope = self.resume.save(&envelope, &wait)?;
904 Ok(json!({
905 "status": "pending",
906 "reason": wait.reason,
907 "resume": wait.snapshot,
908 "reply_scope": reply_scope,
909 "response": execution.output,
910 }))
911 }
912 }
913 }
914}
915
916#[derive(Clone, Debug, Serialize, Deserialize)]
917pub struct IngressEnvelope {
918 pub tenant: String,
919 #[serde(default, skip_serializing_if = "Option::is_none")]
920 pub env: Option<String>,
921 #[serde(default, skip_serializing_if = "Option::is_none")]
922 pub pack_id: Option<String>,
923 pub flow_id: String,
924 #[serde(default, skip_serializing_if = "Option::is_none")]
925 pub flow_type: Option<String>,
926 #[serde(default, skip_serializing_if = "Option::is_none")]
927 pub action: Option<String>,
928 #[serde(default, skip_serializing_if = "Option::is_none")]
929 pub session_hint: Option<String>,
930 #[serde(default, skip_serializing_if = "Option::is_none")]
931 pub provider: Option<String>,
932 #[serde(default, skip_serializing_if = "Option::is_none")]
937 pub messaging_endpoint_id: Option<String>,
938 #[serde(default, skip_serializing_if = "Option::is_none")]
939 pub channel: Option<String>,
940 #[serde(default, skip_serializing_if = "Option::is_none")]
941 pub conversation: Option<String>,
942 #[serde(default, skip_serializing_if = "Option::is_none")]
943 pub user: Option<String>,
944 #[serde(default, skip_serializing_if = "Option::is_none")]
945 pub activity_id: Option<String>,
946 #[serde(default, skip_serializing_if = "Option::is_none")]
947 pub timestamp: Option<String>,
948 #[serde(default)]
949 pub payload: Value,
950 #[serde(default, skip_serializing_if = "Option::is_none")]
951 pub metadata: Option<Value>,
952 #[serde(default, skip_serializing_if = "Option::is_none")]
953 pub reply_scope: Option<ReplyScope>,
954}
955
956fn endpoint_id_is_valid(raw: &str) -> bool {
972 if raw.is_empty() || raw.len() > 128 {
973 return false;
974 }
975 raw.bytes()
976 .all(|b| b.is_ascii_alphanumeric() || matches!(b, b'-' | b'_' | b'.'))
977}
978
979impl IngressEnvelope {
980 pub fn canonicalize(mut self) -> Self {
981 if self.provider.is_none() {
982 self.provider = Some("provider".into());
983 }
984 if self.channel.is_none() {
985 self.channel = Some(self.flow_id.clone());
986 }
987 if self.conversation.is_none() {
988 self.conversation = self.channel.clone();
989 }
990 if self.user.is_none() {
991 if let Some(ref hint) = self.session_hint {
992 self.user = Some(hint.clone());
993 } else if let Some(ref activity) = self.activity_id {
994 self.user = Some(activity.clone());
995 } else {
996 self.user = Some("user".into());
997 }
998 }
999 if let Some(eid) = self.messaging_endpoint_id.as_deref()
1011 && !endpoint_id_is_valid(eid)
1012 {
1013 tracing::warn!(
1014 tenant = %self.tenant,
1015 messaging_endpoint_id = %eid,
1016 "M1.4: invalid messaging_endpoint_id dropped at runner-host canonicalize \
1017 — request will run unscoped (legacy session bucket). \
1018 Producer bug or non-HTTP caller bypassing the boundary validator."
1019 );
1020 self.messaging_endpoint_id = None;
1021 }
1022 let base = self
1026 .session_hint
1027 .clone()
1028 .unwrap_or_else(|| self.canonical_session_hint());
1029 self.session_hint = Some(match &self.messaging_endpoint_id {
1030 Some(eid) => {
1031 let prefix = format!("ep={eid}::");
1032 if base.starts_with(&prefix) {
1033 base
1034 } else {
1035 format!("{prefix}{base}")
1036 }
1037 }
1038 None => base,
1039 });
1040 if self.reply_scope.is_none()
1041 && let Some(conversation) = self.conversation.clone()
1042 {
1043 self.reply_scope = Some(ReplyScope {
1044 conversation,
1045 thread: None,
1046 reply_to: None,
1047 correlation: None,
1048 });
1049 }
1050 self
1051 }
1052
1053 pub fn canonical_session_hint(&self) -> String {
1054 format!(
1059 "{}:{}:{}:{}:{}",
1060 self.tenant,
1061 self.provider.as_deref().unwrap_or("provider"),
1062 self.channel.as_deref().unwrap_or("channel"),
1063 self.conversation.as_deref().unwrap_or("conversation"),
1064 self.user.as_deref().unwrap_or("user")
1065 )
1066 }
1067
1068 pub fn tenant_ctx(&self) -> TenantCtx {
1069 let env_raw = self.env.clone().unwrap_or_else(|| DEFAULT_ENV.into());
1070 let env = EnvId::from_str(env_raw.as_str())
1071 .unwrap_or_else(|_| EnvId::from_str(DEFAULT_ENV).expect("default env must be valid"));
1072 let tenant_id = TenantId::from_str(self.tenant.as_str()).unwrap_or_else(|_| {
1073 TenantId::from_str("tenant.default").expect("tenant fallback must be valid")
1074 });
1075 let mut ctx = TenantCtx::new(env, tenant_id).with_flow(self.flow_id.clone());
1076 if let Some(provider) = &self.provider {
1077 ctx = ctx.with_provider(provider.clone());
1078 }
1079 if let Some(session) = &self.session_hint {
1080 ctx = ctx.with_session(session.clone());
1081 }
1082 if let Some(eid) = &self.messaging_endpoint_id {
1086 ctx.attributes
1087 .insert(attr_keys::MESSAGING_ENDPOINT_ID.to_string(), eid.clone());
1088 }
1089 ctx
1090 }
1091}