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
24 changes: 14 additions & 10 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -225,10 +225,12 @@ The daemon owns:
- object path: `/com/nutanix/io_thread_controller1`
- interface: `com.nutanix.io_thread_controller1`

Debug builds expose `SetThreadCount(vm, threads, sticky)` to change a tracked
instance's thread count. Setting `sticky=true` suppresses automatic scaling for
that instance until a later call clears it. This debug-only override is kept in
memory and is lost when the daemon restarts.
Debug builds expose `SetThreadCount(vm, threads, sticky)` and
`SetAllThreadCounts(threads, sticky)`. The latter attempts every VM in the
daemon inventory and reports partial failures without undoing successful
updates. Setting `sticky=true` suppresses automatic scaling for each successful
VM until a later debug call clears it. The override is kept in memory and is
lost when the daemon restarts.

`GetStats()` returns a JSON fleet snapshot. For example:

Expand Down Expand Up @@ -423,17 +425,19 @@ it never sees a backend directly.
| `Esc` | Clear focus (repaint every vm) |
| `+`, `=` | Ask the daemon to add one thread to the focused vm |
| `-`, `_` | Ask the daemon to remove one thread from the focused vm |
| `A` | Set every VM in the daemon inventory to an exact count |
| `s` | Toggle sticky mode for `+` / `-` (title bar shows `[STICKY]`) |

**`+` / `-` always route through the daemon.** Both
**`+` / `-` / `A` always route through the daemon.** Both
per-frame stats (via `GetSnapshot`) and manual actuations
(via `SetThreadCount(vm, threads, sticky)`) go over
(via `SetThreadCount` or `SetAllThreadCounts`) go over
`com.nutanix.io_thread_controller1` on
`/com/nutanix/io_thread_controller1`, so the daemon is the
single source of truth for `manual_scaling_sticky`. A wire
error on either verb logs at WARN on the `iothread-tui`
tracing target (and shows up in the daemon journal too) but
otherwise leaves the UI running so the operator can retry.
single source of truth for `manual_scaling_sticky`. `A` targets the complete
daemon inventory, including VMs hidden by `--select`. A wire error on either
verb logs at WARN on the `iothread-tui` tracing target (and shows up in the
daemon journal too) but otherwise leaves the UI running so the operator can
retry.
There is no direct-backend fallback path — a missing daemon
is a fatal condition the binary refuses to start with.

Expand Down
108 changes: 108 additions & 0 deletions src/controller.rs
Original file line number Diff line number Diff line change
Expand Up @@ -272,6 +272,14 @@ impl Controller {
.map_err(|error| error.to_string());
let _ = reply.send(result);
}
DbusRequest::SetAllThreadCounts {
threads,
sticky,
reply,
} => {
let result = self.handle_set_all_thread_counts(threads, sticky).await;
let _ = reply.send(result);
}
DbusRequest::GetSnapshot { reply } => {
let _ = reply.send(self.handle_get_snapshot().await);
}
Expand Down Expand Up @@ -409,6 +417,36 @@ impl Controller {
Ok(())
}

/// Apply a debug-only manual count to every VM known to the daemon.
///
/// The daemon inventory, rather than the consumer's filtered view, defines
/// the target set. Every VM is attempted in stable identifier order so
/// one rejection does not hide successful updates to other VMs.
async fn handle_set_all_thread_counts(&self, threads: u32, sticky: bool) -> (u32, String) {
let mut ids: Vec<_> = self.instances.keys().cloned().collect();
ids.sort_unstable();
let attempted = ids.len();
let mut succeeded = 0usize;
let mut failures = Vec::new();
for id in ids {
match self.handle_set_thread_count(&id, threads, sticky).await {
Ok(()) => succeeded += 1,
Err(error) => failures.push(format!("{id}: {error}")),
}
}

let failed = failures.len();
let mut summary = format!(
"set thread count to {threads}{} on {succeeded}/{attempted} VMs",
if sticky { " (sticky)" } else { "" }
);
if !failures.is_empty() {
summary.push_str("; failures: ");
summary.push_str(&failures.join("; "));
}
(u32::try_from(failed).unwrap_or(u32::MAX), summary)
}

