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 asquirrel-built repo layer over Postgres. WhenKAFKA_BROKERSis 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 indocker-compose.yml, topic filter<prefix>/+/devices/+/+), turns each message into anEventRecord, and calls the API’sEventIngestionService.PublishEventover 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.Devicemessages callDeviceService.Create), and โ whenFDS_DB_URLis 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 vialib/fds.- root
flamingoCLI (main.go+cmd/) โ a Cobra-based CLI client, not a server.list/get/create/update/deletesubcommands per entity (device, device-status, space, event-definition, event-record, task-definition, task-execution) that callcmd/apiover ConnectRPC (gRPC-Web) atFLAMINGO_URL(defaulthttp://localhost:9000);create/updateprompt interactively viacmd/form.
External dependencies:
- PostgreSQL โ system of record, accessed via
lib/pq(registered throughotelsqlfor tracing) behind asquirrel-built repo layer, migrations underlib/repo/migrations(managed withgoose). (jackc/pgxis also a dependency, but only for itsErrNoRowssentinel 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/apiis a producer (optional, gated onKAFKA_BROKERS);cmd/workeris 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 inlib/fds/mapping.go(e.g. flamingobattery_level0โ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, thespace_idit currently belongs to, tags, a free-formdevice_type_specificpayload, registration/production timestamps, andmobility_type(static|dynamic). Dynamic devices get their space changes tracked in a separatedevice_location_historytable (assigned_at/unassigned_at) rather than just overwritingspace_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 adevice_type_specificpayload for anything not covered by the fixed fields. - Space (
space.proto) โ a physical location: name,space_type, an optionalparent_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’seventCategorymap 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 rawpayload_json,recorded_at. This is the append-only stream that both MQTT telemetry and any other event source ultimately write to viaEventIngestionService. - 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 (arbitrarygoogle.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:
- Looks up
event_codein a localeventCategorymap to assign a category (status,diagnostic,alert; defaults toalertif unknown). - Resolves the device’s current space by calling
DeviceService.GetDeviceover ConnectRPC, cached in-memory perdevice_idfor 5 minutes (short enough that a dynamic device’s room move shows up within one patrol cycle). - Builds an
EventRecord(new UUID, device id, resolved space id, event code/category,source: "mqtt", the raw payload aspayload_json,recorded_at: now) and callsEventIngestionService.PublishEventoncmd/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
EventRecordstream, each one categorized (status/diagnostic/alert) and tied back to both the device and the space it happened in.EventDefinitionrows describe what event codes can occur per device type;EventRecordrows are what did occur. - Tasks are modeled as a definition/execution split: a
TaskDefinitiondescribes the recurring work and what it targets (spaces, devices, assigned users), while eachTaskExecutionis 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/apiand picked back up bycmd/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.