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
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

6 changes: 4 additions & 2 deletions crates/agentos-sidecar/src/acp/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1308,8 +1308,10 @@ mod tests {

#[test]
fn configured_acp_limit_errors_preserve_stable_wire_codes() {
let mut limits = AcpLimits::default();
limits.max_prompt_bytes = 3;
let mut limits = AcpLimits {
max_prompt_bytes: 3,
..AcpLimits::default()
};
let bytes_error = parse_content_blocks("[{}]", "main", &limits)
.expect_err("prompt bytes must be bounded");
assert_eq!(error_code(&bytes_error), "acp_prompt_bytes_limit");
Expand Down
2 changes: 1 addition & 1 deletion crates/agentos-sidecar/src/acp/restore.rs
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,7 @@ impl AcpExtension {
.restore_acp_runtime(
ctx,
RestoreRuntimeRequest {
acp_session_id: acp_session_id,
acp_session_id,
agent_type: session.agent.clone(),
cwd: session.cwd.clone(),
env,
Expand Down
13 changes: 7 additions & 6 deletions crates/agentos-sidecar/src/acp/runtime.rs
Original file line number Diff line number Diff line change
Expand Up @@ -355,9 +355,9 @@ impl AcpExtension {
.await;
matches!(&sigterm, Err(error) if is_process_already_gone_error(error))
};
let terminated = if adapter_already_gone {
true
} else if wait_for_process_exit(ctx, &session.process_id, SESSION_CLOSE_TIMEOUT).await {
let terminated = if adapter_already_gone
|| wait_for_process_exit(ctx, &session.process_id, SESSION_CLOSE_TIMEOUT).await
{
true
} else {
let sigkill = ctx
Expand Down Expand Up @@ -404,12 +404,13 @@ impl AcpExtension {
Ok(())
}

#[allow(clippy::needless_option_as_deref)]
pub(super) async fn send_runtime_request_with_sink(
&self,
ctx: &mut ExtensionContext<'_>,
request: AcpSessionRequest,
mut durable_sink: Option<&mut DurableUpdateSink>,
mut cancellation: Option<&mut tokio::sync::watch::Receiver<bool>>,
cancellation: Option<&mut tokio::sync::watch::Receiver<bool>>,
) -> AcpHandlerOutput {
let params = match request
.params
Expand Down Expand Up @@ -488,7 +489,7 @@ impl AcpExtension {
&mut stdout_buffer,
Some(&acp_session_id),
durable_sink.as_deref_mut(),
cancellation.as_deref_mut(),
cancellation,
)
.await
{
Expand Down Expand Up @@ -564,7 +565,7 @@ impl AcpExtension {
};
match synthetic {
Ok(Some(notification)) => {
let handled = if let Some(sink) = durable_sink.as_deref_mut() {
let handled = if let Some(sink) = durable_sink {
match serde_json::from_str::<Value>(&notification) {
Ok(notification) => match sink
.handle_notification(ctx, &notification, &mut exchange.events)
Expand Down
32 changes: 32 additions & 0 deletions crates/agentos-sidecar/src/acp/turn.rs
Original file line number Diff line number Diff line change
Expand Up @@ -553,6 +553,36 @@ impl AcpExtension {
if output.response.is_err() {
return output;
}
let response = match &output.response {
Ok(AcpResponse::AcpSessionRpcResponse(response)) => {
match serde_json::from_str::<Value>(&response.response) {
Ok(response) => response,
Err(error) => {
return AcpHandlerOutput {
response: Err(SidecarError::InvalidState(format!(
"invalid ACP config update response JSON: {error}"
))),
events: output.events,
};
}
}
}
Ok(_) => {
return AcpHandlerOutput {
response: Err(SidecarError::InvalidState(String::from(
"invalid ACP config update response",
))),
events: output.events,
};
}
Err(_) => unreachable!("ACP transport errors returned above"),
};
if let Err(error) = response_result(response, "ACP session/set_config_option") {
return AcpHandlerOutput {
response: Err(error),
events: output.events,
};
}
if let Err(error) = sink.flush(ctx, &mut output.events).await {
return AcpHandlerOutput {
response: Err(error),
Expand Down Expand Up @@ -723,6 +753,7 @@ impl DurableUpdateSink {
})
}

#[allow(clippy::too_many_arguments)]
pub(super) async fn handle_permission_request(
&mut self,
ctx: &mut ExtensionContext<'_>,
Expand Down Expand Up @@ -1144,6 +1175,7 @@ pub(super) fn decode_durable_event(event_json: &str) -> Result<AcpDurableEvent,
}
}

#[allow(clippy::too_many_arguments)]
async fn wait_for_permission_signal(
ctx: &mut ExtensionContext<'_>,
process_id: &str,
Expand Down
86 changes: 49 additions & 37 deletions crates/agentos-sidecar/src/session_store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -262,7 +262,12 @@ impl SessionStore {
vec![SqlValue::SqlText(session_id.to_owned())],
))
.await?;
match result.rows.first().map(decode_session).transpose()? {
match result
.rows
.first()
.map(|row| decode_session(row))
.transpose()?
{
Some(mut session) => {
self.hydrate_pending_state(&mut session).await?;
Ok(Some(session))
Expand All @@ -271,6 +276,7 @@ impl SessionStore {
}
}

#[allow(clippy::too_many_arguments)]
pub async fn create(
&self,
session_id: &str,
Expand Down Expand Up @@ -409,7 +415,7 @@ impl SessionStore {
let mut sessions = result
.rows
.iter()
.map(decode_session_summary)
.map(|row| decode_session_summary(row))
.collect::<Result<Vec<_>, _>>()?;
self.hydrate_pending_summaries(&mut sessions).await?;
Ok(sessions)
Expand Down Expand Up @@ -457,7 +463,11 @@ impl SessionStore {
vec![text(session_id), text(idempotency_key)],
))
.await?;
result.rows.first().map(decode_prompt).transpose()
result
.rows
.first()
.map(|row| decode_prompt(row))
.transpose()
}

pub async fn accept_prompt(
Expand Down Expand Up @@ -550,9 +560,7 @@ impl SessionStore {
let results = match self.database.transaction(statements).await {
Ok(results) => results,
Err(error @ VmSqliteError::UnexpectedChanges { .. }) => {
if let Err(limit_error) = self.ensure_prompt_capacity(session_id).await {
return Err(limit_error);
}
self.ensure_prompt_capacity(session_id).await?;
return Err(error);
}
Err(error) => return Err(error),
Expand Down Expand Up @@ -763,9 +771,7 @@ impl SessionStore {
{
Ok(results) => results,
Err(error @ VmSqliteError::UnexpectedChanges { .. }) => {
if let Err(limit_error) = self.ensure_pending_capacity(session_id).await {
return Err(limit_error);
}
self.ensure_pending_capacity(session_id).await?;
return Err(error);
}
Err(error) => return Err(error),
Expand Down Expand Up @@ -1840,15 +1846,13 @@ fn validate_event(event: &Value) -> Result<&str, VmSqliteError> {
event.get("reason").and_then(Value::as_str),
) {
(Some("accepted"), None) => {}
(Some("not_pending"), Some(reason))
if matches!(
reason,
"already_resolved"
| "prompt_cancelled"
| "adapter_exited"
| "session_deleted"
| "vm_shutdown"
) => {}
(
Some("not_pending"),
Some(
"already_resolved" | "prompt_cancelled" | "adapter_exited"
| "session_deleted" | "vm_shutdown",
),
) => {}
_ => {
return Err(VmSqliteError::InvalidResult(
"permission_response event status/reason combination is invalid".to_owned(),
Expand Down Expand Up @@ -2070,7 +2074,7 @@ fn derive_state_json(
serde_json::to_string(&value).map_err(|error| VmSqliteError::InvalidResult(error.to_string()))
}

fn decode_session(row: &Vec<SqlValue>) -> Result<StoredSession, VmSqliteError> {
fn decode_session(row: &[SqlValue]) -> Result<StoredSession, VmSqliteError> {
if row.len() != 26 {
return Err(VmSqliteError::InvalidResult(format!(
"session row has {} columns, expected 26",
Expand Down Expand Up @@ -2127,7 +2131,7 @@ fn decode_session(row: &Vec<SqlValue>) -> Result<StoredSession, VmSqliteError> {
})
}

fn decode_session_summary(row: &Vec<SqlValue>) -> Result<StoredSessionSummary, VmSqliteError> {
fn decode_session_summary(row: &[SqlValue]) -> Result<StoredSessionSummary, VmSqliteError> {
if row.len() != 12 {
return Err(VmSqliteError::InvalidResult(format!(
"session summary row has {} columns, expected 12",
Expand Down Expand Up @@ -2229,7 +2233,7 @@ fn decode_pending_resolution(
Ok(PendingRequestResolution::Terminal { reason, event })
}

fn decode_prompt(row: &Vec<SqlValue>) -> Result<StoredPrompt, VmSqliteError> {
fn decode_prompt(row: &[SqlValue]) -> Result<StoredPrompt, VmSqliteError> {
if row.len() != 5 {
return Err(VmSqliteError::InvalidResult(format!(
"prompt row has {} columns, expected 5",
Expand Down Expand Up @@ -2559,9 +2563,11 @@ mod tests {
.await
.expect("database");

let mut limits = AcpLimits::default();
limits.max_session_history_events = 3;
limits.max_session_history_bytes = 1_000_000;
let mut limits = AcpLimits {
max_session_history_events: 3,
max_session_history_bytes: 1_000_000,
..AcpLimits::default()
};
let store = SessionStore::open(database.clone())
.await
.expect("store")
Expand Down Expand Up @@ -2749,8 +2755,10 @@ mod tests {
.expect("session");
assert_eq!(stale.oldest_retained_sequence, 1);

let mut limits = AcpLimits::default();
limits.max_session_history_events = 2;
let limits = AcpLimits {
max_session_history_events: 2,
..AcpLimits::default()
};
let request_store = SessionStore::from_database(database).with_limits(&limits);
let refreshed = request_store
.enforce_history_retention("main")
Expand Down Expand Up @@ -3613,14 +3621,16 @@ mod tests {
)
.await
.expect("database");
let mut limits = AcpLimits::default();
limits.max_sessions_per_vm = 1;
limits.max_prompts_per_session = 1;
limits.max_prompts_per_vm = 1;
limits.max_pending_permissions_per_session = 1;
limits.max_pending_permissions_per_vm = 1;
limits.max_permission_outcomes_per_session = 1;
limits.max_permission_outcomes_per_vm = 1;
let limits = AcpLimits {
max_sessions_per_vm: 1,
max_prompts_per_session: 1,
max_prompts_per_vm: 1,
max_pending_permissions_per_session: 1,
max_pending_permissions_per_vm: 1,
max_permission_outcomes_per_session: 1,
max_permission_outcomes_per_vm: 1,
..AcpLimits::default()
};
let store = SessionStore::open(database).await.expect("store").with_limits(&limits);
store.create("main", "agent", "native", "/workspace", r#"{"permissionPolicy":"ask"}"#, None, None, "[]").await.expect("create");
assert!(matches!(
Expand Down Expand Up @@ -3660,9 +3670,11 @@ mod tests {
)
.await
.expect("database");
let mut limits = AcpLimits::default();
limits.max_prompts_per_session = 1;
limits.max_prompts_per_vm = 1;
let mut limits = AcpLimits {
max_prompts_per_session: 1,
max_prompts_per_vm: 1,
..AcpLimits::default()
};
let prompt_store = SessionStore::open(database.clone())
.await
.expect("store")
Expand Down
8 changes: 5 additions & 3 deletions crates/agentos-sidecar/src/session_store/performance_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -448,9 +448,11 @@ async fn exercise_near_limit_append(
populate_history(&database, HISTORY_COMPLEXITY_LIMIT - 1).await;

let recording = RecordingDatabase::wrap(database);
let mut limits = AcpLimits::default();
limits.max_session_history_events = HISTORY_COMPLEXITY_LIMIT;
limits.max_session_history_bytes = 16 * 1024 * 1024;
let limits = AcpLimits {
max_session_history_events: HISTORY_COMPLEXITY_LIMIT,
max_session_history_bytes: 16 * 1024 * 1024,
..AcpLimits::default()
};
let store: SharedVmSqliteDatabase = recording.clone();
let store = SessionStore::from_database(store).with_limits(&limits);
if let Some(metrics) = wire_metrics {
Expand Down
3 changes: 2 additions & 1 deletion crates/client/src/fs.rs
Original file line number Diff line number Diff line change
Expand Up @@ -926,12 +926,13 @@ impl AgentOs {
async fn reconfigure_dynamic_mounts(&self) -> Result<()> {
let inner = self.inner();
let config = &inner.config;
let mounts = inner.dynamic_mounts.lock().clone();
let response = self
.transport()
.request_wire(
self.fs_vm_scope(),
wire::RequestPayload::ConfigureVmRequest(wire::ConfigureVmRequest {
mounts: inner.dynamic_mounts.lock().clone(),
mounts,
software: Vec::new(),
permissions: Some(crate::agent_os::permissions_policy(config)),
module_access_cwd: None,
Expand Down
3 changes: 2 additions & 1 deletion crates/kernel/src/kernel.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2284,6 +2284,7 @@ impl<F: VirtualFileSystem + 'static> KernelVm<F> {
.set_xattr(path, name, value, flags, follow_symlinks)?)
}

#[allow(clippy::too_many_arguments)]
pub fn set_xattr_for_process(
&mut self,
requester_driver: &str,
Expand Down Expand Up @@ -9078,7 +9079,7 @@ struct PosixAcl {

impl PosixAcl {
fn parse(value: &[u8], path: &str) -> KernelResult<Self> {
if value.len() < 4 || (value.len() - 4) % 8 != 0 {
if value.len() < 4 || !(value.len() - 4).is_multiple_of(8) {
return Err(invalid_acl(path, "invalid xattr length"));
}
let entry_count = (value.len() - 4) / 8;
Expand Down
3 changes: 2 additions & 1 deletion crates/native-sidecar-core/src/limits.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1928,7 +1928,8 @@ mod tests {
.to_string()
.contains("limits.resources.maxSocketBufferedBytes"));

let acp_relationship_cases: [(&str, &str, fn(&mut VmLimits)); 3] = [
type AcpRelationshipCase = (&'static str, &'static str, fn(&mut VmLimits));
let acp_relationship_cases: [AcpRelationshipCase; 3] = [
(
"limits.acp.maxPromptsPerSession",
"limits.acp.maxPromptsPerVm",
Expand Down
Loading
Loading