arcbox-api 0.9.0

API server for ArcBox (gRPC + REST)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
//! The daemon's API, served over Connect.
//!
//! One set of handlers answers Connect (HTTP POST, JSON or binary
//! protobuf), gRPC, and gRPC-Web on a single endpoint (CORE-53).
//!
//! The sandbox half keeps the proto's control-plane / data-plane split.
//! The local daemon serves both and forwards guest operations to the
//! System VM's agent:
//!
//! - [`control`] — sandbox lifecycle, events, published ports
//! - [`template`] — the template catalog (control plane, CORE-21)
//! - [`process`] — executions (data plane)
//! - [`filesystem`] — file transfer (data plane)
//! - [`snapshot`] — checkpoint / restore
//!
//! The rest are the daemon's own services, all migrated off tonic
//! (CORE-68): [`icon`], [`kubernetes`], [`machine`], [`migration`],
//! [`stats`], [`system`].
//!
//! Request and response types are buffa-generated (`arcbox-connect`) — the
//! one Rust representation of the ArcBox protos since CORE-73. The same
//! types ride the host↔guest vsock wire, so handlers hand messages between
//! the two surfaces without any twin-codegen re-decode.

mod control;
mod filesystem;
mod icon;
mod kubernetes;
mod machine;
#[cfg(target_os = "macos")]
mod macos;
mod migration;
mod process;
mod sandbox_errors;
mod sandbox_resume;
mod snapshot;
mod stats;
mod stream_input;
mod system;
mod template;

use std::sync::Arc;
use std::time::Duration;

use std::sync::OnceLock;

use arcbox_core::Runtime;
use arcbox_core::vm_lifecycle::DEFAULT_MACHINE_NAME;
use connectrpc::{ConnectError, RequestContext};
use tokio_stream::{Stream, StreamExt as _};

pub use control::SandboxServiceImpl;
pub use filesystem::SandboxFilesystemServiceImpl;
pub use icon::IconServiceImpl;
pub use kubernetes::KubernetesServiceImpl;
pub use machine::MachineServiceImpl;
#[cfg(target_os = "macos")]
pub use macos::MacosServiceImpl;
pub use migration::MigrationServiceImpl;
pub use process::SandboxProcessServiceImpl;
pub use snapshot::SandboxSnapshotServiceImpl;
pub use stats::StatsServiceImpl;
pub use system::{SetupState, SystemServiceImpl};
pub use template::TemplateServiceImpl;

/// Shared handle to a runtime that may not be initialized yet.
///
/// Services are registered before the runtime exists, so each RPC calls
/// `runtime.ready()`, which answers `Unavailable` while the daemon is still
/// downloading assets or starting the VM.
pub type SharedRuntime = Arc<OnceLock<Arc<Runtime>>>;

/// Idle interval after which a server stream emits a keepalive frame, so
/// proxies and load balancers never see a silent connection (CORE-55).
const KEEPALIVE_INTERVAL: Duration = Duration::from_secs(15);

/// Interleave keepalive items whenever `stream` stays idle for
/// [`KEEPALIVE_INTERVAL`].
fn with_keepalive<S, T>(
    stream: S,
    keepalive: fn() -> T,
) -> impl Stream<Item = Result<T, ConnectError>>
where
    S: Stream<Item = Result<T, ConnectError>>,
{
    stream
        .timeout(KEEPALIVE_INTERVAL)
        .map(move |item| item.unwrap_or_else(|_elapsed| Ok(keepalive())))
}

/// Drives a `!Send` macOS VM future to completion on a dedicated blocking
/// thread.
///
/// Virtualization.framework operations hold ObjC handles (and the VM's
/// dispatch queue) across await and are not `Send`, but handler futures must
/// be `Send` — that was true under tonic and is equally true under
/// connectrpc. Running the future via a transient current-thread runtime
/// inside `spawn_blocking` keeps that `!Send` state off the server's worker
/// threads; the booted VM (which is `Send + Sync`) outlives the transient
/// runtime.
#[cfg(target_os = "macos")]
pub(crate) async fn run_macos_blocking<T, Fut, F>(f: F) -> Result<T, ConnectError>
where
    T: Send + 'static,
    Fut: std::future::Future<Output = arcbox_core::Result<T>>,
    F: FnOnce() -> Fut + Send + 'static,
{
    tokio::task::spawn_blocking(move || {
        tokio::runtime::Builder::new_current_thread()
            .enable_all()
            .build()
            .map_err(|e| ConnectError::internal(format!("macOS runtime: {e}")))?
            .block_on(f())
            .map_err(|e| ConnectError::internal(e.to_string()))
    })
    .await
    .map_err(|e| ConnectError::internal(format!("macOS task join: {e}")))?
}

