Stellar ControlMission control · by Stellar Systems v0.1.0

Operations

Output Connectors

Persist the data of the MCS in a database or relay it to another system, with durable consumers, at-least-once delivery and lag monitoring.

The MCS keeps its data in JetStream with a bounded retention. Long-term storage, and the integration with other systems, go through output connectors: external processes that subscribe, read-only, to the telecommands, telemetry, measures, run events, alarms, files and streams of the MCS, to store them in a database or relay them elsewhere.

A connector is plug and play, like a driver or a gateway: a process outside the MCS, in any language — typically a short script — that implements a small NATS contract. The Rust SDK (run_connector) and the Python SDK (Connector) implement it; the repository ships stellar-timescale, a reference connector to TimescaleDB.

Declaring a connector#

Connectors are declared in the topology, by name, with the data they take and their targets:

YAML
connectors:
  timeseries-db:
    data: [tm_raw, measures, tc_events, run_events, alarms]
    targets: [sat1-fm, flatsat-1, sim-1, psu-sim-1]
  partner-relay:
    data: [files]
    targets: [sat1-fm]
  csv-export:
    data: [measures, tc_events]
    targets: [psu-lab-2]

The instance that registers with the kind connector under that name is that connector.

Data#

KindStreamSubjects of target TPayload
tm_rawTM_RAWstellar.tm.raw.*.TRaw frame, with Stellar-Ground-Time
measuresPARAMSstellar.param.T.>Sample (JSON): value, raw, time, ground time, link, delivery
tc_eventsTC_EVENTSstellar.tc.evt.T.>Telecommand event (JSON)
run_eventsRUNSstellar.run.evt.>Run event (JSON), of every run: their subjects do not carry the target
alarmsALARMSstellar.alarm.evt.T.>Alarm transition (JSON)
filesKV_stellar_transfers$KV.stellar_transfers.T.>Transfer (JSON); the files themselves in the stellar_files object store
streamsSTREAMSstellar.stream.T.>Raw segment, with Stellar-Ground-Time and Stellar-Stream-Seq

Subjects carry the target, the component, the instance and the measure (stellar.param.<target>.<component>.<instance>.<measure>, _ for a component without instances). Headers carry Stellar-Config, the configuration revision that produced the message. The payloads are described in the NATS Contract Reference.

Consumers#

For each kind of data of a connector, the leader of the reconciler keeps a durable JetStream consumer named connector_<name>_<kind>, filtered on the targets of the connector:

  • it is created when the connector appears in the topology, and delivers everything its stream still holds;
  • its filter follows the targets of the connector when the topology changes;
  • it is deleted when the connector leaves the topology.

Delivery#

  • Pull and acknowledge. The connector pulls its consumers and acknowledges each message once written. A message refused, or not acknowledged within 30 seconds, comes back. At most 1000 messages are unacknowledged at once.
  • At least once, in order. Delivery is at least once, in order within one consumer.
  • Stable identifiers. A message delivered again keeps its identifier: its Nats-Msg-Id when it has one (telecommand events, run events, alarm transitions), else <stream>:<sequence>. The connector deduplicates when it writes.
  • Decoupled from real time. A slow or stopped connector never slows the MCS down. Restarted, it takes up where it stopped, as long as the stream still holds the messages.

Credentials#

A connector registers like a driver or a gateway: a request on stellar.ctl.register with kind: "connector" and its name as instance, then heartbeats on stellar.ctl.hb.connector.<name>, and the control verbs status, credentials and reregister.

Once registered, a connector of the topology receives through credentials a short-lived NATS JWT limited to:

  • publishing its registration, its heartbeats and replies (_INBOX.>);
  • pulling, reading the state of and acknowledging its own consumers;
  • reading the stellar_files object store, when it takes files.

It can publish nothing on the subjects of the MCS. See NATS Accounts and Credentials.

Lag and data at risk#

