From b0965e80d619cb95fb12620abb69714ad23c0667 Mon Sep 17 00:00:00 2001 From: John Basrai <95644651+JohnBasrai@users.noreply.github.com> Date: Sat, 14 Feb 2026 20:05:06 -0800 Subject: [PATCH] chore: bump mom-rpc to 0.8.2 Upgrade to mom-rpc v0.8.2. Includes fix for serialized MQTT SUBSCRIBE handling, preventing concurrent subscribe race during full-duplex transport initialization. --- .github/workflows/ci.yml | 13 ++- CHANGELOG.md | 23 ++++- CONTRIBUTING.md | 4 +- Cargo.toml | 6 +- README.md | 201 +++++++++++++++++++++++++++++--------- scripts/demo.sh | 6 +- scripts/service-start.sh | 6 -- scripts/service-stop.sh | 7 -- scripts/start-services.sh | 11 +++ scripts/stop-services.sh | 10 ++ src/agent/run.rs | 71 ++++++++------ src/bin/device_sim.rs | 157 ++++++++++++++++------------- src/lib.rs | 5 +- src/main.rs | 3 + src/messaging/mod.rs | 5 - src/messaging/mom_rpc.rs | 77 --------------- 16 files changed, 350 insertions(+), 255 deletions(-) delete mode 100755 scripts/service-start.sh delete mode 100755 scripts/service-stop.sh create mode 100755 scripts/start-services.sh create mode 100755 scripts/stop-services.sh delete mode 100644 src/messaging/mom_rpc.rs diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 100b764..bc65d05 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -101,7 +101,7 @@ jobs: steps: - name: Checkout repository uses: actions/checkout@v4 - + - name: Install QEMU run: | sudo apt-get update @@ -109,12 +109,19 @@ jobs: qemu-user \ libc6-arm64-cross \ libstdc++6-arm64-cross - + - name: Download AArch64 binary uses: actions/download-artifact@v4 with: name: rust-edge-agent-aarch64 path: target/aarch64-unknown-linux-gnu/release - + + - name: Start infrastructure services + run: scripts/start-services.sh + - name: QEMU AArch64 smoke test run: scripts/ci-qemu-smoke.sh + + - name: Stop infrastructure services + if: always() + run: scripts/stop-services.sh diff --git a/CHANGELOG.md b/CHANGELOG.md index e4414c8..9c3bd9f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,25 @@ ## [Unreleased] + +## [0.3.2] - 2026-02-15 + +### Fixed +- Add missing services startup to CI qemu-smoke job +- Fix demo.sh signal handling for clean process cleanup +- Fix MQTT_PORT interpolation in start-services.sh + +### Changed +- Update mom-rpc to 0.7.3 +- Replace env_logger with tracing-subscriber +- Disable ANSI codes in logs for better script output +- Redirect device_sim output to /dev/null in demo + +### Documentation +- Remove manual cleanup steps from README +- Add smoke test output example +- Rename service scripts throughout documentation + + ## [0.3.1] - 2026-02-07 ### Changed @@ -8,6 +28,7 @@ - Replace `std::process::exit()` with idiomatic `Result`/`bail!()` error handling - Add package metadata (keywords, categories, description) + ## [0.3.0] - 2026-02-06 ### Changed @@ -27,6 +48,7 @@ - `service-start.sh` now starts both NATS and MQTT brokers - Docker compose includes mosquitto container + ### Migration from v0.2.x - MQTT broker required (port 1883): `./scripts/service-start.sh` - Set `MQTT_BROKER_URL` environment variable if not using default @@ -36,7 +58,6 @@ ### Known Issues - demo.sh cleanup trap needs manual `pkill -9 rust-edge-agent device_sim` after Ctrl+C - ## [0.2.0] – 2026-01-31 ### Added diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index 52dbd61..85aa69b 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -14,10 +14,10 @@ Before running the agent locally: ```bash # Start NATS broker -./scripts/service-start.sh +./scripts/start-services.sh # When done -./scripts/service-stop.sh +./scripts/stop-services.sh ``` ### Running the Demo diff --git a/Cargo.toml b/Cargo.toml index 58e9723..9e4859f 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "rust-edge-agent" -version = "0.3.1" +version = "0.3.2" edition = "2021" license = "MIT OR Apache-2.0" keywords = ["mqtt", "nats", "edge-gateway", "iot", "rpc", "arm64"] @@ -11,14 +11,14 @@ description = "Edge gateway for coordinating heterogeneous devices via dual-prot anyhow = "1.0" async-nats = "0.35" clap = { version = "4", features = ["derive", "env"] } -env_logger = "0.11" futures = "0.3" -mom-rpc = { version = "0.3", features = ["transport_rumqttc"] } +mom-rpc = { version = "0.8", features = ["transport_rumqttc"] } rand = "0.8" serde = { version = "1", features = ["derive"] } serde_json = "1" tokio = { version = "1", features = ["macros", "rt-multi-thread", "time"] } tracing = "0.1" +tracing-subscriber = { version = "0.3", features = ["env-filter"] } [[bin]] name = "rust-edge-agent" diff --git a/README.md b/README.md index 2e0edc8..d79b6bb 100644 --- a/README.md +++ b/README.md @@ -6,7 +6,9 @@ This repository explores design patterns for an edge gateway responsible for coo Current state: Working dual-protocol edge agent with NATS backend and MQTT device communication. Demonstrates sensor/actuator/hybrid device patterns, dynamic device registration via RPC, command routing, and telemetry aggregation. Protocol bridge validated under AArch64 cross-compilation and QEMU emulation. -## Scope and Intent +--- + +## 1.0 Scope and Intent This project is intentionally scoped to explore edge gateway coordination patterns rather than end-device protocols. @@ -14,7 +16,9 @@ The gateway coordinates heterogeneous devices through a brokered control plane, Cross-compilation to `AArch64` and execution under `QEMU` are used to validate that the system behaves correctly as a long-running service across architectures. -## What this is +--- + +## 2.0 What this is - A **Linux-based edge agent** intended for embedded or gateway-class systems - Built in **Rust**, targeting **`AArch64` (`ARM64`)** via cross-compilation @@ -31,7 +35,9 @@ Cross-compilation to `AArch64` and execution under `QEMU` are used to validate t - protocol bridging (NATS ↔ MQTT) - build reproducibility -## What this is not +--- + +## 3.0 What this is not - Not bare-metal firmware - Not a BSP, bootloader, or kernel project @@ -42,7 +48,9 @@ Cross-compilation to `AArch64` and execution under `QEMU` are used to validate t This repository intentionally avoids expanding into those areas in order to keep the scope constrained. -## What This Demonstrates +--- + +## 4.0 What This Demonstrates This project is intentionally scoped to demonstrate edge gateway patterns that commonly appear in real deployments: @@ -53,7 +61,9 @@ This project is intentionally scoped to demonstrate edge gateway patterns that c - Designing long-running edge services that can be cross-compiled and validated under emulation (`AArch64` + `QEMU`) -## Architecture Overview +--- + +## 5.0 Architecture Overview This project demonstrates a **dual-protocol edge gateway** architecture: @@ -94,7 +104,9 @@ Device → [MQTT RPC] → Agent (register-device method) See [docs/architecture.md](docs/architecture.md) for complete details including message formats, RPC methods, and sequence diagrams. -## Important Files +--- + +## 6.0 Important Files * [docs/architecture.md](docs/architecture.md) High-level diagram + text: @@ -105,7 +117,7 @@ See [docs/architecture.md](docs/architecture.md) for complete details including --- -## Quick Start +## 7.0 Quick Start ### Prerequisites @@ -145,70 +157,131 @@ sudo systemctl enable docker ### Running the Demo **1. Start infrastructure services** + ```bash -./scripts/service-start.sh +./scripts/stop-services.sh # stop them if already running (harmless if not running) +./scripts/start-services.sh ``` -This starts NATS (port 4222) and MQTT broker (port 1883) in Docker containers. + +This launches NATS on port 4222 and an MQTT broker on port 1883 inside Docker containers. **2. Build the project** + ```bash ./scripts/ci-build-native.sh ``` -**3. Clean up any stale processes from previous runs** -```bash -pkill -9 rust-edge-agent device_sim 2>/dev/null || true -``` +**3. Run the demo** (starts edge agent + 3 device simulators) -**4. Run the demo** (starts edge agent + 3 device simulators) ```bash ./scripts/demo.sh ``` -The demo starts: -- Edge agent (NATS ↔ MQTT bridge) listening for device registrations -- 3 device simulators (actuator, hybrid, sensor modes) -- Devices register via MQTT RPC -- Agent polls sensors every 5 seconds and forwards telemetry to NATS backend +Let it run for about 30 seconds, then press Control-C to stop the demo. -**Expected output:** -``` -agent:: Edge agent running -agent:: Device registered: device-001 (mode: Actuator, type: Valve) -agent:: Device registered: device-002 (mode: Hybrid, type: Propulsion) -agent:: Device registered: device-003 (mode: Sensor, type: Temperature) -agent:: Telemetry: device-003 (Temperature) = 20.94 -agent:: Telemetry: device-002 (Propulsion) = 19.14 -``` +**4. Stop the demo** -**Monitor telemetry in a separate terminal:** ```bash -nats sub 'backend.telemetry' +Control-C # Stops demo ``` -**Send command to actuator/hybrid device:** -```bash -# Command to actuator -nats req 'backend.command.device-001' '{"target_value": 75.0}' +**5. Stop infrastructure services** -# Command to hybrid device -nats req 'backend.command.device-002' '{"target_value": 50.0}' +```bash +./scripts/stop-services.sh ``` +The demo starts the edge agent (acting as a NATS ↔ MQTT bridge), listening for device registrations. Three device simulators will launch (actuator, hybrid, sensor), register via RPC, and the agent will poll sensors every 5 seconds, forwarding telemetry to the NATS backend. -**5. Stop the demo** +Expected output includes device registrations and telemetry updates, and you can monitor or send commands via NATS as shown. + +
+Expected test output -Press `Ctrl+C` in the terminal running demo.sh, then manually clean up processes: -```bash -pkill -9 rust-edge-agent device_sim ``` + $ ./scripts/stop-services.sh +Stopping NATS broker... +Stopping MQTT broker... + $ ./scripts/start-services.sh +Starting NATS broker on localhost:4222 +de1472c8c7e4f1088f204a1831ce2d034d587bad7ed527d2f087750aa371783b +Starting mosquitto broker on localhost:MQTT_PORT +f92d3ff7f5385c889fca944f0333335f2eecc97a8bfaa81c428d2a7312ed0cce + $ ./scripts/ci-build-native.sh + Finished `release` profile [optimized] target(s) in 0.09s + $ ./scripts/demo.sh +./scripts/demo.sh: NUM_DEVICES : 3 +./scripts/demo.sh: DEVICE_INTERVAL : 5 +./scripts/demo.sh: MODE : interactive +./scripts/demo.sh: Checking for stale processes... +./scripts/demo.sh: Started agent (PID: 726495) +agent:: Starting edge agent... +agent:: NATS URL: nats://localhost:4222 +agent:: MQTT Broker: mqtt://localhost:1883 +agent:: Device timeout: 30s +agent:: Poll interval: 5s +Connected to NATS at nats://localhost:4222 +agent:: Connecting to MQTT broker... +2026-02-15T03:26:35.304702Z INFO async_nats: event: connected +Connected to MQTT broker at mqtt://localhost:1883 (client_id: agent-transport) +agent:: MQTT transport ready +agent:: Creating RPC server... +agent:: Creating RPC client... +2026-02-15T03:26:35.305268Z INFO mom_rpc::transport::rumqttc::transport: rumqttc: connected to broker +2026-02-15T03:26:35.305372Z INFO mom_rpc::transport::rumqttc::transport: rumqttc: successfully subscribed to topic responses/agent-client +agent:: RPC client ready +agent:: RPC server listening for device registrations +agent:: Edge agent running +2026-02-15T03:26:35.305489Z INFO mom_rpc::transport::rumqttc::transport: rumqttc: successfully subscribed to topic requests/agent +./scripts/demo.sh: Starting device: 1 mode:actuator TYPE:valve +./scripts/demo.sh: Starting device: 2 mode:hybrid TYPE:propulsion +./scripts/demo.sh: Starting device: 3 mode:sensor TYPE:temp +Edge agent and 3 devices running. -**Note:** The demo cleanup trap has known issues with signal handling. Manual process cleanup is required after stopping the demo. This will be addressed in a future update. +Monitor telemetry: + nats sub 'backend.telemetry' -**6. Stop infrastructure services** -```bash -./scripts/service-stop.sh +Send command to actuator: + nats req 'backend.command.device-002' '{"target_value": 75.0}' + +Press Ctrl+C to stop. +agent:: Device registered: device-001 (mode: Actuator, type: Valve) +agent:: Device registered: device-002 (mode: Hybrid, type: Propulsion) +agent:: Device registered: device-003 (mode: Sensor, type: Temperature) +agent:: Telemetry: device-002 (Propulsion) = 20.26 +agent:: Telemetry: device-003 (Temperature) = 20.21 +agent:: Telemetry: device-002 (Propulsion) = 21.05 +agent:: Telemetry: device-003 (Temperature) = 21.20 +agent:: Telemetry: device-002 (Propulsion) = 21.75 +agent:: Telemetry: device-003 (Temperature) = 21.65 +agent:: Telemetry: device-002 (Propulsion) = 21.99 +agent:: Telemetry: device-003 (Temperature) = 21.75 +agent:: Telemetry: device-002 (Propulsion) = 21.85 +agent:: Telemetry: device-003 (Temperature) = 21.17 +agent:: Telemetry: device-002 (Propulsion) = 21.51 +agent:: Telemetry: device-003 (Temperature) = 21.82 +^C +./scripts/demo.sh: Cleaning up... +./scripts/demo.sh: Stopping agent (PID: 726495)... +[1] Terminated ./target/release/rust-edge-agent +./scripts/demo.sh: Stopping device (PID: 726525)... +[2] Terminated ./target/release/device_sim --id "device-$(printf "%03d" $i)" --mode $MODE_ARG --type $TYPE_ARG --interval $DEVICE_INTERVAL >&/dev/null +./scripts/demo.sh: Stopping device (PID: 726526)... +[3]- Terminated ./target/release/device_sim --id "device-$(printf "%03d" $i)" --mode $MODE_ARG --type $TYPE_ARG --interval $DEVICE_INTERVAL >&/dev/null +./scripts/demo.sh: Stopping device (PID: 726528)... +[4]+ Terminated ./target/release/device_sim --id "device-$(printf "%03d" $i)" --mode $MODE_ARG --type $TYPE_ARG --interval $DEVICE_INTERVAL >&/dev/null +./scripts/demo.sh: Cleanup complete + $ ./scripts/stop-services.sh +Stopping NATS broker... +nats_svc +nats_svc +Stopping MQTT broker... +mqtt_svc +mqtt_svc + $ ``` +
+ ### Troubleshooting **Problem: Demo hangs or devices can't register** @@ -237,12 +310,12 @@ ps aux | grep -E 'rust-edge-agent|device_sim' | grep -v grep docker ps | grep mosquitto # If not, start services -./scripts/service-start.sh +./scripts/start-services.sh ``` --- -## Cross-compilation smoke test +## 8.0 Cross-compilation smoke test Before introducing agent logic, this repository verifies that an `AArch64` (`ARM64`) binary can be built on an x86_64 host and executed using `QEMU` user-mode emulation. @@ -284,7 +357,9 @@ This smoke test verifies: Subsequent development assumes this baseline. -## Build and validation workflow +--- + +## 9.0 Build and validation workflow This repository uses small, explicit shell scripts to encode build and validation steps. These scripts are used both locally and in CI to avoid divergence between developer workflows and automated checks. @@ -313,5 +388,37 @@ The following scripts live under `scripts/` and are invoked directly by GitHub A - Executes the `ARM64` binary using `QEMU` user-mode emulation - Validates runtime correctness against an explicit `ARM64` sysroot - Fails if the binary does not successfully execute + - **Note:** When running locally, start services first with `./scripts/start-services.sh` These scripts are designed to be runnable locally and are used directly by CI. + +**Expected output of smoke test** + +``` + $ ./scripts/start-services.sh +Starting NATS broker on localhost:4222 +b269c86b88431ddee1cda182fea52e6198d19bd4d812ede692c2f9df02cc4a6b +Starting mosquitto broker on localhost:MQTT_PORT +7c89fea656922989e19bc969507bd32873501b78f78ba1e805636096b0f43229 + $ ./scripts/ci-qemu-smoke.sh +agent:: Starting edge agent... +agent:: NATS URL: nats://localhost:4222 +agent:: MQTT Broker: mqtt://localhost:1883 +agent:: Device timeout: 30s +agent:: Poll interval: 5s +Connected to NATS at nats://localhost:4222 +agent:: Connecting to MQTT broker... +Connected to MQTT broker at mqtt://localhost:1883 (client_id: agent-transport) +agent:: MQTT transport ready +agent:: Creating RPC server... +agent:: Creating RPC client... +agent:: RPC client ready +agent:: RPC server listening for device registrations +agent:: Edge agent running +2026-02-15T03:35:00.364392Z INFO async_nats: event: connected +2026-02-15T03:35:00.384464Z INFO mom_rpc::transport::rumqttc::transport: rumqttc: connected to broker +2026-02-15T03:35:00.387569Z INFO mom_rpc::transport::rumqttc::transport: rumqttc: successfully subscribed to topic responses/agent-client +2026-02-15T03:35:00.393665Z INFO mom_rpc::transport::rumqttc::transport: rumqttc: successfully subscribed to topic requests/agent +./scripts/ci-qemu-smoke.sh: test passed + $ +``` diff --git a/scripts/demo.sh b/scripts/demo.sh index a80b0b2..6f8707f 100755 --- a/scripts/demo.sh +++ b/scripts/demo.sh @@ -1,5 +1,5 @@ #!/usr/bin/env bash -set -euo pipefail +set -meuo pipefail : ${NUM_DEVICES:=3} : ${DEVICE_INTERVAL:=5} @@ -50,7 +50,7 @@ for i in $(seq 1 $NUM_DEVICES); do --id "device-$(printf "%03d" $i)" \ --mode $MODE_ARG \ --type $TYPE_ARG \ - --interval $DEVICE_INTERVAL & + --interval $DEVICE_INTERVAL >& /dev/null & DEVICE_PIDS+=($!) done @@ -90,7 +90,7 @@ cleanup() { } # Trap EXIT, INT (Ctrl+C), and TERM signals -trap cleanup EXIT INT TERM +trap cleanup INT TERM if [ "$MODE" = "ci" ]; then # CI mode: run automated checks, then exit diff --git a/scripts/service-start.sh b/scripts/service-start.sh deleted file mode 100755 index b44d184..0000000 --- a/scripts/service-start.sh +++ /dev/null @@ -1,6 +0,0 @@ -#!/usr/bin/env bash -set -euo pipefail - -echo "Starting NATS broker..." -docker run -d --name nats -p 4222:4222 nats:latest -echo "NATS broker started on localhost:4222" diff --git a/scripts/service-stop.sh b/scripts/service-stop.sh deleted file mode 100755 index 2d2760b..0000000 --- a/scripts/service-stop.sh +++ /dev/null @@ -1,7 +0,0 @@ -#!/usr/bin/env bash -set -euo pipefail - -echo "Stopping NATS broker..." -docker stop nats 2>/dev/null || true -docker rm nats 2>/dev/null || true -echo "NATS broker stopped" diff --git a/scripts/start-services.sh b/scripts/start-services.sh new file mode 100755 index 0000000..09b1957 --- /dev/null +++ b/scripts/start-services.sh @@ -0,0 +1,11 @@ +#!/usr/bin/env bash +set -euo pipefail + +NAT_PORT=4222 +MQTT_PORT=1883 + +echo "Starting NATS broker on localhost:${NAT_PORT}" +docker run -d --name nats_svc -p ${NAT_PORT}:${NAT_PORT} nats:latest + +echo "Starting mosquitto broker on localhost:MQTT_PORT" +docker run -d --name mqtt_svc -p ${MQTT_PORT}:${MQTT_PORT} eclipse-mosquitto diff --git a/scripts/stop-services.sh b/scripts/stop-services.sh new file mode 100755 index 0000000..4e2b460 --- /dev/null +++ b/scripts/stop-services.sh @@ -0,0 +1,10 @@ +#!/usr/bin/env bash +set -euo pipefail + +echo "Stopping NATS broker..." +docker stop nats_svc 2>/dev/null || true +docker rm nats_svc 2>/dev/null || true + +echo "Stopping MQTT broker..." +docker stop mqtt_svc 2>/dev/null || true +docker rm mqtt_svc 2>/dev/null || true diff --git a/src/agent/run.rs b/src/agent/run.rs index de9adc5..a341f05 100644 --- a/src/agent/run.rs +++ b/src/agent/run.rs @@ -1,7 +1,6 @@ use super::registry::DeviceRegistry; use crate::messaging::{ // - self, CommandRequest, CommandResponse, DeviceMode, @@ -13,7 +12,7 @@ use anyhow::{anyhow, Result}; use async_nats::Client; use clap::Parser; use futures::StreamExt; -use mom_rpc::{RpcClient, RpcServer}; +use mom_rpc::{RpcBroker, RpcBrokerBuilder, TransportBuilder}; use std::collections::HashMap; use std::sync::Arc; use std::time::Duration; @@ -23,6 +22,7 @@ use tokio::sync::Mutex; #[derive(Parser, Debug)] #[command(version, about)] struct Args { + // --- /// NATS server URL (for backend communication) #[arg(long, env = "NATS_URL", default_value = "nats://localhost:4222")] nats_url: String, @@ -43,10 +43,13 @@ struct Args { /// Registered device information. #[derive(Debug, Clone)] struct RegisteredDevice { + // --- #[allow(dead_code)] device_id: String, + #[allow(dead_code)] device_type: crate::messaging::DeviceType, + mode: DeviceMode, service_name: String, } @@ -62,6 +65,8 @@ struct RegisteredDevice { /// - Routes backend commands to actuator devices /// - Forwards telemetry to backend via NATS pub async fn run() -> Result<()> { + // --- + let args = Args::parse(); eprintln!("agent:: Starting edge agent..."); @@ -71,30 +76,34 @@ pub async fn run() -> Result<()> { eprintln!("agent:: Poll interval: {}s", args.poll_interval); // Connect to NATS (for backend communication) - let nats_client = messaging::connect_nats_with_retry(&args.nats_url).await; + let nats_client = crate::connect_nats_with_retry(&args.nats_url).await; // Connect to MQTT (for device communication) eprintln!("agent:: Connecting to MQTT broker..."); - let mqtt_transport = - messaging::create_transport_with_retry(&args.mqtt_broker, "agent-transport").await; + let transport = TransportBuilder::new() + .uri(&args.mqtt_broker) + .node_id("agent") + .full_duplex() + .build() + .await?; eprintln!("agent:: MQTT transport ready"); - // Create RPC server for device registrations - eprintln!("agent:: Creating RPC server..."); - let agent_server = RpcServer::with_transport(mqtt_transport.clone(), "agent"); - - // Create RPC client for polling devices - eprintln!("agent:: Creating RPC client..."); - let agent_client = RpcClient::with_transport(mqtt_transport.clone(), "agent-client").await?; - eprintln!("agent:: RPC client ready"); + // Create bidirectional RPC broker (handles both registration requests and + // outbound device polling/commands on a single MQTT connection) + let broker = RpcBrokerBuilder::new(transport) + .retry_max_attempts(5) + .request_total_timeout(Duration::from_secs(args.device_timeout)) + .build()?; let timeout = Duration::from_secs(args.device_timeout); let registry = Arc::new(Mutex::new(DeviceRegistry::new(timeout))); let devices = Arc::new(Mutex::new(HashMap::::new())); - // Register device registration handler + // Register device registration handler (server-side) let devices_clone = devices.clone(); - agent_server.register("register-device", move |req: RegisterRequest| { + broker.register_rpc_handler("register-device", move |req: RegisterRequest| { + // --- + let devices = devices_clone.clone(); async move { let service_name = match req.mode { @@ -122,17 +131,14 @@ pub async fn run() -> Result<()> { message: Some("Registration successful".to_string()), }) } - }); + })?; - // Spawn RPC server - let _server_handle = agent_server.spawn(); - eprintln!("agent:: RPC server listening for device registrations"); + // Spawn broker receive loop + let _broker_handle = broker.clone().spawn()?; + eprintln!("agent:: RPC broker listening for device registrations"); // Subscribe to backend commands via NATS - let mut command_sub = match nats_client - .subscribe(messaging::all_backend_commands()) - .await - { + let mut command_sub = match nats_client.subscribe(crate::all_backend_commands()).await { Ok(c) => c, Err(e) => { return Err(anyhow!( @@ -146,14 +152,14 @@ pub async fn run() -> Result<()> { // Spawn telemetry polling task let devices_poll = devices.clone(); let nats_poll = nats_client.clone(); - let mqtt_poll = agent_client.clone(); + let broker_poll = broker.clone(); let registry_poll = registry.clone(); let poll_interval = args.poll_interval; tokio::spawn(async move { poll_device_telemetry( devices_poll, nats_poll, - mqtt_poll, + broker_poll, registry_poll, poll_interval, ) @@ -165,7 +171,7 @@ pub async fn run() -> Result<()> { tokio::select! { Some(msg) = command_sub.next() => { if let Err(e) = handle_backend_command( - &nats_client, &agent_client, &devices, msg).await { + &nats_client, &broker, &devices, msg).await { eprintln!("agent:: Command error: {e}"); } } @@ -177,10 +183,11 @@ pub async fn run() -> Result<()> { async fn poll_device_telemetry( devices: Arc>>, nats_client: Client, - mqtt_client: RpcClient, + broker: RpcBroker, registry: Arc>, interval_secs: u64, ) { + // --- let mut interval = tokio::time::interval(Duration::from_secs(interval_secs)); loop { @@ -195,7 +202,7 @@ async fn poll_device_telemetry( } // Poll device for telemetry - let result: Result = mqtt_client + let result: Result = broker .request_to(&device.service_name, "read-telemetry", ()) .await; @@ -212,7 +219,7 @@ async fn poll_device_telemetry( // Forward to backend via NATS if let Ok(payload) = serde_json::to_vec(&telemetry) { let _ = nats_client - .publish(messaging::backend_telemetry(), payload.into()) + .publish(crate::backend_telemetry(), payload.into()) .await; } } @@ -227,10 +234,11 @@ async fn poll_device_telemetry( /// Handle backend command and route to appropriate device. async fn handle_backend_command( nats_client: &Client, - mqtt_client: &RpcClient, + broker: &RpcBroker, devices: &Arc>>, msg: async_nats::Message, ) -> Result<()> { + // --- // Extract device ID from NATS subject let device_id = match msg.subject.strip_prefix("backend.command.") { Some(dev) => dev, @@ -270,11 +278,12 @@ async fn handle_backend_command( let command: CommandRequest = serde_json::from_slice(&msg.payload)?; // Route to device via MQTT RPC - let result: Result = mqtt_client + let result: Result = broker .request_to(&device.service_name, "execute-command", command) .await; match result { + // --- Ok(response) => { eprintln!("agent:: Command success: {response:?}"); if let Some(reply) = msg.reply { diff --git a/src/bin/device_sim.rs b/src/bin/device_sim.rs index 33d7ea1..488e380 100644 --- a/src/bin/device_sim.rs +++ b/src/bin/device_sim.rs @@ -5,10 +5,17 @@ use anyhow::{bail, Result}; use clap::Parser; -use mom_rpc::{RpcClient, RpcServer}; +use mom_rpc::{RpcBroker, RpcBrokerBuilder, TransportBuilder}; use rust_edge_agent::messaging::{ - create_transport_with_retry, ActuatorState, CommandRequest, CommandResponse, DeviceMode, - DeviceType, RegisterRequest, RegisterResponse, TelemetryMessage, + // --- + ActuatorState, + CommandRequest, + CommandResponse, + DeviceMode, + DeviceType, + RegisterRequest, + RegisterResponse, + TelemetryMessage, }; use std::sync::Arc; use std::time::{Duration, SystemTime, UNIX_EPOCH}; @@ -17,6 +24,7 @@ use tokio::sync::Mutex; #[derive(Parser, Debug)] #[command(version, about)] struct Args { + // --- /// Unique device identifier (numeric) #[arg(long)] id: String, @@ -40,6 +48,7 @@ struct Args { /// Shared device state for sensor simulation. struct DeviceState { + // --- device_id: String, device_type: DeviceType, current_value: f64, @@ -47,7 +56,8 @@ struct DeviceState { #[tokio::main] async fn main() -> Result<()> { - env_logger::init(); + // --- + tracing_subscriber::fmt().with_ansi(false).init(); let args = Args::parse(); let device_type = parse_device_type(&args.r#type)?; @@ -66,15 +76,24 @@ async fn main() -> Result<()> { }; // Create shared MQTT transport - let transport_id = format!("{service_name}-transport"); - let transport = create_transport_with_retry(&args.mqtt_broker, &transport_id).await; - - // Create RPC server (receives calls from agent) - let server = RpcServer::with_transport(transport.clone(), &service_name); - - // Create RPC client (calls agent for registration) - let client_id = format!("{service_name}-client"); - let client = RpcClient::with_transport(transport.clone(), &client_id).await?; + let transport = TransportBuilder::new() + .uri(&args.mqtt_broker) + .node_id(&service_name) + .full_duplex() + .build() + .await?; + + // Create bidirectional RPC broker (handles both incoming agent calls + // and outbound registration request on a single MQTT connection) + let broker = RpcBrokerBuilder::new(transport) + .retry_max_attempts(1000) + .retry_initial_delay(Duration::from_millis(200)) + .retry_max_delay(Duration::from_secs(5)) + .request_total_timeout(Duration::from_secs(3600)) + .build()?; + // max delay is 1000 x 5 => 5000 seconds, but request_total_timeout will + // shorten it to 3600. The request_to_with_timeout will increase + // request_to_with_timeout. // Shared state for sensor value simulation let state = Arc::new(Mutex::new(DeviceState { @@ -86,19 +105,19 @@ async fn main() -> Result<()> { // Register methods based on mode match mode { DeviceMode::Sensor => { - register_sensor_methods(&server, state.clone()); + register_sensor_methods(&broker, state.clone())?; } DeviceMode::Actuator => { - register_actuator_methods(&server, state.clone()); + register_actuator_methods(&broker, state.clone())?; } DeviceMode::Hybrid => { - register_sensor_methods(&server, state.clone()); - register_actuator_methods(&server, state.clone()); + register_sensor_methods(&broker, state.clone())?; + register_actuator_methods(&broker, state.clone())?; } } - // Spawn server to handle incoming RPC calls - let server_handle = server.spawn(); + // Spawn broker receive loop + let broker_handle = broker.clone().spawn()?; // Register with agent (with retry) eprintln!("[{service_name}] Registering with agent..."); @@ -108,45 +127,38 @@ async fn main() -> Result<()> { mode, }; - let mut retry_count = 0; - let max_retries = 10; - let mut retry_delay = Duration::from_millis(500); - - loop { - match client - .request_to("agent", "register-device", register_req.clone()) - .await - { - Ok(resp) => { - let response: RegisterResponse = resp; - if response.accepted { - eprintln!("[{service_name}] Registration accepted"); - break; - } else { - eprintln!( - "[{}] Registration rejected: {:?}", - service_name, response.message - ); - return Ok(()); - } - } - Err(e) => { - retry_count += 1; - if retry_count > max_retries { - eprintln!( - "[{service_name}] Registration failed after {max_retries} attempts: {e}", - ); - return Ok(()); - } - eprintln!( - "[{service_name}] Registration attempt {retry_count}/{max_retries}\ - failed: {e}. Retrying in {retry_delay:?}...", - ); - tokio::time::sleep(retry_delay).await; - retry_delay = (retry_delay * 2).min(Duration::from_secs(5)); - } - } - } + // Register with agent using a very long timeout. In real deployments, the + // agent and devices start independently with no guaranteed ordering. For + // example, in a vehicle telematics system, sensor nodes on the CAN bus may + // power up before the gateway agent has finished booting or reconnected + // after a network flap. Rather than failing fast, devices wait patiently + // for the agent to become reachable. The broker handles retries internally; + // the long ceiling here covers realistic startup delays without requiring + // external orchestration. + eprintln!("[{service_name}] Registering with agent..."); + let timeout_seconds = 5000; + let response: RegisterResponse = broker + // Overriding default timeout. + .request_to_with_timeout( + "agent", + "register-device", + register_req.clone(), + Duration::from_secs(timeout_seconds), + ) + .await + .map_err(|err| { + anyhow::anyhow!( + "[{service_name}] Registration failed after {timeout_seconds} seconds: {err}" + ) + })?; + + if !response.accepted { + bail!( + "[{service_name}] Registration rejected: {:?}", + response.message + ); + }; + eprintln!("[{service_name}] Registration accepted"); // Spawn sensor value simulation task (if sensor mode) if matches!(mode, DeviceMode::Sensor | DeviceMode::Hybrid) { @@ -158,17 +170,19 @@ async fn main() -> Result<()> { eprintln!("[{service_name}] Device running (Ctrl+C to stop)"); - // Wait for server shutdown - if let Err(e) = server_handle.await { - eprintln!("[{service_name}] Server error: {e}"); + // Wait for broker shutdown + if let Err(e) = broker_handle.await { + eprintln!("[{service_name}] Broker error: {e}"); } Ok(()) } /// Register RPC methods for sensor mode. -fn register_sensor_methods(server: &RpcServer, state: Arc>) { - server.register("read-telemetry", move |_req: ()| { +fn register_sensor_methods(broker: &RpcBroker, state: Arc>) -> Result<()> { + // --- + + broker.register_rpc_handler("read-telemetry", move |_req: ()| { let state = state.clone(); async move { let state = state.lock().await; @@ -184,14 +198,16 @@ fn register_sensor_methods(server: &RpcServer, state: Arc>) { value: state.current_value, }) } - }); + })?; + Ok(()) } /// Register RPC methods for actuator mode. -fn register_actuator_methods(server: &RpcServer, state: Arc>) { +fn register_actuator_methods(broker: &RpcBroker, state: Arc>) -> Result<()> { + // --- // execute-command method let state_cmd = state.clone(); - server.register("execute-command", move |req: CommandRequest| { + broker.register_rpc_handler("execute-command", move |req: CommandRequest| { let state = state_cmd.clone(); async move { let mut state = state.lock().await; @@ -207,10 +223,10 @@ fn register_actuator_methods(server: &RpcServer, state: Arc>) message: Some(format!("Set to {:.2}", req.target_value)), }) } - }); + })?; // read-state method - server.register("read-state", move |_req: ()| { + broker.register_rpc_handler("read-state", move |_req: ()| { let state = state.clone(); async move { let state = state.lock().await; @@ -224,11 +240,14 @@ fn register_actuator_methods(server: &RpcServer, state: Arc>) timestamp, }) } - }); + })?; + Ok(()) } /// Simulate sensor value changes over time. async fn simulate_sensor_values(state: Arc>, interval: u64) { + // --- + loop { { let mut state = state.lock().await; diff --git a/src/lib.rs b/src/lib.rs index c501165..0389826 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -15,7 +15,10 @@ pub use runtime::{Lifecycle, LifecycleState}; // Re-export commonly used messaging types for convenience pub use messaging::{ - // + // --- + all_backend_commands, + backend_telemetry, + connect_nats_with_retry, ActuatorState, CommandRequest, CommandResponse, diff --git a/src/main.rs b/src/main.rs index 4cf8514..777b0f7 100644 --- a/src/main.rs +++ b/src/main.rs @@ -4,6 +4,9 @@ use rust_edge_agent::run; #[tokio::main] async fn main() -> Result<()> { // --- + // --- + tracing_subscriber::fmt().with_ansi(false).init(); + match run().await { Ok(_) => Ok(()), Err(err) => Err(anyhow!("Error:{err}")), diff --git a/src/messaging/mod.rs b/src/messaging/mod.rs index 2b8cc69..312e379 100644 --- a/src/messaging/mod.rs +++ b/src/messaging/mod.rs @@ -8,12 +8,10 @@ //! The agent uses dual messaging protocols: //! //! - **NATS**: Backend ↔ Agent communication (unchanged) -//! - **MQTT + mom-rpc**: Agent ↔ Device communication (new) //! //! NATS subjects are kept for backward compatibility but are marked as legacy //! since device communication now uses MQTT RPC methods. -mod mom_rpc; mod nats; mod subjects; mod types; @@ -21,9 +19,6 @@ mod types; // NATS connection helpers (for Backend ↔ Agent) pub use nats::connect_with_retry as connect_nats_with_retry; -// MQTT transport helpers (for Agent ↔ Devices) -pub use mom_rpc::{create_transport_once, create_transport_with_retry}; - // NATS subjects (legacy - used only for Backend ↔ Agent communication) pub use subjects::{ // diff --git a/src/messaging/mom_rpc.rs b/src/messaging/mom_rpc.rs deleted file mode 100644 index 914c81a..0000000 --- a/src/messaging/mom_rpc.rs +++ /dev/null @@ -1,77 +0,0 @@ -//! mom-rpc transport helpers for MQTT-based device communication. -//! -//! Provides connection management, retry logic, and transport creation -//! for both RPC clients and servers in the edge agent system. - -use anyhow::Result; -use mom_rpc::{create_transport, RpcConfig, TransportPtr}; -use std::time::Duration; -use tokio::time::sleep; - -/// Create MQTT transport with exponential backoff retry. -/// -/// # Arguments -/// -/// * `broker_url` - MQTT broker URL (e.g., "mqtt://localhost:1883") -/// * `client_id` - Unique client identifier for MQTT connection -/// -/// # Behavior -/// -/// - Retries with exponential backoff: 100ms, 200ms, 400ms, ..., 30s (max) -/// - Continues retrying indefinitely until connection succeeds -/// - Logs connection failures but does not crash -/// -/// # Returns -/// -/// A shared transport instance that can be cloned for both RpcClient and RpcServer. -/// -/// # Example -/// -/// ```no_run -/// use rust_edge_agent::messaging::create_transport_with_retry; -/// -/// let transport = create_transport_with_retry( -/// "mqtt://localhost:1883", -/// "device-sensor-1" -/// ).await; -/// -/// // Share transport between client and server -/// let server = RpcServer::with_transport(transport.clone(), "sensor-1"); -/// let client = RpcClient::with_transport(transport.clone(), "sensor-1-client").await?; -/// ``` -pub async fn create_transport_with_retry(broker_url: &str, client_id: &str) -> TransportPtr { - let mut backoff = Duration::from_millis(100); - - loop { - let config = RpcConfig::with_broker(broker_url, client_id); - - match create_transport(&config).await { - Ok(transport) => { - eprintln!("Connected to MQTT broker at {broker_url} (client_id: {client_id})",); - return transport; - } - Err(e) => { - eprintln!("MQTT connection failed: {e}, retrying in {backoff:?}...",); - sleep(backoff).await; - backoff = (backoff * 2).min(Duration::from_secs(30)); - } - } - } -} - -/// Create MQTT transport without retry. -/// -/// # Arguments -/// -/// * `broker_url` - MQTT broker URL (e.g., "mqtt://localhost:1883") -/// * `client_id` - Unique client identifier for MQTT connection -/// -/// # Errors -/// -/// Returns error if initial connection fails. Use `create_transport_with_retry` -/// for resilient connection establishment. -pub async fn create_transport_once(broker_url: &str, client_id: &str) -> Result { - let config = RpcConfig::with_broker(broker_url, client_id); - let transport = create_transport(&config).await?; - Ok(transport) -}