/// Extension trait for obtaining the runtime from a deferred handle.
///
/// The message and the `Unavailable` code match what the tonic services
/// answered before CORE-68, so clients that predate the migration see the
/// same not-ready behaviour they always did.
pub(crate) trait ConnectRuntimeExt {
    /// Returns the runtime, or `Unavailable` if it hasn't been initialized.
    fn ready(&self) -> Result<&Arc<Runtime>, ConnectError>;

    /// Returns the runtime after admitting a write to the target machine's storage.
    fn ready_for_write(&self, machine: &str) -> Result<&Arc<Runtime>, ConnectError>;
}

impl ConnectRuntimeExt for SharedRuntime {
    fn ready(&self) -> Result<&Arc<Runtime>, ConnectError> {
        self.get()
            .ok_or_else(|| ConnectError::unavailable("daemon is starting, runtime not ready yet"))
    }

    fn ready_for_write(&self, machine: &str) -> Result<&Arc<Runtime>, ConnectError> {
        let runtime = self.ready()?;
        runtime
            .ensure_storage_writes_available(machine)
            .map_err(crate::ApiError::from)?;
        Ok(runtime)
    }
}

/// Extension trait for extracting routing metadata from a Connect request.
pub(crate) trait ContextExt {
    /// Returns the target machine name, defaulting to the System VM.
    fn machine_id(&self) -> Result<String, ConnectError>;

    /// Returns the System VM for Sandbox V1, rejecting every other machine.
    fn sandbox_machine_id(&self) -> Result<String, ConnectError>;
}

impl ContextExt for RequestContext {
    /// Reads the optional `x-machine` header.
    ///
    /// Clients may select a local VM through this transport header.
    /// An absent or empty header selects the System VM. Sandbox V1 accepts
    /// only the System VM, as enforced by `sandbox_machine_id`.
    fn machine_id(&self) -> Result<String, ConnectError> {
        match self.header("x-machine") {
            None => Ok(DEFAULT_MACHINE_NAME.to_owned()),
            Some(value) => match value.to_str() {
                Ok("") => Ok(DEFAULT_MACHINE_NAME.to_owned()),
                Ok(s) => Ok(s.to_string()),
                Err(_) => Err(ConnectError::invalid_argument(
                    "invalid x-machine header: must be valid UTF-8",
                )),
            },
        }
    }

    fn sandbox_machine_id(&self) -> Result<String, ConnectError> {
        let machine = self.machine_id()?;
        if machine != DEFAULT_MACHINE_NAME {
            return Err(ConnectError::invalid_argument(
                "Sandbox V1 is available only on the System VM",
            ));
        }
        Ok(machine)
    }
}

/// Map the public port protocol onto the computer layer's enum.
fn port_protocol(
    protocol: arcbox_connect::sandbox_v1::PortProtocol,
) -> arcbox_computer::ports::SandboxPortProtocol {
    use arcbox_computer::ports::SandboxPortProtocol;
    match protocol {
        arcbox_connect::sandbox_v1::PortProtocol::Udp => SandboxPortProtocol::Udp,
        _ => SandboxPortProtocol::Tcp,
    }
}

fn exposed_port(
    mapping: arcbox_core::SandboxPortMapping,
) -> arcbox_connect::sandbox_v1::ExposedPort {
    use arcbox_connect::sandbox_v1::{ExposedPort, PortProtocol};
    use arcbox_core::SandboxPortProtocol;

    ExposedPort {
        sandbox_port: u32::from(mapping.sandbox_port),
        host_port: u32::from(mapping.host_port),
        protocol: match mapping.protocol {
            SandboxPortProtocol::Tcp => PortProtocol::Tcp,
            SandboxPortProtocol::Udp => PortProtocol::Udp,
        }
        .into(),
        ..Default::default()
    }
}

