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:
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#
| Kind | Stream | Subjects of target T | Payload |
|---|---|---|---|
tm_raw | TM_RAW | stellar.tm.raw.*.T | Raw frame, with Stellar-Ground-Time |
measures | PARAMS | stellar.param.T.> | Sample (JSON): value, raw, time, ground time, link, delivery |
tc_events | TC_EVENTS | stellar.tc.evt.T.> | Telecommand event (JSON) |
run_events | RUNS | stellar.run.evt.> | Run event (JSON), of every run: their subjects do not carry the target |
alarms | ALARMS | stellar.alarm.evt.T.> | Alarm transition (JSON) |
files | KV_stellar_transfers | $KV.stellar_transfers.T.> | Transfer (JSON); the files themselves in the stellar_files object store |
streams | STREAMS | stellar.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-Idwhen 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_filesobject store, when it takesfiles.
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:
| Gauge | Meaning |
|---|---|
stellar_connector_pending | Messages not delivered or not acknowledged |
stellar_connector_lag_seconds | Age of the oldest message not acknowledged |
stellar_connector_at_risk | 1 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_ageof 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:
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.
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-timescaleSTELLAR_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.
| Table | Columns |
|---|---|
measures | time, target, component, instance (null for a single instance), measure, value_num / value_bool / value_text, raw, ground_time, link, delivery (realtime or deferred) |
tm_raw | time (ground reception), target, gateway, deferred, frame |
tc_events | time, target, tc, state, detail, event (JSON) |
run_events | time, run, seq, event_kind, event (JSON) |
alarms | time, 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.
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()cd examples/python
uv run connector.py --output exportA 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#
- Writing an Output Connector: the contract, and the SDKs.
- API reference: monitoring (
GET /v1/connectors).