LazyBoy2/crates/lazyboy-core/src/box_runtime.rs

1475 lines
52 KiB
Rust

//! Named-agent Linux computers (Docker boxes).
//! Session CLI keeps `lazyboy-box`. Each named agent gets its own container.
use crate::computer::{
settle_before, steps_for, validate_bounds, ComputerAction, ComputerStep, ScreenSize,
};
use anyhow::{anyhow, Context, Result};
use serde_json::{json, Value};
use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::process::Stdio;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, OnceLock};
use std::time::{Duration, Instant};
use tokio::process::Command;
use tokio::sync::Mutex;
const IMAGE: &str = "lazyboy-box:local";
const SESSION_CONTAINER: &str = "lazyboy-box";
const SESSION_WORKSPACE_VOLUME: &str = "grokboy-box-workspace";
const SESSION_HOME_VOLUME: &str = "grokboy-box-home";
const SESSION_VIEWER_PORT: u16 = 6080;
pub const SESSION_SEAT: &str = "session";
const BOX_REVISION: &str = "computer-use-7";
const BROWSER_PROFILE: &str = "/home/box/chrome-profile";
/// Per-stream cap on shell output returned to the model; the middle is elided.
const JOB_OUTPUT_HEAD_CHARS: usize = 16_000;
const JOB_OUTPUT_TAIL_CHARS: usize = 32_000;
const DEFAULT_MAX_RUNNING: usize = 3;
const COMPUTER_MEMORY: &str = "1536m";
const COMPUTER_CPUS: &str = "2";
const COMPUTER_PIDS: &str = "512";
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct BoxSeat {
pub id: String,
pub container: String,
pub workspace_volume: String,
pub home_volume: String,
pub bind_port: Option<u16>,
}
impl BoxSeat {
pub fn for_id(id: &str) -> Self {
let id = id.trim();
if id.is_empty() || id == SESSION_SEAT {
return Self {
id: SESSION_SEAT.into(),
container: SESSION_CONTAINER.into(),
workspace_volume: SESSION_WORKSPACE_VOLUME.into(),
home_volume: SESSION_HOME_VOLUME.into(),
bind_port: Some(SESSION_VIEWER_PORT),
};
}
let safe = sanitize_seat_id(id);
Self {
id: id.to_string(),
container: format!("lazyboy-box-{safe}"),
workspace_volume: format!("lazyboy-box-workspace-{safe}"),
home_volume: format!("lazyboy-box-home-{safe}"),
bind_port: None,
}
}
pub fn web_viewer_url(&self) -> String {
if self.id == SESSION_SEAT {
"/novnc/vnc.html?autoconnect=true&resize=off&reconnect=true&show_dot=false&path=novnc/websockify"
.into()
} else {
let id = &self.id;
format!(
"/novnc/{id}/vnc.html?autoconnect=true&resize=off&reconnect=true&show_dot=false&path=novnc/{id}/websockify"
)
}
}
}
pub fn sanitize_seat_id(id: &str) -> String {
let mapped: String = id
.chars()
.map(|c| {
if c.is_ascii_alphanumeric() || c == '-' || c == '_' {
c
} else {
'-'
}
})
.collect();
let trimmed = mapped.trim_matches('-');
if trimmed.is_empty() {
"agent".into()
} else {
trimmed.chars().take(64).collect()
}
}
pub struct BoxHub {
seat: BoxSeat,
inner: Mutex<BoxState>,
computer_use: AtomicBool,
}
struct BoxState {
ready: bool,
starting: bool,
last_error: Option<String>,
last_used: Instant,
host_port: Option<u16>,
}
impl BoxHub {
pub fn new() -> Arc<Self> {
Self::for_seat(BoxSeat::for_id(SESSION_SEAT))
}
pub fn for_seat(seat: BoxSeat) -> Arc<Self> {
Arc::new(Self {
seat,
inner: Mutex::new(BoxState {
ready: false,
starting: false,
last_error: None,
last_used: Instant::now(),
host_port: None,
}),
computer_use: AtomicBool::new(false),
})
}
pub fn seat_id(&self) -> &str {
&self.seat.id
}
pub fn container(&self) -> &str {
&self.seat.container
}
pub fn web_viewer_url(&self) -> String {
self.seat.web_viewer_url()
}
pub fn viewer_url() -> String {
format!("http://127.0.0.1:{SESSION_VIEWER_PORT}/vnc.html?autoconnect=true&resize=scale")
}
pub fn host_viewer_url_for_port(port: u16) -> String {
format!("http://127.0.0.1:{port}/vnc.html?autoconnect=true&resize=scale")
}
pub fn try_begin_computer_use(&self) -> Result<()> {
self.computer_use
.compare_exchange(false, true, Ordering::SeqCst, Ordering::SeqCst)
.map_err(|_| {
anyhow!(
"A computerUse subagent is already using this agent's desktop. Only one can run at a time on the same computer."
)
})?;
Ok(())
}
pub fn end_computer_use(&self) {
self.computer_use.store(false, Ordering::SeqCst);
}
pub fn has_computer_use(&self) -> bool {
self.computer_use.load(Ordering::SeqCst)
}
fn ready_payload(&self) -> Value {
json!({
"ready": true,
"state": "ready",
"viewer_url": self.web_viewer_url(),
"host_viewer_url": Self::host_viewer_url_for_port(
self.seat.bind_port.unwrap_or(SESSION_VIEWER_PORT)
),
"browser_surface": "docker",
"workspace": "/workspace",
"browser_profile": BROWSER_PROFILE,
"revision": BOX_REVISION,
"seat": self.seat.id,
"container": self.seat.container,
})
}
fn status_payload(&self, state: &str, ready: bool, error: Option<String>, crowded: bool) -> Value {
let mut payload = json!({
"ready": ready,
"state": state,
"viewer_url": self.web_viewer_url(),
"browser_surface": "docker",
"workspace": "/workspace",
"browser_profile": BROWSER_PROFILE,
"revision": BOX_REVISION,
"seat": self.seat.id,
"container": self.seat.container,
"crowded": crowded,
});
if let Some(error) = error {
payload["error"] = json!(error);
}
payload
}
/// Read-only. Never starts the container.
pub async fn inspect(&self) -> Value {
let (starting, last_error) = {
let state = self.inner.lock().await;
(state.starting, state.last_error.clone())
};
if starting {
return self.status_payload("starting", false, None, false);
}
match docker_running(&self.seat.container).await {
Ok(true) => {
if desktop_up(&self.seat.container).await {
let mut state = self.inner.lock().await;
state.ready = true;
if let Some(port) = host_port_of(&self.seat.container).await {
state.host_port = Some(port);
}
let mut payload = self.ready_payload();
if let Some(port) = state.host_port {
payload["host_viewer_url"] = json!(Self::host_viewer_url_for_port(port));
}
payload
} else {
self.status_payload("starting", false, None, false)
}
}
Ok(false) => {
if let Some(error) = last_error {
self.status_payload("error", false, Some(error), false)
} else {
self.status_payload("stopped", false, None, false)
}
}
Err(error) => self.status_payload("error", false, Some(error.to_string()), false),
}
}
pub async fn ensure_ready(&self) -> Result<Value> {
{
let mut state = self.inner.lock().await;
if state.ready && docker_running(&self.seat.container).await.unwrap_or(false) {
if !desktop_up(&self.seat.container).await {
state.starting = true;
drop(state);
if let Err(error) = wait_desktop(&self.seat.container).await {
let mut state = self.inner.lock().await;
state.starting = false;
state.ready = false;
state.last_error = Some(error.to_string());
return Err(error);
}
let mut state = self.inner.lock().await;
state.starting = false;
state.ready = true;
state.last_used = Instant::now();
state.last_error = None;
return Ok(self.finish_payload(&mut state));
}
state.last_used = Instant::now();
return Ok(self.finish_payload(&mut state));
}
state.starting = true;
state.last_error = None;
}
let result = self.bring_up().await;
let mut state = self.inner.lock().await;
state.starting = false;
match result {
Ok(()) => {
state.ready = true;
state.last_used = Instant::now();
state.last_error = None;
Ok(self.finish_payload(&mut state))
}
Err(error) => {
state.ready = false;
state.last_error = Some(error.to_string());
Err(error)
}
}
}
async fn bring_up(&self) -> Result<()> {
docker_info().await?;
let _provision = ProvisionLock::acquire().await?;
ensure_image().await?;
let pool = BoxPool::global();
pool.reconcile().await?;
pool.reclaim(self.seat_id()).await?;
let port = ensure_container(&self.seat).await?;
wait_desktop(&self.seat.container).await?;
self.inner.lock().await.host_port = Some(port);
Ok(())
}
fn finish_payload(&self, state: &mut BoxState) -> Value {
let mut payload = self.ready_payload();
payload["instruction"] = json!(
"This is my computer. Paths here are /workspace and /home/box, not the user's machine."
);
if let Some(port) = state.host_port.or(self.seat.bind_port) {
payload["host_viewer_url"] = json!(Self::host_viewer_url_for_port(port));
}
payload
}
/// Reboot the existing container. Volumes (workspace + Chrome profile) stay.
pub async fn restart(&self) -> Result<Value> {
{
let mut state = self.inner.lock().await;
state.ready = false;
state.starting = true;
state.last_error = None;
}
let result = async {
docker_info().await?;
let _provision = ProvisionLock::acquire().await?;
let pool = BoxPool::global();
pool.reconcile().await?;
pool.reclaim(self.seat_id()).await?;
if docker(["container", "inspect", self.container()]).await.is_ok() {
limit_container(self.container()).await?;
docker(["restart", "-t", "20", self.container()]).await?;
} else {
ensure_image().await?;
let port = ensure_container(&self.seat).await?;
self.inner.lock().await.host_port = Some(port);
}
wait_desktop(self.container()).await?;
if let Some(port) = host_port_of(self.container()).await {
self.inner.lock().await.host_port = Some(port);
}
Ok(())
}
.await;
self.finish_action(result, "restart").await
}
/// Rebuild the image from the repo Dockerfile. Recreate this seat's
/// container only when the image id actually changed. Volumes stay.
pub async fn update(&self) -> Result<Value> {
{
let mut state = self.inner.lock().await;
state.ready = false;
state.starting = true;
state.last_error = None;
}
let result = async {
docker_info().await?;
let _provision = ProvisionLock::acquire().await?;
let pool = BoxPool::global();
pool.reconcile().await?;
pool.reclaim(self.seat_id()).await?;
let before = local_image_id().await;
build_image().await?;
let after = local_image_id().await;
let image_changed = before != after;
if image_changed {
if docker(["container", "inspect", self.container()]).await.is_ok() {
remove_stale_container(self.container()).await?;
}
let port = ensure_container(&self.seat).await?;
self.inner.lock().await.host_port = Some(port);
} else if !docker_running(self.container()).await? {
let port = ensure_container(&self.seat).await?;
self.inner.lock().await.host_port = Some(port);
}
wait_desktop(self.container()).await?;
Ok(image_changed)
}
.await;
match result {
Ok(image_changed) => {
let mut payload = self.finish_action(Ok(()), "update").await?;
payload["updated"] = json!(image_changed);
Ok(payload)
}
Err(error) => self.finish_action(Err(error), "update").await,
}
}
async fn finish_action(&self, result: Result<()>, action: &str) -> Result<Value> {
let mut state = self.inner.lock().await;
state.starting = false;
match result {
Ok(()) => {
state.ready = true;
state.last_used = Instant::now();
state.last_error = None;
let mut payload = self.finish_payload(&mut state);
payload["action"] = json!(action);
Ok(payload)
}
Err(error) => {
state.ready = false;
state.last_error = Some(error.to_string());
Err(error)
}
}
}
pub async fn stop_keep_volumes(&self) -> Result<()> {
let mut state = self.inner.lock().await;
state.ready = false;
state.starting = false;
drop(state);
if docker_running(self.container()).await.unwrap_or(false) {
docker(["stop", "-t", "20", self.container()]).await?;
}
Ok(())
}
pub async fn wipe(&self) -> Result<()> {
let _ = self.stop_keep_volumes().await;
if docker(["container", "inspect", self.container()]).await.is_ok() {
let _ = docker(["rm", "-f", self.container()]).await;
}
let _ = docker(["volume", "rm", "-f", &self.seat.workspace_volume]).await;
let _ = docker(["volume", "rm", "-f", &self.seat.home_volume]).await;
let mut state = self.inner.lock().await;
state.ready = false;
state.host_port = None;
state.last_error = None;
Ok(())
}
async fn last_used(&self) -> Instant {
self.inner.lock().await.last_used
}
pub async fn host_port(&self) -> Option<u16> {
if let Some(port) = self.inner.lock().await.host_port {
return Some(port);
}
host_port_of(self.container()).await
}
pub async fn ensure_browser_ready(&self) -> Result<()> {
self.ensure_ready().await?;
let probe = ["python3", "-c", "import urllib.request; urllib.request.urlopen('http://127.0.0.1:9222/json/version', timeout=1).close()"];
if self.exec(&probe).await?.status == 0 { return Ok(()); }
// The desktop autostart may already be launching Chromium; a second
// launcher would race it for the profile. box-chrome serialises
// launches, but skipping the launch entirely is cheaper.
let already_launching = self.exec(&[
"pgrep", "-f", "--", &format!("--user-data-dir={BROWSER_PROFILE}"),
]).await?.status == 0;
if !already_launching {
self.exec(&["bash", "-lc", "nohup box-chrome </dev/null >/tmp/lazyboy-browser.log 2>&1 &"]).await?;
}
for _ in 0..80 {
if self.exec(&probe).await?.status == 0 { return Ok(()); }
tokio::time::sleep(Duration::from_millis(250)).await;
}
Err(anyhow!("Box Chromium did not start. Inspect /tmp/lazyboy-browser.log on the box; no local browser was opened."))
}
pub async fn shell(&self, cmd: &str, block_until_ms: u64) -> Result<Value> {
self.ensure_ready().await?;
let id = uuid::Uuid::new_v4().to_string();
// All commands run detached in the box. A foreground wait expiring must
// never kill an install or silently move it onto the user's computer.
// The exit file is renamed into place so a poller never sees it empty.
self.exec(&["bash", "-lc", &format!(
"mkdir -p /tmp/gb-jobs/{id} && {{ nohup bash -lc {q} </dev/null >/tmp/gb-jobs/{id}/stdout 2>/tmp/gb-jobs/{id}/stderr & }}",
q = shell_quote(&format!(
"cd /workspace && bash -lc {}; code=$?; printf '%s\\n' \"$code\" >/tmp/gb-jobs/{id}/exit.tmp && mv -f /tmp/gb-jobs/{id}/exit.tmp /tmp/gb-jobs/{id}/exit",
shell_quote(cmd)
))
)]).await?;
self.wait_job(&id, block_until_ms).await
}
async fn wait_job(&self, id: &str, block_until_ms: u64) -> Result<Value> {
let deadline = tokio::time::Instant::now() + Duration::from_millis(block_until_ms.min(30_000));
loop {
let mut status = self.read_job_status(id).await?;
status["session_id"] = json!(id);
status["surface"] = json!("box");
if status["running"] != true || tokio::time::Instant::now() >= deadline {
if status["running"] == true {
status["instruction"] = json!("Still running on my computer. Use await_shell with this session_id when the result is needed. Do not rerun the command locally.");
}
return Ok(status);
}
tokio::time::sleep(Duration::from_millis(200)).await;
}
}
pub async fn await_job(&self, session_id: &str) -> Result<Value> {
self.ensure_ready().await?;
validate_job_id(session_id)?;
self.wait_job(session_id, 30_000).await
}
async fn read_job_status(&self, session_id: &str) -> Result<Value> {
validate_job_id(session_id)?;
let py = r#"
import json, pathlib, sys
root = pathlib.Path('/tmp/gb-jobs') / sys.argv[1]
head, tail = int(sys.argv[2]), int(sys.argv[3])
def stream(name):
p = root / name
if not p.exists():
return '', False
text = p.read_text(errors='replace')
if len(text) <= head + tail:
return text, False
dropped = len(text) - head - tail
return text[:head] + f'\n... [{dropped} chars elided] ...\n' + text[-tail:], True
if not root.exists():
print(json.dumps({'error': 'unknown session'}))
else:
exitp = root / 'exit'
raw = exitp.read_text(errors='replace').strip() if exitp.exists() else ''
code = int(raw) if raw.lstrip('-').isdigit() else None
out, out_cut = stream('stdout')
err, err_cut = stream('stderr')
status = {'running': code is None, 'exit_code': code, 'stdout': out, 'stderr': err}
if out_cut or err_cut:
status['truncated'] = True
status['full_output'] = {'stdout': str(root / 'stdout'), 'stderr': str(root / 'stderr')}
print(json.dumps(status))
"#;
let out = self.exec(&[
"python3",
"-c",
py,
session_id,
&JOB_OUTPUT_HEAD_CHARS.to_string(),
&JOB_OUTPUT_TAIL_CHARS.to_string(),
])
.await?;
if out.status != 0 {
return Err(anyhow!("box job status failed: {}", out.stderr.trim()));
}
serde_json::from_str(out.stdout.trim()).context("parse box_await json")
}
pub async fn read(&self, path: &str, offset: Option<i64>, limit: Option<i64>) -> Result<Value> {
self.ensure_ready().await?;
let path = resolve_box_path(path)?;
let offset = offset.unwrap_or(1);
let limit = limit.unwrap_or(400);
let py = format!(
r#"
import json, pathlib, sys
p = pathlib.Path(sys.argv[1])
lines = p.read_text(errors='replace').splitlines()
start = max({offset}-1, 0)
chunk = lines[start:start+{limit}]
print(json.dumps({{'path': str(p), 'offset': start+1, 'lines': len(lines), 'content': chr(10).join(chunk)}}))
"#
);
let out = self.exec(&["python3", "-c", &py, &path]).await?;
if out.status != 0 {
return Err(anyhow!(out.stderr.trim().to_string()));
}
serde_json::from_str(out.stdout.trim()).context("parse box_read")
}
pub async fn copy_to_box(&self, host_path: &Path, box_path: Option<&str>) -> Result<Value> {
self.ensure_ready().await?;
if !host_path.is_file() {
return Err(anyhow!("source is not a file: {}", host_path.display()));
}
let dest = match box_path {
Some(p) => resolve_box_path(p)?,
None => format!(
"/workspace/uploads/{}",
host_path
.file_name()
.and_then(|s| s.to_str())
.unwrap_or("file")
),
};
docker(["exec", self.container(), "mkdir", "-p", parent_dir(&dest)]).await?;
docker([
"cp",
&host_path.to_string_lossy(),
&format!("{}:{dest}", self.container()),
])
.await?;
Ok(json!({"box_path": dest, "bytes": std::fs::metadata(host_path)?.len()}))
}
pub async fn copy_from_box(&self, box_path: &str, host_path: &Path) -> Result<Value> {
self.ensure_ready().await?;
let src = resolve_box_path(box_path)?;
if let Some(parent) = host_path.parent() {
std::fs::create_dir_all(parent)?;
}
docker([
"cp",
&format!("{}:{src}", self.container()),
&host_path.to_string_lossy(),
])
.await?;
let is_dir = host_path.is_dir();
Ok(json!({
"computer_path": host_path,
"is_dir": is_dir,
"bytes": if is_dir { dir_size(host_path) } else { std::fs::metadata(host_path).map(|m| m.len()).unwrap_or(0) },
}))
}
pub async fn screenshot(&self, dest: &Path) -> Result<Value> {
self.ensure_ready().await?;
let shot = ShotFile::new();
self.capture_root_png(&shot).await?;
shot.fetch(self.container(), dest).await?;
Ok(json!({
"path": dest,
"viewer_url": self.web_viewer_url(),
}))
}
pub async fn computer(&self, actions: &[ComputerAction], dest: &Path) -> Result<Value> {
self.ensure_ready().await?;
let screen = self.display_geometry().await?;
validate_bounds(actions, screen)?;
let names: Vec<&str> = actions.iter().map(ComputerAction::name).collect();
let shot = ShotFile::new();
let mut snapshot = None;
for (index, action) in actions.iter().enumerate() {
let settle = settle_before(actions, index);
if settle > 0 {
tokio::time::sleep(Duration::from_millis(u64::from(settle))).await;
}
self.run_computer_action(action, &shot, screen, &mut snapshot).await?;
}
if !actions.iter().any(ComputerAction::is_screenshot) {
self.capture_root_png(&shot).await?;
}
shot.fetch(self.container(), dest).await?;
let cursor = self.cursor_position().await.ok();
let mut payload = json!({
"ok": true,
"actions": names,
"path": dest,
"viewer_url": self.web_viewer_url(),
"cursor_position": cursor,
"screen": {"width": screen.width, "height": screen.height},
"instruction": "Computer action ran on my computer. Read this screenshot before the next action. A batched then sequence returns only this final screen."
});
if let Some(text) = snapshot {
payload["snapshot"] = json!(text);
payload["instruction"] = json!(
"Computer action ran on my computer. `snapshot` lists the on-screen elements as `[id] role \"name\" @ (x,y) WxH`; click the (x,y) centre of the element you mean. Read the screenshot too before the next action."
);
}
Ok(payload)
}
async fn run_computer_action(
&self,
action: &ComputerAction,
shot: &ShotFile,
screen: ScreenSize,
snapshot: &mut Option<String>,
) -> Result<()> {
for step in steps_for(action)? {
match step {
ComputerStep::Screenshot => self.capture_root_png(shot).await?,
ComputerStep::A11ySnapshot => {
let width = screen.width.to_string();
let height = screen.height.to_string();
let out = self
.exec_env(
&[
("DISPLAY", ":1"),
("LAZYBOY_SCREEN_W", width.as_str()),
("LAZYBOY_SCREEN_H", height.as_str()),
],
&["box-a11y"],
)
.await?;
if out.status != 0 {
let detail = out.stderr.trim();
if detail.contains("not found") {
return Err(anyhow!(
"My computer image is missing box-a11y. Rebuild the box image and retry."
));
}
return Err(anyhow!(
"computer snapshot failed: {}",
if detail.is_empty() { out.stdout.trim() } else { detail }
));
}
*snapshot = Some(out.stdout.trim_end().to_string());
}
ComputerStep::SleepMs(0) => {}
ComputerStep::SleepMs(ms) => {
let secs = format!("{:.3}", f64::from(ms) / 1000.0);
let out = self.exec(&[
"python3",
"-c",
"import sys,time; time.sleep(float(sys.argv[1]))",
&secs,
])
.await?;
if out.status != 0 {
return Err(anyhow!("computer wait failed: {}", out.stderr.trim()));
}
}
ComputerStep::Xdotool(args) => {
let mut argv: Vec<&str> = vec!["xdotool"];
argv.extend(args.iter().map(String::as_str));
let out = self.exec_env(&[("DISPLAY", ":1")], &argv).await?;
if out.status != 0 {
let detail = out.stderr.trim();
if detail.contains("xdotool") && detail.contains("not found") {
return Err(anyhow!(
"My computer image is missing xdotool. Rebuild the box image and retry."
));
}
return Err(anyhow!(
"computer {} failed: {}",
action.name(),
if detail.is_empty() {
out.stdout.trim()
} else {
detail
}
));
}
}
}
}
Ok(())
}
async fn display_geometry(&self) -> Result<ScreenSize> {
let out = self.exec_env(&[("DISPLAY", ":1")], &["xdotool", "getdisplaygeometry"]).await?;
if out.status == 0 {
if let Ok(size) = ScreenSize::parse(&out.stdout) {
return Ok(size);
}
}
let info = self.exec(&["bash", "-lc", "xdpyinfo -display :1 | awk '/dimensions:/{print $2}'"]).await?;
let raw = info.stdout.trim().replace('x', " ");
ScreenSize::parse(&raw).map_err(|e| {
anyhow!("could not read my computer's display size ({e}). Is the desktop up?")
})
}
async fn cursor_position(&self) -> Result<Value> {
let out = self.exec_env(&[("DISPLAY", ":1")], &["xdotool", "getmouselocation", "--shell"]).await?;
if out.status != 0 {
return Err(anyhow!(out.stderr.trim().to_string()));
}
let mut x = 0i32;
let mut y = 0i32;
for line in out.stdout.lines() {
if let Some(v) = line.strip_prefix("X=") {
x = v.trim().parse().unwrap_or(0);
}
if let Some(v) = line.strip_prefix("Y=") {
y = v.trim().parse().unwrap_or(0);
}
}
Ok(json!({"x": x, "y": y}))
}
async fn capture_root_png(&self, shot: &ShotFile) -> Result<()> {
let out = self.exec(&[
"import",
"-display",
":1",
"-window",
"root",
&shot.box_path,
])
.await?;
if out.status != 0 {
return Err(anyhow!(
"Could not capture my computer's screen (desktop not up yet). {}",
out.stderr.trim()
));
}
Ok(())
}
async fn exec(&self, args: &[&str]) -> Result<CmdOut> {
docker_exec(self.container(), args).await
}
async fn exec_env(&self, env: &[(&str, &str)], args: &[&str]) -> Result<CmdOut> {
docker_exec_env(self.container(), env, args).await
}
}
/// A per-call screenshot file on the box. Concurrent callers (a parent
/// checking in on a computerUse child) must never read each other's frame.
struct ShotFile {
box_path: String,
}
impl ShotFile {
fn new() -> Self {
Self {
box_path: format!("/tmp/lazyboy-shot-{}.png", uuid::Uuid::new_v4()),
}
}
async fn fetch(&self, container: &str, dest: &Path) -> Result<()> {
if let Some(parent) = dest.parent() {
std::fs::create_dir_all(parent)?;
}
let copied = docker([
"cp",
&format!("{container}:{}", self.box_path),
&dest.to_string_lossy(),
])
.await;
let _ = docker_exec(container, &["rm", "-f", &self.box_path]).await;
copied.map(|_| ())
}
}
/// Host-wide advisory lock held while the image/container are provisioned.
/// Released when dropped (the fd closes).
struct ProvisionLock {
_file: std::fs::File,
}
impl ProvisionLock {
async fn acquire() -> Result<Self> {
let path = std::env::temp_dir().join("lazyboy-box-provision.lock");
let file = tokio::task::spawn_blocking(move || -> Result<std::fs::File> {
let file = std::fs::OpenOptions::new()
.create(true)
.truncate(false)
.write(true)
.open(&path)
.with_context(|| format!("open {}", path.display()))?;
#[cfg(unix)]
{
use std::os::fd::AsRawFd;
if unsafe { libc::flock(file.as_raw_fd(), libc::LOCK_EX) } != 0 {
return Err(std::io::Error::last_os_error())
.context("lock box provisioning");
}
}
Ok(file)
})
.await
.context("box provisioning lock task")??;
Ok(Self { _file: file })
}
}
fn dir_size(path: &Path) -> u64 {
let Ok(entries) = std::fs::read_dir(path) else {
return 0;
};
entries
.flatten()
.map(|entry| {
let path = entry.path();
if path.is_dir() {
dir_size(&path)
} else {
std::fs::metadata(&path).map(|m| m.len()).unwrap_or(0)
}
})
.sum()
}
pub fn resolve_box_path(path: &str) -> Result<String> {
let path = path.trim();
if path.is_empty() {
return Err(anyhow!("box path is empty"));
}
let abs = if path.starts_with('/') {
path.to_string()
} else {
format!("/workspace/{path}")
};
// Normalise by component so `a//b` and `./x` are accepted while only a
// real `..` segment (not a name like `report..final.txt`) is rejected.
let mut parts: Vec<&str> = Vec::new();
for part in abs.split('/') {
match part {
"" | "." => {}
".." => return Err(anyhow!("box path must not contain ..")),
other => parts.push(other),
}
}
let clean = format!("/{}", parts.join("/"));
let allowed = ["/workspace", "/home/box", "/tmp/gb-jobs"];
if !allowed
.iter()
.any(|root| clean == *root || clean.starts_with(&format!("{root}/")))
{
return Err(anyhow!(
"box paths must be under /workspace or /home/box (got {clean}). This is my computer, not the user's."
));
}
Ok(clean)
}
fn parent_dir(path: &str) -> &str {
path.rsplit_once('/')
.map(|(p, _)| if p.is_empty() { "/" } else { p })
.unwrap_or("/workspace")
}
fn shell_quote(s: &str) -> String {
format!("'{}'", s.replace('\'', r#"'"'"'"#))
}
fn validate_job_id(id: &str) -> Result<()> {
if id.chars().all(|c| c.is_ascii_hexdigit() || c == '-') && (8..80).contains(&id.len()) {
Ok(())
} else {
Err(anyhow!("invalid session_id"))
}
}
fn box_context_dir() -> PathBuf {
if let Some(dir) = std::env::var_os("LAZYBOY_BOX_DIR").filter(|s| !s.is_empty()) {
return PathBuf::from(dir);
}
PathBuf::from(env!("CARGO_MANIFEST_DIR"))
.join("../..")
.join("box")
}
async fn docker_info() -> Result<()> {
if docker(["info"]).await.is_err() {
return Err(anyhow!(
"My computer needs Docker. Start Docker Desktop (or the docker daemon) and try again. I will not run this on your computer."
));
}
Ok(())
}
async fn image_revision() -> Option<String> {
docker([
"inspect",
"-f",
"{{index .Config.Labels \"lazyboy.box.revision\"}}",
IMAGE,
])
.await
.ok()
.map(|out| out.stdout.trim().to_string())
.filter(|s| !s.is_empty() && s != "<no value>")
}
async fn ensure_image() -> Result<()> {
if image_revision().await.as_deref() == Some(BOX_REVISION) {
return Ok(());
}
build_image().await
}
async fn build_image() -> Result<()> {
let ctx = box_context_dir();
if !ctx.join("Dockerfile").is_file() {
return Err(anyhow!(
"Box image {IMAGE} is missing and Dockerfile was not found at {}",
ctx.display()
));
}
// Stage the helper in the build context; never mount the user's workspace.
let helper_dir = ctx.join("playwright");
tokio::fs::create_dir_all(&helper_dir).await?;
for file in ["package.json", "browser_helper.mjs"] {
let src = ctx.join("../tools/playwright").join(file);
tokio::fs::copy(&src, helper_dir.join(file))
.await
.with_context(|| format!("stage browser helper {}", src.display()))?;
}
eprintln!(
"box: building image {IMAGE} (revision {BOX_REVISION}) from {}; this can take several minutes on first run",
ctx.display()
);
let started = std::time::Instant::now();
docker(["build", "-t", IMAGE, &ctx.to_string_lossy()]).await?;
eprintln!("box: image {IMAGE} ready in {:?}", started.elapsed());
Ok(())
}
async fn container_image_id(container: &str) -> Option<String> {
docker(["inspect", "-f", "{{.Image}}", container])
.await
.ok()
.map(|out| out.stdout.trim().to_string())
.filter(|s| !s.is_empty())
}
async fn host_port_of(container: &str) -> Option<u16> {
let out = docker([
"inspect",
"-f",
"{{(index (index .NetworkSettings.Ports \"6080/tcp\") 0).HostPort}}",
container,
])
.await
.ok()?;
out.stdout.trim().parse().ok()
}
async fn local_image_id() -> Option<String> {
docker(["inspect", "-f", "{{.Id}}", IMAGE])
.await
.ok()
.map(|out| out.stdout.trim().to_string())
.filter(|s| !s.is_empty())
}
// Apply to existing computers too, without recreating them or discarding data.
async fn limit_container(container: &str) -> Result<()> {
docker(["update", "--memory", COMPUTER_MEMORY, "--memory-swap", COMPUTER_MEMORY,
"--cpus", COMPUTER_CPUS, "--pids-limit", COMPUTER_PIDS, container]).await?;
Ok(())
}
async fn ensure_container(seat: &BoxSeat) -> Result<u16> {
if docker(["container", "inspect", &seat.container]).await.is_ok() {
limit_container(&seat.container).await?;
}
if docker_running(&seat.container).await?
&& container_image_id(&seat.container).await == local_image_id().await
{
if let Some(port) = host_port_of(&seat.container).await.or(seat.bind_port) {
return Ok(port);
}
}
if docker(["container", "inspect", &seat.container]).await.is_ok()
&& container_image_id(&seat.container).await != local_image_id().await
{
remove_stale_container(&seat.container).await?;
}
if docker_running(&seat.container).await? {
if let Some(port) = host_port_of(&seat.container).await.or(seat.bind_port) {
return Ok(port);
}
}
if seat.id == SESSION_SEAT && docker(["container", "inspect", "grokboy-box"]).await.is_ok() {
let _ = docker(["rm", "-f", "grokboy-box"]).await;
}
let publish = match seat.bind_port {
Some(port) => format!("127.0.0.1:{port}:6080"),
None => "127.0.0.1::6080".into(),
};
// The container name is Docker's cross-process creation lock. Losing a
// create race must reuse the winner, never remove its container.
let created = docker([
"create",
"--name",
&seat.container,
"--memory", COMPUTER_MEMORY,
"--memory-swap", COMPUTER_MEMORY,
"--cpus", COMPUTER_CPUS,
"--pids-limit", COMPUTER_PIDS,
"--shm-size=256m",
"-p",
&publish,
"--stop-timeout=20",
"-v",
&format!("{}:/workspace", seat.workspace_volume),
"-v",
&format!("{}:/home/box", seat.home_volume),
IMAGE,
])
.await;
if let Err(error) = created {
// Docker reserves the name before inspect can see the new container.
let mut visible = false;
for _ in 0..50 {
if docker(["container", "inspect", &seat.container]).await.is_ok() {
visible = true;
break;
}
tokio::time::sleep(Duration::from_millis(100)).await;
}
if !visible {
return Err(error);
}
}
// Starting an existing container also preserves its writable layer.
docker(["start", &seat.container]).await?;
host_port_of(&seat.container)
.await
.or(seat.bind_port)
.ok_or_else(|| anyhow!("my computer started but the viewer port is missing"))
}
/// X server, window manager, VNC and noVNC all answering. The WM check
/// matters: xdotool against a bare Xvfb has no focus or window stacking.
const DESKTOP_PROBE: &str = "xdpyinfo -display :1 >/dev/null 2>&1 \
&& xprop -display :1 -root _NET_SUPPORTING_WM_CHECK 2>/dev/null | grep -q 'window id' \
&& python3 -c 'import urllib.request, socket; urllib.request.urlopen(\"http://127.0.0.1:6080/vnc.html\", timeout=1).close(); socket.create_connection((\"127.0.0.1\", 5900), timeout=1).close()'";
async fn desktop_up(container: &str) -> bool {
docker_exec(container, &["bash", "-c", DESKTOP_PROBE])
.await
.map(|o| o.status == 0)
.unwrap_or(false)
}
/// Stop and remove the container built from an older image. Another process
/// (outside our provisioning lock, e.g. an older binary) may be doing the
/// same; "already in progress" or "no such container" both mean it is gone.
async fn remove_stale_container(container: &str) -> Result<()> {
if docker_running(container).await? {
if let Err(error) = docker(["stop", "-t", "20", container]).await {
if !is_concurrent_removal(&error) {
return Err(error);
}
}
}
if let Err(error) = docker(["rm", container]).await {
if !is_concurrent_removal(&error) {
return Err(error);
}
}
for _ in 0..100 {
// Gone, or already replaced by a container on the current image.
if docker(["container", "inspect", container]).await.is_err()
|| container_image_id(container).await == local_image_id().await
{
return Ok(());
}
tokio::time::sleep(Duration::from_millis(200)).await;
}
Err(anyhow!(
"stale container {container} is still being removed; retry in a moment"
))
}
fn is_concurrent_removal(error: &anyhow::Error) -> bool {
let text = error.to_string();
text.contains("already in progress")
|| text.contains("No such container")
|| text.contains("is already stopped")
}
async fn wait_desktop(container: &str) -> Result<()> {
for _ in 0..120 {
if desktop_up(container).await {
return Ok(());
}
tokio::time::sleep(Duration::from_millis(250)).await;
}
Err(anyhow!("My computer did not become ready: desktop or viewer startup timed out. Check the container desktop logs and try again."))
}
async fn docker_running(name: &str) -> Result<bool> {
let out = docker(["inspect", "-f", "{{.State.Running}}", name]).await;
Ok(matches!(out, Ok(s) if s.stdout.trim() == "true"))
}
struct CmdOut {
status: i32,
stdout: String,
stderr: String,
}
async fn docker<I, S>(args: I) -> Result<CmdOut>
where
I: IntoIterator<Item = S>,
S: AsRef<std::ffi::OsStr>,
{
let mut cmd = Command::new("docker");
cmd.args(args)
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.kill_on_drop(true);
let out = cmd.output().await.context("run docker")?;
let result = CmdOut {
status: out.status.code().unwrap_or(-1),
stdout: String::from_utf8_lossy(&out.stdout).into_owned(),
stderr: String::from_utf8_lossy(&out.stderr).into_owned(),
};
if result.status != 0 {
return Err(anyhow!(
"docker failed ({}): {}",
result.status,
result.stderr.trim()
));
}
Ok(result)
}
fn max_running() -> usize {
std::env::var("LAZYBOY_COMPUTER_MAX")
.ok()
.and_then(|s| s.parse().ok())
.filter(|&n: &usize| n > 0)
.unwrap_or(DEFAULT_MAX_RUNNING)
}
pub struct BoxPool {
provisioning: Mutex<()>,
hubs: std::sync::Mutex<HashMap<String, Arc<BoxHub>>>,
}
impl BoxPool {
pub fn new() -> Arc<Self> {
Arc::new(Self {
provisioning: Mutex::new(()),
hubs: std::sync::Mutex::new(HashMap::new()),
})
}
pub fn global() -> Arc<Self> {
static POOL: OnceLock<Arc<BoxPool>> = OnceLock::new();
POOL.get_or_init(Self::new).clone()
}
pub fn hub(&self, seat_id: &str) -> Arc<BoxHub> {
let seat = BoxSeat::for_id(seat_id);
let key = seat.id.clone();
self.hubs
.lock()
.unwrap()
.entry(key)
.or_insert_with(|| BoxHub::for_seat(seat))
.clone()
}
pub async fn inspect(&self, seat_id: &str) -> Value {
self.hub(seat_id).inspect().await
}
/// Recover computers left by a previous server, so limits survive make start.
pub async fn reconcile(&self) -> Result<()> {
let names = docker(["ps", "-a", "--format", "{{.Names}}"]).await?;
for name in names.stdout.lines().rev() {
let id = if name == SESSION_CONTAINER { SESSION_SEAT }
else if let Some(id) = name.strip_prefix("lazyboy-box-") { id }
else { continue };
self.hub(id);
limit_container(name).await?;
}
Ok(())
}
pub async fn recover(&self) -> Result<()> {
self.reconcile().await?;
self.reclaim("").await?;
Ok(())
}
pub async fn start(&self, seat_id: &str) -> Result<Value> {
// Serialize admission and launch; concurrent tabs cannot exceed the cap.
let _guard = self.provisioning.lock().await;
self.reconcile().await?;
let crowded = self.reclaim(seat_id).await?;
let mut payload = self.hub(seat_id).ensure_ready().await?;
if crowded {
payload["crowded"] = json!(true);
}
Ok(payload)
}
pub async fn restart(&self, seat_id: &str) -> Result<Value> {
self.hub(seat_id).restart().await
}
pub async fn update(&self, seat_id: &str) -> Result<Value> {
self.hub(seat_id).update().await
}
pub async fn drop_seat(&self, seat_id: &str) -> Result<()> {
if BoxSeat::for_id(seat_id).id == SESSION_SEAT {
return Ok(());
}
let hub = self.hub(seat_id);
hub.wipe().await?;
self.hubs
.lock()
.unwrap()
.remove(&BoxSeat::for_id(seat_id).id);
Ok(())
}
async fn reclaim(&self, keep: &str) -> Result<bool> {
let max = max_running();
let reserve_slot = !keep.is_empty();
let keep = BoxSeat::for_id(keep).id;
loop {
let hubs: Vec<Arc<BoxHub>> = self.hubs.lock().unwrap().values().cloned().collect();
let mut running = Vec::new();
for hub in hubs {
if docker_running(hub.container()).await.unwrap_or(false) {
running.push((hub.last_used().await, hub));
}
}
let already_running = running.iter().any(|(_, hub)| hub.seat_id() == keep);
if running.len() < max || ((!reserve_slot || already_running) && running.len() <= max) {
return Ok(false);
}
running.sort_by_key(|(used, _)| *used);
let victim = running.into_iter().find(|(_, hub)| {
hub.seat_id() != keep && hub.seat_id() != SESSION_SEAT && !hub.has_computer_use()
});
match victim {
Some((_, hub)) => {
if hub.stop_keep_volumes().await.is_err() {
return Err(anyhow!("Could not suspend an idle computer to free memory"));
}
}
None => return Err(anyhow!("All computer slots are busy; stop an idle computer before starting another")),
}
}
}
}
async fn docker_exec(container: &str, args: &[&str]) -> Result<CmdOut> {
docker_exec_env(container, &[], args).await
}
async fn docker_exec_env(container: &str, env: &[(&str, &str)], args: &[&str]) -> Result<CmdOut> {
let mut all = vec!["exec".to_string()];
for (key, value) in env {
all.push("-e".into());
all.push(format!("{key}={value}"));
}
all.push(container.to_string());
all.extend(args.iter().map(|s| (*s).to_string()));
let mut cmd = Command::new("docker");
cmd.args(&all)
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.kill_on_drop(true);
let out = cmd.output().await.context("docker exec")?;
Ok(CmdOut {
status: out.status.code().unwrap_or(-1),
stdout: String::from_utf8_lossy(&out.stdout).into_owned(),
stderr: String::from_utf8_lossy(&out.stderr).into_owned(),
})
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn box_paths_stay_on_my_computer() {
assert_eq!(
resolve_box_path("notes.txt").unwrap(),
"/workspace/notes.txt"
);
assert_eq!(resolve_box_path("/workspace/a").unwrap(), "/workspace/a");
assert!(resolve_box_path("/etc/passwd").is_err());
assert!(resolve_box_path("../secret").is_err());
assert!(resolve_box_path("/workspace/../home").is_err());
assert!(resolve_box_path("/workspace/../../etc/passwd").is_err());
assert!(resolve_box_path("/workspaceX/a").is_err());
assert!(resolve_box_path("/tmp/other").is_err());
}
#[test]
fn box_paths_normalise_without_false_positives() {
assert_eq!(
resolve_box_path("/workspace///a/./b").unwrap(),
"/workspace/a/b"
);
assert_eq!(
resolve_box_path("report..final.txt").unwrap(),
"/workspace/report..final.txt"
);
assert_eq!(resolve_box_path("/home/box/").unwrap(), "/home/box");
assert_eq!(
resolve_box_path("/tmp/gb-jobs/abc/stdout").unwrap(),
"/tmp/gb-jobs/abc/stdout"
);
}
#[test]
fn dir_size_sums_nested_files() {
let dir = std::env::temp_dir().join(format!("gb-dirsize-{}", uuid::Uuid::new_v4()));
std::fs::create_dir_all(dir.join("inner")).unwrap();
std::fs::write(dir.join("a.bin"), [0u8; 10]).unwrap();
std::fs::write(dir.join("inner/b.bin"), [0u8; 5]).unwrap();
assert_eq!(dir_size(&dir), 15);
std::fs::remove_dir_all(&dir).unwrap();
}
#[tokio::test]
async fn missing_docker_is_explicit() {
if docker(["info"]).await.is_ok() {
return;
}
let hub = BoxHub::new();
let err = hub.ensure_ready().await.unwrap_err().to_string();
assert!(err.contains("Docker"), "{err}");
assert!(
err.contains("your computer") || err.contains("My computer"),
"{err}"
);
}
#[tokio::test]
async fn restart_and_update_need_docker() {
if docker(["info"]).await.is_ok() {
return;
}
let hub = BoxHub::new();
for err in [
hub.restart().await.unwrap_err().to_string(),
hub.update().await.unwrap_err().to_string(),
] {
assert!(err.contains("Docker"), "{err}");
}
}
#[test]
fn ready_payload_names_the_box() {
let hub = BoxHub::new();
let payload = hub.ready_payload();
assert_eq!(payload["ready"], true);
assert_eq!(payload["workspace"], "/workspace");
assert_eq!(payload["revision"], BOX_REVISION);
assert_eq!(payload["seat"], SESSION_SEAT);
assert!(payload["viewer_url"].as_str().unwrap().contains("/novnc/"));
}
#[test]
fn named_agents_get_isolated_seats() {
let a = BoxSeat::for_id("11111111-1111-1111-1111-111111111111");
let b = BoxSeat::for_id("22222222-2222-2222-2222-222222222222");
assert_ne!(a.container, b.container);
assert_ne!(a.home_volume, b.home_volume);
assert_ne!(a.workspace_volume, b.workspace_volume);
assert!(a.container.starts_with("lazyboy-box-"));
assert_eq!(BoxSeat::for_id("session").container, SESSION_CONTAINER);
assert_eq!(BoxSeat::for_id("").container, SESSION_CONTAINER);
assert_eq!(sanitize_seat_id("../etc"), "etc");
}
#[test]
fn pool_reuses_the_same_hub() {
let pool = BoxPool::new();
let one = pool.hub("agent-a");
let two = pool.hub("agent-a");
assert!(Arc::ptr_eq(&one, &two));
assert!(!Arc::ptr_eq(&one, &pool.hub("agent-b")));
}
#[tokio::test]
async fn inspect_does_not_start_a_container() {
let hub = BoxHub::for_seat(BoxSeat::for_id("missing-agent"));
let status = hub.inspect().await;
assert_eq!(status["ready"], false);
assert!(
status["state"] == "stopped" || status["state"] == "error",
"{status}"
);
assert_eq!(status["seat"], "missing-agent");
}
#[test]
fn computer_use_lock_is_per_hub() {
let a = BoxHub::for_seat(BoxSeat::for_id("a"));
let b = BoxHub::for_seat(BoxSeat::for_id("b"));
a.try_begin_computer_use().unwrap();
assert!(a.try_begin_computer_use().is_err());
b.try_begin_computer_use().unwrap();
a.end_computer_use();
a.try_begin_computer_use().unwrap();
}
}