Update public backend API surface

Private source commit: ba3f7ce2b6d9
This commit is contained in:
Clusterflux Release 2026-07-26 17:22:33 +02:00
parent 784463622c
commit 26fdcb9d84
34 changed files with 5473 additions and 72 deletions

View file

@ -0,0 +1,500 @@
use std::collections::BTreeSet;
use std::time::Instant;
use clusterflux_core::{
ArtifactId, ArtifactMetadata, NodeId, ProcessId, ProjectId, TenantId, UserId,
};
use super::keys::{process_control_key, ProcessControlKey};
use super::{
ArtifactAvailability, ArtifactRetentionState, ArtifactSummary, CoordinatorResponse,
CoordinatorService, CoordinatorServiceError, DebugAcknowledgementState, DebugEpochSummary,
NodeSummary, ProcessActivityState, ProcessFinalResult, ProcessLifecycleState, ProcessSummary,
TaskAttemptState,
};
#[derive(Clone, Debug, PartialEq, Eq)]
pub(super) struct StoredProcessSummary {
pub(super) started_at_epoch_seconds: u64,
pub(super) ended_at_epoch_seconds: Option<u64>,
pub(super) final_result: Option<ProcessFinalResult>,
pub(super) connected_nodes: Vec<NodeId>,
pub(super) order: u64,
}
impl CoordinatorService {
pub(super) fn handle_list_node_summaries(
&mut self,
tenant: String,
project: String,
actor_user: String,
cursor: Option<String>,
limit: u32,
) -> Result<CoordinatorResponse, CoordinatorServiceError> {
let tenant = TenantId::new(tenant);
let project = ProjectId::new(project);
let actor = UserId::new(actor_user);
let cursor = cursor.as_deref();
let mut nodes = self
.node_descriptors
.iter()
.filter(|(scope, descriptor)| {
scope.tenant == tenant
&& scope.project == project
&& cursor.is_none_or(|cursor| descriptor.id.as_str() > cursor)
})
.map(|(scope, descriptor)| {
let online = self.node_is_live(scope);
NodeSummary {
id: descriptor.id.clone(),
display_name: descriptor.id.as_str().to_owned(),
online,
stale: !online,
last_seen_epoch_seconds: self.node_last_seen_epoch_seconds.get(scope).copied(),
capabilities: descriptor.capabilities.clone(),
direct_connectivity: descriptor.direct_connectivity,
}
})
.collect::<Vec<_>>();
nodes.sort_by(|left, right| left.id.cmp(&right.id));
let has_more = nodes.len() > limit as usize;
nodes.truncate(limit as usize);
let next_cursor = has_more
.then(|| nodes.last().map(|node| node.id.as_str().to_owned()))
.flatten();
Ok(CoordinatorResponse::NodeSummaries {
nodes,
next_cursor,
actor,
})
}
pub(super) fn record_process_started(
&mut self,
tenant: &TenantId,
project: &ProjectId,
process: &ProcessId,
now_epoch_seconds: u64,
) {
let key = process_control_key(tenant, project, process);
if let Some(logs) = self.recent_logs.get_mut(&(tenant.clone(), project.clone())) {
logs.retain(|entry| &entry.process != process);
if logs.is_empty() {
self.recent_logs.remove(&(tenant.clone(), project.clone()));
}
}
self.recent_log_dropped_through.remove(&key);
self.recent_log_offsets
.retain(|(entry_tenant, entry_project, entry_process, _, _), _| {
entry_tenant != tenant || entry_project != project || entry_process != process
});
self.process_summary_order
.retain(|retained| retained != &key);
self.evict_process_summaries_for_project(tenant, project);
self.evict_process_summaries_total();
let order = self.next_process_summary_order;
self.next_process_summary_order = self.next_process_summary_order.saturating_add(1);
self.process_summaries.insert(
key.clone(),
StoredProcessSummary {
started_at_epoch_seconds: now_epoch_seconds,
ended_at_epoch_seconds: None,
final_result: None,
connected_nodes: Vec::new(),
order,
},
);
self.process_summary_order.push_back(key);
}
pub(super) fn record_process_terminal(
&mut self,
tenant: &TenantId,
project: &ProjectId,
process: &ProcessId,
final_result: ProcessFinalResult,
now_epoch_seconds: u64,
) {
let key = process_control_key(tenant, project, process);
let connected_nodes = self
.coordinator
.active_process(tenant, project, process)
.map(|active| active.connected_nodes.iter().cloned().collect())
.unwrap_or_default();
let order = self.next_process_summary_order;
let entry = self
.process_summaries
.entry(key.clone())
.or_insert_with(|| {
self.next_process_summary_order = self.next_process_summary_order.saturating_add(1);
self.process_summary_order.push_back(key.clone());
StoredProcessSummary {
started_at_epoch_seconds: now_epoch_seconds,
ended_at_epoch_seconds: None,
final_result: None,
connected_nodes: Vec::new(),
order,
}
});
entry.ended_at_epoch_seconds = Some(now_epoch_seconds);
entry.final_result = Some(final_result);
entry.connected_nodes = connected_nodes;
}
pub(super) fn handle_list_process_summaries(
&mut self,
tenant: String,
project: String,
actor_user: String,
cursor: Option<String>,
limit: u32,
) -> Result<CoordinatorResponse, CoordinatorServiceError> {
let tenant = TenantId::new(tenant);
let project = ProjectId::new(project);
let actor = UserId::new(actor_user);
let cursor = parse_order_cursor(cursor.as_deref(), "process")?;
let mut stored = self
.process_summaries
.iter()
.filter(|((entry_tenant, entry_project, _), summary)| {
entry_tenant == &tenant
&& entry_project == &project
&& cursor.is_none_or(|cursor| summary.order < cursor)
})
.map(|(key, summary)| (key.clone(), summary.clone()))
.collect::<Vec<_>>();
stored.sort_by(|(_, left), (_, right)| right.order.cmp(&left.order));
let has_more = stored.len() > limit as usize;
stored.truncate(limit as usize);
let processes = stored
.into_iter()
.map(|(key, stored)| self.process_summary_from_stored(&key, stored))
.collect::<Vec<_>>();
let next_cursor = has_more
.then(|| processes.last().map(|process| process.order_cursor.clone()))
.flatten();
Ok(CoordinatorResponse::ProcessSummaries {
processes,
next_cursor,
actor,
})
}
fn process_summary_from_stored(
&self,
key: &ProcessControlKey,
stored: StoredProcessSummary,
) -> ProcessSummary {
let (tenant, project, process) = key;
let active = self.coordinator.active_process(tenant, project, process);
let process_key = process_control_key(tenant, project, process);
let main_wait_state = active.and_then(|_| {
if self.pending_task_launches.iter().any(|pending| {
&pending.tenant == tenant
&& &pending.project == project
&& &pending.process == process
}) {
Some("waiting_for_node".to_owned())
} else if self
.main_runtime
.is_waiting_for_task(tenant, project, process)
{
Some("waiting_for_task".to_owned())
} else {
None
}
});
let current_debug_epoch = self.debug_epoch_summary(&process_key);
let awaiting_action = self.task_attempts.iter().any(
|((attempt_tenant, attempt_project, attempt_process, _), attempts)| {
attempt_tenant == tenant
&& attempt_project == project
&& attempt_process == process
&& attempts.iter().any(|attempt| {
attempt.current && attempt.state == TaskAttemptState::FailedAwaitingAction
})
},
);
let activity = if let Some(result) = &stored.final_result {
match result {
ProcessFinalResult::Completed => ProcessActivityState::Completed,
ProcessFinalResult::Failed => ProcessActivityState::Failed,
ProcessFinalResult::Cancelled => ProcessActivityState::Cancelled,
}
} else if self.process_cancellations.contains(&process_key) {
ProcessActivityState::Cancelling
} else if current_debug_epoch
.as_ref()
.is_some_and(|epoch| epoch.partially_frozen)
{
ProcessActivityState::DebugEpochPartial
} else if awaiting_action {
ProcessActivityState::AwaitingAction
} else {
match main_wait_state.as_deref() {
Some("waiting_for_node") => ProcessActivityState::WaitingForNode,
Some("waiting_for_task") => ProcessActivityState::WaitingForTask,
_ => ProcessActivityState::Running,
}
};
let connected_nodes = active
.map(|active| active.connected_nodes.iter().cloned().collect())
.unwrap_or(stored.connected_nodes);
ProcessSummary {
process: process.clone(),
lifecycle: if active.is_some() {
ProcessLifecycleState::Active
} else {
ProcessLifecycleState::RecentTerminal
},
activity,
main_wait_state,
started_at_epoch_seconds: stored.started_at_epoch_seconds,
ended_at_epoch_seconds: stored.ended_at_epoch_seconds,
final_result: stored.final_result,
connected_nodes,
current_debug_epoch,
order_cursor: format!("process:{}", stored.order),
}
}
fn debug_epoch_summary(&self, key: &ProcessControlKey) -> Option<DebugEpochSummary> {
let runtime = self.debug_epoch_runtime.get(key)?;
let acknowledgements = runtime.acknowledgements.values().collect::<Vec<_>>();
let all_acknowledged = !runtime.expected.is_empty()
&& runtime
.expected
.iter()
.all(|participant| runtime.acknowledgements.contains_key(participant));
let fully_frozen = runtime.command == "freeze"
&& all_acknowledged
&& acknowledgements
.iter()
.all(|ack| ack.state == DebugAcknowledgementState::Frozen);
let freeze_deadline_elapsed =
runtime.command == "freeze" && Instant::now() >= runtime.deadline;
let frozen_count = acknowledgements
.iter()
.filter(|ack| ack.state == DebugAcknowledgementState::Frozen)
.count();
let partially_frozen = freeze_deadline_elapsed && frozen_count > 0 && !fully_frozen;
let fully_resumed = runtime.command == "resume"
&& all_acknowledged
&& acknowledgements
.iter()
.all(|ack| ack.state == DebugAcknowledgementState::Running);
let failed = acknowledgements
.iter()
.any(|ack| ack.state == DebugAcknowledgementState::Failed)
|| (freeze_deadline_elapsed
&& runtime
.expected
.iter()
.any(|participant| !runtime.acknowledgements.contains_key(participant)));
Some(DebugEpochSummary {
epoch: runtime.epoch,
command: runtime.command.clone(),
fully_frozen,
partially_frozen,
fully_resumed,
failed,
})
}
pub(super) fn handle_list_artifacts(
&mut self,
tenant: String,
project: String,
actor_user: String,
process: Option<String>,
cursor: Option<String>,
limit: u32,
) -> Result<CoordinatorResponse, CoordinatorServiceError> {
let tenant = TenantId::new(tenant);
let project = ProjectId::new(project);
let _actor = UserId::new(actor_user);
let process = process.map(ProcessId::new);
if let Some(process) = &process {
self.authorize_task_event_process_scope(&tenant, &project, process)?;
}
let cursor = parse_order_cursor(cursor.as_deref(), "artifact")?;
let mut metadata = self
.artifact_registry
.metadata_for_project(&tenant, &project)
.filter(|metadata| {
process
.as_ref()
.is_none_or(|process| &metadata.process == process)
&& cursor.is_none_or(|cursor| metadata.flushed_epoch < cursor)
})
.cloned()
.collect::<Vec<_>>();
metadata.sort_by(|left, right| right.flushed_epoch.cmp(&left.flushed_epoch));
let has_more = metadata.len() > limit as usize;
metadata.truncate(limit as usize);
let artifacts = metadata
.into_iter()
.map(|metadata| self.artifact_summary(metadata))
.collect::<Vec<_>>();
let next_cursor = has_more
.then(|| {
artifacts
.last()
.map(|artifact| artifact.order_cursor.clone())
})
.flatten();
Ok(CoordinatorResponse::Artifacts {
artifacts,
next_cursor,
})
}
pub(super) fn handle_get_artifact(
&mut self,
tenant: String,
project: String,
actor_user: String,
artifact: String,
) -> Result<CoordinatorResponse, CoordinatorServiceError> {
let tenant = TenantId::new(tenant);
let project = ProjectId::new(project);
let _actor = UserId::new(actor_user);
let artifact = ArtifactId::new(artifact);
let metadata = self
.artifact_registry
.metadata(&tenant, &project, &artifact)
.cloned()
.ok_or(clusterflux_core::DownloadError::NotFound)?;
Ok(CoordinatorResponse::Artifact {
artifact: self.artifact_summary(metadata),
})
}
fn artifact_summary(&self, metadata: ArtifactMetadata) -> ArtifactSummary {
let live_retaining_nodes = metadata
.retaining_nodes
.iter()
.filter(|node| {
self.node_is_live(&crate::NodeScopeKey::from_refs(
&metadata.tenant,
&metadata.project,
node,
))
})
.cloned()
.collect::<BTreeSet<_>>();
let safe_node = live_retaining_nodes
.iter()
.next()
.cloned()
.or_else(|| metadata.retaining_nodes.iter().next().cloned());
let explicit_storage = !metadata.explicit_locations.is_empty();
let downloadable_now = !live_retaining_nodes.is_empty() || explicit_storage;
let availability = if downloadable_now {
ArtifactAvailability::Available
} else if !metadata.retaining_nodes.is_empty() {
ArtifactAvailability::NodeOffline
} else {
ArtifactAvailability::Unavailable
};
let retention_state = if explicit_storage {
ArtifactRetentionState::ExplicitStorage
} else if !metadata.retaining_nodes.is_empty() {
ArtifactRetentionState::NodeRetained
} else {
ArtifactRetentionState::Lost
};
let display_suffix = metadata.id.as_str().replace(':', "/");
let display_path = format!("/vfs/artifacts/{display_suffix}");
let display_name = display_suffix
.rsplit('/')
.next()
.unwrap_or(metadata.id.as_str())
.to_owned();
ArtifactSummary {
id: metadata.id,
display_path,
display_name,
process: metadata.process,
producer_task: metadata.producer_task,
safe_node,
digest: metadata.digest,
size_bytes: metadata.size,
availability,
downloadable_now,
retention_state,
explicit_storage,
order_cursor: format!("artifact:{}", metadata.flushed_epoch),
}
}
fn evict_process_summaries_for_project(&mut self, tenant: &TenantId, project: &ProjectId) {
while self
.process_summaries
.keys()
.filter(|(entry_tenant, entry_project, _)| {
entry_tenant == tenant && entry_project == project
})
.count()
>= super::MAX_RECENT_PROCESS_SUMMARIES_PER_PROJECT
{
let candidate = self.process_summary_order.iter().find(|key| {
&key.0 == tenant
&& &key.1 == project
&& self
.process_summaries
.get(*key)
.is_some_and(|summary| summary.final_result.is_some())
});
let Some(candidate) = candidate.cloned() else {
break;
};
self.process_summaries.remove(&candidate);
self.recent_log_dropped_through.remove(&candidate);
self.process_summary_order
.retain(|retained| retained != &candidate);
}
}
fn evict_process_summaries_total(&mut self) {
while self.process_summaries.len() >= super::MAX_RECENT_PROCESS_SUMMARIES_TOTAL {
let candidate = self.process_summary_order.iter().find(|key| {
self.process_summaries
.get(*key)
.is_some_and(|summary| summary.final_result.is_some())
});
let Some(candidate) = candidate.cloned() else {
break;
};
self.process_summaries.remove(&candidate);
self.recent_log_dropped_through.remove(&candidate);
self.process_summary_order
.retain(|retained| retained != &candidate);
}
}
}
fn parse_order_cursor(
cursor: Option<&str>,
expected_kind: &str,
) -> Result<Option<u64>, CoordinatorServiceError> {
cursor
.map(|cursor| {
let (kind, order) = cursor.split_once(':').ok_or_else(|| {
CoordinatorServiceError::Protocol(format!(
"invalid {expected_kind} pagination cursor"
))
})?;
if kind != expected_kind {
return Err(CoordinatorServiceError::Protocol(format!(
"invalid {expected_kind} pagination cursor"
)));
}
order.parse::<u64>().map_err(|_| {
CoordinatorServiceError::Protocol(format!(
"invalid {expected_kind} pagination cursor"
))
})
})
.transpose()
}