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
15 changes: 15 additions & 0 deletions Cargo.lock

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

35 changes: 35 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -157,6 +157,41 @@ dwd udp <IP:PORT> --no-ui --generator profile.yaml --api-addr 127.0.0.1:8080
curl http://127.0.0.1:8080/api/v1/metrics
```

Alternatively (or additionally), expose the gRPC control API with `--grpc-addr` and drive the run remotely — see below.

## Remote control (gRPC API)

The interactive TUI talks to the load-generating core through a gRPC API (`dwd.v1.Dwd`, see `dwd-proto/dwdpb/dwd.proto`). Passing the global `--grpc-addr` flag exposes that same API on a network address, giving remote clients full control over the run — everything the TUI can do:

- **`Control`** — set a fixed RPS, suspend/resume the load profile, or stop the run entirely (the engine drains and the process shuts down gracefully, exactly like a TUI exit or `SIGTERM`);
- **`StreamStats`** — a live stream of cumulative statistic snapshots (requests, bytes, response codes, latency histogram buckets);
- **`Describe`** — the engine kind and which statistic groups it provides.

The flag works with and without the TUI; combined with `--no-ui` it makes a fully API-driven headless run:

```bash
dwd udp <IP:PORT> --no-ui --grpc-addr 127.0.0.1:8081
```

Server reflection is enabled, so `grpcurl` works without any proto files:

```bash
# Set the load to 5000 RPS:
grpcurl -plaintext -d '{"set":{"rps":"5000"}}' 127.0.0.1:8081 dwd.v1.Dwd/Control

# Suspend / resume the load profile:
grpcurl -plaintext -d '{"suspend":{}}' 127.0.0.1:8081 dwd.v1.Dwd/Control
grpcurl -plaintext -d '{"resume":{}}' 127.0.0.1:8081 dwd.v1.Dwd/Control

# Watch live statistics:
grpcurl -plaintext -d '{}' 127.0.0.1:8081 dwd.v1.Dwd/StreamStats