/// Assemble the machine-readable fleet snapshot.
async fn handle_get_snapshot(&self) -> String {
let mut payload = SnapshotPayload {
Expand Down Expand Up @@ -1075,6 +1113,76 @@ mod tests {
assert_eq!(target.load(Ordering::Relaxed), 2);
}

struct BulkClient {
fail: bool,
target: Arc<AtomicU32>,
}

#[async_trait]
impl InstanceClient for BulkClient {
async fn set_thread_count(&self, n: u32) -> Result<(), BackendClientError> {
if self.fail {
return Err(BackendClientError::Protocol("injected failure".to_string()));
}
self.target.store(n, Ordering::Relaxed);
Ok(())
}

async fn get_thread_pool_snapshot(&self) -> Result<ThreadPoolSnapshot, BackendClientError> {
Ok(ThreadPoolSnapshot {
thread_count: self.target.load(Ordering::Relaxed),
vcpu_count: 16,
perf: None,
per_thread_util: Some(0.0),
})
}

async fn close(&self) {}
}

/// Test that `set_all` applies a target to every VM and reports how
/// many failed when some clients error.
#[tokio::test]
async fn set_all_attempts_every_vm_and_reports_partial_failure() {
let state_dir = tempfile::tempdir().unwrap();
let cfg = Config {
vm_state_path: Path::new(&state_dir.path().join("ownership.json")),
..Default::default()
};
let engine = Box::new(ThresholdEngine::new(ThresholdConfig::default()));
let mut controller = Controller::new(cfg, engine).unwrap();
let good_target = Arc::new(AtomicU32::new(1));
let good = BulkClient {
fail: false,
target: Arc::clone(&good_target),
};
let bad_target = Arc::new(AtomicU32::new(1));
let bad = BulkClient {
fail: true,
target: Arc::clone(&bad_target),
};
let good_instance = Arc::new(Instance::new("visible".to_string(), Path::new(""), 0, good));
let bad_instance = Arc::new(Instance::new("hidden".to_string(), Path::new(""), 0, bad));
good_instance.status.write().await.vcpu_count = 16;
bad_instance.status.write().await.vcpu_count = 16;
controller
.instances
.insert("visible".to_string(), good_instance.clone());
controller
.instances
.insert("hidden".to_string(), bad_instance.clone());

let (failed, summary) = controller.handle_set_all_thread_counts(4, true).await;

assert_eq!(failed, 1);
assert_eq!(good_target.load(Ordering::Relaxed), 4);
assert_eq!(bad_target.load(Ordering::Relaxed), 1);
assert!(good_instance.status.read().await.manual_scaling_sticky);
assert!(!bad_instance.status.read().await.manual_scaling_sticky);
assert!(summary.contains("1/2"));
assert!(summary.contains("hidden: injected failure"));
}

struct SnapshotClient {
threads: u32,
closed: Arc<AtomicUsize>,
Expand Down
50 changes: 44 additions & 6 deletions src/dbus.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,8 @@
//! target VM and dispatches through its [`crate::instance::InstanceClient`].
//!
//! Verbs currently exposed:
//! * debug-only `SetThreadCount(vm, threads, sticky)`;
//! * debug-only `SetThreadCount(vm, threads, sticky)` and
//! `SetAllThreadCounts(threads, sticky)`;
//! * read-only `GetStats()`;
//! * read-only `GetVersion()`;
//! * named-IOThread and virtqueue-mapping operations for clients that support
Expand Down Expand Up @@ -49,14 +50,28 @@ pub enum DbusRequest {
/// One-shot channel used to complete the D-Bus method call.
reply: oneshot::Sender<Result<(), String>>,
},
/// Structured snapshot of every tracked instance. The reply is a JSON
/// string so the wire signature stays a bare `s` and the payload shape can
/// evolve without D-Bus IDL churn. Field contract is documented on
/// Operator-driven fleet-wide thread-count update.
SetAllThreadCounts {
/// Target thread count for every currently tracked VM.
threads: u32,
/// Apply or clear the manual sticky override on every successful VM.
sticky: bool,
/// Reply is `(number of failures, human-readable summary)`.
reply: oneshot::Sender<(u32, String)>,
},
/// JSON snapshot of all tracked VMs.
GetStats {
/// Reply channel carrying the serialized snapshot.
reply: oneshot::Sender<String>,
},
/// Structured snapshot of every tracked instance -- the
/// machine-readable sibling of `GetStats`. Consumed by
/// `iothread-tui` (and any other operator UI) to render
/// live plots + dashboards without having to peek at
/// backend sockets directly. The reply is a JSON string
/// so the wire signature stays a bare `s` and the payload
/// shape can evolve without D-Bus IDL churn. Field
/// contract is documented on
/// [`crate::controller::SnapshotPayload`].
GetSnapshot {
/// Reply channel; JSON-encoded
Expand Down Expand Up @@ -148,7 +163,11 @@ impl Service {
vm: String,
threads: u32,
sticky: bool,
) -> zbus::fdo::Result<()> {
) -> zbus::fdo::Result<(u32, String)> {
let summary = format!(
"set thread count on VM {vm} to {threads}{}",
if sticky { " (sticky)" } else { "" }
);
let (tx, rx) = oneshot::channel();
self.tx
.send(DbusRequest::SetThreadCount {
Expand All @@ -160,11 +179,30 @@ impl Service {
.await
.map_err(|_| zbus::fdo::Error::Failed("controller channel closed".into()))?;
match await_with_timeout(rx).await? {
Ok(()) => Ok(()),
Ok(()) => Ok((0, summary)),
Err(e) => Err(zbus::fdo::Error::Failed(e)),
}
}

/// Apply one exact worker count to every VM in the daemon inventory.
#[cfg(debug_assertions)]
async fn set_all_thread_counts(
&self,
threads: u32,
sticky: bool,
) -> zbus::fdo::Result<(u32, String)> {
let (tx, rx) = oneshot::channel();
self.tx
.send(DbusRequest::SetAllThreadCounts {
threads,
sticky,
reply: tx,
})
.await
.map_err(|_| zbus::fdo::Error::Failed("controller channel closed".into()))?;
await_with_timeout(rx).await
}

/// Return the controller's machine-readable fleet snapshot.
async fn get_stats(&self) -> zbus::fdo::Result<String> {
let (tx, rx) = oneshot::channel();
Expand Down
Loading