use axum::http::StatusCode;
use axum::response::{IntoResponse, Response};
use axum::Json;
use serde::{Deserialize, Serialize};
use utoipa::ToSchema;
pub const API_VERSION: &str = "v2";
pub const MAX_EVENTS: usize = 1000;
pub const MAX_LOGS: usize = 1000;
pub const MAX_COMMAND_RECORDS: usize = 1000;
pub const COMMAND_RECORD_TTL_SECS: i64 = 24 * 60 * 60;
pub const MAX_CORRELATION_ID_LEN: usize = 64;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
#[serde(rename_all = "snake_case")]
pub enum ErrorCode {
Unauthorized,
Forbidden,
NotFound,
StaleRevision,
LifecycleConflict,
CommandExecutorUnbound,
TargetIneligible,
RootBusy,
IdempotencyMismatch,
RegistryCapacity,
ValidationFailed,
WorktreeExists,
WorktreeNotFound,
WorktreeDirty,
WorktreeDirtyUnknown,
MergeConflict,
ExecutionBindingMismatch,
TransportNotPermitted,
InternalError,
}
impl ErrorCode {
pub fn http_status(self) -> StatusCode {
match self {
Self::Unauthorized => StatusCode::UNAUTHORIZED,
Self::Forbidden => StatusCode::FORBIDDEN,
Self::NotFound | Self::WorktreeNotFound => StatusCode::NOT_FOUND,
Self::StaleRevision
| Self::LifecycleConflict
| Self::CommandExecutorUnbound
| Self::TargetIneligible
| Self::RootBusy
| Self::IdempotencyMismatch
| Self::WorktreeExists
| Self::WorktreeDirty
| Self::WorktreeDirtyUnknown
| Self::MergeConflict
| Self::ExecutionBindingMismatch => StatusCode::CONFLICT,
Self::TransportNotPermitted => StatusCode::FORBIDDEN,
Self::RegistryCapacity => StatusCode::SERVICE_UNAVAILABLE,
Self::ValidationFailed => StatusCode::UNPROCESSABLE_ENTITY,
Self::InternalError => StatusCode::INTERNAL_SERVER_ERROR,
}
}
pub fn as_str(self) -> &'static str {
match self {
Self::Unauthorized => "unauthorized",
Self::Forbidden => "forbidden",
Self::NotFound => "not_found",
Self::StaleRevision => "stale_revision",
Self::LifecycleConflict => "lifecycle_conflict",
Self::CommandExecutorUnbound => "command_executor_unbound",
Self::TargetIneligible => "target_ineligible",
Self::RootBusy => "root_busy",
Self::IdempotencyMismatch => "idempotency_mismatch",
Self::RegistryCapacity => "registry_capacity",
Self::ValidationFailed => "validation_failed",
Self::WorktreeExists => "worktree_exists",
Self::WorktreeNotFound => "worktree_not_found",
Self::WorktreeDirty => "worktree_dirty",
Self::WorktreeDirtyUnknown => "worktree_dirty_unknown",
Self::MergeConflict => "merge_conflict",
Self::ExecutionBindingMismatch => "execution_binding_mismatch",
Self::TransportNotPermitted => "transport_not_permitted",
Self::InternalError => "internal_error",
}
}
}
pub const ALL_ERROR_CODES: [ErrorCode; 19] = [
ErrorCode::Unauthorized,
ErrorCode::Forbidden,
ErrorCode::NotFound,
ErrorCode::StaleRevision,
ErrorCode::LifecycleConflict,
ErrorCode::CommandExecutorUnbound,
ErrorCode::TargetIneligible,
ErrorCode::RootBusy,
ErrorCode::IdempotencyMismatch,
ErrorCode::RegistryCapacity,
ErrorCode::ValidationFailed,
ErrorCode::WorktreeExists,
ErrorCode::WorktreeNotFound,
ErrorCode::WorktreeDirty,
ErrorCode::WorktreeDirtyUnknown,
ErrorCode::MergeConflict,
ErrorCode::ExecutionBindingMismatch,
ErrorCode::TransportNotPermitted,
ErrorCode::InternalError,
];
#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
pub struct ApiError {
pub error_code: ErrorCode,
pub message: String,
pub correlation_id: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub current_revision: Option<u64>,
}
impl ApiError {
pub fn new(error_code: ErrorCode, message: impl Into<String>, correlation_id: &str) -> Self {
Self {
error_code,
message: message.into(),
correlation_id: correlation_id.to_string(),
current_revision: None,
}
}
pub fn with_revision(mut self, revision: u64) -> Self {
self.current_revision = Some(revision);
self
}
}
impl IntoResponse for ApiError {
fn into_response(self) -> Response {
(self.error_code.http_status(), Json(self)).into_response()
}
}
pub fn is_valid_correlation_id(value: &str) -> bool {
!value.is_empty()
&& value.len() <= MAX_CORRELATION_ID_LEN
&& value
.bytes()
.all(|b| b.is_ascii_alphanumeric() || matches!(b, b'.' | b'_' | b':' | b'-'))
}
pub use crate::ids::new_hex_id;
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
#[serde(tag = "type", rename_all = "snake_case", deny_unknown_fields)]
pub enum CommandSpec {
Start,
Stop,
CancelStop,
ForceStop,
SetExecutionMark {
change_id: String,
marked: bool,
},
SetQueueIntent {
change_id: String,
queued: bool,
},
RetryChange {
change_id: String,
},
RetryErrors {
#[serde(default)]
change_ids: Vec<String>,
},
StopAndDequeue {
change_id: String,
},
ForceStopChange {
change_id: String,
},
ResolveMerge {
change_id: String,
},
SetAllExecutionMarks {},
CreateWorktree {
target: ChangeTarget,
#[serde(default)]
params: EmptyParams,
},
DeleteWorktree {
target: WorktreeTarget,
#[serde(default)]
params: EmptyParams,
},
MergeWorktree {
target: WorktreeTarget,
#[serde(default)]
params: EmptyParams,
},
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
#[serde(deny_unknown_fields)]
pub struct ChangeTarget {
pub change_id: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
#[serde(deny_unknown_fields)]
pub struct WorktreeTarget {
pub worktree_id: String,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
#[serde(deny_unknown_fields)]
pub struct EmptyParams {}
impl CommandSpec {
pub fn type_name(&self) -> &'static str {
match self {
Self::Start => "start",
Self::Stop => "stop",
Self::CancelStop => "cancel_stop",
Self::ForceStop => "force_stop",
Self::SetExecutionMark { .. } => "set_execution_mark",
Self::SetQueueIntent { .. } => "set_queue_intent",
Self::RetryChange { .. } => "retry_change",
Self::RetryErrors { .. } => "retry_errors",
Self::StopAndDequeue { .. } => "stop_and_dequeue",
Self::ForceStopChange { .. } => "force_stop_change",
Self::ResolveMerge { .. } => "resolve_merge",
Self::SetAllExecutionMarks { .. } => "set_all_execution_marks",
Self::CreateWorktree { .. } => "create_worktree",
Self::DeleteWorktree { .. } => "delete_worktree",
Self::MergeWorktree { .. } => "merge_worktree",
}
}
#[cfg_attr(not(test), allow(dead_code))]
pub fn target(&self) -> Option<&str> {
match self {
Self::SetExecutionMark { change_id, .. }
| Self::SetQueueIntent { change_id, .. }
| Self::RetryChange { change_id }
| Self::StopAndDequeue { change_id }
| Self::ForceStopChange { change_id }
| Self::ResolveMerge { change_id } => Some(change_id),
Self::CreateWorktree { target, .. } => Some(&target.change_id),
Self::Start | Self::Stop | Self::CancelStop | Self::ForceStop => None,
Self::RetryErrors { .. } => None,
Self::SetAllExecutionMarks { .. } => None,
Self::DeleteWorktree { .. } | Self::MergeWorktree { .. } => None,
}
}
}
pub const SUPPORTED_COMMANDS: [&str; 15] = [
"start",
"stop",
"cancel_stop",
"force_stop",
"set_execution_mark",
"set_queue_intent",
"retry_change",
"retry_errors",
"stop_and_dequeue",
"force_stop_change",
"resolve_merge",
"set_all_execution_marks",
"create_worktree",
"delete_worktree",
"merge_worktree",
];
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
pub struct CommandRequest {
#[serde(flatten)]
pub command: CommandSpec,
pub expected_revision: u64,
pub idempotency_key: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub correlation_id: Option<String>,
}
impl CommandRequest {
pub fn identity(&self) -> CommandIdentity {
CommandIdentity {
command: self.command.clone(),
expected_revision: self.expected_revision,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CommandIdentity {
pub command: CommandSpec,
pub expected_revision: u64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
#[serde(rename_all = "snake_case")]
pub enum CommandState {
Running,
Succeeded,
NoOp,
Failed,
}
impl CommandState {
pub fn is_in_progress(self) -> bool {
matches!(self, Self::Running)
}
}
#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
pub struct CommandRecord {
pub command_id: String,
pub instance_id: String,
#[serde(rename = "type")]
pub command_type: String,
pub state: CommandState,
pub expected_revision: u64,
#[serde(skip_serializing_if = "Option::is_none")]
pub result_revision: Option<u64>,
pub correlation_id: String,
pub idempotency_key: String,
pub created_at: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub completed_at: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub detail: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub error_code: Option<ErrorCode>,
#[serde(skip_serializing_if = "Option::is_none")]
pub result: Option<CommandResult>,
}
impl CommandRecord {
pub fn http_status(&self) -> StatusCode {
match self.state {
CommandState::Running => StatusCode::ACCEPTED,
CommandState::Succeeded | CommandState::NoOp => StatusCode::OK,
CommandState::Failed => self
.error_code
.map(ErrorCode::http_status)
.unwrap_or(StatusCode::INTERNAL_SERVER_ERROR),
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
pub struct HealthResponse {
pub status: String,
pub api_version: String,
pub version: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
pub struct CapabilityLimits {
pub max_events: usize,
pub max_logs: usize,
pub max_commands: usize,
pub max_idempotency_records: usize,
pub command_record_ttl_secs: i64,
pub max_correlation_id_len: usize,
}
#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
pub struct CapabilitiesResponse {
pub api_version: String,
pub instance_id: String,
pub commands: Vec<String>,
pub transports: Vec<TransportDescriptor>,
pub error_codes: Vec<String>,
pub limits: CapabilityLimits,
pub authentication_required: bool,
pub command_execution: CommandExecutionCapability,
#[serde(default)]
pub execution_sinks: ExecutionSinkCapability,
#[serde(default)]
pub proposal_subscriptions: ProposalSubscriptionCapability,
pub worktrees: super::worktrees::WorktreeCapabilities,
pub parallel: ParallelCapabilities,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
pub struct CommandExecutionCapability {
pub available: bool,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, ToSchema)]
pub struct ExecutionSinkCapability {
pub available: bool,
pub max_command_args: usize,
pub max_command_arg_len: usize,
pub callback_timeout_ms: u64,
pub max_callback_output_bytes: usize,
}
#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
pub struct TransportDescriptor {
pub name: String,
pub path: String,
pub client: String,
pub browser_native_supported: bool,
}
#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
pub struct InstanceResponse {
pub instance_id: String,
pub started_at: String,
pub pid: u32,
pub version: String,
pub api_version: String,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
#[serde(rename_all = "snake_case")]
pub enum QueueIntent {
#[default]
NotQueued,
Queued,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
#[serde(rename_all = "snake_case")]
pub enum AttentionState {
#[default]
None,
New,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
#[serde(rename_all = "snake_case")]
pub enum BlockerKind {
#[default]
None,
Dependency,
External,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
pub struct ChangeBlocker {
pub status: String,
pub kind: BlockerKind,
pub category: Option<String>,
pub detail: Option<String>,
pub unblock_condition: Option<String>,
pub prerequisite_owner: Option<String>,
pub origin: Option<String>,
pub resumable: bool,
#[serde(default)]
pub dependencies: Vec<String>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
#[serde(rename_all = "snake_case")]
pub enum ActionBlockedReason {
FinalStatus,
RetryRequired,
StopPending,
StatusImmutable,
ModeHasNoQueue,
NoRetryableEvidence,
HoldNotResumable,
ChangeActive,
NotMergeWaiting,
ParallelIneligible,
NotAdmitted,
NoManagedProcess,
ApplyIterationLimitActive,
}
impl ActionBlockedReason {
pub fn from_force_stop_exclusion(
exclusion: crate::orchestration::operator_command::ForceStopExclusion,
) -> Self {
use crate::orchestration::operator_command::ForceStopExclusion as E;
match exclusion {
E::TerminalTarget => Self::FinalStatus,
E::UnknownTarget | E::NotAdmitted => Self::NotAdmitted,
E::MergeWait | E::ResolveWait | E::NoLiveProcess => Self::NoManagedProcess,
}
}
pub fn from_mark_exclusion(
exclusion: crate::orchestration::operator_command::MarkExclusion,
) -> Self {
use crate::orchestration::operator_command::MarkExclusion as E;
match exclusion {
E::ArchiveComplete | E::FinalStatus => Self::FinalStatus,
E::RetryRequired => Self::RetryRequired,
E::StopPending => Self::StopPending,
E::ChangeActive => Self::ChangeActive,
E::StatusImmutable => Self::StatusImmutable,
E::ParallelIneligible | E::ParallelProposalAbsent => Self::ParallelIneligible,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
pub struct ActionEligibility {
pub allowed: bool,
pub blocked_reason: Option<ActionBlockedReason>,
}
impl ActionEligibility {
pub fn allowed() -> Self {
Self {
allowed: true,
blocked_reason: None,
}
}
pub fn blocked(reason: ActionBlockedReason) -> Self {
Self {
allowed: false,
blocked_reason: Some(reason),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
pub struct ChangeActions {
pub set_execution_mark: ActionEligibility,
pub set_queue_intent: ActionEligibility,
pub retry_change: ActionEligibility,
pub stop_and_dequeue: ActionEligibility,
pub force_stop_change: ActionEligibility,
pub resolve_merge: ActionEligibility,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
#[serde(rename_all = "snake_case")]
pub enum ParallelBlockedReason {
NotCommitted,
UncommittedChanges,
DependencyBlocked,
}
pub const ALL_PARALLEL_BLOCKED_REASONS: [&str; 3] =
["not_committed", "uncommitted_changes", "dependency_blocked"];
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
pub struct ParallelRuntimeState {
pub max_concurrent: usize,
pub vcs_backend: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
pub struct ParallelCapabilities {
pub max_concurrent: usize,
pub vcs_backend: String,
pub blocked_reasons: Vec<String>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
pub struct ParallelEligibility {
pub eligible: bool,
pub blocked_reason: Option<ParallelBlockedReason>,
}
impl Default for ParallelEligibility {
fn default() -> Self {
Self {
eligible: true,
blocked_reason: None,
}
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
pub struct ChangeTiming {
pub started_at: Option<String>,
pub completed_at: Option<String>,
pub elapsed_ms: Option<u64>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
pub struct ChangeActivity {
pub event_type: String,
pub timestamp: String,
pub detail: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
pub struct ChangeWorktree {
pub path: String,
pub branch: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, ToSchema)]
pub struct ChangeResource {
pub id: String,
pub display_status: String,
pub progress_status: String,
pub completed_tasks: u32,
pub total_tasks: u32,
pub progress_percent: f32,
pub dependencies: Vec<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub iteration_number: Option<u32>,
pub execution_marked: bool,
pub queue_intent: QueueIntent,
pub attention: AttentionState,
pub blocker: Option<ChangeBlocker>,
pub error_detail: Option<String>,
pub actions: ChangeActions,
pub parallel: ParallelEligibility,
pub timing: ChangeTiming,
pub latest_activity: Option<ChangeActivity>,
pub worktree: Option<ChangeWorktree>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
pub struct SnapshotTotals {
pub total: usize,
pub completed: usize,
pub in_progress: usize,
pub pending: usize,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, ToSchema)]
pub struct InstanceSnapshot {
pub app_mode: String,
#[serde(default)]
pub persistent_scheduler_idle: bool,
pub is_resolving: bool,
pub process_error: Option<String>,
pub parallel: ParallelRuntimeState,
pub changes: Vec<ChangeResource>,
pub totals: SnapshotTotals,
}
impl InstanceSnapshot {
pub fn empty() -> Self {
Self {
app_mode: "select".to_string(),
persistent_scheduler_idle: false,
is_resolving: false,
process_error: None,
parallel: ParallelRuntimeState::default(),
changes: Vec::new(),
totals: SnapshotTotals {
total: 0,
completed: 0,
in_progress: 0,
pending: 0,
},
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
pub struct StateResponse {
pub instance_id: String,
pub state_revision: u64,
pub event_sequence: u64,
pub snapshot: InstanceSnapshot,
}
#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
pub struct ChangesResponse {
pub instance_id: String,
pub state_revision: u64,
pub changes: Vec<ChangeResource>,
}
#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
pub struct ChangeResponse {
pub instance_id: String,
pub state_revision: u64,
pub change: ChangeResource,
}
#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
pub struct LogsResponse {
pub instance_id: String,
pub state_revision: u64,
pub event_sequence: u64,
pub logs: Vec<crate::events::LogEntry>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
#[serde(rename_all = "snake_case")]
pub enum ExecutionPhase {
Preparing,
Apply,
Acceptance,
RejectionReview,
Archive,
Resolve,
Push,
Merge,
None,
Unknown,
}
impl ExecutionPhase {
pub fn from_shared(phase: crate::orchestration::execution_facts::ExecutionPhase) -> Self {
use crate::orchestration::execution_facts::ExecutionPhase as P;
match phase {
P::Preparing => Self::Preparing,
P::Apply => Self::Apply,
P::Acceptance => Self::Acceptance,
P::RejectionReview => Self::RejectionReview,
P::Archive => Self::Archive,
P::Resolve => Self::Resolve,
P::Push => Self::Push,
P::Merge => Self::Merge,
P::None => Self::None,
P::Unknown => Self::Unknown,
}
}
}
#[allow(dead_code)] pub const ALL_EXECUTION_PHASES: [&str; 10] = [
"preparing",
"apply",
"acceptance",
"rejection_review",
"archive",
"resolve",
"push",
"merge",
"none",
"unknown",
];
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
#[serde(rename_all = "snake_case")]
pub enum ChangeExecutionState {
Queued,
Active,
Waiting,
Stopping,
Stopped,
Failed,
Completed,
Unknown,
}
impl ChangeExecutionState {
pub fn from_shared(state: crate::orchestration::execution_facts::ChangeExecutionState) -> Self {
use crate::orchestration::execution_facts::ChangeExecutionState as S;
match state {
S::Queued => Self::Queued,
S::Active => Self::Active,
S::Waiting => Self::Waiting,
S::Stopping => Self::Stopping,
S::Stopped => Self::Stopped,
S::Failed => Self::Failed,
S::Completed => Self::Completed,
S::Unknown => Self::Unknown,
}
}
}
#[allow(dead_code)] pub const ALL_CHANGE_EXECUTION_STATES: [&str; 8] = [
"queued",
"active",
"waiting",
"stopping",
"stopped",
"failed",
"completed",
"unknown",
];
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
pub struct LatestLogProjection {
pub message: String,
pub level: crate::events::LogLevel,
pub operation: Option<String>,
pub iteration: Option<u32>,
pub created_at: String,
}
impl LatestLogProjection {
pub fn from_entry(entry: &crate::events::LogEntry) -> Self {
Self {
message: entry.message.clone(),
level: entry.level,
operation: entry.operation.clone(),
iteration: entry.iteration,
created_at: entry.created_at.to_rfc3339(),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
pub struct ProcessExecutionStatus {
pub app_mode: String,
pub scheduler_running: bool,
pub has_active_work: bool,
pub active_activities: Vec<String>,
pub latest_log: Option<LatestLogProjection>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
pub struct ChangeExecutionStatus {
pub id: String,
#[serde(default)]
pub execution_id: Option<String>,
pub execution_state: ChangeExecutionState,
pub current_phase: ExecutionPhase,
pub last_completed_phase: Option<ExecutionPhase>,
pub iteration: Option<u32>,
pub phase_started_at: Option<String>,
pub last_completed_at: Option<String>,
pub run_started_at: Option<String>,
pub run_completed_at: Option<String>,
pub latest_activity: Option<ChangeActivity>,
pub latest_log: Option<LatestLogProjection>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
pub struct ExecutionStatusResponse {
pub instance_id: String,
pub state_revision: u64,
pub event_sequence: u64,
pub observed_at: String,
pub process: ProcessExecutionStatus,
pub changes: Vec<ChangeExecutionStatus>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
#[serde(rename_all = "snake_case")]
pub enum ExecutionEventType {
Completed,
Failed,
Stopped,
Blocked,
OwnerStopping,
}
impl ExecutionEventType {
pub fn as_str(self) -> &'static str {
match self {
Self::Completed => "completed",
Self::Failed => "failed",
Self::Stopped => "stopped",
Self::Blocked => "blocked",
Self::OwnerStopping => "owner_stopping",
}
}
pub fn is_terminal(self) -> bool {
matches!(self, Self::Completed | Self::Failed | Self::Stopped)
}
}
#[allow(dead_code)] pub const ALL_EXECUTION_EVENT_TYPES: [ExecutionEventType; 5] = [
ExecutionEventType::Completed,
ExecutionEventType::Failed,
ExecutionEventType::Stopped,
ExecutionEventType::Blocked,
ExecutionEventType::OwnerStopping,
];
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
pub struct ExecutionSinkSpec {
pub command: Vec<String>,
pub notify_blocked: bool,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
#[serde(deny_unknown_fields)]
pub struct ExecutionSinkRequest {
pub instance_id: String,
pub change_id: String,
pub command: Vec<String>,
#[serde(default)]
pub notify_blocked: bool,
}
#[derive(Debug, Clone, Default, Deserialize)]
pub struct ExecutionSinkParams {
pub instance_id: Option<String>,
pub change_id: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
pub struct ExecutionSinkResponse {
pub instance_id: String,
pub execution_id: String,
pub change_id: String,
pub sink: Option<ExecutionSinkSpec>,
#[serde(default)]
pub sink_registered: bool,
pub execution_state: ChangeExecutionState,
pub terminal_dispatched: bool,
pub delivered_events: Vec<ExecutionEventType>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
#[serde(deny_unknown_fields)]
pub struct ProposalSubscriptionRequest {
pub instance_id: String,
pub command: Vec<String>,
#[serde(default)]
pub notify_blocked: bool,
}
#[derive(Debug, Clone, Default, Deserialize)]
pub struct ProposalSubscriptionParams {
pub instance_id: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
pub struct ProposalSubscriptionResponse {
pub instance_id: String,
pub change_id: String,
pub sink: Option<ExecutionSinkSpec>,
#[serde(default)]
pub subscribed: bool,
#[serde(skip_serializing_if = "Option::is_none")]
pub execution_id: Option<String>,
pub execution_state: ChangeExecutionState,
pub terminal_dispatched: bool,
pub delivered_events: Vec<ExecutionEventType>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, ToSchema)]
pub struct ProposalSubscriptionCapability {
pub available: bool,
pub max_targets: usize,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ExecutionEventFile {
pub schema_version: u32,
pub event_type: ExecutionEventType,
pub instance_id: String,
pub execution_id: String,
pub change_id: String,
pub emitted_at: String,
pub terminal: bool,
#[serde(skip_serializing_if = "Option::is_none")]
pub terminal_mode: Option<TerminalMode>,
#[serde(skip_serializing_if = "Option::is_none")]
pub evidence: Option<String>,
}
pub const EXECUTION_EVENT_SCHEMA_VERSION: u32 = 1;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
#[serde(rename_all = "snake_case")]
pub enum TerminalMode {
Merged,
BasePublished,
BranchPushed,
}
impl TerminalMode {
#[allow(dead_code)] pub fn as_str(self) -> &'static str {
match self {
Self::Merged => "merged",
Self::BasePublished => "base_published",
Self::BranchPushed => "branch_pushed",
}
}
}
#[allow(dead_code)] pub const ALL_TERMINAL_MODES: [&str; 3] = ["merged", "base_published", "branch_pushed"];
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
pub struct OwnerExecutionContract {
pub base_branch: String,
pub terminal_mode: TerminalMode,
#[serde(skip_serializing_if = "Option::is_none")]
pub remote: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub pushed_branch: Option<String>,
}
impl OwnerExecutionContract {
pub fn resolve(
base_branch: impl Into<String>,
push_remote: Option<&str>,
upstream_remote: Option<&str>,
) -> Self {
let (terminal_mode, remote) = match (upstream_remote, push_remote) {
(Some(remote), _) => (TerminalMode::BasePublished, Some(remote.to_string())),
(None, Some(remote)) => (TerminalMode::BranchPushed, Some(remote.to_string())),
(None, None) => (TerminalMode::Merged, None),
};
Self {
base_branch: base_branch.into(),
terminal_mode,
remote,
pushed_branch: None,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
pub struct ExecutionContractResponse {
pub instance_id: String,
pub state_revision: u64,
pub contract: Option<OwnerExecutionContract>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
pub struct ApplyCommitEvidence {
pub present: Option<bool>,
pub oid: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
#[serde(tag = "kind", rename_all = "snake_case")]
pub enum CommandResult {
StopAndDequeue {
cancelled_phase: ExecutionPhase,
last_completed_phase: Option<ExecutionPhase>,
apply_commit: ApplyCommitEvidence,
effects_rolled_back: bool,
},
ForceStopChange {
change_id: String,
execution_id: Option<String>,
cancelled_phase: ExecutionPhase,
last_completed_phase: Option<ExecutionPhase>,
terminated: bool,
apply_commit: ApplyCommitEvidence,
effects_rolled_back: bool,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
#[serde(rename_all = "snake_case")]
pub enum EventCategory {
State,
Log,
Gap,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, ToSchema)]
pub struct EventEnvelope {
pub instance_id: String,
pub event_sequence: u64,
pub state_revision: u64,
pub category: EventCategory,
pub event_type: String,
pub timestamp: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub change_id: Option<String>,
pub payload: serde_json::Value,
}