diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index b4dbd2552d..b3793323bd 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -154,6 +154,7 @@ jobs: AGENTOS_SIDECAR_BIN: ${{ github.workspace }}/target/debug/agentos-native-sidecar run: | cargo test -p agentos-client \ + --features service-internals \ --lib \ --test scaffold \ --test e2e_smoke \ @@ -163,6 +164,7 @@ jobs: --test sidecar_pool_e2e \ --test cron_e2e \ --test cron_grammar_e2e \ + --test vm_config_compare_e2e \ -- --test-threads=1 - name: Stage stripped sidecar artifacts run: | diff --git a/crates/client/tests/vm_config_compare_e2e.rs b/crates/client/tests/vm_config_compare_e2e.rs new file mode 100644 index 0000000000..24e3648cac --- /dev/null +++ b/crates/client/tests/vm_config_compare_e2e.rs @@ -0,0 +1,45 @@ +//! Config comparison against a running VM through a real `agentos-native-sidecar`. +#![cfg(feature = "service-internals")] + +mod common; + +use std::collections::BTreeMap; + +use agentos_client::config::AgentOsConfig; +use agentos_client::service_internals::vm_config_equivalent; +use agentos_client::AgentOs; + +fn config(environment: Option>) -> AgentOsConfig { + AgentOsConfig { + environment, + ..Default::default() + } +} + +#[tokio::test] +async fn running_vm_compares_configs_through_the_sidecar() { + if !common::require_sidecar("running_vm_compares_configs_through_the_sidecar") { + return; + } + let vm = AgentOs::create(config(None)) + .await + .expect("create VM against real sidecar"); + + let same = vm_config_equivalent(&vm, &config(None), &config(None), Vec::new(), Vec::new()) + .await + .expect("compare an unchanged config"); + let changed_environment = Some(BTreeMap::from([(String::from("DEBUG"), String::from("1"))])); + let changed = vm_config_equivalent( + &vm, + &config(None), + &config(changed_environment), + Vec::new(), + Vec::new(), + ) + .await + .expect("compare a config with a changed environment"); + vm.shutdown().await.expect("shutdown VM"); + + assert!(same); + assert!(!changed); +} diff --git a/crates/native-sidecar/src/plugins/chunked_sqlite.rs b/crates/native-sidecar/src/plugins/chunked_sqlite.rs index d77861c17e..c12a45ddb7 100644 --- a/crates/native-sidecar/src/plugins/chunked_sqlite.rs +++ b/crates/native-sidecar/src/plugins/chunked_sqlite.rs @@ -24,6 +24,7 @@ use vfs::engine::CachedMetadataStore; #[cfg(test)] mod persistence_tests { use super::*; + use std::sync::atomic::{AtomicU8, Ordering}; #[test] fn metadata_and_blocks_survive_local_database_reopen() { @@ -96,14 +97,150 @@ mod persistence_tests { database.close().await.unwrap(); }); } + + const NO_FAULT: u8 = 0; + const COMMIT_THEN_FAIL: u8 = 1; + const FAIL_BEFORE_COMMIT: u8 = 2; + + /// Real local VM SQLite with one armed fault on the next transaction that + /// moves the index head. + struct HeadFaultDatabase { + inner: SharedVmSqliteDatabase, + fault: AtomicU8, + } + + impl HeadFaultDatabase { + fn arm(&self, fault: u8) { + self.fault.store(fault, Ordering::SeqCst); + } + } + + #[async_trait] + impl crate::vm_sqlite::VmSqliteDatabase for HeadFaultDatabase { + async fn query( + &self, + statement: SqlStatement, + ) -> Result { + self.inner.query(statement).await + } + + async fn transaction( + &self, + statements: Vec, + ) -> Result, crate::vm_sqlite::VmSqliteError> { + let moves_head = statements.iter().any(|statement| { + statement + .sql + .starts_with("INSERT INTO agentos_fs_metadata_heads") + }); + let fault = match moves_head { + true => self.fault.swap(NO_FAULT, Ordering::SeqCst), + false => NO_FAULT, + }; + match fault { + COMMIT_THEN_FAIL => { + self.inner.transaction(statements).await?; + Err(crate::vm_sqlite::VmSqliteError::Callback( + "reply lost after commit".to_owned(), + )) + } + FAIL_BEFORE_COMMIT => Err(crate::vm_sqlite::VmSqliteError::Callback( + "transaction failed".to_owned(), + )), + _ => self.inner.transaction(statements).await, + } + } + + async fn close(&self) -> Result<(), crate::vm_sqlite::VmSqliteError> { + self.inner.close().await + } + } + + #[test] + fn index_survives_a_lost_save_reply_followed_by_a_failed_save() { + let runtime = + agentos_runtime::SidecarRuntime::process(&agentos_runtime::RuntimeConfig::default()) + .unwrap(); + runtime.block_on(async { + let directory = tempfile::tempdir().unwrap(); + let database = crate::vm_sqlite::open_local_vm_sqlite( + directory.path().join("state.sqlite"), + runtime.context(), + 128 * 1024 * 1024, + ) + .await + .unwrap(); + bootstrap_schema(database.as_ref()).await.unwrap(); + let faults = std::sync::Arc::new(HeadFaultDatabase { + inner: database.clone(), + fault: AtomicU8::new(NO_FAULT), + }); + { + let metadata = SqliteMetadataStore::new( + faults.clone(), + "test".to_owned(), + DEFAULT_MAX_METADATA_BYTES, + ); + let root = metadata.resolve("/").await.unwrap(); + // A 64 KiB xattr on each of 40 directories makes the index larger + // than one save transaction can carry. + for index in 0..40u8 { + let mut attrs = CreateInodeAttrs::directory(0o755, 1000, 1000); + attrs.xattrs.insert( + "user.fill".to_owned(), + vec![index; vfs::engine::types::XATTR_SIZE_MAX], + ); + metadata + .create(root.ino, &format!("dir-{index}"), attrs) + .await + .unwrap(); + } + + faults.arm(COMMIT_THEN_FAIL); + metadata + .create( + root.ino, + "saved", + CreateInodeAttrs::directory(0o755, 1000, 1000), + ) + .await + .unwrap_err(); + faults.arm(FAIL_BEFORE_COMMIT); + metadata + .create( + root.ino, + "not-saved", + CreateInodeAttrs::directory(0o755, 1000, 1000), + ) + .await + .unwrap_err(); + } + + let metadata = SqliteMetadataStore::new( + database.clone(), + "test".to_owned(), + DEFAULT_MAX_METADATA_BYTES, + ); + metadata.resolve("/saved").await.unwrap(); + metadata.resolve("/dir-39").await.unwrap(); + assert!(metadata.resolve("/not-saved").await.is_err()); + database.close().await.unwrap(); + }); + } } const DEFAULT_METADATA_CACHE_ENTRIES: usize = 4096; const MAX_METADATA_CACHE_ENTRIES: usize = 1_000_000; -const METADATA_CHUNK_SIZE: usize = 256 * 1024; +/// Remote actor SQLite rejects a statement whose bound values exceed 128 KiB. +/// The other values in a metadata chunk write are at most 272 bytes, so 64 KiB +/// leaves ample room. +const METADATA_CHUNK_SIZE: usize = 64 * 1024; +/// VM SQLite requests are JSON, which writes each blob byte as up to 4 +/// characters, and a sidecar frame is at most 16 MiB. 32 statements carry at +/// most 2 MiB of index, which encodes to at most about 8 MiB. +const METADATA_SAVE_STATEMENTS_PER_TRANSACTION: usize = 32; const DEFAULT_MAX_METADATA_BYTES: usize = 64 * 1024 * 1024; const MAX_METADATA_BYTES: usize = 1024 * 1024 * 1024; -const METADATA_CLEANUP_BATCH_SIZE: i64 = 64; const MAX_CHUNK_SIZE: u32 = 16 * 1024 * 1024; const VFS_MIGRATION_1: &[&str] = &[ "CREATE TABLE agentos_fs_metadata_heads ( @@ -355,48 +492,54 @@ async fn persist_metadata( .checked_add(1) .ok_or_else(|| VfsError::eio("SQLite VFS metadata generation overflow"))?; - let chunks = dump.chunks(METADATA_CHUNK_SIZE).collect::>(); - for (chunk_index, content) in chunks.iter().enumerate() { - database - .query(SqlStatement::new( + let chunk_count = i64::try_from(dump.len().div_ceil(METADATA_CHUNK_SIZE)) + .map_err(|_| VfsError::eio("SQLite VFS metadata chunk count overflow"))?; + let byte_length = i64::try_from(dump.len()) + .map_err(|_| VfsError::eio("SQLite VFS metadata byte length overflow"))?; + // A load keeps reading the previous generation until the transaction that + // moves the head commits. Each batch is built only when it is sent. + let mut statements = (0_i64..) + .zip(dump.chunks(METADATA_CHUNK_SIZE)) + .map(|(chunk_index, content)| { + SqlStatement::new( "INSERT INTO agentos_fs_metadata_chunks (namespace, generation, chunk_index, content) VALUES (?, ?, ?, ?)", vec![ SqlValue::SqlText(namespace.to_owned()), SqlValue::SqlInteger(generation), - SqlValue::SqlInteger(i64::try_from(chunk_index).map_err(|_| { - VfsError::eio("SQLite VFS metadata chunk index overflow") - })?), + SqlValue::SqlInteger(chunk_index), SqlValue::SqlBlob(content.to_vec()), ], - )) - .await - .map_err(actor_sql_error)?; - } - - database - .query(SqlStatement::new( - "INSERT INTO agentos_fs_metadata_heads (namespace, generation, chunk_count, byte_length) VALUES (?, ?, ?, ?) \ - ON CONFLICT(namespace) DO UPDATE SET generation = excluded.generation, chunk_count = excluded.chunk_count, byte_length = excluded.byte_length", - vec![ - SqlValue::SqlText(namespace.to_owned()), - SqlValue::SqlInteger(generation), - SqlValue::SqlInteger(i64::try_from(chunks.len()).map_err(|_| { - VfsError::eio("SQLite VFS metadata chunk count overflow") - })?), - SqlValue::SqlInteger(i64::try_from(dump.len()).map_err(|_| { - VfsError::eio("SQLite VFS metadata byte length overflow") - })?), - ], - )) - .await - .map_err(actor_sql_error)?; - - if let Err(error) = cleanup_old_metadata(database, namespace, generation).await { - eprintln!( - "agentos chunked_sqlite failed to clean superseded metadata generations: {error}" - ); + ) + }) + .chain([ + SqlStatement::new( + "INSERT INTO agentos_fs_metadata_heads (namespace, generation, chunk_count, byte_length) VALUES (?, ?, ?, ?) \ + ON CONFLICT(namespace) DO UPDATE SET generation = excluded.generation, chunk_count = excluded.chunk_count, byte_length = excluded.byte_length", + vec![ + SqlValue::SqlText(namespace.to_owned()), + SqlValue::SqlInteger(generation), + SqlValue::SqlInteger(chunk_count), + SqlValue::SqlInteger(byte_length), + ], + ), + SqlStatement::new( + "DELETE FROM agentos_fs_metadata_chunks WHERE namespace = ? AND generation <> ?", + vec![ + SqlValue::SqlText(namespace.to_owned()), + SqlValue::SqlInteger(generation), + ], + ), + ]); + loop { + let batch = statements + .by_ref() + .take(METADATA_SAVE_STATEMENTS_PER_TRANSACTION) + .collect::>(); + if batch.is_empty() { + return Ok(()); + } + database.transaction(batch).await.map_err(actor_sql_error)?; } - Ok(()) } async fn load_metadata( @@ -474,49 +617,6 @@ async fn load_metadata( Ok(Some(dump)) } -async fn cleanup_old_metadata( - database: &SharedVmSqliteDatabase, - namespace: &str, - current_generation: i64, -) -> VfsResult<()> { - loop { - let result = database - .query(SqlStatement::new( - "SELECT generation, chunk_index FROM agentos_fs_metadata_chunks WHERE namespace = ? AND generation <> ? LIMIT ?", - vec![ - SqlValue::SqlText(namespace.to_owned()), - SqlValue::SqlInteger(current_generation), - SqlValue::SqlInteger(METADATA_CLEANUP_BATCH_SIZE), - ], - )) - .await - .map_err(actor_sql_error)?; - if result.rows.is_empty() { - return Ok(()); - } - for row in result.rows { - if row.len() != 2 { - return Err(VfsError::eio( - "SQLite returned malformed stale metadata row", - )); - } - let generation = sql_nonnegative_integer(&row[0], "metadata generation")?; - let chunk_index = sql_nonnegative_integer(&row[1], "metadata chunk index")?; - database - .query(SqlStatement::new( - "DELETE FROM agentos_fs_metadata_chunks WHERE namespace = ? AND generation = ? AND chunk_index = ?", - vec![ - SqlValue::SqlText(namespace.to_owned()), - SqlValue::SqlInteger(generation), - SqlValue::SqlInteger(chunk_index), - ], - )) - .await - .map_err(actor_sql_error)?; - } - } -} - fn first_integer(result: QueryResult, description: &str) -> VfsResult { let row = result .rows diff --git a/crates/native-sidecar/src/service.rs b/crates/native-sidecar/src/service.rs index 35e1a09e2a..22adffb08f 100644 --- a/crates/native-sidecar/src/service.rs +++ b/crates/native-sidecar/src/service.rs @@ -2607,6 +2607,7 @@ where | RequestRoute::GetZombieTimerCount(_) | RequestRoute::ProvidedCommands(_) | RequestRoute::ListMounts(_) + | RequestRoute::CompareVmConfig(_) | RequestRoute::GuestFilesystemCall(_) | RequestRoute::GuestKernelCall(_) | RequestRoute::BootstrapRootFilesystem(_) @@ -2801,6 +2802,12 @@ where let future = self.list_mounts(&request, payload); Ok(Some(PreparedRequest::from_vm_command(request, future))) } + RequestRoute::CompareVmConfig(payload) => { + let result = self.compare_vm_config(&request, payload); + Ok(Some(PreparedRequest::from_future(request, async move { + result + }))) + } RequestRoute::BootstrapRootFilesystem(payload) => { let future = self.bootstrap_root_filesystem(&request, payload.entries); Ok(Some(PreparedRequest::from_vm_command(request, future))) @@ -3032,9 +3039,7 @@ where }) }))) } - RequestRoute::CreateVm(_) - | RequestRoute::CompareVmConfig(_) - | RequestRoute::DisposeVm(_) => { + RequestRoute::CreateVm(_) | RequestRoute::DisposeVm(_) => { unreachable!("VM creation and disposal use dedicated prepared routes") } }