11 Commits
Author SHA1 Message Date
eSlider 2003ff69fb chore: ignore local tooling and prototype-3d WIP
CI & Release / Verify simulator (push) Successful in 18s
CI & Release / Trivy scan (push) Successful in 17s
CI & Release / Semantic Release (push) Successful in 6s
Keep .cursor, .grok, and prototype-3d on disk without polluting git status.
2026-07-17 12:17:06 +01:00
eSlider 2a6b50929c docs: align wire-spec and mark implemented vs proposed
CI & Release / Verify simulator (push) Successful in 18s
CI & Release / Trivy scan (push) Successful in 18s
CI & Release / Semantic Release (push) Successful in 6s
Keep the 45-byte pose frame consistent across docs, simulator, and
prototype, clarify MinIO is not peer sync, and state clearly that the
repo is a from-scratch platform sketch with options rather than mandates.
2026-07-17 12:16:39 +01:00
eSlider ed531f9fb6 Link license to readme.md
CI & Release / Verify simulator (push) Skipped
CI & Release / Trivy scan (push) Skipped
CI & Release / Semantic Release (push) Skipped
2026-07-09 15:55:04 +01:00
eSlider 336f9c0428 Add license
CI & Release / Verify simulator (push) Successful in 11s
CI & Release / Trivy scan (push) Successful in 8s
CI & Release / Semantic Release (push) Successful in 5s
2026-07-09 15:52:16 +01:00
eSlider 3d13f744e6 docs(journey): link the live visual prototype at swarm.produktor.io
CI & Release / Verify simulator (push) Skipped
CI & Release / Trivy scan (push) Skipped
CI & Release / Semantic Release (push) Skipped
2026-07-09 12:47:12 +01:00
eSlider dfb0556eda docs: add design journey, a narrative walk-through linking into the code
CI & Release / Verify simulator (push) Skipped
CI & Release / Trivy scan (push) Skipped
CI & Release / Semantic Release (push) Skipped
A chronological companion to the ADRs: how the design came together, from
the first mental model and 2D prototype through the three data tiers, the
relative-pose broadcast, the layered WireGuard/SSH trust model, SQL as the
peer contract, offload and replication, and the end-to-end proof of concept.
Each chapter links straight into the relevant source and decision records.
Flagged as the recommended first read in the README.
2026-07-09 12:30:40 +01:00
eSlider d3ccd6f94e docs: record architecture decisions (ADR-0001..0009) and operations knowledge
CI & Release / Verify simulator (push) Successful in 11s
CI & Release / Trivy scan (push) Successful in 8s
CI & Release / Semantic Release (push) Successful in 5s
Add docs/adr/ with a chronological ADR set reconstructed from the project
history (no orchestrator, one storage format, derived-state sync,
SQL-over-SSH contract, WireGuard beneath SSH, IaC boundaries, non-root
image, runtime data paths, bounded lake scans), an index README, and a
template. Distinguish repo-local ADRs from ecosystem-wide ASRs.

Expand AGENTS.md with a decision-records section and a runtime/operations
guide (service URLs, rebuild workflow, CPU budget knobs). Reference the
ADRs from README, CONTRIBUTING, and the affected design docs, and add an
open question for T1 partition pruning.
2026-07-09 11:41:05 +01:00
eSlider ea062af86a perf: bound lake scans and add platform Grafana dashboard
CI & Release / Verify simulator (push) Successful in 11s
CI & Release / Trivy scan (push) Successful in 9s
CI & Release / Semantic Release (push) Successful in 5s
Explorer and exporter were scanning the full Parquet lake every 1–5s
(~2000 files, 200MB+), driving ~2.2 CPU cores. Limit metrics to the
last N flights, cache DuckDB views and the partition tree, slow live
polls to 2s, and keep drones alive after seal to stop restart churn.

