use std::collections::HashMap;
use std::path::Path;
use std::sync::Arc;
use anyhow::{Context, Result, anyhow, bail};
use serde_json::Value;
use crate::activity::{Activity, WelcomeFlowHint};
use crate::boot;
use crate::component_api::node::{ExecCtx as ComponentExecCtx, TenantCtx as ComponentTenantCtx};
use crate::config::{Fast2FlowRoutingConfig, HostConfig};
use crate::engine::host::{SessionHost, StateHost};
use crate::engine::runtime::{FlowResumeStore, IngressEnvelope};
#[cfg(feature = "greentic-x-provider")]
use crate::greentic_x_provider::RunnerPackFast2FlowRoutingProvider;
use crate::http::health::HealthState;
use crate::pack::{IdentifyOutcome, PackRuntime};
use crate::runner::adapt_timer;
use crate::runner::engine::FlowEngine;
use crate::runtime::{ActivePacks, TenantRuntime};
use crate::secrets::{DynSecretsManager, default_manager};
use crate::storage::{
DynSessionStore, DynStateStore, new_session_store, new_state_store, session_host_from,
state_host_from,
};
use crate::wasi::RunnerWasiPolicy;
use greentic_deploy_spec::ids::{BundleId, DeploymentId, RevisionId};
#[cfg(feature = "greentic-x-provider")]
use greentic_x_runtime::{
Fast2FlowDirective, Fast2FlowMessageEnvelope, Fast2FlowRouteRequest, Fast2FlowRoutingProvider,
};
#[derive(Clone, Debug)]
pub struct TelemetryCfg {
pub config: greentic_telemetry::TelemetryConfig,
pub export: greentic_telemetry::export::ExportConfig,
}
pub struct HostBuilder {
configs: HashMap<String, HostConfig>,
telemetry: Option<TelemetryCfg>,
wasi_policy: RunnerWasiPolicy,
secrets: Option<DynSecretsManager>,
}
impl HostBuilder {
pub fn new() -> Self {
Self {
configs: HashMap::new(),
telemetry: None,
wasi_policy: RunnerWasiPolicy::default(),
secrets: None,
}
}
pub fn with_config(mut self, config: HostConfig) -> Self {
self.configs.insert(config.tenant.clone(), config);
self
}
pub fn with_telemetry(mut self, telemetry: TelemetryCfg) -> Self {
self.telemetry = Some(telemetry);
self
}
pub fn with_wasi_policy(mut self, policy: RunnerWasiPolicy) -> Self {
self.wasi_policy = policy;
self
}
pub fn with_secrets_manager(mut self, manager: DynSecretsManager) -> Self {
self.secrets = Some(manager);
self
}
pub fn build(self) -> Result<RunnerHost> {
if self.configs.is_empty() {
bail!("at least one tenant configuration is required");
}
let wasi_policy = Arc::new(self.wasi_policy);
let configs = self
.configs
.into_iter()
.map(|(tenant, cfg)| (tenant, Arc::new(cfg)))
.collect();
let session_store = new_session_store();
let session_host = session_host_from(Arc::clone(&session_store));
let state_store = new_state_store();
let state_host = state_host_from(Arc::clone(&state_store));
let secrets = match self.secrets {
Some(manager) => manager,
None => default_manager().context("failed to initialise default secrets backend")?,
};
Ok(RunnerHost {
configs,
active: Arc::new(ActivePacks::new()),
health: Arc::new(HealthState::new()),
session_store,
state_store,
session_host,
state_host,
wasi_policy,
secrets_manager: secrets,
telemetry: self.telemetry,
})
}
}
impl Default for HostBuilder {
fn default() -> Self {
Self::new()
}
}
pub struct RunnerHost {
configs: HashMap<String, Arc<HostConfig>>,
active: Arc<ActivePacks>,
health: Arc<HealthState>,
session_store: DynSessionStore,
state_store: DynStateStore,
session_host: Arc<dyn SessionHost>,
state_host: Arc<dyn StateHost>,
wasi_policy: Arc<RunnerWasiPolicy>,
secrets_manager: DynSecretsManager,
telemetry: Option<TelemetryCfg>,
}
#[derive(Clone)]
pub struct TenantHandle {
runtime: Arc<TenantRuntime>,
}
impl RunnerHost {
pub async fn start(&self) -> Result<()> {
boot::init(&self.health, self.telemetry.as_ref())?;
Ok(())
}
pub async fn stop(&self) -> Result<()> {
self.active.replace(HashMap::new());
Ok(())
}
pub async fn load_pack(&self, tenant: &str, pack_path: &Path) -> Result<()> {
let archive_source = if is_pack_archive(pack_path) {
Some(pack_path)
} else {
None
};
let runtime = self
.prepare_runtime(tenant, pack_path, archive_source)
.await
.with_context(|| format!("failed to load tenant {tenant}"))?;
self.active.insert_pack(tenant, runtime);
tracing::info!(tenant, pack = %pack_path.display(), "pack loaded");
Ok(())
}
pub async fn handle_activity(&self, tenant: &str, activity: Activity) -> Result<Vec<Activity>> {
let runtime = self
.active
.load_pack(tenant)
.with_context(|| format!("tenant {tenant} not loaded"))?;
self.dispatch_activity(&runtime, tenant, activity).await
}
pub async fn handle_activity_for_revision(
&self,
tenant: &str,
deployment_id: DeploymentId,
bundle_id: BundleId,
revision_id: RevisionId,
activity: Activity,
) -> Result<Vec<Activity>> {
let runtime = self
.active
.load_revision(tenant, deployment_id, bundle_id, revision_id)
.with_context(|| {
format!(
"revision runtime not loaded for tenant {tenant} \
(deployment {deployment_id}, revision {revision_id})"
)
})?;
self.dispatch_activity(&runtime, tenant, activity).await
}
fn load_revision_runtime(
&self,
tenant: &str,
deployment_id: DeploymentId,
bundle_id: BundleId,
revision_id: RevisionId,
) -> Result<Arc<crate::runtime::TenantRuntime>> {
self.active
.load_revision(tenant, deployment_id, bundle_id, revision_id)
.with_context(|| {
format!(
"revision runtime not loaded for tenant {tenant} \
(deployment {deployment_id}, revision {revision_id})"
)
})
}
pub async fn identify_messaging_endpoints_for_revision(
&self,
tenant: &str,
deployment_id: DeploymentId,
bundle_id: BundleId,
revision_id: RevisionId,
provider_types: &[&str],
payload: &[u8],
) -> Result<HashMap<String, IdentifyOutcome>> {
if provider_types.is_empty() {
return Ok(HashMap::new());
}
let runtime = self.load_revision_runtime(tenant, deployment_id, bundle_id, revision_id)?;
let mut merged: HashMap<String, IdentifyOutcome> = provider_types
.iter()
.map(|ty| ((*ty).to_string(), IdentifyOutcome::Unsupported))
.collect();
for pack in runtime.all_packs() {
let remaining: Vec<&str> = provider_types
.iter()
.copied()
.filter(|ty| !matches!(merged.get(*ty), Some(IdentifyOutcome::Identified(_))))
.collect();
if remaining.is_empty() {
break;
}
let probe = pack
.identify_endpoints_by_provider_type(&remaining, payload)
.await?;
for (ty, outcome) in probe {
if let Some(existing) = merged.get_mut(&ty) {
existing.merge_in(outcome);
}
}
}
Ok(merged)
}
#[allow(clippy::too_many_arguments)]
pub async fn identify_messaging_endpoints_for_revision_scoped(
&self,
tenant: &str,
deployment_id: DeploymentId,
bundle_id: BundleId,
revision_id: RevisionId,
provider_types: &[&str],
headers: &[(String, String)],
body: &Value,
) -> Result<HashMap<String, IdentifyOutcome>> {
if provider_types.is_empty() {
return Ok(HashMap::new());
}
let runtime = self.load_revision_runtime(tenant, deployment_id, bundle_id, revision_id)?;
let mut merged: HashMap<String, IdentifyOutcome> = provider_types
.iter()
.map(|ty| ((*ty).to_string(), IdentifyOutcome::Unsupported))
.collect();
for pack in runtime.all_packs() {
let remaining: Vec<&str> = provider_types
.iter()
.copied()
.filter(|ty| !matches!(merged.get(*ty), Some(IdentifyOutcome::Identified(_))))
.collect();
if remaining.is_empty() {
break;
}
let probe = pack
.identify_endpoints_by_provider_type_scoped(&remaining, headers, body)
.await?;
for (ty, outcome) in probe {
if let Some(existing) = merged.get_mut(&ty) {
existing.merge_in(outcome);
}
}
}
Ok(merged)
}
pub async fn describe_identify_instances_for_revision(
&self,
tenant: &str,
deployment_id: DeploymentId,
bundle_id: BundleId,
revision_id: RevisionId,
provider_types: &[&str],
) -> Result<HashMap<String, Option<crate::identify_hint::IdentifyInstanceHint>>> {
if provider_types.is_empty() {
return Ok(HashMap::new());
}
let runtime = self.load_revision_runtime(tenant, deployment_id, bundle_id, revision_id)?;
let mut merged: HashMap<String, Option<crate::identify_hint::IdentifyInstanceHint>> =
provider_types
.iter()
.map(|ty| ((*ty).to_string(), None))
.collect();
for pack in runtime.all_packs() {
let remaining: Vec<&str> = provider_types
.iter()
.copied()
.filter(|ty| !matches!(merged.get(*ty), Some(Some(_))))
.collect();
if remaining.is_empty() {
break;
}
let probe = pack
.describe_identify_hints_by_provider_type(&remaining)
.await?;
for (ty, hint) in probe {
if let Some(slot) = merged.get_mut(&ty)
&& slot.is_none()
{
*slot = hint;
}
}
}
Ok(merged)
}
#[allow(clippy::too_many_arguments)]
pub async fn invoke_provider_for_revision(
&self,
tenant: &str,
deployment_id: DeploymentId,
bundle_id: BundleId,
revision_id: RevisionId,
provider_type: &str,
op: &str,
input_json: Vec<u8>,
correlation_id: Option<String>,
trace_id: Option<String>,
) -> Result<Value> {
let runtime = self.load_revision_runtime(tenant, deployment_id, bundle_id, revision_id)?;
let mut matched = None;
for pack in runtime.all_packs() {
let Some(registry) = pack.provider_registry_optional()? else {
continue;
};
let Some((binding, declared_ops)) =
registry.try_resolve_with_ops(None, Some(provider_type))?
else {
continue;
};
if matched.is_some() {
bail!(
"ambiguous provider_type `{provider_type}` in revision \
(deployment {deployment_id}, revision {revision_id}): \
multiple packs bind the same type; pack manifests must \
declare each provider_type at most once across main + overlays"
);
}
matched = Some((Arc::clone(pack), binding, declared_ops));
}
let Some((pack, binding, declared_ops)) = matched else {
bail!(
"no pack in revision binds provider_type `{provider_type}` \
(deployment {deployment_id}, revision {revision_id})"
);
};
if !declared_ops.iter().any(|d| d == op) {
bail!(
"op `{op}` is not declared for provider_type `{provider_type}` \
in revision (deployment {deployment_id}, revision {revision_id}); \
declared ops: {declared_ops:?}"
);
}
let exec_ctx = ComponentExecCtx {
tenant: ComponentTenantCtx {
tenant: tenant.to_string(),
team: None,
user: None,
trace_id,
i18n_id: None,
correlation_id: correlation_id.clone(),
deadline_unix_ms: None,
attempt: 1,
idempotency_key: correlation_id,
},
i18n_id: None,
flow_id: format!("provider-webhook/{provider_type}"),
node_id: None,
};
pack.invoke_provider(&binding, exec_ctx, op, input_json)
.await
}
async fn dispatch_activity(
&self,
runtime: &TenantRuntime,
tenant: &str,
activity: Activity,
) -> Result<Vec<Activity>> {
let activity = apply_fast2flow_routing(runtime, tenant, activity)?;
if activity.action() == Some("response") && activity.flow_id().is_none() {
return Ok(vec![activity]);
}
let (pack_id, flow_id) = resolve_flow_id(runtime, &activity)?;
let action = activity.action().map(|value| value.to_string());
let session = activity.session_id().map(|value| value.to_string());
let provider = activity.provider_id().map(|value| value.to_string());
let messaging_endpoint_id = activity
.messaging_endpoint_id()
.map(|value| value.to_string());
let channel = activity.channel().map(|value| value.to_string());
let conversation = activity.conversation().map(|value| value.to_string());
let user = activity.user().map(|value| value.to_string());
let welcome_flow_hint = activity.welcome_flow_hint().cloned();
let resolved_flow_type =
activity
.flow_type()
.map(|value| value.to_string())
.or_else(|| {
runtime
.engine()
.flow_by_key(&pack_id, &flow_id)
.map(|desc| desc.flow_type.clone())
});
let payload = activity.into_payload();
let mut envelope = IngressEnvelope {
tenant: tenant.to_string(),
env: std::env::var("GREENTIC_ENV").ok(),
pack_id: Some(pack_id.clone()),
flow_id: flow_id.clone(),
flow_type: resolved_flow_type,
action,
session_hint: session,
provider,
messaging_endpoint_id,
channel,
conversation,
user,
activity_id: None,
timestamp: None,
payload,
metadata: None,
reply_scope: None,
}
.canonicalize();
let hint_flow_type = welcome_flow_hint.as_ref().and_then(|hint| {
runtime
.engine()
.flow_by_key(&hint.pack_id, &hint.flow_id)
.map(|desc| desc.flow_type.clone())
});
apply_welcome_flow_override(
runtime.session_store(),
&mut envelope,
welcome_flow_hint.as_ref(),
hint_flow_type,
)?;
let result = runtime.state_machine().handle(envelope).await?;
Ok(normalize_replies(result, tenant))
}
pub async fn tenant(&self, tenant: &str) -> Option<TenantHandle> {
self.active
.load_pack(tenant)
.map(|runtime| TenantHandle { runtime })
}
pub fn active_packs(&self) -> Arc<ActivePacks> {
Arc::clone(&self.active)
}
pub fn health_state(&self) -> Arc<HealthState> {
Arc::clone(&self.health)
}
pub fn wasi_policy(&self) -> Arc<RunnerWasiPolicy> {
Arc::clone(&self.wasi_policy)
}
pub fn session_store(&self) -> DynSessionStore {
Arc::clone(&self.session_store)
}
pub fn state_store(&self) -> DynStateStore {
Arc::clone(&self.state_store)
}
pub fn session_host(&self) -> Arc<dyn SessionHost> {
Arc::clone(&self.session_host)
}
pub fn state_host(&self) -> Arc<dyn StateHost> {
Arc::clone(&self.state_host)
}
pub fn secrets_manager(&self) -> DynSecretsManager {
Arc::clone(&self.secrets_manager)
}
pub fn tenant_configs(&self) -> HashMap<String, Arc<HostConfig>> {
self.configs.clone()
}
#[cfg(test)]
pub(crate) fn for_test() -> Arc<Self> {
use crate::config::{
FlowRetryConfig, OperatorPolicy, RateLimits, SecretsPolicy, StateStorePolicy,
WebhookPolicy,
};
let config = Arc::new(HostConfig {
tenant: "test".into(),
bindings_path: std::path::PathBuf::from("<test>"),
flow_type_bindings: std::collections::HashMap::new(),
rate_limits: RateLimits::default(),
retry: FlowRetryConfig::default(),
http_enabled: false,
secrets_policy: SecretsPolicy::allow_all(),
state_store_policy: StateStorePolicy::default(),
webhook_policy: WebhookPolicy::default(),
timers: Vec::new(),
oauth: None,
mocks: None,
pack_bindings: Vec::new(),
env_passthrough: Vec::new(),
trace: crate::trace::TraceConfig::from_env(),
validation: crate::validate::ValidationConfig::from_env(),
operator_policy: OperatorPolicy::allow_all(),
fast2flow: Default::default(),
#[cfg(feature = "agentic-worker")]
agents: std::collections::HashMap::new(),
#[cfg(feature = "agentic-worker")]
graphs: std::collections::HashMap::new(),
});
let session_store = new_session_store();
let session_host = session_host_from(Arc::clone(&session_store));
let state_store = new_state_store();
let state_host = state_host_from(Arc::clone(&state_store));
let secrets_manager = default_manager().expect("test secrets manager");
Arc::new(Self {
configs: std::collections::HashMap::from([("test".to_string(), config)]),
active: Arc::new(ActivePacks::new()),
health: Arc::new(HealthState::new()),
session_store,
state_store,
session_host,
state_host,
wasi_policy: Arc::new(RunnerWasiPolicy::default()),
secrets_manager,
telemetry: None,
})
}
async fn prepare_runtime(
&self,
tenant: &str,
pack_path: &Path,
archive_source: Option<&Path>,
) -> Result<Arc<TenantRuntime>> {
let config = self
.configs
.get(tenant)
.cloned()
.with_context(|| format!("tenant {tenant} not registered"))?;
if config.tenant != tenant {
bail!(
"tenant mismatch: config declares '{}' but '{tenant}' was requested",
config.tenant
);
}
let runtime = TenantRuntime::load(
pack_path,
Arc::clone(&config),
None,
archive_source,
None,
self.wasi_policy(),
self.session_host(),
self.session_store(),
self.state_store(),
self.state_host(),
self.secrets_manager(),
)
.await?;
let timers = adapt_timer::spawn_timers(Arc::clone(&runtime))?;
runtime.register_timers(timers);
Ok(runtime)
}
}
impl TenantHandle {
pub fn config(&self) -> Arc<HostConfig> {
Arc::clone(self.runtime.config())
}
pub fn pack(&self) -> Arc<PackRuntime> {
self.runtime.pack()
}
pub fn engine(&self) -> Arc<FlowEngine> {
Arc::clone(self.runtime.engine())
}
pub fn overlays(&self) -> Vec<Arc<PackRuntime>> {
self.runtime.overlays()
}
pub fn overlay_digests(&self) -> Vec<Option<String>> {
self.runtime.overlay_digests()
}
}
fn apply_welcome_flow_override(
session_store: &DynSessionStore,
envelope: &mut IngressEnvelope,
hint: Option<&WelcomeFlowHint>,
hint_flow_type: Option<String>,
) -> Result<()> {
let Some(hint) = hint else {
return Ok(());
};
if envelope.messaging_endpoint_id.is_none() {
return Ok(());
}
if !try_mark_welcome_first_contact(session_store, envelope)? {
return Ok(());
}
let resume = FlowResumeStore::new(Arc::clone(session_store));
let snapshot = resume
.fetch(envelope)
.map_err(|err| anyhow!("welcome-flow first-contact probe failed: {err}"))?;
if snapshot.is_some() {
return Ok(());
}
envelope.pack_id = Some(hint.pack_id.clone());
envelope.flow_id = hint.flow_id.clone();
envelope.flow_type = hint_flow_type;
Ok(())
}
fn try_mark_welcome_first_contact(
store: &DynSessionStore,
envelope: &IngressEnvelope,
) -> Result<bool> {
let Some(scope) = welcome_marker_scope(envelope) else {
return Ok(false);
};
let (ctx, user) = FlowResumeStore::contact_identity(envelope)
.map_err(|e| anyhow!("welcome marker identity probe failed: {e}"))?;
if store
.find_wait_by_scope(&ctx, &user, &scope)
.map_err(|e| anyhow!("welcome marker probe failed: {e}"))?
.is_some()
{
return Ok(false);
}
let data = marker_session_data(&ctx, &user);
let session_key = marker_session_key(&ctx, &user, &scope);
store
.register_wait(&ctx, &user, &scope, &session_key, data, None)
.map_err(|e| anyhow!("welcome marker register failed: {e}"))?;
Ok(true)
}
fn marker_session_key(
ctx: &greentic_types::TenantCtx,
user: &greentic_types::UserId,
scope: &greentic_types::ReplyScope,
) -> greentic_session::SessionKey {
use sha2::{Digest, Sha256};
let team = match ctx.team_id.as_ref().or(ctx.team.as_ref()) {
Some(t) => t.as_str(),
None => "<none>",
};
let digest = Sha256::digest(
format!(
"welcome-marker:v1\0{}\0{}\0team={team}\0{}\0{}",
ctx.env.as_str(),
ctx.tenant_id.as_str(),
user.as_str(),
scope.conversation,
)
.as_bytes(),
);
greentic_session::SessionKey::new(format!("welcome-marker::{}", hex::encode(digest)))
}
fn welcome_marker_scope(envelope: &IngressEnvelope) -> Option<greentic_types::ReplyScope> {
let eid = envelope.messaging_endpoint_id.as_deref()?;
Some(greentic_types::ReplyScope {
conversation: format!("welcome-seen::ep={eid}"),
thread: None,
reply_to: None,
correlation: None,
})
}
fn marker_session_data(
ctx: &greentic_types::TenantCtx,
user: &greentic_types::UserId,
) -> greentic_session::SessionData {
use std::str::FromStr;
use std::sync::LazyLock;
static FLOW_ID: LazyLock<greentic_types::FlowId> =
LazyLock::new(|| greentic_types::FlowId::from_str("welcome-marker").expect("valid id"));
static PACK_ID: LazyLock<greentic_types::PackId> =
LazyLock::new(|| greentic_types::PackId::from_str("welcome-marker").expect("valid id"));
let cursor = greentic_types::SessionCursor::new("marker".to_string());
let ctx = ctx.clone().with_user(Some(user.clone()));
greentic_session::SessionData {
tenant_ctx: ctx,
flow_id: FLOW_ID.clone(),
pack_id: Some(PACK_ID.clone()),
cursor,
context_json: "{}".to_string(),
}
}
fn apply_fast2flow_routing(
runtime: &TenantRuntime,
tenant: &str,
activity: Activity,
) -> Result<Activity> {
let config = &runtime.config().fast2flow;
if !config.enabled || activity.flow_id().is_some() {
return Ok(activity);
}
apply_fast2flow_routing_enabled(runtime, tenant, activity, config)
}
#[cfg(feature = "greentic-x-provider")]
fn apply_fast2flow_routing_enabled(
runtime: &TenantRuntime,
tenant: &str,
activity: Activity,
config: &Fast2FlowRoutingConfig,
) -> Result<Activity> {
let Some(text) = activity.payload().get("text").and_then(Value::as_str) else {
return Ok(activity);
};
if text.trim().is_empty() {
return Ok(activity);
}
let mut envelope = Fast2FlowMessageEnvelope::new(text.trim().to_owned());
if let Some(channel) = activity.channel() {
envelope = envelope.with_channel(channel.to_owned());
}
if let Some(provider) = activity.provider_id() {
envelope = envelope.with_provider(provider.to_owned());
}
let request = Fast2FlowRouteRequest {
scope: config.scope.clone().unwrap_or_else(|| tenant.to_owned()),
envelope,
session_active: activity.session_id().is_some(),
input_locale: "en".to_owned(),
time_budget_ms: config.time_budget_ms,
registry_path: config.registry_path.clone(),
indexes_path: config.indexes_path.clone(),
now_unix_ms: chrono::Utc::now().timestamp_millis().max(0) as u64,
metadata: Default::default(),
};
let provider = RunnerPackFast2FlowRoutingProvider::new(runtime.pack())
.map_err(|err| anyhow!(err.to_string()))?
.with_component_ref(config.component_ref.clone())
.with_operation(config.operation.clone())
.with_tenant(tenant.to_owned());
let route = provider
.route_intent(request)
.map_err(|err| anyhow!(err.to_string()))?;
match route.directive {
Fast2FlowDirective::Continue => Ok(activity),
Fast2FlowDirective::Dispatch {
target, entities, ..
} => apply_fast2flow_target(activity, &target, entities),
Fast2FlowDirective::Respond { message } => Ok(Activity::custom(
"response",
serde_json::json!({ "messages": [{ "text": message }] }),
)
.ensure_tenant(tenant)),
Fast2FlowDirective::Deny { reason } => Ok(Activity::custom(
"response",
serde_json::json!({ "messages": [{ "text": reason }] }),
)
.ensure_tenant(tenant)),
}
}
#[cfg(not(feature = "greentic-x-provider"))]
fn apply_fast2flow_routing_enabled(
_runtime: &TenantRuntime,
_tenant: &str,
_activity: Activity,
_config: &Fast2FlowRoutingConfig,
) -> Result<Activity> {
bail!("fast2flow routing requires the greentic-x-provider feature")
}
#[cfg(feature = "greentic-x-provider")]
fn apply_fast2flow_target(
activity: Activity,
target: &str,
entities: Vec<greentic_x_runtime::Fast2FlowRoutingEntity>,
) -> Result<Activity> {
let target = target.trim();
if target.is_empty() {
bail!("fast2flow dispatch target is empty");
}
if let Some((pack_id, flow_id)) = target.split_once('/') {
if pack_id.trim().is_empty() || flow_id.trim().is_empty() {
bail!("fast2flow dispatch target `{target}` must be `pack_id/flow_id` or `flow_id`");
}
return Ok(attach_fast2flow_entities(
activity.with_pack(pack_id.trim()).with_flow(flow_id.trim()),
entities,
));
}
Ok(attach_fast2flow_entities(
activity.with_flow(target),
entities,
))
}
#[cfg(feature = "greentic-x-provider")]
fn attach_fast2flow_entities(
activity: Activity,
entities: Vec<greentic_x_runtime::Fast2FlowRoutingEntity>,
) -> Activity {
if entities.is_empty() {
return activity;
}
activity.with_payload_field(
"fast2flow",
serde_json::json!({
"entities": entities,
}),
)
}
fn resolve_flow_id(runtime: &TenantRuntime, activity: &Activity) -> Result<(String, String)> {
let engine = runtime.engine();
if let Some(flow_id) = activity.flow_id() {
if let Some(pack_id) = activity.pack_id() {
if engine.flow_by_key(pack_id, flow_id).is_none() {
bail!("flow {flow_id} not registered for pack {pack_id}");
}
return Ok((pack_id.to_string(), flow_id.to_string()));
}
if let Some(flow) = engine.flow_by_id(flow_id) {
return Ok((flow.pack_id.clone(), flow.id.clone()));
}
bail!("flow {flow_id} is ambiguous; pack_id is required");
}
if let Some(flow_type) = activity.flow_type() {
if let Some(pack_id) = activity.pack_id() {
if let Some(flow) = engine
.flows()
.iter()
.find(|flow| flow.pack_id == pack_id && flow.flow_type == flow_type)
{
return Ok((pack_id.to_string(), flow.id.clone()));
}
bail!("flow type {flow_type} not registered for pack {pack_id}");
}
if let Some(flow) = engine.flow_by_type(flow_type) {
return Ok((flow.pack_id.clone(), flow.id.clone()));
}
bail!("flow type {flow_type} is ambiguous; pack_id is required");
}
let pack = runtime.pack();
let flow_id = pack
.metadata()
.entry_flows
.first()
.cloned()
.ok_or_else(|| anyhow!("no entry flows registered for tenant {}", runtime.tenant()))?;
Ok((pack.metadata().pack_id.clone(), flow_id))
}
fn normalize_replies(result: Value, tenant: &str) -> Vec<Activity> {
result
.as_array()
.cloned()
.unwrap_or_else(|| vec![result])
.into_iter()
.map(|payload| Activity::from_output(payload, tenant))
.collect()
}
fn is_pack_archive(path: &Path) -> bool {
path.extension()
.and_then(|ext| ext.to_str())
.map(|ext| ext.eq_ignore_ascii_case("gtpack"))
.unwrap_or(false)
}
#[cfg(test)]
mod welcome_flow_tests {
use super::*;
use crate::engine::runtime::IngressEnvelope;
use crate::runner::engine::{ExecutionState, FlowSnapshot, FlowWait};
use crate::storage::new_session_store;
use greentic_types::ReplyScope;
use serde_json::json;
fn sample_envelope(endpoint_id: Option<&str>) -> IngressEnvelope {
sample_envelope_for_user(endpoint_id, "user-1")
}
fn sample_envelope_for_user(endpoint_id: Option<&str>, user: &str) -> IngressEnvelope {
IngressEnvelope {
tenant: "demo".into(),
env: Some("local".into()),
pack_id: Some("pack.default".into()),
flow_id: "flow.default".into(),
flow_type: Some("messaging".into()),
action: Some("messaging".into()),
session_hint: None,
provider: Some("teams".into()),
messaging_endpoint_id: endpoint_id.map(String::from),
channel: Some("chan".into()),
conversation: Some(format!("conv-{user}")),
user: Some(user.to_string()),
activity_id: None,
timestamp: None,
payload: json!({}),
metadata: None,
reply_scope: Some(ReplyScope {
conversation: format!("conv-{user}"),
thread: None,
reply_to: None,
correlation: None,
}),
}
.canonicalize()
}
fn hint() -> WelcomeFlowHint {
WelcomeFlowHint {
pack_id: "pack.welcome".into(),
flow_id: "flow.welcome".into(),
}
}
fn seed_resume(store: &DynSessionStore, envelope: &IngressEnvelope) {
let resume = FlowResumeStore::new(Arc::clone(store));
let state: ExecutionState = serde_json::from_value(json!({
"input": { "text": "hi" },
"nodes": {},
"egress": []
}))
.expect("state");
let wait = FlowWait {
reason: Some("await-user".into()),
snapshot: FlowSnapshot {
pack_id: envelope.pack_id.clone().expect("pack_id"),
flow_id: envelope.flow_id.clone(),
next_flow: None,
next_node: "node-2".into(),
state,
},
};
resume.save(envelope, &wait).expect("seed save");
}
#[test]
fn override_is_no_op_when_hint_absent() {
let store = new_session_store();
let mut envelope = sample_envelope(Some("teams-legal"));
let before = envelope.clone();
apply_welcome_flow_override(&store, &mut envelope, None, None).expect("ok");
assert_eq!(envelope.pack_id, before.pack_id);
assert_eq!(envelope.flow_id, before.flow_id);
assert_eq!(envelope.flow_type, before.flow_type);
}
#[test]
fn override_is_no_op_when_endpoint_id_absent() {
let store = new_session_store();
let mut envelope = sample_envelope(None);
let before = envelope.clone();
apply_welcome_flow_override(&store, &mut envelope, Some(&hint()), Some("welcome".into()))
.expect("ok");
assert_eq!(envelope.pack_id, before.pack_id);
assert_eq!(envelope.flow_id, before.flow_id);
}
#[test]
fn override_swaps_pack_flow_and_threads_flow_type_through() {
for hint_flow_type in [Some("welcome".to_string()), None] {
let store = new_session_store();
let mut envelope = sample_envelope(Some("teams-legal"));
apply_welcome_flow_override(
&store,
&mut envelope,
Some(&hint()),
hint_flow_type.clone(),
)
.expect("ok");
assert_eq!(envelope.pack_id.as_deref(), Some("pack.welcome"));
assert_eq!(envelope.flow_id, "flow.welcome");
assert_eq!(envelope.flow_type, hint_flow_type);
}
}
#[test]
fn override_is_no_op_on_repeat_turn_with_existing_session() {
let store = new_session_store();
let envelope_template = sample_envelope(Some("teams-legal"));
seed_resume(&store, &envelope_template);
let mut envelope = envelope_template.clone();
apply_welcome_flow_override(&store, &mut envelope, Some(&hint()), Some("welcome".into()))
.expect("ok");
assert_eq!(envelope.pack_id, envelope_template.pack_id);
assert_eq!(envelope.flow_id, envelope_template.flow_id);
assert_eq!(envelope.flow_type, envelope_template.flow_type);
}
#[test]
fn override_is_no_op_post_completion_when_marker_present() {
let store = new_session_store();
let mut first = sample_envelope(Some("teams-legal"));
apply_welcome_flow_override(&store, &mut first, Some(&hint()), Some("welcome".into()))
.expect("first turn ok");
assert_eq!(
first.pack_id.as_deref(),
Some("pack.welcome"),
"first turn fires welcome"
);
let resume = FlowResumeStore::new(Arc::clone(&store));
resume.clear(&first).expect("clear post-completion wait");
let mut second = sample_envelope(Some("teams-legal"));
apply_welcome_flow_override(&store, &mut second, Some(&hint()), Some("welcome".into()))
.expect("second turn ok");
assert_eq!(
second.pack_id.as_deref(),
Some("pack.default"),
"second turn must NOT re-fire welcome"
);
assert_eq!(second.flow_id, "flow.default");
}
#[test]
fn override_is_no_op_on_second_turn_after_marker_set() {
let store = new_session_store();
let mut first = sample_envelope(Some("teams-legal"));
apply_welcome_flow_override(&store, &mut first, Some(&hint()), Some("welcome".into()))
.expect("first turn ok");
assert_eq!(first.pack_id.as_deref(), Some("pack.welcome"));
let mut second = sample_envelope(Some("teams-legal"));
apply_welcome_flow_override(&store, &mut second, Some(&hint()), Some("welcome".into()))
.expect("second turn ok");
assert_eq!(
second.pack_id.as_deref(),
Some("pack.default"),
"second turn must NOT re-fire welcome"
);
}
#[test]
fn override_partitions_marker_per_endpoint() {
let store = new_session_store();
let mut legal = sample_envelope(Some("teams-legal"));
apply_welcome_flow_override(&store, &mut legal, Some(&hint()), Some("welcome".into()))
.expect("legal first turn ok");
assert_eq!(legal.pack_id.as_deref(), Some("pack.welcome"));
let mut accounting = sample_envelope(Some("teams-accounting"));
apply_welcome_flow_override(
&store,
&mut accounting,
Some(&hint()),
Some("welcome".into()),
)
.expect("accounting first turn ok");
assert_eq!(
accounting.pack_id.as_deref(),
Some("pack.welcome"),
"different endpoint = independent first contact"
);
}
#[test]
fn override_partitions_marker_per_user_on_same_endpoint() {
let store = new_session_store();
let mut a1 = sample_envelope_for_user(Some("teams-legal"), "user-a");
apply_welcome_flow_override(&store, &mut a1, Some(&hint()), Some("welcome".into()))
.expect("user-a first ok");
assert_eq!(a1.pack_id.as_deref(), Some("pack.welcome"));
let mut b1 = sample_envelope_for_user(Some("teams-legal"), "user-b");
apply_welcome_flow_override(&store, &mut b1, Some(&hint()), Some("welcome".into()))
.expect("user-b first must not collide with user-a marker");
assert_eq!(
b1.pack_id.as_deref(),
Some("pack.welcome"),
"user-b is independent first contact"
);
let mut a2 = sample_envelope_for_user(Some("teams-legal"), "user-a");
apply_welcome_flow_override(&store, &mut a2, Some(&hint()), Some("welcome".into()))
.expect("user-a second ok");
assert_eq!(
a2.pack_id.as_deref(),
Some("pack.default"),
"user-a must not be re-welcomed after user-b joined"
);
let mut b2 = sample_envelope_for_user(Some("teams-legal"), "user-b");
apply_welcome_flow_override(&store, &mut b2, Some(&hint()), Some("welcome".into()))
.expect("user-b second ok");
assert_eq!(b2.pack_id.as_deref(), Some("pack.default"));
}
#[test]
fn marker_is_not_written_when_hint_absent() {
let store = new_session_store();
let mut envelope = sample_envelope(Some("teams-legal"));
apply_welcome_flow_override(&store, &mut envelope, None, None).expect("ok");
let mut next = sample_envelope(Some("teams-legal"));
apply_welcome_flow_override(&store, &mut next, Some(&hint()), Some("welcome".into()))
.expect("ok");
assert_eq!(
next.pack_id.as_deref(),
Some("pack.welcome"),
"no marker leaked from hint-absent path"
);
}
}
#[cfg(test)]
mod identify_endpoints_tests {
use super::*;
fn dummy_runner_host() -> RunnerHost {
let session_store = new_session_store();
let state_store = new_state_store();
RunnerHost {
configs: HashMap::new(),
active: Arc::new(ActivePacks::new()),
health: Arc::new(HealthState::new()),
session_host: session_host_from(session_store.clone()),
state_host: state_host_from(state_store.clone()),
session_store,
state_store,
wasi_policy: Arc::new(RunnerWasiPolicy::new()),
secrets_manager: default_manager().expect("default secrets manager"),
telemetry: None,
}
}
#[tokio::test]
async fn empty_provider_types_returns_empty_map_without_loading_revision() {
let host = dummy_runner_host();
let map = host
.identify_messaging_endpoints_for_revision(
"demo",
DeploymentId::new(),
BundleId::new("anything"),
RevisionId::new(),
&[],
b"{}",
)
.await
.expect("empty types is the cheap fast path");
assert!(map.is_empty());
}
#[tokio::test]
async fn missing_revision_surfaces_clear_error() {
let host = dummy_runner_host();
let deployment = DeploymentId::new();
let revision = RevisionId::new();
let err = host
.identify_messaging_endpoints_for_revision(
"demo",
deployment,
BundleId::new("missing"),
revision,
&["teams"],
b"{}",
)
.await
.expect_err("missing revision must fail closed");
let msg = format!("{err:#}");
assert!(
msg.contains("revision runtime not loaded"),
"error chain should name the failure mode, got: {msg}"
);
assert!(
msg.contains(&deployment.to_string()),
"error chain should name the deployment id, got: {msg}"
);
assert!(
msg.contains(&revision.to_string()),
"error chain should name the revision id, got: {msg}"
);
}
#[tokio::test]
async fn scoped_empty_provider_types_returns_empty_map() {
let host = dummy_runner_host();
let map = host
.identify_messaging_endpoints_for_revision_scoped(
"demo",
DeploymentId::new(),
BundleId::new("anything"),
RevisionId::new(),
&[],
&[],
&Value::Null,
)
.await
.expect("empty types is the cheap fast path");
assert!(map.is_empty());
}
#[tokio::test]
async fn scoped_missing_revision_surfaces_clear_error() {
let host = dummy_runner_host();
let deployment = DeploymentId::new();
let revision = RevisionId::new();
let err = host
.identify_messaging_endpoints_for_revision_scoped(
"demo",
deployment,
BundleId::new("missing"),
revision,
&["teams"],
&[],
&Value::Null,
)
.await
.expect_err("missing revision must fail closed");
let msg = format!("{err:#}");
assert!(
msg.contains("revision runtime not loaded"),
"error chain should name the failure mode, got: {msg}"
);
assert!(
msg.contains(&deployment.to_string()),
"error chain should name the deployment id, got: {msg}"
);
assert!(
msg.contains(&revision.to_string()),
"error chain should name the revision id, got: {msg}"
);
}
#[test]
fn identify_futures_are_send() {
fn assert_send<F: Send>(_: F) {}
let host = dummy_runner_host();
assert_send(host.identify_messaging_endpoints_for_revision(
"demo",
DeploymentId::new(),
BundleId::new("anything"),
RevisionId::new(),
&["teams"],
b"{}",
));
assert_send(host.identify_messaging_endpoints_for_revision_scoped(
"demo",
DeploymentId::new(),
BundleId::new("anything"),
RevisionId::new(),
&["teams"],
&[("x-telegram-bot-api-secret-token".into(), "tok".into())],
&Value::Null,
));
assert_send(host.describe_identify_instances_for_revision(
"demo",
DeploymentId::new(),
BundleId::new("anything"),
RevisionId::new(),
&["teams"],
));
assert_send(host.invoke_provider_for_revision(
"demo",
DeploymentId::new(),
BundleId::new("anything"),
RevisionId::new(),
"messaging.telegram.bot",
"ingest_http",
b"{}".to_vec(),
None,
None,
));
}
#[tokio::test]
async fn invoke_provider_missing_revision_surfaces_clear_error() {
let host = dummy_runner_host();
let deployment = DeploymentId::new();
let revision = RevisionId::new();
let err = host
.invoke_provider_for_revision(
"demo",
deployment,
BundleId::new("missing"),
revision,
"messaging.telegram.bot",
"ingest_http",
b"{}".to_vec(),
None,
None,
)
.await
.expect_err("missing revision must fail closed");
let msg = format!("{err:#}");
assert!(
msg.contains("revision runtime not loaded"),
"error chain should name the failure mode, got: {msg}"
);
assert!(
msg.contains(&deployment.to_string()),
"error chain should name the deployment id, got: {msg}"
);
assert!(
msg.contains(&revision.to_string()),
"error chain should name the revision id, got: {msg}"
);
}
}
#[cfg(all(test, feature = "greentic-x-provider"))]
mod fast2flow_tests {
use greentic_x_runtime::Fast2FlowRoutingEntity;
use super::*;
#[test]
fn dispatch_target_attaches_prefill_entities_to_payload() {
let activity = Activity::text("show traffic tomorrow");
let routed = apply_fast2flow_target(
activity,
"telco-x/prefix-traffic",
vec![Fast2FlowRoutingEntity::new("date", "20260611").with_format("iso", "2026-06-11")],
)
.expect("target should route");
assert_eq!(routed.pack_id(), Some("telco-x"));
assert_eq!(routed.flow_id(), Some("prefix-traffic"));
assert_eq!(
routed.payload()["fast2flow"]["entities"][0]["normalized"],
"20260611"
);
assert_eq!(
routed.payload()["fast2flow"]["entities"][0]["formats"]["iso"],
"2026-06-11"
);
}
}