use std::sync::Arc;
use crate::{
MicrosandboxError, MicrosandboxResult,
backend::{
Backend, CloudCreateSandboxResponse, SandboxCloudState, SandboxHandleCloudState,
SandboxHandleInner, SandboxHandleLocalState,
},
db::entity::sandbox as sandbox_entity,
error::Operation,
};
use super::{Sandbox, SandboxConfig, SandboxModificationBuilder, SandboxStatus, SandboxStopResult};
pub const DEFAULT_CONNECT_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(10);
pub const DEFAULT_STOP_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(10);
pub const DEFAULT_KILL_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5);
pub struct SandboxHandle {
backend: Arc<dyn Backend>,
inner: SandboxHandleInner,
name: String,
}
impl SandboxHandle {
pub(crate) fn from_local_model(
backend: Arc<dyn Backend>,
model: sandbox_entity::Model,
pid: Option<i32>,
) -> Self {
let name = model.name.clone();
Self {
backend,
inner: SandboxHandleInner::Local(SandboxHandleLocalState {
db_id: model.id,
status: model.status,
config_json: model.config,
active_config_json: model.active_config,
created_at: model.created_at.map(|dt| dt.and_utc()),
updated_at: model.updated_at.map(|dt| dt.and_utc()),
pid,
}),
name,
}
}
pub(crate) fn from_cloud(
backend: Arc<dyn Backend>,
cloud: CloudCreateSandboxResponse,
) -> MicrosandboxResult<Self> {
let status = crate::backend::sandbox::cloud_status_to_sandbox_status(cloud.status);
let config_json = serde_json::to_string(&cloud.spec)?;
let name = cloud.name.clone();
Ok(Self {
backend,
inner: SandboxHandleInner::Cloud(SandboxHandleCloudState {
id: cloud.id,
org_id: cloud.org_id,
status,
config_json,
created_at: Some(cloud.created_at),
started_at: cloud.started_at,
stopped_at: cloud.stopped_at,
last_failure_message: cloud.last_failure_message,
}),
name,
})
}
pub fn name(&self) -> &str {
&self.name
}
pub fn backend_kind(&self) -> crate::backend::BackendKind {
self.backend.kind()
}
pub fn local(&self) -> Option<&SandboxHandleLocalState> {
match &self.inner {
SandboxHandleInner::Local(s) => Some(s),
SandboxHandleInner::Cloud(_) => None,
}
}
pub fn cloud(&self) -> Option<&SandboxHandleCloudState> {
match &self.inner {
SandboxHandleInner::Cloud(s) => Some(s),
SandboxHandleInner::Local(_) => None,
}
}
pub fn status_snapshot(&self) -> SandboxStatus {
match &self.inner {
SandboxHandleInner::Local(s) => s.status,
SandboxHandleInner::Cloud(s) => s.status,
}
}
pub fn last_failure_message_snapshot(&self) -> Option<String> {
match &self.inner {
SandboxHandleInner::Cloud(s) => s.last_failure_message.clone(),
SandboxHandleInner::Local(_) => None,
}
}
pub fn config_json(&self) -> &str {
match &self.inner {
SandboxHandleInner::Local(s) => &s.config_json,
SandboxHandleInner::Cloud(s) => &s.config_json,
}
}
pub fn active_config_json(&self) -> Option<&str> {
match &self.inner {
SandboxHandleInner::Local(s) => s.active_config_json.as_deref(),
SandboxHandleInner::Cloud(_) => None,
}
}
pub fn config(&self) -> MicrosandboxResult<SandboxConfig> {
match &self.inner {
SandboxHandleInner::Local(s) => Ok(serde_json::from_str(&s.config_json)?),
SandboxHandleInner::Cloud(_) => Err(MicrosandboxError::local_only(
Operation::SandboxHandleConfig,
)),
}
}
pub fn active_config(&self) -> MicrosandboxResult<Option<SandboxConfig>> {
self.active_config_json()
.map(serde_json::from_str)
.transpose()
.map_err(Into::into)
}
pub fn modify(&self) -> SandboxModificationBuilder {
SandboxModificationBuilder::new(self.backend.clone(), self.name.clone())
}
fn require_running(&self, operation: &str) -> MicrosandboxResult<()> {
let status = self.status_snapshot();
if matches!(
status,
super::SandboxStatus::Running | super::SandboxStatus::Draining
) {
return Ok(());
}
Err(MicrosandboxError::SandboxNotRunning(format!(
"'{}' is not running (status: {status:?}); cannot {operation}",
self.name
)))
}
pub async fn refresh(&self) -> MicrosandboxResult<SandboxHandle> {
self.backend
.sandboxes()
.get(self.backend.clone(), &self.name)
.await
}
pub fn created_at(&self) -> Option<chrono::DateTime<chrono::Utc>> {
match &self.inner {
SandboxHandleInner::Local(s) => s.created_at,
SandboxHandleInner::Cloud(s) => s.created_at,
}
}
pub fn updated_at(&self) -> Option<chrono::DateTime<chrono::Utc>> {
match &self.inner {
SandboxHandleInner::Local(s) => s.updated_at,
SandboxHandleInner::Cloud(s) => s.stopped_at.or(s.started_at).or(s.created_at),
}
}
pub async fn logs(
&self,
opts: &crate::logs::LogOptions,
) -> MicrosandboxResult<Vec<crate::logs::LogEntry>> {
self.backend
.sandboxes()
.logs(self.backend.clone(), &self.name, opts)
.await
}
pub async fn log_stream(
&self,
opts: &crate::logs::LogStreamOptions,
) -> MicrosandboxResult<crate::backend::sandbox::LogStream> {
self.backend
.sandboxes()
.log_stream(self.backend.clone(), &self.name, opts)
.await
}
pub async fn metrics(&self) -> MicrosandboxResult<super::SandboxMetrics> {
let local = self
.local()
.ok_or_else(|| MicrosandboxError::local_only(Operation::SandboxHandleMetrics))?;
if local.status != SandboxStatus::Running && local.status != SandboxStatus::Draining {
return Err(MicrosandboxError::SandboxNotRunning(format!(
"'{}' is not running (status: {:?})",
self.name, local.status
)));
}
let config = self.config()?;
if config.effective_metrics_interval().is_none() {
return Err(MicrosandboxError::MetricsDisabled(self.name.clone()));
}
let local_backend = self
.backend
.as_local()
.ok_or_else(|| MicrosandboxError::local_only(Operation::SandboxHandleMetrics))?;
let db = local_backend.db().await?.read();
super::metrics::metrics_for_sandbox(db, local_backend, local.db_id, &config).await
}
pub async fn start(&self) -> MicrosandboxResult<Sandbox> {
self.backend
.sandboxes()
.start(self.backend.clone(), &self.name)
.await
}
pub async fn start_detached(&self) -> MicrosandboxResult<Sandbox> {
self.backend
.sandboxes()
.start_detached(self.backend.clone(), &self.name)
.await
}
pub async fn connect(&self) -> MicrosandboxResult<Sandbox> {
self.connect_with_timeout(DEFAULT_CONNECT_TIMEOUT).await
}
pub async fn connect_with_timeout(
&self,
timeout: std::time::Duration,
) -> MicrosandboxResult<Sandbox> {
if !matches!(
self.status_snapshot(),
SandboxStatus::Running | SandboxStatus::Draining
) {
return Err(MicrosandboxError::SandboxNotRunning(format!(
"'{}' is not running (status: {:?})",
self.name,
self.status_snapshot()
)));
}
match &self.inner {
SandboxHandleInner::Local(local) => {
let local_backend = self.backend.as_local().ok_or_else(|| {
MicrosandboxError::local_only(Operation::SandboxHandleConnect)
})?;
let client = crate::sandbox::fs::agent::connect_agent_with_timeout(
local_backend,
&self.name,
timeout,
)
.await?;
let config: SandboxConfig = serde_json::from_str(&local.config_json)?;
Ok(Sandbox::from_local(
self.backend.clone(),
crate::backend::SandboxLocalState {
db_id: local.db_id,
handle: None,
client: Arc::new(client),
},
config,
))
}
SandboxHandleInner::Cloud(cloud) => {
let spec = serde_json::from_str(&cloud.config_json)?;
let config =
crate::backend::sandbox::sandbox_config_from_cloud_spec(&self.name, spec);
let created_at = cloud.created_at.ok_or_else(|| {
MicrosandboxError::Runtime(format!(
"cloud sandbox {:?} is missing its creation timestamp",
self.name
))
})?;
Ok(Sandbox::from_cloud_state(
self.backend.clone(),
SandboxCloudState {
id: cloud.id.clone(),
org_id: cloud.org_id.clone(),
created_at,
},
self.name.clone(),
config,
))
}
}
}
pub async fn ping(&self) -> MicrosandboxResult<super::SandboxPingResult> {
self.require_running("ping")?;
self.connect().await?.ping().await
}
pub async fn touch(&self) -> MicrosandboxResult<super::SandboxTouchResult> {
self.require_running("touch")?;
self.connect().await?.touch().await
}
pub async fn snapshot(
&self,
name: &str,
) -> MicrosandboxResult<super::super::snapshot::Snapshot> {
if self.local().is_none() {
return Err(MicrosandboxError::local_only(
Operation::SandboxHandleSnapshot,
));
}
use super::super::snapshot::Snapshot;
Snapshot::builder(name)
.from_sandbox(&self.name)
.create()
.await
}
pub async fn stop(&self) -> MicrosandboxResult<()> {
self.stop_with_timeout(DEFAULT_STOP_TIMEOUT).await
}
pub async fn stop_with_timeout(&self, timeout: std::time::Duration) -> MicrosandboxResult<()> {
let current = self.refresh().await?;
if sandbox_status_is_terminal(current.status_snapshot()) {
return Ok(());
}
if timeout.is_zero() {
current.kill_with_timeout(DEFAULT_KILL_TIMEOUT).await?;
return Ok(());
}
current.request_stop().await?;
match tokio::time::timeout(timeout, current.wait_until_stopped()).await {
Ok(Ok(_)) => {
#[cfg(windows)]
current.reap_leaked_local_runtime().await?;
return Ok(());
}
Ok(Err(error)) => return Err(error),
Err(_) => {}
}
tracing::warn!(
sandbox = %current.name,
timeout_secs = timeout.as_secs(),
"graceful stop exceeded timeout, escalating to kill"
);
current.request_kill().await?;
match tokio::time::timeout(DEFAULT_KILL_TIMEOUT, current.wait_until_stopped()).await {
Ok(result) => {
result?;
Ok(())
}
Err(_) => Err(MicrosandboxError::Runtime(format!(
"timed out observing stopped state for sandbox '{}'",
current.name
))),
}
}
pub async fn request_stop(&self) -> MicrosandboxResult<()> {
let current = self.refresh().await?;
if sandbox_status_is_terminal(current.status_snapshot()) {
return Ok(());
}
current
.backend
.sandboxes()
.stop(current.backend.clone(), ¤t.name)
.await
}
pub async fn kill(&self) -> MicrosandboxResult<()> {
self.kill_with_timeout(DEFAULT_KILL_TIMEOUT).await
}
pub async fn request_kill(&self) -> MicrosandboxResult<()> {
let current = self.refresh().await?;
if sandbox_status_is_terminal(current.status_snapshot()) {
return Ok(());
}
current
.backend
.sandboxes()
.kill(current.backend.clone(), ¤t.name)
.await
}
pub async fn kill_with_timeout(&self, timeout: std::time::Duration) -> MicrosandboxResult<()> {
let current = self.refresh().await?;
if sandbox_status_is_terminal(current.status_snapshot()) {
return Ok(());
}
current.request_kill().await?;
match tokio::time::timeout(timeout, current.wait_until_stopped()).await {
Ok(result) => {
result?;
Ok(())
}
Err(_) => Err(MicrosandboxError::Runtime(format!(
"timed out observing stopped state for sandbox '{}'",
current.name
))),
}
}
pub async fn request_drain(&self) -> MicrosandboxResult<()> {
let current = self.refresh().await?;
if sandbox_status_is_terminal(current.status_snapshot()) {
return Ok(());
}
current
.backend
.sandboxes()
.drain(current.backend.clone(), ¤t.name)
.await
}
pub async fn wait_until_stopped(&self) -> MicrosandboxResult<SandboxStopResult> {
loop {
let current = match self.refresh().await {
Ok(current) => current,
Err(error)
if self.is_local_ephemeral()
&& super::sandbox_not_found_for_name(&error, &self.name) =>
{
return Ok(super::ephemeral_cleanup_stop_result(&self.name));
}
Err(error) => return Err(error),
};
let status = current.status_snapshot();
if sandbox_status_is_terminal(status) {
return Ok(SandboxStopResult {
name: current.name,
status,
exit_code: None,
signal: None,
observed_at: chrono::Utc::now(),
source: Some("refreshed backend state".to_string()),
});
}
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
}
}
pub async fn remove(&self) -> MicrosandboxResult<()> {
match &self.inner {
SandboxHandleInner::Local(_) => {
let refreshed = self.refresh().await?;
let local = refreshed
.local()
.ok_or_else(|| MicrosandboxError::local_only(Operation::SandboxHandleRemove))?;
if matches!(
local.status,
SandboxStatus::Running | SandboxStatus::Draining | SandboxStatus::Paused
) {
return Err(MicrosandboxError::SandboxStillRunning(format!(
"cannot remove sandbox '{}': still running",
self.name
)));
}
let local_backend = self
.backend
.as_local()
.ok_or_else(|| MicrosandboxError::local_only(Operation::SandboxHandleRemove))?;
#[cfg(windows)]
super::reap_leaked_runtime_process(local_backend, local.db_id, &self.name).await?;
super::remove_local_persisted_sandbox(local_backend, &self.name, local.db_id).await
}
SandboxHandleInner::Cloud(_) => {
self.backend
.sandboxes()
.remove(self.backend.clone(), &self.name)
.await
}
}
}
#[cfg(windows)]
async fn reap_leaked_local_runtime(&self) -> MicrosandboxResult<()> {
let Some(local) = self.local() else {
return Ok(());
};
let Some(local_backend) = self.backend.as_local() else {
return Ok(());
};
super::reap_leaked_runtime_process(local_backend, local.db_id, &self.name)
.await
.map(|_| ())
}
fn is_local_ephemeral(&self) -> bool {
is_local_ephemeral_handle(&self.inner)
}
}
fn is_local_ephemeral_handle(inner: &SandboxHandleInner) -> bool {
let SandboxHandleInner::Local(state) = inner else {
return false;
};
serde_json::from_str::<SandboxConfig>(&state.config_json)
.map(|config| config.spec.lifecycle.ephemeral)
.unwrap_or(false)
}
fn sandbox_status_is_terminal(status: SandboxStatus) -> bool {
matches!(status, SandboxStatus::Stopped | SandboxStatus::Crashed)
}
impl std::fmt::Debug for SandboxHandle {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("SandboxHandle")
.field("name", &self.name)
.field("backend_kind", &self.backend.kind())
.field("status", &self.status_snapshot())
.finish()
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::backend::{BackendKind, CloudBackend, CloudSandboxStatus};
#[tokio::test]
async fn cloud_connect_rebuilds_live_sandbox_without_http_request() {
let handle = cloud_handle(CloudSandboxStatus::Running);
let sandbox = handle.connect().await.unwrap();
assert_eq!(sandbox.name(), "cloud-connect-test");
assert_eq!(sandbox.backend_kind(), BackendKind::Cloud);
assert_eq!(sandbox.cloud().unwrap().id, "sandbox-id");
assert_eq!(sandbox.config().spec.name, "cloud-connect-test");
}
#[tokio::test]
async fn cloud_connect_rejects_stopped_sandbox() {
let handle = cloud_handle(CloudSandboxStatus::Stopped);
let result = handle.connect().await;
assert!(matches!(
result,
Err(MicrosandboxError::SandboxNotRunning(_))
));
}
fn cloud_handle(status: CloudSandboxStatus) -> SandboxHandle {
let backend: Arc<dyn Backend> =
Arc::new(CloudBackend::new("https://unused.invalid", "msb_test_connect").unwrap());
SandboxHandle::from_cloud(
backend,
CloudCreateSandboxResponse {
id: "sandbox-id".into(),
org_id: "org-id".into(),
name: "cloud-connect-test".into(),
slug: "cloud-connect-test".into(),
status,
status_reason: None,
spec: None,
ephemeral: false,
created_at: chrono::Utc::now(),
started_at: None,
stopped_at: None,
last_failure_message: None,
},
)
.unwrap()
}
}