diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index ae31c6835..009f2622e 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -102,7 +102,7 @@ jobs: - name: Native lint including Windows backend run: cargo clippy -p buzz-foundation -p buzz-agent-controller -p buzz-credential-store --locked --all-targets -- -D warnings - name: All native package tests - run: cargo test -p buzz-foundation -p buzz-agent-controller -p buzz-credential-store --locked + run: cargo test -p buzz-foundation -p buzz-agent-controller -p buzz-credential-store --locked --no-fail-fast measurements: if: github.event_name != 'workflow_dispatch' diff --git a/Cargo.lock b/Cargo.lock index 50fd7dcb1..2b8b1be3a 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -586,6 +586,7 @@ dependencies = [ "strip-ansi-escapes", "tempfile", "url", + "windows-sys 0.61.2", "zeroize", ] diff --git a/crates/agent-controller/Cargo.toml b/crates/agent-controller/Cargo.toml index 0d2af2507..7899d1f55 100644 --- a/crates/agent-controller/Cargo.toml +++ b/crates/agent-controller/Cargo.toml @@ -30,5 +30,8 @@ libc = "0.2" security-framework = "3.7" nix = { version = "0.31", features = ["user"] } +[target.'cfg(windows)'.dependencies] +windows-sys = { version = "0.61", features = ["Win32_Foundation", "Win32_Security", "Win32_Storage_FileSystem", "Win32_System_Diagnostics_ToolHelp", "Win32_System_IO", "Win32_System_JobObjects", "Win32_System_Pipes", "Win32_System_SystemServices", "Win32_System_Threading"] } + [target.'cfg(any(target_os = "windows", target_os = "linux"))'.dependencies] buzz-credential-store = { path = "../credential-store" } diff --git a/crates/agent-controller/src/connection.rs b/crates/agent-controller/src/connection.rs index 641cf1a94..3f98b9592 100644 --- a/crates/agent-controller/src/connection.rs +++ b/crates/agent-controller/src/connection.rs @@ -23,6 +23,10 @@ impl DatabricksSettings { Ok(()) } } +/// The pinned engine keeps non-Unix OAuth tokens in memory only, so neither a +/// later request nor a worker could reuse a sign-in. Deferred on Windows. +pub const DATABRICKS_WINDOWS: &str = + "Databricks sign-in is not supported on Windows yet. Choose OpenAI for this agent"; pub fn origin(raw: &str) -> Result { if raw.is_empty() { return Err("Databricks workspace is not configured. Edit the agent, open Advanced → Model, and set Databricks workspace (HTTPS origin).".into()); diff --git a/crates/agent-controller/src/lib.rs b/crates/agent-controller/src/lib.rs index 65bf1a027..96e32f24f 100644 --- a/crates/agent-controller/src/lib.rs +++ b/crates/agent-controller/src/lib.rs @@ -19,9 +19,7 @@ mod restart; mod runtime; mod secret; mod store; -#[cfg(unix)] mod supervisor; -#[cfg(unix)] pub use supervisor::dispatch as dispatch_agent_supervisor; pub use agent_defaults::{AgentDefaultsEdit, AgentDefaultsView}; diff --git a/crates/agent-controller/src/logs.rs b/crates/agent-controller/src/logs.rs index 52dc4e54f..fc4d755cb 100644 --- a/crates/agent-controller/src/logs.rs +++ b/crates/agent-controller/src/logs.rs @@ -2,9 +2,7 @@ use crate::Result; use sha2::{Digest, Sha256}; use std::fs::{File, OpenOptions}; -#[cfg(any(unix, test))] -use std::io::Write; -use std::io::{Read, Seek, SeekFrom}; +use std::io::{Read, Seek, SeekFrom, Write}; use std::path::{Path, PathBuf}; /// Domain-separated proof input; only the native-generated single-use nonce is @@ -94,12 +92,10 @@ fn open(path: &Path, write: bool) -> std::io::Result { Ok(file) } -// Only the Unix supervisor writes logs; keep the cross-platform retention tests. -#[cfg(any(unix, test))] +// Only the supervisor writes logs. pub(crate) struct Writer { file: File, } -#[cfg(any(unix, test))] impl Writer { pub(crate) fn new(path: &Path) -> Result { let file = open(path, true).map_err(|_| "Could not open harness log")?; diff --git a/crates/agent-controller/src/process.rs b/crates/agent-controller/src/process.rs index 3bef8cf57..44110ff5d 100644 --- a/crates/agent-controller/src/process.rs +++ b/crates/agent-controller/src/process.rs @@ -1,5 +1,6 @@ //! Every listener gets its own Unix session. ACP's worker process groups stay //! inside that session, so teardown is not limited to the listener's group. +//! On Windows, a kill-on-close Job Object contains the listener before it runs. use crate::Result; #[cfg(unix)] use std::process::Stdio; @@ -11,6 +12,8 @@ pub(crate) struct Process { child: Child, #[cfg(unix)] session: u32, + #[cfg(windows)] + job: job::Job, stopped: bool, } impl Process { @@ -28,10 +31,15 @@ impl Process { }); } } - #[cfg(not(unix))] + #[cfg(windows)] { - let _ = command; - Err("Agent process containment is not supported on this platform yet".into()) + let job = job::Job::create()?; + let child = job.spawn(command)?; + Ok(Self { + child, + job, + stopped: false, + }) } #[cfg(unix)] { @@ -50,6 +58,8 @@ impl Process { if self.stopped { return Ok(false); } + #[cfg(windows)] + self.job.sweep(); match self .child .try_wait() @@ -107,8 +117,16 @@ impl Process { self.stopped = true; Ok(()) } - #[cfg(not(unix))] - Err("Agent process containment is not supported on this platform yet".into()) + #[cfg(windows)] + { + // A windowless listener has no cooperative stop signal. + self.job.stop()?; + self.child + .wait() + .map_err(|_| "Could not reap agent listener")?; + self.stopped = true; + Ok(()) + } } } impl Drop for Process { @@ -164,3 +182,397 @@ fn session_members(session: u32) -> Result> { } Ok(members) } + +#[cfg(windows)] +mod job { + use crate::Result; + use std::collections::HashSet; + use std::os::windows::io::{AsRawHandle, FromRawHandle, OwnedHandle}; + use std::os::windows::process::CommandExt; + use std::process::{Child, Command}; + use std::sync::{mpsc, Arc}; + use std::time::{Duration, Instant}; + use windows_sys::Win32::Foundation::{ + GetLastError, ERROR_NO_MORE_FILES, FILETIME, INVALID_HANDLE_VALUE, WAIT_OBJECT_0, + }; + use windows_sys::Win32::System::Diagnostics::ToolHelp::{ + CreateToolhelp32Snapshot, Thread32First, Thread32Next, TH32CS_SNAPTHREAD, THREADENTRY32, + }; + use windows_sys::Win32::System::JobObjects::{ + AssignProcessToJobObject, CreateJobObjectW, IsProcessInJob, + JobObjectAssociateCompletionPortInformation, JobObjectBasicAccountingInformation, + JobObjectBasicProcessIdList, JobObjectExtendedLimitInformation, QueryInformationJobObject, + SetInformationJobObject, TerminateJobObject, JOBOBJECT_ASSOCIATE_COMPLETION_PORT, + JOBOBJECT_BASIC_ACCOUNTING_INFORMATION, JOBOBJECT_EXTENDED_LIMIT_INFORMATION, + JOB_OBJECT_LIMIT_KILL_ON_JOB_CLOSE, + }; + use windows_sys::Win32::System::SystemServices::JOB_OBJECT_MSG_NEW_PROCESS; + use windows_sys::Win32::System::Threading::{ + GetProcessTimes, OpenProcess, OpenThread, ResumeThread, WaitForSingleObject, + CREATE_NO_WINDOW, CREATE_SUSPENDED, INFINITE, PROCESS_QUERY_LIMITED_INFORMATION, + PROCESS_SYNCHRONIZE, THREAD_SUSPEND_RESUME, + }; + use windows_sys::Win32::System::IO::{ + CreateIoCompletionPort, GetQueuedCompletionStatus, PostQueuedCompletionStatus, + }; + + /// Completion key of job messages; the watcher stops on any other. + const JOB: usize = 1; + /// A member's process object, keyed by ID and creation time. + type Member = ((u32, u64), OwnedHandle); + + /// Only a member's own signaled process object proves that it has exited. + /// A watcher opens each process as it joins; `seen` counts the processes + /// opened and `held` keeps those not yet signaled, across Stop attempts. + pub(super) struct Job { + handle: Arc, + port: Arc, + joined: mpsc::Receiver, + held: Vec, + seen: HashSet<(u32, u64)>, + } + + impl Job { + /// Unnamed and without breakaway: members cannot leave, and closing the + /// last handle (including on owner death) terminates every member. + pub(super) fn create() -> Result { + let handle = unsafe { CreateJobObjectW(std::ptr::null(), std::ptr::null()) }; + if handle.is_null() { + return Err("Could not create agent process container".into()); + } + let handle = Arc::new(unsafe { OwnedHandle::from_raw_handle(handle) }); + let mut limits = JOBOBJECT_EXTENDED_LIMIT_INFORMATION::default(); + limits.BasicLimitInformation.LimitFlags = JOB_OBJECT_LIMIT_KILL_ON_JOB_CLOSE; + if unsafe { + SetInformationJobObject( + handle.as_raw_handle(), + JobObjectExtendedLimitInformation, + (&limits as *const JOBOBJECT_EXTENDED_LIMIT_INFORMATION).cast(), + std::mem::size_of_val(&limits) as u32, + ) + } == 0 + { + return Err("Could not configure agent process container".into()); + } + // Associated before the listener joins, so no member predates the port. + let port = + unsafe { CreateIoCompletionPort(INVALID_HANDLE_VALUE, std::ptr::null_mut(), 0, 1) }; + if port.is_null() { + return Err("Could not configure agent process container".into()); + } + let port = Arc::new(unsafe { OwnedHandle::from_raw_handle(port) }); + let association = JOBOBJECT_ASSOCIATE_COMPLETION_PORT { + CompletionKey: JOB as *mut _, + CompletionPort: port.as_raw_handle(), + }; + if unsafe { + SetInformationJobObject( + handle.as_raw_handle(), + JobObjectAssociateCompletionPortInformation, + (&association as *const JOBOBJECT_ASSOCIATE_COMPLETION_PORT).cast(), + std::mem::size_of_val(&association) as u32, + ) + } == 0 + { + return Err("Could not configure agent process container".into()); + } + let (found, joined) = mpsc::channel(); + let (job, queue) = (handle.clone(), port.clone()); + std::thread::Builder::new() + .spawn(move || watch(&job, &queue, &found)) + .map_err(|_| "Could not configure agent process container")?; + Ok(Self { + handle, + port, + joined, + held: Vec::new(), + seen: HashSet::new(), + }) + } + + pub(super) fn spawn(&self, command: &mut Command) -> Result { + spawn(&self.handle, command) + } + + /// Collect opened members and retire only signaled handles. + pub(super) fn sweep(&mut self) { + while let Ok(member) = self.joined.try_recv() { + self.add(member); + } + self.held.retain( + |held| unsafe { WaitForSingleObject(held.as_raw_handle(), 0) } != WAIT_OBJECT_0, + ); + } + + fn add(&mut self, (key, handle): Member) { + if self.seen.insert(key) { + self.held.push(handle); + } + } + + /// Terminating only requests exit. Succeed once every process that ever + /// joined was opened and each opened process is signaled; a lost + /// message or a member gone before it was opened fails closed. + pub(super) fn stop(&mut self) -> Result<()> { + // Open live members directly in case their messages are still queued. + for id in listed(&self.handle) { + if let Some(member) = member(&self.handle, id) { + self.add(member); + } + } + if unsafe { TerminateJobObject(self.handle.as_raw_handle(), 1) } == 0 { + return Err("Could not stop agent processes".into()); + } + let deadline = Instant::now() + Duration::from_secs(5); + loop { + self.sweep(); + if self.held.is_empty() && self.seen.len() == joined(&self.handle)? as usize { + return Ok(()); + } + if Instant::now() >= deadline { + return Err("Agent descendants have not exited; shutdown is incomplete".into()); + } + std::thread::sleep(Duration::from_millis(25)); + } + } + } + + impl Drop for Job { + /// Wake the watcher so it releases its job handle. + fn drop(&mut self) { + unsafe { + PostQueuedCompletionStatus(self.port.as_raw_handle(), 0, 0, std::ptr::null()) + }; + } + } + + /// Try to open each process as it joins. Delivery is not guaranteed and a + /// fast member can exit first; a missed process stays uncounted, never inferred. + fn watch(job: &OwnedHandle, port: &OwnedHandle, found: &mpsc::Sender) { + loop { + let (mut message, mut key, mut id) = (0, 0, std::ptr::null_mut()); + if unsafe { + GetQueuedCompletionStatus( + port.as_raw_handle(), + &mut message, + &mut key, + &mut id, + INFINITE, + ) + } == 0 + || key != JOB + { + return; + } + if message == JOB_OBJECT_MSG_NEW_PROCESS { + // For job messages the overlapped pointer carries the process ID. + if let Some(member) = member(job, id as usize as u32) { + if found.send(member).is_err() { + return; + } + } + } + } + } + + /// A reused ID names a process outside the job. Creation time tells a + /// member reopened for a later message from one that is not yet counted. + fn member(job: &OwnedHandle, id: u32) -> Option { + let access = PROCESS_SYNCHRONIZE | PROCESS_QUERY_LIMITED_INFORMATION; + let handle = unsafe { OpenProcess(access, 0, id) }; + if handle.is_null() { + return None; + } + let handle = unsafe { OwnedHandle::from_raw_handle(handle) }; + let mut inside = 0; + let mut times = [FILETIME::default(); 4]; + let [created, exited, kernel, user] = &mut times; + if unsafe { IsProcessInJob(handle.as_raw_handle(), job.as_raw_handle(), &mut inside) } == 0 + || inside == 0 + || unsafe { GetProcessTimes(handle.as_raw_handle(), created, exited, kernel, user) } + == 0 + { + return None; + } + let created = (created.dwHighDateTime as u64) << 32 | created.dwLowDateTime as u64; + Some(((id, created), handle)) + } + + /// Fail closed: a child that is not contained never runs its first instruction. + fn spawn(job: &OwnedHandle, command: &mut Command) -> Result { + command.creation_flags(CREATE_SUSPENDED | CREATE_NO_WINDOW); + let mut child = command.spawn().map_err(|_| { + "Could not start bundled agent listener; check the runtime installation" + })?; + if let Err(error) = contain(job, &child) { + let _ = child.kill(); + let _ = child.wait(); + return Err(error); + } + Ok(child) + } + + /// Assign a suspended child, then resume its only thread. + fn contain(job: &OwnedHandle, child: &Child) -> Result<()> { + let thread = primary_thread(child.id()).ok_or("Could not contain agent listener")?; + if unsafe { AssignProcessToJobObject(job.as_raw_handle(), child.as_raw_handle()) } == 0 { + return Err("Could not contain agent listener".into()); + } + if unsafe { ResumeThread(thread.as_raw_handle()) } != 1 { + return Err("Could not resume contained agent listener".into()); + } + Ok(()) + } + + /// Every process ever associated with the job, including exited ones. + fn joined(job: &OwnedHandle) -> Result { + let mut info = JOBOBJECT_BASIC_ACCOUNTING_INFORMATION::default(); + if unsafe { + QueryInformationJobObject( + job.as_raw_handle(), + JobObjectBasicAccountingInformation, + (&mut info as *mut JOBOBJECT_BASIC_ACCOUNTING_INFORMATION).cast(), + std::mem::size_of_val(&info) as u32, + std::ptr::null_mut(), + ) + } == 0 + { + return Err("Could not inspect agent descendants".into()); + } + Ok(info.TotalProcesses) + } + + /// IDs of the active members. Any it misses is left to the joined count. + fn listed(job: &OwnedHandle) -> Vec { + // JOBOBJECT_BASIC_PROCESS_ID_LIST with room for any listener tree. + #[repr(C)] + struct List { + assigned: u32, + listed: u32, + ids: [usize; 1024], + } + let mut list = List { + assigned: 0, + listed: 0, + ids: [0; 1024], + }; + let queried = unsafe { + QueryInformationJobObject( + job.as_raw_handle(), + JobObjectBasicProcessIdList, + (&mut list as *mut List).cast(), + std::mem::size_of_val(&list) as u32, + std::ptr::null_mut(), + ) + }; + if queried == 0 { + return Vec::new(); + } + list.ids + .iter() + .take(list.listed as usize) + .map(|&id| id as u32) + .collect() + } + + // Same exactly-one-thread check as the host command container. + fn primary_thread(process_id: u32) -> Option { + let snapshot = unsafe { CreateToolhelp32Snapshot(TH32CS_SNAPTHREAD, 0) }; + if snapshot == INVALID_HANDLE_VALUE { + return None; + } + let snapshot = unsafe { OwnedHandle::from_raw_handle(snapshot) }; + let mut entry = THREADENTRY32 { + dwSize: std::mem::size_of::() as u32, + ..Default::default() + }; + if unsafe { Thread32First(snapshot.as_raw_handle(), &mut entry) } == 0 { + return None; + } + let mut thread_id = None; + loop { + if entry.th32OwnerProcessID == process_id + && thread_id.replace(entry.th32ThreadID).is_some() + { + return None; + } + entry.dwSize = std::mem::size_of::() as u32; + if unsafe { Thread32Next(snapshot.as_raw_handle(), &mut entry) } == 0 { + break; + } + } + if unsafe { GetLastError() } != ERROR_NO_MORE_FILES { + return None; + } + let handle = unsafe { OpenThread(THREAD_SUSPEND_RESUME, 0, thread_id?) }; + if handle.is_null() { + return None; + } + Some(unsafe { OwnedHandle::from_raw_handle(handle) }) + } + + #[cfg(test)] + #[test] + fn uncontained_child_never_runs() { + let dir = tempfile::tempdir().unwrap(); + let marker = dir.path().join("ran"); + let mut command = Command::new(std::env::var_os("ComSpec").unwrap()); + command.raw_arg(format!("/d /c type nul > \"{}\"", marker.display())); + // A handle that is not a job: assignment fails after the suspended spawn. + let not_a_job = std::fs::File::create(dir.path().join("not-a-job")) + .unwrap() + .into(); + assert_eq!( + spawn(¬_a_job, &mut command).unwrap_err(), + "Could not contain agent listener" + ); + assert!(!marker.exists(), "uncontained child ran"); + let job = Job::create().unwrap(); + let mut child = job.spawn(&mut command).unwrap(); + assert!(child.wait().unwrap().success()); + assert!(marker.exists(), "contained child was not resumed"); + } + + #[cfg(test)] + #[test] + fn member_gone_before_it_was_opened_fails_stop() { + let shell = std::env::var_os("ComSpec").unwrap(); + let mut job = Job::create().unwrap(); + let mut root = job + .spawn(Command::new(&shell).raw_arg("/d /c ping -n 600 127.0.0.1 >nul")) + .unwrap(); + let mut gone = job + .spawn(Command::new(&shell).raw_arg("/d /c exit")) + .unwrap(); + // Its own signaled handle proves the exit; holding it pins the ID. + assert!(gone.wait().unwrap().success()); + let deadline = Instant::now() + Duration::from_secs(30); + // Once received, a sweep after the exit also retires the watcher's handle. + let key = loop { + job.sweep(); + if let Some(&key) = job.seen.iter().find(|&&(id, _)| id == gone.id()) { + break key; + } + assert!(Instant::now() < deadline, "watcher did not open the member"); + std::thread::sleep(Duration::from_millis(10)); + }; + drop(gone); + // Wait until Stop could no longer reopen that exact process. + while member(&job.handle, key.0).is_some_and(|(found, _)| found == key) { + assert!(Instant::now() < deadline, "exited member was not released"); + std::thread::sleep(Duration::from_millis(10)); + } + // As if its notification were lost. + job.seen.remove(&key); + assert_eq!( + job.stop().unwrap_err(), + "Agent descendants have not exited; shutdown is incomplete" + ); + assert!(root.try_wait().unwrap().is_some(), "root survived Stop"); + assert!(job.held.is_empty(), "an opened member survived Stop"); + assert!( + !job.seen.contains(&key), + "Stop recaptured the missing member" + ); + } +} diff --git a/crates/agent-controller/src/runtime.rs b/crates/agent-controller/src/runtime.rs index 67886a284..4ec3c6c54 100644 --- a/crates/agent-controller/src/runtime.rs +++ b/crates/agent-controller/src/runtime.rs @@ -1,8 +1,5 @@ use crate::bundle::RuntimeBundle; use crate::config::Agent; -#[cfg(not(unix))] -use crate::process::Process; -#[cfg(unix)] use crate::supervisor::Supervised; use crate::{AgentEdit, ControlSnapshot, Credentials, ProcessStatus, Result, Store}; use serde::Deserialize; @@ -69,6 +66,36 @@ impl RuntimeBundle { .stdin(Stdio::null()) .stdout(Stdio::null()) .stderr(Stdio::null()); + // Windows system, profile and tool-discovery locations; none are secrets. + #[cfg(windows)] + let platform = [ + "SystemRoot", + "windir", + "SystemDrive", + "ComSpec", + "PATHEXT", + "USERPROFILE", + "HOMEDRIVE", + "HOMEPATH", + "USERNAME", + "APPDATA", + "LOCALAPPDATA", + "ProgramData", + "ProgramFiles", + "ProgramFiles(x86)", + "ProgramW6432", + "CommonProgramFiles", + "CommonProgramFiles(x86)", + "CommonProgramW6432", + "PROCESSOR_ARCHITECTURE", + "NUMBER_OF_PROCESSORS", + "OS", + // Runtime shell overrides; Git Bash is otherwise discovered from PATH. + "BUZZ_SHELL", + "GIT_BASH", + ]; + #[cfg(not(windows))] + let platform: [&str; 0] = []; for name in [ "HOME", "TMPDIR", @@ -78,7 +105,10 @@ impl RuntimeBundle { "SSH_AUTH_SOCK", "SSL_CERT_FILE", "SSL_CERT_DIR", - ] { + ] + .into_iter() + .chain(platform) + { if let Some(value) = std::env::var_os(name) { command.env(name, value); } @@ -98,7 +128,7 @@ impl RuntimeBundle { ( agent.harness.args.clone(), &agent.environment, - "/usr/bin:/bin:/usr/sbin:/sbin".into(), + tools_path()?, ) }; let path = std::env::join_paths( @@ -220,6 +250,10 @@ fn databricks_with_defaults( ) { return Ok(None); } + // Before any OAuth-setting check: Windows never asks for a workspace. + if cfg!(windows) { + return Err(crate::connection::DATABRICKS_WINDOWS.into()); + } if agent.environment.contains_key("DATABRICKS_TOKEN") { return Err("Remove DATABRICKS_TOKEN to use this app's persistent OAuth connection".into()); } @@ -234,6 +268,25 @@ fn databricks_with_defaults( settings.validate()?; Ok(Some(settings)) } +/// PATH after the runtime bundle for non-Pi harnesses. Windows keeps its native +/// PATH, where Git Bash and user tools are installed; Unix uses a fixed floor +/// plus, on Linux, common user-level install locations. +fn tools_path() -> Result { + if cfg!(windows) { + return Ok(std::env::var_os("PATH").unwrap_or_default()); + } + let mut dirs = Vec::new(); + if cfg!(target_os = "linux") { + let home = std::env::var_os("HOME").map(PathBuf::from); + dirs.extend( + home.filter(|h| h.is_absolute()) + .map(|h| h.join(".local/bin")), + ); + dirs.push(PathBuf::from("/usr/local/bin")); + } + dirs.extend(["/usr/bin", "/bin", "/usr/sbin", "/sbin"].map(PathBuf::from)); + std::env::join_paths(dirs).map_err(|_| "Invalid runtime tools path".into()) +} /// App-owned npm shims and the pinned Node binary are separate from user-global tools. /// `app_data` is Tauri's resolved app-data directory, never browser input. pub fn managed_tool(app_data: &Path, name: &str) -> Option { @@ -296,20 +349,13 @@ pub enum Action { Restart, } struct Running { - #[cfg(unix)] process: Supervised, - #[cfg(not(unix))] - process: Process, revision: u64, /// Native-only: holds environment values and is never serialized. spawned: serde_json::Value, databricks_host: Option, #[cfg(all(test, unix))] temporary: Option, - #[cfg(not(unix))] - _temporary: Option, - #[cfg(not(unix))] - _ownership: crate::ownership::Ownership, } impl Drop for Running { fn drop(&mut self) { @@ -394,7 +440,6 @@ impl Controller { ); } Err(error) => { - #[cfg(unix)] if run.process.stopped() { self.running.remove(&agent.id); } @@ -774,8 +819,6 @@ impl Controller { // Blank fields inherit agent defaults at each start; never saved back. let agent = crate::agent_defaults::effective(&agent, &self.store.defaults()?); let bundle = self.bundle.as_ref().map_err(Clone::clone)?; - #[cfg(not(unix))] - let ownership = crate::ownership::Ownership::acquire(&self.ownership_root, &agent.id)?; let stored; let key = match supplied { Some(key) => key, @@ -816,11 +859,8 @@ impl Controller { } // Disarm app-side deletion before a child can use this directory. The // supervisor deletes it only after confirmed whole-session teardown. - #[cfg(unix)] let log_path = crate::logs::path(config, &agent.id)?; - #[cfg(unix)] let temporary = temporary.keep(); - #[cfg(unix)] let process = Supervised::spawn( &command, &self.ownership_root, @@ -828,8 +868,6 @@ impl Controller { &temporary, &log_path, )?; - #[cfg(not(unix))] - let process = Process::spawn(&mut command)?; self.running.insert( id.into(), Running { @@ -839,10 +877,6 @@ impl Controller { databricks_host: settings.map(|s| s.host), #[cfg(all(test, unix))] temporary: Some(temporary), - #[cfg(not(unix))] - _temporary: Some(temporary), - #[cfg(not(unix))] - _ownership: ownership, }, ); Ok(()) @@ -850,7 +884,6 @@ impl Controller { fn stop(&mut self, id: &str) -> Result<()> { if let Some(run) = self.running.get_mut(id) { let result = run.process.stop(); - #[cfg(unix)] if run.process.stopped() { // E confirms worker exit even when private-dir removal failed. self.running.remove(id); diff --git a/crates/agent-controller/src/runtime/tests.rs b/crates/agent-controller/src/runtime/tests.rs index 2b6282712..a68454668 100644 --- a/crates/agent-controller/src/runtime/tests.rs +++ b/crates/agent-controller/src/runtime/tests.rs @@ -4,9 +4,7 @@ use crate::config::{agent_id, HarnessEdit}; use crate::process::Process; use crate::Secret; use serde_json::json; -#[cfg(unix)] use std::fs; -#[cfg(unix)] use std::time::{Duration, Instant}; const KEY: &str = "0000000000000000000000000000000000000000000000000000000000000001"; const PUB: &str = "79be667ef9dcbbac55a06295ce870b07029bfcdb2dce28d959f2815b16f81798"; @@ -157,17 +155,24 @@ fn denied_credential_deletion_keeps_a_disabled_card_for_retry() { assert_eq!(remaining.len(), 1); assert!(!remaining[0].enabled); } -#[cfg(unix)] fn bundle(directory: &Path) -> RuntimeBundle { - use std::os::unix::fs::PermissionsExt; - for name in [ + let names = [ "buzz-acp", "buzz-agent", "buzz-dev-mcp", "buzz", "git-credential-nostr", - ] { + ] + .map(|name| format!("{name}{}", std::env::consts::EXE_SUFFIX)); + for name in &names { let path = directory.join(name); + // Windows cannot run this sh fixture; the listener is `windows_listener` + // and the other tools are only integrity-checked. + #[cfg(windows)] + if name == "buzz-acp.exe" { + fs::copy(std::env::current_exe().unwrap(), &path).unwrap(); + continue; + } fs::write(&path, r#"#!/bin/sh printf '%s\n' "$BUZZ_ACP_LAZY_POOL" "$BUZZ_ACP_IDLE_POOL_SLEEP" "$BUZZ_ACP_SYSTEM_PROMPT" "$BUZZ_ACP_MODEL" "$BUZZ_ACP_AGENT_ARGS" "$BUZZ_RELAY_URL" "$BUZZ_ACP_RESPOND_TO" "$BUZZ_MANAGED_AGENT" "$BUZZ_ACP_REPLAY_FLOOR" "$PROVIDER_TEST_SETTING" >> starts printf '%s' "$BUZZ_ACP_TEAM_INSTRUCTIONS" > team-instructions @@ -176,33 +181,30 @@ printf 'harness fixture output\n' trap 'exit 0' TERM INT while :; do [ -f "$BUZZ_AGENT_CONFIG_DIR/exit-listener" ] && exit 0; /bin/sleep 0.1; done "#).unwrap(); - fs::set_permissions(path, fs::Permissions::from_mode(0o700)).unwrap(); + #[cfg(unix)] + { + use std::os::unix::fs::PermissionsExt; + fs::set_permissions(path, fs::Permissions::from_mode(0o700)).unwrap(); + } } - let files: BTreeMap<_, _> = [ - "buzz-acp", - "buzz-agent", - "buzz-dev-mcp", - "buzz", - "git-credential-nostr", - ] - .into_iter() - .map(|name| { - use sha2::{Digest, Sha256}; - ( - name, - format!( - "{:x}", - Sha256::digest(fs::read(directory.join(name)).unwrap()) - ), - ) - }) - .collect(); + let files: BTreeMap<_, _> = names + .iter() + .map(|name| { + use sha2::{Digest, Sha256}; + ( + name, + format!( + "{:x}", + Sha256::digest(fs::read(directory.join(name)).unwrap()) + ), + ) + }) + .collect(); let source: serde_json::Value = serde_json::from_str(include_str!("../../../../runtime/agent-runtime.json")).unwrap(); fs::write(directory.join("manifest.json"), serde_json::to_vec(&json!({"version":1,"revision":source["revision"],"target":env!("BUZZ_RUNTIME_TARGET"),"files":files})).unwrap()).unwrap(); RuntimeBundle::new(directory.into()).unwrap() } -#[cfg(unix)] fn wait_for_contents(path: &Path, parse: impl Fn(&str) -> Option) -> T { let deadline = Instant::now() + Duration::from_secs(5); loop { @@ -645,7 +647,6 @@ fn stop_reports_cleanup_before_durable_disable_failure() { } #[test] -#[cfg(unix)] fn exact_command_has_no_ambient_identity_and_launch_failure_is_truthful() { let dir = tempfile::tempdir().unwrap(); let tools = tempfile::tempdir().unwrap(); @@ -699,6 +700,135 @@ fn exact_command_has_no_ambient_identity_and_launch_failure_is_truthful() { let stopped = controller.action(&a.id, Action::Stop).unwrap(); assert!(!stopped.agents[0].enabled); } +/// Windows listener stand-in: `bundle` installs this test binary as +/// buzz-acp.exe, and the guardian runs only this test from that copy. +#[test] +#[cfg(windows)] +fn windows_listener() { + let exe = std::env::current_exe().unwrap(); + if exe.file_stem() != Some("buzz-acp".as_ref()) { + return; + } + let env: BTreeMap<_, _> = std::env::vars() + .map(|(key, value)| (key.to_ascii_uppercase(), value)) + .collect(); + let mut line = serde_json::to_vec(&env).unwrap(); + line.push(b'\n'); + let mut starts = fs::OpenOptions::new() + .create(true) + .append(true) + .open("starts") + .unwrap(); + std::io::Write::write_all(&mut starts, &line).unwrap(); + // Bounded, so a containment failure cannot leave it running indefinitely. + std::thread::sleep(Duration::from_secs(300)); +} +#[test] +#[cfg(windows)] +fn windows_start_restart_shutdown_hold_custody_and_isolate_the_listener() { + let dir = tempfile::tempdir().unwrap(); + // Spaces reach the guardian's argv, the workspace and PATH. + let root = dir.path().join("Buzz profile"); + let tools = root.join("agent runtime"); + fs::create_dir_all(&tools).unwrap(); + let mut store = Store::open(root.join("config")).unwrap(); + let a = agent(&root); + store.insert(vec![a.clone()]).unwrap(); + let ownership = root.join("ownership"); + let mut controller = Controller::new( + store, + Arc::new(Memory), + Ok(bundle(&tools)), + ownership.clone(), + ); + let mut runs = Vec::new(); + for action in [Action::Start, Action::Restart] { + let snapshot = controller.action(&a.id, action).unwrap(); + assert!( + matches!(snapshot.agents[0].status, ProcessStatus::Running), + "{:?}", + snapshot.agents[0].error + ); + assert!(crate::ownership::Ownership::acquire(&ownership, &a.id).is_err()); + // Running is acknowledged at spawn, before the listener records its + // environment; wait so the next action cannot kill it before that row. + let count = runs.len() + 1; + runs = wait_for_contents(&root.join("starts"), |text| { + let runs = text + .lines() + .map(|line| serde_json::from_str(line).ok()) + .collect::>>>()?; + (runs.len() == count).then_some(runs) + }); + } + // Restart confirmed the first session's full teardown before starting again. + assert!(!Path::new(&runs[0]["TEMP"]).exists()); + controller.shutdown().unwrap(); + assert!(controller.running.is_empty()); + assert!(!Path::new(&runs[1]["TEMP"]).exists()); + let _released = crate::ownership::Ownership::acquire(&ownership, &a.id).unwrap(); + assert!(controller.store.agents().unwrap()[0].enabled); + // The listener got the launch environment, not the test process's. + let env = &runs[1]; + assert_eq!(env["BUZZ_PRIVATE_KEY"], KEY); + assert_eq!(env["SYSTEMROOT"], std::env::var("SystemRoot").unwrap()); + assert!(std::env::var_os("CARGO_MANIFEST_DIR").is_some()); + assert!(!env.contains_key("CARGO_MANIFEST_DIR")); + let native = std::env::var_os("PATH").unwrap(); + assert_eq!( + std::env::split_paths(&env["PATH"]).collect::>(), + std::iter::once(tools) + .chain(std::env::split_paths(&native)) + .collect::>() + ); +} +#[test] +#[cfg(windows)] +fn windows_refuses_databricks_before_workspace_validation_or_start() { + let refused = Some(crate::connection::DATABRICKS_WINDOWS); + let dir = tempfile::tempdir().unwrap(); + let mut a = agent(dir.path()); + a.harness.provider = "databricks_v2".into(); + let mut valid = a.clone(); + valid.harness.databricks = Some(crate::connection::DatabricksSettings { + host: "https://agent.example".into(), + filter: "agent-*".into(), + }); + // Empty, build-default and valid workspaces are all refused, never validated. + for (saved, defaults) in [ + (&a, crate::BuildDefaults::default()), + (&a, deployment_defaults()), + (&valid, crate::BuildDefaults::default()), + ] { + assert_eq!( + databricks_with_defaults(saved, &defaults).err().as_deref(), + refused + ); + } + let mut openai = a.clone(); + openai.harness.provider = "openai".into(); + assert!(databricks_with_defaults(&openai, &deployment_defaults()) + .unwrap() + .is_none()); + let tools = dir.path().join("tools"); + fs::create_dir_all(&tools).unwrap(); + let mut store = Store::open(dir.path().join("config")).unwrap(); + store.insert(vec![a.clone()]).unwrap(); + let mut controller = Controller::new( + store, + Arc::new(Memory), + Ok(bundle(&tools)), + dir.path().join("ownership"), + ); + assert_eq!( + controller.credential_request(&a.id).err().as_deref(), + refused + ); + let snapshot = controller.action(&a.id, Action::Start).unwrap(); + assert_eq!(snapshot.agents[0].error.as_deref(), refused); + assert!(controller.running.is_empty()); + assert!(!dir.path().join("starts").exists()); +} #[test] #[cfg(unix)] fn teardown_reaps_a_worker_in_a_separate_process_group() { @@ -936,6 +1066,31 @@ fn blank_selectors_without_overrides_leave_harness_defaults_intact() { } } +#[test] +fn launch_path_puts_bundled_tools_before_platform_tools() { + let dir = tempfile::tempdir().unwrap(); + let tools = tempfile::tempdir().unwrap(); + let command = bundle(tools.path()) + .command(&agent(dir.path()), &Secret::parse(KEY, PUB).unwrap()) + .unwrap(); + let env: BTreeMap<_, _> = command.get_envs().collect(); + let path = env[std::ffi::OsStr::new("PATH")].unwrap(); + let mut expected = vec![tools.path().to_owned()]; + if cfg!(windows) { + // Native PATH is where Git Bash and user tools live; system keys come along. + expected.extend(std::env::split_paths(&std::env::var_os("PATH").unwrap())); + assert!(env.contains_key(std::ffi::OsStr::new("SystemRoot"))); + } else { + if cfg!(target_os = "linux") { + let home = std::env::var_os("HOME").map(std::path::PathBuf::from); + expected.extend(home.map(|home| home.join(".local/bin"))); + expected.push("/usr/local/bin".into()); + } + expected.extend(["/usr/bin", "/bin", "/usr/sbin", "/sbin"].map(Into::into)); + } + assert_eq!(std::env::split_paths(path).collect::>(), expected); +} + #[test] #[cfg(unix)] fn start_and_restart_reject_missing_saved_identities() { @@ -1300,6 +1455,7 @@ fn build_floor_agrees_at_command_oauth_and_discovery_without_rewriting_saved_age } #[test] +#[cfg(unix)] // Windows: see windows_refuses_databricks_before_workspace_validation_or_start fn databricks_workspace_errors_distinguish_missing_configuration_from_invalid_origins() { let dir = tempfile::tempdir().unwrap(); let mut agent = agent(dir.path()); @@ -1340,6 +1496,7 @@ fn databricks_workspace_errors_distinguish_missing_configuration_from_invalid_or } #[test] +#[cfg(unix)] // Windows: see windows_refuses_databricks_before_workspace_validation_or_start fn saved_selectors_and_environment_override_build_floor_including_empty() { let dir = tempfile::tempdir().unwrap(); let mut agent = agent(dir.path()); @@ -1500,6 +1657,7 @@ fn goose_model_context_uses_effective_draft_provider_without_projecting_secrets( } #[test] +#[cfg(unix)] // Windows: see windows_refuses_databricks_before_workspace_validation_or_start fn discovery_accepts_only_v2_from_saved_environment_or_build_provider() { let dir = tempfile::tempdir().unwrap(); for provider in ["databricks_v2", "databricks-v2", "databricks"] { diff --git a/crates/agent-controller/src/supervisor.rs b/crates/agent-controller/src/supervisor.rs index e15b7ab75..ae9e70e07 100644 --- a/crates/agent-controller/src/supervisor.rs +++ b/crates/agent-controller/src/supervisor.rs @@ -4,11 +4,15 @@ use crate::ownership::Ownership; use crate::process::Process; use crate::Result; use std::io::{Read, Write}; -use std::os::fd::{FromRawFd, OwnedFd}; -use std::os::unix::net::UnixStream; use std::path::Path; use std::process::{Child, Command, Stdio}; use std::time::Duration; +#[cfg(windows)] +mod pipe; +#[cfg(windows)] +use pipe::{channel, inherited, log_channel, peek, Channel}; +#[cfg(unix)] +use unix::{channel, inherited, log_channel, peek, Channel}; pub const MODE: &str = "--buzz-agent-session-supervisor"; @@ -29,8 +33,11 @@ pub fn dispatch() -> bool { ) else { std::process::exit(1); }; - // The connected socket occupies stdin; stdout/stderr stay closed. - let socket = unsafe { UnixStream::from_raw_fd(0) }; + // The connected channel occupies stdin; stdout/stderr stay closed. + let Ok(socket) = inherited() else { + let _ = std::fs::remove_dir_all(&temp); + std::process::exit(1); + }; let result = serve( socket, Path::new(&root), @@ -43,7 +50,7 @@ pub fn dispatch() -> bool { } pub(crate) fn serve( - mut socket: UnixStream, + mut socket: Channel, root: &Path, id: &str, temp: &Path, @@ -73,19 +80,14 @@ pub(crate) fn serve( let mut spawned = false; let result = (|| { let mut writer = crate::logs::Writer::new(log)?; - let (mut reader, output) = - UnixStream::pair().map_err(|_| "Could not capture harness log")?; - reader - .set_nonblocking(true) - .map_err(|_| "Could not capture harness log")?; - let stderr = output - .try_clone() - .map_err(|_| "Could not capture harness log")?; + let (mut reader, stdout, stderr) = + log_channel().map_err(|_| "Could not capture harness log")?; let mut command = Command::new(runtime); - command - .stdin(Stdio::null()) - .stdout(Stdio::from(OwnedFd::from(output))) - .stderr(Stdio::from(OwnedFd::from(stderr))); + command.stdin(Stdio::null()).stdout(stdout).stderr(stderr); + // Windows cannot run the sh listener fixture; its stand-in is a copy of + // the test binary running only runtime::tests::windows_listener. + #[cfg(all(test, windows))] + command.args(["--exact", "runtime::tests::windows_listener", "--nocapture"]); let mut process = match Process::spawn(&mut command) { Ok(process) => process, Err(error) => { @@ -163,7 +165,7 @@ pub(crate) fn serve( } // A noisy child cannot starve Stop or app-death detection. -fn drain_log(reader: &mut UnixStream, writer: &mut crate::logs::Writer, live: bool) -> Result<()> { +fn drain_log(reader: &mut Channel, writer: &mut crate::logs::Writer, live: bool) -> Result<()> { let mut bytes = [0u8; 8192]; let mut chunks = 0; loop { @@ -192,7 +194,7 @@ fn reap_failed_start(mut child: Child) { }); } -fn confirm_start(mut socket: UnixStream, child: Child) -> Result<(UnixStream, Child)> { +fn confirm_start(mut socket: Channel, child: Child) -> Result<(Channel, Child)> { let mut state = [0u8]; let read = socket.read_exact(&mut state); if read.is_err() || state != *b"R" { @@ -217,7 +219,7 @@ fn confirm_start(mut socket: UnixStream, child: Child) -> Result<(UnixStream, Ch } pub(crate) struct Supervised { - socket: UnixStream, + socket: Channel, child: Child, stopped: bool, cleanup_failed: bool, @@ -232,8 +234,7 @@ impl Supervised { log: &Path, ) -> Result { let preflight = (|| { - let pair = - UnixStream::pair().map_err(|_| "Could not create agent supervision channel")?; + let pair = channel().map_err(|_| "Could not create agent supervision channel")?; let program = std::env::current_exe().map_err(|_| "Could not locate agent supervisor")?; let cwd = command.get_current_dir().ok_or("Missing agent workspace")?; @@ -276,13 +277,14 @@ impl Supervised { ); guardian .current_dir(cwd) - .stdin(Stdio::from(OwnedFd::from(other))) + .stdin(other) .stdout(Stdio::null()) .stderr(Stdio::null()); // A separate session keeps terminal/launcher group signals from killing // this lock owner, and is distinct from the listener's session. - use std::os::unix::process::CommandExt; + #[cfg(unix)] unsafe { + use std::os::unix::process::CommandExt; guardian.pre_exec(|| { if libc::setsid() == -1 { return Err(std::io::Error::last_os_error()); @@ -290,6 +292,13 @@ impl Supervised { Ok(()) }); } + // Windows: its own hidden console, outside a terminal's Ctrl+C group. + // An enclosing kill-on-close launcher job can still end the guardian. + #[cfg(windows)] + { + use std::os::windows::process::CommandExt; + guardian.creation_flags(windows_sys::Win32::System::Threading::CREATE_NO_WINDOW); + } let child = guardian.spawn().map_err(|_| { let _ = std::fs::remove_dir_all(temp); // no guardian can use it "Could not start agent supervisor" @@ -339,23 +348,8 @@ impl Supervised { Ok(true) } fn failed(&self) -> Result { - use std::os::fd::AsRawFd; - let mut byte = [0u8]; - let n = unsafe { - libc::recv( - self.socket.as_raw_fd(), - byte.as_mut_ptr().cast(), - 1, - libc::MSG_PEEK | libc::MSG_DONTWAIT, - ) - }; - if n < 0 { - if std::io::Error::last_os_error().kind() == std::io::ErrorKind::WouldBlock { - return Ok(false); - } - return Err("Could not inspect agent supervisor".into()); - } - Ok(n == 1 && byte == *b"F") + let byte = peek(&self.socket).map_err(|_| "Could not inspect agent supervisor")?; + Ok(byte == Some(b'F')) } pub(crate) fn stop(&mut self) -> Result<()> { if self.stopped { @@ -397,6 +391,55 @@ impl Drop for Supervised { } } +/// Unix transport: a socketpair on the guardian's stdin. +#[cfg(unix)] +mod unix { + use std::io; + use std::os::fd::{AsRawFd, FromRawFd, OwnedFd}; + pub(super) use std::os::unix::net::UnixStream as Channel; + use std::process::Stdio; + + pub(super) fn channel() -> io::Result<(Channel, Stdio)> { + let (ours, theirs) = Channel::pair()?; + Ok((ours, OwnedFd::from(theirs).into())) + } + + pub(super) fn inherited() -> io::Result { + Ok(unsafe { Channel::from_raw_fd(0) }) + } + + pub(super) fn log_channel() -> io::Result<(Channel, Stdio, Stdio)> { + let (reader, output) = Channel::pair()?; + reader.set_nonblocking(true)?; + let stderr = output.try_clone()?; + Ok(( + reader, + OwnedFd::from(output).into(), + OwnedFd::from(stderr).into(), + )) + } + + pub(super) fn peek(channel: &Channel) -> io::Result> { + let mut byte = [0u8]; + let n = unsafe { + libc::recv( + channel.as_raw_fd(), + byte.as_mut_ptr().cast(), + 1, + libc::MSG_PEEK | libc::MSG_DONTWAIT, + ) + }; + if n < 0 { + let error = io::Error::last_os_error(); + if error.kind() == io::ErrorKind::WouldBlock { + return Ok(None); + } + return Err(error); + } + Ok((n == 1).then_some(byte[0])) + } +} + #[cfg(test)] mod tests { #[test] @@ -405,7 +448,7 @@ mod tests { return; }; let args: Vec = serde_json::from_str(&args.to_string_lossy()).unwrap(); - let socket = unsafe { std::os::unix::net::UnixStream::from_raw_fd(0) }; + let socket = super::inherited().unwrap(); super::serve( socket, std::path::Path::new(&args[0]), @@ -416,13 +459,16 @@ mod tests { ) .unwrap(); } - use std::os::fd::FromRawFd; } -#[cfg(test)] +#[cfg(all(test, windows))] +mod windows_tests; + +#[cfg(all(test, unix))] mod lifecycle_tests { use super::*; use std::fs; + use std::os::unix::net::UnixStream; use std::time::Instant; const ID: &str = "79be667ef9dcbbac55a06295ce870b07029bfcdb2dce28d959f2815b16f81798-79be667ef9dcbbac55a06295ce870b07029bfcdb2dce28d959f2815b16f81798"; diff --git a/crates/agent-controller/src/supervisor/pipe.rs b/crates/agent-controller/src/supervisor/pipe.rs new file mode 100644 index 000000000..6dc05bcb5 --- /dev/null +++ b/crates/agent-controller/src/supervisor/pipe.rs @@ -0,0 +1,170 @@ +//! Windows stand-in for the Unix socketpair: one private duplex named-pipe +//! instance per guardian. Reads poll `PeekNamedPipe`, so a timeout never +//! leaves a blocking read queued ahead of a write on the synchronous handle. +use std::cell::Cell; +use std::fs::File; +use std::io::{self, Read, Write}; +use std::os::windows::io::{AsRawHandle, FromRawHandle, OwnedHandle}; +use std::process::Stdio; +use std::time::{Duration, Instant}; +use windows_sys::Win32::Foundation::{ + SetHandleInformation, ERROR_BROKEN_PIPE, ERROR_PIPE_NOT_CONNECTED, GENERIC_READ, GENERIC_WRITE, + HANDLE_FLAG_INHERIT, INVALID_HANDLE_VALUE, +}; +use windows_sys::Win32::Storage::FileSystem::{ + CreateFileW, FILE_FLAG_FIRST_PIPE_INSTANCE, OPEN_EXISTING, PIPE_ACCESS_DUPLEX, +}; +use windows_sys::Win32::System::Pipes::{ + CreateNamedPipeW, PeekNamedPipe, PIPE_REJECT_REMOTE_CLIENTS, PIPE_TYPE_BYTE, +}; + +pub(crate) struct Channel { + file: File, + timeout: Cell>, +} + +impl Channel { + fn new(handle: OwnedHandle) -> Self { + Self { + file: File::from(handle), + timeout: Cell::new(None), + } + } + + pub(crate) fn set_read_timeout(&self, timeout: Option) -> io::Result<()> { + self.timeout.set(timeout); + Ok(()) + } + + pub(crate) fn set_nonblocking(&self, nonblocking: bool) -> io::Result<()> { + self.timeout.set(nonblocking.then_some(Duration::ZERO)); + Ok(()) + } + + /// Readable byte count and the next byte, unconsumed. `None` once the peer + /// has closed and buffered bytes are drained. + fn poll(&self) -> io::Result> { + let (mut byte, mut read, mut available) = (0u8, 0u32, 0u32); + let peeked = unsafe { + PeekNamedPipe( + self.file.as_raw_handle(), + (&mut byte as *mut u8).cast(), + 1, + &mut read, + &mut available, + std::ptr::null_mut(), + ) + }; + if peeked == 0 { + let error = io::Error::last_os_error(); + return match error.raw_os_error().map(|code| code as u32) { + Some(ERROR_BROKEN_PIPE | ERROR_PIPE_NOT_CONNECTED) => Ok(None), + _ => Err(error), + }; + } + Ok(Some((available, byte))) + } +} + +impl Read for Channel { + fn read(&mut self, buf: &mut [u8]) -> io::Result { + let deadline = self.timeout.get().map(|timeout| Instant::now() + timeout); + loop { + match self.poll()? { + None => return Ok(0), + Some((0, _)) => {} + Some((available, _)) => { + let n = buf.len().min(available as usize); + return self.file.read(&mut buf[..n]); + } + } + if deadline.is_some_and(|deadline| Instant::now() >= deadline) { + return Err(io::ErrorKind::WouldBlock.into()); + } + std::thread::sleep(Duration::from_millis(10)); + } + } +} + +impl Write for Channel { + fn write(&mut self, buf: &[u8]) -> io::Result { + self.file.write(buf) + } + fn flush(&mut self) -> io::Result<()> { + Ok(()) + } +} + +/// Both ends start non-inheritable; std duplicates the guardian's end only +/// for the duration of that one spawn. +pub(crate) fn channel() -> io::Result<(Channel, Stdio)> { + let mut nonce = [0u8; 16]; + getrandom::fill(&mut nonce).map_err(|_| io::Error::other("no randomness"))?; + let name = format!( + r"\\.\pipe\buzz-agent-supervisor-{:032x}", + u128::from_ne_bytes(nonce) + ); + let name: Vec = name.encode_utf16().chain([0]).collect(); + // A single first instance: a squatter or competing client fails closed. + let server = unsafe { + CreateNamedPipeW( + name.as_ptr(), + PIPE_ACCESS_DUPLEX | FILE_FLAG_FIRST_PIPE_INSTANCE, + PIPE_TYPE_BYTE | PIPE_REJECT_REMOTE_CLIENTS, + 1, + 4096, + 4096, + 0, + std::ptr::null(), + ) + }; + if server == INVALID_HANDLE_VALUE { + return Err(io::Error::last_os_error()); + } + let server = unsafe { OwnedHandle::from_raw_handle(server) }; + let client = unsafe { + CreateFileW( + name.as_ptr(), + GENERIC_READ | GENERIC_WRITE, + 0, + std::ptr::null(), + OPEN_EXISTING, + 0, + std::ptr::null_mut(), + ) + }; + if client == INVALID_HANDLE_VALUE { + return Err(io::Error::last_os_error()); + } + let client = unsafe { OwnedHandle::from_raw_handle(client) }; + Ok((Channel::new(server), Stdio::from(client))) +} + +/// The guardian's stdin. Clear inheritance before anything is spawned so the +/// listener's tree cannot hold the app channel open past app death. +pub(crate) fn inherited() -> io::Result { + let handle = io::stdin().as_raw_handle(); + if unsafe { SetHandleInformation(handle, HANDLE_FLAG_INHERIT, 0) } == 0 { + return Err(io::Error::last_os_error()); + } + Ok(Channel::new(unsafe { + OwnedHandle::from_raw_handle(handle) + })) +} + +/// Listener stdout/stderr; the guardian-private read end never blocks. +pub(crate) fn log_channel() -> io::Result<(Channel, Stdio, Stdio)> { + let (reader, writer) = io::pipe()?; + let stderr = writer.try_clone()?; + let reader = Channel::new(reader.into()); + reader.set_nonblocking(true)?; + Ok((reader, writer.into(), stderr.into())) +} + +/// The next unread byte without consuming it; `None` when empty or closed. +pub(crate) fn peek(channel: &Channel) -> io::Result> { + Ok(channel + .poll()? + .filter(|(n, _)| *n > 0) + .map(|(_, byte)| byte)) +} diff --git a/crates/agent-controller/src/supervisor/windows_tests.rs b/crates/agent-controller/src/supervisor/windows_tests.rs new file mode 100644 index 000000000..724a1e069 --- /dev/null +++ b/crates/agent-controller/src/supervisor/windows_tests.rs @@ -0,0 +1,217 @@ +//! Native lifecycle: the listener tree lives in the guardian's job, and the +//! guardian keeps ownership and temp storage until that job is empty. +use super::*; +use std::fs; +use std::os::windows::io::{AsRawHandle, FromRawHandle, OwnedHandle}; +use std::time::Instant; +use windows_sys::Win32::Foundation::{INVALID_HANDLE_VALUE, WAIT_OBJECT_0}; +use windows_sys::Win32::System::Diagnostics::ToolHelp::{ + CreateToolhelp32Snapshot, Thread32First, Thread32Next, TH32CS_SNAPTHREAD, THREADENTRY32, +}; +use windows_sys::Win32::System::Threading::{ + GetExitCodeProcess, OpenProcess, OpenThread, ResumeThread, SuspendThread, TerminateProcess, + WaitForSingleObject, PROCESS_QUERY_LIMITED_INFORMATION, PROCESS_SYNCHRONIZE, PROCESS_TERMINATE, + THREAD_SUSPEND_RESUME, +}; + +const ID: &str = "79be667ef9dcbbac55a06295ce870b07029bfcdb2dce28d959f2815b16f81798-79be667ef9dcbbac55a06295ce870b07029bfcdb2dce28d959f2815b16f81798"; + +/// cmd root -> PowerShell listener -> ping worker. The root first runs short +/// cmd children whose handles it closes as they exit, so Stop must have opened +/// them as they joined. The listener logs a line, records both IDs, and exits +/// only once `exit` appears; the worker outlives it. +fn listener(root: &Path) -> Command { + fs::write( + root.join("listener.ps1"), + r#"[Console]::Out.WriteLine('listener-ready'); [Console]::Out.Flush() +$worker = Start-Process -PassThru -NoNewWindow "$env:SystemRoot\System32\PING.EXE" '-n 600 127.0.0.1' +Set-Content -LiteralPath "$PSScriptRoot\worker.tmp" -Value "$PID $($worker.Id)" +Rename-Item -LiteralPath "$PSScriptRoot\worker.tmp" -NewName worker +while (-not (Test-Path -LiteralPath "$PSScriptRoot\exit")) { Start-Sleep -Milliseconds 20 } +"#, + ) + .unwrap(); + fs::write( + root.join("listener.cmd"), + "@for /l %%i in (1,1,5) do @\"%SystemRoot%\\System32\\cmd.exe\" /d /c exit\r\n@\"%SystemRoot%\\System32\\WindowsPowerShell\\v1.0\\powershell.exe\" -NoProfile -NonInteractive -ExecutionPolicy Bypass -File \"%~dp0listener.ps1\"\r\n", + ) + .unwrap(); + let mut command = Command::new(root.join("listener.cmd")); + command.current_dir(root).envs(std::env::vars_os()); + command +} + +/// A process handle; dropping it terminates the process so failures never leak. +struct Proc(OwnedHandle); +impl Proc { + fn open(pid: u32) -> Self { + let access = PROCESS_SYNCHRONIZE | PROCESS_TERMINATE | PROCESS_QUERY_LIMITED_INFORMATION; + let handle = unsafe { OpenProcess(access, 0, pid) }; + assert!(!handle.is_null(), "process {pid} is not running"); + Self(unsafe { OwnedHandle::from_raw_handle(handle) }) + } + fn exited(&self) -> bool { + unsafe { WaitForSingleObject(self.0.as_raw_handle(), 0) == WAIT_OBJECT_0 } + } + fn exit_code(&self) -> u32 { + let waited = unsafe { WaitForSingleObject(self.0.as_raw_handle(), 30_000) }; + assert_eq!(waited, WAIT_OBJECT_0, "process did not exit"); + let mut code = 0; + assert_ne!( + unsafe { GetExitCodeProcess(self.0.as_raw_handle(), &mut code) }, + 0 + ); + code + } +} +impl Drop for Proc { + fn drop(&mut self) { + unsafe { TerminateProcess(self.0.as_raw_handle(), 1) }; + } +} + +fn read(path: &Path) -> String { + let deadline = Instant::now() + Duration::from_secs(30); + loop { + if let Ok(text) = fs::read_to_string(path) { + return text; + } + assert!(Instant::now() < deadline, "fixture did not reach {path:?}"); + std::thread::sleep(Duration::from_millis(10)); + } +} + +fn tree(root: &Path) -> (Proc, Proc) { + let text = read(&root.join("worker")); + let (listener, worker) = text.trim().split_once(' ').unwrap(); + ( + Proc::open(listener.parse().unwrap()), + Proc::open(worker.parse().unwrap()), + ) +} + +/// Suspend every guardian thread so fast cleanup cannot pass the order check by luck. +fn freeze(pid: u32) -> Vec { + let snapshot = unsafe { CreateToolhelp32Snapshot(TH32CS_SNAPTHREAD, 0) }; + assert_ne!(snapshot, INVALID_HANDLE_VALUE); + let snapshot = unsafe { OwnedHandle::from_raw_handle(snapshot) }; + let mut entry = THREADENTRY32 { + dwSize: std::mem::size_of::() as u32, + ..Default::default() + }; + let mut threads = Vec::new(); + let mut more = unsafe { Thread32First(snapshot.as_raw_handle(), &mut entry) } != 0; + while more { + if entry.th32OwnerProcessID == pid { + let thread = unsafe { OpenThread(THREAD_SUSPEND_RESUME, 0, entry.th32ThreadID) }; + assert!(!thread.is_null()); + let thread = unsafe { OwnedHandle::from_raw_handle(thread) }; + assert_ne!(unsafe { SuspendThread(thread.as_raw_handle()) }, u32::MAX); + threads.push(thread); + } + more = unsafe { Thread32Next(snapshot.as_raw_handle(), &mut entry) } != 0; + } + assert!(!threads.is_empty()); + threads +} + +#[test] +fn parent_entrypoint() { + let Some(dir) = std::env::var_os("BUZZ_SUPERVISOR_PARENT_TEST") else { + return; + }; + let dir = Path::new(&dir); + let run = Supervised::spawn( + &listener(dir), + &dir.join("locks"), + ID, + &dir.join("temp"), + &dir.join("harness.log"), + ) + .unwrap(); + fs::write(dir.join("ready.tmp"), run.child.id().to_string()).unwrap(); + fs::rename(dir.join("ready.tmp"), dir.join("ready")).unwrap(); + loop { + std::thread::park(); + } +} + +#[test] +fn app_death_keeps_lock_through_listener_tree_cleanup() { + let dir = tempfile::tempdir().unwrap(); + let root = dir.path(); + fs::create_dir(root.join("temp")).unwrap(); + let mut parent = Command::new(std::env::current_exe().unwrap()) + .args([ + "--exact", + "supervisor::windows_tests::parent_entrypoint", + "--nocapture", + ]) + .env("BUZZ_SUPERVISOR_PARENT_TEST", root) + .stdin(Stdio::null()) + .stdout(Stdio::null()) + .stderr(Stdio::null()) + .spawn() + .unwrap(); + let app = Proc::open(parent.id()); + let guardian_pid: u32 = read(&root.join("ready")).parse().unwrap(); + let guardian = Proc::open(guardian_pid); + let (listener, worker) = tree(root); + let frozen = freeze(guardian_pid); + drop(app); // TerminateProcess: the app dies without Stop. + parent.wait().unwrap(); + assert!(Ownership::acquire(&root.join("locks"), ID).is_err()); + assert!( + !guardian.exited(), + "guardian died with the tested launch context" + ); + assert!( + !worker.exited(), + "listener tree did not belong to the guardian" + ); + for thread in frozen { + assert_ne!(unsafe { ResumeThread(thread.as_raw_handle()) }, u32::MAX); + } + // Exit 0 means serve() confirmed an empty job and removed temp storage. + assert_eq!(guardian.exit_code(), 0, "guardian did not complete cleanup"); + assert!(listener.exited() && worker.exited()); + assert!(!root.join("temp").exists()); + let _lock = Ownership::acquire(&root.join("locks"), ID).unwrap(); +} + +#[test] +fn stop_and_root_exit_end_the_listener_tree_before_release() { + for stop in [true, false] { + let dir = tempfile::tempdir().unwrap(); + let root = dir.path(); + let temp = root.join("temp"); + fs::create_dir(&temp).unwrap(); + let mut run = Supervised::spawn( + &listener(root), + &root.join("locks"), + ID, + &temp, + &root.join("harness.log"), + ) + .unwrap(); + let (listener, worker) = tree(root); + if stop { + run.stop().unwrap(); + } else { + fs::write(root.join("exit"), b"").unwrap(); + let deadline = Instant::now() + Duration::from_secs(30); + while run.alive().unwrap() { + assert!(Instant::now() < deadline, "exited listener was not reaped"); + std::thread::sleep(Duration::from_millis(10)); + } + } + assert!(listener.exited() && worker.exited(), "descendant survived"); + assert!(!temp.exists()); + let log = fs::read_to_string(root.join("harness.log")).unwrap(); + assert!( + log.contains("listener-ready"), + "listener output was not captured" + ); + let _lock = Ownership::acquire(&root.join("locks"), ID).unwrap(); + } +} diff --git a/docs/agent-control.md b/docs/agent-control.md index 74859638c..7f846da0d 100644 --- a/docs/agent-control.md +++ b/docs/agent-control.md @@ -38,8 +38,9 @@ Stop. Late completion never closes a subsequently opened dialog. If an operation cannot be confirmed, refresh status before repeating it. Create is blocked with an explanation if this app’s runtime is unavailable; existing agents and profile retry remain intact. -The dev broker and native host must both support this flow. Packaged human -signing remains unavailable. +Without the dev broker, the native identity signs the owner authorization only +for the key this host prepared for the pending Create; other broker-only helpers +remain unavailable. **Not imported from old Buzz** is a separate collapsible section. Expanding it loads installed identities for the connected community; already-managed exact @@ -352,14 +353,30 @@ spawn, and exposes verification failure on the agent. Native Start projects one final snapshot after recording its outcome. This adds no incoming wake service for fully stopped listeners and no durable interrupted-turn recovery. -On Unix, an execed supervisor in the same app binary owns each agent's shared +An execed supervisor in the same app binary owns each agent's shared identity lock, isolated listener session, and temporary runtime directory. App Stop/Quit and kernel EOF on forced app death both trigger the existing whole-session teardown; ownership is released only after listener and workers have exited. Unconfirmed teardown keeps the lock and private directory, and normal Quit remains fail-closed. This does not clean up listeners orphaned before this fix, does not -contain a worker that deliberately escapes its session, and does not add process -containment on non-Unix platforms. +contain a worker that deliberately escapes its session, and force-killing the +supervisor itself releases ownership without confirmed teardown. + +On Windows the session is a kill-on-close job object. The listener starts +suspended and runs only after joining it; a failed assignment aborts Start. A +watcher opens each process as it joins. Stop terminates the job immediately, +without Unix's two-second cooperative cancel, and succeeds only once every process +the job ever admitted was opened and has exited. A lost job notification, or a +process gone before the watcher opened it, fails Stop closed and keeps ownership. +The windowless supervisor can still +be ended by an enclosing kill-on-close launcher job. Profile, config and temporary +directories inherit Windows ACLs; they are not verified to match Unix 0700/0600 modes. + +Non-Pi harnesses get the bundled tools first on PATH, then Windows' native PATH; +on Linux `~/.local/bin` and `/usr/local/bin` precede the system directories, which +are macOS's only entries. On Windows the shell tool needs Git Bash from Git for +Windows, or a `BUZZ_SHELL`/`GIT_BASH` override under Advanced → Environment. +Settings says **Shell setup not verified**; Buzz does not check it before Start. ## Ownership and handoff @@ -425,7 +442,11 @@ containment on non-Unix platforms. host must validate launch configuration and unsupported imported semantics before execution. - Harness and Provider choices come from native `harnessOptions` through the - injected Core snapshot. Buzz Agent offers Databricks v2. Goose appears with an + injected Core snapshot. Buzz Agent offers Databricks v2 and OpenAI; Windows + offers only OpenAI, which a new agent without a default uses, and explains an + inherited or saved Databricks provider as unsupported. OpenAI + uses the masked key field below as `OPENAI_COMPAT_API_KEY`, per agent or from + Agent defaults. Goose appears with an absolute executable path when the local CLI is installed, and offers common Goose providers plus a custom ID. A missing CLI leaves Goose disabled; the **Check again** action re-detects it after installation without an app @@ -513,7 +534,9 @@ containment on non-Unix platforms. native catalog, worker catalog and inference share this exact engine layout. Unix directories are owner-only; helper token files are owner-only. They are **not Keychain-encrypted**; other code running as your OS user can access them. - Non-Unix helper persistence remains memory-only. No old Buzz cache/Keychain or + Non-Unix helper persistence remains memory-only, so Windows refuses Databricks + Connect, catalog, credential open and Start with an explicit unsupported error + (OpenAI is unaffected). No old Buzz cache/Keychain or ambient `DATABRICKS_HOST`/`DATABRICKS_TOKEN` is read. - Disconnect requires Stop for all owned workers using the displayed workspace, retires pending starts for it, and removes only its app cache (retaining the lock diff --git a/docs/identity.md b/docs/identity.md index 7251e6bc1..4b1f32640 100644 --- a/docs/identity.md +++ b/docs/identity.md @@ -32,10 +32,9 @@ isolation from other programs running as the same OS user. There is no file/environment fallback, automatic legacy migration, human key replacement or human delete command. The existing **explicit agent import** may read only the selected old Buzz service/account and copy the selected agent key -into this app's separate agent namespace; it never writes the old blob. The -Windows/Linux agent adapters are backend groundwork: normal agent import/create -UI remains macOS-only (see [local agent controls](agent-control.md)). This storage -change does not enable those actions or Windows agent execution. +into this app's separate agent namespace; it never writes the old blob. Agent +import remains macOS-only; Create also saves new agent keys through the +Windows/Linux adapters (see [local agent controls](agent-control.md)). No user key belongs in release configuration. A shared credential blob can reduce repeated OS prompts by caching many credentials diff --git a/src-tauri/build.rs b/src-tauri/build.rs index eebe8283f..0664e6561 100644 --- a/src-tauri/build.rs +++ b/src-tauri/build.rs @@ -44,6 +44,7 @@ fn main() { "plugin_host_run_command", "plugin_host_request", "agent_control_create_prepare", + "agent_control_create_authorize", "agent_control_create_commit", "agent_control_creation_profile", "agent_control_snapshot", diff --git a/src-tauri/capabilities/default.json b/src-tauri/capabilities/default.json index dfd865970..919f8cc9c 100644 --- a/src-tauri/capabilities/default.json +++ b/src-tauri/capabilities/default.json @@ -22,6 +22,7 @@ "allow-plugin-host-run-command", "allow-plugin-host-request", "allow-agent-control-create-prepare", + "allow-agent-control-create-authorize", "allow-agent-control-create-commit", "allow-agent-control-creation-profile", "allow-agent-control-snapshot", diff --git a/src-tauri/src/agent_models.rs b/src-tauri/src/agent_models.rs index a2203e4af..550034977 100644 --- a/src-tauri/src/agent_models.rs +++ b/src-tauri/src/agent_models.rs @@ -499,6 +499,9 @@ impl RuntimeConnection { cache: &std::path::Path, opener: Arc, ) -> Result { + if cfg!(windows) { + return Err(buzz_agent_controller::connection::DATABRICKS_WINDOWS.into()); + } // Match the pinned runtime's discovery/client/scopes/namespace exactly. // Do not call the convenience wrapper: its default opener logs the URL. let auth = PkceOAuthTokenSource::new_with( @@ -609,5 +612,6 @@ async fn execute( #[cfg(test)] mod tests; -#[cfg(test)] +// Windows refuses this OAuth engine: see databricks_oauth_is_unsupported_on_windows. +#[cfg(all(test, unix))] mod bundled_tests; diff --git a/src-tauri/src/agent_models/tests.rs b/src-tauri/src/agent_models/tests.rs index 633aead0b..42743025c 100644 --- a/src-tauri/src/agent_models/tests.rs +++ b/src-tauri/src/agent_models/tests.rs @@ -485,6 +485,28 @@ fn native_discovery_preserves_absolute_harness_and_saved_or_draft_provider_overr } #[test] +#[cfg(windows)] +fn databricks_oauth_is_unsupported_on_windows() { + struct NoBrowser; + impl BrowserOpener for NoBrowser { + fn open(&self, _: &str) -> Result<(), String> { + panic!("Unsupported sign-in opened a browser"); + } + } + let dir = tempfile::tempdir().unwrap(); + let connection = RuntimeFactory.open( + "https://workspace.example.invalid", + dir.path(), + Arc::new(NoBrowser), + ); + assert_eq!( + connection.err().as_deref(), + Some(buzz_agent_controller::connection::DATABRICKS_WINDOWS) + ); + assert_eq!(std::fs::read_dir(dir.path()).unwrap().count(), 0); +} +#[test] +#[cfg(unix)] fn runtime_factory_no_ambient_auth_on_construction_or_empty_headless_refresh() { struct NoBrowser; impl BrowserOpener for NoBrowser { diff --git a/src-tauri/src/agents.rs b/src-tauri/src/agents.rs index 131fa4e0b..081714436 100644 --- a/src-tauri/src/agents.rs +++ b/src-tauri/src/agents.rs @@ -41,7 +41,12 @@ impl Snapshot { data, inventory_warnings: Vec::new(), import_available, - create_available: import_available, + // Platforms with a native credential store for the new identity. + create_available: cfg!(any( + target_os = "macos", + target_os = "windows", + target_os = "linux" + )), avatar_editing_available: true, local_inventory_actions: true, default_workspace: workspace.to_string_lossy().into_owned(), @@ -192,10 +197,17 @@ fn harness_options(app_data: &std::path::Path) -> Vec { status: "ready", install_supported: None, default_args: &[], - providers: &[ProviderOption { - value: "databricks_v2", - label: "Databricks v2", - }], + // Windows refuses Databricks sign-in (DATABRICKS_WINDOWS): omit it. + providers: &[ + ProviderOption { + value: "databricks_v2", + label: "Databricks v2", + }, + ProviderOption { + value: "openai", + label: "OpenAI", + }, + ][usize::from(cfg!(windows))..], }, HarnessOption { command: goose.as_ref().map_or_else( @@ -1173,6 +1185,29 @@ pub(crate) async fn agent_control_create_prepare( }) .await } +/// Owner attestation for the pending create's generated key only. The +/// identity may read OS credentials, so the agent host stays unlocked. +#[tauri::command] +pub(crate) async fn agent_control_create_authorize( + state: tauri::State<'_, AgentHost>, + identity: tauri::State<'_, crate::identity::IdentityHost>, + destination: String, + owner: String, + pubkey: String, +) -> Result, String> { + let (owner, pubkey) = run(state.inner().clone(), move |host| { + let (_, prepared) = host + .creating + .as_ref() + .ok_or("Create request expired; reopen Add agent")?; + if prepared.key.pubkey() != pubkey || !prepared.matches(&destination, &owner)? { + return Err("Authorization does not match the pending create request".into()); + } + Ok((owner, pubkey)) + }) + .await?; + identity.inner().authorize_agent(owner, pubkey).await +} #[tauri::command] pub(crate) async fn agent_control_create_commit( state: tauri::State<'_, AgentHost>, diff --git a/src-tauri/src/agents/tests.rs b/src-tauri/src/agents/tests.rs index b53b142e6..630c14682 100644 --- a/src-tauri/src/agents/tests.rs +++ b/src-tauri/src/agents/tests.rs @@ -68,6 +68,7 @@ pub(crate) fn fixture_with_models( .manage(host.clone()) .manage(crate::harness_setup::HarnessSetup::default()) .manage(model_host) + .manage(crate::identity::IdentityHost::fixture()) .invoke_handler(crate::commands()) .build(crate::app_context()) .unwrap(); @@ -419,12 +420,20 @@ fn real_ipc_snapshot_save_cas_stop_and_launch_gate() { let before = invoke(&view, "agent_control_snapshot", json!({})).unwrap(); assert_eq!(before["runtimeAvailable"], false); assert_eq!(before["importAvailable"], cfg!(target_os = "macos")); + assert_eq!(before["createAvailable"], true); + // Windows omits Databricks, whose sign-in it refuses; Unix lists it first. + let openai = json!({"value":"openai", "label":"OpenAI"}); + let providers = if cfg!(windows) { + json!([openai]) + } else { + json!([{"value":"databricks_v2", "label":"Databricks v2"}, openai]) + }; assert_eq!( before["harnessOptions"][0], json!({ "command":"buzz-agent", "label":"Buzz Agent", "available":true, "status":"ready", "defaultArgs":[], - "providers":[{"value":"databricks_v2", "label":"Databricks v2"}] + "providers": providers }) ); assert_eq!(before["harnessOptions"].as_array().unwrap().len(), 3); @@ -559,10 +568,11 @@ fn real_ipc_preview_source_no_import_and_shutdown_fence() { json!({"source":"development","destination":"wss://chosen.example"}), ) .unwrap(); - assert!(preview["sourcePath"] - .as_str() - .unwrap() - .ends_with("xyz.block.buzz.app.dev/agents/managed-agents.json")); + // Compare path components: Windows joins with `\`, Unix with `/`. + assert!( + std::path::Path::new(preview["sourcePath"].as_str().unwrap()) + .ends_with("xyz.block.buzz.app.dev/agents/managed-agents.json") + ); assert!(invoke( &view, "agent_control_import_preview", @@ -2175,6 +2185,66 @@ async fn native_create_waits_for_a_snapshot_and_keeps_its_prepared_identity() { assert_eq!(prepared, retried); } +#[test] +fn native_create_authorization_binds_the_prepared_key_owner_and_identity() { + let (dir, _host, _app, view) = fixture(); + let identity = invoke(&view, "identity_restore", json!({})).unwrap(); + let identity = identity.as_str().unwrap(); + let other = "cd".repeat(32); + let prepare = |owner: &str| { + let request = uuid::Uuid::new_v4().to_string(); + let prepared = invoke( + &view, + "agent_control_create_prepare", + json!({"requestId": request, "destination": "https://relay.example", "owner": owner}), + ) + .unwrap(); + (request, prepared["pubkey"].as_str().unwrap().to_owned()) + }; + let authorize = |owner: &str, pubkey: &str| { + invoke( + &view, + "agent_control_create_authorize", + json!({"destination": "https://relay.example", "owner": owner, "pubkey": pubkey}), + ) + }; + let mismatch = json!("Authorization does not match the pending create request"); + assert_eq!( + authorize(identity, &"ab".repeat(32)).unwrap_err(), + json!("Create request expired; reopen Add agent") + ); + // The request matches its prepared owner, but that owner is not this identity. + let (_, pubkey) = prepare(&other); + assert_eq!( + authorize(&other, &pubkey).unwrap_err(), + json!("The agent owner is not your signed-in identity") + ); + let (request, pubkey) = prepare(identity); + // A requested owner or key other than the prepared one is never signed. + assert_eq!(authorize(&other, &pubkey).unwrap_err(), mismatch); + assert_eq!(authorize(identity, &"ab".repeat(32)).unwrap_err(), mismatch); + let auth = authorize(identity, &pubkey).unwrap(); + assert_eq!(auth[1], identity); + let commit = |auth: &Value| { + let edit = json!({"name":"Created","systemPrompt":"","workspace":dir.path().to_str().unwrap(), + "harness":{"command":"buzz-agent","args":[],"model":"chosen","provider":"databricks_v2"},"environment":{}}); + invoke( + &view, + "agent_control_create_commit", + json!({"requestId": request, "edit": edit, "auth": auth.to_string()}), + ) + .unwrap_err() + }; + let mut forged = auth.clone(); + forged[3] = json!("00".repeat(64)); + assert_eq!( + commit(&forged), + json!("Owner attestation does not authorize this agent key") + ); + // The unchanged verifier accepts the real attestation; only synthetic custody refuses. + assert_eq!(commit(&auth), json!(IMPORT_GATE)); +} + #[tokio::test] async fn dropped_caller_does_not_release_a_running_native_operation() { let (_dir, host, _app, _view) = fixture(); diff --git a/src-tauri/src/browser_permissions_tests.rs b/src-tauri/src/browser_permissions_tests.rs index aa3997728..0c41a5de0 100644 --- a/src-tauri/src/browser_permissions_tests.rs +++ b/src-tauri/src/browser_permissions_tests.rs @@ -79,6 +79,7 @@ fn native_command_permissions_allow_only_main_webview() { "plugin_host_run_command", "plugin_host_request", "agent_control_create_prepare", + "agent_control_create_authorize", "agent_control_create_commit", "agent_control_creation_profile", "agent_control_snapshot", diff --git a/src-tauri/src/identity.rs b/src-tauri/src/identity.rs index 2ffaaaf0f..c7a99a31a 100644 --- a/src-tauri/src/identity.rs +++ b/src-tauri/src/identity.rs @@ -45,22 +45,39 @@ impl Key { return Err("Relay event is too large".into()); } let hash = Sha256::digest(serialized); + let signature = self.schnorr(&hash).ok_or("Could not sign relay event")?; + Ok(serde_json::json!({ + "id": format!("{hash:x}"), "pubkey": pubkey, + "created_at": event.created_at, "kind": event.kind, + "tags": event.tags, "content": event.content, "sig": signature + })) + } + /// Unconditional NIP-OA owner attestation for exactly this agent key. + fn authorize(&self, agent: &str) -> Result> { + let digest = Sha256::digest(format!("nostr:agent-auth:{agent}:")); + let signature = self + .schnorr(&digest) + .ok_or("Could not authorize the agent")?; + Ok(vec![ + "auth".into(), + self.viewer()?, + String::new(), + signature, + ]) + } + fn schnorr(&self, digest: &[u8]) -> Option { let secp = Secp256k1::signing_only(); - let mut secret = SecretKey::from_byte_array(*self.0).map_err(|_| INVALID)?; + let mut secret = SecretKey::from_byte_array(*self.0).ok()?; let mut pair = Keypair::from_secret_key(&secp, &secret); secret.non_secure_erase(); let mut random = Zeroizing::new([0; 32]); if getrandom::fill(random.as_mut()).is_err() { pair.non_secure_erase(); - return Err("Could not sign relay event".into()); + return None; } - let signature = secp.sign_schnorr_with_aux_rand(&hash, &pair, &random); + let signature = secp.sign_schnorr_with_aux_rand(digest, &pair, &random); pair.non_secure_erase(); - Ok(serde_json::json!({ - "id": format!("{hash:x}"), "pubkey": pubkey, - "created_at": event.created_at, "kind": event.kind, - "tags": event.tags, "content": event.content, "sig": signature.to_string() - })) + Some(signature.to_string()) } fn parse(text: &str) -> Result { let text = text.trim(); @@ -257,6 +274,24 @@ impl IdentityHost { }) .await } + + /// Callers bind `agent` to a key the app generated; the owner must be this identity. + pub(crate) async fn authorize_agent( + &self, + owner: String, + agent: String, + ) -> Result> { + with_identity(self.clone(), move |identity| { + if identity.restore()?.as_deref() != Some(owner.as_str()) { + return Err("The agent owner is not your signed-in identity".into()); + } + match &identity.state { + State::Ready(key) => key.authorize(&agent), + _ => Err("Set up your identity first".into()), + } + }) + .await + } } impl Default for IdentityHost { fn default() -> Self { diff --git a/src-tauri/src/lib.rs b/src-tauri/src/lib.rs index e470fc885..3b3a43235 100644 --- a/src-tauri/src/lib.rs +++ b/src-tauri/src/lib.rs @@ -27,11 +27,11 @@ mod managed_pi; mod pi_models; use agents::{ agent_control_action, agent_control_attach_mention, agent_control_clone_settings, - agent_control_create_commit, agent_control_create_prepare, agent_control_creation_profile, - agent_control_delete, agent_control_import_commit, agent_control_import_preview, - agent_control_log_challenge, agent_control_read_log, agent_control_save, - agent_control_save_defaults, agent_control_snapshot, agent_control_start_on_app_launch, - agent_control_use_here, AgentHost, + agent_control_create_authorize, agent_control_create_commit, agent_control_create_prepare, + agent_control_creation_profile, agent_control_delete, agent_control_import_commit, + agent_control_import_preview, agent_control_log_challenge, agent_control_read_log, + agent_control_save, agent_control_save_defaults, agent_control_snapshot, + agent_control_start_on_app_launch, agent_control_use_here, AgentHost, }; use buzzodz_plugins::{ imports::{prepare_folder, prepare_git, PreparedImport, Preview}, @@ -387,6 +387,7 @@ fn commands() -> impl Fn(tauri::ipc::Invoke) -> bool + Sen plugin_host_run_command, plugin_host_request, agent_control_create_prepare, + agent_control_create_authorize, agent_control_create_commit, agent_control_creation_profile, agent_control_snapshot, diff --git a/src-tauri/src/main.rs b/src-tauri/src/main.rs index d4c3f1ae9..a2650f6d4 100644 --- a/src-tauri/src/main.rs +++ b/src-tauri/src/main.rs @@ -1,7 +1,6 @@ #![cfg_attr(not(debug_assertions), windows_subsystem = "windows")] fn main() { - #[cfg(unix)] buzz_agent_controller::dispatch_agent_supervisor(); buzz_foundation_lib::run(); } diff --git a/src/bundled/agents/AgentCreateDialog.tsx b/src/bundled/agents/AgentCreateDialog.tsx index 76e14e550..8552d229f 100644 --- a/src/bundled/agents/AgentCreateDialog.tsx +++ b/src/bundled/agents/AgentCreateDialog.tsx @@ -40,7 +40,7 @@ function newAgentDraft(state: AgentControlState): AgentDraft { inherits || state.data?.agentDefaults?.provider ? "" - : "databricks_v2", + : (chosen?.providers[0]?.value ?? "databricks_v2"), environment: {}, }; } diff --git a/src/bundled/agents/AgentSettingsFields.test.tsx b/src/bundled/agents/AgentSettingsFields.test.tsx index 863108c91..afa3071f3 100644 --- a/src/bundled/agents/AgentSettingsFields.test.tsx +++ b/src/bundled/agents/AgentSettingsFields.test.tsx @@ -705,3 +705,95 @@ it("hides the key for an unknown saved provider until its override is removed", control.dispose(); } }); + +it("gives Buzz Agent OpenAI a write-only key unless a hidden override decides the provider", async () => { + const fixture = controlFixture(); + fixture.data.harnessOptions = [ + { + command: "buzz-agent", + label: "Buzz Agent", + providers: [ + { value: "databricks_v2", label: "Databricks v2" }, + { value: "openai", label: "OpenAI" }, + ], + }, + ]; + const control = createAgentControl(fixture.host); + const user = userEvent.setup(); + let draft!: AgentDraft; + function Editor({ savedKeys = [] }: { savedKeys?: string[] }) { + const [value, setValue] = useState(() => ({ + ...agentDraft(fixture.agent), + command: "buzz-agent", + provider: "openai", + model: "gpt-5", + environment: {}, + })); + draft = value; + return ( + setValue((current) => ({ ...current, ...patch }))} + /> + ); + } + const view = render(); + try { + const key = screen.getByLabelText("OpenAI API key"); + expect(key).toHaveAttribute("type", "password"); + await user.type(key, "sk-test"); + expect(agentEdit(draft).environment).toEqual({ + OPENAI_COMPAT_API_KEY: "sk-test", + }); + await user.click(screen.getByRole("combobox", { name: "Provider" })); + await user.click( + await screen.findByRole("option", { name: "Databricks v2" }), + ); + expect(screen.queryByLabelText("OpenAI API key")).toBeNull(); + expect(draft.environment).toEqual({}); + view.unmount(); + render(); + expect(screen.getByLabelText("OpenAI API key")).toHaveAttribute( + "placeholder", + "Saved key unchanged", + ); + expect(agentEdit(draft).environment).toEqual({}); + cleanup(); + render(); + expect(screen.queryByLabelText("OpenAI API key")).toBeNull(); + } finally { + cleanup(); + control.dispose(); + } +}); + +it("labels Windows Buzz Agent shell setup as unverified", () => { + const platform = vi.spyOn(navigator, "platform", "get"); + platform.mockReturnValue("Win32"); + const f = controlFixture(); + const control = createAgentControl(f.host); + try { + render( + , + ); + expect(screen.getByText(/Shell setup not verified/)).toBeVisible(); + } finally { + platform.mockRestore(); + control.dispose(); + } +}); diff --git a/src/bundled/agents/AgentSettingsFields.tsx b/src/bundled/agents/AgentSettingsFields.tsx index 0583ed16a..ad32b02ce 100644 --- a/src/bundled/agents/AgentSettingsFields.tsx +++ b/src/bundled/agents/AgentSettingsFields.tsx @@ -30,9 +30,42 @@ function effectiveGooseProvider(draft: AgentDraft, savedKeys: string[]) { return draft.provider; } -function providerApiKey(draft: AgentDraft, savedKeys: string[]) { +// Draft → Agent defaults → build floor, as native resolves it; null when a +// saved or global BUZZ_AGENT_PROVIDER override hides the effective value. +function effectiveBuzzProvider( + draft: AgentDraft, + savedKeys: string[], + data: AgentControlState["data"], +) { + const key = "BUZZ_AGENT_PROVIDER"; + const override = draft.environment[key]; + if (typeof override === "string") return override; + if ( + data?.defaultSettings?.environmentKeys.includes(key) || + (override === undefined && savedKeys.includes(key)) + ) + return null; + const inherited = data?.defaultSettings?.harness === "buzz-agent"; + return ( + draft.provider || + (inherited ? data?.defaultSettings?.provider : "") || + data?.agentDefaults?.provider + ); +} + +function providerApiKey( + draft: AgentDraft, + savedKeys: string[], + data: AgentControlState["data"], +) { if (draft.command.split("/").at(-1) === "buzz-pi-acp") return PI_API_KEYS[draft.provider]; + if (harnessKind(draft.command) === "buzz-agent") + return ["openai", "openai-compat"].includes( + effectiveBuzzProvider(draft, savedKeys, data) ?? "", + ) + ? { label: "OpenAI", env: "OPENAI_COMPAT_API_KEY" } + : undefined; if (!isGoose(draft.command)) return undefined; const provider = effectiveGooseProvider(draft, savedKeys); return provider ? gooseApiKey(provider) : undefined; @@ -80,12 +113,11 @@ export function AgentSettingsFields({ state.data?.defaultSettings?.harness === harnessKind(draft.command) ? state.data?.defaultSettings : undefined; - // Draft → same-harness Agent default → build floor, as native resolves it. - const buzzProvider = - draft.environment.BUZZ_AGENT_PROVIDER ?? - (draft.provider || - inherited?.provider || - state.data?.agentDefaults?.provider); + const buzzProvider = effectiveBuzzProvider( + draft, + environmentKeys, + state.data, + ); const modelDefaultKnown = (typeof draft.environment.BUZZ_AGENT_PROVIDER === "string" || !overridden("BUZZ_AGENT_PROVIDER")) && @@ -97,6 +129,7 @@ export function AgentSettingsFields({ ? effectiveGooseProvider(draft, environmentKeys) : null; const buzzAgent = harnessKind(draft.command) === "buzz-agent"; + const windows = /Win/i.test(globalThis.navigator?.platform ?? ""); // An environment selector can override the visible scalar default. const [modelKey, providerKey] = buzzAgent ? ["BUZZ_AGENT_MODEL", "BUZZ_AGENT_PROVIDER"] @@ -124,7 +157,7 @@ export function AgentSettingsFields({ host: inheritedKey("DATABRICKS_HOST"), filter: inheritedKey("DATABRICKS_MODEL_FILTER"), }; - const apiKey = providerApiKey(draft, environmentKeys); + const apiKey = providerApiKey(draft, environmentKeys, state.data); const savedKey = !!apiKey && environmentKeys.includes(apiKey.env); // Saved keys are write-only; only a key typed for this provider can be shown. const typedKey = apiKey && draft.environment[apiKey.env] ? apiKey.env : null; @@ -137,7 +170,8 @@ export function AgentSettingsFields({ // A typed key belongs to the provider it was entered for. if ( key && - providerApiKey({ ...draft, ...patch }, environmentKeys)?.env !== key + providerApiKey({ ...draft, ...patch }, environmentKeys, state.data) + ?.env !== key ) { const environment = { ...(patch.environment ?? draft.environment) }; if (typeof environment[key] === "string") { @@ -182,6 +216,19 @@ export function AgentSettingsFields({ onOpenHarnesses={onOpenHarnesses} discardEdits={discardEdits} /> + {buzzAgent && windows && ( +