/// Registers every Connect-served service on one router.
///
/// Kept here rather than in the daemon so that adding a service is one edit
/// next to the impls, not a silent omission at the call site — and because
/// registration is what decides whether a path is served at all: the router
/// answers exactly what it was given, and an omission 404s at runtime.
#[must_use]
pub fn router(runtime: SharedRuntime) -> connectrpc::Router {
    let clone = || Arc::clone(&runtime);
    let sandbox_operations = Arc::new(arcbox_computer::locks::SandboxOperationLocks::default());
    let router = connectrpc::Router::new()
        .add_service(Arc::new(SandboxServiceImpl::new(
            clone(),
            Arc::clone(&sandbox_operations),
        )))
        .add_service(Arc::new(SandboxProcessServiceImpl::new(
            clone(),
            Arc::clone(&sandbox_operations),
        )))
        .add_service(Arc::new(SandboxFilesystemServiceImpl::new(
            clone(),
            Arc::clone(&sandbox_operations),
        )))
        .add_service(Arc::new(SandboxSnapshotServiceImpl::new(
            clone(),
            Arc::clone(&sandbox_operations),
        )))
        .add_service(Arc::new(TemplateServiceImpl::new(clone())))
        // The daemon's own services, all on Connect since CORE-68.
        .add_service(Arc::new(IconServiceImpl::new()))
        .add_service(Arc::new(StatsServiceImpl::new(clone())))
        .add_service(Arc::new(KubernetesServiceImpl::new(clone())))
        .add_service(Arc::new(MigrationServiceImpl::new(clone())))
        .add_service(Arc::new(MachineServiceImpl::new(clone())));
    #[cfg(target_os = "macos")]
    let router = router.add_service(Arc::new(MacosServiceImpl::new(clone())));
    router
}

/// Registers the services plus the ones needing extra state the router
/// helper above cannot reach.
///
/// `SystemService` is separate because it also observes the setup state and
/// the early runtime handle — it must answer while `shared_runtime` is still
/// empty, which is the whole point of a diagnostics RPC.
#[must_use]
pub fn router_with_system(runtime: SharedRuntime, system: SystemServiceImpl) -> connectrpc::Router {
    router(runtime).add_service(Arc::new(system))
}

#[cfg(test)]
mod tests {
    use super::*;

    pub(super) fn storage_runtime() -> (std::path::PathBuf, SharedRuntime) {
        let directory = std::env::temp_dir().join(format!(
            "arcbox-api-storage-admission-{}",
            uuid::Uuid::new_v4()
        ));
        std::fs::create_dir_all(&directory).unwrap();
        let runtime = Arc::new(
            Runtime::new(arcbox_core::config::Config {
                data_dir: directory.clone(),
                ..Default::default()
            })
            .unwrap(),
        );
        (directory, Arc::new(OnceLock::from(runtime)))
    }

    #[tokio::test]
    async fn storage_admission_rejects_reserved_and_held_writes_but_preserves_reads() {
        let (directory, shared) = storage_runtime();
        let runtime = shared.ready().unwrap();
        for machine in [DEFAULT_MACHINE_NAME, "rosetta", "dev"] {
            assert!(shared.ready_for_write(machine).is_ok());
        }

        let reservation = runtime.machine_manager().reserve_storage().unwrap();
        for machine in [DEFAULT_MACHINE_NAME, "rosetta"] {
            let error = shared.ready_for_write(machine).err().unwrap();
            assert_eq!(error.code, connectrpc::ErrorCode::FailedPrecondition);
        }
        assert!(shared.ready().is_ok());
        assert!(shared.ready_for_write("dev").is_ok());
        drop(reservation);
        assert!(shared.ready_for_write(DEFAULT_MACHINE_NAME).is_ok());

        let hold = runtime.machine_manager().storage_hold_path();
        std::fs::create_dir_all(hold.parent().unwrap()).unwrap();
        std::fs::write(&hold, "offline-check").unwrap();
        assert!(!runtime.storage_writes_protected());
        for machine in [DEFAULT_MACHINE_NAME, "rosetta"] {
            let error = shared.ready_for_write(machine).err().unwrap();
            assert_eq!(error.code, connectrpc::ErrorCode::FailedPrecondition);
        }
        assert!(shared.ready().is_ok());
        assert!(shared.ready_for_write("dev").is_ok());
        std::fs::remove_file(hold).unwrap();
        assert!(shared.ready_for_write(DEFAULT_MACHINE_NAME).is_ok());
        std::fs::remove_dir_all(directory).unwrap();
    }

