use async_trait::async_trait;
use concepts::ExecutionFailureKind;
use concepts::ExecutionId;
use concepts::ExecutionMetadata;
use concepts::FunctionMetadata;
use concepts::JoinSetId;
use concepts::TrapKind;
use concepts::component_id::ComponentDigest;
use concepts::storage::DbErrorWrite;
use concepts::storage::HistoryEvent;
use concepts::storage::Locked;
use concepts::storage::ResponseWithCursor;
use concepts::storage::Version;
use concepts::storage::http_client_trace::HttpClientTrace;
use concepts::{FinishedExecutionFailure, StrVariant};
use concepts::{FunctionFqn, ParamsParsingError, ResultParsingError};
use concepts::{Params, SupportedFunctionReturnValue};
use tracing::Span;
#[async_trait]
pub trait Worker: Send + Sync + 'static {
async fn run(&self, ctx: WorkerContext) -> WorkerResult;
fn exported_functions_noext(&self) -> &[FunctionMetadata];
}
pub type WorkerResult = Result<WorkerResultOk, WorkerError>;
#[derive(Debug, derive_more::Display)]
pub enum WorkerResultOk {
#[display("db updated by worker or watcher")]
DbUpdatedByWorkerOrWatcher,
#[display("{_0}")]
RunFinished(RunFinished),
}
#[derive(Debug, derive_more::Display)]
#[display("{retval}")]
pub struct RunFinished {
pub retval: SupportedFunctionReturnValue,
pub version: Version,
pub http_client_traces: Option<Vec<HttpClientTrace>>,
}
#[derive(Debug)]
pub struct WorkerContext {
pub execution_id: ExecutionId,
pub metadata: ExecutionMetadata,
pub component_digest: ComponentDigest,
pub ffqn: FunctionFqn,
pub params: Params,
pub event_history: Vec<(HistoryEvent, Version)>,
pub responses: Vec<ResponseWithCursor>,
pub parent: Option<(ExecutionId, JoinSetId)>,
pub version: Version,
pub can_be_retried: bool,
pub worker_span: Span,
pub locked_event: Locked,
pub executor_close_watcher: tokio::sync::watch::Receiver<bool>,
}
#[derive(Debug, thiserror::Error)]
pub enum WorkerError {
#[error("activity {trap_kind}: {reason}")]
ActivityTrap {
reason: String,
trap_kind: TrapKind,
detail: Option<String>,
version: Version,
http_client_traces: Option<Vec<HttpClientTrace>>,
},
#[error("limit reached: {reason}")]
LimitReached { reason: String, version: Version },
#[error("temporary timeout")]
TemporaryTimeout {
http_client_traces: Option<Vec<HttpClientTrace>>,
version: Version,
},
#[error("executor closing")]
ExecutorClosing(Version),
#[error(transparent)]
DbError(DbErrorWrite),
#[error("fatal error: {0}")]
FatalError(FatalError, Version),
}
#[derive(Debug, thiserror::Error)]
pub enum FatalError {
#[error("nondeterminism detected")]
NondeterminismDetected { detail: String },
#[error(transparent)]
ParamsParsingError(ParamsParsingError),
#[error("{reason}")]
CannotInstantiate {
reason: String,
detail: Option<String>,
},
#[error(transparent)]
ResultParsingError(ResultParsingError),
#[error("error calling imported function {ffqn} : {reason}")]
ImportedFunctionCallError {
ffqn: FunctionFqn,
reason: StrVariant,
detail: Option<String>,
},
#[error("workflow {trap_kind}: {reason}")]
WorkflowTrap {
reason: String,
trap_kind: TrapKind,
detail: Option<String>,
},
#[error("out of fuel: {reason}")]
OutOfFuel { reason: String },
#[error("constraint violation: {reason}")]
ConstraintViolation { reason: StrVariant },
#[error("cancelled")]
Cancelled,
}
impl From<FatalError> for FinishedExecutionFailure {
fn from(err: FatalError) -> Self {
let reason_generic = err.to_string(); match err {
FatalError::NondeterminismDetected { detail } => FinishedExecutionFailure {
reason: None,
kind: ExecutionFailureKind::NondeterminismDetected,
detail: Some(detail),
},
FatalError::OutOfFuel { reason } => FinishedExecutionFailure {
reason: Some(reason),
kind: ExecutionFailureKind::OutOfFuel,
detail: None,
},
FatalError::ParamsParsingError(err) => FinishedExecutionFailure {
reason: Some(reason_generic),
kind: ExecutionFailureKind::Uncategorized,
detail: err.detail(),
},
FatalError::CannotInstantiate { reason, detail } => FinishedExecutionFailure {
reason: Some(reason),
kind: ExecutionFailureKind::Uncategorized,
detail,
},
FatalError::ResultParsingError(_) | FatalError::ConstraintViolation { reason: _ } => {
FinishedExecutionFailure {
reason: Some(reason_generic),
kind: ExecutionFailureKind::Uncategorized,
detail: None,
}
}
FatalError::ImportedFunctionCallError { detail, .. }
| FatalError::WorkflowTrap { detail, .. } => FinishedExecutionFailure {
reason: Some(reason_generic),
kind: ExecutionFailureKind::Uncategorized,
detail,
},
FatalError::Cancelled => FinishedExecutionFailure {
kind: ExecutionFailureKind::Cancelled,
reason: None,
detail: None,
},
}
}
}
impl From<&FatalError> for FinishedExecutionFailure {
fn from(err: &FatalError) -> Self {
let reason_generic = err.to_string(); match err {
FatalError::NondeterminismDetected { detail } => FinishedExecutionFailure {
reason: None,
kind: ExecutionFailureKind::NondeterminismDetected,
detail: Some(detail.clone()),
},
FatalError::OutOfFuel { reason } => FinishedExecutionFailure {
reason: Some(reason.clone()),
kind: ExecutionFailureKind::OutOfFuel,
detail: None,
},
FatalError::ParamsParsingError(err) => FinishedExecutionFailure {
reason: Some(reason_generic),
kind: ExecutionFailureKind::Uncategorized,
detail: err.detail(),
},
FatalError::CannotInstantiate { reason, detail } => FinishedExecutionFailure {
reason: Some(reason.clone()),
kind: ExecutionFailureKind::Uncategorized,
detail: detail.clone(),
},
FatalError::ResultParsingError(_) | FatalError::ConstraintViolation { reason: _ } => {
FinishedExecutionFailure {
reason: Some(reason_generic),
kind: ExecutionFailureKind::Uncategorized,
detail: None,
}
}
FatalError::ImportedFunctionCallError { detail, .. }
| FatalError::WorkflowTrap { detail, .. } => FinishedExecutionFailure {
reason: Some(reason_generic),
kind: ExecutionFailureKind::Uncategorized,
detail: detail.clone(),
},
FatalError::Cancelled => FinishedExecutionFailure {
kind: ExecutionFailureKind::Cancelled,
reason: None,
detail: None,
},
}
}
}