From c90100a021e67db639ba3117d111e08977a47da0 Mon Sep 17 00:00:00 2001 From: sdairs Date: Fri, 2 Oct 2026 19:26:20 +0100 Subject: [PATCH 1/3] Add bounded C++ device readings with confirmed async commit --- .github/workflows/device-heartbeats.yml | 41 +++ applications/device-heartbeats/.clang-format | 3 + applications/device-heartbeats/.env.example | 12 + applications/device-heartbeats/.gitignore | 4 + applications/device-heartbeats/CMakeLists.txt | 17 + applications/device-heartbeats/README.md | 152 +++++++++ .../device-heartbeats/sql/bootstrap.sql | 9 + .../device-heartbeats/sql/cleanup.sql | 7 + applications/device-heartbeats/sql/grants.sql | 6 + .../device-heartbeats/sql/migrate.sql | 36 +++ applications/device-heartbeats/sql/seed.sql | 8 + applications/device-heartbeats/src/server.cpp | 205 ++++++++++++ .../device-heartbeats/src/service.cpp | 216 +++++++++++++ .../device-heartbeats/src/service.hpp | 22 ++ .../device-heartbeats/src/validation.cpp | 102 ++++++ .../device-heartbeats/src/validation.hpp | 24 ++ applications/device-heartbeats/test/cloud.py | 301 ++++++++++++++++++ .../device-heartbeats/test/preflight.py | 29 ++ .../device-heartbeats/test/validation.cpp | 59 ++++ 19 files changed, 1253 insertions(+) create mode 100644 .github/workflows/device-heartbeats.yml create mode 100644 applications/device-heartbeats/.clang-format create mode 100644 applications/device-heartbeats/.env.example create mode 100644 applications/device-heartbeats/.gitignore create mode 100644 applications/device-heartbeats/CMakeLists.txt create mode 100644 applications/device-heartbeats/README.md create mode 100644 applications/device-heartbeats/sql/bootstrap.sql create mode 100644 applications/device-heartbeats/sql/cleanup.sql create mode 100644 applications/device-heartbeats/sql/grants.sql create mode 100644 applications/device-heartbeats/sql/migrate.sql create mode 100644 applications/device-heartbeats/sql/seed.sql create mode 100644 applications/device-heartbeats/src/server.cpp create mode 100644 applications/device-heartbeats/src/service.cpp create mode 100644 applications/device-heartbeats/src/service.hpp create mode 100644 applications/device-heartbeats/src/validation.cpp create mode 100644 applications/device-heartbeats/src/validation.hpp create mode 100644 applications/device-heartbeats/test/cloud.py create mode 100644 applications/device-heartbeats/test/preflight.py create mode 100644 applications/device-heartbeats/test/validation.cpp diff --git a/.github/workflows/device-heartbeats.yml b/.github/workflows/device-heartbeats.yml new file mode 100644 index 00000000..c947e9f0 --- /dev/null +++ b/.github/workflows/device-heartbeats.yml @@ -0,0 +1,41 @@ +name: Device Heartbeats checks +on: + pull_request: + paths: + - 'applications/device-heartbeats/**' + - '.github/workflows/device-heartbeats.yml' + push: + branches: [main] + paths: + - 'applications/device-heartbeats/**' + - '.github/workflows/device-heartbeats.yml' + workflow_dispatch: +permissions: + contents: read +jobs: + check: + runs-on: ubuntu-24.04 + defaults: + run: + working-directory: applications/device-heartbeats + steps: + - uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1 + - name: Install native prerequisites + run: | + sudo apt-get update + sudo apt-get install -y build-essential cmake git pkg-config libjsoncpp-dev libssl-dev libpq-dev uuid-dev zlib1g-dev libbrotli-dev clang-format + - name: Build pinned Drogon and Trantor + run: | + git clone --branch v1.9.13 --depth 1 --recurse-submodules https://github.com/drogonframework/drogon.git /tmp/drogon-source + test "$(git -C /tmp/drogon-source rev-parse HEAD)" = 4c5430757ea5451a7c38fbbef4b4bef7dbb47f2f + test "$(git -C /tmp/drogon-source/trantor rev-parse HEAD)" = 63a4e5e164e219dc3bf30cdbfa1462ae5602fa97 + cmake -S /tmp/drogon-source -B /tmp/drogon-build -DCMAKE_BUILD_TYPE=Release -DCMAKE_INSTALL_PREFIX=/tmp/drogon -DBUILD_TESTING=OFF -DBUILD_EXAMPLES=OFF -DBUILD_CTL=OFF -DBUILD_ORM=ON -DUSE_POSTGRESQL=ON -DUSE_MYSQL=OFF -DUSE_SQLITE3=OFF -DUSE_REDIS=OFF + cmake --build /tmp/drogon-build --parallel 2 + cmake --install /tmp/drogon-build + - name: Format, build and check + run: | + clang-format --dry-run --Werror src/*.cpp src/*.hpp test/validation.cpp + cmake -S . -B build -DCMAKE_BUILD_TYPE=Release -DCMAKE_PREFIX_PATH=/tmp/drogon + cmake --build build --parallel 2 + ctest --test-dir build --output-on-failure + python3 test/preflight.py "$PWD/build/device_heartbeats" diff --git a/applications/device-heartbeats/.clang-format b/applications/device-heartbeats/.clang-format new file mode 100644 index 00000000..75a203bd --- /dev/null +++ b/applications/device-heartbeats/.clang-format @@ -0,0 +1,3 @@ +BasedOnStyle: LLVM +IndentWidth: 4 +ColumnLimit: 100 diff --git a/applications/device-heartbeats/.env.example b/applications/device-heartbeats/.env.example new file mode 100644 index 00000000..67bbee5d --- /dev/null +++ b/applications/device-heartbeats/.env.example @@ -0,0 +1,12 @@ +PGHOST=your-cloud-host +PGPORT=5432 +PGDATABASE=postgres +PGUSER=heartbeats_app +PGPASSWORD=replace-with-runtime-password +PGSSLROOTCERT=/private/path/ca.pem +ADMIN_USER=cloud-created-user +ADMIN_PASSWORD=cloud-created-password +MIGRATION_PASSWORD=replace-with-separate-password +NORTH_TOKEN=replace-with-at-least-32-random-characters +SOUTH_TOKEN=replace-with-another-random-token +PORT=4000 diff --git a/applications/device-heartbeats/.gitignore b/applications/device-heartbeats/.gitignore new file mode 100644 index 00000000..f4ff1afb --- /dev/null +++ b/applications/device-heartbeats/.gitignore @@ -0,0 +1,4 @@ +build*/ +.env +*.pem +__pycache__/ diff --git a/applications/device-heartbeats/CMakeLists.txt b/applications/device-heartbeats/CMakeLists.txt new file mode 100644 index 00000000..c6fc69f0 --- /dev/null +++ b/applications/device-heartbeats/CMakeLists.txt @@ -0,0 +1,17 @@ +cmake_minimum_required(VERSION 3.22) +project(device_heartbeats VERSION 1.0.0 LANGUAGES CXX) +set(CMAKE_CXX_STANDARD 20) +set(CMAKE_CXX_STANDARD_REQUIRED ON) +set(CMAKE_CXX_EXTENSIONS OFF) +find_package(Drogon 1.9.13 REQUIRED CONFIG) +add_library(heartbeat_validation src/validation.cpp) +target_link_libraries(heartbeat_validation PUBLIC Drogon::Drogon) +target_include_directories(heartbeat_validation PUBLIC src) +target_compile_options(heartbeat_validation PRIVATE -Wall -Wextra -Wpedantic) +add_executable(device_heartbeats src/server.cpp src/service.cpp) +target_link_libraries(device_heartbeats PRIVATE heartbeat_validation) +target_compile_options(device_heartbeats PRIVATE -Wall -Wextra -Wpedantic) +enable_testing() +add_executable(validation_tests test/validation.cpp) +target_link_libraries(validation_tests PRIVATE heartbeat_validation) +add_test(NAME validation COMMAND validation_tests) diff --git a/applications/device-heartbeats/README.md b/applications/device-heartbeats/README.md new file mode 100644 index 00000000..4c559b96 --- /dev/null +++ b/applications/device-heartbeats/README.md @@ -0,0 +1,152 @@ +# Device Heartbeats + +A small C++20/Drogon API saves synthetic readings and advances each device's latest sequence in one transaction on [ClickHouse Managed Postgres (public beta)](https://clickhouse.com/docs/products/managed-postgres). Project tokens select seeded devices; matching sample UUIDs replay retained results. This is an instructional ingestion protocol, without hardware integration or monitoring/alerting behavior. + +## Transaction and trust model + +A write locks the `(project, device)` parent, checks a retained UUID **before** checking the latest sequence or 200-sample cap, inserts history, and updates the parent's latest sequence/count. Every query uses the same Drogon transaction. A retained request replays only when sequence, integer reading and canonical optional observation time match. A different UUID must carry a strictly greater sequence. No history is silently pruned. + +Drogon 1.9.13 queues COMMIT when its transaction is destroyed. The final SQL callback releases our transaction pointer, breaking the callback/state ownership cycle. Only the COMMIT callback can send 201/200. Rollback paths resolve HTTP independently because Drogon does not invoke that callback on rollback. Failed or uncertain commits return 503/`retry_sample_id`; retry the same UUID and payload to resolve the durable outcome. + +Tokens are distinct 32–128 character printable ASCII secrets for trusted project operators. They are not individual device identities. The server derives the project; no project field is accepted in JSON. The shared runtime database role is also trusted: it can directly insert history or change permitted parent metadata, bypassing the application's monotonic-sequence, true-count and history-cap protocol. Grants deny history UPDATE/DELETE, key edits and DDL; constraints enforce UUID/sequence uniqueness, positive sequence, reading bounds, bounded parent counters and the composite device foreign key. They do not independently enforce consistency between history and parent metadata. + +## Build on native Linux + +Tested 2 October 2026 on Ubuntu 24.04 ARM64: GCC 13.3.0, CMake 3.28.3, Drogon 1.9.13, Trantor revision below, JsonCpp 1.9.5, OpenSSL 3.0.13 and libpq 16.15. PostgreSQL server was 18.6. Framework sources are pinned; distro packages provide native libraries. Build outside shared host mounts: + +```sh +sudo apt-get update +sudo apt-get install -y build-essential cmake git pkg-config libjsoncpp-dev \ + libssl-dev libpq-dev uuid-dev zlib1g-dev libbrotli-dev postgresql-client \ + curl ca-certificates clang-format +git clone --branch v1.9.13 --depth 1 --recurse-submodules \ + https://github.com/drogonframework/drogon.git "$HOME/drogon-source" +test "$(git -C "$HOME/drogon-source" rev-parse HEAD)" = 4c5430757ea5451a7c38fbbef4b4bef7dbb47f2f +test "$(git -C "$HOME/drogon-source/trantor" rev-parse HEAD)" = 63a4e5e164e219dc3bf30cdbfa1462ae5602fa97 +cmake -S "$HOME/drogon-source" -B "$HOME/drogon-build" \ + -DCMAKE_BUILD_TYPE=Release -DCMAKE_INSTALL_PREFIX="$HOME/drogon" \ + -DBUILD_TESTING=OFF -DBUILD_EXAMPLES=OFF -DBUILD_CTL=OFF -DBUILD_ORM=ON \ + -DUSE_POSTGRESQL=ON -DUSE_MYSQL=OFF -DUSE_SQLITE3=OFF -DUSE_REDIS=OFF +cmake --build "$HOME/drogon-build" --parallel 2 +cmake --install "$HOME/drogon-build" +cd /path/to/examples/applications/device-heartbeats +cmake -S . -B build -DCMAKE_BUILD_TYPE=Release -DCMAKE_PREFIX_PATH="$HOME/drogon" +cmake --build build --parallel 2 +clang-format --dry-run --Werror src/*.cpp src/*.hpp test/validation.cpp +ctest --test-dir build --output-on-failure +python3 test/preflight.py "$PWD/build/device_heartbeats" +``` + +The preflight mode exposes only `/health` without a database. Normal startup requires and verifies the configured runtime connection before opening the listener. CTest runs one executable containing three validation groups; the separate preflight exercises the compiled listener and waits for successful SIGTERM exit. CI repeats these local checks without Cloud credentials. + +## Create your dedicated Cloud fixture + +Use an authenticated [clickhousectl](https://github.com/ClickHouse/clickhousectl) CLI. Creation starts billing; delete your fixture after testing. Check the current supported region/shape for your organization. This modest no-HA shape was supported by the tested service: + +```sh +umask 077 +export ORG_ID=your-clickhouse-organization-id +clickhousectl cloud postgres create --org-id "$ORG_ID" \ + --name device-heartbeats-demo --provider aws --region us-east-1 \ + --size c6gd.large --pg-version 18 --ha-type none --json > /private/path/create.json +export PG_ID=your-created-service-id +clickhousectl cloud postgres get "$PG_ID" --org-id "$ORG_ID" +# Repeat get until state is running, then fetch the official PEM: +clickhousectl cloud postgres certs get "$PG_ID" --org-id "$ORG_ID" \ + --output /private/path/ca.pem +``` + +Keep the returned ID, hostname, username and one-time password privately. `--output` writes PEM even when CLI stdout is automatically JSON. Refresh the official bundle after certificate rotation. The app passes libpq `sslmode=verify-full`, the explicit hostname and `sslrootcert`; the full official bundle worked unchanged. Never replace verification with `require`. + +## Bootstrap, migrate and seed + +Copy `.env.example` to `/private/path/setup.env`, mode 600, outside the checkout. Fill in Cloud connection fields, `ADMIN_USER`/`ADMIN_PASSWORD` from creation, a separate `MIGRATION_PASSWORD`, `PGPASSWORD` for `heartbeats_app`, and two distinct project tokens. Paths must be absolute. Shell-safe quote private values when necessary. Export fields to child processes: + +```sh +set -a; source /private/path/setup.env; set +a +export PGSSLMODE=verify-full +APP_PASSWORD=$PGPASSWORD +export PGUSER=$ADMIN_USER PGPASSWORD=$ADMIN_PASSWORD +psql -X -v MIGRATION_PASSWORD="$MIGRATION_PASSWORD" -v APP_PASSWORD="$APP_PASSWORD" \ + -f sql/bootstrap.sql +export PGUSER=heartbeats_migration PGPASSWORD=$MIGRATION_PASSWORD +psql -X -f sql/migrate.sql +psql -X -f sql/grants.sql +psql -X -f sql/seed.sql +``` + +Bootstrap creates the roles/schema and revokes PUBLIC database CREATE/TEMP and public-schema CREATE on this dedicated fixture. The schema owner needs no database CREATE. Versioned SQL uses an advisory transaction lock; migration and seed are repeatable. Seed adds north/south projects and meter-a/b/c without resetting stored readings. The server never runs migrations or runtime DDL. + +The runtime gets SELECT on devices/samples, INSERT on samples, identity-sequence USAGE and UPDATE only on latest_sequence/sample_count. It cannot read migration history. Both project scopes share this role; project authorization lives in the API. + +## Run and try + +```sh +export PGUSER=heartbeats_app PGPASSWORD=$APP_PASSWORD +unset ADMIN_USER ADMIN_PASSWORD MIGRATION_PASSWORD +./build/device_heartbeats +``` + +The listener binds `127.0.0.1:4000`; PORT can change the local port. `/health` reports process readiness after startup, without continuously probing the database. SIGTERM stops the server and closes the pool. Create a private mode-600 runtime.env containing only PGHOST, PGPORT, PGDATABASE, PGUSER=heartbeats_app, its PGPASSWORD, PGSSLROOTCERT, NORTH_TOKEN, SOUTH_TOKEN and optional PORT. In a second shell: + +```sh +set -a; source /private/path/runtime.env; set +a +curl -sS http://127.0.0.1:4000/devices -H "Authorization: Bearer $NORTH_TOKEN" +curl -sS http://127.0.0.1:4000/devices/meter-a/samples \ + -H "Authorization: Bearer $NORTH_TOKEN" -H 'Content-Type: application/json' \ + --data '{"sample_id":"aaaaaaaa-1234-4567-890a-123456789abc","sequence":"9007199254740993","reading":12,"observed_at":"2026-10-02T12:13:14Z"}' +curl -sS 'http://127.0.0.1:4000/devices/meter-a/samples?limit=25' \ + -H "Authorization: Bearer $NORTH_TOKEN" +``` + +New samples return 201; the same canonical payload returns 200/`replay=true` with identical retained fields. Changed payload, stale sequence or full history return 409. Unknown project-scoped device writes return 404; history reads return an empty page. Invalid credentials return 401. JSON media type is compared exactly, case-insensitively, before optional `;` parameters; duplicate keys, comments and trailing commas reject. + +| Route | Bounds | +| --- | --- | +| GET `/devices?limit=25&after=meter-a` | ASCII device key ascending, scoped to token | +| POST `/devices/{device}/samples` | Required sample_id, sequence, reading; optional observed_at; no other fields | +| GET `/devices/{device}/samples?limit=25&before=123` | Numeric identity descending; before/IDs are decimal strings | + +Lists default to 25, cap at 100 and return last-row `next_after`/`next_before`; finish on an empty page. Numeric ordering qualifies `heartbeats.samples.id`, avoiding the `id::text` output alias. Identity order is allocation order, not commit chronology; gaps are normal. Pages are live reads, not a snapshot. + +Sequence is a decimal string 1–9223372036854775807, at most 19 digits, with no leading zero. Reading is an exact JSON integer from -1,000,000 through 1,000,000; booleans/floating-point JSON values reject. UUIDs normalize to lowercase. Observation time is null/omitted or whole-second UTC `YYYY-MM-DDTHH:MM:SSZ`, years 1970–2100 with a valid calendar date. Database `clock_timestamp()` records receipt when history is inserted after the parent lock; it is not device observation time or a commit timestamp. + +Limits: 4 KiB body, 32 connections, two HTTP event-loop threads, ten-second connection idle timeout, 32 requests per keepalive connection, four pipelined requests, 12 admitted authenticated operations and eight admitted writes. The DB pool has four connections; Drogon's client acquisition/query timeout is five seconds, statement timeout four seconds, lock timeout two seconds and idle-in-transaction timeout six seconds. The connection string also sets `connect_timeout=5`, but libpq ignores that option during asynchronous `PQconnectPoll`; it is not an additional connection deadline here. These are individual controls, not a guaranteed total HTTP deadline. Runtime SQL is asynchronous; only startup connectivity is synchronous. + +## Real Cloud acceptance + +Use only a disposable dedicated fixture: tests add devices, triggers and owner-controlled history/counters. Stop the server. Restore setup credentials, reset, then rerun bootstrap/migration/grants/seed above (migration/seed can be repeated): + +```sh +set -a; source /private/path/setup.env; set +a +export PGSSLMODE=verify-full +APP_PASSWORD=$PGPASSWORD +export PGUSER=$ADMIN_USER PGPASSWORD=$ADMIN_PASSWORD +psql -X -f sql/cleanup.sql +# Repeat bootstrap, migration twice, grants and seed twice as above. +export PGUSER=heartbeats_app PGPASSWORD=$APP_PASSWORD +openssl req -x509 -newkey rsa:2048 -nodes -days 2 -subj /CN=UnrelatedAcceptanceCA \ + -keyout /private/path/wrong-ca.key -out /private/path/wrong-ca.pem +export WRONG_CA=/private/path/wrong-ca.pem +export HEARTBEATS_EXECUTABLE="$PWD/build/device_heartbeats" +export EVIDENCE_DIR=/private/path/acceptance +python3 test/cloud.py +``` + +The standard-library Python helper uses `psql` for independent controls. Its C++ child receives only runtime fields. Nine cases cover scoped validation/64-bit boundaries, actual two-session lock waits and UUID replay, competing sequences, ten post-history write failures followed by recovery, deferred COMMIT failure followed by same-ID 201/200, all 200 history rows in numeric pages, grants/constraints and the trusted-role bypass, actual wrong-CA/wrong-host certificate diagnostics with a positive DNS control, and old-process exit before durable restart/replay. Temporary fault triggers are owner-only and removed without widening runtime grants. Owner-accelerated full-history fixtures test the cap without implying device throughput. + +## Cleanup + +Stop the app. Optional schema destruction requires administrator credentials restored explicitly after the runtime steps unset them: + +```sh +set -a; source /private/path/setup.env; set +a +export PGUSER=$ADMIN_USER PGPASSWORD=$ADMIN_PASSWORD PGSSLMODE=verify-full +psql -X -f sql/cleanup.sql +clickhousectl cloud postgres delete "$PG_ID" --org-id "$ORG_ID" +clickhousectl cloud postgres list --org-id "$ORG_ID" +``` + +Confirm your exact ID is absent. Cleanup leaves PUBLIC revocations in place and is fixture destruction, not a production downgrade. Retained history is bounded rather than evicted; longer-lived ingestion needs an explicit archival/replay policy. + +Primary references: [Drogon 1.9.13 DbClient API](https://github.com/drogonframework/drogon/blob/v1.9.13/orm_lib/inc/drogon/orm/DbClient.h), [transaction implementation](https://github.com/drogonframework/drogon/blob/v1.9.13/orm_lib/src/TransactionImpl.cc), [libpq verified connection options](https://www.postgresql.org/docs/current/libpq-connect.html). diff --git a/applications/device-heartbeats/sql/bootstrap.sql b/applications/device-heartbeats/sql/bootstrap.sql new file mode 100644 index 00000000..65fb28d6 --- /dev/null +++ b/applications/device-heartbeats/sql/bootstrap.sql @@ -0,0 +1,9 @@ +\set ON_ERROR_STOP on +SELECT format('CREATE ROLE heartbeats_migration LOGIN PASSWORD %L', :'MIGRATION_PASSWORD') +WHERE NOT EXISTS (SELECT FROM pg_roles WHERE rolname='heartbeats_migration') \gexec +SELECT format('CREATE ROLE heartbeats_app LOGIN PASSWORD %L', :'APP_PASSWORD') +WHERE NOT EXISTS (SELECT FROM pg_roles WHERE rolname='heartbeats_app') \gexec +CREATE SCHEMA IF NOT EXISTS heartbeats AUTHORIZATION heartbeats_migration; +REVOKE CREATE ON SCHEMA public FROM PUBLIC; +SELECT format('REVOKE CREATE,TEMP ON DATABASE %I FROM PUBLIC',current_database()) \gexec +SELECT format('GRANT CONNECT ON DATABASE %I TO heartbeats_migration,heartbeats_app',current_database()) \gexec diff --git a/applications/device-heartbeats/sql/cleanup.sql b/applications/device-heartbeats/sql/cleanup.sql new file mode 100644 index 00000000..ef65bba3 --- /dev/null +++ b/applications/device-heartbeats/sql/cleanup.sql @@ -0,0 +1,7 @@ +\set ON_ERROR_STOP on +DROP SCHEMA IF EXISTS heartbeats CASCADE; +SELECT format('DROP OWNED BY %I',rolname) FROM pg_roles +WHERE rolname IN('heartbeats_app','heartbeats_migration') \gexec +DROP ROLE IF EXISTS heartbeats_app; +DROP ROLE IF EXISTS heartbeats_migration; +-- Dedicated fixture reset; PUBLIC revocations remain. Not a production downgrade. diff --git a/applications/device-heartbeats/sql/grants.sql b/applications/device-heartbeats/sql/grants.sql new file mode 100644 index 00000000..6334425d --- /dev/null +++ b/applications/device-heartbeats/sql/grants.sql @@ -0,0 +1,6 @@ +\set ON_ERROR_STOP on +GRANT USAGE ON SCHEMA heartbeats TO heartbeats_app; +GRANT SELECT ON heartbeats.devices,heartbeats.samples TO heartbeats_app; +GRANT INSERT ON heartbeats.samples TO heartbeats_app; +GRANT UPDATE(latest_sequence,sample_count) ON heartbeats.devices TO heartbeats_app; +GRANT USAGE ON SEQUENCE heartbeats.samples_id_seq TO heartbeats_app; diff --git a/applications/device-heartbeats/sql/migrate.sql b/applications/device-heartbeats/sql/migrate.sql new file mode 100644 index 00000000..8f1e008a --- /dev/null +++ b/applications/device-heartbeats/sql/migrate.sql @@ -0,0 +1,36 @@ +\set ON_ERROR_STOP on +BEGIN; +SELECT pg_advisory_xact_lock(721149); +CREATE TABLE IF NOT EXISTS heartbeats.schema_versions ( + version integer PRIMARY KEY, + applied_at timestamptz NOT NULL DEFAULT now() +); +SELECT NOT EXISTS(SELECT FROM heartbeats.schema_versions WHERE version=1) AS apply_v1 \gset +\if :apply_v1 +CREATE TABLE heartbeats.projects ( + id text PRIMARY KEY CHECK(id ~ '^[a-z][a-z0-9-]{0,31}$') +); +CREATE TABLE heartbeats.devices ( + project text NOT NULL REFERENCES heartbeats.projects(id), + id text NOT NULL CHECK(id ~ '^[a-z][a-z0-9-]{0,31}$'), + latest_sequence bigint NOT NULL DEFAULT 0 CHECK(latest_sequence>=0), + sample_count integer NOT NULL DEFAULT 0 CHECK(sample_count BETWEEN 0 AND 200), + PRIMARY KEY(project,id) +); +CREATE TABLE heartbeats.samples ( + id bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY, + project text NOT NULL, + device text NOT NULL, + sample_id uuid NOT NULL, + sequence bigint NOT NULL CHECK(sequence>0), + reading integer NOT NULL CHECK(reading BETWEEN -1000000 AND 1000000), + observed_at timestamptz, + received_at timestamptz NOT NULL DEFAULT clock_timestamp(), + CONSTRAINT sample_device FOREIGN KEY(project,device) REFERENCES heartbeats.devices(project,id), + CONSTRAINT retained_sample UNIQUE(project,device,sample_id), + CONSTRAINT device_sequence UNIQUE(project,device,sequence) +); +CREATE INDEX samples_device_id ON heartbeats.samples(project,device,id DESC); +INSERT INTO heartbeats.schema_versions(version) VALUES(1); +\endif +COMMIT; diff --git a/applications/device-heartbeats/sql/seed.sql b/applications/device-heartbeats/sql/seed.sql new file mode 100644 index 00000000..5e3fa60d --- /dev/null +++ b/applications/device-heartbeats/sql/seed.sql @@ -0,0 +1,8 @@ +\set ON_ERROR_STOP on +BEGIN; +INSERT INTO heartbeats.projects VALUES('north'),('south') ON CONFLICT DO NOTHING; +INSERT INTO heartbeats.devices(project,id) +SELECT project,device FROM (VALUES('north'),('south')) AS p(project), + (VALUES('meter-a'),('meter-b'),('meter-c')) AS d(device) +ON CONFLICT DO NOTHING; +COMMIT; diff --git a/applications/device-heartbeats/src/server.cpp b/applications/device-heartbeats/src/server.cpp new file mode 100644 index 00000000..9afd1be9 --- /dev/null +++ b/applications/device-heartbeats/src/server.cpp @@ -0,0 +1,205 @@ +#include "service.hpp" +#include +#include +#include +#include +#include +#include + +static std::string required(const char *name) { + const auto value = std::getenv(name); + if (!value || !*value) + throw std::runtime_error(std::string("Missing ") + name); + return value; +} +static std::string quoted(const std::string &value) { + std::string result = "'"; + for (char c : value) { + if (c == '\\' || c == '\'') + result += '\\'; + result += c; + } + return result + "'"; +} +static std::string connectionInfo() { + std::string result; + for (const auto &[option, env] : + std::map{{"host", "PGHOST"}, + {"port", "PGPORT"}, + {"dbname", "PGDATABASE"}, + {"user", "PGUSER"}, + {"password", "PGPASSWORD"}, + {"sslrootcert", "PGSSLROOTCERT"}}) + result += option + "=" + quoted(required(env.c_str())) + " "; + return result + "sslmode=verify-full connect_timeout=5 application_name=device-heartbeats " + "options='-c statement_timeout=4000 -c lock_timeout=2000 -c " + "idle_in_transaction_session_timeout=6000'"; +} +static Json::Value error(const std::string &code) { + Json::Value body; + body["error"] = code; + return body; +} +static void queryFields(const drogon::HttpRequestPtr &request, + const std::set &allowed) { + for (const auto &[key, value] : request->parameters()) { + (void)value; + if (!allowed.count(key)) + throw ApiError(400, "invalid_query"); + } +} +int main(int argc, char **argv) { + try { + const bool preflight = argc == 2 && std::string(argv[1]) == "--preflight"; + const bool check = argc == 2 && std::string(argv[1]) == "--check-db"; + if (argc > 1 && !preflight && !check) + throw std::runtime_error("Unknown option"); + auto &app = drogon::app(); + app.setLogLevel(trantor::Logger::kWarn); + drogon::orm::DbClientPtr database; + std::shared_ptr service; + if (!preflight) { + if (required("PGUSER") != "heartbeats_app") + throw std::runtime_error("Runtime requires heartbeats_app"); + database = drogon::orm::DbClient::newPgClient(connectionInfo(), 4); + database->setTimeout(5.0); + // Synchronous startup only; no event loop is blocked by runtime handlers. + database->execSqlSync("SELECT 1"); + if (check) { + database->closeAll(); + std::cout << "Verified database connection\n"; + return 0; + } + service = std::make_shared(database); + } + int listenPort = 4000; + if (const auto configured = std::getenv("PORT")) { + const std::string text = configured; + if (text.empty() || text.size() > 5 || + !std::all_of(text.begin(), text.end(), [](char c) { return c >= '0' && c <= '9'; })) + throw std::runtime_error("Invalid PORT"); + listenPort = std::stoi(text); + if (listenPort < 1024 || listenPort > 65535) + throw std::runtime_error("Invalid PORT"); + } + std::map tokens; + if (!preflight) { + tokens = {{"north", required("NORTH_TOKEN")}, {"south", required("SOUTH_TOKEN")}}; + const std::regex pattern("[!-~]{32,128}"); + for (const auto &[project, token] : tokens) { + (void)project; + if (!std::regex_match(token, pattern)) + throw std::runtime_error("Invalid project token"); + } + if (tokens.at("north") == tokens.at("south")) + throw std::runtime_error("Tokens must differ"); + } + auto active = std::make_shared>(0); + auto guarded = [tokens, active](const drogon::HttpRequestPtr &request, Reply reply, + auto operation) { + std::string project; + for (const auto &[id, token] : tokens) + if (equalToken(request->getHeader("authorization"), "Bearer " + token)) + project = id; + if (project.empty()) { + reply(jsonResponse(401, error("unauthorized"))); + return; + } + if (active->fetch_add(1) >= 12) { + --*active; + reply(jsonResponse(503, error("busy"))); + return; + } + auto completed = std::make_shared>(false); + Reply finish = [active, completed, + reply = std::move(reply)](const drogon::HttpResponsePtr &response) { + if (!completed->exchange(true)) { + --*active; + reply(response); + } + }; + try { + operation(project, finish); + } catch (const ApiError &e) { + finish(jsonResponse(e.status, error(e.what()))); + } catch (const std::exception &) { + finish(jsonResponse(500, error("request_failure"))); + } + }; + app.registerHandler("/health", + [](const drogon::HttpRequestPtr &, Reply &&reply) { + Json::Value body; + body["status"] = "up"; + reply(jsonResponse(200, body)); + }, + {drogon::Get}); + if (service) { + app.registerHandler( + "/devices", + [service, guarded](const drogon::HttpRequestPtr &request, Reply &&reply) { + guarded(request, std::move(reply), + [&](const std::string &project, Reply finish) { + queryFields(request, {"limit", "after"}); + const auto after = request->getParameter("after"); + if (!after.empty()) + identifier(after); + service->devices(project, pageLimit(request->getParameter("limit")), + after, std::move(finish)); + }); + }, + {drogon::Get}); + app.registerHandler( + "/devices/{1}/samples", + [service, guarded](const drogon::HttpRequestPtr &request, Reply &&reply, + const std::string &device) { + guarded( + request, std::move(reply), [&](const std::string &project, Reply finish) { + identifier(device); + if (request->method() == drogon::Post) { + queryFields(request, {}); + auto mediaType = request->getHeader("content-type"); + mediaType = mediaType.substr(0, mediaType.find(';')); + while (!mediaType.empty() && + (mediaType.back() == ' ' || mediaType.back() == '\t')) + mediaType.pop_back(); + while (!mediaType.empty() && + (mediaType.front() == ' ' || mediaType.front() == '\t')) + mediaType.erase(0, 1); + std::transform(mediaType.begin(), mediaType.end(), + mediaType.begin(), [](unsigned char c) { + return static_cast( + c >= 'A' && c <= 'Z' ? c + ('a' - 'A') : c); + }); + if (mediaType != "application/json") + throw ApiError(415, "json_required"); + service->submit(project, device, parseSample(request->body()), + std::move(finish)); + } else { + queryFields(request, {"limit", "before"}); + const auto before = request->getParameter("before"); + service->history( + project, device, pageLimit(request->getParameter("limit")), + before.empty() ? 0 : decimal(before), std::move(finish)); + } + }); + }, + {drogon::Get, drogon::Post}); + } + app.addListener("127.0.0.1", listenPort) + .setThreadNum(2) + .setMaxConnectionNum(32) + .setMaxConnectionNumPerIP(32) + .setClientMaxBodySize(4096) + .setClientMaxMemoryBodySize(4096) + .setIdleConnectionTimeout(10) + .setKeepaliveRequestsNumber(32) + .setPipeliningRequestsNumber(4); + app.run(); + if (database) + database->closeAll(); + return 0; + } catch (const std::exception &) { + std::cerr << "Startup failed; check private configuration and connection diagnostics\n"; + return 1; + } +} diff --git a/applications/device-heartbeats/src/service.cpp b/applications/device-heartbeats/src/service.cpp new file mode 100644 index 00000000..4bec35af --- /dev/null +++ b/applications/device-heartbeats/src/service.cpp @@ -0,0 +1,216 @@ +#include "service.hpp" +#include +using namespace drogon::orm; + +drogon::HttpResponsePtr jsonResponse(int status, const Json::Value &value) { + auto response = drogon::HttpResponse::newHttpJsonResponse(value); + response->setStatusCode(static_cast(status)); + response->addHeader("Cache-Control", "no-store"); + return response; +} +static Json::Value errorBody(const std::string &code) { + Json::Value value; + value["error"] = code; + return value; +} +static const std::string projection = + "id::text, sample_id::text, sequence::text, reading, " + "to_char(observed_at AT TIME ZONE 'UTC','YYYY-MM-DD\"T\"HH24:MI:SS\"Z\"') AS observed_at, " + "to_char(received_at AT TIME ZONE 'UTC','YYYY-MM-DD\"T\"HH24:MI:SS.US\"Z\"') AS received_at"; +static Json::Value sampleRow(const Row &row) { + Json::Value value; + for (const auto &name : {"id", "sample_id", "sequence", "received_at"}) + value[name] = row[name].as(); + value["reading"] = row["reading"].as(); + value["observed_at"] = row["observed_at"].isNull() + ? Json::Value() + : Json::Value(row["observed_at"].as()); + return value; +} +struct Lease { + std::shared_ptr> active; + ~Lease() { --*active; } +}; +class Write : public std::enable_shared_from_this { + std::string project_, device_; + SampleInput input_; + Reply reply_; + std::shared_ptr lease_; + std::shared_ptr transaction_; + bool responded_ = false; + Json::Value result_; + int status_ = 201; + int64_t latest_ = 0; + int count_ = 0; + + void finish(int status, const Json::Value &body) { + if (responded_) + return; + responded_ = true; + reply_(jsonResponse(status, body)); + } + void reject(int status, const std::string &code) { + if (transaction_) { + transaction_->rollback(); + transaction_.reset(); + } + // Drogon's commit callback is NOT invoked on rollback. Resolve this path independently. + finish(status, errorBody(code)); + } + auto onError() { + return [self = shared_from_this()](const DrogonDbException &) { + self->reject(503, "retry_sample_id"); + }; + } + void retained() { + auto self = shared_from_this(); + transaction_->execSqlAsync( + "SELECT " + projection + + " FROM heartbeats.samples WHERE project=$1 AND device=$2 AND sample_id=$3::uuid", + [self](const Result &rows) { + if (!rows.empty()) { + const auto saved = sampleRow(rows[0]); + const auto observed = self->input_.observedAt + ? Json::Value(*self->input_.observedAt) + : Json::Value(); + if (saved["sequence"].asString() != std::to_string(self->input_.sequence) || + saved["reading"].asInt() != self->input_.reading || + saved["observed_at"] != observed) { + self->reject(409, "sample_conflict"); + return; + } + self->result_["sample"] = saved; + self->result_["replay"] = true; + self->status_ = 200; + self->releaseForCommit(); + } else if (self->input_.sequence <= self->latest_) { + self->reject(409, "stale_sequence"); + } else if (self->count_ >= 200) { + self->reject(409, "history_full"); + } else { + self->insert(); + } + }, + onError(), project_, device_, input_.sampleId); + } + void insert() { + auto self = shared_from_this(); + transaction_->execSqlAsync( + "INSERT INTO heartbeats.samples(project,device,sample_id,sequence,reading,observed_at) " + "VALUES($1,$2,$3::uuid,$4,$5,NULLIF($6,'')::timestamptz) RETURNING " + + projection, + [self](const Result &rows) { + self->result_["sample"] = sampleRow(rows[0]); + self->result_["replay"] = false; + self->update(); + }, + onError(), project_, device_, input_.sampleId, input_.sequence, input_.reading, + input_.observedAt.value_or("")); + } + void update() { + auto self = shared_from_this(); + transaction_->execSqlAsync( + "UPDATE heartbeats.devices SET latest_sequence=$3,sample_count=sample_count+1 " + "WHERE project=$1 AND id=$2", + [self](const Result &) { self->releaseForCommit(); }, onError(), project_, device_, + input_.sequence); + } + void releaseForCommit() { + // Break the ownership cycle: transaction owns its commit callback, callback owns this + // state. Its destructor queues COMMIT. Only that callback can send successful HTTP status. + transaction_.reset(); + } + + public: + Write(std::string project, std::string device, SampleInput input, Reply reply, + std::shared_ptr lease) + : project_(std::move(project)), device_(std::move(device)), input_(std::move(input)), + reply_(std::move(reply)), lease_(std::move(lease)) {} + void begin(const std::shared_ptr &transaction) { + if (!transaction) { + finish(503, errorBody("pool_timeout")); + return; + } + transaction_ = transaction; + auto self = shared_from_this(); + transaction_->setCommitCallback([self](bool committed) { + if (committed) + self->finish(self->status_, self->result_); + else + self->finish(503, errorBody("retry_sample_id")); + }); + transaction_->execSqlAsync( + "SELECT latest_sequence::text,sample_count FROM heartbeats.devices WHERE project=$1 " + "AND id=$2 FOR UPDATE", + [self](const Result &rows) { + if (rows.empty()) { + self->reject(404, "unknown_device"); + return; + } + self->latest_ = std::stoll(rows[0]["latest_sequence"].as()); + self->count_ = rows[0]["sample_count"].as(); + self->retained(); + }, + onError(), project_, device_); + } +}; +Service::Service(DbClientPtr client) + : client_(std::move(client)), active_(std::make_shared>(0)) {} +void Service::submit(const std::string &project, const std::string &device, SampleInput input, + Reply reply) { + if (active_->fetch_add(1) >= 8) { + --*active_; + reply(jsonResponse(503, errorBody("busy"))); + return; + } + auto lease = std::make_shared(); + lease->active = active_; + auto state = std::make_shared(project, device, std::move(input), std::move(reply), + std::move(lease)); + client_->newTransactionAsync( + [state](const std::shared_ptr &transaction) { state->begin(transaction); }); +} +void Service::devices(const std::string &project, int limit, const std::string &after, + Reply reply) { + client_->execSqlAsync( + "SELECT id,latest_sequence::text,sample_count FROM heartbeats.devices " + "WHERE project=$1 AND id>$2 ORDER BY id LIMIT $3::integer", + [reply](const Result &rows) { + Json::Value result, values(Json::arrayValue); + for (const auto &row : rows) { + Json::Value value; + value["device"] = row["id"].as(); + value["latest_sequence"] = row["latest_sequence"].as(); + value["sample_count"] = row["sample_count"].as(); + values.append(value); + } + result["rows"] = values; + result["next_after"] = + values.empty() ? Json::Value() : values[values.size() - 1]["device"]; + reply(jsonResponse(200, result)); + }, + [reply](const DrogonDbException &) { + reply(jsonResponse(503, errorBody("database_unavailable"))); + }, + project, after, limit); +} +void Service::history(const std::string &project, const std::string &device, int limit, + int64_t before, Reply reply) { + client_->execSqlAsync( + "SELECT " + projection + + " FROM heartbeats.samples WHERE project=$1 AND device=$2 " + "AND ($3::bigint=0 OR id<$3) ORDER BY heartbeats.samples.id DESC LIMIT $4::integer", + [reply](const Result &rows) { + Json::Value result, values(Json::arrayValue); + for (const auto &row : rows) + values.append(sampleRow(row)); + result["rows"] = values; + result["next_before"] = + values.empty() ? Json::Value() : values[values.size() - 1]["id"]; + reply(jsonResponse(200, result)); + }, + [reply](const DrogonDbException &) { + reply(jsonResponse(503, errorBody("database_unavailable"))); + }, + project, device, before, limit); +} diff --git a/applications/device-heartbeats/src/service.hpp b/applications/device-heartbeats/src/service.hpp new file mode 100644 index 00000000..4566a99f --- /dev/null +++ b/applications/device-heartbeats/src/service.hpp @@ -0,0 +1,22 @@ +#pragma once +#include "validation.hpp" +#include +#include +#include +#include +#include + +using Reply = std::function; +drogon::HttpResponsePtr jsonResponse(int status, const Json::Value &value); +class Service { + drogon::orm::DbClientPtr client_; + std::shared_ptr> active_; + + public: + explicit Service(drogon::orm::DbClientPtr client); + void submit(const std::string &project, const std::string &device, SampleInput input, + Reply reply); + void devices(const std::string &project, int limit, const std::string &after, Reply reply); + void history(const std::string &project, const std::string &device, int limit, int64_t before, + Reply reply); +}; diff --git a/applications/device-heartbeats/src/validation.cpp b/applications/device-heartbeats/src/validation.cpp new file mode 100644 index 00000000..622febb3 --- /dev/null +++ b/applications/device-heartbeats/src/validation.cpp @@ -0,0 +1,102 @@ +#include "validation.hpp" +#include +#include +#include +#include +#include +#include +#include + +std::string identifier(const std::string &text) { + static const std::regex expression("[a-z][a-z0-9-]{0,31}"); + if (text.size() > 32 || !std::regex_match(text, expression)) + throw ApiError(400, "invalid_identifier"); + return text; +} +std::string uuid(const std::string &text) { + static const std::regex expression( + "[0-9a-fA-F]{8}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{12}"); + if (text.size() != 36 || !std::regex_match(text, expression)) + throw ApiError(400, "invalid_sample_id"); + auto result = text; + std::transform(result.begin(), result.end(), result.begin(), [](unsigned char c) { + return static_cast(c >= 'A' && c <= 'F' ? c + ('a' - 'A') : c); + }); + return result; +} +int64_t decimal(const std::string &text) { + if (text.empty() || text.size() > 19 || text[0] == '0' || + !std::all_of(text.begin(), text.end(), [](char c) { return c >= '0' && c <= '9'; })) + throw ApiError(400, "invalid_sequence"); + int64_t result = 0; + auto converted = std::from_chars(text.data(), text.data() + text.size(), result); + if (converted.ec != std::errc() || converted.ptr != text.data() + text.size()) + throw ApiError(400, "invalid_sequence"); + return result; +} +static std::string observation(const std::string &text) { + static const std::regex expression("[0-9]{4}-[0-9]{2}-[0-9]{2}T[0-9]{2}:[0-9]{2}:[0-9]{2}Z"); + if (text.size() != 20 || !std::regex_match(text, expression)) + throw ApiError(400, "invalid_observed_at"); + auto number = [&text](size_t at, size_t count) { return std::stoi(text.substr(at, count)); }; + const int y = number(0, 4); + const auto date = std::chrono::year_month_day{ + std::chrono::year{y}, std::chrono::month{static_cast(number(5, 2))}, + std::chrono::day{static_cast(number(8, 2))}}; + if (y < 1970 || y > 2100 || !date.ok() || number(11, 2) > 23 || number(14, 2) > 59 || + number(17, 2) > 59) + throw ApiError(400, "invalid_observed_at"); + return text; +} +SampleInput parseSample(std::string_view body) { + if (body.size() > 4096) + throw ApiError(413, "body_too_large"); + Json::CharReaderBuilder builder; + builder["rejectDupKeys"] = true; + builder["failIfExtra"] = true; + builder["allowComments"] = false; + builder["allowTrailingCommas"] = false; + builder["strictRoot"] = true; + std::unique_ptr reader(builder.newCharReader()); + Json::Value value; + std::string errors; + if (!reader->parse(body.data(), body.data() + body.size(), &value, &errors) || + !value.isObject()) + throw ApiError(400, "invalid_json"); + const std::set fields{"sample_id", "sequence", "reading", "observed_at"}; + for (const auto &name : value.getMemberNames()) + if (!fields.count(name)) + throw ApiError(400, "invalid_fields"); + if (!value["sample_id"].isString() || !value["sequence"].isString()) + throw ApiError(400, "invalid_fields"); + const auto &reading = value["reading"]; + if (reading.type() != Json::intValue && reading.type() != Json::uintValue) + throw ApiError(400, "invalid_reading"); + if (!reading.isInt64() || reading.asInt64() < -1000000 || reading.asInt64() > 1000000) + throw ApiError(400, "invalid_reading"); + std::optional observed; + if (value.isMember("observed_at") && !value["observed_at"].isNull()) { + if (!value["observed_at"].isString()) + throw ApiError(400, "invalid_observed_at"); + observed = observation(value["observed_at"].asString()); + } + return {uuid(value["sample_id"].asString()), decimal(value["sequence"].asString()), + static_cast(reading.asInt64()), observed}; +} +int pageLimit(const std::string &text) { + if (text.empty()) + return 25; + if (text.size() > 3) + throw ApiError(400, "invalid_limit"); + int result = 0; + auto converted = std::from_chars(text.data(), text.data() + text.size(), result); + if (converted.ec != std::errc() || converted.ptr != text.data() + text.size() || result < 1 || + result > 100) + throw ApiError(400, "invalid_limit"); + return result; +} +bool equalToken(const std::string &left, const std::string &right) { + if (left.size() != right.size()) + return false; + return CRYPTO_memcmp(left.data(), right.data(), left.size()) == 0; +} diff --git a/applications/device-heartbeats/src/validation.hpp b/applications/device-heartbeats/src/validation.hpp new file mode 100644 index 00000000..578752b9 --- /dev/null +++ b/applications/device-heartbeats/src/validation.hpp @@ -0,0 +1,24 @@ +#pragma once +#include +#include +#include +#include +#include +#include + +struct ApiError : std::runtime_error { + int status; + ApiError(int status, const std::string &code) : std::runtime_error(code), status(status) {} +}; +struct SampleInput { + std::string sampleId; + int64_t sequence; + int reading; + std::optional observedAt; +}; +std::string identifier(const std::string &text); +std::string uuid(const std::string &text); +int64_t decimal(const std::string &text); +SampleInput parseSample(std::string_view body); +int pageLimit(const std::string &text); +bool equalToken(const std::string &left, const std::string &right); diff --git a/applications/device-heartbeats/test/cloud.py b/applications/device-heartbeats/test/cloud.py new file mode 100644 index 00000000..cc758f54 --- /dev/null +++ b/applications/device-heartbeats/test/cloud.py @@ -0,0 +1,301 @@ +"""Destructive acceptance against your own dedicated Cloud fixture; standard-library tools only.""" +import concurrent.futures +import json +import os +from pathlib import Path +import re +import signal +import socket +import subprocess +import time +import urllib.error +import urllib.request +import uuid + +EXECUTABLE = os.environ['HEARTBEATS_EXECUTABLE'] +OUTPUT = Path(os.environ.get('EVIDENCE_DIR', '.')) +OUTPUT.mkdir(parents=True, exist_ok=True) +RUNTIME_FIELDS = ['PGHOST', 'PGPORT', 'PGDATABASE', 'PGUSER', 'PGPASSWORD', + 'PGSSLROOTCERT', 'NORTH_TOKEN', 'SOUTH_TOKEN', 'PORT'] +RUNTIME = {field: os.environ[field] for field in RUNTIME_FIELDS} +server = None +server_log = None + + +def database(sql, role='owner', check=True): + env = dict(os.environ, PGSSLMODE='verify-full') + if role == 'owner': + env.update(PGUSER='heartbeats_migration', PGPASSWORD=os.environ['MIGRATION_PASSWORD']) + elif role == 'admin': + env.update(PGUSER=os.environ['ADMIN_USER'], PGPASSWORD=os.environ['ADMIN_PASSWORD']) + elif role != 'runtime': + raise ValueError('unknown role') + result = subprocess.run(['psql', '-X', '-qAt', '-v', 'ON_ERROR_STOP=1', '-v', 'VERBOSITY=verbose'], + input=sql, text=True, capture_output=True, env=env, timeout=20) + if check and result.returncode: + raise RuntimeError(result.stderr) + return result + + +def call(device=None, body=None, path=None, south=False, raw=None, media='application/json', auth=True): + target = path or f'/devices/{device}/samples' + data = raw if raw is not None else (json.dumps(body).encode() if body is not None else None) + headers = {'Content-Type': media} + if auth: + headers['Authorization'] = 'Bearer ' + RUNTIME['SOUTH_TOKEN' if south else 'NORTH_TOKEN'] + request = urllib.request.Request('http://127.0.0.1:' + RUNTIME['PORT'] + target, + data=data, headers=headers) + try: + response = urllib.request.urlopen(request, timeout=12) + except urllib.error.HTTPError as error: + response = error + text = response.read() + try: + value = json.loads(text) + except json.JSONDecodeError: + value = {} + return response.status, value + + +def payload(sequence='1', reading=12, sample=None, observed=None): + value = {'sample_id': sample or str(uuid.uuid4()), 'sequence': sequence, 'reading': reading} + if observed is not None: + value['observed_at'] = observed + return value + + +def device(name, project='north'): + assert re.fullmatch(r'[a-z][a-z0-9-]{0,31}', name) and project in ('north', 'south') + database(f"INSERT INTO heartbeats.devices(project,id) VALUES('{project}','{name}')") + return name + + +def state(name): + return json.loads(database(f"""SELECT json_build_object('latest',latest_sequence::text,'count',sample_count, + 'rows',(SELECT count(*) FROM heartbeats.samples WHERE project='north' AND device='{name}')) + FROM heartbeats.devices WHERE project='north' AND id='{name}'""").stdout) + + +def start(): + global server, server_log + server_log = open(OUTPUT / f'server-{time.time_ns()}.log', 'w') + server = subprocess.Popen([EXECUTABLE], env=RUNTIME, stdout=server_log, stderr=subprocess.STDOUT) + for _ in range(100): + if server.poll() is not None: + raise RuntimeError('runtime startup failed; inspect private server log') + try: + if call(path='/health', auth=False)[0] == 200: + return + except OSError: + pass + time.sleep(.1) + raise RuntimeError('runtime did not become ready') + + +def stop(): + if server is not None and server.poll() is None: + server.send_signal(signal.SIGTERM) + assert server.wait(timeout=15) == 0 + if server_log is not None: + server_log.close() + + +def held_pair(name, first, second): + # An independent administrator psql process holds the parent row. Runtime requests + # use separate C++ pool connections; pg_blocking_pids proves actual database waits. + env = dict(os.environ, PGUSER=os.environ['ADMIN_USER'], PGPASSWORD=os.environ['ADMIN_PASSWORD'], + PGSSLMODE='verify-full') + holder = subprocess.Popen(['psql', '-X', '-qAt', '-v', 'ON_ERROR_STOP=1'], stdin=subprocess.PIPE, + stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True, env=env) + holder.stdin.write(f"BEGIN; SELECT id FROM heartbeats.devices WHERE project='north' AND id='{name}' FOR UPDATE;\n") + holder.stdin.flush() + assert holder.stdout.readline().strip() == name + with concurrent.futures.ThreadPoolExecutor(max_workers=2) as executor: + a = executor.submit(call, name, first) + b = executor.submit(call, name, second) + try: + for _ in range(60): + rows = database("""SELECT pg_stat_clear_snapshot(); + SELECT count(*) FROM pg_stat_activity WHERE application_name='device-heartbeats' + AND cardinality(pg_blocking_pids(pid))>0""", role='admin').stdout.splitlines() + if rows and rows[-1] == '2': + break + time.sleep(.02) + else: + raise AssertionError('two runtime database sessions were not observed blocked') + finally: + holder.stdin.write('COMMIT;\n'); holder.stdin.flush(); holder.stdin.close() + assert holder.wait(timeout=10) == 0 + return a.result(timeout=12), b.result(timeout=12) + + +def scoped_and_input(): + name = device('scope-case') + original = payload('9007199254740993', -12, observed='2026-10-02T12:13:14Z') + status, created = call(name, original) + assert status == 201 and created['sample']['sequence'] == original['sequence'] + assert created['sample']['observed_at'] == original['observed_at'] + assert created['sample']['received_at'] != original['observed_at'] + upper = dict(original, sample_id=original['sample_id'].upper()) + assert call(name, upper) == (200, dict(created, replay=True)) + assert call(name, dict(original, reading=-11))[0] == 409 + assert call(name, original, south=True)[0] == 404 + assert call(name, south=True)[1]['rows'] == [] + foreign = device('foreign-only', 'south') + assert call(foreign, payload())[0] == 404 + assert call(name, dict(payload(), project='south'))[0] == 400 + for reading in [True, 1.0, 1.5, 1000001, -1000001, None]: + assert call(name, payload(reading=reading))[0] == 400 + for sequence in ['0', '01', '-1', '9223372036854775808', '9' * 1000, 1, True]: + assert call(name, payload(sequence))[0] == 400 + for observed in ['2025-02-29T00:00:00Z', '2026-10-02T24:00:00Z', '2026-10-02T12:00:00+00:00', '\ud800']: + assert call(name, payload(observed=observed))[0] == 400 + assert call(name, payload(sample='\ud800'))[0] == 400 + assert call(name, raw=b'{bad')[0] == 400 + assert call(name, raw=b'{"sample_id":"a","sequence":"1","reading":1,}')[0] == 400 + assert call(name, raw=b'{}' + b' ' * 5000)[0] == 413 + assert call(name, payload(), auth=False)[0] == 401 + assert call(name, payload(), media='application/jsonp')[0] == 415 + assert call(name, payload(), media='application/jsonjunk')[0] == 415 + assert call(name, payload('9223372036854775807'), media='Application/JSON; charset=utf-8')[0] == 201 + assert call(path='/devices?limit=101')[0] == 400 + assert call(path=f'/devices/{name}/samples?before=9223372036854775808')[0] == 400 + assert state(name) == {'latest': '9223372036854775807', 'count': 2, 'rows': 2} + + +def concurrent_replay(): + name = device('replay-case') + original = payload('10') + a, b = held_pair(name, original, original) + assert sorted([a[0], b[0]]) == [200, 201] + assert a[1]['sample'] == b[1]['sample'] + assert state(name) == {'latest': '10', 'count': 1, 'rows': 1} + assert call(name, payload('9'))[0] == 409 + assert call(name, original)[0] == 200 + + +def competing_sequence(): + name = device('sequence-case') + a, b = held_pair(name, payload('42', 1), payload('42', 2)) + assert sorted([a[0], b[0]]) == [201, 409] + assert state(name) == {'latest': '42', 'count': 1, 'rows': 1} + assert call(name, payload('41'))[0] == 409 + assert call(name, payload('43'))[0] == 201 + + +def rollback_and_release(): + name = device('rollback-case') + database("""CREATE FUNCTION heartbeats.reject_latest() RETURNS trigger LANGUAGE plpgsql AS $$ BEGIN + IF NEW.id='rollback-case' THEN RAISE EXCEPTION 'acceptance latest failure'; END IF; RETURN NEW; END $$; + CREATE TRIGGER acceptance_latest BEFORE UPDATE ON heartbeats.devices + FOR EACH ROW EXECUTE FUNCTION heartbeats.reject_latest();""") + try: + for number in range(10): + assert call(name, payload(str(number + 1)))[0] == 503 + assert state(name) == {'latest': '0', 'count': 0, 'rows': 0} + finally: + database('DROP TRIGGER acceptance_latest ON heartbeats.devices; DROP FUNCTION heartbeats.reject_latest()') + assert call(name, payload())[0] == 201 # More errors than admission slots; no retained cycle/capacity leak. + assert state(name) == {'latest': '1', 'count': 1, 'rows': 1} + + +def deferred_commit(): + name = device('commit-case') + original = payload() + database("""CREATE FUNCTION heartbeats.reject_commit() RETURNS trigger LANGUAGE plpgsql AS $$ BEGIN + IF NEW.device='commit-case' THEN RAISE EXCEPTION 'acceptance deferred commit failure'; END IF; RETURN NEW; END $$; + CREATE CONSTRAINT TRIGGER acceptance_commit AFTER INSERT ON heartbeats.samples + DEFERRABLE INITIALLY DEFERRED FOR EACH ROW EXECUTE FUNCTION heartbeats.reject_commit();""") + try: + assert call(name, original)[0] == 503 # Final SQL succeeded, actual COMMIT failed; never201. + assert state(name) == {'latest': '0', 'count': 0, 'rows': 0} + finally: + database('DROP TRIGGER acceptance_commit ON heartbeats.samples; DROP FUNCTION heartbeats.reject_commit()') + assert call(name, original)[0] == 201 + assert call(name, original)[0] == 200 + + +def cap_and_history(): + name = device('cap-case') + database("""INSERT INTO heartbeats.samples(project,device,sample_id,sequence,reading) + SELECT 'north','cap-case',('60000000-0000-4000-8000-' || lpad(n::text,12,'0'))::uuid,n,n + FROM generate_series(1,200) AS g(n); + UPDATE heartbeats.devices SET latest_sequence=200,sample_count=200 WHERE project='north' AND id='cap-case';""") + saved = payload('1', 1, '60000000-0000-4000-8000-000000000001') + assert call(name, saved)[0] == 200 + assert call(name, payload('201'))[1]['error'] == 'history_full' + seen = [] + before = None + for _ in range(5): + path = f'/devices/{name}/samples?limit=100' + (f'&before={before}' if before else '') + status, value = call(path=path) + assert status == 200 and len(value['rows']) <= 100 + if not value['rows']: + break + seen.extend(row['id'] for row in value['rows']) + before = value['next_before'] + assert len(seen) == len(set(seen)) == 200 + assert seen == sorted(seen, key=int, reverse=True) + assert state(name) == {'latest': '200', 'count': 200, 'rows': 200} + + +def permissions(): + for statement in ["UPDATE heartbeats.samples SET reading=1", "DELETE FROM heartbeats.samples", + "UPDATE heartbeats.devices SET id='x'", 'CREATE TABLE heartbeats.forbidden(id int)', + 'CREATE SCHEMA forbidden', 'CREATE TEMP TABLE forbidden(id int)', + 'SELECT * FROM heartbeats.schema_versions']: + result = database(statement, role='runtime', check=False) + assert result.returncode and '42501' in result.stderr + for statement, code in [ + ("UPDATE heartbeats.devices SET latest_sequence=-1", '23514'), + ("UPDATE heartbeats.devices SET sample_count=201", '23514'), + ("INSERT INTO heartbeats.samples(project,device,sample_id,sequence,reading) VALUES('north','foreign-only','70000000-0000-4000-8000-000000000001',1,1)", '23503'), + ("INSERT INTO heartbeats.samples(project,device,sample_id,sequence,reading) VALUES('north','meter-a','70000000-0000-4000-8000-000000000002',0,1)", '23514')]: + result = database(statement, role='runtime', check=False) + assert result.returncode and code in result.stderr + first = database("INSERT INTO heartbeats.samples(project,device,sample_id,sequence,reading) VALUES('north','meter-a','70000000-0000-4000-8000-000000000003',1,1)", role='runtime') + assert first.returncode == 0 + duplicate = database("INSERT INTO heartbeats.samples(project,device,sample_id,sequence,reading) VALUES('north','meter-a','70000000-0000-4000-8000-000000000003',2,1)", role='runtime', check=False) + assert duplicate.returncode and '23505' in duplicate.stderr and 'retained_sample' in duplicate.stderr + # This direct trusted-role fixture deliberately does not update parent metadata. + + +def tls(): + def probe(env, filename): + result = subprocess.run([EXECUTABLE, '--check-db'], env=env, capture_output=True, text=True, timeout=15) + (OUTPUT / filename).write_text(result.stdout + result.stderr) + return result + wrong_ca = probe(dict(RUNTIME, PGSSLROOTCERT=os.environ['WRONG_CA']), 'tls-wrong-ca.log') + assert wrong_ca.returncode and re.search(r'certificate verify failed', wrong_ca.stdout + wrong_ca.stderr, re.I) + address = socket.gethostbyname(RUNTIME['PGHOST']) + wrong_host = probe(dict(RUNTIME, PGHOST=address), 'tls-wrong-host.log') + assert wrong_host.returncode and re.search(r'does not match host name', wrong_host.stdout + wrong_host.stderr, re.I) + positive = probe(RUNTIME, 'tls-positive.log') + assert positive.returncode == 0 and 'Verified database connection' in positive.stdout + + +def restart(): + name = device('restart-case') + original = payload('123', 55) + status, saved = call(name, original) + assert status == 201 + before = state(name) + old_pid = server.pid + stop() # wait/exact exit0 BEFORE starting replacement + start() + assert server.pid != old_pid + assert call(name, original) == (200, dict(saved, replay=True)) and state(name) == before + + +if __name__ == '__main__': + print('Actual PostgreSQL:', database('SHOW server_version').stdout.strip()) + start() + try: + cases = [scoped_and_input, concurrent_replay, competing_sequence, rollback_and_release, + deferred_commit, cap_and_history, permissions, tls, restart] + for case in cases: + case() + print('PASS', case.__name__, flush=True) + print(f'All {len(cases)} Cloud cases passed', flush=True) + finally: + stop() diff --git a/applications/device-heartbeats/test/preflight.py b/applications/device-heartbeats/test/preflight.py new file mode 100644 index 00000000..d7420f36 --- /dev/null +++ b/applications/device-heartbeats/test/preflight.py @@ -0,0 +1,29 @@ +"""Exercise the compiled listener and orderly exit without database credentials.""" +import os +import signal +import subprocess +import sys +import time +import urllib.request + +process = subprocess.Popen([sys.argv[1], '--preflight'], env=dict(os.environ, PORT='4100')) +try: + for attempt in range(100): + if process.poll() is not None: + raise RuntimeError('preflight listener exited before readiness') + try: + with urllib.request.urlopen('http://127.0.0.1:4100/health', timeout=1) as response: + assert response.status == 200 + assert response.read() == b'{"status":"up"}' + break + except OSError: + time.sleep(.05) + else: + raise RuntimeError('preflight listener did not become ready') + process.send_signal(signal.SIGTERM) + assert process.wait(timeout=10) == 0 + print('PASS compiled loopback readiness and SIGTERM exit') +finally: + if process.poll() is None: + process.kill() + process.wait(timeout=5) diff --git a/applications/device-heartbeats/test/validation.cpp b/applications/device-heartbeats/test/validation.cpp new file mode 100644 index 00000000..b2f8bd39 --- /dev/null +++ b/applications/device-heartbeats/test/validation.cpp @@ -0,0 +1,59 @@ +#include "validation.hpp" +#include +#include +#include +#include +static void require(bool truth) { + if (!truth) + throw std::runtime_error("assertion failed"); +} +static void rejects(const std::function &fn) { + try { + fn(); + } catch (const ApiError &) { + return; + } + throw std::runtime_error("expected validation rejection"); +} +int main() { + const std::string valid = + R"({"sample_id":"ABCDEF00-1234-4567-890A-123456789ABC","sequence":"9007199254740993","reading":-12,"observed_at":"2026-10-02T12:13:14Z"})"; + auto value = parseSample(valid); + require(value.sampleId == "abcdef00-1234-4567-890a-123456789abc"); + require(value.sequence == 9007199254740993LL && value.reading == -12); + for (const std::string reading : {"true", "1.0", "1.5", "1000001", "-1000001", "null"}) + rejects([&] { + parseSample("{\"sample_id\":\"abcdef00-1234-4567-890a-123456789abc\",\"sequence\":" + "\"1\",\"reading\":" + + reading + "}"); + }); + for (const std::string sequence : + std::vector{"0", "01", "-1", "9223372036854775808", std::string(1000, '9')}) + rejects([&] { decimal(sequence); }); + require(decimal("9223372036854775807") == INT64_MAX); + rejects([&] { + parseSample( + R"({"sample_id":"abcdef00-1234-4567-890a-123456789abc","sequence":"1","reading":1,"reading":2})"); + }); + rejects([&] { + parseSample( + R"({"sample_id":"abcdef00-1234-4567-890a-123456789abc","sequence":"1","reading":1,"project":"south"})"); + }); + rejects([&] { parseSample(R"({"sample_id":"\ud800","sequence":"1","reading":1})"); }); + rejects([&] { + parseSample( + R"({"sample_id":"abcdef00-1234-4567-890a-123456789abc","sequence":"1","reading":1,"observed_at":"2025-02-29T00:00:00Z"})"); + }); + rejects([&] { + parseSample( + R"({"sample_id":"abcdef00-1234-4567-890a-123456789abc","sequence":"1","reading":1,})"); + }); + require(equalToken("same-token", "same-token")); + require(!equalToken("same-token", "bad-token!")); + require(pageLimit("100") == 100); + rejects([&] { pageLimit("101"); }); + rejects([&] { identifier("device\n"); }); + rejects([&] { identifier(std::string("device\0", 7)); }); + std::cout << "3 validation groups passed: canonical payload, numeric bounds, parser/calendar " + "boundaries\n"; +} From ec08a0ce6a85947ce66524788b454cb6890b2c06 Mon Sep 17 00:00:00 2001 From: sdairs Date: Fri, 2 Oct 2026 19:29:20 +0100 Subject: [PATCH 2/3] Use the verified final Managed Postgres documentation URL --- applications/device-heartbeats/README.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/applications/device-heartbeats/README.md b/applications/device-heartbeats/README.md index 4c559b96..d9854d2a 100644 --- a/applications/device-heartbeats/README.md +++ b/applications/device-heartbeats/README.md @@ -1,6 +1,6 @@ # Device Heartbeats -A small C++20/Drogon API saves synthetic readings and advances each device's latest sequence in one transaction on [ClickHouse Managed Postgres (public beta)](https://clickhouse.com/docs/products/managed-postgres). Project tokens select seeded devices; matching sample UUIDs replay retained results. This is an instructional ingestion protocol, without hardware integration or monitoring/alerting behavior. +A small C++20/Drogon API saves synthetic readings and advances each device's latest sequence in one transaction on [ClickHouse Managed Postgres (public beta)](https://clickhouse.com/docs/products/managed-postgres/overview). Project tokens select seeded devices; matching sample UUIDs replay retained results. This is an instructional ingestion protocol, without hardware integration or monitoring/alerting behavior. ## Transaction and trust model From 644fa4f820d442924524853860a663ba4bab327f Mon Sep 17 00:00:00 2001 From: sdairs Date: Fri, 2 Oct 2026 21:31:17 +0100 Subject: [PATCH 3/3] docs: remove public beta wording --- applications/device-heartbeats/README.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/applications/device-heartbeats/README.md b/applications/device-heartbeats/README.md index d9854d2a..e8be5274 100644 --- a/applications/device-heartbeats/README.md +++ b/applications/device-heartbeats/README.md @@ -1,6 +1,6 @@ # Device Heartbeats -A small C++20/Drogon API saves synthetic readings and advances each device's latest sequence in one transaction on [ClickHouse Managed Postgres (public beta)](https://clickhouse.com/docs/products/managed-postgres/overview). Project tokens select seeded devices; matching sample UUIDs replay retained results. This is an instructional ingestion protocol, without hardware integration or monitoring/alerting behavior. +A small C++20/Drogon API saves synthetic readings and advances each device's latest sequence in one transaction on [ClickHouse Managed Postgres](https://clickhouse.com/docs/products/managed-postgres/overview). Project tokens select seeded devices; matching sample UUIDs replay retained results. This is an instructional ingestion protocol, without hardware integration or monitoring/alerting behavior. ## Transaction and trust model