Add node-exporter, scan-duration metrics, and a Swarm Platform Grafana
dashboard for node CPU/memory and scan health.
2026-07-08 21:19:53 +01:00
eSlider bcd956d11a fix: wire live mode and Grafana to the generated data stream
CI & Release / Verify simulator (push) Successful in 11s
CI & Release / Trivy scan (push) Successful in 9s
CI & Release / Semantic Release (push) Successful in 5s
Use the real hive column name (`drone`) in the prototype's live query
and pin the Grafana datasource uid to `swarm-prom` so provisioned
dashboards resolve their Prometheus source in k3d.
2026-07-08 20:34:37 +01:00
eSlider 9f99d132e3 fix: chown var/t1 and var/t3 to simulator uid 10001 for k3d
CI & Release / Verify simulator (push) Successful in 12s
CI & Release / Trivy scan (push) Successful in 9s
CI & Release / Semantic Release (push) Successful in 4s
Non-root drone pods need write access on the hostPath lake. Ansible now
creates var/t1 and var/t3 owned by uid 10001 instead of relying on 0777
alone after legacy root-owned parquet trees.
2026-07-08 20:27:25 +01:00
eSlider 669e8ff005 fix: run simulator image as non-root user (Trivy DS-0002)
CI & Release / Verify simulator (push) Successful in 11s
CI & Release / Trivy scan (push) Successful in 9s
CI & Release / Semantic Release (push) Successful in 4s
Add swarm uid 10001 in both build stages so Trivy config scan passes
and containers do not run as root. Mounted /data volumes should be
world-writable or owned by uid 10001 (var/t1 from Ansible is 0777).
2026-07-08 20:20:19 +01:00
40 changed files with 1281 additions and 128 deletions
+5
View File
@@ -13,3 +13,8 @@ dist/
.terraform/ .terraform/
*.tfstate *.tfstate
*.tfstate.* *.tfstate.*
# Local tooling / WIP (keep on disk, never commit)
.cursor/
.grok/
prototype-3d/
+49
View File
@@ -29,3 +29,52 @@ See [`CONTRIBUTING.md`](CONTRIBUTING.md) for the development workflow and TDD re
- TypeScript: strict mode, functional components and hooks, no extra UI libraries. - TypeScript: strict mode, functional components and hooks, no extra UI libraries.
- YAML: 2-space indent. - YAML: 2-space indent.
- Diagrams: Mermaid inside Markdown. - Diagrams: Mermaid inside Markdown.
## Decision records
Significant architectural decisions are recorded chronologically as **ADRs** in
[`docs/adr/`](docs/adr/) — one immutable file per decision, in the
Status / Context / Options / Decision / Consequences format. Start from
[`docs/adr/README.md`](docs/adr/README.md).
- **ADR** = a point-in-time decision for *this* repo → `docs/adr/`.
- **ASR** = a standing, ecosystem-wide standard (Go-first, secrets, layout) →
lives in the `inventar` repo, **not** here.
When a change alters a boundary, contract, or trade-off, add an ADR (copy
[`docs/adr/TEMPLATE.md`](docs/adr/TEMPLATE.md), take the next number, update the
index) and cross-link it from the affected design doc. A later reversal gets a
new ADR that supersedes the old one — never edit an Accepted record.
## Runtime & operations (sim environment)
The k3d ground miniature is `swarm-sim`. Services are exposed as NodePorts:
| Service | URL |
| --- | --- |
| Explorer (T1 live lake) | `http://localhost:30088` |
| Warehouse explorer (T3) | `http://localhost:30089` |
| Grafana (dashboards `swarm-fleet`, `swarm-platform`) | `http://localhost:30300` |
| Prometheus | `http://localhost:30990` |
| MinIO console | `http://localhost:30901` |
| Prototype dev server (`npm run dev`) | `http://localhost:5173` |
Rebuild the simulator image and roll it out:
```bash
cd simulator && docker build -t swarm-house/simulator:dev .
k3d image import swarm-house/simulator:dev -c swarm-sim
kubectl rollout restart -n swarm statefulset/drone deployment/explorer deployment/exporter
cd infra/terraform/sim-env && terraform apply
```
**CPU budget knobs** (see [ADR-0009](docs/adr/ADR-0009-bounded-lake-scans.md)):
`METRIC_FLIGHT_WINDOW` (scan only recent flights), `SCAN_INTERVAL_S`,
`VIEW_REFRESH_S`, `TREE_CACHE_S`, `KEEP_ALIVE=1` (drones idle after seal instead
of restarting), `IMU_HZ`. Never scan the full lake on a hot path — it grows with
pod churn. Long-horizon analysis is a T3 warehouse query, not an exporter job.
Runtime data lives under `var/t1` / `var/t3`, owned by uid **10001**
([ADR-0007](docs/adr/ADR-0007-non-root-image.md),
[ADR-0008](docs/adr/ADR-0008-runtime-data-paths.md)); both are gitignored.
Terraform state is local and never committed.
+9
View File
@@ -81,9 +81,18 @@ identifying names of people, companies, or locations.
- [ ] Tests added or updated; `pytest tests/ -q` passes locally - [ ] Tests added or updated; `pytest tests/ -q` passes locally
- [ ] `docker build` in `simulator/` succeeds (test stage included) - [ ] `docker build` in `simulator/` succeeds (test stage included)
- [ ] Documentation updated if behaviour or boundaries changed - [ ] Documentation updated if behaviour or boundaries changed
- [ ] **ADR added** if the change alters an architectural boundary, contract, or trade-off (see [`docs/adr/`](docs/adr/))
- [ ] English only; no identifying references - [ ] English only; no identifying references
- [ ] Conventional commit message - [ ] Conventional commit message
## Decisions
Architectural decisions are recorded chronologically as ADRs in
[`docs/adr/`](docs/adr/). Read [`docs/adr/README.md`](docs/adr/README.md) before
changing a boundary or contract; add a new ADR (from
[`docs/adr/TEMPLATE.md`](docs/adr/TEMPLATE.md)) rather than editing an Accepted
one when you reverse course.
## Questions ## Questions
Open design questions live in [`docs/09-open-questions.md`](docs/09-open-questions.md). Open design questions live in [`docs/09-open-questions.md`](docs/09-open-questions.md).
+19
View File
@@ -0,0 +1,19 @@
PROPRIETARY LICENSE
Copyright (c) 2026 Andriy Oblivantsev. All rights reserved.
This material is proprietary and confidential.
No person or organization may, without prior written permission from the copyright holder:
- Publish or republish this material, in whole or in part.
- Copy, reproduce, distribute, transmit, forward, or share this material by any means.
- Use, modify, adapt, translate, create derivative works from, or incorporate this material into any other work.
- Make this material available to any third party, whether publicly or privately.
- Store this material in any repository, archive, or knowledge base accessible to others.
Permission must be obtained in writing before any of the above activities are undertaken.
Unauthorized use, publication, forwarding, reproduction, distribution, or disclosure of this material is strictly prohibited and may result in civil and criminal penalties under applicable copyright and intellectual property laws.
This license grants no rights except the right to view the material for its intended purpose by an authorized recipient.
+25
View File
@@ -17,6 +17,25 @@ This repository describes how to build, deliver, test, and operate that platform
7. **SQL is the contract.** Peer data access is read-only DuckDB SQL over SSH forced commands — the query language already lives on both ends, so no service, port, or protocol is invented for it. 7. **SQL is the contract.** Peer data access is read-only DuckDB SQL over SSH forced commands — the query language already lives on both ends, so no service, port, or protocol is invented for it.
8. **No debug API.** The bench and the flight use the same channel: an engineer debugging on the ground runs the identical query through the identical wrapper, permissions, and output format a peer drone would use. What you test is what flies. 8. **No debug API.** The bench and the flight use the same channel: an engineer debugging on the ground runs the identical query through the identical wrapper, permissions, and output format a peer drone would use. What you test is what flies.
None of the concrete tool picks above are mandates. This repo is a **from-scratch platform sketch**: enough structure to hire and build against, with every decision recorded so the team can replace a piece when a better fit appears.
## Implemented now vs proposed next
Honest map so a reader knows what runs today versus what is design intent.
| Area | Implemented in this repo (runnable PoC) | Proposed for a production air-gapped fleet |
| --- | --- | --- |
| On-board layout | Hive-partitioned Parquet writer, seal step, DuckDB views | Same contract; Compose services under systemd |
| Pose path | Fixed **45-byte** UDP frame + 2D bandwidth visualisation | Zenoh pub/sub (UDP kept as degraded minimal profile) |
| Peer query | Read-only SQL **gate** over HTTP explorer (keyword allow-list) | Same gate idea via **SSH forced command** + OS/engine hardening |
| Bulk sync | Visualised opportunistic transfer volume | rsync/rclone over persistent SSH between peers |
| Mesh trust | Ansible templates for WireGuard + ed25519 forced commands | Provisioned per-device keys; nothing joins at runtime |
| Ground segment | k3d/Terraform sim: lake, Grafana, offload CronJob, optional MinIO | k3s warehouse, GitOps overlays, post-flight mirror |
| CI / delivery | GitHub Actions: pytest, smoke flight, Trivy, semver release + fleet manifest artifact | Self-hosted GitLab + registry inside the air gap (same stages) |
| Docs | Problem, architecture, ADRs, design journey, open questions | Living ADRs owned by the team |
Start from the [design journey](docs/12-design-journey.md) for the story; use the table above when reviewing scope.
## Documentation ## Documentation
| Document | Contents | | Document | Contents |
@@ -33,6 +52,10 @@ This repository describes how to build, deliver, test, and operate that platform
| [09 — Open questions](docs/09-open-questions.md) | Known unknowns and proposed answers | | [09 — Open questions](docs/09-open-questions.md) | Known unknowns and proposed answers |
| [10 — Domain context](docs/10-domain-context.md) | Swarm autonomy principles this design builds on | | [10 — Domain context](docs/10-domain-context.md) | Swarm autonomy principles this design builds on |
| [11 — CI/CD & delivery](docs/11-cicd-delivery.md) | Pipelines, registry, dev containers, fleet releases | | [11 — CI/CD & delivery](docs/11-cicd-delivery.md) | Pipelines, registry, dev containers, fleet releases |
| [12 — Design journey](docs/12-design-journey.md) | **Start here** — a narrative walk-through of how the design came together, linking into the code |
| [ADRs](docs/adr/) | Architecture Decision Records — the decisions behind the above, in the order they were made |
New here? Read the [design journey](docs/12-design-journey.md) first: it tells the story chronologically and links straight into the code and decisions. The design principles below are the *what*; the [ADRs](docs/adr/) are the *why and when* — each principle traces to a dated, immutable decision record.
## Runnable parts ## Runnable parts
@@ -73,3 +96,5 @@ Contributing (TDD-first workflow): [`CONTRIBUTING.md`](CONTRIBUTING.md).
## CI & releases ## CI & releases
Every push runs the pipeline in [`.github/workflows/release.yaml`](.github/workflows/release.yaml): **pytest unit tests**, a compile check, a 10-second virtual smoke flight verified with DuckDB, and a Trivy vulnerability/misconfiguration scan. The simulator Docker build also runs `pytest` — a failing test blocks the image. Pushes to `main` then cut a semantic version from the conventional-commit type (`feat:` → minor, `!`/`BREAKING CHANGE` → major, else patch), tag the commit, and publish a release with a **fleet release manifest** (versioned artifact list with checksums — the delivery unit described in [11 — CI/CD & delivery](docs/11-cicd-delivery.md)) plus a source bundle. Every push runs the pipeline in [`.github/workflows/release.yaml`](.github/workflows/release.yaml): **pytest unit tests**, a compile check, a 10-second virtual smoke flight verified with DuckDB, and a Trivy vulnerability/misconfiguration scan. The simulator Docker build also runs `pytest` — a failing test blocks the image. Pushes to `main` then cut a semantic version from the conventional-commit type (`feat:` → minor, `!`/`BREAKING CHANGE` → major, else patch), tag the commit, and publish a release with a **fleet release manifest** (versioned artifact list with checksums — the delivery unit described in [11 — CI/CD & delivery](docs/11-cicd-delivery.md)) plus a source bundle.
[LICENSE](LICENSE)
+11 -8
View File
@@ -20,8 +20,9 @@ graph TB
subgraph serving [Layer 3 — Serving and sync] subgraph serving [Layer 3 — Serving and sync]
HOOK["event hook<br/>fires on new derived data"] HOOK["event hook<br/>fires on new derived data"]
PUB["state publisher<br/>pub/sub broadcast"] PUB["state publisher<br/>pub/sub broadcast"]
MINIO["MinIO<br/>derived datasets bucket"] BULK["bulk sync<br/>rsync over SSH"]
QAPI["query API<br/>peer data requests"] MINIO["MinIO<br/>optional derived bucket"]
QAPI["query API<br/>SQL-over-SSH"]
end end
end end
@@ -33,12 +34,14 @@ graph TB
WRITER -->|derived rows| HOOK WRITER -->|derived rows| HOOK
HOOK --> PUB HOOK --> PUB
HOOK --> MINIO HOOK --> MINIO
HOOK --> BULK
QAPI --> DUCK QAPI --> DUCK
PUB -.->|mesh| PEERS["peer drones"] PUB -.->|mesh pose| PEERS["peer drones"]
MINIO -.->|replication| PEERS BULK -.->|sealed partitions| PEERS
QAPI -.->|on demand| PEERS QAPI -.->|on demand| PEERS
``` ```
MinIO stays available where an S3 API helps (on-board derived datasets, ground warehouse). **In-flight peer bulk sync is SSH/rsync**, not object-store replication — see [04 — Swarm sync](04-swarm-sync.md).
### Layer 1 — Ingestion ### Layer 1 — Ingestion
- `sensor-ingest` subscribes to sensor sources (ROS 2 topics where available, raw drivers otherwise) and normalizes them into typed streams: IMU, barometer, temperature, LiDAR, RSSI, power, and so on. - `sensor-ingest` subscribes to sensor sources (ROS 2 topics where available, raw drivers otherwise) and normalizes them into typed streams: IMU, barometer, temperature, LiDAR, RSSI, power, and so on.
@@ -53,10 +56,10 @@ graph TB
### Layer 3 — Serving and sync ### Layer 3 — Serving and sync
- The **event hook** is the on-board "lambda": when the writer lands new *derived* rows (state, detections), it triggers registered actions — broadcast, MinIO upload, or a local mission-logic callback. Nothing polls. - The **event hook** is the on-board "lambda": when the writer lands new *derived* rows (state, detections), it triggers registered actions — broadcast, optional MinIO put, bulk-sync hint, or a local mission-logic callback. Nothing polls.
- The **state publisher** broadcasts compact position/attitude/detection payloads over the mesh pub/sub (transport analysis in [04 — Swarm sync](04-swarm-sync.md)). - The **state publisher** broadcasts compact position/attitude/detection payloads over the mesh pub/sub (UDP in the PoC; Zenoh proposed — [04](04-swarm-sync.md)).
- **Bulk sync** pulls sealed derived partitions from peers over persistent SSH (rsync delta transfer); MinIO remains optional where an S3 API is wanted. - **Bulk sync** pulls sealed derived partitions from peers over persistent SSH (rsync delta transfer). MinIO is optional where an S3 API is wanted; it is **not** the in-flight peer replication path.
- **Peer queries** are read-only DuckDB SQL over SSH forced commands — SELECT-only gate, read-only OS user, columnar responses (the design decision and its reasoning are in [04](04-swarm-sync.md)). - **Peer queries** are read-only DuckDB SQL — HTTP explorer gate in the PoC; **SSH forced commands** proposed for flight (SELECT-only gate, read-only OS user, columnar responses [04](04-swarm-sync.md)).
## Communication planes ## Communication planes
+23 -13
View File
@@ -16,19 +16,27 @@ Bandwidth is the scarcest resource in the system. Every message class gets an ex
## Broadcast payload: small on the wire, precise at rest ## Broadcast payload: small on the wire, precise at rest
The `state` broadcast is a fixed compact frame: The `state` broadcast is a fixed compact frame. The runnable PoC encodes it as
**45 bytes** little-endian (`simulator/virtual_drone/broadcast.py`); that size is
what the prototype uses for bandwidth estimates.
| Field | Type | Notes | | Field | Wire type (PoC) | Notes |
| --- | --- | --- | | --- | --- | --- |
| `drone_id` | uint16 | Fleet-scoped registry | | `magic` + `version` | 2s + uint8 | `b"SH"`, version `1` |
| `drone_id` | 8s ascii | Zero-padded; a fleet `uint16` registry id is a natural production swap |
| `ts_ns` | int64 | Epoch nanoseconds, same clock domain as storage | | `ts_ns` | int64 | Epoch nanoseconds, same clock domain as storage |
| `pos_x/y/z` | int32 | **Millimeters** in the mission frame — quantized only here, storage keeps full float precision | | `pos_x/y/z` | 3 × int32 | **Millimeters** in the **mission frame** — quantized on the wire only; storage keeps full float precision |
| `att_roll/pitch/yaw` | int16 | Centi-degrees | | `att_roll/pitch/yaw` | 3 × int16 | Centi-degrees — attitude stays on the wire so peers need no local shape model |
| `vel_x/y/z` | int16 | cm/s | | `vel_x/y/z` | 3 × int16 | cm/s |
| `frame_ref` | uint8 | Frame of reference id (GPS-denied: local/visual-odometry frames must be explicit) | | `frame_ref` | uint8 | Frame of reference id (GPS-denied: local/visual-odometry frames must be explicit) |
| `flags` | uint8 | Battery-low, returning, degraded-sensors, … | | `flags` | uint8 | Battery-low, returning, degraded-sensors, … |
~40 bytes per frame → a 50-drone swarm at 5 Hz is ~10 KB/s of pose traffic before transport overhead. Trivial even on a congested mesh. 45 bytes × 5 Hz × 50 drones ≈ 11 KB/s of pose traffic before transport overhead — still trivial on a congested mesh.
An earlier sketch used relative coordinates and a bounding sphere (no attitude).
It was dropped: mission-frame pose + attitude is simpler to fuse post-flight and
costs almost nothing at this frame size. Relative localization remains a
consumer concern when `frame_ref` differs across peers ([09](09-open-questions.md)).
`detections` events are slightly larger (class, confidence, bounding volume, ego-pose) but event-shaped and rare by comparison. `detections` events are slightly larger (class, confidence, bounding volume, ego-pose) but event-shaped and rare by comparison.
@@ -51,25 +59,27 @@ graph LR
SB -->|"store-and-forward relay"| SC SB -->|"store-and-forward relay"| SC
``` ```
### Recommended: Zenoh ### Recommended (proposal): Zenoh
- Designed exactly for constrained, dynamic networks: built-in peer discovery, brokerless peer-to-peer mode, store-and-forward, and a query layer on top of pub/sub. - Designed exactly for constrained, dynamic networks: built-in peer discovery, brokerless peer-to-peer mode, store-and-forward, and a query layer on top of pub/sub.
- First-class robotics citizenship: an official ROS 2 RMW implementation exists, so the ingestion side and the sync side can share one middleware. - First-class robotics citizenship: an official ROS 2 RMW implementation exists, so the ingestion side and the sync side can share one middleware.
- Tiny footprint, ARM64-native. - Tiny footprint, ARM64-native.
**PoC today:** the simulator and the 2D prototype exercise the **raw UDP** pose path only — the minimal degraded profile below. Zenoh is the proposed production pub/sub, not yet wired into the runnable stack.
### Alternatives considered ### Alternatives considered
| Option | Verdict | | Option | Verdict |
| --- | --- | | --- | --- |
| **DDS multicast** (ROS 2 default) | Works, battle-tested; but discovery storms and tuning pain on lossy wireless meshes are well documented. Keep as fallback since ROS 2 speaks it natively | | **DDS multicast** (ROS 2 default) | Works, battle-tested; but discovery storms and tuning pain on lossy wireless meshes are well documented. Keep as fallback since ROS 2 speaks it natively |
| **MQTT** | Needs a broker — a per-drone broker bridge is possible but adds moving parts for no gain over Zenoh | | **MQTT** | Needs a broker — a per-drone broker bridge is possible but adds moving parts for no gain over Zenoh |
| **Raw UDP multicast** | Perfect as a last-resort minimal profile for the pose broadcast alone (fixed frame, no discovery); no query layer, no reliability — documented as the degraded mode | | **Raw UDP multicast** | **Implemented in the PoC** for pose broadcast (fixed frame, no discovery); no query layer, no reliability — also the documented degraded mode |
| **MinIO bucket replication** | Wrong tool for the 5 Hz pose path, right tool for bulk derived datasets — see below | | **MinIO bucket replication** | Wrong tool for the 5 Hz pose path; optional on board for derived datasets and primary on the ground warehouse — not the in-flight bulk path |
## Two sync mechanisms, deliberately separate ## Two sync mechanisms, deliberately separate
1. **Fast path — pub/sub (Zenoh):** pose frames and detection events. Fire-and-forget with bounded staleness; consumers keep a peer-state cache. 1. **Fast path — pub/sub (Zenoh proposed; UDP in the PoC):** pose frames and detection events. Fire-and-forget with bounded staleness; consumers keep a peer-state cache.
2. **Bulk path — rsync over persistent SSH:** sealed `detections`/`state` Parquet partitions are pulled opportunistically between drones when links allow. This is how a drone that was out of range catches up on mission history without anyone re-sending events. 2. **Bulk path — rsync over persistent SSH (proposed):** sealed `detections`/`state` Parquet partitions are pulled opportunistically between drones when links allow. This is how a drone that was out of range catches up on mission history without anyone re-sending events. The PoC visualises bulk volume; it does not yet run real rsync between virtual drones.
Why SSH-based bulk sync over object-store replication: Why SSH-based bulk sync over object-store replication:
@@ -105,7 +115,7 @@ SQL access must not become a write channel. A single "read-only connection" flag
| Layer | Mechanism | What it stops | | Layer | Mechanism | What it stops |
| --- | --- | --- | | --- | --- | --- |
| 1. Key = operation | Forced command: the query key can only invoke the query wrapper, nothing else | Arbitrary exec, lateral movement | | 1. Key = operation | Forced command: the query key can only invoke the query wrapper, nothing else | Arbitrary exec, lateral movement |
| 2. Statement gate | Wrapper accepts a single statement, parses it, rejects anything but `SELECT` (no `COPY`, `ATTACH`, `INSTALL`, `SET`, multi-statements); parameters bound, not interpolated | SQL-as-a-write-channel, config tampering | | 2. Statement gate | Wrapper accepts a single statement, rejects anything but `SELECT`/`WITH`/… (no `COPY`, `ATTACH`, `INSTALL`, `SET`, multi-statements); production should prefer a real parser + bound parameters, not keywords alone | SQL-as-a-write-channel, config tampering |
| 3. OS permissions | Wrapper runs as a dedicated user with **read-only filesystem access** to the data root and write access to nothing | Any write that slips past layer 2 | | 3. OS permissions | Wrapper runs as a dedicated user with **read-only filesystem access** to the data root and write access to nothing | Any write that slips past layer 2 |
| 4. Engine hardening | `:memory:` database, external access disabled except the data-root glob, extension loading off | Reaching outside the store | | 4. Engine hardening | `:memory:` database, external access disabled except the data-root glob, extension loading off | Reaching outside the store |
| 5. Resource caps | Timeout, memory cap, niced CPU (flight software always wins), response size budget | Denial of service via expensive queries | | 5. Resource caps | Timeout, memory cap, niced CPU (flight software always wins), response size budget | Denial of service via expensive queries |
+2
View File
@@ -2,6 +2,8 @@
Zero trust, fully air-gapped, everything provisioned before wheels-up. Zero trust, fully air-gapped, everything provisioned before wheels-up.
> **Decision:** [ADR-0005 — WireGuard beneath SSH](adr/ADR-0005-wireguard-beneath-ssh.md).
## Trust model ## Trust model
- **Nothing joins the swarm at runtime.** Every device receives its identity (key pair + certificate) during ground provisioning, signed by the fleet's offline CA. There is no trust-on-first-use, no runtime enrollment endpoint, no exception path. - **Nothing joins the swarm at runtime.** Every device receives its identity (key pair + certificate) during ground provisioning, signed by the fleet's offline CA. There is no trust-on-first-use, no runtime enrollment endpoint, no exception path.
+2
View File
@@ -2,6 +2,8 @@
Three infrastructures, one data contract. The layout, schemas, and pipelines are identical everywhere; only scale and lifetime differ. Three infrastructures, one data contract. The layout, schemas, and pipelines are identical everywhere; only scale and lifetime differ.
> **Decisions:** [ADR-0006 — IaC boundaries](adr/ADR-0006-iac-boundaries.md), [ADR-0008 — runtime data paths](adr/ADR-0008-runtime-data-paths.md).
```mermaid ```mermaid
graph LR graph LR
subgraph dev [3 — Dev and simulation] subgraph dev [3 — Dev and simulation]
+10 -1
View File
@@ -2,6 +2,8 @@
The platform watches itself with the same discipline it applies to sensor data — and largely through the same pipeline. The platform watches itself with the same discipline it applies to sensor data — and largely through the same pipeline.
> **Decision:** [ADR-0009 — bound observability lake scans to a flight window and cache](adr/ADR-0009-bounded-lake-scans.md).
## Two kinds of signals, one storage ## Two kinds of signals, one storage
| Signal | Examples | Where it goes | | Signal | Examples | Where it goes |
@@ -53,6 +55,13 @@ graph LR
| Which link degraded first, and was it distance or interference? | RSSI telemetry joined with pose distance | | Which link degraded first, and was it distance or interference? | RSSI telemetry joined with pose distance |
| Is the new writer version flushing slower on ARM64? | CI simulation runs emit the same metrics; diff across fleet releases | | Is the new writer version flushing slower on ARM64? | CI simulation runs emit the same metrics; diff across fleet releases |
A working miniature of this doctrine ships in [`simulator/`](../simulator/): the `monitoring` Compose profile runs a DuckDB-based exporter over the generated Parquet plus Prometheus and a provisioned Grafana dashboard — fleet statistics derived from the data platform itself, with no agent on the (virtual) drones. A working miniature of this doctrine ships in [`simulator/`](../simulator/): the `monitoring` Compose profile runs a DuckDB-based exporter over the generated Parquet plus Prometheus and two provisioned Grafana dashboards:
| Dashboard | Focus |
| --- | --- |
| **Swarm Fleet — generated data** | Row counts, detections, pose frames, Parquet bytes, battery, RSSI |
| **Swarm Platform — CPU & scan health** | Node CPU/memory (node-exporter), exporter/explorer scan durations, flight-window gauge |
The exporter and explorer **scan only the most recent flight partitions** (`METRIC_FLIGHT_WINDOW`, default 5) and cache tree/views — otherwise CPU climbs as every pod restart appends a new `flight=` tree to the shared lake. Drones in k3d set `KEEP_ALIVE=1` after sealing so they idle instead of exiting and spawning another flight.
Alerting on the ground follows standard practice (Prometheus alert rules for infrastructure, CI gates for regression in simulated staleness/throughput budgets). In flight there is nobody to page — the platform's job is to degrade in the documented order and record everything for the post-mortem. Alerting on the ground follows standard practice (Prometheus alert rules for infrastructure, CI gates for regression in simulated staleness/throughput budgets). In flight there is nobody to page — the platform's job is to degrade in the documented order and record everything for the post-mortem.
+6
View File
@@ -61,3 +61,9 @@ Who owns the device registry (drone ids, key issuance, revocation) organizationa
Sub-250 g class units cannot run the full stack (no GPU, minimal CPU/storage). Sub-250 g class units cannot run the full stack (no GPU, minimal CPU/storage).
*Proposal:* define a minimal profile early — state publisher + UDP-multicast pose broadcast, no local Parquet store, no video — so the swarm protocol never assumes full-stack peers. *Proposal:* define a minimal profile early — state publisher + UDP-multicast pose broadcast, no local Parquet store, no video — so the swarm protocol never assumes full-stack peers.
## 11 — Lake retention and partition pruning
The observability path now scans only the most recent flights ([ADR-0009](adr/ADR-0009-bounded-lake-scans.md)), but old `flight=` partitions still accumulate on T1 disk after offload to T3.
*Proposal:* a scheduled prune of T1 partitions whose flights are confirmed present in the T3 warehouse (offload as the retention gate), configurable by age and free-space watermark. On the drone the same policy is bounded by the NVMe quota ladder ([03](03-data-platform.md)); in the sim environment it is a ground CronJob alongside the offload job.
+2
View File
@@ -2,6 +2,8 @@
How software, models, and configuration reach the fleet — reproducibly, scanned, versioned, and atomically rollback-able. Everything lives inside the air gap. How software, models, and configuration reach the fleet — reproducibly, scanned, versioned, and atomically rollback-able. Everything lives inside the air gap.
> **Decisions:** [ADR-0006 — IaC boundaries & fleet manifest](adr/ADR-0006-iac-boundaries.md), [ADR-0007 — non-root image](adr/ADR-0007-non-root-image.md).
## Source and pipelines: GitLab ## Source and pipelines: GitLab
GitLab (self-managed, on-prem) is the backbone: repositories, CI/CD, and the container registry in one system. GitLab (self-managed, on-prem) is the backbone: repositories, CI/CD, and the container registry in one system.
+199
View File
@@ -0,0 +1,199 @@
# 12 — Design journey
How this design came together, told in the order the thinking actually happened.
It is a walk-through, not a report — and **not a mandate**. The goal is a
from-scratch platform sketch that shows how the pieces fit; every concrete
choice is an option with trade-offs recorded as
[Architecture Decision Records](adr/README.md). Replace any piece if a better
fit shows up. This is the story behind the decisions, and each chapter links
straight into the code it produced.
---
## 1. Starting point
At the start there was only a rough brief and a few adjacent interests to lean
on. Key rotation matters in every pipeline: a cipher that is safe today may not
be safe next year, so the infrastructure has to leave room to react, without
spending too much latency or bandwidth to do it. That balance, security against
speed and volume, is a recurring trade-off. Geospatial data is older ground:
routing, indexing, fast spatial lookups at scale. So the first instinct was that
the value might be less about new storage and more about optimizing access to
data that already exists.
## 2. From words to a model, then to a picture
Once the shape was clear, an autonomous swarm where each unit talks to any peer
in range and [no controller spans the swarm](adr/ADR-0001-no-swarm-orchestrator.md),
the next step was to make the model visible. A small 2D prototype was enough to
check it: units patrolling, static and moving obstacles, links forming and
breaking, and a rough count of how much data crosses the air. Speed was not the
question yet, only volume, because volume is what you size the infrastructure
for. Starting from a picture gives an anchor before building anything.
> **Read the code**
> - [`prototype/src/sim.ts`](https://git.produktor.io/eSlider/swarm-house/src/branch/main/prototype/src/sim.ts#L219-L287) — the tick: links form and break by range, pose broadcasts and bulk sync accumulate bytes
> - [`prototype/src/sim.ts`](https://git.produktor.io/eSlider/swarm-house/src/branch/main/prototype/src/sim.ts#L57-L63) — the pose-frame size and rate that drive the volume estimate
The picture is running, not just described: the visual prototype is live at
[swarm.produktor.io](https://swarm.produktor.io/) (same access as this
repository). It draws units patrolling, links forming and breaking by range, and
a running estimate of how much data crosses the air.
## 3. Naming the limits
Before any code, the limits had to be named. Maximum distance between units in
the air. The fidelity data is stored at. The fidelity it can be transmitted at.
Those drive volume, and volume drives the whole shape of the system. The
conclusion: constrain the transport, but do not constrain on-board capture. Keep
everything locally. Send only what peers truly need.
> **Read more**
> - [01 — Problem statement](01-problem-statement.md) and [03 — Data platform](03-data-platform.md) — the constraints and the storage tiers
## 4. Three tiers inside each unit
That gave [three layers](03-data-platform.md) on every drone:
1. **Raw capture**, never transformed.
2. **ETL / reduction**, where data is cut down to what has to be shared:
millimetres to centimetres, thinning time series where that is enough,
dropping what nobody downstream reads.
3. **A bidirectional interface**: on one side it broadcasts position into the
shared channel, on the other it lets peers pull data.
The ETL is deliberately not over-specified. How sensor data is transformed is a
data-engineering decision. The platform's job is the contract, meaning what gets
stored and in which structure, not to step into that work.
> **Read the code**
> - [`simulator/virtual_drone/writer.py`](https://git.produktor.io/eSlider/swarm-house/src/branch/main/simulator/virtual_drone/writer.py#L23-L80) — the Hive-partition layout and the seal step that compacts a flight
> - [`simulator/virtual_drone/sensors.py`](https://git.produktor.io/eSlider/swarm-house/src/branch/main/simulator/virtual_drone/sensors.py) — the raw capture that feeds tier one
## 5. What actually needs to sync
The most time-critical item is where each peer is, so every unit has time to
react. The wire frame carries **mission-frame position** (millimetres on the
wire, full float at rest), **attitude**, **velocity**, a `frame_ref`, and flags —
45 bytes at 5 Hz. That is cheap enough that shrinking further is not worth the
fusion pain. An earlier sketch used relative coordinates and a bounding sphere
with no orientation; it was dropped. When peers disagree on frames,
`frame_ref` makes the mismatch explicit for consumers
([09 — Open questions](09-open-questions.md)).
> **Read the code**
> - [`simulator/virtual_drone/broadcast.py`](https://git.produktor.io/eSlider/swarm-house/src/branch/main/simulator/virtual_drone/broadcast.py#L1-L45) — the 45-byte pose frame, position quantized on the wire only
> - [`prototype/src/sim.ts`](https://git.produktor.io/eSlider/swarm-house/src/branch/main/prototype/src/sim.ts#L60-L63) — `POSE_BYTES = 45` drives the volume estimate
> - [`prototype/src/sim.ts`](https://git.produktor.io/eSlider/swarm-house/src/branch/main/prototype/src/sim.ts#L260-L277) — pose broadcasts at 5 Hz and opportunistic bulk sync, made visible
## 6. Security as nested layers, not one wall
The assumption was an ad-hoc wireless mesh. Broadcasting positions in the clear,
even over encrypted Wi-Fi, is not acceptable: anyone who joins can reconstruct
trajectories over time, and wireless encryption has been broken by ordinary
people for two decades. Relying on that layer alone is risky, and pushing heavy
encryption there is wasteful. So the trust model is [layered](05-network-security.md):
1. Wireless transport (WPA class) as the outer shell.
2. **WireGuard over UDP** as the mesh overlay: proven, peer-authenticated by
static keys, cheap.
3. **SSH inside the tunnel** as the application layer, instead of HTTP/HTTPS. SSH
is one of the most tested remote-access technologies we have. HTTPS would pull
in a certificate authority and extra parts, and it is the first surface a
security researcher probes. [Fewer moving parts, on purpose](adr/ADR-0005-wireguard-beneath-ssh.md).
Identity and encryption solve together at provisioning time: a per-flight,
per-unit key pair recorded in a flight registry. Each unit holds its own key pair
plus the public keys of its peers. A separate key registry for SSH means that if
one WireGuard key leaks, SSH access does not follow automatically.
> **Read the code**
> - [`infra/ansible/drone-provision.yml`](https://git.produktor.io/eSlider/swarm-house/src/branch/main/infra/ansible/drone-provision.yml#L20-L70) — ed25519 identity, peers authorized only through a forced command, WireGuard brought up
> - [05 — Network & security](05-network-security.md) and [ADR-0005](adr/ADR-0005-wireguard-beneath-ssh.md)
## 7. SQL as the contract
Beyond the broadcast, people and peers want richer access. The whole progression
is familiar: fully dynamic APIs, then JSON, OpenAPI/Swagger, GraphQL, then heavier
schemes. Each time the same pain: hard to debug in the moment, and a lot of
infrastructure to support, test, and agree on. The simpler conclusion:
[**SQL is the contract**](adr/ADR-0004-sql-over-ssh-contract.md). Everyone who
touches data already speaks SQL, and we do not know in advance which slice each
of them needs, so give them a declarative, extensible way in rather than guessing
endpoints ahead of time.
Concretely: an SSH connection whose forced command runs a read-only **DuckDB**
query directly. No extra auth handshake between the engine and the data, the way
a long-running database needs. No new port or protocol. When access control is
needed, it rides on Linux file permissions, because the data is discrete
Hive-partitioned files and permissions can be set per file. Not new invention:
coming back to foundations DevOps and sysadmins already trust.
> **Read the code**
> - [`simulator/explorer/server.py`](https://git.produktor.io/eSlider/swarm-house/src/branch/main/simulator/explorer/server.py#L43-L53) — the read-only gate that rejects anything that writes
> - [`simulator/explorer/server.py`](https://git.produktor.io/eSlider/swarm-house/src/branch/main/simulator/explorer/server.py#L111-L131) — the query path a peer or an engineer hits identically
> - [`simulator/tests/test_sql_gate.py`](https://git.produktor.io/eSlider/swarm-house/src/branch/main/simulator/tests/test_sql_gate.py#L13-L35) — the gate is safety-critical, so it is tested first
## 8. Replication and offload
Moving accumulated data uses rsync-class tooling (rclone) over the same SSH
channel. It transfers diffs of the Hive tree efficiently. This matters because a
unit can be lost, and if it is, we want its data to already exist elsewhere.
When peers have a stable link the partition tree is pulled between them, which
gives redundancy and a way to cross-check later. On the ground, when the fleet
returns, the same path unifies every unit's data into a
[local warehouse](06-environments.md). The warehouse is read-mostly for
analytics, so transaction contention is not a concern, which opens a
DuckDB / DuckLake approach with MinIO underneath, the
[same storage model on both ends](adr/ADR-0002-one-storage-format.md). One
uniform structure is what lets an analyst trust a single picture even when a
unit's data has a gap.
> **Read the code**
> - [`infra/terraform/ground/main.tf`](https://git.produktor.io/eSlider/swarm-house/src/branch/main/infra/terraform/ground/main.tf#L124-L160) — the post-flight T1 to T3 offload CronJob
> - [06 — Environments](06-environments.md) and [ADR-0002](adr/ADR-0002-one-storage-format.md)
## 9. A proof of concept, not slides
A concept can be handed over and be wrong. Instead there is a generator: a
virtual drone producing real sensor data into the Hive layout, a small frontend
to watch it fill, and a **Grafana** dashboard on top, not to get ahead of the
analysts, but to see the infrastructure work end to end. Closing the loop
(Terraform, connector, exporter, dashboard) is how a misconfiguration shows up
early. A panel that renders wrong points at exactly which layer broke, before an
analyst has to report that the infra is up but nothing shows. Better to exercise
every layer now and drop any layer or technology later if there is a sound
argument against it.
Watching that loop is also where the [CPU cost surfaced](adr/ADR-0009-bounded-lake-scans.md):
scanning the whole lake on every query does not scale, so the scan is
[bounded to a flight window and cached](07-observability.md).
> **Read the code**
> - [`simulator/monitoring/exporter.py`](https://git.produktor.io/eSlider/swarm-house/src/branch/main/simulator/monitoring/exporter.py#L38-L107) — the DuckDB scan that turns Parquet into Prometheus metrics
> - [`simulator/monitoring/lake.py`](https://git.produktor.io/eSlider/swarm-house/src/branch/main/simulator/monitoring/lake.py#L13-L33) — the flight-window scan that keeps CPU bounded as the lake grows
> - [`simulator/monitoring/grafana/dashboards/`](https://git.produktor.io/eSlider/swarm-house/src/branch/main/simulator/monitoring/grafana/dashboards) — the fleet and platform dashboards
> - [`.github/workflows/release.yaml`](https://git.produktor.io/eSlider/swarm-house/src/branch/main/.github/workflows/release.yaml) — tests, smoke flight, Trivy, semantic release
## 10. Working method
A note on how this is meant to be built, not just what. Working together needs a
shared language, and when the working language is nobody's first language,
answering live is expensive. The form that scales is written-first: lay the idea
down carefully, let others read and answer in their own time, keep the history. A
new colleague, or an agent, can then load the full context of why a decision was
made by reading the discussion. That is exactly why the decisions here are
[ADRs](adr/README.md) and why the [contribution flow](../CONTRIBUTING.md) is
test-first: the reasoning is captured where the next person will look for it.
The mental image is a city planner. Draw perfectly straight roads and people
still cross the grass where it is faster. The job is not to decide for them, but
to understand the real paths, the pain points and limits of each team, and pave
those. So constraints get checked with the other DevOps, data engineers, and
analysts before infrastructure is laid, instead of building roads nobody uses.
---
The platform is air-gapped and on-prem by design: nothing leaves the swarm, and
nothing joins it at runtime. Start from the [README](../README.md) for the map,
or the [ADR index](adr/README.md) for the decisions in full.
@@ -0,0 +1,38 @@
# ADR-0001: No swarm-wide orchestrator; coordinate through data
## Status
Accepted (2026-07-08)
## Context
The fleet flies fully autonomously with no internet uplink and only
intermittent mesh connectivity between drones. A conventional control plane
(Kubernetes across the swarm, a leader election, a central scheduler) assumes
stable membership and low-latency links — exactly what a swarm does not have.
Partitions are the normal case, not the exception.
## Options
| Option | Pros | Cons |
| --- | --- | --- |
| A — Cluster orchestrator spanning the swarm | Familiar tooling | Assumes stable membership; a partition stalls coordination; single point of failure in the air |
| B — Each drone autonomous, coordinating through exchanged data | Survives partitions; no air-side control plane to fail | No global view; behaviour must be derivable from local state |
## Decision
Every drone is an autonomous node. There is **no orchestrator spanning the
swarm**. Coordination happens by exchanging derived state (see
[ADR-0003](ADR-0003-sync-derived-state-only.md)), and each drone makes its own
in-flight decisions against its local data window. Orchestration tooling
(Kubernetes, Flux) is confined to the **ground** segment
([ADR-0006](ADR-0006-iac-boundaries.md)).
## Consequences
- The platform degrades gracefully under partition: a drone that loses peers
keeps flying and recording.
- There is no single "fleet state" to query in the air; health questions are
answered from each drone's own store and reconciled on the ground.
- Design principle 1 in [`../../README.md`](../../README.md) and the on-board
model in [`../02-architecture.md`](../02-architecture.md) follow from this.
+32
View File
@@ -0,0 +1,32 @@
# ADR-0002: One storage format on every tier — Parquet + DuckDB, Hive layout
## Status
Accepted (2026-07-08)
## Context
Data lives in three places: on the drone's NVMe in flight (T1), in the ground
warehouse after landing (T3), and in the dev/simulation environment. If each
tier used a different format or layout, offload would need transformation code,
schemas would drift, and a query proven on the bench would not run unchanged on
flight data.
## Decision
Use **one storage format on every floor**: columnar **Parquet** files in an
identical **Hive-partitioned** layout (`dataset=…/flight=…/drone=…/…`), queried
with **DuckDB**. Flight offload T1 → T3 is therefore a plain **mirror**
(copy/rsync of partition directories), not an ETL step.
## Consequences
- Offload is a file operation, testable with a checksum; no format converters
to maintain or version.
- The same SQL runs against T1, T3, and simulated data — the simulator
([`../../simulator/`](../../simulator/)) generates the *real* layout, so tests
exercise the production format.
- Partitioning choices are load-bearing; changing them is a schema migration.
Detailed in [`../03-data-platform.md`](../03-data-platform.md).
- DuckDB's single-file, no-server model fits an air-gapped drone with no room
for a database daemon.
@@ -0,0 +1,33 @@
# ADR-0003: Sync derived state only; raw telemetry stays local
## Status
Accepted (2026-07-08)
## Context
Each drone produces high-rate sensor telemetry and on-board video detections —
far more than the mesh can carry. But peers only need enough to coordinate
flight: where everyone is, where they are heading, what was detected. Pushing
raw feeds across the air would saturate the radio and drain batteries for data
nobody reads in flight.
## Decision
Only **derived state** crosses the air — position, attitude, detections — as
small pub/sub broadcasts (~5 Hz). **Raw telemetry stays on local NVMe** until
the drone lands, then travels to the warehouse during offload
([ADR-0002](ADR-0002-one-storage-format.md)). The design is **event-driven**:
new derived data triggers downstream action through hooks; nothing polls asking
"anything new yet?".
## Consequences
- Air bandwidth scales with the number of peers and the broadcast rate, not
with sensor resolution.
- The warehouse is the only place with the full picture; in-flight decisions use
the local window plus peer broadcasts.
- Broadcast staleness per peer becomes a first-class health signal
([`../07-observability.md`](../07-observability.md)).
- What syncs and what does not is specified in
[`../04-swarm-sync.md`](../04-swarm-sync.md).
@@ -0,0 +1,46 @@
# ADR-0004: Read-only SQL-over-SSH as the peer query contract; no debug API
## Status
Accepted (2026-07-08). Supersedes the earlier sketch of a bespoke query service.
## Context
Beyond the 5 Hz broadcast, a peer (or an engineer on the bench) sometimes needs
to *ask* a drone a richer question: "show me your detections in the last minute
near this position". Inventing a service for this means a new port, a new
protocol, an auth layer, and a second code path that only runs during
debugging — the classic "works on the bench, untested in flight" trap.
## Options
| Option | Pros | Cons |
| --- | --- | --- |
| A — Custom query microservice + REST API | Flexible | New port/protocol/auth to secure; a debug-only path that never flies |
| B — Read-only DuckDB SQL over an SSH forced command | Reuses SSH trust + the query language already on both ends; identical path on bench and in flight | SQL must be gated to read-only |
## Decision
Peer data access is **read-only DuckDB SQL executed through an SSH forced
command**. SQL is already the query language on both ends
([ADR-0002](ADR-0002-one-storage-format.md)), and SSH already carries the trust
model ([ADR-0005](ADR-0005-wireguard-beneath-ssh.md)), so no new service, port,
or protocol is invented. A **SQL gate** rejects anything that writes
(`COPY`, `INSERT`, `ATTACH`, pragmas). There is **no separate debug API**: an
engineer debugging on the ground runs the identical query through the identical
wrapper, permissions, and output format a peer drone would use.
## Consequences
- One access path is designed for flight and bench — "what you test is what
flies".
- The statement gate in the PoC (`simulator/explorer/server.py`, covered by
`simulator/tests/test_sql_gate.py`) is a **first layer**: keyword allow/deny
over HTTP for the explorer. Production still needs the remaining layers in
[04 — Swarm sync](../04-swarm-sync.md) (forced-command key, read-only OS user,
engine hardening, resource caps) and a real SQL parser rather than keywords
alone.
- The explorer and the prototype's live mode are *stand-in consumers* of the
same read-only contract over HTTP today; the proposed flight path is
SQL-over-SSH ([`../04-swarm-sync.md`](../04-swarm-sync.md)).
- Design principles 7 and 8 in [`../../README.md`](../../README.md) restate this.
@@ -0,0 +1,40 @@
# ADR-0005: WireGuard beneath SSH for the mesh transport
## Status
Accepted (2026-07-08)
## Context
Peer bulk sync and SQL-over-SSH ([ADR-0004](ADR-0004-sql-over-ssh-contract.md))
need an encrypted, authenticated channel across an ad-hoc radio mesh. SSH alone
would work at the application layer, but it leaves the drones' listening
surface (sshd, any other service port) exposed on the raw mesh network, and it
gives no cheap, uniform way to authenticate *the peer* before the application
handshake.
## Options
| Option | Pros | Cons |
| --- | --- | --- |
| A — SSH only | One daemon, familiar | Service ports exposed on the raw mesh; per-service auth; no network-layer peer identity |
| B — WireGuard tunnel, SSH inside it | Network-layer peer auth via static keys; only the WG port is exposed; SSH speaks over a private, encrypted subnet | Two layers to provision |
## Decision
Run **WireGuard as the network layer** and **SSH on top of it**. WireGuard
authenticates peers by pre-shared static public keys and presents a private
encrypted subnet; SSH forced commands then carry the SQL and bulk-sync contract
over that subnet. Keys are issued **per device before deployment** — nothing
joins the mesh at runtime (design principle 5).
## Consequences
- The only thing exposed on the raw radio is the WireGuard port; sshd and the
data services listen only on the WG interface.
- Peer identity is established once, at the network layer, and reused by every
application above it.
- Provisioning must place both WG and SSH keys during host prep — handled by
`infra/ansible/drone-provision.yml` ([ADR-0006](ADR-0006-iac-boundaries.md)).
- Rationale and the "why not SSH alone" argument live in
[`../05-network-security.md`](../05-network-security.md).
+41
View File
@@ -0,0 +1,41 @@
# ADR-0006: IaC split — Ansible hosts, Terraform ground, Flux ground-only, fleet manifest for drones
## Status
Accepted (2026-07-08)
## Context
Three very different things need provisioning: (1) drone hosts, which are
air-gapped and never reachable by a controller in flight; (2) the ground
segment (warehouse, observability, offload), a normal k3s cluster; (3) the
simulation environment that mirrors the ground segment locally. Using one tool
for all three would force a GitOps controller or a Terraform apply loop onto the
drone — impossible for an autonomous, disconnected node
([ADR-0001](ADR-0001-no-swarm-orchestrator.md)).
## Decision
Split infrastructure ownership by lifecycle:
| Layer | Tool | Scope |
| --- | --- | --- |
| Host preparation | **Ansible** | Drone bench prep (keys, forced commands, WireGuard, Compose bundle) and the local k3d sim cluster |
| Ground workloads | **Terraform** | k3s/k3d resources: virtual fleet + observability (`sim-env`), warehouse + T1→T3 offload (`ground`) |
| Ground continuous delivery | **Flux (GitOps)** | Ground segment **only** — policy labels, dashboard bundles |
| Drone delivery | **Fleet release manifest** | A versioned, digest-pinned artifact list; no controller pulls to the drone |
Drones are delivered by the **fleet release manifest** (one version for the
whole fleet, atomic rollback), never by a controller reaching into the air.
## Consequences
- No GitOps agent or `terraform apply` ever targets a drone; the air side stays
controller-free.
- Terraform state is environment-local and kept out of Git
([ADR-0008](ADR-0008-runtime-data-paths.md)).
- The sim cluster is a faithful miniature: the same Terraform describes it and
the ground segment.
- Boundaries are documented in
[`../06-environments.md`](../06-environments.md#iac-boundaries); delivery in
[`../11-cicd-delivery.md`](../11-cicd-delivery.md).
+27
View File
@@ -0,0 +1,27 @@
# ADR-0007: Run the simulator image as a non-root user
## Status
Accepted (2026-07-08)
## Context
The Trivy scan in CI flagged **DS-0002**: the simulator image ran as `root`
because no `USER` was set. A container writing the data lake as root also
produces root-owned Parquet on bind mounts, which then blocks a later non-root
process from the same files.
## Decision
Create a dedicated unprivileged user **`swarm` (uid/gid 10001)** in the
`simulator/Dockerfile` and switch to it with `USER swarm` before the entrypoint.
The uid is fixed (not auto-assigned) so host-side ownership of the shared data
paths ([ADR-0008](ADR-0008-runtime-data-paths.md)) is deterministic.
## Consequences
- Trivy DS-0002 clears; the image follows least-privilege.
- Bind-mounted data is owned by `10001:10001`; host prep (`sim-cluster.yml`) and
Ansible `chown` the `var/t1` / `var/t3` paths to this uid so k3d pods can write.
- Any future service image reusing this data must run as the same uid or share
the group.
+40
View File
@@ -0,0 +1,40 @@
# ADR-0008: Shared `var/t1` / `var/t3` runtime paths; state out of Git
## Status
Accepted (2026-07-08)
## Context
Compose, the k3d fleet, and the ground offload job all need to read and write
the same data lake, so that "go live" in the prototype shows the data the
simulator just generated. Early on, Terraform state and the provider cache were
accidentally committed (~6 MB of `*.tfstate` and a vendored provider binary),
polluting history and risking stale/secret data in Git.
## Decision
Standardise two **gitignored** runtime directories at the repo root:
| Path | Tier | Contents |
| --- | --- | --- |
| `var/t1/` | T1 live lake | In-flight Parquet from Compose or the k3d fleet |
| `var/t3/` | T3 warehouse | Offloaded historical partitions |
They are created automatically and overridable via
`SWARM_T1_DIR` / `SWARM_WAREHOUSE_DATA` / `SWARM_SIM_DATA`. Treat `var/` as
**runtime-only** per the ecosystem canonical layout. `.gitignore` excludes
`*.tfstate`, `*.tfstate.backup`, and `.terraform/` while **keeping**
`.terraform.lock.hcl` (the lock is source, the state is not).
## Consequences
- One lake feeds every component; the live prototype reflects real generated
data.
- Git history carries no machine state or large binaries; `terraform` treats
its state as local/remote-backend concern, not versioned source.
- Ownership of these paths is fixed to uid 10001
([ADR-0007](ADR-0007-non-root-image.md)).
- Recovery from earlier state drift is documented in
[`../../infra/terraform/README.md`](../../infra/terraform/README.md)
(`terraform import`).
+52
View File
@@ -0,0 +1,52 @@
# ADR-0009: Bound observability lake scans to a flight window and cache
## Status
Accepted (2026-07-08)
## Context
The k3d sim node was running at ~245% CPU. Three causes compounded:
1. **Explorer** rebuilt a DuckDB view over the *entire* lake on every HTTP
request (~2000 Parquet files, 200 MB+) — roughly 1.5 cores.
2. **Exporter** scanned the full lake every 5 s — roughly 0.7 cores.
3. **Drones** exited after sealing a flight; Kubernetes restarted each pod,
and every restart appended a fresh `flight=` partition tree, so the lake grew
without bound and every scan got more expensive.
The cost is inherent to scanning an append-only lake whose size grows with pod
churn, not to any single slow query. Metrics and the live view only ever need
recent flights.
## Decision
Bound and cache the scans instead of reading the whole lake:
- **Flight window.** The exporter and explorer scan only the most recent
`METRIC_FLIGHT_WINDOW` flights (default 5) via a shared helper
(`simulator/monitoring/lake.py`). Older partitions stay on disk but out of the
hot path.
- **Caching.** Explorer caches DuckDB views (`VIEW_REFRESH_S=30`) and the
partition tree (`TREE_CACHE_S=15`); the exporter scans every
`SCAN_INTERVAL_S=30` instead of 5 s.
- **Stop restart churn.** Drones set `KEEP_ALIVE=1` and idle after sealing
instead of exiting, so no new flight tree is created per restart. In k3d,
`IMU_HZ=20` and resource limits cap per-drone load; the live prototype polls
every 2 s.
- **Self-observability.** Add node-exporter, a
`swarm_exporter_scan_duration_seconds` metric, and a **Swarm Platform**
Grafana dashboard (node CPU/memory + scan health) alongside the existing fleet
dashboard.
## Consequences
- Steady state dropped to explorer ~29 m, exporter ~16 m CPU; the k3d node to
~56% (from ~245%).
- Metrics reflect recent flights only; long-horizon analysis is a warehouse
(T3) query, not an exporter concern — consistent with
[ADR-0003](ADR-0003-sync-derived-state-only.md).
- Old flight partitions accumulate on disk; pruning them is a follow-up
(tracked in [`../09-open-questions.md`](../09-open-questions.md)).
- Observability design is documented in
[`../07-observability.md`](../07-observability.md).
+40
View File
@@ -0,0 +1,40 @@
# Architecture Decision Records
This directory records the significant decisions behind Swarm House, in the
order they were made. Each record is immutable once **Accepted** — a later
decision that changes course gets a new number and supersedes the old one
rather than editing it.
## ADR vs ASR
| Kind | What it captures | Where it lives |
| --- | --- | --- |
| **ADR** (Architecture Decision Record) | A point-in-time decision for *this* repository, with the options considered and their consequences | Here, in `docs/adr/` |
| **ASR** (Architecture Standard Record) | A standing, cross-repository policy every project must obey (Go-first, secrets, layout) | The ecosystem `inventar` repo, not here |
The rules in [`../../AGENTS.md`](../../AGENTS.md) are this repo's local standing
constraints; anything ecosystem-wide belongs in `inventar` as an ASR. Design
decisions specific to the platform are ADRs and belong here.
## Index
| ADR | Decision | Status |
| --- | --- | --- |
| [ADR-0001](ADR-0001-no-swarm-orchestrator.md) | No swarm-wide orchestrator; coordinate through data | Accepted |
| [ADR-0002](ADR-0002-one-storage-format.md) | One storage format on every tier — Parquet + DuckDB, Hive layout | Accepted |
| [ADR-0003](ADR-0003-sync-derived-state-only.md) | Sync derived state only; raw telemetry stays local | Accepted |
| [ADR-0004](ADR-0004-sql-over-ssh-contract.md) | Read-only SQL-over-SSH as the peer query contract; no debug API | Accepted |
| [ADR-0005](ADR-0005-wireguard-beneath-ssh.md) | WireGuard beneath SSH for the mesh transport | Accepted |
| [ADR-0006](ADR-0006-iac-boundaries.md) | IaC split: Ansible hosts, Terraform ground, Flux ground-only, fleet manifest for drones | Accepted |
| [ADR-0007](ADR-0007-non-root-image.md) | Run the simulator image as a non-root user | Accepted |
| [ADR-0008](ADR-0008-runtime-data-paths.md) | Shared `var/t1` / `var/t3` runtime paths; state out of Git | Accepted |
| [ADR-0009](ADR-0009-bounded-lake-scans.md) | Bound observability lake scans to a flight window and cache | Accepted |
## Writing a new ADR
1. Copy [`TEMPLATE.md`](TEMPLATE.md) to `ADR-NNNN-short-slug.md` (next free number).
2. Fill in Context, Options (if more than one was real), Decision, Consequences.
3. Add a row to the index above.
4. Cross-link the ADR from the design doc it affects (and vice versa).
Keep the tone laconic: state the decision and why, skip the narrative.
+29
View File
@@ -0,0 +1,29 @@
# ADR-NNNN: Short decision title
## Status
Proposed | Accepted | Superseded by ADR-XXXX
## Context
What forces are in play? What problem or constraint triggered the decision?
Keep it to the facts that actually shaped the choice.
## Options
| Option | Pros | Cons |
| --- | --- | --- |
| A — … | … | … |
| B — … | … | … |
(Drop this section when only one option was ever real.)
## Decision
The choice, stated plainly. One paragraph.
## Consequences
- What becomes easier.
- What becomes harder or is now forbidden.
- Follow-ups, links to affected docs.
+4
View File
@@ -29,12 +29,16 @@
path: "{{ data_dir }}" path: "{{ data_dir }}"
state: directory state: directory
mode: "0777" mode: "0777"
owner: "10001"
group: "10001"
- name: Create warehouse directory (T3 host path inside the k3d node) - name: Create warehouse directory (T3 host path inside the k3d node)
ansible.builtin.file: ansible.builtin.file:
path: "{{ warehouse_dir }}" path: "{{ warehouse_dir }}"
state: directory state: directory
mode: "0777" mode: "0777"
owner: "10001"
group: "10001"
- name: List existing clusters - name: List existing clusters
ansible.builtin.command: k3d cluster list -o json ansible.builtin.command: k3d cluster list -o json
+131 -1
View File
@@ -63,6 +63,14 @@ resource "kubernetes_stateful_set" "drone" {
name = "DATA_DIR" name = "DATA_DIR"
value = "/data" value = "/data"
} }
env {
name = "KEEP_ALIVE"
value = "1"
}
env {
name = "IMU_HZ"
value = "20"
}
volume_mount { volume_mount {
name = "lake" name = "lake"
@@ -118,9 +126,21 @@ resource "kubernetes_deployment" "exporter" {
name = "EXPORTER_PORT" name = "EXPORTER_PORT"
value = "9105" value = "9105"
} }
env {
name = "SCAN_INTERVAL_S"
value = "30"
}
env {
name = "METRIC_FLIGHT_WINDOW"
value = "5"
}
port { port {
container_port = 9105 container_port = 9105
} }
resources {
requests = { cpu = "50m", memory = "64Mi" }
limits = { cpu = "500m", memory = "256Mi" }
}
volume_mount { volume_mount {
name = "lake" name = "lake"
mount_path = "/data" mount_path = "/data"
@@ -182,9 +202,25 @@ resource "kubernetes_deployment" "explorer" {
name = "EXPLORER_PORT" name = "EXPLORER_PORT"
value = "8088" value = "8088"
} }
env {
name = "VIEW_REFRESH_S"
value = "30"
}
env {
name = "TREE_CACHE_S"
value = "15"
}
env {
name = "METRIC_FLIGHT_WINDOW"
value = "5"
}
port { port {
container_port = 8088 container_port = 8088
} }
resources {
requests = { cpu = "50m", memory = "64Mi" }
limits = { cpu = "750m", memory = "256Mi" }
}
volume_mount { volume_mount {
name = "lake" name = "lake"
mount_path = "/data" mount_path = "/data"
@@ -226,11 +262,17 @@ resource "kubernetes_config_map" "prometheus" {
data = { data = {
"prometheus.yml" = <<-EOT "prometheus.yml" = <<-EOT
global: global:
scrape_interval: 5s scrape_interval: 15s
scrape_configs: scrape_configs:
- job_name: swarm - job_name: swarm
static_configs: static_configs:
- targets: ["exporter:9105"] - targets: ["exporter:9105"]
- job_name: explorer
static_configs:
- targets: ["explorer:8088"]
- job_name: node
static_configs:
- targets: ["node-exporter:9100"]
EOT EOT
} }
} }
@@ -298,8 +340,12 @@ resource "kubernetes_config_map" "grafana_provisioning" {
data = { data = {
"datasource.yml" = <<-EOT "datasource.yml" = <<-EOT
apiVersion: 1 apiVersion: 1
deleteDatasources:
- name: Prometheus
orgId: 1
datasources: datasources:
- name: Prometheus - name: Prometheus
uid: swarm-prom
type: prometheus type: prometheus
access: proxy access: proxy
url: http://prometheus:9090 url: http://prometheus:9090
@@ -324,6 +370,90 @@ resource "kubernetes_config_map" "grafana_dashboard" {
} }
data = { data = {
"swarm.json" = file("${path.module}/../../../simulator/monitoring/grafana/dashboards/swarm.json") "swarm.json" = file("${path.module}/../../../simulator/monitoring/grafana/dashboards/swarm.json")
"swarm-platform.json" = file("${path.module}/../../../simulator/monitoring/grafana/dashboards/swarm-platform.json")
}
}
# Host metrics for the k3d node (CPU / memory on the platform dashboard)
resource "kubernetes_deployment" "node_exporter" {
metadata {
name = "node-exporter"
namespace = kubernetes_namespace.swarm.metadata[0].name
labels = local.labels
}
spec {
replicas = 1
selector {
match_labels = { app = "node-exporter" }
}
template {
metadata {
labels = merge(local.labels, { app = "node-exporter" })
}
spec {
host_network = true
host_pid = true
container {
name = "node-exporter"
image = "prom/node-exporter:v1.8.2"
args = [
"--path.procfs=/host/proc",
"--path.sysfs=/host/sys",
"--path.rootfs=/host/root",
"--collector.filesystem.mount-points-exclude=^/(sys|proc|dev|host|etc)($$|/)",
]
port {
container_port = 9100
}
volume_mount {
name = "proc"
mount_path = "/host/proc"
read_only = true
}
volume_mount {
name = "sys"
mount_path = "/host/sys"
read_only = true
}
volume_mount {
name = "root"
mount_path = "/host/root"
read_only = true
}
resources {
requests = { cpu = "20m", memory = "32Mi" }
limits = { cpu = "200m", memory = "64Mi" }
}
}
volume {
name = "proc"
host_path { path = "/proc" }
}
volume {
name = "sys"
host_path { path = "/sys" }
}
volume {
name = "root"
host_path { path = "/" }
}
}
}
}
}
resource "kubernetes_service" "node_exporter" {
metadata {
name = "node-exporter"
namespace = kubernetes_namespace.swarm.metadata[0].name
}
spec {
selector = { app = "node-exporter" }
port {
port = 9100
target_port = 9100
}
} }
} }
+2 -2
View File
@@ -25,7 +25,7 @@ function displayAspect(): number {
return Math.max(1, Math.min(MAX_ASPECT, q)); return Math.max(1, Math.min(MAX_ASPECT, q));
} }
const LIVE_POLL_MS = 1000; const LIVE_POLL_MS = 2000;
export default function App(): JSX.Element { export default function App(): JSX.Element {
const [droneCount, setDroneCount] = useState(8); const [droneCount, setDroneCount] = useState(8);
@@ -268,7 +268,7 @@ export default function App(): JSX.Element {
{busiest.length === 0 && <div style={styles.dim}>no links in range</div>} {busiest.length === 0 && <div style={styles.dim}>no links in range</div>}
</Section> </Section>
<Section title="Legend"> <Section title="Legend">
<div style={styles.dim}>blue line pose broadcast (5 Hz, ~46 B)</div> <div style={styles.dim}>blue line pose broadcast (5 Hz, 45 B)</div>
<div style={styles.dim}>green line bulk sync (sealed partitions)</div> <div style={styles.dim}>green line bulk sync (sealed partitions)</div>
<div style={styles.dim}>flash broadcast event delivered</div> <div style={styles.dim}>flash broadcast event delivered</div>
<div style={styles.dim}>red circle transit object crossing the area</div> <div style={styles.dim}>red circle transit object crossing the area</div>
+7 -7
View File
@@ -20,24 +20,24 @@ export interface LivePose {
// latest battery reading from telemetry. arg_max keeps it a single scan. // latest battery reading from telemetry. arg_max keeps it a single scan.
const LIVE_SQL = ` const LIVE_SQL = `
WITH pose AS ( WITH pose AS (
SELECT drone_id, SELECT drone,
arg_max(pos_x, ts_ns) AS x, arg_max(pos_x, ts_ns) AS x,
arg_max(pos_y, ts_ns) AS y, arg_max(pos_y, ts_ns) AS y,
arg_max(yaw, ts_ns) AS yaw, arg_max(yaw, ts_ns) AS yaw,
max(ts_ns) AS ts_ns max(ts_ns) AS ts_ns
FROM state FROM state
WHERE direction = 'sent' WHERE direction = 'sent'
GROUP BY drone_id GROUP BY drone
), ),
batt AS ( batt AS (
SELECT drone_id, arg_max(level_pct, ts_ns) AS battery SELECT drone, arg_max(level_pct, ts_ns) AS battery
FROM telemetry FROM telemetry
WHERE sensor = 'battery' WHERE sensor = 'battery'
GROUP BY drone_id GROUP BY drone
) )
SELECT p.drone_id, p.x, p.y, p.yaw, p.ts_ns, b.battery SELECT p.drone, p.x, p.y, p.yaw, p.ts_ns, b.battery
FROM pose p LEFT JOIN batt b USING (drone_id) FROM pose p LEFT JOIN batt b USING (drone)
ORDER BY p.drone_id`; ORDER BY p.drone`;
interface QueryResponse { interface QueryResponse {
columns?: string[]; columns?: string[];
+1 -1
View File
@@ -57,7 +57,7 @@ export function areaWidth(): number {
export const LINK_RANGE = 320; export const LINK_RANGE = 320;
const CRUISE = 28; const CRUISE = 28;
const SEPARATION = 55; const SEPARATION = 55;
const POSE_BYTES = 46; const POSE_BYTES = 45;
const POSE_HZ = 5; const POSE_HZ = 5;
// Steady per-link broadcast throughput (both directions), bytes/s // Steady per-link broadcast throughput (both directions), bytes/s
export const POSE_RATE_BYTES = POSE_BYTES * POSE_HZ; export const POSE_RATE_BYTES = POSE_BYTES * POSE_HZ;
+9
View File
@@ -7,6 +7,10 @@ COPY virtual_drone ./virtual_drone
COPY monitoring ./monitoring COPY monitoring ./monitoring
COPY explorer ./explorer COPY explorer ./explorer
COPY tests ./tests COPY tests ./tests
RUN groupadd --gid 10001 swarm \
&& useradd --uid 10001 --gid swarm --home-dir /app --shell /usr/sbin/nologin swarm \
&& chown -R swarm:swarm /app
USER swarm
RUN pytest tests/ -q RUN pytest tests/ -q
FROM python:3.12-slim FROM python:3.12-slim
@@ -19,6 +23,11 @@ COPY virtual_drone ./virtual_drone
# them without bind mounts (Compose overrides these with live mounts) # them without bind mounts (Compose overrides these with live mounts)
COPY monitoring ./monitoring COPY monitoring ./monitoring
COPY explorer ./explorer COPY explorer ./explorer
RUN groupadd --gid 10001 swarm \
&& useradd --uid 10001 --gid swarm --home-dir /app --shell /usr/sbin/nologin swarm \
&& mkdir -p /data \
&& chown -R swarm:swarm /app /data
ENV DATA_DIR=/data ENV DATA_DIR=/data
USER swarm
CMD ["python", "-m", "virtual_drone.main"] CMD ["python", "-m", "virtual_drone.main"]
+5 -1
View File
@@ -24,7 +24,8 @@ services:
command: ["python", "monitoring/exporter.py"] command: ["python", "monitoring/exporter.py"]
environment: environment:
DATA_DIR: /data DATA_DIR: /data
SCAN_INTERVAL_S: "5" SCAN_INTERVAL_S: "30"
METRIC_FLIGHT_WINDOW: "5"
volumes: volumes:
- ${SWARM_T1_DIR:-../var/t1}:/data:ro - ${SWARM_T1_DIR:-../var/t1}:/data:ro
- ./monitoring:/app/monitoring:ro - ./monitoring:/app/monitoring:ro
@@ -37,6 +38,9 @@ services:
command: ["python", "explorer/server.py"] command: ["python", "explorer/server.py"]
environment: environment:
DATA_DIR: /data DATA_DIR: /data
VIEW_REFRESH_S: "30"
TREE_CACHE_S: "15"
METRIC_FLIGHT_WINDOW: "5"
volumes: volumes:
- ${SWARM_T1_DIR:-../var/t1}:/data:ro - ${SWARM_T1_DIR:-../var/t1}:/data:ro
- ./explorer:/app/explorer:ro - ./explorer:/app/explorer:ro
+87 -36
View File
@@ -1,33 +1,28 @@
"""Data-plane explorer: a web view over the Hive-partitioned Parquet lake. """Data-plane explorer: a web view over the Hive-partitioned Parquet lake."""
Serves three things:
/ single-page UI (partition tree + read-only SQL console)
/api/tree partition hierarchy with file counts and bytes, live
/api/query gated read-only DuckDB SQL, same statement rules as the
peer query channel (SELECT/WITH only, single statement)
The point is doctrinal, not just convenient: the explorer reuses the exact
read-only SQL contract that drones expose to each other, so "looking at the
data plane" on the bench exercises the same path a peer would use in flight.
"""
from __future__ import annotations from __future__ import annotations
import json import json
import os import os
import re import re
import sys
import threading
import time
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
from pathlib import Path from pathlib import Path
import duckdb import duckdb
sys.path.insert(0, str(Path(__file__).resolve().parents[1] / "monitoring"))
from lake import flight_window, iter_parquet_files, parquet_reader # noqa: E402
DATA_DIR = Path(os.environ.get("DATA_DIR", "./data")) DATA_DIR = Path(os.environ.get("DATA_DIR", "./data"))
PORT = int(os.environ.get("EXPLORER_PORT", "8088")) PORT = int(os.environ.get("EXPLORER_PORT", "8088"))
ROW_LIMIT = int(os.environ.get("ROW_LIMIT", "500")) ROW_LIMIT = int(os.environ.get("ROW_LIMIT", "500"))
VIEW_REFRESH_S = float(os.environ.get("VIEW_REFRESH_S", "30"))
TREE_CACHE_S = float(os.environ.get("TREE_CACHE_S", "15"))
STATIC_DIR = Path(__file__).parent STATIC_DIR = Path(__file__).parent
# Same spirit as the forced-command gate on a real drone: one statement,
# must be a read, no statement that could write, configure, or reach out.
_ALLOWED_START = re.compile(r"^\s*(SELECT|WITH|DESCRIBE|SUMMARIZE|SHOW)\b", re.IGNORECASE) _ALLOWED_START = re.compile(r"^\s*(SELECT|WITH|DESCRIBE|SUMMARIZE|SHOW)\b", re.IGNORECASE)
_FORBIDDEN = re.compile( _FORBIDDEN = re.compile(
r"\b(INSERT|UPDATE|DELETE|CREATE|DROP|ALTER|ATTACH|DETACH|COPY|EXPORT" r"\b(INSERT|UPDATE|DELETE|CREATE|DROP|ALTER|ATTACH|DETACH|COPY|EXPORT"
@@ -35,9 +30,17 @@ _FORBIDDEN = re.compile(
re.IGNORECASE, re.IGNORECASE,
) )
_lock = threading.Lock()
_con: duckdb.DuckDBPyConnection | None = None
_views_at = 0.0
_tree_cache: dict | None = None
_tree_at = 0.0
_last_tree_s = 0.0
_last_query_s = 0.0
_query_count = 0
def gate(sql: str) -> str | None: def gate(sql: str) -> str | None:
"""Return a rejection reason, or None if the statement passes."""
stripped = re.sub(r"--[^\n]*|/\*.*?\*/", " ", sql, flags=re.DOTALL).strip().rstrip(";") stripped = re.sub(r"--[^\n]*|/\*.*?\*/", " ", sql, flags=re.DOTALL).strip().rstrip(";")
if not stripped: if not stripped:
return "empty statement" return "empty statement"
@@ -50,27 +53,39 @@ def gate(sql: str) -> str | None:
return None return None
def connect() -> duckdb.DuckDBPyConnection: def _refresh_views() -> None:
"""Fresh connection with the three datasets pre-registered as views.""" global _con, _views_at
con = duckdb.connect() con = duckdb.connect()
for ds in ("telemetry", "detections", "state"): for ds in ("telemetry", "detections", "state"):
pattern = f"{DATA_DIR}/dataset={ds}/**/*.parquet" reader = parquet_reader(DATA_DIR, ds)
try: try:
con.execute( con.execute(f"CREATE OR REPLACE VIEW {ds} AS SELECT * FROM read_parquet({reader})")
f"CREATE VIEW {ds} AS SELECT * FROM read_parquet("
f"'{pattern}', hive_partitioning=true, union_by_name=true)"
)
except duckdb.Error: except duckdb.Error:
pass # dataset not written yet; view simply won't exist pass
return con with _lock:
if _con is not None:
_con.close()
_con = con
_views_at = time.monotonic()
def connect() -> duckdb.DuckDBPyConnection:
if _con is None or time.monotonic() - _views_at > VIEW_REFRESH_S:
_refresh_views()
assert _con is not None
return _con
def tree() -> dict: def tree() -> dict:
"""Partition hierarchy: dataset -> flight -> drone -> leafs, with sizes.""" global _tree_cache, _tree_at, _last_tree_s
now = time.monotonic()
if _tree_cache is not None and now - _tree_at < TREE_CACHE_S:
return _tree_cache
started = now
root: dict = {} root: dict = {}
total_bytes = 0 total_bytes = 0
total_files = 0 total_files = 0
for f in sorted(DATA_DIR.rglob("*.parquet")): for f in sorted(iter_parquet_files(DATA_DIR)):
rel = f.relative_to(DATA_DIR) rel = f.relative_to(DATA_DIR)
size = f.stat().st_size size = f.stat().st_size
total_bytes += size total_bytes += size
@@ -80,7 +95,6 @@ def tree() -> dict:
node = node.setdefault("children", {}).setdefault(part, {}) node = node.setdefault("children", {}).setdefault(part, {})
leaf = node.setdefault("children", {}).setdefault(rel.parts[-1], {}) leaf = node.setdefault("children", {}).setdefault(rel.parts[-1], {})
leaf["bytes"] = size leaf["bytes"] = size
# roll sizes up the tree
node = root node = root
node["bytes"] = node.get("bytes", 0) + size node["bytes"] = node.get("bytes", 0) + size
node["files"] = node.get("files", 0) + 1 node["files"] = node.get("files", 0) + 1
@@ -88,18 +102,26 @@ def tree() -> dict:
node = node["children"][part] node = node["children"][part]
node["bytes"] = node.get("bytes", 0) + size node["bytes"] = node.get("bytes", 0) + size
node["files"] = node.get("files", 0) + 1 node["files"] = node.get("files", 0) + 1
return {"tree": root, "total_bytes": total_bytes, "total_files": total_files} _last_tree_s = time.monotonic() - started
_tree_cache = {"tree": root, "total_bytes": total_bytes, "total_files": total_files}
_tree_at = now
return _tree_cache
def run_query(sql: str) -> dict: def run_query(sql: str) -> dict:
global _last_query_s, _query_count
reason = gate(sql) reason = gate(sql)
if reason: if reason:
return {"error": f"rejected by read-only gate: {reason}"} return {"error": f"rejected by read-only gate: {reason}"}
started = time.monotonic()
with _lock:
con = connect() con = connect()
try: try:
cur = con.sql(sql) cur = con.sql(sql)
columns = cur.columns columns = cur.columns
rows = cur.fetchmany(ROW_LIMIT) rows = cur.fetchmany(ROW_LIMIT)
_query_count += 1
_last_query_s = time.monotonic() - started
return { return {
"columns": columns, "columns": columns,
"rows": [[repr(v) if isinstance(v, bytes) else v for v in row] for row in rows], "rows": [[repr(v) if isinstance(v, bytes) else v for v in row] for row in rows],
@@ -107,8 +129,23 @@ def run_query(sql: str) -> dict:
} }
except duckdb.Error as exc: except duckdb.Error as exc:
return {"error": str(exc)} return {"error": str(exc)}
finally:
con.close()
def metrics_text() -> str:
return "\n".join([
"# HELP swarm_explorer_query_duration_seconds Wall time of the last SQL query",
"# TYPE swarm_explorer_query_duration_seconds gauge",
f"swarm_explorer_query_duration_seconds {_last_query_s:.4f}",
"# HELP swarm_explorer_tree_duration_seconds Wall time of the last partition tree build",
"# TYPE swarm_explorer_tree_duration_seconds gauge",
f"swarm_explorer_tree_duration_seconds {_last_tree_s:.4f}",
"# HELP swarm_explorer_queries_total Read-only queries served",
"# TYPE swarm_explorer_queries_total counter",
f"swarm_explorer_queries_total {_query_count}",
"# HELP swarm_metric_flight_window Flight partitions in DuckDB views",
"# TYPE swarm_metric_flight_window gauge",
f"swarm_metric_flight_window {flight_window()}",
]) + "\n"
class Handler(BaseHTTPRequestHandler): class Handler(BaseHTTPRequestHandler):
@@ -116,13 +153,11 @@ class Handler(BaseHTTPRequestHandler):
self.send_response(code) self.send_response(code)
self.send_header("Content-Type", ctype) self.send_header("Content-Type", ctype)
self.send_header("Content-Length", str(len(body))) self.send_header("Content-Length", str(len(body)))
# Dev CORS: lets the prototype's live mode poll the API from another
# origin. Everything behind this is read-only by construction.
self.send_header("Access-Control-Allow-Origin", "*") self.send_header("Access-Control-Allow-Origin", "*")
self.end_headers() self.end_headers()
self.wfile.write(body) self.wfile.write(body)
def do_OPTIONS(self) -> None: # noqa: N802 — http.server API def do_OPTIONS(self) -> None: # noqa: N802
self.send_response(204) self.send_response(204)
self.send_header("Access-Control-Allow-Origin", "*") self.send_header("Access-Control-Allow-Origin", "*")
self.send_header("Access-Control-Allow-Methods", "GET, POST, OPTIONS") self.send_header("Access-Control-Allow-Methods", "GET, POST, OPTIONS")
@@ -132,15 +167,17 @@ class Handler(BaseHTTPRequestHandler):
def _json(self, payload: dict, code: int = 200) -> None: def _json(self, payload: dict, code: int = 200) -> None:
self._send(code, json.dumps(payload, default=str).encode(), "application/json") self._send(code, json.dumps(payload, default=str).encode(), "application/json")
def do_GET(self) -> None: # noqa: N802 — http.server API def do_GET(self) -> None: # noqa: N802
if self.path in ("/", "/index.html"): if self.path in ("/", "/index.html"):
self._send(200, (STATIC_DIR / "index.html").read_bytes(), "text/html; charset=utf-8") self._send(200, (STATIC_DIR / "index.html").read_bytes(), "text/html; charset=utf-8")
elif self.path == "/api/tree": elif self.path == "/api/tree":
self._json(tree()) self._json(tree())
elif self.path == "/metrics":
self._send(200, metrics_text().encode(), "text/plain; version=0.0.4")
else: else:
self._send(404, b"not found", "text/plain") self._send(404, b"not found", "text/plain")
def do_POST(self) -> None: # noqa: N802 — http.server API def do_POST(self) -> None: # noqa: N802
if self.path != "/api/query": if self.path != "/api/query":
self._send(404, b"not found", "text/plain") self._send(404, b"not found", "text/plain")
return return
@@ -156,6 +193,20 @@ class Handler(BaseHTTPRequestHandler):
pass pass
def _view_loop() -> None:
while True:
try:
_refresh_views()
except Exception:
pass
time.sleep(VIEW_REFRESH_S)
if __name__ == "__main__": if __name__ == "__main__":
print(f"data-plane explorer on :{PORT}, reading {DATA_DIR}") _refresh_views()
threading.Thread(target=_view_loop, daemon=True).start()
print(
f"data-plane explorer on :{PORT}, views refresh every {VIEW_REFRESH_S}s, "
f"last {flight_window()} flights, reading {DATA_DIR}"
)
ThreadingHTTPServer(("", PORT), Handler).serve_forever() ThreadingHTTPServer(("", PORT), Handler).serve_forever()
+49 -34
View File
@@ -1,13 +1,14 @@
"""Prometheus exporter over the simulator's Parquet output. """Prometheus exporter over the simulator's Parquet output.
Periodically scans DATA_DIR with DuckDB and exposes fleet statistics as Periodically scans DATA_DIR with DuckDB and exposes fleet statistics as
/metrics. Zero dependencies beyond duckdb: the exposition format is plain /metrics. Scans only the most recent flight partitions by default so CPU
text, served with the standard library HTTP server. stays bounded as the lake grows across pod restarts.
""" """
from __future__ import annotations from __future__ import annotations
import os import os
import sys
import threading import threading
import time import time
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
@@ -15,27 +16,29 @@ from pathlib import Path
import duckdb import duckdb
sys.path.insert(0, str(Path(__file__).resolve().parent))
from lake import flight_window, iter_parquet_files, parquet_reader # noqa: E402
DATA_DIR = Path(os.environ.get("DATA_DIR", "./data")) DATA_DIR = Path(os.environ.get("DATA_DIR", "./data"))
PORT = int(os.environ.get("EXPORTER_PORT", "9105")) PORT = int(os.environ.get("EXPORTER_PORT", "9105"))
SCAN_INTERVAL_S = float(os.environ.get("SCAN_INTERVAL_S", "5")) SCAN_INTERVAL_S = float(os.environ.get("SCAN_INTERVAL_S", "15"))
_lock = threading.Lock() _lock = threading.Lock()
_payload = "# swarm exporter starting\n" _payload = "# swarm exporter starting\n"
_last_scan_s = 0.0
def _q(con: duckdb.DuckDBPyConnection, sql: str) -> list[tuple]: def _q(con: duckdb.DuckDBPyConnection, sql: str) -> list[tuple]:
try: try:
return con.sql(sql).fetchall() return con.sql(sql).fetchall()
except duckdb.Error: except duckdb.Error:
return [] # partitions may not exist yet while the swarm warms up return []
def collect() -> str: def collect() -> str:
global _last_scan_s
started = time.monotonic()
con = duckdb.connect() con = duckdb.connect()
# union_by_name: sensors have different schemas under one dataset glob
glob = lambda ds: ( # noqa: E731
f"'{DATA_DIR}/dataset={ds}/**/*.parquet', hive_partitioning=true, union_by_name=true"
)
lines: list[str] = [] lines: list[str] = []
def metric(name: str, help_text: str, mtype: str, rows: list[str]) -> None: def metric(name: str, help_text: str, mtype: str, rows: list[str]) -> None:
@@ -43,27 +46,38 @@ def collect() -> str:
lines.append(f"# TYPE {name} {mtype}") lines.append(f"# TYPE {name} {mtype}")
lines.extend(rows) lines.extend(rows)
for ds, name in (
("telemetry", "swarm_rows_total"),
("detections", "swarm_detections_total"),
("state", "swarm_state_frames_total"),
):
reader = parquet_reader(DATA_DIR, ds)
if ds == "telemetry":
metric( metric(
"swarm_rows_total", "Telemetry rows written per drone and sensor", "gauge", name, f"Rows in dataset={ds} (recent {flight_window()} flights)", "gauge",
[f'swarm_rows_total{{drone="{d}",sensor="{s}"}} {n}' [f'{name}{{drone="{d}",sensor="{s}"}} {n}'
for d, s, n in _q(con, f"SELECT drone, sensor, count(*) FROM read_parquet({glob('telemetry')}) GROUP BY 1,2")], for d, s, n in _q(con, f"SELECT drone, sensor, count(*) FROM read_parquet({reader}) GROUP BY 1,2")],
) )
elif ds == "detections":
metric( metric(
"swarm_detections_total", "Detection events per drone and class", "gauge", name, "Detection events per drone and class", "gauge",
[f'swarm_detections_total{{drone="{d}",cls="{c}"}} {n}' [f'{name}{{drone="{d}",cls="{c}"}} {n}'
for d, c, n in _q(con, f"SELECT drone, cls, count(*) FROM read_parquet({glob('detections')}) GROUP BY 1,2")], for d, c, n in _q(con, f"SELECT drone, cls, count(*) FROM read_parquet({reader}) GROUP BY 1,2")],
) )
else:
metric( metric(
"swarm_state_frames_total", "Pose broadcast frames per drone and direction", "gauge", name, "Pose broadcast frames per drone and direction", "gauge",
[f'swarm_state_frames_total{{drone="{d}",direction="{dr}"}} {n}' [f'{name}{{drone="{d}",direction="{dr}"}} {n}'
for d, dr, n in _q(con, f"SELECT drone, direction, count(*) FROM read_parquet({glob('state')}) GROUP BY 1,2")], for d, dr, n in _q(con, f"SELECT drone, direction, count(*) FROM read_parquet({reader}) GROUP BY 1,2")],
) )
telem = parquet_reader(DATA_DIR, "telemetry")
metric( metric(
"swarm_battery_pct", "Latest battery level per drone", "gauge", "swarm_battery_pct", "Latest battery level per drone", "gauge",
[f'swarm_battery_pct{{drone="{d}"}} {v}' [f'swarm_battery_pct{{drone="{d}"}} {v}'
for d, v in _q(con, f""" for d, v in _q(con, f"""
SELECT drone, arg_max(level_pct, ts_ns) SELECT drone, arg_max(level_pct, ts_ns)
FROM read_parquet({glob('telemetry')}) FROM read_parquet({telem})
WHERE sensor='battery' GROUP BY drone""")], WHERE sensor='battery' GROUP BY drone""")],
) )
metric( metric(
@@ -71,26 +85,24 @@ def collect() -> str:
[f'swarm_rssi_dbm{{drone="{d}",peer="{p}"}} {v}' [f'swarm_rssi_dbm{{drone="{d}",peer="{p}"}} {v}'
for d, p, v in _q(con, f""" for d, p, v in _q(con, f"""
SELECT drone, peer_id, arg_max(rssi_dbm, ts_ns) SELECT drone, peer_id, arg_max(rssi_dbm, ts_ns)
FROM read_parquet({glob('telemetry')}) FROM read_parquet({telem})
WHERE sensor='rssi' GROUP BY drone, peer_id""")],
)
metric(
"swarm_peer_distance_m", "Latest inter-drone distance estimate", "gauge",
[f'swarm_peer_distance_m{{drone="{d}",peer="{p}"}} {v}'
for d, p, v in _q(con, f"""
SELECT drone, peer_id, arg_max(distance_m, ts_ns)
FROM read_parquet({glob('telemetry')})
WHERE sensor='rssi' GROUP BY drone, peer_id""")], WHERE sensor='rssi' GROUP BY drone, peer_id""")],
) )
files = list(DATA_DIR.rglob("*.parquet")) files = iter_parquet_files(DATA_DIR)
metric( metric(
"swarm_parquet_bytes", "Bytes on disk per dataset", "gauge", "swarm_parquet_bytes", "Bytes on disk per dataset (recent flights)", "gauge",
[f'swarm_parquet_bytes{{dataset="{ds}"}} {sum(f.stat().st_size for f in files if f"dataset={ds}" in str(f))}' [f'swarm_parquet_bytes{{dataset="{ds}"}} {sum(f.stat().st_size for f in files if f"dataset={ds}" in str(f))}'
for ds in ("telemetry", "detections", "state")], for ds in ("telemetry", "detections", "state")],
) )
metric("swarm_parquet_files", "Parquet files on disk", "gauge", metric("swarm_parquet_files", "Parquet files scanned (recent flights)", "gauge",
[f"swarm_parquet_files {len(files)}"]) [f"swarm_parquet_files {len(files)}"])
metric("swarm_metric_flight_window", "Flight partitions included per dataset", "gauge",
[f"swarm_metric_flight_window {flight_window()}"])
_last_scan_s = time.monotonic() - started
metric("swarm_exporter_scan_duration_seconds", "Wall time of the last metrics scan", "gauge",
[f"swarm_exporter_scan_duration_seconds {_last_scan_s:.4f}"])
con.close() con.close()
return "\n".join(lines) + "\n" return "\n".join(lines) + "\n"
@@ -101,11 +113,11 @@ def scanner() -> None:
started = time.monotonic() started = time.monotonic()
try: try:
payload = collect() payload = collect()
except Exception as exc: # keep serving stale metrics over dying except Exception as exc:
payload = f"# collect error: {exc}\n" payload = f"# collect error: {exc}\n"
with _lock: with _lock:
_payload = payload _payload = payload
time.sleep(max(0.5, SCAN_INTERVAL_S - (time.monotonic() - started))) time.sleep(max(1.0, SCAN_INTERVAL_S - (time.monotonic() - started)))
class Handler(BaseHTTPRequestHandler): class Handler(BaseHTTPRequestHandler):
@@ -123,10 +135,13 @@ class Handler(BaseHTTPRequestHandler):
self.wfile.write(body) self.wfile.write(body)
def log_message(self, *_args: object) -> None: def log_message(self, *_args: object) -> None:
pass # scrapes every few seconds; keep the log quiet pass
if __name__ == "__main__": if __name__ == "__main__":
threading.Thread(target=scanner, daemon=True).start() threading.Thread(target=scanner, daemon=True).start()
print(f"swarm exporter on :{PORT}/metrics, scanning {DATA_DIR} every {SCAN_INTERVAL_S}s") print(
f"swarm exporter on :{PORT}/metrics, scanning last {flight_window()} flights "
f"every {SCAN_INTERVAL_S}s under {DATA_DIR}"
)
ThreadingHTTPServer(("", PORT), Handler).serve_forever() ThreadingHTTPServer(("", PORT), Handler).serve_forever()
@@ -0,0 +1,90 @@
{
"uid": "swarm-platform",
"title": "Swarm Platform — CPU & scan health",
"tags": ["swarm", "platform"],
"timezone": "utc",
"schemaVersion": 39,
"version": 1,
"refresh": "10s",
"time": { "from": "now-30m", "to": "now" },
"panels": [
{
"id": 1, "type": "timeseries", "title": "Node CPU %",
"gridPos": { "h": 8, "w": 12, "x": 0, "y": 0 },
"datasource": { "type": "prometheus", "uid": "swarm-prom" },
"targets": [{
"expr": "100 - (avg by (instance) (rate(node_cpu_seconds_total{mode=\"idle\"}[1m])) * 100)",
"legendFormat": "{{instance}}", "refId": "A"
}],
"fieldConfig": { "defaults": { "unit": "percent", "min": 0, "max": 100 }, "overrides": [] }
},
{
"id": 2, "type": "timeseries", "title": "Node memory used %",
"gridPos": { "h": 8, "w": 12, "x": 12, "y": 0 },
"datasource": { "type": "prometheus", "uid": "swarm-prom" },
"targets": [{
"expr": "(1 - node_memory_MemAvailable_bytes / node_memory_MemTotal_bytes) * 100",
"legendFormat": "used", "refId": "A"
}],
"fieldConfig": { "defaults": { "unit": "percent", "min": 0, "max": 100 }, "overrides": [] }
},
{
"id": 3, "type": "timeseries", "title": "Exporter scan duration",
"gridPos": { "h": 7, "w": 8, "x": 0, "y": 8 },
"datasource": { "type": "prometheus", "uid": "swarm-prom" },
"targets": [{
"expr": "swarm_exporter_scan_duration_seconds",
"legendFormat": "scan seconds", "refId": "A"
}],
"fieldConfig": { "defaults": { "unit": "s" }, "overrides": [] }
},
{
"id": 4, "type": "timeseries", "title": "Explorer query duration",
"gridPos": { "h": 7, "w": 8, "x": 8, "y": 8 },
"datasource": { "type": "prometheus", "uid": "swarm-prom" },
"targets": [{
"expr": "swarm_explorer_query_duration_seconds",
"legendFormat": "last query", "refId": "A"
}],
"fieldConfig": { "defaults": { "unit": "s" }, "overrides": [] }
},
{
"id": 5, "type": "timeseries", "title": "Explorer tree build duration",
"gridPos": { "h": 7, "w": 8, "x": 16, "y": 8 },
"datasource": { "type": "prometheus", "uid": "swarm-prom" },
"targets": [{
"expr": "swarm_explorer_tree_duration_seconds",
"legendFormat": "tree seconds", "refId": "A"
}],
"fieldConfig": { "defaults": { "unit": "s" }, "overrides": [] }
},
{
"id": 6, "type": "stat", "title": "Parquet files scanned",
"gridPos": { "h": 5, "w": 6, "x": 0, "y": 15 },
"datasource": { "type": "prometheus", "uid": "swarm-prom" },
"targets": [{ "expr": "swarm_parquet_files", "instant": true, "refId": "A" }],
"options": { "reduceOptions": { "calcs": ["lastNotNull"] } }
},
{
"id": 7, "type": "stat", "title": "Flight window",
"gridPos": { "h": 5, "w": 6, "x": 6, "y": 15 },
"datasource": { "type": "prometheus", "uid": "swarm-prom" },
"targets": [{ "expr": "swarm_metric_flight_window", "instant": true, "refId": "A" }],
"options": { "reduceOptions": { "calcs": ["lastNotNull"] } }
},
{
"id": 8, "type": "stat", "title": "Telemetry rows (window)",
"gridPos": { "h": 5, "w": 6, "x": 12, "y": 15 },
"datasource": { "type": "prometheus", "uid": "swarm-prom" },
"targets": [{ "expr": "sum(swarm_rows_total)", "instant": true, "refId": "A" }],
"options": { "reduceOptions": { "calcs": ["lastNotNull"] } }
},
{
"id": 9, "type": "stat", "title": "Min fleet battery %",
"gridPos": { "h": 5, "w": 6, "x": 18, "y": 15 },
"datasource": { "type": "prometheus", "uid": "swarm-prom" },
"targets": [{ "expr": "min(swarm_battery_pct)", "instant": true, "refId": "A" }],
"fieldConfig": { "defaults": { "unit": "percent", "min": 0, "max": 100 }, "overrides": [] }
}
]
}
+43
View File
@@ -0,0 +1,43 @@
"""Scan the Hive-partitioned lake without re-reading every historical flight."""
from __future__ import annotations
import os
from pathlib import Path
def flight_window() -> int:
return max(1, int(os.environ.get("METRIC_FLIGHT_WINDOW", "5")))
def recent_flights(data_dir: Path, dataset: str, limit: int | None = None) -> list[Path]:
"""Most recently modified flight= partitions for a dataset."""
limit = limit or flight_window()
base = data_dir / f"dataset={dataset}"
if not base.is_dir():
return []
flights = [p for p in base.iterdir() if p.is_dir() and p.name.startswith("flight=")]
flights.sort(key=lambda p: p.stat().st_mtime, reverse=True)
return flights[:limit]
def parquet_reader(data_dir: Path, dataset: str, *, limit: int | None = None) -> str:
"""DuckDB read_parquet() source limited to recent flights."""
flights = recent_flights(data_dir, dataset, limit)
if not flights:
path = data_dir / f"dataset={dataset}" / "**" / "*.parquet"
return f"'{path}', hive_partitioning=true, union_by_name=true"
if len(flights) == 1:
return f"'{flights[0]}/**/*.parquet', hive_partitioning=true, union_by_name=true"
inner = ", ".join(f"'{f}/**/*.parquet'" for f in flights)
return f"[{inner}], hive_partitioning=true, union_by_name=true"
def iter_parquet_files(data_dir: Path, *, flight_limit: int | None = None) -> list[Path]:
"""Parquet paths under recent flights only — avoids full-lake rglob."""
limit = flight_limit or flight_window()
out: list[Path] = []
for ds in ("telemetry", "detections", "state"):
for flight in recent_flights(data_dir, ds, limit):
out.extend(flight.rglob("*.parquet"))
return out
+4 -1
View File
@@ -1,7 +1,10 @@
global: global:
scrape_interval: 5s scrape_interval: 15s
scrape_configs: scrape_configs:
- job_name: swarm - job_name: swarm
static_configs: static_configs:
- targets: ["exporter:9105"] - targets: ["exporter:9105"]
- job_name: explorer
static_configs:
- targets: ["explorer:8088"]
+1 -1
View File
@@ -9,7 +9,7 @@ from virtual_drone.flight import Pose
def test_frame_size_matches_wire_spec() -> None: def test_frame_size_matches_wire_spec() -> None:
# Documented as 46 B in docs/04; struct packs to FRAME.size (45 B on this layout). # Keep docs/04, the journey, and prototype/src/sim.ts POSE_BYTES in lockstep.
assert FRAME.size == 45 assert FRAME.size == 45
+40
View File
@@ -0,0 +1,40 @@
"""Lake scan helpers — bounded to recent flight partitions."""
from __future__ import annotations
from pathlib import Path
from monitoring.lake import flight_window, iter_parquet_files, parquet_reader, recent_flights
def test_parquet_reader_limits_to_recent_flights(tmp_path: Path) -> None:
for i, name in enumerate(("flight=aaa", "flight=bbb", "flight=ccc")):
hour = tmp_path / f"dataset=telemetry/{name}/drone=dr-01/sensor=imu/year=2026/month=07/day=08/hour=10"
hour.mkdir(parents=True)
f = hour / "data.parquet"
f.write_bytes(b"x" * (i + 1))
# Make later names newer
import os
import time
os.utime(f, (time.time() + i, time.time() + i))
reader = parquet_reader(tmp_path, "telemetry", limit=1)
assert "flight=ccc" in reader
assert "flight=bbb" not in reader
def test_iter_parquet_files_skips_old_flights(tmp_path: Path) -> None:
old = tmp_path / "dataset=state/flight=old/drone=dr-01/year=2026/month=07/day=08/hour=09"
old.mkdir(parents=True)
(old / "data.parquet").write_bytes(b"old")
new = tmp_path / "dataset=state/flight=new/drone=dr-01/year=2026/month=07/day=08/hour=10"
new.mkdir(parents=True)
new_file = new / "data.parquet"
new_file.write_bytes(b"new")
import os
import time
os.utime(new_file, (time.time() + 10, time.time() + 10))
files = iter_parquet_files(tmp_path, flight_limit=1)
assert len(files) == 1
assert "flight=new" in str(files[0])
+4 -3
View File
@@ -1,12 +1,13 @@
"""State broadcast over UDP: the compact pose frame from the sync design. """State broadcast over UDP: the compact pose frame from the sync design.
Frame layout (little-endian, 46 bytes): Frame layout (little-endian, 45 bytes) — source of truth for docs and the
prototype bandwidth estimate:
magic 2s b"SH" magic 2s b"SH"
version B version B
drone_id 8s zero-padded ascii drone_id 8s zero-padded ascii (PoC; production may switch to uint16 registry)
ts_ns q epoch nanoseconds ts_ns q epoch nanoseconds
pos_mm 3i position, millimeters (quantized on the wire only) pos_mm 3i position in the mission frame, millimeters (wire quantization only)
att_cdeg 3h roll/pitch/yaw, centi-degrees att_cdeg 3h roll/pitch/yaw, centi-degrees
vel_cms 3h velocity, cm/s vel_cms 3h velocity, cm/s
frame_ref B frame_ref B
+5
View File
@@ -2,6 +2,7 @@
from __future__ import annotations from __future__ import annotations
import os
import random import random
import select import select
import time import time
@@ -108,6 +109,10 @@ def run(cfg: Config) -> None:
writer.seal() writer.seal()
print(f"[{cfg.drone_id}] done: {frames_sent} state frames sent, " print(f"[{cfg.drone_id}] done: {frames_sent} state frames sent, "
f"{len(peers_seen)} peers seen {sorted(peers_seen)}; sealed to {root}") f"{len(peers_seen)} peers seen {sorted(peers_seen)}; sealed to {root}")
if os.environ.get("KEEP_ALIVE", "0") == "1":
print(f"[{cfg.drone_id}] KEEP_ALIVE=1 — idle after seal (no pod restart churn)")
while True:
time.sleep(3600)
if __name__ == "__main__": if __name__ == "__main__":