use std::sync::Arc;
use aion_core::{Event, Payload, RunId, TimerId, WorkflowId};
use chrono::Utc;
use crate::EngineError;
use crate::durability::{
ContinuationTerminal, ContinuedGeneration, OpeningGeneration, WorkflowStartRecord,
};
use crate::lifecycle::start::{
StartWorkflowContext, abort_unmonitored_start, arm_declared_deadline_or_fail_start,
deadline_fire_at, install_started_monitor,
};
use crate::registry::{SucceedingRunParts, WorkflowHandle};
use crate::supervision::spawn_workflow_with_policy;
pub(crate) enum ContinuationOrigin {
RecordsTheTerminal {
input: Payload,
workflow_type: Option<String>,
},
TerminalAlreadyRecorded,
}
pub(crate) struct ContinuationRequest {
pub predecessor_run: RunId,
pub origin: ContinuationOrigin,
pub workflow_type: String,
pub input: Payload,
}
pub(crate) enum ContinuationOutcome {
Opened(Box<WorkflowHandle>),
AlreadyOpen,
}
pub(crate) async fn open_successor_generation(
context: &StartWorkflowContext,
predecessor: &WorkflowHandle,
request: ContinuationRequest,
) -> Result<ContinuationOutcome, EngineError> {
let ContinuationRequest {
predecessor_run,
origin,
workflow_type,
input,
} = request;
let pinned = crate::lifecycle::start::admission::resolve_contract(
&context.catalog,
&workflow_type,
None,
)?;
let loaded = pinned.workflow();
let workflow_id = predecessor.workflow_id().clone();
let successor_run = RunId::new_v4();
let recorder = predecessor.recorder();
let mut recorder = recorder.lock().await;
let recorded_at = Utc::now();
let history = context.store.read_history(&workflow_id).await?;
let TransitionAdmission::Append(terminal) =
admit_transition(context, &history, &workflow_id, &predecessor_run, origin)?
else {
return Ok(ContinuationOutcome::AlreadyOpen);
};
let deadline = successor_deadline(&successor_run, loaded.declared_timeout(), recorded_at)?;
let armed_deadline = recorder
.record_continue_as_new_boundary(
recorded_at,
ContinuedGeneration {
run_id: predecessor_run.clone(),
terminal,
outstanding_deadline: crate::time::outstanding_deadline_timer(
&history,
&predecessor_run,
),
},
OpeningGeneration {
start: WorkflowStartRecord {
workflow_type: loaded.workflow_type().to_owned(),
input: input.clone(),
run_id: successor_run.clone(),
parent_run_id: Some(predecessor_run.clone()),
parent_workflow_id: None,
package_version: crate::loader::package_version_of(loaded.version()),
},
deadline,
},
)
.await?;
let successor = publish_successor(
context,
predecessor,
loaded,
SuccessorProcess {
input: &input,
predecessor_run: &predecessor_run,
successor_run: successor_run.clone(),
},
)?;
drop(recorder);
arm_declared_deadline_or_fail_start(context, &workflow_id, &successor, armed_deadline).await?;
install_started_monitor(context, &successor)?;
super::visibility::upsert_workflow_visibility(
Arc::clone(&context.store),
Arc::clone(&context.visibility_store),
&workflow_id,
&successor_run,
)
.await?;
Ok(ContinuationOutcome::Opened(Box::new(successor)))
}
enum TransitionAdmission {
Append(Option<ContinuationTerminal>),
AlreadyOpen,
}
struct SuccessorProcess<'a> {
input: &'a Payload,
predecessor_run: &'a RunId,
successor_run: RunId,
}
fn admit_transition(
context: &StartWorkflowContext,
history: &[Event],
workflow_id: &WorkflowId,
predecessor_run: &RunId,
origin: ContinuationOrigin,
) -> Result<TransitionAdmission, EngineError> {
if successor_already_started(history, predecessor_run) {
return Ok(TransitionAdmission::AlreadyOpen);
}
if let Some(incumbent) = context.registry.sole_handle(workflow_id)?
&& incumbent.run_id() != predecessor_run
{
return Err(EngineError::WorkflowWriterHeld {
workflow_id: workflow_id.to_string(),
holder_run_id: incumbent.run_id().to_string(),
holder_pid: incumbent.pid(),
});
}
match origin {
ContinuationOrigin::RecordsTheTerminal {
input,
workflow_type,
} => {
refuse_continuation_from_a_terminal_run(history, workflow_id, predecessor_run)?;
super::continue_as_new::guard_no_pending_work(aion_core::run_segment(
history,
predecessor_run,
))?;
Ok(TransitionAdmission::Append(Some(ContinuationTerminal {
input,
workflow_type,
})))
}
ContinuationOrigin::TerminalAlreadyRecorded => Ok(TransitionAdmission::Append(None)),
}
}
fn publish_successor(
context: &StartWorkflowContext,
predecessor: &WorkflowHandle,
loaded: &crate::loader::LoadedWorkflow,
process: SuccessorProcess<'_>,
) -> Result<WorkflowHandle, EngineError> {
let SuccessorProcess {
input,
predecessor_run,
successor_run,
} = process;
context
.supervision
.ensure_type_supervisor(loaded.workflow_type())?;
let runtime_input = crate::runtime::RuntimeInput::from_payload(input)?;
let pid = spawn_workflow_with_policy(
&context.runtime,
loaded.deployed_entry_module(),
loaded.entry_function(),
runtime_input,
)?;
if let Err(error) = context
.supervision
.place_workflow(loaded.workflow_type(), pid)
{
return Err(abort_unmonitored_start(&context.runtime, pid, error));
}
let successor = predecessor.succeeding_run(SucceedingRunParts {
run_id: successor_run.clone(),
pid,
loaded_version: loaded.version().clone(),
});
if let Err(error) = context.registry.rekey_generation(
predecessor.workflow_id(),
predecessor_run,
successor_run,
successor.clone(),
) {
tracing::error!(
workflow_id = %predecessor.workflow_id(),
predecessor_run = %predecessor_run,
successor_run = %successor.run_id(),
error = %error,
"the continue-as-new successor is durably started but could not be published; its \
process is being aborted and startup recovery re-installs the run"
);
return Err(abort_unmonitored_start(&context.runtime, pid, error));
}
Ok(successor)
}
fn successor_already_started(history: &[Event], predecessor_run: &RunId) -> bool {
history.iter().any(|event| {
matches!(
event,
Event::WorkflowStarted {
parent_run_id: Some(existing),
..
} if existing == predecessor_run
)
})
}
fn refuse_continuation_from_a_terminal_run(
history: &[Event],
workflow_id: &WorkflowId,
run: &RunId,
) -> Result<(), EngineError> {
if super::completion::terminal_outcome_from_history(history, run).is_some() {
return Err(EngineError::Runtime {
reason: format!(
"continue_as_new rejected: workflow {workflow_id} run {run} already recorded a \
terminal event"
),
});
}
Ok(())
}
fn successor_deadline(
run_id: &RunId,
declared_timeout: Option<std::time::Duration>,
started_at: chrono::DateTime<Utc>,
) -> Result<Option<(TimerId, chrono::DateTime<Utc>)>, EngineError> {
let Some(timeout) = declared_timeout else {
return Ok(None);
};
let fire_at = deadline_fire_at(started_at, timeout)?;
let deadline_id =
crate::time::deadline_timer_id(run_id).map_err(|error| EngineError::Runtime {
reason: format!("failed to mint deadline timer id: {error}"),
})?;
Ok(Some((deadline_id, fire_at)))
}