use std::collections::VecDeque;
use std::sync::{Mutex, MutexGuard};
use chrono::{DateTime, Utc};
use serde_json::json;
use tokio::sync::broadcast;
use crate::events::{ExecutionEvent, LogEntry};
use crate::orchestration::operator_command::{
classify_force_stop_change, classify_queue_intent_route, classify_retry_route,
is_active_status, is_final_status, is_markable_status, ForceStopAdmission, MarkExclusion,
OperatorMode, QueueIntentRoute, RetryRoute,
};
use crate::orchestration::state::BlockerKind;
use crate::web::state::OrchestratorStateSnapshot;
use super::dto::{
new_hex_id, ActionBlockedReason, ActionEligibility, BlockerKind as DtoBlockerKind,
ChangeActions, ChangeBlocker, ChangeResource, CommandRecord, CommandRequest, CommandState,
EventCategory, EventEnvelope, InstanceSnapshot, SnapshotTotals, MAX_EVENTS, MAX_LOGS,
};
use super::registry::{CommandOutcome, CommandRegistry, IdempotencyLookup, ReserveError};
const EVENT_CHANNEL_CAPACITY: usize = 512;
#[derive(Debug, Clone, PartialEq)]
pub enum EventsSince {
Replay(Vec<EventEnvelope>),
Gap,
}
#[derive(Debug, Clone)]
pub enum Admission {
Admitted(Box<CommandRecord>),
Replay(Box<CommandRecord>),
IdempotencyMismatch,
Stale(u64),
Capacity,
}
#[derive(Debug, Clone)]
pub struct ExecutionObservation {
pub snapshot: InstanceSnapshot,
pub state_revision: u64,
pub event_sequence: u64,
pub process_log: Option<LogEntry>,
pub change_logs: std::collections::HashMap<String, LogEntry>,
}
struct Inner {
state_revision: u64,
event_sequence: u64,
snapshot: InstanceSnapshot,
events: VecDeque<EventEnvelope>,
logs: VecDeque<LogEntry>,
registry: CommandRegistry,
}
pub struct Projection {
instance_id: String,
started_at: String,
inner: Mutex<Inner>,
tx: broadcast::Sender<EventEnvelope>,
}
impl Default for Projection {
fn default() -> Self {
Self::new()
}
}
impl Projection {
pub fn new() -> Self {
Self::with_registry(CommandRegistry::default())
}
pub fn with_registry(registry: CommandRegistry) -> Self {
let (tx, _) = broadcast::channel(EVENT_CHANNEL_CAPACITY);
Self {
instance_id: new_hex_id(),
started_at: Utc::now().to_rfc3339(),
inner: Mutex::new(Inner {
state_revision: 0,
event_sequence: 0,
snapshot: InstanceSnapshot::empty(),
events: VecDeque::new(),
logs: VecDeque::new(),
registry,
}),
tx,
}
}
pub fn instance_id(&self) -> &str {
&self.instance_id
}
pub fn started_at(&self) -> &str {
&self.started_at
}
pub fn subscribe(&self) -> broadcast::Receiver<EventEnvelope> {
self.tx.subscribe()
}
fn lock(&self) -> MutexGuard<'_, Inner> {
self.inner
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
}
pub fn snapshot(&self) -> (InstanceSnapshot, u64, u64) {
let inner = self.lock();
(
inner.snapshot.clone(),
inner.state_revision,
inner.event_sequence,
)
}
pub fn revision(&self) -> u64 {
self.lock().state_revision
}
pub fn logs(&self) -> (Vec<LogEntry>, u64, u64) {
let inner = self.lock();
(
inner.logs.iter().cloned().collect(),
inner.state_revision,
inner.event_sequence,
)
}
pub fn execution_observation(&self) -> ExecutionObservation {
let inner = self.lock();
let wanted: std::collections::HashSet<&str> = inner
.snapshot
.changes
.iter()
.map(|change| change.id.as_str())
.collect();
let mut change_logs: std::collections::HashMap<String, LogEntry> =
std::collections::HashMap::new();
for entry in inner.logs.iter().rev() {
let Some(change_id) = entry.change_id.as_deref() else {
continue;
};
if !wanted.contains(change_id) || change_logs.contains_key(change_id) {
continue;
}
change_logs.insert(change_id.to_string(), entry.clone());
if change_logs.len() == wanted.len() {
break;
}
}
ExecutionObservation {
snapshot: inner.snapshot.clone(),
state_revision: inner.state_revision,
event_sequence: inner.event_sequence,
process_log: inner.logs.back().cloned(),
change_logs,
}
}
pub fn events_after(&self, after: u64) -> EventsSince {
let inner = self.lock();
if after > inner.event_sequence {
return EventsSince::Gap;
}
if after == inner.event_sequence {
return EventsSince::Replay(Vec::new());
}
match inner.events.front() {
Some(oldest) if oldest.event_sequence <= after + 1 => EventsSince::Replay(
inner
.events
.iter()
.filter(|event| event.event_sequence > after)
.cloned()
.collect(),
),
_ => EventsSince::Gap,
}
}
pub fn gap_envelope(&self, requested_after: u64) -> EventEnvelope {
let inner = self.lock();
EventEnvelope {
instance_id: self.instance_id.clone(),
event_sequence: inner.event_sequence,
state_revision: inner.state_revision,
category: EventCategory::Gap,
event_type: "replay_gap".to_string(),
timestamp: Utc::now().to_rfc3339(),
change_id: None,
payload: json!({
"requested_after": requested_after,
"oldest_retained": inner.events.front().map(|e| e.event_sequence),
"recover_with": "GET /api/v2/state",
}),
}
}
pub fn apply_state(
&self,
event_type: &str,
change_id: Option<String>,
payload: serde_json::Value,
candidate: InstanceSnapshot,
) -> EventEnvelope {
let mut inner = self.lock();
if inner.snapshot != candidate {
inner.state_revision += 1;
inner.snapshot = candidate;
}
inner.event_sequence += 1;
let envelope = EventEnvelope {
instance_id: self.instance_id.clone(),
event_sequence: inner.event_sequence,
state_revision: inner.state_revision,
category: EventCategory::State,
event_type: event_type.to_string(),
timestamp: Utc::now().to_rfc3339(),
change_id,
payload,
};
self.store_and_publish(&mut inner, envelope.clone());
envelope
}
pub fn apply_state_if_changed(
&self,
event_type: &str,
candidate: InstanceSnapshot,
) -> Option<EventEnvelope> {
{
let inner = self.lock();
if inner.snapshot == candidate {
return None;
}
}
Some(self.apply_state(event_type, None, json!({}), candidate))
}
pub fn apply_presentation(
&self,
event_type: &str,
change_id: Option<String>,
payload: serde_json::Value,
) -> EventEnvelope {
let mut inner = self.lock();
inner.event_sequence += 1;
let envelope = EventEnvelope {
instance_id: self.instance_id.clone(),
event_sequence: inner.event_sequence,
state_revision: inner.state_revision,
category: EventCategory::State,
event_type: event_type.to_string(),
timestamp: Utc::now().to_rfc3339(),
change_id,
payload,
};
self.store_and_publish(&mut inner, envelope.clone());
envelope
}
pub fn apply_log(&self, log: LogEntry) -> EventEnvelope {
let mut inner = self.lock();
inner.event_sequence += 1;
let envelope = EventEnvelope {
instance_id: self.instance_id.clone(),
event_sequence: inner.event_sequence,
state_revision: inner.state_revision,
category: EventCategory::Log,
event_type: "log".to_string(),
timestamp: Utc::now().to_rfc3339(),
change_id: log.change_id.clone(),
payload: serde_json::to_value(&log).unwrap_or_else(|_| json!({})),
};
inner.logs.push_back(log);
while inner.logs.len() > MAX_LOGS {
inner.logs.pop_front();
}
self.store_and_publish(&mut inner, envelope.clone());
envelope
}
fn store_and_publish(&self, inner: &mut Inner, envelope: EventEnvelope) {
inner.events.push_back(envelope.clone());
while inner.events.len() > MAX_EVENTS {
inner.events.pop_front();
}
let _ = self.tx.send(envelope);
}
pub fn admit(
&self,
request: &CommandRequest,
correlation_id: &str,
now: DateTime<Utc>,
) -> Admission {
let identity = request.identity();
let mut inner = self.lock();
match inner
.registry
.lookup(&request.idempotency_key, &identity, now)
{
IdempotencyLookup::Replay(record) => return Admission::Replay(record),
IdempotencyLookup::Mismatch => return Admission::IdempotencyMismatch,
IdempotencyLookup::Unknown => {}
}
self.reserve_admitted(&mut inner, request, &identity, correlation_id, now)
}
pub fn resolve_replay(
&self,
request: &CommandRequest,
now: DateTime<Utc>,
) -> Option<Admission> {
let identity = request.identity();
let mut inner = self.lock();
match inner
.registry
.lookup(&request.idempotency_key, &identity, now)
{
IdempotencyLookup::Replay(record) => Some(Admission::Replay(record)),
IdempotencyLookup::Mismatch => Some(Admission::IdempotencyMismatch),
IdempotencyLookup::Unknown => None,
}
}
fn reserve_admitted(
&self,
inner: &mut MutexGuard<'_, Inner>,
request: &CommandRequest,
identity: &super::dto::CommandIdentity,
correlation_id: &str,
now: DateTime<Utc>,
) -> Admission {
if request.expected_revision != inner.state_revision {
return Admission::Stale(inner.state_revision);
}
let record = CommandRecord {
command_id: new_hex_id(),
instance_id: self.instance_id.clone(),
command_type: request.command.type_name().to_string(),
state: CommandState::Running,
expected_revision: request.expected_revision,
result_revision: None,
correlation_id: correlation_id.to_string(),
idempotency_key: request.idempotency_key.clone(),
created_at: now.to_rfc3339(),
completed_at: None,
detail: None,
error_code: None,
result: None,
};
match inner.registry.reserve(
&request.idempotency_key,
identity.clone(),
record.clone(),
now,
) {
Ok(()) => Admission::Admitted(Box::new(record)),
Err(ReserveError::Capacity) => Admission::Capacity,
}
}
pub fn complete_command(
&self,
command_id: &str,
state: CommandState,
detail: Option<String>,
error_code: Option<super::dto::ErrorCode>,
result_revision: Option<u64>,
result: Option<super::dto::CommandResult>,
) -> Option<CommandRecord> {
let mut inner = self.lock();
let admitted = inner
.registry
.get(command_id)
.map(|record| record.expected_revision);
let result_revision =
result_revision.unwrap_or_else(|| admitted.unwrap_or(inner.state_revision));
inner.registry.complete(
command_id,
CommandOutcome {
state,
result_revision,
detail,
error_code,
result,
},
)
}
pub fn command(&self, command_id: &str) -> Option<CommandRecord> {
self.lock().registry.get(command_id).cloned()
}
#[cfg_attr(not(test), allow(dead_code))]
pub fn registry_sizes(&self) -> (usize, usize) {
let inner = self.lock();
(
inner.registry.command_len(),
inner.registry.idempotency_len(),
)
}
#[cfg(test)]
pub fn reserve_for_test(
&self,
key: &str,
identity: super::dto::CommandIdentity,
record: CommandRecord,
now: DateTime<Utc>,
) -> Result<(), ReserveError> {
self.lock().registry.reserve(key, identity, record, now)
}
}
pub fn project_snapshot(source: &OrchestratorStateSnapshot) -> InstanceSnapshot {
let mode = OperatorMode::from_app_mode(&source.app_mode);
let changes: Vec<ChangeResource> = source
.changes
.iter()
.map(|change| {
let display_status = change
.queue_status
.clone()
.unwrap_or_else(|| "not queued".to_string());
ChangeResource {
id: change.id.clone(),
actions: classify_actions(
mode,
&display_status,
change.blocker.as_ref(),
change.parallel.eligible,
change.managed_process_live,
),
display_status,
progress_status: change.status.clone(),
completed_tasks: change.completed_tasks,
total_tasks: change.total_tasks,
progress_percent: change.progress_percent,
dependencies: change.dependencies.clone(),
iteration_number: change.iteration_number,
execution_marked: change.execution_marked,
queue_intent: change.queue_intent,
attention: change.attention,
blocker: change.blocker.clone(),
error_detail: change.error_detail.clone(),
parallel: change.parallel,
timing: change.timing.clone(),
latest_activity: change.latest_activity.clone(),
worktree: change.worktree.clone(),
}
})
.collect();
InstanceSnapshot {
app_mode: source.app_mode.clone(),
persistent_scheduler_idle: source.persistent_scheduler_idle,
is_resolving: source.is_resolving,
process_error: source.process_error.clone(),
parallel: source.parallel.clone(),
totals: SnapshotTotals {
total: changes.len(),
completed: source.completed_changes,
in_progress: source.in_progress_changes,
pending: source.pending_changes,
},
changes,
}
}
fn classify_actions(
mode: OperatorMode,
display_status: &str,
blocker: Option<&ChangeBlocker>,
parallel_eligible: bool,
managed_process_live: bool,
) -> ChangeActions {
let final_status = is_final_status(display_status);
let route = classify_queue_intent_route(mode, display_status);
let immutable_reason = if final_status {
ActionBlockedReason::FinalStatus
} else if matches!(mode, OperatorMode::Stopping) {
ActionBlockedReason::StopPending
} else {
ActionBlockedReason::StatusImmutable
};
let parallel_blocked = !parallel_eligible && !final_status;
let parallel_reason =
ActionBlockedReason::from_mark_exclusion(MarkExclusion::ParallelIneligible);
let set_execution_mark = if is_markable_status(display_status, false) {
ActionEligibility::allowed()
} else {
ActionEligibility::blocked(ActionBlockedReason::FinalStatus)
};
let set_queue_intent = match route {
_ if parallel_blocked => ActionEligibility::blocked(parallel_reason),
QueueIntentRoute::Mutable => ActionEligibility::allowed(),
QueueIntentRoute::NoQueueEffect
if matches!(mode, OperatorMode::Select | OperatorMode::Stopped) =>
{
ActionEligibility::blocked(ActionBlockedReason::ModeHasNoQueue)
}
QueueIntentRoute::NoQueueEffect => {
ActionEligibility::blocked(ActionBlockedReason::StatusImmutable)
}
QueueIntentRoute::RetryRequired => {
ActionEligibility::blocked(ActionBlockedReason::RetryRequired)
}
QueueIntentRoute::Immutable => ActionEligibility::blocked(immutable_reason),
};
let retry_change = classify_retry_action(display_status, blocker, final_status);
let stop_and_dequeue = if final_status {
ActionEligibility::blocked(ActionBlockedReason::FinalStatus)
} else {
ActionEligibility::allowed()
};
let resolve_merge = match display_status {
"merge wait" | "resolve pending" | "archived" => ActionEligibility::allowed(),
status if is_active_status(status) => {
ActionEligibility::blocked(ActionBlockedReason::ChangeActive)
}
status if is_final_status(status) => {
ActionEligibility::blocked(ActionBlockedReason::FinalStatus)
}
_ => ActionEligibility::blocked(ActionBlockedReason::NotMergeWaiting),
};
let force_stop_change =
match classify_force_stop_change(display_status, true, managed_process_live) {
ForceStopAdmission::KillAndDequeue | ForceStopAdmission::DequeueOnly => {
ActionEligibility::allowed()
}
ForceStopAdmission::Refused(reason) => {
ActionEligibility::blocked(ActionBlockedReason::from_force_stop_exclusion(reason))
}
};
ChangeActions {
set_execution_mark,
set_queue_intent,
retry_change,
stop_and_dequeue,
force_stop_change,
resolve_merge,
}
}
#[doc(hidden)]
#[allow(dead_code)] pub fn change_actions_for_test(
app_mode: &str,
display_status: &str,
blocker: Option<&ChangeBlocker>,
) -> ChangeActions {
classify_actions(
OperatorMode::from_app_mode(app_mode),
display_status,
blocker,
true,
false,
)
}
#[doc(hidden)]
#[allow(dead_code)] pub fn change_actions_with_live_process_for_test(
app_mode: &str,
display_status: &str,
managed_process_live: bool,
) -> ChangeActions {
classify_actions(
OperatorMode::from_app_mode(app_mode),
display_status,
None,
true,
managed_process_live,
)
}
#[cfg(test)]
pub fn parallel_change_actions_for_test(
app_mode: &str,
display_status: &str,
parallel_eligible: bool,
) -> ChangeActions {
classify_actions(
OperatorMode::from_app_mode(app_mode),
display_status,
None,
parallel_eligible,
false,
)
}
fn classify_retry_action(
display_status: &str,
blocker: Option<&ChangeBlocker>,
final_status: bool,
) -> ActionEligibility {
if final_status {
return ActionEligibility::blocked(ActionBlockedReason::FinalStatus);
}
let blocker_kind = match blocker.map(|blocker| blocker.kind) {
Some(DtoBlockerKind::Dependency) => BlockerKind::Dependency,
Some(DtoBlockerKind::External) => BlockerKind::External,
_ => BlockerKind::None,
};
let Some(route) = classify_retry_route(display_status, blocker_kind) else {
return ActionEligibility::blocked(ActionBlockedReason::NoRetryableEvidence);
};
if matches!(route, RetryRoute::AcceptanceStall)
&& blocker.is_some_and(|blocker| !blocker.resumable)
{
return ActionEligibility::blocked(ActionBlockedReason::HoldNotResumable);
}
ActionEligibility::allowed()
}
pub fn describe_event(event: &ExecutionEvent) -> (&'static str, Option<String>, serde_json::Value) {
use ExecutionEvent as E;
fn detail(text: &str) -> serde_json::Value {
json!(crate::events::sanitize_detail(text))
}
fn optional_detail(text: &Option<String>) -> serde_json::Value {
match text {
Some(text) => detail(text),
None => serde_json::Value::Null,
}
}
fn command_facts(command: &str) -> serde_json::Value {
json!(crate::events::command_log_summary(command))
}
match event {
E::ProcessingStarted(id) => ("processing_started", Some(id.clone()), json!({})),
E::ProcessingError { id, error } => (
"processing_error",
Some(id.clone()),
json!({ "detail": detail(error) }),
),
E::ApplyStarted { change_id, command } => (
"apply_started",
Some(change_id.clone()),
json!({ "command_summary": command_facts(command) }),
),
E::ApplyCompleted {
change_id,
revision,
} => (
"apply_completed",
Some(change_id.clone()),
json!({ "revision": revision }),
),
E::ApplyFailed { change_id, error } => (
"apply_failed",
Some(change_id.clone()),
json!({ "detail": detail(error) }),
),
E::ApplyOutput {
change_id,
iteration,
..
} => (
"apply_output",
Some(change_id.clone()),
json!({ "iteration": iteration }),
),
E::ApplyCommitPhase {
change_id,
phase,
attempt,
} => (
"apply_commit_phase",
Some(change_id.clone()),
json!({ "phase": phase.as_str(), "attempt": attempt }),
),
E::ApplyCommitOutput {
change_id,
attempt,
stream,
..
} => (
"apply_commit_output",
None,
json!({
"change_id": change_id,
"attempt": attempt,
"stream": stream.as_str(),
}),
),
E::ArchiveStarted { change_id, command } => (
"archive_started",
Some(change_id.clone()),
json!({ "command_summary": command_facts(command) }),
),
E::ArchiveResumed {
change_id,
reason,
summary,
} => (
"archive_resumed",
Some(change_id.clone()),
json!({ "detail": optional_detail(reason), "summary": optional_detail(summary) }),
),
E::ArchiveRetryScheduled {
change_id,
attempt,
max_attempts,
reason,
summary,
} => (
"archive_retry_scheduled",
Some(change_id.clone()),
json!({
"detail": optional_detail(reason),
"summary": optional_detail(summary),
"attempt": attempt,
"max_attempts": max_attempts,
}),
),
E::ChangeArchived(id) => ("change_archived", Some(id.clone()), json!({})),
E::ArchiveFailed {
change_id,
error,
reason,
summary,
} => (
"archive_failed",
Some(change_id.clone()),
json!({
"detail": detail(error),
"reason": optional_detail(reason),
"summary": optional_detail(summary),
}),
),
E::ArchiveOutput {
change_id,
iteration,
..
} => (
"archive_output",
Some(change_id.clone()),
json!({ "iteration": iteration }),
),
E::AcceptanceStarted { change_id, command } => (
"acceptance_started",
Some(change_id.clone()),
json!({ "command_summary": command_facts(command) }),
),
E::AcceptanceCompleted { change_id } => {
("acceptance_completed", Some(change_id.clone()), json!({}))
}
E::AcceptanceFailed { change_id, error } => (
"acceptance_failed",
Some(change_id.clone()),
json!({ "detail": detail(error) }),
),
E::ChangeRejected { change_id, reason } => (
"change_rejected",
Some(change_id.clone()),
json!({ "detail": detail(reason) }),
),
E::RejectionReviewCompleted { change_id, outcome } => (
"rejection_review_completed",
Some(change_id.clone()),
json!({ "outcome": rejection_outcome_name(*outcome) }),
),
E::RejectionReviewFailed { change_id, error } => (
"rejection_review_failed",
Some(change_id.clone()),
json!({ "detail": detail(error) }),
),
E::AcceptanceOutput {
change_id,
iteration,
..
} => (
"acceptance_output",
Some(change_id.clone()),
json!({ "iteration": iteration }),
),
E::ProgressUpdated {
change_id,
completed,
total,
} => (
"progress_updated",
Some(change_id.clone()),
json!({ "completed": completed, "total": total }),
),
E::WorkspacePreparationStarted { change_id } => (
"workspace_preparation_started",
Some(change_id.clone()),
json!({}),
),
E::WorkspacePreparationEnded { change_id } => (
"workspace_preparation_ended",
Some(change_id.clone()),
json!({}),
),
E::WorkspaceCreated {
change_id,
workspace,
} => (
"workspace_created",
Some(change_id.clone()),
json!({ "workspace": detail(workspace) }),
),
E::WorkspaceResumed {
change_id,
workspace,
} => (
"workspace_resumed",
Some(change_id.clone()),
json!({ "workspace": detail(workspace) }),
),
E::WorkspacePreserved {
change_id,
workspace_name,
} => (
"workspace_preserved",
Some(change_id.clone()),
json!({ "workspace": detail(workspace_name) }),
),
E::WorkspaceStatusUpdated {
change_id,
workspace_name,
status,
} => (
"workspace_status_updated",
Some(change_id.clone()),
json!({
"workspace": detail(workspace_name),
"status": format!("{status:?}"),
}),
),
E::ResolveStarted { change_id, command } => (
"resolve_started",
Some(change_id.clone()),
json!({ "command_summary": command_facts(command) }),
),
E::ResolveCompleted { change_id, .. } => {
("resolve_completed", Some(change_id.clone()), json!({}))
}
E::ResolveFailed { change_id, error } => (
"resolve_failed",
Some(change_id.clone()),
json!({ "detail": detail(error) }),
),
E::ResolveOutput {
change_id,
iteration,
..
} => (
"resolve_output",
Some(change_id.clone()),
json!({ "iteration": iteration }),
),
E::MergeCompleted {
change_id,
revision,
} => (
"merge_completed",
Some(change_id.clone()),
json!({ "revision": revision }),
),
E::MergeDeferred {
change_id,
reason,
auto_resumable,
} => (
"merge_deferred",
Some(change_id.clone()),
json!({ "detail": detail(reason), "auto_resumable": auto_resumable }),
),
E::PushStarted {
change_id,
remote,
branch,
} => (
"push_started",
Some(change_id.clone()),
json!({ "remote": detail(remote), "branch": detail(branch) }),
),
E::PushCompleted {
change_id,
remote,
branch,
} => (
"push_completed",
Some(change_id.clone()),
json!({ "remote": detail(remote), "branch": detail(branch) }),
),
E::PushFailed {
change_id,
remote,
branch,
error,
} => (
"push_failed",
Some(change_id.clone()),
json!({
"detail": detail(error),
"remote": detail(remote),
"branch": detail(branch),
}),
),
E::ChangeSkipped { change_id, reason } => (
"change_skipped",
Some(change_id.clone()),
json!({ "detail": detail(reason) }),
),
E::DependencyBlocked {
change_id,
dependency_ids,
} => (
"dependency_blocked",
Some(change_id.clone()),
json!({ "dependency_ids": dependency_ids }),
),
E::DependencyResolved { change_id } => {
("dependency_resolved", Some(change_id.clone()), json!({}))
}
E::AcceptanceGated { change_id, blocker } => (
"acceptance_gated",
Some(change_id.clone()),
describe_blocker(blocker),
),
E::ExecutionBlocked { change_id, blocker } => (
"execution_blocked",
Some(change_id.clone()),
describe_blocker(blocker),
),
E::HookStarted {
change_id,
hook_type,
} => (
"hook_started",
Some(change_id.clone()),
json!({ "hook_type": detail(hook_type) }),
),
E::HookCompleted {
change_id,
hook_type,
} => (
"hook_completed",
Some(change_id.clone()),
json!({ "hook_type": detail(hook_type) }),
),
E::HookFailed {
change_id,
hook_type,
error,
} => (
"hook_failed",
Some(change_id.clone()),
json!({ "hook_type": detail(hook_type), "detail": detail(error) }),
),
E::OperatorCommandApplied { effect } => (
"operator_command_applied",
effect.change_id().map(str::to_string),
describe_operator_effect(effect),
),
E::ChangeDequeued { change_id } => ("change_dequeued", Some(change_id.clone()), json!({})),
E::ChangeStopped { change_id } => ("change_stopped", Some(change_id.clone()), json!({})),
E::ChangeStopFailed { change_id, error } => (
"change_stop_failed",
Some(change_id.clone()),
json!({ "detail": detail(error) }),
),
E::Stopping => ("stopping", None, json!({})),
E::Stopped => ("stopped", None, json!({})),
E::AllCompleted => ("all_completed", None, json!({})),
E::PersistentSchedulerIdle => ("persistent_scheduler_idle", None, json!({})),
E::Error { message } => ("process_error", None, json!({ "detail": detail(message) })),
E::Log(entry) => ("log", entry.change_id.clone(), json!({})),
E::ChangesRefreshed {
changes,
rejected_changes,
..
} => (
"changes_refreshed",
None,
json!({ "changes": changes.len(), "rejected_changes": rejected_changes.len() }),
),
E::WorktreesRefreshed { worktrees } => (
"worktrees_refreshed",
None,
json!({ "worktrees": worktrees.len() }),
),
E::CleanupStarted { workspace } => (
"cleanup_started",
None,
json!({ "workspace": detail(workspace) }),
),
E::CleanupCompleted { workspace } => (
"cleanup_completed",
None,
json!({ "workspace": detail(workspace) }),
),
E::MergeStarted { revisions } => ("merge_started", None, json!({ "revisions": revisions })),
E::MergeConflict { files } => ("merge_conflict", None, json!({ "files": files })),
E::ConflictResolutionStarted => ("conflict_resolution_started", None, json!({})),
E::ConflictResolutionCompleted => ("conflict_resolution_completed", None, json!({})),
E::ConflictResolutionFailed { error } => (
"conflict_resolution_failed",
None,
json!({ "detail": detail(error) }),
),
E::AnalysisStarted {
remaining_changes,
attempt_id,
} => (
"analysis_started",
None,
json!({ "remaining_changes": remaining_changes, "attempt_id": attempt_id }),
),
E::AnalysisOutput { iteration, .. } => {
("analysis_output", None, json!({ "iteration": iteration }))
}
E::AnalysisCompleted { groups_found } => (
"analysis_completed",
None,
json!({ "groups_found": groups_found }),
),
E::Warning { title, message } => (
"warning",
None,
json!({ "title": detail(title), "detail": detail(message) }),
),
E::ParallelStartRejected { change_ids, reason } => (
"parallel_start_rejected",
None,
json!({ "change_ids": change_ids, "detail": detail(reason) }),
),
E::BranchMergeStarted { branch_name } => (
"branch_merge_started",
None,
json!({ "branch": detail(branch_name) }),
),
E::BranchMergeCompleted { branch_name } => (
"branch_merge_completed",
None,
json!({ "branch": detail(branch_name) }),
),
E::BranchMergeFailed { branch_name, error } => (
"branch_merge_failed",
None,
json!({ "branch": detail(branch_name), "detail": detail(error) }),
),
}
}
fn describe_operator_effect(effect: &crate::events::OperatorCommandEffect) -> serde_json::Value {
use crate::events::OperatorCommandEffect as Effect;
let mut payload = json!({ "effect": effect.as_str() });
let object = payload
.as_object_mut()
.expect("a json object literal is an object");
match effect {
Effect::RunDispatched {
change_ids,
explicit_retry,
scheduler_started,
} => {
object.insert("change_ids".to_string(), json!(change_ids));
object.insert("explicit_retry".to_string(), json!(explicit_retry));
object.insert("scheduler_started".to_string(), json!(scheduler_started));
}
Effect::StopCancelled => {}
Effect::ForceStopAwaitingBoundary { force_stop } => {
object.insert("force_stop".to_string(), json!(force_stop));
object.insert("awaiting_safe_boundary".to_string(), json!(true));
}
Effect::MarkDelta { change_ids, marked } => {
object.insert("change_ids".to_string(), json!(change_ids));
object.insert("marked".to_string(), json!(marked));
}
Effect::QueueDelta { change_id, queued } => {
object.insert("change_id".to_string(), json!(change_id));
object.insert("queued".to_string(), json!(queued));
}
Effect::ResolveReserved { change_id, active } => {
object.insert("change_id".to_string(), json!(change_id));
object.insert("active".to_string(), json!(active));
}
}
payload
}
fn rejection_outcome_name(outcome: crate::events::RejectionOutcome) -> &'static str {
match outcome {
crate::events::RejectionOutcome::Confirm => "confirm",
crate::events::RejectionOutcome::Resume => "resume",
crate::events::RejectionOutcome::Block => "block",
}
}
fn describe_blocker(blocker: &crate::events::StalledBlocker) -> serde_json::Value {
let sanitize = crate::events::sanitize_detail;
json!({
"category": sanitize(&blocker.category),
"phase": sanitize(&blocker.phase),
"gate": sanitize(&blocker.gate),
"detail": sanitize(&blocker.error_summary),
"evidence": blocker
.evidence
.iter()
.map(|item| sanitize(item))
.collect::<Vec<_>>(),
"unblock_condition": blocker.unblock_condition.as_deref().map(sanitize),
"prerequisite_owner": blocker.prerequisite_owner.as_deref().map(sanitize),
"next_action": sanitize(&blocker.next_action),
"resumable": blocker.resumable,
})
}