Skip to main content

greentic_runner_host/engine/
runtime.rs

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    /// Returns the `(tenant_ctx, user_id)` pair this store derives from
131    /// `envelope` for wait bucketing.
132    ///
133    /// Exposed so siblings — currently the M1.5 welcome-seen marker — can
134    /// partition by the SAME identity instead of re-deriving the digest and
135    /// risking drift.
136    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        // canonical_session_hint() does NOT embed messaging_endpoint_id —
367        // namespacing happens at the canonicalize() layer so explicit
368        // producer-supplied hints get the same treatment.
369        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        // Backward compat: pre-M1.4 envelopes (endpoint_id = None) must
385        // produce the same session_hint bytes, or existing Redis sessions
386        // orphan. Both derived and explicit hints stay untouched.
387        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        // Codex M1.4b-iii regression: producer-supplied hints must still
417        // partition by endpoint id, not collapse to the same session key.
418        // Two endpoints with IDENTICAL producer-supplied session_hints must
419        // resolve to distinct effective session keys.
420        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        // Re-canonicalizing must not double-prefix the namespace marker.
446        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        // Defensive depth: an eid that bypassed the http-layer validator
461        // (embedded caller, future WIT-derived `identify_instance`) must
462        // not corrupt the `ep=<eid>::<base>` namespace prefix or the
463        // telemetry attribute. Drop to None ⇒ run unscoped.
464        let cases = [
465            ("empty", ""),
466            ("colon embedded", "teams:legal"), // collides prefix delimiter
467            ("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            // Session hint must be the un-prefixed form when eid was dropped.
481            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        // The validator must accept both the M1.2 ULID form and a
492        // hand-typeable slug, and the dot variant used in display ids.
493        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    /// Construct a runtime from explicit flow definitions (legacy entrypoint used by tests/examples).
539    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    /// Build a state-machine runtime that proxies pack flows through the legacy FlowEngine.
565    #[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    /// Execute the flow associated with the provided ingress event.
615    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    /// NATS client backing the per-node audit-event sink; `None` (the
720    /// default, off path) when `GREENTIC_EVENTS_NATS_URL` is unset or NATS
721    /// could not be reached — see `docs/superpowers/specs/2026-07-03-runner-audit-emitter-design.md`.
722    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            // Build the best-effort audit sink from the threaded NATS client
826            // (EPIC-B B-2). `audit_tenant` is derived from the same
827            // condition as the sink so the two stay coupled: both `Some`
828            // (audit enabled) or both `None` (default, file-trace-only path,
829            // zero behaviour change). `TraceRecorder` construction happens
830            // on the async execution path, so `AuditSink::new`'s
831            // `tokio::spawn` runs inside a live tokio runtime context.
832            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            // Carry the inbound reply scope so async-dispatch nodes can encode
853            // the originating thread/reply_to into the dispatch correlation id
854            // (so a threaded wait can be re-keyed on resume).
855            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    /// Multi-instance messaging endpoint discriminator (M1.4).
938    /// Distinguishes provider instances of the same `provider_type` —
939    /// e.g. `teams-legal` vs `teams-accounting` — so sessions and traces
940    /// don't collide across endpoints that share provider/channel/user.
941    #[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    /// Flow node this run must START at, lifted from a card-navigation target
950    /// that names a node rather than a card asset — see
951    /// [`crate::runner::card_nav`]. `None` starts at the flow's entrypoint.
952    ///
953    /// A parked snapshot still outranks it: resuming a waiting run is never
954    /// the same thing as re-entering the flow at a node.
955    #[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
969/// Validate a producer-asserted messaging endpoint id at the runner-host
970/// boundary. Mirrors the http-layer validator in
971/// `greentic-start::revision_serve::validate_endpoint_id` so the
972/// canonicalize-layer `ep=<eid>::<base>` session prefix is safe regardless
973/// of how the eid reached the envelope — http header (already validated),
974/// embedded callers (`start_embedded_host`), or the upcoming WIT-derived
975/// `identify_instance` path. Invalid → drop to None (run unscoped instead
976/// of corrupting the namespace prefix).
977///
978/// The threat the grammar defends against: a value containing `:` collides
979/// the prefix delimiter (`eid="a"+base="b::c"` and `eid="a::b"+base="c"`
980/// both produce `ep=a::b::c`); empty/whitespace-only values collapse all
981/// malformed traffic into one namespace; control characters or unbounded
982/// length corrupt downstream session-store keys and telemetry attribute
983/// values.
984fn 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        // Defensive: drop an invalid eid BEFORE it reaches the prefix step
1013        // or telemetry stamping. See `endpoint_id_is_valid` for why.
1014        //
1015        // The drop is silent at the http boundary (greentic-start validates
1016        // and drops there too — see Codex review note: that path is fronted
1017        // by request validation, so an invalid eid means the operator never
1018        // asserted one). Embedded callers and the future WIT-derived
1019        // `identify_instance` path have no such validator, so a downgrade
1020        // here usually signals a producer bug. Emit at WARN so operators can
1021        // spot it, then continue with the same drop-to-None semantics the
1022        // http layer uses — staying consistent across boundaries.
1023        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        // Endpoint isolation: prefix the hint (explicit OR derived) with
1036        // `ep=<eid>::` so two endpoints reusing the same producer-supplied
1037        // session key never collide. Idempotent — re-canonicalize is a no-op.
1038        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        // Pre-M1.4 structured form. Endpoint isolation is applied uniformly
1068        // at the `canonicalize()` layer so explicit producer-supplied hints
1069        // get the same namespacing as derived ones — keeping this function
1070        // pure and bytes-identical to pre-M1.4 for `endpoint_id = None`.
1071        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        // M1.4: ride the attributes map so the runner-host local
1096        // `tenant_ctx_to_telemetry` projection copies it onto
1097        // `TelemetryCtx::messaging_endpoint_id` for OTel export.
1098        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}