mod proto_buf_conv;
pub mod spec;
pub mod proto {
#![allow(clippy::all)]
include!(concat!(env!("OUT_DIR"), "/maelstrom_client_base.items.rs"));
}
#[cfg(test)]
extern crate self as maelstrom_client;
pub use proto_buf_conv::{IntoProtoBuf, TryFromProtoBuf};
use derive_more::{Debug, From, Into};
use maelstrom_base::{
stats::JobState, ClientJobId, JobBrokerStatus, JobOutcomeResult, JobWorkerStatus,
};
use maelstrom_container::ContainerImageDepotDir;
use maelstrom_macro::{IntoProtoBuf, TryFromProtoBuf};
use maelstrom_util::{
config::common::{BrokerAddr, CacheSize, InlineLimit, Slots},
root::RootBuf,
};
use serde::Deserialize;
use url::Url;
pub struct ProjectDir;
pub struct CacheDir;
pub const MANIFEST_DIR: &str = "manifests";
pub const STUB_MANIFEST_DIR: &str = "manifests/stubs";
pub const SYMLINK_MANIFEST_DIR: &str = "manifests/symlinks";
pub const SO_LISTINGS_DIR: &str = "so-listings";
impl From<proto::Error> for anyhow::Error {
fn from(e: proto::Error) -> Self {
anyhow::Error::msg(e.message)
}
}
#[derive(Debug, PartialEq, Eq, IntoProtoBuf, TryFromProtoBuf)]
#[proto(proto_buf_type = "proto::RemoteProgress")]
pub struct RemoteProgress {
pub name: String,
pub size: u64,
pub progress: u64,
}
#[derive(Debug, Default, PartialEq, Eq, IntoProtoBuf, TryFromProtoBuf)]
#[proto(proto_buf_type = "proto::IntrospectResponse")]
pub struct IntrospectResponse {
pub artifact_uploads: Vec<RemoteProgress>,
pub image_downloads: Vec<RemoteProgress>,
}
#[derive(Clone, Copy, Debug, Deserialize, From, Into, TryFromProtoBuf, IntoProtoBuf)]
#[proto(proto_buf_type = bool, try_from_into)]
#[serde(transparent)]
#[debug("{_0:?}")]
pub struct AcceptInvalidRemoteContainerTlsCerts(bool);
impl AcceptInvalidRemoteContainerTlsCerts {
pub fn into_inner(self) -> bool {
self.0
}
}
#[derive(Clone, IntoProtoBuf, TryFromProtoBuf)]
#[proto(proto_buf_type = "proto::LogKeyValue")]
pub struct RpcLogKeyValue {
pub key: String,
pub value: String,
}
#[derive(Clone, IntoProtoBuf, TryFromProtoBuf)]
#[proto(proto_buf_type = "proto::LogMessage")]
pub struct RpcLogMessage {
pub message: String,
pub level: slog::Level,
pub tag: String,
pub key_values: Vec<RpcLogKeyValue>,
}
impl RpcLogMessage {
pub fn log_to(self, log: &slog::Logger) {
let location = slog::RecordLocation {
file: "<remote-file>",
line: 0,
column: 0,
function: "",
module: "<remote-module>",
};
let rs = slog::RecordStatic {
location: &location,
level: self.level,
tag: &self.tag,
};
let kv = SimpleKV(
self.key_values
.into_iter()
.map(|e| slog::SingleKV::from((e.key, e.value)))
.collect(),
);
log.log(&slog::Record::new(
&rs,
&format_args!("[client-process] {}", self.message),
slog::BorrowedKV(&kv),
));
}
}
struct SimpleKV(Vec<slog::SingleKV<String>>);
impl slog::KV for SimpleKV {
fn serialize(
&self,
record: &slog::Record,
serializer: &mut dyn slog::Serializer,
) -> slog::Result {
for e in &self.0 {
e.serialize(record, serializer)?;
}
Ok(())
}
}
#[derive(Clone, Debug, PartialEq, Eq, PartialOrd, Ord, IntoProtoBuf, TryFromProtoBuf)]
#[proto(
proto_buf_type = "proto::JobRunningStatus",
enum_type = "proto::job_running_status::Status"
)]
pub enum JobRunningStatus {
AtBroker(JobBrokerStatus),
AtLocalWorker(JobWorkerStatus),
}
impl JobRunningStatus {
pub fn to_state(&self) -> JobState {
match self {
JobRunningStatus::AtBroker(JobBrokerStatus::WaitingForLayers) => {
JobState::WaitingForArtifacts
}
JobRunningStatus::AtBroker(JobBrokerStatus::WaitingForWorker) => JobState::Pending,
JobRunningStatus::AtBroker(JobBrokerStatus::AtWorker(
_,
JobWorkerStatus::WaitingForLayers,
))
| JobRunningStatus::AtLocalWorker(JobWorkerStatus::WaitingForLayers) => {
JobState::WaitingForArtifacts
}
JobRunningStatus::AtBroker(JobBrokerStatus::AtWorker(
_,
JobWorkerStatus::WaitingToExecute,
))
| JobRunningStatus::AtLocalWorker(JobWorkerStatus::WaitingToExecute) => {
JobState::Pending
}
JobRunningStatus::AtBroker(JobBrokerStatus::AtWorker(
_,
JobWorkerStatus::Executing,
))
| JobRunningStatus::AtLocalWorker(JobWorkerStatus::Executing) => JobState::Running,
}
}
}
#[derive(Clone, Debug, From, PartialEq, Eq, PartialOrd, Ord, IntoProtoBuf, TryFromProtoBuf)]
#[proto(
proto_buf_type = "proto::JobStatus",
enum_type = "proto::job_status::Status"
)]
pub enum JobStatus {
Running(JobRunningStatus),
#[proto(proto_buf_type = "proto::JobCompletedStatus")]
Completed {
client_job_id: ClientJobId,
#[proto(option)]
result: JobOutcomeResult,
},
}
#[derive(Clone, Debug, IntoProtoBuf, TryFromProtoBuf)]
#[proto(proto_buf_type = proto::TcpClusterCommunicationStrategy)]
pub struct TcpClusterCommunicationStrategy {
pub broker: BrokerAddr,
}
#[derive(Clone, Debug, IntoProtoBuf, TryFromProtoBuf)]
#[proto(proto_buf_type = proto::GitHubClusterCommunicationStrategy)]
pub struct GitHubClusterCommunicationStrategy {
pub token: String,
pub url: Url,
}
#[derive(Clone, Debug, IntoProtoBuf, TryFromProtoBuf)]
#[proto(
proto_buf_type = "proto::ClusterCommunicationStrategy",
enum_type = "proto::cluster_communication_strategy::Strategy"
)]
pub enum ClusterCommunicationStrategy {
Tcp(TcpClusterCommunicationStrategy),
GitHub(GitHubClusterCommunicationStrategy),
}
#[derive(Clone, Debug, IntoProtoBuf, TryFromProtoBuf)]
#[proto(proto_buf_type = "proto::StartRequest")]
pub struct StartRequest {
pub project_dir: RootBuf<ProjectDir>,
pub container_image_depot_dir: RootBuf<ContainerImageDepotDir>,
pub cache_dir: RootBuf<CacheDir>,
pub cache_size: CacheSize,
pub inline_limit: InlineLimit,
pub slots: Slots,
pub accept_invalid_remote_container_tls_certs: AcceptInvalidRemoteContainerTlsCerts,
pub cluster_communication_strategy: Option<ClusterCommunicationStrategy>,
}
#[derive(IntoProtoBuf, TryFromProtoBuf)]
#[proto(proto_buf_type = "proto::RunJobRequest")]
pub struct RunJobRequest {
#[proto(option)]
pub spec: spec::JobSpec,
}
#[derive(IntoProtoBuf, TryFromProtoBuf)]
#[proto(proto_buf_type = "proto::AddContainerRequest")]
pub struct AddContainerRequest {
pub name: String,
#[proto(option)]
pub container: spec::ContainerSpec,
}