Skip to content
Open
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
7 changes: 7 additions & 0 deletions crates/execution/src/javascript.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1856,6 +1856,13 @@ impl JavascriptExecution {
self.exited.load(Ordering::Acquire)
}

/// Whether the runtime event bridge has already made another event
/// durable for this execution. This is a non-consuming readiness probe;
/// callers still own bounded draining and ordering.
pub fn has_pending_events(&self) -> bool {
!self.events.is_empty()
}

/// Run another sidecar-managed operation in this execution's retained V8
/// context. Public clients submit semantic language requests; only the
/// sidecar calls this executor primitive.
Expand Down
4 changes: 4 additions & 0 deletions crates/execution/src/python.rs
Original file line number Diff line number Diff line change
Expand Up @@ -520,6 +520,10 @@ impl PythonExecution {
self.inner.uses_shared_v8_runtime()
}

pub fn has_pending_events(&self) -> bool {
self.inner.has_pending_events()
}

/// Run another sidecar-managed operation in the retained Pyodide
/// interpreter owned by this execution.
pub fn execute_retained(&mut self, source: String) -> Result<(), PythonExecutionError> {
Expand Down
6 changes: 6 additions & 0 deletions crates/execution/src/wasm.rs
Original file line number Diff line number Diff line change
Expand Up @@ -503,6 +503,12 @@ impl WasmExecution {
self.inner.uses_shared_v8_runtime()
}

pub fn has_pending_events(&self) -> bool {
!self.pending_events.is_empty()
|| !self.internal_sync_rpc.pending_events.is_empty()
|| self.inner.has_pending_events()
}

pub fn start_prepared(&mut self) -> Result<(), WasmExecutionError> {
self.inner.start_prepared().map_err(map_javascript_error)?;
self.execution_started_at = Instant::now();
Expand Down
32 changes: 32 additions & 0 deletions crates/kernel/src/command_registry.rs
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,13 @@ impl CommandRegistry {
pub fn register(&mut self, driver: CommandDriver) -> VfsResult<()> {
driver.validate_commands()?;

// Registering one logical driver is replacement, not an append. This is
// required for transactional registries (for example host callbacks):
// rolling a driver back to its previous command set must make aliases
// introduced by the failed update unresolvable.
self.commands
.retain(|_, existing| existing.name() != driver.name());

for command in &driver.commands {
if let Some(existing) = self.commands.get(command) {
self.warnings.push(format!(
Expand Down Expand Up @@ -142,3 +149,28 @@ fn validate_command_name(command: &str) -> VfsResult<()> {

Ok(())
}

#[cfg(test)]
mod tests {
use super::*;

#[test]
fn registering_same_driver_replaces_its_command_set() {
let mut registry = CommandRegistry::new();
registry
.register(CommandDriver::new("bindings", ["old", "temporary"]))
.expect("register initial driver");
registry
.register(CommandDriver::new("bindings", ["old"]))
.expect("replace driver commands");

assert_eq!(
registry.resolve("old").map(CommandDriver::name),
Some("bindings")
);
assert!(
registry.resolve("temporary").is_none(),
"aliases removed by a driver refresh must not remain executable"
);
}
}
185 changes: 124 additions & 61 deletions crates/native-sidecar/src/bindings.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,7 @@ use crate::protocol::{
RequestFrame, ResponsePayload,
};
use crate::service::{kernel_error, normalize_path, DispatchResult};
use crate::state::{BridgeError, VmState, BINDING_DRIVER_NAME};
use crate::state::{BridgeError, SharedBridge, VmHandle, VmState, BINDING_DRIVER_NAME};
use crate::{NativeSidecar, NativeSidecarBridge, SidecarError};
use agentos_kernel::command_registry::CommandDriver;
use agentos_native_sidecar_core::bindings::{
Expand All @@ -24,6 +24,7 @@ pub(crate) use agentos_native_sidecar_core::bindings::{
use agentos_native_sidecar_core::permissions::{
allow_all_policy, deny_all_policy, evaluate_permissions_policy,
};
use agentos_native_sidecar_core::respond as shared_respond;
use agentos_vm_config::PermissionMode;
use serde_json::{json, Map, Number, Value};
use std::collections::{BTreeMap, BTreeSet};
Expand Down Expand Up @@ -52,73 +53,135 @@ pub(crate) fn register_host_callbacks<B>(
sidecar: &mut NativeSidecar<B>,
request: &RequestFrame,
payload: RegisterHostCallbacksRequest,
) -> Result<DispatchResult, SidecarError>
) -> impl std::future::Future<Output = Result<DispatchResult, SidecarError>> + 'static
where
B: NativeSidecarBridge + Send + 'static,
BridgeError<B>: fmt::Debug + Send + Sync + 'static,
{
let (connection_id, session_id, vm_id) = sidecar.vm_scope_for(&request.ownership)?;
sidecar.require_owned_vm(&connection_id, &session_id, &vm_id)?;

validate_bindings_registration(&payload)?;

let registered_name = payload.name.clone();
let (original_permissions, original_bindings, original_command_guest_paths) = {
let vm = sidecar.vms.get(&vm_id).expect("owned VM should exist");
(
vm.configuration.permissions.clone(),
vm.bindings.clone(),
vm.command_guest_paths.clone(),
)
};
sidecar
.bridge
.set_vm_permissions(&vm_id, &allow_all_policy())?;
let registration_result = (|| -> Result<_, SidecarError> {
let vm = sidecar.vms.get_mut(&vm_id).expect("owned VM should exist");
ensure_collection_name_available(&vm.bindings, &registered_name)?;
ensure_command_aliases_available(&vm.bindings, &payload)?;
ensure_binding_registry_capacity(&vm.bindings, &payload)?;
vm.bindings.insert(registered_name.clone(), payload);
refresh_binding_registry(vm)?;
Ok::<_, SidecarError>(binding_command_names(vm).len() as u32)
})();
let command_count = match registration_result {
Ok(result) => {
sidecar
.bridge
.set_vm_permissions(&vm_id, &original_permissions)?;
result
}
Err(error) => {
let vm = sidecar.vms.get_mut(&vm_id).expect("owned VM should exist");
vm.bindings = original_bindings;
vm.command_guest_paths = original_command_guest_paths;
match sidecar.bridge.restore_vm_permissions_fail_closed(
&vm_id,
&original_permissions,
"binding collection registration rollback",
&error,
) {
Ok(()) => return Err(error),
Err(rollback_error) => {
vm.configuration.permissions = deny_all_policy();
return Err(rollback_error);
let input = sidecar.prepare_owned_vm_route(request);
let bridge = sidecar.bridge.clone();
async move {
let input = input?;
validate_bindings_registration(&payload)?;

let registered_name = payload.name.clone();
let (original_permissions, original_bindings, original_command_guest_paths) = input
.vm
.try_read("snapshot host callback registration", |vm| {
(
vm.configuration.permissions.clone(),
vm.bindings.clone(),
vm.command_guest_paths.clone(),
)
})?;
bridge.set_vm_permissions(&input.vm_id, &allow_all_policy())?;
let registration_result = input.vm.try_command("register host callbacks", |vm| {
ensure_collection_name_available(&vm.bindings, &registered_name)?;
ensure_command_aliases_available(&vm.bindings, &payload)?;
ensure_binding_registry_capacity(&vm.bindings, &payload)?;
vm.bindings.insert(registered_name.clone(), payload);
refresh_binding_registry(vm)?;
Ok(binding_command_names(vm).len() as u32)
});
let command_count = match registration_result {
Ok(result) => {
if let Err(restore_error) =
bridge.set_vm_permissions(&input.vm_id, &original_permissions)
{
if let Err(rollback_error) = rollback_host_callback_registration(
&input.vm,
&bridge,
&input.vm_id,
&original_permissions,
original_bindings,
original_command_guest_paths,
&restore_error,
) {
return Err(rollback_error);
}
return Err(restore_error);
}
result
}
}
};
Err(error) => {
match rollback_host_callback_registration(
&input.vm,
&bridge,
&input.vm_id,
&original_permissions,
original_bindings,
original_command_guest_paths,
&error,
) {
Ok(()) => return Err(error),
Err(rollback_error) => return Err(rollback_error),
}
}
};

Ok(DispatchResult {
response: sidecar.respond(
request,
ResponsePayload::HostCallbacksRegistered(HostCallbacksRegisteredResponse {
registration: registered_name,
command_count,
}),
),
events: Vec::new(),
})
Ok(DispatchResult {
response: shared_respond(
&input.request,
ResponsePayload::HostCallbacksRegistered(HostCallbacksRegisteredResponse {
registration: registered_name,
command_count,
}),
),
events: Vec::new(),
})
}
}

fn rollback_host_callback_registration<B>(
vm: &VmHandle,
bridge: &SharedBridge<B>,
vm_id: &str,
original_permissions: &agentos_vm_config::PermissionsPolicy,
original_bindings: BTreeMap<String, RegisterHostCallbacksRequest>,
original_command_guest_paths: BTreeMap<String, String>,
operation_error: &SidecarError,
) -> Result<(), SidecarError>
where
B: NativeSidecarBridge + Send + 'static,
BridgeError<B>: fmt::Debug + Send + Sync + 'static,
{
// These two rollback legs are independent. In particular, a failed VM
// registry refresh must never skip restoring (or fail-closing) the trusted
// bridge permission policy.
let state_rollback = vm.try_command("rollback host callback registration", |vm| {
vm.bindings = original_bindings;
vm.command_guest_paths = original_command_guest_paths;
refresh_binding_registry(vm)
});
let permission_rollback = bridge.restore_vm_permissions_fail_closed(
vm_id,
original_permissions,
"binding collection registration rollback",
operation_error,
);

match (state_rollback, permission_rollback) {
(Ok(()), Ok(())) => Ok(()),
(Err(state_error), Ok(())) => Err(SidecarError::InvalidState(format!(
"binding collection registration state rollback failed after {operation_error}: {state_error}; original permissions were restored"
))),
(Ok(()), Err(permission_error)) => {
vm.try_command("record fail-closed binding permissions", |vm| {
vm.configuration.permissions = deny_all_policy();
Ok(())
})?;
Err(permission_error)
}
(Err(state_error), Err(permission_error)) => {
let fail_closed_state = vm.try_command("record fail-closed binding permissions", |vm| {
vm.configuration.permissions = deny_all_policy();
Ok(())
});
Err(SidecarError::InvalidState(format!(
"binding collection registration rollback failed after {operation_error}: state rollback: {state_error}; permission rollback: {permission_error}; fail-closed VM state: {fail_closed_state:?}"
)))
}
}
}

fn refresh_binding_registry(vm: &mut VmState) -> Result<(), SidecarError> {
Expand Down
Loading