Architecture

System Context

Flamingo is four Go binaries sharing one domain layer (lib/flamingo/) and one Postgres database:

  • cmd/api โ€” the ConnectRPC server (:9000). Exposes a service per entity (Device, DeviceStatus, Space, EventDefinition, EventRecord, EventIngestion, TaskDefinition, TaskExecution, SpaceMetricSnapshot), backed by a squirrel-built repo layer over Postgres. When KAFKA_BROKERS is set, writes to every service except SpaceMetricSnapshot are also published as protobuf messages onto Kafka so other systems can react without polling the API.
  • cmd/mqtt โ€” a bridge process. Subscribes to an MQTT broker (Mosquitto in docker-compose.yml, topic filter <prefix>/+/devices/+/+), turns each message into an EventRecord, and calls the API’s EventIngestionService.PublishEvent over ConnectRPC (gRPC-Web).
  • cmd/worker โ€” a Kafka consumer process. Runs a consumer per entity topic that imports/replicates incoming messages into the domain services (e.g. Device messages call DeviceService.Create), and โ€” when FDS_DB_URL is set โ€” a second set of consumers that maps EventRecord/TaskDefinition/TaskExecution/Device/Space/DeviceStatus messages into an external “FDS” (facility data sync) read-model database via lib/fds.
  • root flamingo CLI (main.go + cmd/) โ€” a Cobra-based CLI client, not a server. list/get/create/update/delete subcommands per entity (device, device-status, space, event-definition, event-record, task-definition, task-execution) that call cmd/api over ConnectRPC (gRPC-Web) at FLAMINGO_URL (default http://localhost:9000); create/update prompt interactively via cmd/form.

External dependencies:

  • PostgreSQL โ€” system of record, accessed via lib/pq (registered through otelsql for tracing) behind a squirrel-built repo layer, migrations under lib/repo/migrations (managed with goose). (jackc/pgx is also a dependency, but only for its ErrNoRows sentinel in unit tests โ€” it isn’t the runtime driver.)
  • MQTT broker โ€” device telemetry ingress, consumed only by cmd/mqtt.
  • Kafka (Redpanda in dev) โ€” fan-out bus. cmd/api is a producer (optional, gated on KAFKA_BROKERS); cmd/worker is the consumer, both for re-importing entities and for the FDS sync path.
  • FDS database โ€” a second Postgres-compatible store, written to only by cmd/worker’s FDS consumers, with field-level mapping/translation in lib/fds/mapping.go (e.g. flamingo battery_level 0โ€“100 โ†’ FDS 0.0โ€“1.0, flamingo space types collapsed onto the FDS space-type CHECK constraint).
  • ui/ โ€” a Vite/React/TypeScript frontend (separate from the Go service, talks to the API over ConnectRPC).
        MQTT broker                         Kafka / Redpanda
        (telemetry)                         (fan-out bus)
            |                                    ^  |
            v                                    |  v
       cmd/mqtt  ---PublishEvent (ConnectRPC)---> cmd/api <---optional publish---
       (bridge)                                    |  ^
                                                    |  |
                                              squirrel repo layer
                                                    |
                                                Postgres
                                                    ^
                                                    |
                                               cmd/worker
                                          (entity import consumers +
                                           FDS sync consumers) --> FDS DB

(The flamingo CLI isn’t shown above โ€” it’s an operator tool that talks to cmd/api the same way cmd/mqtt does, over ConnectRPC, and doesn’t participate in the telemetry/fan-out data flow.)

Data Model

Core entities live in schemas/flamingo/v1/*.proto, generated into lib/schemas/flamingo/v1, persisted via lib/repo:

  • Device (device.proto) โ€” a physical device: type, manufacturer, model, serial, lifecycle_status, the space_id it currently belongs to, tags, a free-form device_type_specific payload, registration/production timestamps, and mobility_type (static | dynamic). Dynamic devices get their space changes tracked in a separate device_location_history table (assigned_at/unassigned_at) rather than just overwriting space_id.
  • DeviceStatus (device_status.proto) โ€” a point-in-time telemetry snapshot for a device: health, battery_level, signal_strength, geolocation, last_contact, recorded_at, plus a device_type_specific payload for anything not covered by the fixed fields.
  • Space (space.proto) โ€” a physical location: name, space_type, an optional parent_id (spaces nest, e.g. room โ†’ floor โ†’ building), area_sqft, capacity, tags.
  • EventDefinition (event_definition.proto) โ€” the catalog of known event codes per device type: event_code, category, severity, description. This is the reference table; cmd/mqtt’s eventCategory map mirrors a subset of it for fast local categorization before the event ever reaches the API.
  • EventRecord (event_record.proto) โ€” an actual occurrence: device_id, space_id, event_code, category, source (e.g. "mqtt"), correlation_id, a raw payload_json, recorded_at. This is the append-only stream that both MQTT telemetry and any other event source ultimately write to via EventIngestionService.
  • TaskDefinition (task_definition.proto) โ€” a recurring or ad-hoc unit of work: name, activation_type, the spaces/devices/users it targets, and a schedule (schedule_type, start time/date, time-of-day, period).
  • TaskExecution (task_execution.proto) โ€” one run of a task: status (scheduled โ†’ started โ†’ completed/cancelled), planned vs. actual start/end, who executed it, notes, and the spaces/devices it touched.
  • SpaceMetricSnapshot (space_metric_snapshot.proto) โ€” a periodic rollup of metrics (arbitrary google.protobuf.Struct) for a space and/or device over an aggregation window โ€” the basis for things like occupancy or utilization reporting.

MQTT Bridge (Device Telemetry)

cmd/mqtt is the only piece that talks to the MQTT broker. It subscribes to <MQTT_TOPIC_PREFIX>/+/devices/+/+ (default prefix fds) and parses each topic as prefix/space/devices/{device_id}/{event_code}. For every message it:

  1. Looks up event_code in a local eventCategory map to assign a category (status, diagnostic, alert; defaults to alert if unknown).
  2. Resolves the device’s current space by calling DeviceService.GetDevice over ConnectRPC, cached in-memory per device_id for 5 minutes (short enough that a dynamic device’s room move shows up within one patrol cycle).
  3. Builds an EventRecord (new UUID, device id, resolved space id, event code/category, source: "mqtt", the raw payload as payload_json, recorded_at: now) and calls EventIngestionService.PublishEvent on cmd/api.

The bridge never talks to Postgres or Kafka directly โ€” everything flows through the API’s ingestion RPC, so all the normal validation/publish/audit behavior of EventRecord writes applies uniformly regardless of source.

Task & Event Tracking

  • Events are tracked as an append-only EventRecord stream, each one categorized (status / diagnostic / alert) and tied back to both the device and the space it happened in. EventDefinition rows describe what event codes can occur per device type; EventRecord rows are what did occur.
  • Tasks are modeled as a definition/execution split: a TaskDefinition describes the recurring work and what it targets (spaces, devices, assigned users), while each TaskExecution is a concrete instance with its own lifecycle status and planned vs. actual timing โ€” so a single recurring cleaning task produces a trail of individual, auditable executions rather than one row that gets mutated in place.
  • Event, task-definition, and task-execution activity, along with device/space/device-status changes, can be published to Kafka by cmd/api and picked back up by cmd/worker’s FDS consumers, which translate Flamingo’s model into an external facility-data-sync schema (lib/fds/mapping.go) โ€” this is how event/task state propagates out of Flamingo without those external systems needing to know Flamingo’s internal representation.