use aion_core::{ActivityErrorKind, ActivityId, ContentType, Payload, WorkflowId};
use crate::error::ServerError;
use crate::worker::dispatch::{ActivityCompletion, ActivityCompletionOutcome};
use super::PendingActivities;
impl PendingActivities {
#[must_use]
pub const fn transport_losses(&self) -> &crate::worker::transport_loss::TransportLossLedger {
&self.transport_losses
}
pub(super) fn classify_worker_loss(
&self,
workflow_id: &WorkflowId,
activity_id: &ActivityId,
worker_id: crate::worker::registry::WorkerId,
) -> String {
let detail = crate::worker::transport_loss::worker_lost_detail(worker_id);
match self
.transport_losses
.record_loss(workflow_id, activity_id, &detail)
{
Ok(verdict) => {
if verdict.exhausted {
tracing::error!(
operation = "activity_complete",
workflow_id = %workflow_id,
activity_id = %activity_id,
worker_id = ?worker_id,
error_type = "TransportExhausted",
losses = verdict.losses,
budget_ms = self.transport_losses.budget().as_millis(),
"activity abandoned: the transport kept losing its worker past the \
transport-loss budget"
);
} else {
tracing::warn!(
operation = "activity_complete",
workflow_id = %workflow_id,
activity_id = %activity_id,
worker_id = ?worker_id,
error_type = "WorkerLost",
losses = verdict.losses,
budget_ms = self.transport_losses.budget().as_millis(),
"worker lost before reporting an activity result; the activity never ran \
and will be re-dispatched attempt-neutrally"
);
}
verdict.reason
}
Err(error) => {
tracing::error!(
workflow_id = %workflow_id,
activity_id = %activity_id,
%error,
"transport-loss ledger is unreadable; abandoning the activity rather than \
re-dispatching it without a budget"
);
format!(
"{}{detail} (transport-loss budget unreadable: {error})",
crate::worker::transport_loss::TRANSPORT_EXHAUSTED_REASON_PREFIX
)
}
}
}
fn classify_admission_refusal(
&self,
workflow_id: &WorkflowId,
activity_id: &ActivityId,
worker_id: crate::worker::registry::WorkerId,
reason: &str,
) -> String {
let detail = format!("worker {worker_id:?} refused the dispatch: {reason}");
match self
.transport_losses
.record_loss(workflow_id, activity_id, &detail)
{
Ok(verdict) if verdict.exhausted => {
tracing::error!(
operation = "activity_complete",
workflow_id = %workflow_id,
activity_id = %activity_id,
worker_id = ?worker_id,
error_type = "TransportExhausted",
refusals = verdict.losses,
budget_ms = self.transport_losses.budget().as_millis(),
reason,
"activity abandoned: workers kept refusing it past the transport-loss budget. The server tracks every worker's advertised concurrency and does not select a full one, so a refusal this persistent means the server's count and the worker's admission disagree"
);
verdict.reason
}
Ok(verdict) => {
tracing::info!(
operation = "activity_complete",
workflow_id = %workflow_id,
activity_id = %activity_id,
worker_id = ?worker_id,
refusals = verdict.losses,
reason,
"worker refused a dispatch it had no slot for; re-parking it clock-free and re-selecting, with no attempt consumed and nothing recorded"
);
verdict.reason
}
Err(error) => {
tracing::error!(
workflow_id = %workflow_id,
activity_id = %activity_id,
%error,
"transport-loss ledger is unreadable; abandoning the refused activity rather than re-dispatching it without a budget"
);
format!(
"{}{detail} (transport-loss budget unreadable: {error})",
crate::worker::transport_loss::TRANSPORT_EXHAUSTED_REASON_PREFIX
)
}
}
}
pub(crate) fn complete_activity_after_accept(
&self,
completion: ActivityCompletion,
after_accept: impl FnOnce() -> Result<(), ServerError>,
) -> Result<(), ServerError> {
let result = match completion.outcome {
ActivityCompletionOutcome::Succeeded(payload) => {
payload_to_string(&payload).map_err(|reason| {
tracing::error!(
operation = "activity_complete",
workflow_id = %completion.workflow_id,
activity_id = %completion.activity_id,
error_type = "ActivityResultDecode",
%reason,
"activity completion failed"
);
ServerError::worker_dispatch("", "", format!("payload decode: {reason}"))
})?
}
ActivityCompletionOutcome::Failed(error) => {
let prefix = match error.kind {
ActivityErrorKind::Retryable => "retryable",
ActivityErrorKind::PolicyRefused => "policy_refused",
ActivityErrorKind::Terminal => "terminal",
};
tracing::error!(
operation = "activity_complete",
workflow_id = %completion.workflow_id,
activity_id = %completion.activity_id,
error_type = "ActivityFailed",
error_kind = prefix,
reason = %error.message,
"activity completion failed"
);
Err(format!("{prefix}:{}", error.message))
}
ActivityCompletionOutcome::WorkerLost { worker_id } => Err(self.classify_worker_loss(
&completion.workflow_id,
&completion.activity_id,
worker_id,
)),
ActivityCompletionOutcome::Refused { worker_id, reason } => Err(self
.classify_admission_refusal(
&completion.workflow_id,
&completion.activity_id,
worker_id,
&reason,
)),
};
let accepted_settlement = || after_accept().map(|()| true);
self.complete_fenced_after_accept(
&completion.workflow_id,
&completion.activity_id,
completion.run_id.as_ref(),
&completion.completion_token,
result,
accepted_settlement,
)?;
Ok(())
}
}
fn payload_to_string(payload: &Payload) -> Result<Result<String, String>, String> {
match payload.content_type() {
ContentType::Json => String::from_utf8(payload.bytes().to_vec())
.map(Ok)
.map_err(|_| "activity result payload is not valid UTF-8".to_owned()),
}
}