use crate::{
CloudEvent, Digest, Error, InvocationIdentity, ProgramDescriptor, ProgramOutcome, ProgramRef,
};
use serde::{Deserialize, Serialize};
use serde_json::Value;
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ExecutionRequest {
pub descriptor: ProgramDescriptor,
pub event: CloudEvent,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ExecutionContext {
#[serde(flatten)]
pub identity: InvocationIdentity,
pub program: ProgramRef,
pub digest: Digest,
}
impl From<&ExecutionRequest> for ExecutionContext {
fn from(request: &ExecutionRequest) -> Self {
Self {
identity: InvocationIdentity::from(&request.event),
program: request.descriptor.program.clone(),
digest: request.descriptor.digest.clone(),
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ExecutionReport {
#[serde(flatten)]
pub context: Box<ExecutionContext>,
pub process_id: u32,
pub reused_process: bool,
pub outcome: ProgramOutcome,
pub elapsed_ms: u64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum Phase {
Admission,
Preparation,
Startup,
Execution,
Cleanup,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ExecutionFailure {
#[serde(flatten)]
pub context: Box<ExecutionContext>,
pub error: Error,
#[serde(skip_serializing_if = "Option::is_none")]
pub cleanup_error: Option<Error>,
pub phase: Phase,
pub execution_may_have_started: bool,
}
impl std::fmt::Display for ExecutionFailure {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{:?}: {}", self.phase, self.error)
}
}
impl std::error::Error for ExecutionFailure {}
pub type ExecutionResult = std::result::Result<ExecutionReport, ExecutionFailure>;
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct RuntimeInvocation {
pub event: CloudEvent,
#[serde(skip_serializing_if = "Option::is_none")]
pub processing_context: Option<crate::TraceContext>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub extension: Option<RuntimeExtension>,
}
impl From<CloudEvent> for RuntimeInvocation {
fn from(event: CloudEvent) -> Self {
Self {
event,
processing_context: None,
extension: None,
}
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct RuntimeExtension {
pub schema: String,
pub payload: Value,
}
impl RuntimeExtension {
pub fn validate(&self) -> crate::Result<()> {
validate_name(&self.schema, "runtime extension schema")?;
crate::validate_runtime_payload(&self.payload, crate::RUNTIME_EXTENSION_MAX_BYTES)
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct RuntimeRequest {
pub id: u64,
pub operation: String,
pub payload: Value,
}
impl RuntimeRequest {
pub fn validate(&self) -> crate::Result<()> {
if self.id == 0 {
return Err(Error::new(
crate::ErrorKind::InvalidInput,
"runtime request ID must be positive",
));
}
validate_name(&self.operation, "runtime request operation")?;
crate::validate_runtime_payload(&self.payload, crate::RUNTIME_REQUEST_MAX_BYTES)
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct RuntimeReply {
pub id: u64,
pub result: Value,
}
fn validate_name(value: &str, label: &str) -> crate::Result<()> {
if value.is_empty() || value.len() > 256 || value.chars().any(char::is_control) {
return Err(Error::new(
crate::ErrorKind::InvalidInput,
format!("{label} must contain 1..=256 bytes without control characters"),
));
}
Ok(())
}