use arcbox_connect::sandbox_v1;
use arcbox_vm::file_io::proto as file_proto;
use arcbox_vm::template_catalog::{ReadyProbeSpec, TemplateDefaultsSpec, TemplateEntry};
use arcbox_vm::{
CheckpointInfo, CheckpointSummary, ExecutionChannel, ExecutionSnapshot, ExitStatus, IdleAction,
SandboxEvent as VmSandboxEvent, SandboxInfo, SandboxState, SandboxSummary, StdinState,
};
use buffa_types::google::protobuf::Timestamp;
use chrono::{DateTime, Utc};
const DEFAULT_PAGE_SIZE: u32 = 100;
const MAX_PAGE_SIZE: u32 = 1000;
pub(super) fn timestamp(dt: DateTime<Utc>) -> Timestamp {
Timestamp {
seconds: dt.timestamp(),
nanos: i32::try_from(dt.timestamp_subsec_nanos()).unwrap_or(0),
..Default::default()
}
}
pub(super) fn timestamp_from_rfc3339(s: &str) -> Option<Timestamp> {
DateTime::parse_from_rfc3339(s)
.ok()
.map(|dt| timestamp(dt.with_timezone(&Utc)))
}
pub(super) fn timestamp_from_unix_nanos(ns: i64) -> Timestamp {
Timestamp {
seconds: ns.div_euclid(1_000_000_000),
nanos: i32::try_from(ns.rem_euclid(1_000_000_000)).expect("nanos below 1e9"),
..Default::default()
}
}
pub(super) fn exit_status_to_proto(status: ExitStatus) -> sandbox_v1::ExitStatus {
let status = match status {
ExitStatus::Code(code) => sandbox_v1::exit_status::Status::Code(code),
ExitStatus::Signaled(signal) => sandbox_v1::exit_status::Status::Signal(signal),
};
sandbox_v1::ExitStatus {
status: Some(status),
..Default::default()
}
}
pub(super) fn stdin_to_proto(state: StdinState) -> sandbox_v1::StdinStatus {
sandbox_v1::StdinStatus {
bytes_written: state.bytes_written,
closed: state.closed,
..Default::default()
}
}
pub(super) fn execution_to_proto(snap: &ExecutionSnapshot) -> sandbox_v1::Execution {
let state = if snap.is_running() {
sandbox_v1::ExecutionState::Running
} else {
sandbox_v1::ExecutionState::Exited
};
sandbox_v1::Execution {
id: snap.id.clone(),
sandbox_id: snap.sandbox_id.clone(),
state: state.into(),
tty: snap.tty,
started_at: timestamp(snap.started_at).into(),
exited_at: snap.exited_at.map(timestamp).into(),
exit_status: snap.exit_status.map(exit_status_to_proto).into(),
error: snap.error.clone().unwrap_or_default(),
stdout_len: snap.stdout_len,
stderr_len: snap.stderr_len,
stdin: stdin_to_proto(snap.stdin).into(),
..Default::default()
}
}
fn file_kind_to_proto(kind: &str) -> sandbox_v1::FileKind {
match kind {
file_proto::KIND_FILE => sandbox_v1::FileKind::File,
file_proto::KIND_DIR => sandbox_v1::FileKind::Directory,
file_proto::KIND_SYMLINK => sandbox_v1::FileKind::Symlink,
file_proto::KIND_OTHER => sandbox_v1::FileKind::Other,
_ => sandbox_v1::FileKind::Unspecified,
}
}
pub(super) fn file_stat_to_proto(dto: &file_proto::FileStatDto) -> sandbox_v1::FileStat {
sandbox_v1::FileStat {
name: dto.name.clone(),
kind: file_kind_to_proto(&dto.kind).into(),
size: dto.size,
mode: dto.mode,
modified_at: Timestamp {
seconds: dto.mtime_secs,
nanos: i32::try_from(dto.mtime_nanos).unwrap_or(0),
..Default::default()
}
.into(),
uid: dto.uid,
gid: dto.gid,
symlink_target: dto.symlink_target.clone(),
..Default::default()
}
}
pub(super) fn fs_event_to_proto(dto: file_proto::FsEventDto) -> sandbox_v1::FsEvent {
let kind = match dto.kind.as_str() {
file_proto::EVENT_CREATED => sandbox_v1::FsEventKind::Created,
file_proto::EVENT_MODIFIED => sandbox_v1::FsEventKind::Modified,
file_proto::EVENT_REMOVED => sandbox_v1::FsEventKind::Removed,
file_proto::EVENT_RENAMED => sandbox_v1::FsEventKind::Renamed,
_ => sandbox_v1::FsEventKind::Unspecified,
};
sandbox_v1::FsEvent {
kind: kind.into(),
path: dto.path,
renamed_to: dto.renamed_to,
..Default::default()
}
}
pub(super) fn channel_to_proto(channel: ExecutionChannel, tty: bool) -> sandbox_v1::StdioChannel {
match (channel, tty) {
(ExecutionChannel::Stdout, true) => sandbox_v1::StdioChannel::Pty,
(ExecutionChannel::Stdout, false) => sandbox_v1::StdioChannel::Stdout,
(ExecutionChannel::Stderr, _) => sandbox_v1::StdioChannel::Stderr,
}
}
pub(super) fn state_to_proto(state: SandboxState) -> sandbox_v1::SandboxState {
match state {
SandboxState::Starting => sandbox_v1::SandboxState::Starting,
SandboxState::Ready => sandbox_v1::SandboxState::Ready,
SandboxState::Running => sandbox_v1::SandboxState::Running,
SandboxState::Stopping => sandbox_v1::SandboxState::Stopping,
SandboxState::Stopped => sandbox_v1::SandboxState::Stopped,
SandboxState::Failed => sandbox_v1::SandboxState::Failed,
SandboxState::Pausing => sandbox_v1::SandboxState::Pausing,
SandboxState::Paused => sandbox_v1::SandboxState::Paused,
}
}
pub(super) fn state_filter(state: sandbox_v1::SandboxState) -> Option<&'static str> {
match state {
sandbox_v1::SandboxState::Unspecified => None,
sandbox_v1::SandboxState::Starting => Some("starting"),
sandbox_v1::SandboxState::Ready => Some("ready"),
sandbox_v1::SandboxState::Running => Some("running"),
sandbox_v1::SandboxState::Stopping => Some("stopping"),
sandbox_v1::SandboxState::Stopped => Some("stopped"),
sandbox_v1::SandboxState::Failed => Some("failed"),
sandbox_v1::SandboxState::Pausing => Some("pausing"),
sandbox_v1::SandboxState::Paused => Some("paused"),
}
}
pub(super) fn event_kind(action: &str) -> sandbox_v1::SandboxEventKind {
match action {
"created" => sandbox_v1::SandboxEventKind::Created,
"ready" => sandbox_v1::SandboxEventKind::Ready,
"running" => sandbox_v1::SandboxEventKind::Running,
"idle" => sandbox_v1::SandboxEventKind::Idle,
"stopping" => sandbox_v1::SandboxEventKind::Stopping,
"stopped" => sandbox_v1::SandboxEventKind::Stopped,
"failed" => sandbox_v1::SandboxEventKind::Failed,
"removed" => sandbox_v1::SandboxEventKind::Removed,
"pausing" => sandbox_v1::SandboxEventKind::Pausing,
"paused" => sandbox_v1::SandboxEventKind::Paused,
"resumed" => sandbox_v1::SandboxEventKind::Resumed,
_ => sandbox_v1::SandboxEventKind::Unspecified,
}
}
pub(super) fn vm_event_to_proto(e: VmSandboxEvent) -> sandbox_v1::SandboxEvent {
sandbox_v1::SandboxEvent {
sandbox_id: e.sandbox_id,
kind: event_kind(&e.action).into(),
time: timestamp_from_unix_nanos(e.timestamp_ns).into(),
attributes: e.attributes.into_iter().collect(),
..Default::default()
}
}
pub(super) fn info_to_proto(info: SandboxInfo) -> sandbox_v1::SandboxInfo {
let network = info.network.map(|n| sandbox_v1::SandboxNetwork {
ip_address: n.ip_address,
gateway: n.gateway,
..Default::default()
});
sandbox_v1::SandboxInfo {
id: info.id,
state: state_to_proto(info.state).into(),
labels: info.labels.into_iter().collect(),
limits: sandbox_v1::ResourceLimits {
vcpus: info.vcpus,
memory_mib: info.memory_mib,
..Default::default()
}
.into(),
network: network.into(),
created_at: timestamp(info.created_at).into(),
ready_at: info.ready_at.map(timestamp).into(),
last_exited_at: info.last_exited_at.map(timestamp).into(),
last_exit_status: info.last_exit_status.map(exit_status_to_proto).into(),
error: info.error.unwrap_or_default(),
paused_at: info.paused_at.map(timestamp).into(),
storage_bytes: info.storage_bytes,
ttl_deadline: info.ttl_deadline.map(timestamp).into(),
idle_timeout_seconds: info.idle_timeout_seconds,
on_idle: idle_action_to_proto(info.on_idle).into(),
..Default::default()
}
}
pub(super) fn idle_action_to_proto(action: IdleAction) -> sandbox_v1::IdleAction {
match action {
IdleAction::Kill => sandbox_v1::IdleAction::Kill,
IdleAction::Pause => sandbox_v1::IdleAction::Pause,
}
}
pub(super) fn summary_to_proto(s: SandboxSummary) -> sandbox_v1::SandboxSummary {
sandbox_v1::SandboxSummary {
id: s.id,
state: state_to_proto(s.state).into(),
labels: s.labels.into_iter().collect(),
ip_address: s.ip_address,
created_at: timestamp(s.created_at).into(),
paused_at: s.paused_at.map(timestamp).into(),
storage_bytes: s.storage_bytes,
..Default::default()
}
}
pub(super) fn checkpoint_to_proto(info: CheckpointInfo) -> sandbox_v1::CheckpointResponse {
sandbox_v1::CheckpointResponse {
snapshot_id: info.snapshot_id,
created_at: timestamp_from_rfc3339(&info.created_at).into(),
..Default::default()
}
}
pub(super) fn checkpoint_summary_to_proto(s: CheckpointSummary) -> sandbox_v1::SnapshotSummary {
sandbox_v1::SnapshotSummary {
id: s.id,
sandbox_id: s.sandbox_id,
name: s.name,
labels: s.labels.into_iter().collect(),
created_at: timestamp_from_rfc3339(&s.created_at).into(),
..Default::default()
}
}
pub(super) fn template_to_proto(name: &str, entry: &TemplateEntry) -> sandbox_v1::Template {
sandbox_v1::Template {
name: name.to_owned(),
version: entry.version.clone(),
digest: entry.digest.clone(),
rootfs_ref: entry.rootfs_path.clone(),
warm_snapshot_id: entry
.warm
.as_ref()
.map(|w| w.snapshot_id.clone())
.unwrap_or_default(),
defaults: template_defaults_to_proto(&entry.defaults).into(),
created_at: timestamp(entry.created_at).into(),
labels: entry.labels.clone().into_iter().collect(),
size_bytes: entry.size_bytes,
..Default::default()
}
}
const MAX_READY_PROBE_TIMEOUT_SECS: u32 = 600;
pub(super) fn template_defaults_from_proto(
pb: &sandbox_v1::TemplateDefaults,
) -> Result<TemplateDefaultsSpec, crate::error::SandboxError> {
use crate::error::SandboxError;
for port in &pb.exposed_ports {
if *port == 0 || *port > u32::from(u16::MAX) {
return Err(SandboxError::InvalidArgument(format!(
"exposed_ports entry {port} is out of range (1-65535)"
)));
}
}
Ok(TemplateDefaultsSpec {
vcpus: pb.limits.vcpus,
memory_mib: pb.limits.memory_mib,
cmd: pb.cmd.clone(),
env: pb.env.clone().into_iter().collect(),
exposed_ports: pb.exposed_ports.clone(),
ready_probe: pb
.ready_probe
.as_option()
.map(ready_probe_from_proto)
.transpose()?,
})
}
fn ready_probe_from_proto(
pb: &sandbox_v1::ReadyProbe,
) -> Result<ReadyProbeSpec, crate::error::SandboxError> {
use crate::error::SandboxError;
use sandbox_v1::ready_probe::Probe;
if pb.timeout_seconds > MAX_READY_PROBE_TIMEOUT_SECS {
return Err(SandboxError::InvalidArgument(format!(
"ready_probe timeout_seconds {} exceeds the {MAX_READY_PROBE_TIMEOUT_SECS}s cap",
pb.timeout_seconds
)));
}
match pb.probe.as_ref() {
Some(Probe::Port(port)) => {
let port = u16::try_from(*port)
.ok()
.filter(|p| *p != 0)
.ok_or_else(|| {
SandboxError::InvalidArgument(format!(
"ready_probe port {port} is out of range (1-65535)"
))
})?;
Ok(ReadyProbeSpec::Port {
port,
timeout_seconds: pb.timeout_seconds,
})
}
Some(Probe::Command(command)) => {
if command.cmd.is_empty() {
return Err(SandboxError::InvalidArgument(
"ready_probe command must not be empty".into(),
));
}
Ok(ReadyProbeSpec::Command {
cmd: command.cmd.clone(),
timeout_seconds: pb.timeout_seconds,
})
}
None => Err(SandboxError::InvalidArgument(
"ready_probe requires a port or a command".into(),
)),
}
}
fn template_defaults_to_proto(spec: &TemplateDefaultsSpec) -> Option<sandbox_v1::TemplateDefaults> {
if *spec == TemplateDefaultsSpec::default() {
return None;
}
let limits = (spec.vcpus != 0 || spec.memory_mib != 0).then(|| sandbox_v1::ResourceLimits {
vcpus: spec.vcpus,
memory_mib: spec.memory_mib,
..Default::default()
});
Some(sandbox_v1::TemplateDefaults {
limits: limits.into(),
cmd: spec.cmd.clone(),
env: spec.env.clone().into_iter().collect(),
exposed_ports: spec.exposed_ports.clone(),
ready_probe: spec.ready_probe.as_ref().map(ready_probe_to_proto).into(),
..Default::default()
})
}
fn ready_probe_to_proto(spec: &ReadyProbeSpec) -> sandbox_v1::ReadyProbe {
use sandbox_v1::ready_probe::Probe;
let (probe, timeout_seconds) = match spec {
ReadyProbeSpec::Port {
port,
timeout_seconds,
} => (Probe::Port(u32::from(*port)), *timeout_seconds),
ReadyProbeSpec::Command {
cmd,
timeout_seconds,
} => (
Probe::Command(Box::new(sandbox_v1::CommandProbe {
cmd: cmd.clone(),
..Default::default()
})),
*timeout_seconds,
),
};
sandbox_v1::ReadyProbe {
timeout_seconds,
probe: Some(probe),
..Default::default()
}
}
pub(super) fn paginate<T>(
mut items: Vec<T>,
id_of: impl Fn(&T) -> &str,
page_size: u32,
page_token: &str,
) -> (Vec<T>, String) {
items.sort_by(|a, b| id_of(a).cmp(id_of(b)));
let size = if page_size == 0 {
DEFAULT_PAGE_SIZE
} else {
page_size.min(MAX_PAGE_SIZE)
} as usize;
let start = if page_token.is_empty() {
0
} else {
items.partition_point(|it| id_of(it) <= page_token)
};
let has_more = items.len() > start.saturating_add(size);
let page: Vec<T> = items.into_iter().skip(start).take(size).collect();
let next = if has_more {
page.last()
.map(|it| id_of(it).to_owned())
.unwrap_or_default()
} else {
String::new()
};
(page, next)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn paginate_orders_by_id_and_chains_tokens() {
let items = vec!["c", "a", "d", "b", "e"];
let (page1, token1) = paginate(items.clone(), |s| s, 2, "");
assert_eq!(page1, vec!["a", "b"]);
assert_eq!(token1, "b");
let (page2, token2) = paginate(items.clone(), |s| s, 2, &token1);
assert_eq!(page2, vec!["c", "d"]);
assert_eq!(token2, "d");
let (page3, token3) = paginate(items, |s| s, 2, &token2);
assert_eq!(page3, vec!["e"]);
assert_eq!(token3, "");
}
#[test]
fn paginate_exact_boundary_has_no_next_token() {
let items = vec!["a", "b"];
let (page, token) = paginate(items, |s| s, 2, "");
assert_eq!(page, vec!["a", "b"]);
assert_eq!(token, "");
}
#[test]
fn paginate_defaults_and_caps_page_size() {
let items: Vec<String> = (0..250).map(|i| format!("{i:04}")).collect();
let (page, token) = paginate(items.clone(), |s| s, 0, "");
assert_eq!(page.len(), 100);
assert!(!token.is_empty());
let (page, _) = paginate(items, |s| s, 5000, "");
assert_eq!(page.len(), 250);
}
#[test]
fn unix_nanos_round_to_timestamp() {
let ts = timestamp_from_unix_nanos(1_700_000_001_500_000_000);
assert_eq!(ts.seconds, 1_700_000_001);
assert_eq!(ts.nanos, 500_000_000);
let ts = timestamp_from_unix_nanos(-1_500_000_000);
assert_eq!(ts.seconds, -2);
assert_eq!(ts.nanos, 500_000_000);
}
#[test]
fn rfc3339_parses_to_timestamp() {
let ts = timestamp_from_rfc3339("2026-08-01T00:00:00Z").unwrap();
assert!(ts.seconds > 1_700_000_000);
assert!(timestamp_from_rfc3339("not a date").is_none());
}
}