+ Shell setup not verified. On Windows, the shell tool needs Git + Bash from Git for Windows, or a shell set with BUZZ_SHELL under + Advanced → Environment. Buzz does not check this before starting. +

+ )} + {buzzAgent && windows && databricks && ( +

+ Databricks sign-in is not supported on Windows yet. Choose OpenAI + for this agent. +

+ )} {goose && gooseProvider === null && (

This agent has a saved GOOSE_PROVIDER override whose value is @@ -222,9 +269,11 @@ export function AgentSettingsFields({ ? "Will remove on save" : savedKey ? "Saved key unchanged" - : pi - ? "Paste API key or use an existing Pi sign-in" - : "Paste API key or use existing Goose credentials" + : buzzAgent + ? "Paste API key" + : pi + ? "Paste API key or use an existing Pi sign-in" + : "Paste API key or use existing Goose credentials" } onChange={(event) => { const environment = { ...draft.environment }; @@ -237,11 +286,12 @@ export function AgentSettingsFields({

- {apiKey.env} is used for this agent and model lookup. Leave - blank to keep a saved key, if present, or use{" "} - {pi ? "your Pi sign-in" : "Goose credentials"}. Keys exported in - your shell profile are not used. Saved keys are stored in this - device’s local agent settings files. + {apiKey.env}{" "} + {buzzAgent + ? "is required for OpenAI. Leave blank to keep a saved key, if present, or use one from Agent defaults." + : `is used for this agent and model lookup. Leave blank to keep a saved key, if present, or use ${pi ? "your Pi sign-in" : "Goose credentials"}.`}{" "} + Keys exported in your shell profile are not used. Saved keys are + stored in this device’s local agent settings files.

)} diff --git a/src/bundled/agents/AgentsPage.test.tsx b/src/bundled/agents/AgentsPage.test.tsx index 7d464e54e..158db3f64 100644 --- a/src/bundled/agents/AgentsPage.test.tsx +++ b/src/bundled/agents/AgentsPage.test.tsx @@ -1846,6 +1846,86 @@ it("create copies only the default harness and shows inherited defaults", async sessionPolicy: "channel", }); }); +for (const platform of ["MacIntel", "Win32"] as const) { + it(`new Buzz Agent on ${platform} starts with the first listed provider and explains inherited Databricks`, async () => { + const windows = platform === "Win32"; + vi.spyOn(navigator, "platform", "get").mockReturnValue(platform); + const user = userEvent.setup(); + const create = vi.fn(); + vi.spyOn(communityApi, "communityRequest").mockResolvedValue({ auth: [] }); + const { f, control } = setup("connected", (fixture) => { + // Native omits Databricks on Windows, whose sign-in it refuses. + fixture.data.harnessOptions = [ + { + command: "buzz-agent", + label: "Buzz Agent", + providers: [ + { value: "databricks_v2", label: "Databricks v2" }, + { value: "openai", label: "OpenAI" }, + ].slice(windows ? 1 : 0), + }, + ]; + fixture.data.createAvailable = true; + fixture.data.defaultWorkspace = "/fixture/workspace"; + fixture.host.prepareCreate = async () => ({ + id: "created", + pubkey: "cd".repeat(32), + }); + fixture.host.commitCreate = create.mockImplementation(async () => { + fixture.data.agents.push({ + ...structuredClone(fixture.agent), + id: "created", + }); + return structuredClone(fixture.data); + }); + }); + const unsupported = + "Databricks sign-in is not supported on Windows yet. Choose OpenAI for this agent."; + const add = async () => { + await user.click( + await screen.findByRole("button", { name: "Add agent" }), + ); + const dialog = screen.getByRole("dialog"); + fireEvent.change(within(dialog).getByLabelText("Name"), { + target: { value: "Provider agent" }, + }); + return dialog; + }; + const created = async (dialog: HTMLElement, calls: number) => { + await user.click( + within(dialog).getByRole("button", { name: "Create agent" }), + ); + await waitFor(() => expect(create).toHaveBeenCalledTimes(calls)); + await waitFor(() => expect(screen.queryByRole("dialog")).toBeNull()); + return create.mock.calls[calls - 1]?.[1].harness.provider; + }; + // Without a default, the first listed provider is chosen explicitly. + let dialog = await add(); + expect(within(dialog).queryByText(unsupported)).toBeNull(); + expect(await created(dialog, 1)).toBe(windows ? "openai" : "databricks_v2"); + // An inherited Databricks default stays inherited, explained before Create + // on Windows, with OpenAI still selectable. + f.data.agentDefaults = { + provider: "databricks_v2", + model: "", + ownerOnly: true, + }; + await act(async () => control.refresh()); + dialog = await add(); + if (!windows) { + expect(within(dialog).queryByText(unsupported)).toBeNull(); + expect(await created(dialog, 2)).toBe(""); + return; + } + expect(within(dialog).getByText(unsupported)).toBeVisible(); + await user.click( + within(dialog).getByRole("combobox", { name: "Provider" }), + ); + await user.click(await screen.findByRole("option", { name: "OpenAI" })); + expect(within(dialog).queryByText(unsupported)).toBeNull(); + expect(await created(dialog, 2)).toBe("openai"); + }); +} it("qualifies management identities while keeping configured names and edit targets exact", async () => { const { f } = setup("ready", (fixture) => { fixture.data.agents.push({ diff --git a/src/features/communities/native-api.test.ts b/src/features/communities/native-api.test.ts index 4d3fc5b4c..1efcaf704 100644 --- a/src/features/communities/native-api.test.ts +++ b/src/features/communities/native-api.test.ts @@ -185,6 +185,23 @@ it("restores a signed community profile and preserves extra fields when publishi ).rejects.toThrow("not confirmed"); }); +it("authorizes a new agent only through the native pending-create command", async () => { + const auth = ["auth", key.pubkey, "", "ab".repeat(64)]; + vi.mocked(invoke).mockResolvedValueOnce(auth); + const request = { pubkey: "ba".repeat(32), owner: key.pubkey }; + await expect( + communityRequest(community, "authorize-agent", request), + ).resolves.toEqual({ auth }); + expect(invoke).toHaveBeenCalledExactlyOnceWith( + "agent_control_create_authorize", + { destination: community, ...request }, + ); + await expect( + communityRequest(community, "authorize-agent", { owner: key.pubkey }), + ).rejects.toThrow("Invalid agent owner authorization"); + expect(invoke).toHaveBeenCalledOnce(); +}); + it("keeps development requests on the existing broker even inside Tauri", async () => { vi.stubEnv("VITE_BUZZ_LIVE", "1"); vi.mocked(fetch).mockResolvedValue( diff --git a/src/features/communities/native-api.ts b/src/features/communities/native-api.ts index 41000bd92..160a54bcc 100644 --- a/src/features/communities/native-api.ts +++ b/src/features/communities/native-api.ts @@ -1,3 +1,4 @@ +import { invoke } from "@tauri-apps/api/core"; import { avatarPictureError } from "../profiles/avatar-upload"; import { connectNativeTransport, @@ -137,6 +138,18 @@ export async function nativeCommunityRequest( throw new Error("Profile publication was not confirmed"); return result; } + if (route === "authorize-agent") { + // Native signs only the key it generated for the pending create request. + const input = body as { pubkey?: unknown; owner?: unknown } | undefined; + if (typeof input?.pubkey !== "string" || typeof input.owner !== "string") + throw new Error("Invalid agent owner authorization"); + const auth = await invoke("agent_control_create_authorize", { + destination: community, + owner: input.owner, + pubkey: input.pubkey, + }); + return { auth }; + } throw new Error("This operation is unavailable on the packaged connection"); } diff --git a/tests/browser/settings.spec.mjs b/tests/browser/settings.spec.mjs index 4e800bab9..40267f402 100644 --- a/tests/browser/settings.spec.mjs +++ b/tests/browser/settings.spec.mjs @@ -396,8 +396,8 @@ confirmedPresence( await profile.click(); for (const width of [1280, 390]) { await page.setViewportSize({ width, height: 844 }); - if (await button(page, "Show navigation").isVisible()) - await button(page, "Show navigation").click(); + // The narrow toggle appears only after the resize's media-query change renders. + if (width === 390) await button(page, "Show navigation").click(); await expect(settingsSidebar).toBeVisible(); await profile.focus(); await tab(); diff --git a/tests/integration/browser-ci.test.mjs b/tests/integration/browser-ci.test.mjs index 3dc5390be..4ef43a6b4 100644 --- a/tests/integration/browser-ci.test.mjs +++ b/tests/integration/browser-ci.test.mjs @@ -230,7 +230,7 @@ test("automatic CI stays on Linux and manual dispatch runs only Windows", () => "-p buzz-foundation -p buzz-agent-controller -p buzz-credential-store"; assert.ok( windows.steps.some( - (step) => step.run === `cargo test ${packages} --locked`, + (step) => step.run === `cargo test ${packages} --locked --no-fail-fast`, ), "on-demand Windows validation retains complete tests for all native identity packages", );