Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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 \
Expand All @@ -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: |
Expand Down
45 changes: 45 additions & 0 deletions crates/client/tests/vm_config_compare_e2e.rs
Original file line number Diff line number Diff line change
@@ -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<BTreeMap<String, String>>) -> 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);
}
260 changes: 180 additions & 80 deletions crates/native-sidecar/src/plugins/chunked_sqlite.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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() {
Expand Down Expand Up @@ -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<QueryResult, crate::vm_sqlite::VmSqliteError> {
self.inner.query(statement).await
}

async fn transaction(
&self,
statements: Vec<SqlStatement>,
) -> Result<Vec<QueryResult>, 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;
Comment thread
eersnington marked this conversation as resolved.
/// 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 (
Expand Down Expand Up @@ -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::<Vec<_>>();
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::<Vec<_>>();
if batch.is_empty() {
return Ok(());
}
database.transaction(batch).await.map_err(actor_sql_error)?;
}
Ok(())
}

async fn load_metadata(
Expand Down Expand Up @@ -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<i64> {
let row = result
.rows
Expand Down
11 changes: 8 additions & 3 deletions crates/native-sidecar/src/service.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2607,6 +2607,7 @@ where
| RequestRoute::GetZombieTimerCount(_)
| RequestRoute::ProvidedCommands(_)
| RequestRoute::ListMounts(_)
| RequestRoute::CompareVmConfig(_)
| RequestRoute::GuestFilesystemCall(_)
| RequestRoute::GuestKernelCall(_)
| RequestRoute::BootstrapRootFilesystem(_)
Expand Down Expand Up @@ -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)))
Expand Down Expand Up @@ -3032,9 +3039,7 @@ where
})
})))
}
RequestRoute::CreateVm(_)
| RequestRoute::CompareVmConfig(_)
| RequestRoute::DisposeVm(_) => {
RequestRoute::CreateVm(_) | RequestRoute::DisposeVm(_) => {
unreachable!("VM creation and disposal use dedicated prepared routes")
}
}
Expand Down