# Stop the run gracefully:
grpcurl -plaintext -d '{"stop":{}}' 127.0.0.1:8081 dwd.v1.Dwd/Control
```

Note: `set` overrides the profile with a fixed value (implicitly suspending it); `resume` returns control to the profile. The API is unauthenticated plaintext HTTP/2 — bind it to localhost or a trusted network only.

## Load profiles
A load profile is a declarative description of how the load should be generated and for how long.

Expand Down
8 changes: 6 additions & 2 deletions dwd-core/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -51,10 +51,14 @@ rand = "0.9"
axum = "0.8"
prometheus-client = "0.24"

# gRPC service (server half of the in-process API seam).
# gRPC service (server half of the API seam; in-process for the TUI and
# optionally exposed on a network address for remote control).
dwd-proto = { path = "../dwd-proto" }
tonic = "0.14"
tokio-stream = "0.1"
# Reflection lets `grpcurl` & co discover the API without local proto files.
tonic-reflection = "0.14"
# `net` provides `TcpListenerStream` for the network gRPC endpoint.
tokio-stream = { version = "0.1", features = ["net"] }

# Optional DPDK FFI backend (Linux only, needs DPDK 24.11 installed).
dpdk = { version = "0.1", path = "../dpdk", package = "dwd-dpdk", optional = true }
Expand Down
47 changes: 43 additions & 4 deletions dwd-core/src/grpc/mod.rs
Original file line number Diff line number Diff line change
@@ -1,8 +1,10 @@
//! Server half of the in-process gRPC API seam.
//! Server half of the gRPC API seam.
//!
//! [`server`] implements the [`Dwd`](dwd_proto::dwd_server::Dwd) service over the
//! running core; [`snapshot`] maps the core statistics traits onto the wire
//! snapshot. The client half (TUI) lives in the `dwd` binary.
//! snapshot. The service is hosted in two ways: [`serve`] over an in-memory pipe
//! for the built-in TUI, and [`serve_tcp`] on a network address for remote
//! clients. The client half (TUI) lives in the `dwd` binary.

pub mod server;
pub mod snapshot;
Expand All @@ -12,8 +14,10 @@ pub use self::{
snapshot::{EngineDescriptor, SnapshotSource},
};

use tokio::{io::DuplexStream, task::JoinHandle};
use tokio_stream::StreamExt as _;
use core::net::SocketAddr;

use tokio::{io::DuplexStream, net::TcpListener, task::JoinHandle};
use tokio_stream::{wrappers::TcpListenerStream, StreamExt as _};
use tonic::transport::Server;

use dwd_proto::dwd_server::DwdServer;
Expand All @@ -35,3 +39,38 @@ pub fn serve(service: DwdService, server_io: DuplexStream) -> JoinHandle<Result<
.await
})
}

/// Hosts the [`DwdService`] on a TCP address for remote control.
///
/// Binds eagerly so misconfiguration (port in use, bad address) fails startup
/// instead of surfacing later in a background task, and returns the bound
/// address (useful with port 0). Alongside the [`Dwd`](dwd_proto::dwd_server::Dwd)
/// service, gRPC reflection (v1 + v1alpha) is served so tools like `grpcurl`
/// can discover the API without proto files. The returned [`JoinHandle`] owns
/// the server task and should be aborted on shutdown.
pub async fn serve_tcp(
service: DwdService,
addr: SocketAddr,
) -> Result<(SocketAddr, JoinHandle<Result<(), tonic::transport::Error>>), anyhow::Error> {
let listener = TcpListener::bind(addr).await?;
let addr = listener.local_addr()?;
log::info!("gRPC API server listening on {addr}");

let reflection_v1 = tonic_reflection::server::Builder::configure()
.register_encoded_file_descriptor_set(dwd_proto::FILE_DESCRIPTOR_SET)
.build_v1()?;
let reflection_v1alpha = tonic_reflection::server::Builder::configure()
.register_encoded_file_descriptor_set(dwd_proto::FILE_DESCRIPTOR_SET)
.build_v1alpha()?;

let handle = tokio::spawn(async move {
Server::builder()
.add_service(reflection_v1)
.add_service(reflection_v1alpha)
.add_service(DwdServer::new(service))
.serve_with_incoming(TcpListenerStream::new(listener))
.await
});

Ok((addr, handle))
}
82 changes: 82 additions & 0 deletions dwd-core/src/grpc/server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,10 @@ const STATS_INTERVAL: Duration = Duration::from_millis(25);
const STATS_CHANNEL_CAP: usize = 4;

/// The [`Dwd`] service implementation over the in-process core.
///
/// Cloning is cheap (channels and `Arc`s), so the same service can back both
/// the in-memory TUI seam and the network endpoint at once.
#[derive(Clone)]
pub struct DwdService {
/// Forwards control RPCs to `run_generator` (mirrors the old UI mpsc).
control_tx: Sender<GeneratorEvent>,
Expand Down Expand Up @@ -65,6 +69,15 @@ impl Dwd for DwdService {
Kind::Suspend(_) => GeneratorEvent::Suspend,
Kind::Resume(_) => GeneratorEvent::Resume,
Kind::Set(set) => GeneratorEvent::Set(set.rps),
Kind::Stop(_) => {
// Stop bypasses the generator channel: flipping the shared
// run flag drains the engine and the generator loop exactly
// as a TUI exit or SIGTERM would.
log::info!("stop requested via the API");
self.is_running.store(false, Ordering::SeqCst);

return Ok(Response::new(ControlResponse {}));
}
};

// Non-blocking, coalescing send — mirrors the TUI's `try_send`: if the
Expand Down Expand Up @@ -103,3 +116,72 @@ impl Dwd for DwdService {
Ok(Response::new(ReceiverStream::new(rx)))
}
}

#[cfg(test)]
mod tests {
use dwd_proto::{SetControl, StatGroup, StopControl, SuspendControl};

use super::*;
use crate::grpc::snapshot::EngineDescriptor;

struct EmptySource;

impl SnapshotSource for EmptySource {
fn snapshot(&self) -> StatsSnapshot {
StatsSnapshot::default()
}
}

fn service() -> (DwdService, mpsc::Receiver<GeneratorEvent>, Arc<AtomicBool>) {
let (control_tx, control_rx) = mpsc::channel(1);
let descriptor = EngineDescriptor {
engine_kind: "udp",
groups: vec![StatGroup::Common],
}
.into_response();
let is_running = Arc::new(AtomicBool::new(true));
let service = DwdService::new(control_tx, Arc::new(EmptySource), descriptor, is_running.clone());

(service, control_rx, is_running)
}

#[tokio::test]
async fn control_forwards_generator_events() {
let (service, mut control_rx, is_running) = service();

service
.control(Request::new(ControlRequest {
kind: Some(Kind::Set(SetControl { rps: 1234 })),
}))
.await
.expect("set control");
assert!(matches!(control_rx.try_recv(), Ok(GeneratorEvent::Set(1234))));

service
.control(Request::new(ControlRequest {
kind: Some(Kind::Suspend(SuspendControl {})),
}))
.await
.expect("suspend control");
assert!(matches!(control_rx.try_recv(), Ok(GeneratorEvent::Suspend)));

// Generator control never touches the run flag.
assert!(is_running.load(Ordering::SeqCst));
}

#[tokio::test]
async fn control_stop_flips_run_flag() {
let (service, mut control_rx, is_running) = service();

service
.control(Request::new(ControlRequest {
kind: Some(Kind::Stop(StopControl {})),
}))
.await
.expect("stop control");

// Stop is not a generator event: it acts on the run flag directly.
assert!(!is_running.load(Ordering::SeqCst));
assert!(control_rx.try_recv().is_err());
}
}
5 changes: 5 additions & 0 deletions dwd-proto/build.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,9 +3,14 @@ fn main() -> Result<(), Box<dyn std::error::Error>> {
// protobuf-compiler (keeps CI and the Docker builds self-contained).
std::env::set_var("PROTOC", protoc_bin_vendored::protoc_bin_path()?);

let out_dir = std::path::PathBuf::from(std::env::var("OUT_DIR")?);

tonic_prost_build::configure()
.build_client(true)
.build_server(true)
// Emitted for the gRPC reflection service, which lets tools like
// `grpcurl` discover the API without a local copy of the proto file.
.file_descriptor_set_path(out_dir.join("dwd_descriptor.bin"))
.compile_protos(&["dwdpb/dwd.proto"], &["dwdpb"])?;
Ok(())
}
11 changes: 8 additions & 3 deletions dwd-proto/dwdpb/dwd.proto
Original file line number Diff line number Diff line change
Expand Up @@ -5,9 +5,10 @@ package dwd.v1;
// The control + observability API seam between the DWD CLI/TUI (client) and the
// load-generating core (server).
//
// The service is consumed in-process over an in-memory transport (no sockets),
// but the contract is a regular gRPC service so the CLI talks to the core
// exclusively through this API.
// The service is consumed in-process over an in-memory transport (no sockets)
// by the built-in TUI, and can additionally be exposed on a network address
// (`--grpc-addr`) for remote control — the contract is a regular gRPC service
// either way, so every client talks to the core exclusively through this API.
service Dwd {
// Describes the running engine so the client knows which statistics groups
// (i.e. which TUI widgets) to render. Replaces the previous compile-time
Expand Down Expand Up @@ -48,6 +49,7 @@ message ControlRequest {
SuspendControl suspend = 1;
ResumeControl resume = 2;
SetControl set = 3;
StopControl stop = 4;
}
}

Expand All @@ -59,6 +61,9 @@ message ResumeControl {}
message SetControl {
uint64 rps = 1;
}
// Stops the run entirely: the engine drains and the process shuts down
// gracefully, exactly as if the TUI had exited or SIGTERM was received.
message StopControl {}

message ControlResponse {}

Expand Down
4 changes: 4 additions & 0 deletions dwd-proto/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,3 +10,7 @@ pub mod pb {
}

pub use pb::*;

/// Encoded `FileDescriptorSet` of the API, consumed by the gRPC reflection
/// service so tools like `grpcurl` can discover the API without proto files.
pub const FILE_DESCRIPTOR_SET: &[u8] = include_bytes!(concat!(env!("OUT_DIR"), "/dwd_descriptor.bin"));
11 changes: 10 additions & 1 deletion dwd/src/cfg.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,8 @@ pub struct Config {
pub generator_fn: BoxedGeneratorNew,
/// Address to expose the Prometheus API on.
pub api_addr: Option<SocketAddr>,
/// Address to expose the gRPC control API on.
pub grpc_addr: Option<SocketAddr>,
/// Run headless, without the interactive TUI.
pub no_ui: bool,
}
Expand All @@ -32,6 +34,7 @@ impl TryFrom<Cmd> for Config {
fn try_from(v: Cmd) -> Result<Self, Self::Error> {
let mode = v.mode.into_config()?;
let api_addr = v.api_addr;
let grpc_addr = v.grpc_addr;
let no_ui = v.no_ui;
let generator_fn = {
let path = v.generator.clone();
Expand All @@ -49,6 +52,12 @@ impl TryFrom<Cmd> for Config {
}))
};

Ok(Self { mode, generator_fn, api_addr, no_ui })
Ok(Self {
mode,
generator_fn,
api_addr,
grpc_addr,
no_ui,
})
}
}
9 changes: 9 additions & 0 deletions dwd/src/cmd.rs
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,15 @@ pub struct Cmd {
/// to observe the generator state and metrics.
#[clap(long, global = true, value_name = "IP:PORT")]
pub api_addr: Option<SocketAddr>,
/// Address to expose the gRPC control API on.
///
/// When specified, serves the same `dwd.v1.Dwd` gRPC service the built-in
/// TUI uses, so remote clients get full control over the run: set the RPS,
/// suspend/resume the generator, stop the run and stream live statistics.
/// Server reflection is enabled, so `grpcurl` works without proto files.
/// Combine with `--no-ui` to drive a headless run entirely over the API.
#[clap(long, global = true, value_name = "IP:PORT")]
pub grpc_addr: Option<SocketAddr>,
/// Run headless, without the interactive TUI.
///
/// The load runs until the profile is exhausted or the process receives
Expand Down
Loading
Loading