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            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        // canonical_session_hint() does NOT embed messaging_endpoint_id —
364        // namespacing happens at the canonicalize() layer so explicit
365        // producer-supplied hints get the same treatment.
366        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        // Backward compat: pre-M1.4 envelopes (endpoint_id = None) must
382        // produce the same session_hint bytes, or existing Redis sessions
383        // orphan. Both derived and explicit hints stay untouched.
384        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        // Codex M1.4b-iii regression: producer-supplied hints must still
414        // partition by endpoint id, not collapse to the same session key.
415        // Two endpoints with IDENTICAL producer-supplied session_hints must
416        // resolve to distinct effective session keys.
417        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        // Re-canonicalizing must not double-prefix the namespace marker.
443        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        // Defensive depth: an eid that bypassed the http-layer validator
458        // (embedded caller, future WIT-derived `identify_instance`) must
459        // not corrupt the `ep=<eid>::<base>` namespace prefix or the
460        // telemetry attribute. Drop to None ⇒ run unscoped.
461        let cases = [
462            ("empty", ""),
463            ("colon embedded", "teams:legal"), // collides prefix delimiter
464            ("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            // Session hint must be the un-prefixed form when eid was dropped.
478            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        // The validator must accept both the M1.2 ULID form and a
489        // hand-typeable slug, and the dot variant used in display ids.
490        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    /// Construct a runtime from explicit flow definitions (legacy entrypoint used by tests/examples).
536    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    /// Build a state-machine runtime that proxies pack flows through the legacy FlowEngine.
562    #[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    /// Execute the flow associated with the provided ingress event.
612    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    /// NATS client backing the per-node audit-event sink; `None` (the
717    /// default, off path) when `GREENTIC_EVENTS_NATS_URL` is unset or NATS
718    /// could not be reached — see `docs/superpowers/specs/2026-07-03-runner-audit-emitter-design.md`.
719    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            // Build the best-effort audit sink from the threaded NATS client
823            // (EPIC-B B-2). `audit_tenant` is derived from the same
824            // condition as the sink so the two stay coupled: both `Some`
825            // (audit enabled) or both `None` (default, file-trace-only path,
826            // zero behaviour change). `TraceRecorder` construction happens
827            // on the async execution path, so `AuditSink::new`'s
828            // `tokio::spawn` runs inside a live tokio runtime context.
829            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            // Carry the inbound reply scope so async-dispatch nodes can encode
850            // the originating thread/reply_to into the dispatch correlation id
851            // (so a threaded wait can be re-keyed on resume).
852            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    /// Multi-instance messaging endpoint discriminator (M1.4).
933    /// Distinguishes provider instances of the same `provider_type` —
934    /// e.g. `teams-legal` vs `teams-accounting` — so sessions and traces
935    /// don't collide across endpoints that share provider/channel/user.
936    #[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
956/// Validate a producer-asserted messaging endpoint id at the runner-host
957/// boundary. Mirrors the http-layer validator in
958/// `greentic-start::revision_serve::validate_endpoint_id` so the
959/// canonicalize-layer `ep=<eid>::<base>` session prefix is safe regardless
960/// of how the eid reached the envelope — http header (already validated),
961/// embedded callers (`start_embedded_host`), or the upcoming WIT-derived
962/// `identify_instance` path. Invalid → drop to None (run unscoped instead
963/// of corrupting the namespace prefix).
964///
965/// The threat the grammar defends against: a value containing `:` collides
966/// the prefix delimiter (`eid="a"+base="b::c"` and `eid="a::b"+base="c"`
967/// both produce `ep=a::b::c`); empty/whitespace-only values collapse all
968/// malformed traffic into one namespace; control characters or unbounded
969/// length corrupt downstream session-store keys and telemetry attribute
970/// values.
971fn 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        // Defensive: drop an invalid eid BEFORE it reaches the prefix step
1000        // or telemetry stamping. See `endpoint_id_is_valid` for why.
1001        //
1002        // The drop is silent at the http boundary (greentic-start validates
1003        // and drops there too — see Codex review note: that path is fronted
1004        // by request validation, so an invalid eid means the operator never
1005        // asserted one). Embedded callers and the future WIT-derived
1006        // `identify_instance` path have no such validator, so a downgrade
1007        // here usually signals a producer bug. Emit at WARN so operators can
1008        // spot it, then continue with the same drop-to-None semantics the
1009        // http layer uses — staying consistent across boundaries.
1010        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        // Endpoint isolation: prefix the hint (explicit OR derived) with
1023        // `ep=<eid>::` so two endpoints reusing the same producer-supplied
1024        // session key never collide. Idempotent — re-canonicalize is a no-op.
1025        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        // Pre-M1.4 structured form. Endpoint isolation is applied uniformly
1055        // at the `canonicalize()` layer so explicit producer-supplied hints
1056        // get the same namespacing as derived ones — keeping this function
1057        // pure and bytes-identical to pre-M1.4 for `endpoint_id = None`.
1058        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        // M1.4: ride the attributes map so the runner-host local
1083        // `tenant_ctx_to_telemetry` projection copies it onto
1084        // `TelemetryCtx::messaging_endpoint_id` for OTel export.
1085        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}