fix(config): coalesce direct and fleet updates per component

Signed-off-by: Piotr Mlocek <pmlocek@nvidia.com>
This commit is contained in:
Piotr Mlocek
2026-10-01 15:58:20 -07:00
parent c01364531e
commit 0e3d78b295
+25 -10
View File
@@ -171,7 +171,6 @@ impl ConfigComponentKind {
struct DeliveryKey {
sandbox_id: String,
component: ConfigComponentKind,
route_remote: bool,
}
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
@@ -417,21 +416,19 @@ pub fn publish_sandbox_components(
if state.config.config_delivery_mode != ConfigDeliveryMode::Push {
return;
}
enqueue_sandbox(state, sandbox_id, components, true);
enqueue_sandbox(state, sandbox_id, components);
}
fn enqueue_sandbox(
state: &Arc<ServerState>,
sandbox_id: &str,
components: ConfigComponents,
route_remote: bool,
) -> bool {
let mut accepted = true;
for component in components.selected() {
let key = DeliveryKey {
sandbox_id: sandbox_id.to_string(),
component,
route_remote,
};
match state.config_delivery_queue.enqueue(key.clone()) {
DeliveryEnqueue::StartWorker(permit) => {
@@ -469,7 +466,6 @@ async fn enqueue_sandbox_from_fanout(
let key = DeliveryKey {
sandbox_id: sandbox_id.to_string(),
component,
route_remote: false,
};
match state
.config_delivery_queue
@@ -483,10 +479,10 @@ async fn enqueue_sandbox_from_fanout(
}
async fn publish_sandbox_component_now(state: &Arc<ServerState>, key: &DeliveryKey) {
if key.route_remote && route_remote_sandbox_component(state, key).await {
if route_remote_sandbox_component(state, key).await {
return;
}
if !key.route_remote && !is_current_local_owner(state, &key.sandbox_id).await {
if !is_current_local_owner(state, &key.sandbox_id).await {
return;
}
if !state
@@ -520,6 +516,9 @@ async fn publish_sandbox_component_now(state: &Arc<ServerState>, key: &DeliveryK
match state.config_delivery_queue.run_bounded_build(build).await {
Ok(Ok(None)) => {}
Ok(Ok(Some(message))) => {
if !is_current_local_owner(state, &key.sandbox_id).await {
return;
}
let disposition = state
.supervisor_config_router()
.deliver(&key.sandbox_id, message)
@@ -815,8 +814,7 @@ pub async fn handle_peer_config_update_hint(
.is_current_session(&target.sandbox_id, &target.session_id)
});
if !response.stale_owner {
response.queue_full =
!enqueue_sandbox(state, &target.sandbox_id, components, false);
response.queue_full = !enqueue_sandbox(state, &target.sandbox_id, components);
}
}
Some(peer_config_update_hint_request::Scope::Workspace(workspace)) => {
@@ -994,7 +992,6 @@ mod tests {
DeliveryKey {
sandbox_id: sandbox_id.to_string(),
component,
route_remote: true,
}
}
@@ -1020,6 +1017,24 @@ mod tests {
drop(permit);
}
#[tokio::test]
async fn direct_and_fanout_updates_share_one_component_worker() {
let queue = ConfigDeliveryQueue::default();
let key = key("sb-1", ConfigComponentKind::SandboxConfig);
let DeliveryEnqueue::StartWorker(permit) = queue.enqueue(key.clone()) else {
panic!("direct update must start a worker");
};
queue.take(&key);
assert!(matches!(
queue.enqueue_from_fanout(key.clone()).await,
DeliveryEnqueue::Coalesced
));
assert!(queue.finish_pass(&key));
queue.take(&key);
assert!(!queue.finish_pass(&key));
drop(permit);
}
#[test]
fn build_bound_is_sized_from_the_database_pool() {
assert_eq!(