use crate::workflow::replay_advance::AdvanceResponse;
use crate::workflow::workflow_js_worker::WorkflowJsWorker;
use crate::workflow::workflow_worker::{
AdvanceError, BacktraceCapture, ReplayAdvanceable, ReplayError, ReplayResponse, WorkflowWorker,
};
use concepts::ComponentId;
use concepts::ComponentType;
use concepts::ExecutionId;
use concepts::FunctionFqn;
use concepts::FunctionMetadata;
use concepts::FunctionRegistry;
use concepts::IfcFqnName;
use concepts::PackageIfcFns;
use concepts::StrVariant;
use concepts::component_id::ComponentDigest;
use hashbrown::HashMap;
use indexmap::IndexMap;
use std::fmt::Debug;
use std::ops::Deref;
use std::sync::Arc;
use tracing::error;
pub use concepts::storage::WitOrigin;
#[derive(Debug, Clone)]
pub struct ComponentConfig {
pub component_id: ComponentId,
pub imports: Vec<FunctionMetadata>,
pub workflow_or_activity_config: Option<ComponentConfigImportable>,
pub wit: String,
pub wit_origin: WitOrigin,
}
#[derive(Clone)]
pub enum ReplayWorker {
Wasm(Arc<WorkflowWorker>),
Js(Arc<WorkflowJsWorker>),
}
impl Debug for ReplayWorker {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
ReplayWorker::Wasm(_) => f.write_str("ReplayWorker::Wasm(..)"),
ReplayWorker::Js(_) => f.write_str("ReplayWorker::Js(..)"),
}
}
}
impl ReplayWorker {
pub async fn replay(
&self,
execution_id: ExecutionId,
backtrace_capture: BacktraceCapture,
) -> Result<ReplayResponse, ReplayError> {
match self {
Self::Wasm(worker) => worker.replay(execution_id, backtrace_capture).await,
Self::Js(worker) => worker.replay(execution_id, backtrace_capture).await,
}
}
pub async fn persist_backtraces(
&self,
execution_id: ExecutionId,
) -> Result<usize, ReplayError> {
match self {
Self::Wasm(worker) => worker.persist_backtraces(execution_id).await,
Self::Js(worker) => worker.persist_backtraces(execution_id).await,
}
}
pub async fn advance(
&self,
execution_id: ExecutionId,
requested: ReplayAdvanceable,
backtrace_capture: BacktraceCapture,
) -> Result<AdvanceResponse, AdvanceError> {
match self {
Self::Wasm(worker) => {
worker
.advance(execution_id, requested, backtrace_capture)
.await
}
Self::Js(worker) => {
worker
.advance(execution_id, requested, backtrace_capture)
.await
}
}
}
}
#[derive(Default, Debug, Clone)]
pub struct ReplayWorkerRegistry {
workers: HashMap<ComponentDigest, (ComponentId, ReplayWorker)>,
}
impl ReplayWorkerRegistry {
pub fn insert(&mut self, component_id: ComponentId, worker: ReplayWorker) {
let digest = component_id.component_digest.clone();
let old = self.workers.insert(digest, (component_id, worker));
assert!(
old.is_none(),
"replay worker already registered for this digest"
);
}
#[must_use]
pub fn get(&self, digest: &ComponentDigest) -> Option<(&ComponentId, &ReplayWorker)> {
self.workers.get(digest).map(|(id, w)| (id, w))
}
}
#[derive(Debug, Clone)]
pub struct ComponentConfigImportable {
pub exports_ext: Vec<FunctionMetadata>,
pub exports_hierarchy_ext: Vec<PackageIfcFns>,
}
#[derive(Default, Debug)]
pub struct ComponentConfigRegistry {
inner: ComponentConfigRegistryInner,
}
#[derive(Default, Debug)]
struct ComponentConfigRegistryInner {
exported_ffqns_ext: IndexMap<FunctionFqn, (ComponentId, FunctionMetadata)>,
export_hierarchy: Vec<PackageIfcFns>,
export_hierarchy_origin: HashMap<IfcFqnName, WitOrigin>,
names_to_components: IndexMap<StrVariant, ComponentConfig>,
digests_to_wit: IndexMap<ComponentDigest, String>,
}
#[derive(Debug, Clone, thiserror::Error)]
#[error("registering component failed: {0}")]
pub struct ComponentInsertionError(StrVariant);
impl ComponentConfigRegistry {
pub fn insert(&mut self, component: ComponentConfig) -> Result<(), ComponentInsertionError> {
let name = &component.component_id.name;
if self.inner.names_to_components.contains_key(name) {
return Err(ComponentInsertionError(
format!("component with name `{name}` is already registered").into(),
));
}
if let Some(workflow_or_activity_config) = &component.workflow_or_activity_config {
if self
.inner
.digests_to_wit
.contains_key(&component.component_id.component_digest)
{
return Err(ComponentInsertionError(
format!(
"component {} is already inserted with the same digest",
component.component_id
)
.into(),
));
}
for exported_ffqn in workflow_or_activity_config
.exports_ext
.iter()
.map(|f| &f.ffqn)
{
if let Some((conflicting_id, _)) = self.inner.exported_ffqns_ext.get(exported_ffqn)
{
return Err(ComponentInsertionError(
format!(
"function {exported_ffqn} is already exported by component {conflicting_id}, cannot insert {}",
component.component_id
).into()));
}
}
for exported_fn_metadata in &workflow_or_activity_config.exports_ext {
let old = self.inner.exported_ffqns_ext.insert(
exported_fn_metadata.ffqn.clone(),
(component.component_id.clone(), exported_fn_metadata.clone()),
);
assert!(old.is_none());
}
for new_ifc_fns in &workflow_or_activity_config.exports_hierarchy_ext {
if !new_ifc_fns.extension {
if let Some(&existing_origin) =
self.inner.export_hierarchy_origin.get(&new_ifc_fns.ifc_fqn)
{
if existing_origin != component.wit_origin {
return Err(ComponentInsertionError(
format!(
"interface `{}` is already exported by a {} component, cannot insert {} which is a {} component",
new_ifc_fns.ifc_fqn,
existing_origin,
component.component_id,
component.wit_origin,
)
.into(),
));
}
} else {
self.inner
.export_hierarchy_origin
.insert(new_ifc_fns.ifc_fqn.clone(), component.wit_origin);
}
}
if let Some(existing) = self.inner.export_hierarchy.iter_mut().find(|e| {
e.ifc_fqn == new_ifc_fns.ifc_fqn && e.extension == new_ifc_fns.extension
}) {
existing.fns.extend(new_ifc_fns.fns.clone());
} else {
self.inner.export_hierarchy.push(new_ifc_fns.clone());
}
}
let old = self.inner.digests_to_wit.insert(
component.component_id.component_digest.clone(),
component.wit.clone(),
);
assert!(old.is_none());
} else if component.component_id.component_type == ComponentType::WebhookEndpoint {
self.inner
.digests_to_wit
.entry(component.component_id.component_digest.clone())
.or_insert(component.wit.clone());
}
let old = self
.inner
.names_to_components
.insert(name.clone(), component);
assert!(old.is_none());
Ok(())
}
pub fn verify_registry(
self,
) -> (
ComponentConfigRegistryRO,
Option<String>, /* supressed_errors */
) {
let mut errors = Vec::new();
for examined_component in self.inner.names_to_components.values() {
self.verify_imports_component(examined_component, &mut errors);
}
let errors = if !errors.is_empty() {
let errors = errors.join("\n");
tracing::warn!("component resolution error: \n{errors}");
Some(errors)
} else {
None
};
(
ComponentConfigRegistryRO {
inner: Arc::new(self.inner),
},
errors,
)
}
fn additional_import_allowlist(
import: &FunctionMetadata,
component_type: ComponentType,
) -> bool {
match component_type {
ComponentType::Activity => {
match import.ffqn.ifc_fqn.namespace() {
"wasi" => true,
"obelisk" => import.ffqn.ifc_fqn.deref() == "obelisk:log/log@1.0.0",
_ => false,
}
}
ComponentType::Workflow => {
matches!(
import.ffqn.ifc_fqn.pkg_fqn_name().to_string().as_str(),
"obelisk:log@1.0.0"
| "obelisk:workflow@6.0.0"
| "obelisk:workflow@5.0.0"
| "obelisk:workflow@5.1.0"
| "obelisk:types@5.0.0"
| "obelisk:types@4.2.0"
)
}
ComponentType::WebhookEndpoint => {
match import.ffqn.ifc_fqn.namespace() {
"wasi" => true,
"obelisk" => matches!(
import.ffqn.ifc_fqn.pkg_fqn_name().to_string().as_str(),
"obelisk:webhook@6.0.0"
| "obelisk:webhook@5.3.0"
| "obelisk:webhook@5.2.0"
| "obelisk:webhook@5.1.0"
| "obelisk:webhook@5.0.0"
| "obelisk:log@1.0.0"
| "obelisk:types@5.0.0"
| "obelisk:types@4.0.0"
| "obelisk:types@4.1.0"
| "obelisk:types@4.2.0"
),
_ => false,
}
}
ComponentType::ActivityStub | ComponentType::Cron => false,
}
}
fn verify_imports_component(&self, component: &ComponentConfig, errors: &mut Vec<String>) {
let component_id = &component.component_id;
for imported_fn_metadata in &component.imports {
if let Some((exported_component_id, exported_fn_metadata)) = self
.inner
.exported_ffqns_ext
.get(&imported_fn_metadata.ffqn)
{
if imported_fn_metadata.parameter_types != exported_fn_metadata.parameter_types {
error!(
"Parameter types do not match: {ffqn} imported by {component_id} , exported by {exported_component_id}",
ffqn = imported_fn_metadata.ffqn
);
error!(
"Import {import}",
import = serde_json::to_string(imported_fn_metadata).unwrap(), );
error!(
"Export {export}",
export = serde_json::to_string(exported_fn_metadata).unwrap(),
);
errors.push(format!("parameter types do not match: {component_id} imports {imported_fn_metadata} , {exported_component_id} exports {exported_fn_metadata}"));
}
if imported_fn_metadata.return_type != exported_fn_metadata.return_type {
error!(
"Return types do not match: {ffqn} imported by {component_id} , exported by {exported_component_id}",
ffqn = imported_fn_metadata.ffqn
);
error!(
"Import {import}",
import = serde_json::to_string(imported_fn_metadata).unwrap(), );
error!(
"Export {export}",
export = serde_json::to_string(exported_fn_metadata).unwrap(),
);
errors.push(format!("return types do not match: {component_id} imports {imported_fn_metadata} , {exported_component_id} exports {exported_fn_metadata}"));
}
} else if !Self::additional_import_allowlist(
imported_fn_metadata,
component_id.component_type,
) {
errors.push(format!(
"function imported by {component_id} not found: {imported_fn_metadata}"
));
}
}
}
}
#[derive(Debug, Clone)]
pub struct ComponentConfigRegistryRO {
inner: Arc<ComponentConfigRegistryInner>,
}
impl ComponentConfigRegistryRO {
#[must_use]
pub fn get_wit(&self, input_digest: &ComponentDigest) -> Option<&str> {
self.inner
.digests_to_wit
.get(input_digest)
.map(std::string::String::as_str)
}
#[must_use]
pub fn find_by_exported_ffqn_submittable(
&self,
ffqn: &FunctionFqn,
) -> Option<(&ComponentId, &FunctionMetadata)> {
self.inner
.exported_ffqns_ext
.get(ffqn)
.and_then(|(component_id, fn_metadata)| {
if fn_metadata.submittable {
Some((component_id, fn_metadata))
} else {
None
}
})
}
#[must_use]
pub fn find_by_exported_ffqn(
&self,
ffqn: &FunctionFqn,
) -> Option<(&ComponentId, &FunctionMetadata)> {
self.inner
.exported_ffqns_ext
.get(ffqn)
.map(|t| (&t.0, &t.1))
}
#[must_use]
pub fn find_by_exported_ffqn_stub(
&self,
ffqn: &FunctionFqn,
) -> Option<(&ComponentId, &FunctionMetadata)> {
self.inner
.exported_ffqns_ext
.get(ffqn)
.and_then(|(component_id, fn_metadata)| {
if component_id.component_type == ComponentType::ActivityStub {
assert!(!ffqn.ifc_fqn.is_extension());
Some((component_id, fn_metadata))
} else {
None
}
})
}
#[must_use]
pub fn list(&self, extensions: bool) -> Vec<ComponentConfig> {
self.inner
.names_to_components
.values()
.cloned()
.map(|mut component| {
if !extensions && let Some(importable) = &mut component.workflow_or_activity_config
{
importable
.exports_ext
.retain(|fn_metadata| !fn_metadata.ffqn.ifc_fqn.is_extension());
importable
.exports_hierarchy_ext
.retain(|ifc_fns| !ifc_fns.extension);
}
component
})
.collect()
}
}
impl FunctionRegistry for ComponentConfigRegistryRO {
fn get_by_exported_function(
&self,
ffqn: &FunctionFqn,
) -> Option<(FunctionMetadata, ComponentId)> {
if ffqn.ifc_fqn.is_extension() {
None
} else {
self.inner
.exported_ffqns_ext
.get(ffqn)
.map(|(id, metadata)| (metadata.clone(), id.clone()))
}
}
fn all_exports(&self) -> &[PackageIfcFns] {
&self.inner.export_hierarchy
}
}