diff --git a/CHANGELOG.md b/CHANGELOG.md index 676c730..7e08105 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,3 +1,54 @@ +## [Unreleased] + +## [0.2.0] – 2026-01-31 + +### Added +- **Device simulator** (`device_sim` binary) supporting three modes: + - Sensor: Publish-only telemetry (temperature, humidity) + - Actuator: Command handling with status publishing (valve, propulsion) + - Hybrid: Both telemetry and command capabilities +- **NATS-based messaging architecture**: + - Request/reply for actuator commands with timeout handling + - Pub/sub for telemetry forwarding + - Subject-based routing (`devices.*`, `backend.*`) +- **Device registry** with state tracking and configurable timeout detection +- **Exponential backoff connection retry** (1s → 2s → 4s → 8s → 16s → 30s cap) + - Uses bit-shift implementation for efficient exponential calculation + - Resilient reconnection for embedded systems without UI +- **Service management scripts**: + - `service-start.sh` - Start NATS broker in Docker + - `service-stop.sh` - Stop and remove NATS broker + - `demo.sh` - Interactive demo with configurable device count +- **CONTRIBUTING.md** with comprehensive guidelines: + - Code formatting conventions (`// ---` separators) + - Documentation standards for messaging/edge systems + - EMBP (Explicit Module Boundary Pattern) architecture reference + - Testing strategy and coverage expectations + +### Changed +- **Applied EMBP architecture pattern** throughout codebase: + - Private module declarations with gateway exports + - Sibling imports via `super::`, external via `crate::` + - Messaging module serves as public API gateway +- **CLI argument parsing** via CLAP with environment variable fallbacks: + - `--nats-url` / `NATS_URL` (default: `nats://localhost:4222`) + - `--device-timeout` / `DEVICE_TIMEOUT` (default: 30s) + - `--interval` / `DEVICE_INTERVAL` for device simulator +- **Async runtime** using Tokio for agent and device simulators + +### Documentation +- **README.md**: Added Quick Start guide and demo instructions +- **docs/architecture.md**: Describes gateway vs leaf device patterns +- **CONTRIBUTING.md**: Production-grade documentation standards and EMBP patterns + +### Architecture Decisions +- **NATS over MQTT**: Chose NATS for request/reply semantics and simpler implementation + - Avoids MQTT correlation ID complexity (deferred to Phase 2 as `mqtt-rpc` library) +- **Forward raw telemetry**: Edge agent forwards telemetry without aggregation + - Backend has compute/storage for aggregation; keeps edge agent simple +- **Gateway pattern**: Demonstrates coordination between devices and backend + - Not a leaf device - maintains state, routes bidirectionally + ## [0.1.0] – 2026-01-27 ### Added diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md new file mode 100644 index 0000000..52dbd61 --- /dev/null +++ b/CONTRIBUTING.md @@ -0,0 +1,307 @@ +# Contributing to rust-edge-agent + +Thanks for considering contributing! + +## Quick Start + +**New to the project?** See README.md for cross-compilation setup and QEMU validation. + +## Local Development + +### Starting Services + +Before running the agent locally: + +```bash +# Start NATS broker +./scripts/service-start.sh + +# When done +./scripts/service-stop.sh +``` + +### Running the Demo + +```bash +# Start 3 devices (default) +./scripts/demo.sh + +# Start 10 devices +NUM_DEVICES=10 ./scripts/demo.sh + +# Use custom telemetry interval +DEVICE_INTERVAL=2 ./scripts/demo.sh +``` + +See `scripts/demo.sh` for monitoring and testing examples. + +**Before submitting a pull request:** + +- Run local CI scripts (includes fmt, clippy, and builds): + ```bash + ./scripts/ci-lint.sh + ./scripts/ci-build-native.sh + ./scripts/ci-build-aarch64.sh + cargo test --release + ``` +- If your change affects behavior, please update `CHANGELOG.md` under the [Unreleased] section +- Keep commits focused and descriptive + +We follow [Keep a Changelog](https://keepachangelog.com/en/1.0.0/) and [Semantic Versioning](https://semver.org/). + +## Code Formatting + +This project uses `rustfmt` for consistent code formatting. All code should be formatted before committing. + +### Visual Separators + +Since `rustfmt` removes blank lines at the start of impl blocks, function bodies, and module blocks, we use comment separators `// ---` for visual clarity: + +```rust +// Module blocks +mod messaging { + // --- + use super::*; + + pub fn start_control_handler() { + // --- + // function body + } +} + +// Struct definitions +pub struct DeviceState { + // --- + device_id: String, + last_seen: Instant, +} + +// Impl blocks +impl ControlHandler for NatsControlHandler { + // --- + async fn handle_command(&self, cmd: Command) -> Result { + // --- + // implementation + } +} + +// Regular functions +pub async fn route_device_command(device_id: &str, cmd: Command) { + // --- + let client = get_device_client(device_id); + // ... +} + +// Struct literals (construction) - NO separator +let config = MqttOptions { + client_id, + broker_addr, + port, +}; + +// Test modules +#[cfg(test)] +mod tests { + // --- + use super::*; + + #[tokio::test] + async fn test_command_routing() { + // --- + // test body + } +} +``` + +**Style Guidelines:** +1) Use `// ---` for visual separation in at a minimum **module blocks**, **impl blocks**, **struct definitions**, and **function bodies** +2) Place separators after the opening brace and before the first meaningful line +3) Between meaningful steps of logic processing (e.g., separating message parsing, routing, and response handling) +4) For modules: place separator after `mod name {` and before imports/content +5) For impl blocks: place separator after `impl ... {` and before the first method +6) For struct definitions: place separator after `struct Name {` and before field declarations +7) For functions: place separator after function signature and before the main logic +8) Do NOT use separators inside struct literals (during construction) +9) Keep separators consistent across the codebase + +**Note:** This project uses rustfmt's default configuration. The `// ---` separator pattern is a formatting convention to work around rustfmt's blank line removal in stable Rust. + +## Documentation and Doc Comments + +This project follows a **production-grade documentation standard** for Rust code, with special attention to embedded systems and messaging patterns. + +### Required Doc Comments + +Use Rust doc comments (`///`) for: + +- Public structs and enums (especially messaging types like `ControlCommand`, `TelemetryMessage`) +- Public functions (especially handlers and control plane methods) +- Public modules that define architectural boundaries +- Critical system behavior (device lifecycle, message routing, failure handling) +- Macros that encode non-obvious behavior or policy decisions + +Doc comments should describe **intent, guarantees, and failure semantics** — +not restate what the code obviously does. + +### Messaging/Edge-Specific Documentation + +For messaging and control plane code, doc comments should explicitly describe: + +- **Failure modes** - What happens when devices disconnect, messages timeout, etc.? +- **Message flow** - Which part of the control/telemetry flow is this? +- **Delivery semantics** - At-most-once, at-least-once, exactly-once? +- **Concurrency** - Can multiple messages be in-flight? How are they handled? + +Example: +```rust +/// Routes a control command to the specified device. +/// +/// This implements the edge agent's command routing logic, translating +/// backend requests into device-specific commands. +/// +/// # Behavior +/// +/// - Uses NATS request/reply for synchronous command execution +/// - Waits up to 5 seconds for device acknowledgment +/// - Returns error if device is offline or command times out +/// +/// # Errors +/// +/// Returns an error if: +/// - The device ID is unknown or offline +/// - The command times out (5s default) +/// - The device returns an error response +pub async fn route_command(device_id: &str, cmd: Command) -> Result { + // --- + // implementation +} +``` + +### Optional (Encouraged) Doc Comments + +Doc comments or short block comments are encouraged for: + +- Internal functions with concurrency or timing implications +- Device state management logic +- Message serialization and validation +- Configuration parsing and validation +- Startup and initialization logic + +### Not Required + +Doc comments are not required for: + +- Trivial helpers +- Simple getters or pass-through functions +- Test code (assert messages should be sufficient) +- Obvious glue code + +### General Guidance + +- Prefer documenting *why* over *how* +- Be explicit about failure behavior and recovery +- Keep comments accurate and up to date +- Avoid over-documenting trivial code +- For messaging patterns, describe delivery semantics clearly + +Well-written doc comments are considered part of the code's correctness, especially for distributed systems and edge infrastructure. + +## Architecture Guidelines + +This project uses the [Explicit Module Boundary Pattern (EMBP)](https://github.com/JohnBasrai/architecture-patterns/blob/main/rust/embp.md) for module organization. Please review the EMBP documentation before making structural changes. + +### Key EMBP Principles + +- Each module's public API is defined in its `mod.rs` gateway file +- Sibling modules import from each other using `super::` +- External modules import through `crate::module::` +- Never bypass module gateways with deep imports + +### Edge Agent Module Structure + +``` +src/ +├── agent/ # Agent lifecycle and coordination +├── messaging/ # NATS/MQTT messaging abstraction +├── runtime/ # Device registry, state management +└── bin/ + └── device_sim.rs # Device simulator for testing +``` + +## Test Coverage + +This project uses a **layered testing approach** optimized for embedded systems: + +### Current Test Strategy + +**QEMU Smoke Tests (Primary):** +- `scripts/ci-qemu-smoke.sh` - Validates ARM64 binary execution +- Ensures cross-compilation correctness +- Tests basic runtime behavior + +**Integration Tests (Planned):** +- Multi-device scenarios with NATS broker +- Command routing and telemetry aggregation +- Failure recovery and reconnect logic + +**Unit Tests:** +- Core logic (device state, message parsing) +- Lifecycle transitions +- Error handling + +### When to Add Tests + +**Add integration tests when:** +- Adding new messaging patterns +- Changing device lifecycle behavior +- Implementing failure recovery logic + +**Add unit tests when:** +- Complex business logic needs isolated testing +- Edge cases are difficult to trigger via integration tests +- Testing device state transitions + +### Test Organization + +``` +scripts/ + ci-lint.sh # Formatting and clippy + ci-build-native.sh # Native x86_64 build + ci-build-aarch64.sh # ARM64 cross-compilation + ci-qemu-smoke.sh # QEMU validation +``` + +**Running tests:** +```bash +# Local build validation +./scripts/ci-lint.sh +./scripts/ci-build-native.sh +./scripts/ci-build-aarch64.sh + +# QEMU smoke test (requires qemu-user and cross-compilation tools) +./scripts/ci-qemu-smoke.sh + +# All workspace tests (use --release to match build configuration) +cargo test --release +``` + +**Note:** This is a portfolio/demo project showcasing embedded Linux patterns and cross-compilation. Production code would include more comprehensive integration tests and hardware-in-the-loop validation. + +## Testing Edge Agent Behavior + +When testing edge agent and messaging flows: + +- Test both connected and disconnected device scenarios +- Verify message routing and aggregation +- Test timeout and retry behavior +- Include tests for edge cases (duplicate messages, out-of-order delivery) +- Validate device lifecycle transitions + +## Cross-Compilation Notes + +This project targets `aarch64-unknown-linux-gnu`. When adding dependencies: + +- Verify they support cross-compilation (check for C dependencies) +- Test on both native and ARM64 targets +- Document any platform-specific behavior +- Update CI scripts if new system dependencies are needed diff --git a/Cargo.toml b/Cargo.toml index 5bb1fd8..2bf67ce 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,12 +1,23 @@ [package] name = "rust-edge-agent" -version = "0.1.0" +version = "0.2.0" edition = "2021" license = "MIT OR Apache-2.0" [dependencies] +anyhow = "1.0" async-nats = "0.35" +clap = { version = "4", features = ["derive", "env"] } futures = "0.3" +rand = "0.8" serde = { version = "1", features = ["derive"] } serde_json = "1" tokio = { version = "1", features = ["macros", "rt-multi-thread"] } + +[[bin]] +name = "rust-edge-agent" +path = "src/main.rs" + +[[bin]] +name = "device_sim" +path = "src/bin/device_sim.rs" diff --git a/README.md b/README.md index b4380ba..10411bc 100644 --- a/README.md +++ b/README.md @@ -2,9 +2,9 @@ [![CI](https://github.com/JohnBasrai/rust-edge-agent/actions/workflows/ci.yml/badge.svg)](https://github.com/JohnBasrai/rust-edge-agent/actions/workflows/ci.yml) -This project focuses on the mechanics of building and validating an embedded Linux edge agent—specifically cross-compilation, runtime correctness, and reproducible CI workflows. +This project demonstrates edge gateway patterns through cross-compilation, NATS messaging, and device coordination. Phase 1 implements device simulators and command routing to establish gateway architecture foundations. -Higher-level distributed-system behaviors (coordination, failure handling, and reconnect semantics) are minimal at this stage and intentionally deferred. +Current state: Working NATS-based edge agent with device simulators demonstrating sensor/actuator patterns, command routing, and telemetry aggregation. ## What this is @@ -42,6 +42,73 @@ This repository intentionally avoids expanding into those areas in order to keep --- +## Quick Start + +### Prerequisites +- Docker (for NATS broker) +- Rust toolchain (see `rust-toolchain.toml`) +- For cross-compilation: `gcc-aarch64-linux-gnu` and `qemu-user` +- nats CLI (`apt-get install nats-io/nats-tools/nats` or equivalent) +- localhost:4222 available + +Required: +- Docker +- Rust (stable) +- QEMU user emulation + +Install on Ubuntu/Debian: +```bash +sudo apt-get update +sudo apt-get install -y \ + qemu-user \ + gcc-aarch64-linux-gnu \ + libc6-arm64-cross +``` +Optional (for local inspection): +```bash +sudo apt-get install -y natscli +``` +macOS users may install equivalents via Homebrew, but CI and official +support assume a Debian-based Linux environment. + +### Running the Demo + +**1. Start NATS broker** +```bash +./scripts/service-start.sh +``` + +**2. Build the project** +```bash +cargo build --release +``` + +**3. Run the demo** (starts edge agent + 3 device simulators) +```bash +./scripts/demo.sh +``` + +The demo starts: +- Edge agent listening for device telemetry and backend commands +- 3 device simulators (sensor, actuator, hybrid modes) + +**Monitor telemetry:** +```bash +nats sub 'backend.telemetry' +``` + +**Send command to actuator:** +```bash +nats req 'backend.command.device-002' '{"target_value": 75.0}' +``` + +**4. Cleanup** +```bash +./scripts/service-stop.sh +``` + +--- + ## 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. diff --git a/docs/architecture.md b/docs/architecture.md index 5aa06a8..9d71b4b 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -63,6 +63,58 @@ The agent is a single long-running process with a clearly defined lifecycle and JetStream is **not required** for this project but is discussed as an optional extension. +### NATS vs MQTT + +This project uses **NATS** rather than MQTT, despite MQTT being common in edge/IoT deployments. + +**Rationale:** +- **Request/reply semantics**: NATS provides native request/reply, avoiding the need to implement correlation IDs manually +- **Simpler implementation**: Focus on edge gateway coordination logic rather than message correlation +- **Phase 2 scope**: The correlation problem inherent in pure MQTT pub/sub will be addressed as a reusable `mqtt-rpc` library + +NATS allows us to demonstrate gateway patterns cleanly while deferring the correlation complexity to a future abstraction layer that can be applied to MQTT, AMQP, or other pub/sub transports. + +--- + +## Phase 1 Implementation + +The current implementation demonstrates core gateway patterns: + +**Device Simulators:** +- Three operational modes: sensor, actuator, hybrid +- Device types: temperature, humidity, valve, propulsion +- Configurable telemetry intervals +- Simulated value generation with random walk + +**NATS Subject Design:** +``` +# Device → Agent +devices..telemetry (pub) +devices..status (pub, reserved for future) + +# Backend → Agent → Device +backend.command. (request/reply, agent subscribes) +devices..command (request/reply, device subscribes) + +# Agent → Backend +backend.telemetry (pub, aggregated forwarding) +``` + +**Message Flow:** +1. Device simulators publish telemetry to `devices..telemetry` +2. Edge agent subscribes to `devices.*.telemetry` (wildcard) +3. Agent updates device registry and forwards to `backend.telemetry` +4. Backend sends commands to `backend.command.` +5. Agent routes to device via `devices..command` (request/reply) +6. Device responds with status, agent forwards response to backend + +**Command Timeout Handling:** +- Agent uses NATS request/reply with default timeout +- Returns error response if device unreachable +- Backend caller receives either device response or timeout error + +The agent is a single long-running process with a clearly defined lifecycle and explicit separation between control and telemetry concerns. + --- ### Control Plane (Request / Reply) diff --git a/make-edge-sync.sh b/make-edge-sync.sh index 2b8b229..a50f6ed 100755 --- a/make-edge-sync.sh +++ b/make-edge-sync.sh @@ -5,6 +5,7 @@ tar cfvz make-edge-sync.gz \ .github \ .gitignore \ CHANGELOG.md \ + CONTRIBUTING.md \ Cargo.lock \ Cargo.toml \ LICENSE \ diff --git a/new-features/phase1.md b/new-features/phase1.md new file mode 100644 index 0000000..a62ffb1 --- /dev/null +++ b/new-features/phase1.md @@ -0,0 +1,45 @@ +I'm working on `rust-edge-agent`, an embedded Linux edge gateway for ARM64 systems. The project demonstrates cross-compilation and QEMU validation. + +**Current state:** Basic "hello world" with NATS dependencies, cross-compiles to aarch64, runs under QEMU. + +**Next milestone (Phase 1):** Add real functionality using NATS messaging to demonstrate edge gateway patterns. + +**Architecture Decision:** +- Use NATS for ALL messaging (devices ↔ edge agent ↔ backend) +- Rationale: Focus on edge agent coordination logic, avoid reimplementing MQTT request/response correlation +- Phase 2 will extract this correlation as a reusable `mqtt-rpc` crate + +**What to build:** + +1. **Device Simulator** (`src/bin/device_sim.rs`): + - Single binary with CLI args: `--id --mode --type ` + - Modes: + - `sensor`: Publishes telemetry periodically (pub only) + - `actuator`: Handles commands, publishes status/acks (sub + pub) + - `hybrid`: Both telemetry and command handling + - Uses NATS request/reply for actuator commands (proper RPC semantics) + - Simple JSON payloads + - Can run multiple instances with different IDs + +2. **Edge Agent** (enhance `src/main.rs` and `src/agent/`): + - Connects to NATS + - Subscribes to device telemetry (aggregates, forwards to backend) + - Routes control commands from backend to devices + - Uses NATS request/reply for commands to actuators + - Maintains basic device registry/state + +**Success criteria:** +- Start 3-4 device simulators (mix of sensors/actuators) +- Edge agent aggregates telemetry, routes commands +- Can send backend commands via `nats` CLI and see device responses +- Demonstrates gateway-class coordination patterns + +**Constraints:** +- Keep it simple - this is foundation for Phase 2 (mqtt-rpc library) +- Focus on architecture, not feature breadth +- Maintain cross-compilation and QEMU validation +- Document the MQTT correlation problem this avoids + +Also let's use `CONTRIBUTING.md` in this repo and we work under `WORKFLOW-v1.2.md` workflow. + +Ready to implement Phase 1. Where should we start? diff --git a/new-features/phase2.md b/new-features/phase2.md new file mode 100644 index 0000000..c684b96 --- /dev/null +++ b/new-features/phase2.md @@ -0,0 +1,151 @@ +I want to create `mqtt-rpc`, a Rust crate that provides RPC semantics over MQTT pub/sub. + +## Problem Statement + +MQTT is a lightweight pub/sub protocol widely used in IoT, but it lacks built-in RPC semantics. Every developer using MQTT for request/response patterns must solve the same problems: + +1. **Request/response correlation** - Matching responses to requests using correlation IDs +2. **Timeout handling** - Detecting when a request will never get a response +3. **Concurrent request handling** - Processing multiple in-flight requests without blocking +4. **Code duplication** - This logic is reimplemented in every MQTT client/server application + +This crate extracts that common pattern into a reusable library. + +## Design Goals + +**Client-side:** +- Send request, get Future that resolves when response arrives +- Automatic correlation ID generation (compact: counter-based, not UUIDs) +- Timeout support +- Concurrent requests (multiple in-flight) + +**Server-side:** +- Register async handlers for request topics +- Handlers execute concurrently (spawned, not blocking event loop) +- Automatic correlation ID handling +- Response publishing + +**Implementation:** +- Built on `rumqttc` (popular async Rust MQTT client) +- Tokio-based async runtime +- Compact correlation IDs: `{device_id}:{counter}` format +- Internal `HashMap>` for pending requests + +## Proposed API + +### Client Side +```rust +use mqtt_rpc::RpcClient; + +let client = RpcClient::new(mqtt_client, "client-01").await?; + +// Send request, wait for response (with timeout) +let response: Response = client + .request("devices/valve-01/command", request_payload) + .timeout(Duration::from_secs(5)) + .await?; +``` + +### Server Side +```rust +use mqtt_rpc::RpcServer; + +let server = RpcServer::new(mqtt_client, "device-01").await?; + +// Register handler - runs concurrently for each request +server.handle("command", |req: Request| async move { + // Long-running async operation + actuator.open().await?; + + Ok(Response { + status: "opened", + position: 1.0, + }) +}).await?; + +server.run().await?; +``` + +## Technical Requirements + +1. **Correlation ID Strategy:** + - Use `AtomicU64` counter for compact IDs + - Format: `"{namespace}:{counter}"` (e.g., "valve-01:42") + - Namespace prevents conflicts in multi-device scenarios + +2. **Concurrent Handler Execution:** + - Spawn handler tasks with `tokio::spawn` + - Don't block MQTT event loop + - Multiple requests can be processed simultaneously + +3. **Request/Response Flow:** +``` + Client: + 1. Generate correlation_id + 2. Store oneshot::Sender in HashMap + 3. Publish to request topic with correlation_id in payload + 4. Await on oneshot::Receiver + + Server: + 1. Receive request with correlation_id + 2. Spawn handler task + 3. Handler completes → publish response with same correlation_id + + Client: + 1. Receive response, extract correlation_id + 2. Find oneshot::Sender in HashMap + 3. Send response through channel + 4. Remove from HashMap +``` + +4. **Timeout Handling:** + - Use `tokio::time::timeout` wrapper + - Clean up HashMap entries on timeout + - Return clear error (not silent drop) + +5. **Topic Convention:** +``` + Request topic: {base_topic}/request + Response topic: {base_topic}/response +``` + +## Success Criteria + +- [ ] Client can send request and receive response +- [ ] Multiple concurrent requests work correctly +- [ ] Timeouts are enforced and HashMap is cleaned up +- [ ] Server handles multiple concurrent requests (spawns tasks) +- [ ] Correlation IDs are compact and collision-free +- [ ] Works with JSON payloads (serde support) +- [ ] Clean error types for timeout, connection loss, etc. +- [ ] Example showing both client and server usage +- [ ] Unit tests for correlation logic +- [ ] Integration test with real MQTT broker + +## Context from Phase 1 + +In `rust-edge-agent` Phase 1, we used NATS everywhere to avoid this correlation problem. Now we're extracting the pattern as a reusable crate that could be used with actual MQTT devices. + +Key learnings: +- Request/response correlation is non-trivial +- Concurrent handler execution is critical +- Compact IDs matter for bandwidth-constrained IoT +- This pattern is repetitive across MQTT applications + +## Constraints + +- Keep API simple and ergonomic +- Minimize allocations (IoT devices are resource-constrained) +- No unsafe code unless absolutely necessary +- Clear error messages for debugging +- Works with both MQTT v4 and v5 (rumqttc supports both) + +## Deliverables + +1. `mqtt-rpc` library crate +2. Example client and server programs +3. README explaining the problem and solution +4. API documentation with examples +5. Basic test suite + +Ready to implement `mqtt-rpc`. Where should we start - API design, correlation logic, or project structure? diff --git a/rust-toolchain.toml b/rust-toolchain.toml index a1cd21b..ee7c48f 100644 --- a/rust-toolchain.toml +++ b/rust-toolchain.toml @@ -1,4 +1,5 @@ [toolchain] +msrv = "1.84.0" # Minimum Supported Rust Version channel = "1.92.0" profile = "minimal" components = ["rustfmt", "clippy"] diff --git a/scripts/ci-lint.sh b/scripts/ci-lint.sh index 6ecb747..441561e 100755 --- a/scripts/ci-lint.sh +++ b/scripts/ci-lint.sh @@ -2,4 +2,9 @@ set -euo pipefail cargo fmt --check -cargo clippy --release --all-targets --all-features +cargo clippy --release --all-targets --all-features -- -D warnings \ + -D warnings \ + -D clippy::unwrap-used \ + -D clippy::expect_used \ + -D clippy::indexing_slicing \ + -D clippy::panic $* diff --git a/scripts/ci-qemu-smoke.sh b/scripts/ci-qemu-smoke.sh index 7786b07..98f898a 100755 --- a/scripts/ci-qemu-smoke.sh +++ b/scripts/ci-qemu-smoke.sh @@ -26,9 +26,16 @@ if [ -n "${DEBUG}" ] ; then echo "$0: === Running ARM64 binary under QEMU ===" fi chmod +x "$BIN" +set +e +OUT="$(timeout 5s qemu-aarch64 -L /usr/aarch64-linux-gnu "$BIN")" -OUT="$(qemu-aarch64 -L /usr/aarch64-linux-gnu "$BIN")" +status=$? echo "$OUT" -echo "$OUT" | grep -q "Hello world from aarch64!" +if [ "$status" != 124 ]; then + echo "$0: test failed" + exit 1 +fi +echo "$0: test passed" +exit 0 diff --git a/scripts/demo.sh b/scripts/demo.sh new file mode 100755 index 0000000..bcb8917 --- /dev/null +++ b/scripts/demo.sh @@ -0,0 +1,65 @@ +#!/usr/bin/env bash +set -euo pipefail + +: ${NUM_DEVICES:=3} +: ${DEVICE_INTERVAL:=5} + +# Detect CI vs interactive mode +if [ -n "${CI:-}" ]; then + MODE="ci" +else + MODE="interactive" +fi + +echo "$0: NUM_DEVICES : ${NUM_DEVICES}" +echo "$0: DEVICE_INTERVAL : ${DEVICE_INTERVAL}" +echo "$0: MODE : ${MODE}" + +# Start edge agent in background +./target/release/rust-edge-agent & +AGENT_PID=$! + +# Start device simulators +DEVICE_PIDS=() +for i in $(seq 1 $NUM_DEVICES); do + case $((i % 3)) in + 0) MODE_ARG="sensor"; TYPE_ARG="temp" ;; + 1) MODE_ARG="actuator"; TYPE_ARG="valve" ;; + 2) MODE_ARG="hybrid"; TYPE_ARG="propulsion" ;; + esac + printf "%s: Starting device:%2d mode:%-10s TYPE:%-12s\n" $0 $i ${MODE_ARG} ${TYPE_ARG} + ./target/release/device_sim \ + --id "device-$(printf "%03d" $i)" \ + --mode $MODE_ARG \ + --type $TYPE_ARG \ + --interval $DEVICE_INTERVAL & + DEVICE_PIDS+=($!) +done + +cleanup() { + echo "Cleaning up..." + kill $AGENT_PID 2>/dev/null || true + for pid in "${DEVICE_PIDS[@]}"; do + kill $pid 2>/dev/null || true + done +} +trap cleanup EXIT + +if [ "$MODE" = "ci" ]; then + # CI mode: run automated checks, then exit + sleep 5 # Let system stabilize + # TODO: Add automated validation (check telemetry, send command, verify response) + echo "CI validation passed" +else + # Interactive mode: show monitoring instructions + echo "Edge agent and $NUM_DEVICES devices running." + echo "" + echo "Monitor telemetry:" + echo " nats sub 'backend.telemetry'" + echo "" + echo "Send command to actuator:" + echo " nats req 'backend.command.device-002' '{\"target_value\": 75.0}'" + echo "" + echo "Press Ctrl+C to stop." + wait +fi diff --git a/scripts/service-start.sh b/scripts/service-start.sh new file mode 100755 index 0000000..b44d184 --- /dev/null +++ b/scripts/service-start.sh @@ -0,0 +1,6 @@ +#!/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 new file mode 100755 index 0000000..2d2760b --- /dev/null +++ b/scripts/service-stop.sh @@ -0,0 +1,7 @@ +#!/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/src/agent/mod.rs b/src/agent/mod.rs index ac3c725..a5aa24b 100644 --- a/src/agent/mod.rs +++ b/src/agent/mod.rs @@ -3,6 +3,8 @@ //! The agent is the primary runtime unit of the system. It owns lifecycle //! transitions and coordinates messaging and runtime components. +mod registry; mod run; +pub use registry::{DeviceRegistry, DeviceState}; pub use run::run; diff --git a/src/agent/registry.rs b/src/agent/registry.rs new file mode 100644 index 0000000..739c4fe --- /dev/null +++ b/src/agent/registry.rs @@ -0,0 +1,78 @@ +//! Device registry and state tracking. + +use crate::messaging::{DeviceType, TelemetryMessage}; +use std::collections::HashMap; +use std::time::{Duration, Instant}; + +/// Device state tracked by the edge agent. +#[derive(Debug, Clone)] +pub struct DeviceState { + // --- + #[allow(dead_code)] + pub device_id: String, + #[allow(dead_code)] + pub device_type: DeviceType, + pub last_seen: Instant, + pub last_value: Option, +} + +/// Registry of known devices and their states. +/// +/// # Behavior +/// +/// - Devices are added on first telemetry message +/// - Last-seen timestamp updated on each message +/// - Devices considered offline after configured timeout +pub struct DeviceRegistry { + // --- + devices: HashMap, + #[allow(dead_code)] + timeout: Duration, +} + +impl DeviceRegistry { + // --- + /// Create a new device registry with the specified timeout. + pub fn new(timeout: Duration) -> Self { + // --- + Self { + devices: HashMap::new(), + timeout, + } + } + + /// Update registry with telemetry message. + /// + /// Adds device if not present, updates last-seen and value. + pub fn update(&mut self, msg: &TelemetryMessage) { + // --- + self.devices + .entry(msg.device_id.clone()) + .and_modify(|state| { + state.last_seen = Instant::now(); + state.last_value = Some(msg.value); + }) + .or_insert_with(|| DeviceState { + device_id: msg.device_id.clone(), + device_type: msg.device_type, + last_seen: Instant::now(), + last_value: Some(msg.value), + }); + } + + /// Check if a device is currently online. + #[allow(dead_code)] + pub fn is_online(&self, device_id: &str) -> bool { + // --- + self.devices + .get(device_id) + .map(|state| state.last_seen.elapsed() < self.timeout) + .unwrap_or(false) + } + + /// Get current device count. + pub fn device_count(&self) -> usize { + // --- + self.devices.len() + } +} diff --git a/src/agent/run.rs b/src/agent/run.rs index 342d2b3..c8b1434 100644 --- a/src/agent/run.rs +++ b/src/agent/run.rs @@ -1,7 +1,140 @@ -use crate::LifecycleState; +use super::registry::DeviceRegistry; +use crate::messaging::{self, CommandResponse, TelemetryMessage}; // CommandRequest, DeviceType +use anyhow::{anyhow, Result}; +use async_nats::Client; +use clap::Parser; +use futures::StreamExt; +use std::time::Duration; -pub fn run() { - let _state = LifecycleState::Init; +/// Edge agent for ARM64 embedded Linux systems. +#[derive(Parser, Debug)] +#[command(version, about)] +struct Args { + // --- + /// NATS server URL + #[arg(long, env = "NATS_URL", default_value = "nats://localhost:4222")] + nats_url: String, - println!("Hello world from {}!", std::env::consts::ARCH); + /// Device offline timeout in seconds + #[arg(long, env = "DEVICE_TIMEOUT", default_value = "30")] + device_timeout: u64, +} + +/// Run the edge agent. +/// +/// # Behavior +/// +/// - Connects to NATS with retry on failure +/// - Subscribes to all device telemetry +/// - Subscribes to backend commands for routing +/// - Forwards telemetry to backend +/// - Routes commands to devices with timeout handling +pub async fn run() -> Result<()> { + // --- + let args = Args::parse(); + + eprintln!("agent:: Starting edge agent..."); + eprintln!("agent:: NATS URL: {}", args.nats_url); + eprintln!("agent:: Device timeout: {}s", args.device_timeout); + + let client = messaging::connect_with_retry(&args.nats_url).await; + let timeout = Duration::from_secs(args.device_timeout); + + let mut registry = DeviceRegistry::new(timeout); + + let mut telemetry_sub = match client.subscribe(messaging::all_device_telemetry()).await { + Ok(t) => t, + Err(e) => { + return Err(anyhow!( + "agent::run: failed to subscribe to device telemetry::{e}" + )); + } + }; + + let mut command_sub = match client.subscribe(messaging::all_backend_commands()).await { + Ok(c) => c, + Err(e) => { + return Err(anyhow!( + "agent::run: failed to subscribe to backend commands:{e}" + )); + } + }; + + eprintln!("agent:: Edge agent running"); + + loop { + tokio::select! { + Some(msg) = telemetry_sub.next() => { + handle_telemetry(&client, &mut registry, msg).await?; + } + Some(msg) = command_sub.next() => { + handle_command(&client, msg).await?; + } + }; + } +} + +async fn handle_telemetry( + client: &Client, + registry: &mut DeviceRegistry, + msg: async_nats::Message, +) -> Result<()> { + // --- + match serde_json::from_slice::(&msg.payload) { + Ok(telemetry) => { + registry.update(&telemetry); + + eprintln!( + "agent::Telemetry: {} ({:?}) = {:.2} [devices: {}]", + telemetry.device_id, + telemetry.device_type, + telemetry.value, + registry.device_count() + ); + + let payload = serde_json::to_vec(&telemetry)?; + let _ = client + .publish(messaging::backend_telemetry(), payload.into()) + .await; + Ok(()) + } + Err(e) => Err(anyhow!("agent::Telemetry: Invalid telemetry payload: {e}")), + } +} + +async fn handle_command(client: &Client, msg: async_nats::Message) -> Result<()> { + // --- + let device_id = match msg.subject.strip_prefix("backend.command.") { + Some(dev) => dev, + None => { + return Err(anyhow!( + "agent::command: Error getting command from:{msg:?}" + )); + } + }; + + eprintln!("agent:: Command for device: {}", device_id); + + let device_subject = messaging::device_command(device_id); + + match client.request(device_subject, msg.payload.clone()).await { + // --- + Ok(response) => { + if let Some(reply) = msg.reply { + let _ = client.publish(reply, response.payload).await; + } + } + Err(e) => { + eprintln!("Command failed: {}", e); + if let Some(reply) = msg.reply { + let error_response = CommandResponse { + status: "error".to_string(), + message: Some(format!("Device unreachable: {}", e)), + }; + let payload = serde_json::to_vec(&error_response)?; + let _ = client.publish(reply, payload.into()).await; + } + } + } + Ok(()) } diff --git a/src/bin/device_sim.rs b/src/bin/device_sim.rs new file mode 100644 index 0000000..e104835 --- /dev/null +++ b/src/bin/device_sim.rs @@ -0,0 +1,161 @@ +//! Device simulator for testing edge agent. +//! +//! Simulates sensors, actuators, or hybrid devices publishing telemetry +//! and responding to commands via NATS. + +use anyhow::Result; +use clap::Parser; +use futures::StreamExt; +use rust_edge_agent::{self, CommandRequest, CommandResponse, DeviceType, TelemetryMessage}; +use std::time::{Duration, SystemTime, UNIX_EPOCH}; + +#[derive(Parser, Debug)] +#[command(version, about)] +struct Args { + // --- + /// Unique device identifier + #[arg(long)] + id: String, + + /// Device mode: sensor, actuator, or hybrid + #[arg(long)] + mode: String, + + /// Device type: temp, humidity, valve, propulsion + #[arg(long, value_name = "TYPE")] + r#type: String, + + /// Telemetry interval in seconds + #[arg(long, env = "DEVICE_INTERVAL", default_value = "5")] + interval: u64, + + /// NATS server URL + #[arg(long, env = "NATS_URL", default_value = "nats://localhost:4222")] + nats_url: String, +} + +#[tokio::main] +async fn main() -> Result<()> { + // --- + let args = Args::parse(); + + let device_type = parse_device_type(&args.r#type); + let mode = args.mode.as_str(); + + eprintln!( + "Device simulator: {} (mode: {mode}, type: {:?})", + args.id, device_type + ); + + let client = rust_edge_agent::connect_with_retry(&args.nats_url).await; + + match mode { + "sensor" => run_sensor(&client, &args.id, device_type, args.interval).await, + "actuator" => run_actuator(&client, &args.id, device_type).await, + "hybrid" => { + let client_clone = client.clone(); + let id_clone = args.id.clone(); + + tokio::spawn(async move { + // --- + let _status = + run_sensor(&client_clone, &id_clone, device_type, args.interval).await; + }); + + run_actuator(&client, &args.id, device_type).await + } + _ => { + eprintln!("Invalid mode: {}", mode); + std::process::exit(1); + } + } +} + +fn parse_device_type(s: &str) -> DeviceType { + // --- + match s { + "temp" => DeviceType::Temperature, + "humidity" => DeviceType::Humidity, + "valve" => DeviceType::Valve, + "propulsion" => DeviceType::Propulsion, + _ => { + eprintln!("Invalid device type: {}", s); + std::process::exit(1); + } + } +} + +async fn run_sensor( + client: &async_nats::Client, + device_id: &str, + device_type: DeviceType, + interval: u64, +) -> Result<()> { + // --- + let subject = rust_edge_agent::device_telemetry(device_id); + let mut value = 20.0; + + loop { + let timestamp = SystemTime::now() + .duration_since(UNIX_EPOCH) + .map(|d| d.as_secs()) + .unwrap_or(0); + + value += rand::random::() * 2.0 - 1.0; + value = value.clamp(15.0, 30.0); + + let telemetry = TelemetryMessage { + device_id: device_id.to_string(), + device_type, + timestamp, + value, + }; + + let payload = serde_json::to_vec(&telemetry)?; + let _ = client.publish(subject.clone(), payload.into()).await; + + eprintln!("[{}] Published: {:.2}", device_id, value); + + tokio::time::sleep(Duration::from_secs(interval)).await; + } +} + +async fn run_actuator( + client: &async_nats::Client, + device_id: &str, + _device_type: DeviceType, +) -> anyhow::Result<()> { + // --- + let subject = rust_edge_agent::device_command(device_id); + + let mut sub = client.subscribe(subject).await?; + + eprintln!("[{}] Listening for commands", device_id); + + while let Some(msg) = sub.next().await { + // --- + match serde_json::from_slice::(&msg.payload) { + Ok(cmd) => { + eprintln!( + "[{}] Received command: target_value = {:.2}", + device_id, cmd.target_value + ); + + let response = CommandResponse { + status: "ok".to_string(), + message: Some(format!("Set to {:.2}", cmd.target_value)), + }; + + if let Some(reply) = msg.reply { + let payload = serde_json::to_vec(&response)?; + let _ = client.publish(reply, payload.into()).await; + } + } + Err(e) => { + eprintln!("[{}] Invalid command: {}", device_id, e); + tokio::time::sleep(Duration::from_secs(1_u64)).await; + } + } + } + Ok(()) +} diff --git a/src/lib.rs b/src/lib.rs new file mode 100644 index 0000000..3b1e4d7 --- /dev/null +++ b/src/lib.rs @@ -0,0 +1,17 @@ +//! Rust edge agent library. +mod agent; +mod messaging; +mod runtime; + +pub use agent::{run, DeviceRegistry, DeviceState}; +pub use messaging::{ + // --- + all_backend_commands, + all_device_telemetry, + backend_telemetry, + connect_with_retry, + device_command, + device_telemetry, +}; +pub use messaging::{CommandRequest, CommandResponse, DeviceType, TelemetryMessage}; +pub use runtime::{Lifecycle, LifecycleState}; diff --git a/src/main.rs b/src/main.rs index fcbcc8a..4cf8514 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1,11 +1,11 @@ -mod agent; -mod messaging; -mod runtime; - -use runtime::LifecycleState; +use anyhow::{anyhow, Result}; +use rust_edge_agent::run; #[tokio::main] -async fn main() { +async fn main() -> Result<()> { // --- - agent::run(); + match run().await { + Ok(_) => Ok(()), + Err(err) => Err(anyhow!("Error:{err}")), + } } diff --git a/src/messaging/mod.rs b/src/messaging/mod.rs index 94ddf71..a1fa018 100644 --- a/src/messaging/mod.rs +++ b/src/messaging/mod.rs @@ -4,5 +4,17 @@ //! messaging details from leaking into higher-level components. mod nats; +mod subjects; +mod types; -//pub use nats::{start_control, start_heartbeat}; +// Public API exports +pub use nats::connect_with_retry; +pub use subjects::{ + // --- + all_backend_commands, + all_device_telemetry, + backend_telemetry, + device_command, + device_telemetry, +}; +pub use types::{CommandRequest, CommandResponse, DeviceType, TelemetryMessage}; diff --git a/src/messaging/nats.rs b/src/messaging/nats.rs index ea6bfa2..cb1d1db 100644 --- a/src/messaging/nats.rs +++ b/src/messaging/nats.rs @@ -1,68 +1,36 @@ use async_nats::Client; -use futures::StreamExt; -use serde::{Deserialize, Serialize}; use std::time::Duration; -/// Control request payload. -#[derive(Deserialize)] -#[allow(dead_code)] -struct ControlRequest { - // --- - command: String, -} - -/// Control response payload. -#[derive(Serialize)] -#[allow(dead_code)] -struct ControlResponse { - // --- - status: &'static str, -} - -/// Start the NATS control request/reply handler. +/// Connect to NATS with exponential backoff retry. +/// +/// # Behavior +/// +/// - Retries with exponential backoff: 1s, 2s, 4s, 8s, 16s, 30s (max) +/// - Continues retrying indefinitely until connection succeeds +/// - Logs connection failures but does not crash /// -/// Listens for control commands and responds synchronously. -#[allow(dead_code)] -pub async fn start_control(client: Client) { +/// This implements resilient edge behavior for UI-less embedded systems. +pub async fn connect_with_retry(url: &str) -> Client { // --- - let mut sub = client - .subscribe("edge.control") - .await - .expect("failed to subscribe to control subject"); + let mut retry_count = 0; - tokio::spawn(async move { - // --- - while let Some(msg) = sub.next().await { - // --- - let _req: ControlRequest = - serde_json::from_slice(&msg.payload).expect("invalid control payload"); - - let resp = ControlResponse { status: "ok" }; - if let Some(reply) = msg.reply { - let _ = client - .publish(reply, serde_json::to_vec(&resp).unwrap().into()) - .await; + loop { + match async_nats::connect(url).await { + Ok(client) => { + eprintln!("Connected to NATS at {}", url); + return client; + } + Err(e) => { + let delay = if retry_count < 5 { + 1 << retry_count + } else { + 30_u64 + }; + + eprintln!("NATS connection failed: {}, retrying in {}s...", e, delay); + tokio::time::sleep(Duration::from_secs(delay)).await; + retry_count += 1; } } - }); -} - -/// Start periodic heartbeat telemetry publication. -#[allow(dead_code)] -pub async fn start_heartbeat(client: Client) { - // --- - tokio::spawn(async move { - // --- - loop { - let payload = serde_json::json!({ - "arch": std::env::consts::ARCH, - }); - - let _ = client - .publish("edge.heartbeat", payload.to_string().into()) - .await; - - tokio::time::sleep(Duration::from_secs(5)).await; - } - }); + } } diff --git a/src/messaging/subjects.rs b/src/messaging/subjects.rs new file mode 100644 index 0000000..821da98 --- /dev/null +++ b/src/messaging/subjects.rs @@ -0,0 +1,50 @@ +//! NATS subject routing and construction. + +/// Construct device telemetry subject. +#[allow(dead_code)] +pub fn device_telemetry(device_id: &str) -> String { + // --- + format!("devices.{}.telemetry", device_id) +} + +/// Construct device status subject. +#[allow(dead_code)] +pub fn device_status(device_id: &str) -> String { + // --- + format!("devices.{}.status", device_id) +} + +/// Construct device command subject (for device to subscribe). +#[allow(dead_code)] +pub fn device_command(device_id: &str) -> String { + // --- + format!("devices.{}.command", device_id) +} + +/// Construct backend command subject (for agent to subscribe). +#[allow(dead_code)] +pub fn backend_command(device_id: &str) -> String { + // --- + format!("backend.command.{}", device_id) +} + +/// Backend telemetry aggregation subject. +#[allow(dead_code)] +pub fn backend_telemetry() -> &'static str { + // --- + "backend.telemetry" +} + +/// Wildcard pattern for all device telemetry. +#[allow(dead_code)] +pub fn all_device_telemetry() -> &'static str { + // --- + "devices.*.telemetry" +} + +/// Wildcard pattern for all backend commands. +#[allow(dead_code)] +pub fn all_backend_commands() -> &'static str { + // --- + "backend.command.*" +} diff --git a/src/messaging/types.rs b/src/messaging/types.rs new file mode 100644 index 0000000..d8536be --- /dev/null +++ b/src/messaging/types.rs @@ -0,0 +1,40 @@ +//! Message type definitions for edge agent and device communication. + +use serde::{Deserialize, Serialize}; + +/// Device type classification. +#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)] +#[serde(rename_all = "lowercase")] +pub enum DeviceType { + // --- + #[serde(rename = "temp")] + Temperature, + Humidity, + Valve, + Propulsion, +} + +/// Telemetry message published by devices. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct TelemetryMessage { + // --- + pub device_id: String, + pub device_type: DeviceType, + pub timestamp: u64, + pub value: f64, +} + +/// Command request sent to actuator devices. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct CommandRequest { + // --- + pub target_value: f64, +} + +/// Command response from actuator devices. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct CommandResponse { + // --- + pub status: String, + pub message: Option, +} diff --git a/src/runtime/lifecycle.rs b/src/runtime/lifecycle.rs index 7be4bc4..bd6e9b0 100644 --- a/src/runtime/lifecycle.rs +++ b/src/runtime/lifecycle.rs @@ -17,6 +17,13 @@ pub struct Lifecycle { state: LifecycleState, } +impl Default for Lifecycle { + // --- + fn default() -> Self { + Self::new() + } +} + impl Lifecycle { /// Create a new lifecycle starting in `Init`. #[allow(dead_code)] diff --git a/src/runtime/mod.rs b/src/runtime/mod.rs index 122efb0..3422167 100644 --- a/src/runtime/mod.rs +++ b/src/runtime/mod.rs @@ -2,4 +2,4 @@ mod lifecycle; -pub use lifecycle::LifecycleState; // Lifecycle +pub use lifecycle::{Lifecycle, LifecycleState};