use std::collections::HashMap;
use std::str::FromStr;
use std::time::Duration;
use auv_api_proto::auv::api::daemon::v1 as proto;
use crate::client::Client;
use crate::error::ClientError;
use crate::resource::{DeviceId, RunnerClassId, RunnerClassSelector, RunnerId, RunnerSelector};
use crate::time::Timestamp;
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum RunnerLifecycle {
Ephemeral,
UnlessIdle,
UnlessShutdown,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum RunnerPhase {
Unspecified,
Starting,
Ready,
Draining,
Stopped,
Failed,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct RunnerClass {
pub id: RunnerClassId,
pub device: Option<DeviceId>,
pub display_name: String,
pub supported_lifecycles: Vec<RunnerLifecycle>,
pub available: bool,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct Runner {
pub id: RunnerId,
pub device: DeviceId,
pub class: RunnerClassId,
pub labels: HashMap<String, String>,
pub lifecycle: RunnerLifecycle,
pub idle_timeout: Option<Duration>,
pub phase: RunnerPhase,
pub created_at: Option<Timestamp>,
pub process_id: Option<u32>,
pub active_operations: u64,
pub idle_deadline: Option<Timestamp>,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct CreateRunner {
pub device: Option<DeviceId>,
pub class: RunnerClassId,
pub labels: HashMap<String, String>,
pub lifecycle: RunnerLifecycle,
pub idle_timeout: Option<Duration>,
}
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub struct StopRunner {
pub grace_period: Option<Duration>,
pub force: bool,
}
#[derive(Debug, thiserror::Error)]
pub enum RunnerError {
#[error(transparent)]
Client(#[from] ClientError),
#[error(transparent)]
Identity(#[from] crate::resource::IdentityError),
#[error("Runner response omitted {0}")]
MissingField(&'static str),
#[error("unknown Runner ID {0:?}")]
NotFound(String),
#[error("ambiguous Runner ID prefix {0:?}; provide more characters")]
Ambiguous(String),
#[error("Runner is not owned by selected Device {0:?}")]
DeviceConflict(String),
#[error("duration exceeds the protocol range")]
DurationRange,
}
#[derive(Clone, Debug)]
pub struct Runners {
client: Client,
}
impl Runners {
pub(crate) fn new(client: Client) -> Self {
Self { client }
}
pub async fn create(&self, request: CreateRunner) -> Result<Runner, RunnerError> {
let response = self
.client
.grpc_client()
.runners()
.create_runner(proto::CreateRunnerRequest {
device: request.device.map(|id| proto::DeviceRef {
device_id: id.to_string(),
}),
runner_class: Some(proto::RunnerClassRef {
runner_class: request.class.to_string(),
}),
labels: request.labels,
lifecycle: proto::RunnerLifecycle::from(request.lifecycle) as i32,
idle_timeout: request.idle_timeout.map(duration_to_proto).transpose()?,
})
.await
.map_err(|status| ClientError::from_status("CreateRunner", status))?;
Runner::try_from(response)
}
pub async fn list(&self) -> Result<Vec<Runner>, RunnerError> {
self
.client
.grpc_client()
.runners()
.list_runners()
.await
.map_err(|status| ClientError::from_status("ListRunners", status))?
.into_iter()
.map(Runner::try_from)
.collect()
}
pub async fn get(&self, selector: &RunnerSelector, device: Option<&DeviceId>) -> Result<Runner, RunnerError> {
let id = resolve(selector, &self.list().await?)?;
let runner = Runner::try_from(
self
.client
.grpc_client()
.runners()
.get_runner(id.to_string())
.await
.map_err(|status| ClientError::from_status("GetRunner", status))?,
)?;
validate_device(&runner, device)?;
Ok(runner)
}
pub async fn stop(&self, selector: &RunnerSelector, device: Option<&DeviceId>, options: StopRunner) -> Result<Runner, RunnerError> {
let existing = self.get(selector, device).await?;
let response = self
.client
.grpc_client()
.runners()
.delete_runner_with_options(existing.id.to_string(), options.grace_period.map(duration_to_proto).transpose()?, options.force)
.await
.map_err(|status| ClientError::from_status("DeleteRunner", status))?;
Runner::try_from(response)
}
pub async fn classes(&self, device: Option<&DeviceId>) -> Result<Vec<RunnerClass>, RunnerError> {
self
.client
.grpc_client()
.runner_classes()
.list_runner_classes(device.map(|id| proto::DeviceRef {
device_id: id.to_string(),
}))
.await
.map_err(|status| ClientError::from_status("ListRunnerClasses", status))?
.into_iter()
.map(RunnerClass::try_from)
.collect()
}
pub async fn class(&self, selector: &RunnerClassSelector, device: Option<&DeviceId>) -> Result<RunnerClass, RunnerError> {
let response = self
.client
.grpc_client()
.runner_classes()
.get_runner_class(
selector.id().to_string(),
device.map(|id| proto::DeviceRef {
device_id: id.to_string(),
}),
)
.await
.map_err(|status| ClientError::from_status("GetRunnerClass", status))?;
RunnerClass::try_from(response)
}
}
fn resolve(selector: &RunnerSelector, runners: &[Runner]) -> Result<RunnerId, RunnerError> {
let matches = runners.iter().filter(|runner| selector.matches(&runner.id)).collect::<Vec<_>>();
match matches.as_slice() {
[] => Err(RunnerError::NotFound(selector.as_str().to_string())),
[runner] => Ok(runner.id.clone()),
_ => Err(RunnerError::Ambiguous(selector.as_str().to_string())),
}
}
fn validate_device(runner: &Runner, expected: Option<&DeviceId>) -> Result<(), RunnerError> {
if expected.is_some_and(|expected| expected != &runner.device) {
return Err(RunnerError::DeviceConflict(expected.expect("checked Some").to_string()));
}
Ok(())
}
fn duration_to_proto(value: Duration) -> Result<prost_types::Duration, RunnerError> {
Ok(prost_types::Duration {
seconds: i64::try_from(value.as_secs()).map_err(|_| RunnerError::DurationRange)?,
nanos: i32::try_from(value.subsec_nanos()).expect("subsecond nanoseconds fit i32"),
})
}
impl TryFrom<proto::Runner> for Runner {
type Error = RunnerError;
fn try_from(runner: proto::Runner) -> Result<Self, Self::Error> {
let lifecycle = proto::RunnerLifecycle::try_from(runner.lifecycle).unwrap_or(proto::RunnerLifecycle::Unspecified);
Ok(Self {
id: RunnerId::from_str(&runner.r#ref.ok_or(RunnerError::MissingField("canonical ID"))?.runner_id)?,
device: DeviceId::from_str(&runner.device.ok_or(RunnerError::MissingField("Device"))?.device_id)?,
class: RunnerClassId::from_str(&runner.runner_class.ok_or(RunnerError::MissingField("RunnerClass"))?.runner_class)?,
labels: runner.labels,
lifecycle: RunnerLifecycle::try_from(lifecycle)?,
idle_timeout: runner.idle_timeout.map(|value| Duration::new(value.seconds.max(0) as u64, value.nanos.max(0) as u32)),
phase: match proto::RunnerPhase::try_from(runner.phase).unwrap_or(proto::RunnerPhase::Unspecified) {
proto::RunnerPhase::Unspecified => RunnerPhase::Unspecified,
proto::RunnerPhase::Starting => RunnerPhase::Starting,
proto::RunnerPhase::Ready => RunnerPhase::Ready,
proto::RunnerPhase::Draining => RunnerPhase::Draining,
proto::RunnerPhase::Stopped => RunnerPhase::Stopped,
proto::RunnerPhase::Failed => RunnerPhase::Failed,
},
created_at: runner.created_at.map(Into::into),
process_id: (runner.process_id != 0).then_some(runner.process_id),
active_operations: runner.active_operations,
idle_deadline: runner.idle_deadline.map(Into::into),
})
}
}
impl TryFrom<proto::RunnerClass> for RunnerClass {
type Error = RunnerError;
fn try_from(class: proto::RunnerClass) -> Result<Self, Self::Error> {
Ok(Self {
id: RunnerClassId::from_str(&class.r#ref.ok_or(RunnerError::MissingField("RunnerClass ID"))?.runner_class)?,
device: class.device.map(|device| DeviceId::from_str(&device.device_id)).transpose()?,
display_name: class.display_name,
supported_lifecycles: class
.supported_lifecycles
.into_iter()
.filter_map(|value| proto::RunnerLifecycle::try_from(value).ok())
.filter_map(|value| RunnerLifecycle::try_from(value).ok())
.collect(),
available: class.available,
})
}
}
impl TryFrom<proto::RunnerLifecycle> for RunnerLifecycle {
type Error = RunnerError;
fn try_from(value: proto::RunnerLifecycle) -> Result<Self, Self::Error> {
match value {
proto::RunnerLifecycle::Ephemeral => Ok(Self::Ephemeral),
proto::RunnerLifecycle::UnlessIdle => Ok(Self::UnlessIdle),
proto::RunnerLifecycle::UnlessShutdown => Ok(Self::UnlessShutdown),
proto::RunnerLifecycle::Unspecified => Err(RunnerError::MissingField("Runner lifecycle")),
}
}
}
impl From<RunnerLifecycle> for proto::RunnerLifecycle {
fn from(value: RunnerLifecycle) -> Self {
match value {
RunnerLifecycle::Ephemeral => Self::Ephemeral,
RunnerLifecycle::UnlessIdle => Self::UnlessIdle,
RunnerLifecycle::UnlessShutdown => Self::UnlessShutdown,
}
}
}