    #[tokio::test]
    async fn failed_hold_write_keeps_public_storage_writes_protected() {
        let (directory, shared) = storage_runtime();
        let runtime = shared.ready().unwrap();
        let hold = runtime.machine_manager().storage_hold_path();
        std::fs::create_dir_all(&hold).unwrap();
        assert!(
            runtime
                .recover_storage(arcbox_connect::v1::recover_storage_request::Action::CheckOnly)
                .await
                .is_err()
        );
        std::fs::remove_dir(&hold).unwrap();
        assert!(!runtime.storage_recovery_active());
        assert!(runtime.storage_writes_protected());
        assert!(
            runtime
                .machine_manager()
                .ensure_storage_available(DEFAULT_MACHINE_NAME)
                .is_err()
        );
        for machine in [DEFAULT_MACHINE_NAME, "rosetta"] {
            let error = shared.ready_for_write(machine).err().unwrap();
            assert_eq!(error.code, connectrpc::ErrorCode::FailedPrecondition);
        }
        assert!(shared.ready().is_ok());
        assert!(shared.ready_for_write("dev").is_ok());
        assert!(matches!(
            runtime.trim_machine_disk(DEFAULT_MACHINE_NAME).await,
            Err(arcbox_core::CoreError::Common(
                arcbox_error::CommonError::InvalidState(_)
            ))
        ));
        std::fs::remove_dir_all(directory).unwrap();
    }

    #[tokio::test]
    async fn resume_rechecks_storage_admission_after_waiting_for_the_operation_lock() {
        let (directory, shared) = storage_runtime();
        let runtime = shared.ready().unwrap();
        let operations = arcbox_computer::locks::SandboxOperationLocks::default();
        let operation = operations.lock(DEFAULT_MACHINE_NAME, "sandbox").await;
        let resumed = sandbox_resume::resume(
            runtime,
            &operations,
            DEFAULT_MACHINE_NAME,
            "sandbox",
            sandbox_resume::REASON_RESUME,
        );
        tokio::pin!(resumed);
        tokio::select! {
            biased;
            result = &mut resumed => panic!("resume must wait for its operation lock: {result:?}"),
            () = tokio::task::yield_now() => {}
        }
        let reservation = runtime.machine_manager().reserve_storage().unwrap();
        drop(operation);

        let error = resumed.await.unwrap_err();
        assert_eq!(error.code, connectrpc::ErrorCode::FailedPrecondition);
        drop(reservation);
        std::fs::remove_dir_all(directory).unwrap();
    }

    fn ctx_with(header: Option<&str>) -> RequestContext {
        let mut headers = http::HeaderMap::new();
        if let Some(value) = header {
            headers.insert("x-machine", value.parse().expect("valid header value"));
        }
        RequestContext::new(headers)
    }

    /// `x-machine` is optional transport metadata (CORE-54): absent or
    /// empty must resolve to the System VM, not an error.
    #[test]
    fn machine_id_defaults_to_the_system_vm() {
        assert_eq!(
            ctx_with(None).machine_id().expect("absent header is valid"),
            DEFAULT_MACHINE_NAME
        );
        assert_eq!(
            ctx_with(Some(""))
                .machine_id()
                .expect("empty header is valid"),
            DEFAULT_MACHINE_NAME
        );
        assert_eq!(
            ctx_with(Some("other-vm"))
                .machine_id()
                .expect("explicit header is valid"),
            "other-vm"
        );
    }

    #[test]
    fn sandbox_machine_id_accepts_only_the_system_vm() {
        assert_eq!(DEFAULT_MACHINE_NAME, "default");
        for header in [None, Some(""), Some(DEFAULT_MACHINE_NAME)] {
            assert_eq!(
                ctx_with(header)
                    .sandbox_machine_id()
                    .expect("System VM routing should be accepted"),
                DEFAULT_MACHINE_NAME
            );
        }
        let error = ctx_with(Some("other-vm"))
            .sandbox_machine_id()
            .expect_err("other machines must be rejected");
        assert_eq!(error.code, connectrpc::ErrorCode::InvalidArgument);
    }

    #[test]
    fn exposed_port_preserves_the_authoritative_mapping() {
        let port = exposed_port(arcbox_core::SandboxPortMapping {
            sandbox_port: 8080,
            host_port: 45_000,
            protocol: arcbox_core::SandboxPortProtocol::Udp,
        });
        assert_eq!(port.sandbox_port, 8080);
        assert_eq!(port.host_port, 45_000);
        assert_eq!(
            port.protocol.as_known(),
            Some(arcbox_connect::sandbox_v1::PortProtocol::Udp)
        );
    }
}