Every 15 seconds, the leader of the reconciler measures each consumer of each connector: messages not delivered yet, messages not acknowledged, and the age of the oldest of them. It publishes the result on stellar.metrics.connector.<name>, and as Prometheus gauges labelled by connector and kind:

GaugeMeaning
stellar_connector_pendingMessages not delivered or not acknowledged
stellar_connector_lag_secondsAge of the oldest message not acknowledged
stellar_connector_at_risk1 when the retention is about to remove unread messages

A consumer is at risk when the retention of its stream is about to remove messages it has not read:

  • its oldest unread message has reached 80 % of the max_age of the stream; or
  • the stream, bounded by max_bytes, is 90 % full while its oldest message is still unread.

The reconciler then logs a warning. Alert on the gauge, for instance with this Prometheus rule:

YAML
groups:
  - name: stellar-connectors
    rules:
      - alert: ConnectorAboutToLoseData
        expr: stellar_connector_at_risk == 1
        labels: {severity: critical}
        annotations:
          summary: "Connector {{ $labels.connector }} ({{ $labels.kind }}) is about to lose data"

GET /v1/connectors gives the same figures on demand — each connector as declared, its instance, and the lag of each of its consumers — and the web console shows them under the topology.

The reference connector: stellar-timescale#

stellar-timescale (examples/rust/connectors/timescale, crate stellar-connector-timescale) writes raw telemetry, measures, telecommand events, run events and alarm transitions into TimescaleDB, a PostgreSQL extension.

Shell
STELLAR_NATS_URL=nats://127.0.0.1:4222 \
STELLAR_INSTANCE=timeseries-db \
STELLAR_TIMESCALE_URL="host=127.0.0.1 user=stellar password=stellar dbname=stellar" \
stellar-timescale

STELLAR_INSTANCE is the name of the connector in the topology; STELLAR_TIMESCALE_URL a PostgreSQL connection string. For development, podman compose up -d timescaledb starts a database with these credentials.

It creates its tables at start-up — one hypertable per kind of data, partitioned on time (the on-board time of a sample when given, else its ground reception time), with the primary key (time, id). Rows are written with INSERT … ON CONFLICT DO NOTHING: id is the stable identifier of the message, so a message delivered again is never written twice.

TableColumns
measurestime, target, component, instance (null for a single instance), measure, value_num / value_bool / value_text, raw, ground_time, link, delivery (realtime or deferred)
tm_rawtime (ground reception), target, gateway, deferred, frame
tc_eventstime, target, tc, state, detail, event (JSON)
run_eventstime, run, seq, event_kind, event (JSON)
alarmstime, target, component, instance, alarm, state, previous, severity, value, event (JSON)

These tables are its contract with the dashboards: Grafana reads them with its PostgreSQL data source. Without the TimescaleDB extension, on plain PostgreSQL, it creates plain tables with the same schema and warns. Compression and retention policies of the hypertables are left to the deployment.

A connector in Python#

examples/python/connector.py is the csv-export connector of the example topology: it writes the measures and telecommand events of psu-lab-2 to CSV files, one row per message identifier.

Python
from stellar_mcs import Connector, Message

connector = Connector("csv-export", software="csv-export", version="1.0.0")

@connector.on("measures")
def measure(message: Message) -> None:
    m = message.measure()
    measures.append(message.id, {"time": m.time.isoformat(), "target": m.target,
                                 "measure": ".".join(filter(None, (m.component, m.instance, m.name))),
                                 "value": m.value, "delivery": m.delivery})

@connector.on("tc_events")
def tc_event(message: Message) -> None:
    event = message.json()
    events.append(message.id, {"time": event["at"], "target": message.target,
                               "tc": event["tc"], "state": event["state"],
                               "detail": event.get("detail", "")})

connector.run()
Shell
cd examples/python
uv run connector.py --output export

A handler that raises gets the message again a bit later: write idempotently. Stop the connector and restart it: it resumes where it stopped, without duplicates. See Writing an Output Connector.

See also#

Stellar Control · v0.1.0

↑↓ to moveEnter to open