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 consume the data of the MCS and write it to a database, an archive, or another system. A connector is written in any language, typically a short script; the SDKs make it a handful of lines.
The MCS never queries the database of a connector: the API reads NATS only. History beyond the retention is read in the system the connector feeds, for instance Grafana or an SQL client on TimescaleDB. See Output Connectors for the operator's view.
Declaration#
A connector is declared in the topology, by name, with the kinds of data it takes and its targets:
connectors:
timeseries-db:
data: [tm_raw, measures, tc_events, run_events, alarms]
targets: [sat1-fm, flatsat-1, sim-1, psu-sim-1]
csv-export:
data: [measures, tc_events]
targets: [psu-lab-2]The instance that registers under that name, with the kind connector, is that connector.
Data#
For each kind of data, 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, follows its targets, and is deleted when the connector leaves
the topology. A new consumer delivers everything its stream still holds.
| Kind | Stream | Subjects of target T | Payload |
|---|---|---|---|
tm_raw | TM_RAW | stellar.tm.raw.*.T | raw frame, header 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.> (all runs: their subjects carry no target) | run event (JSON) |
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, Stellar-Ground-Time, Stellar-Stream-Seq |
Subjects carry the target, component, instance and measure
(stellar.param.<target>.<component>.<instance>.<measure>, _ for a single-instance
component); headers carry Stellar-Config, the configuration revision. The messages are those
of the NATS contract.
Consumption#
- Pull the consumer (
$JS.API.CONSUMER.MSG.NEXT.<stream>.connector_<name>_<kind>) and acknowledge each message once written (+ACKon its reply subject). - A message not acknowledged within 30 s, or refused (
-NAK), is delivered again: delivery is at least once, in order within a consumer. At most 1000 messages are unacknowledged at once. - Deduplicate at writing. A message keeps its identifier when delivered again: its
Nats-Msg-Idheader when it has one (telecommand events, run events, alarm transitions), else<stream>:<stream sequence>. Both SDKs give it as theidof the message. - Resume. A connector stopped takes up where it stopped, as long as its stream still holds the messages. A slow or stopped connector never slows the MCS down.
Registration and credentials#
As any component (see Writing a Driver): a request on
stellar.ctl.register with kind: "connector", the connector name as instance, its software
and version, its heartbeat period and the public key of its nkey; then heartbeats on
stellar.ctl.hb.connector.<name> and control verbs on stellar.ctl.rpc.connector.<name>.<verb>
(status, credentials, reregister).
Once registered, a connector of the topology receives through credentials a 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.
A connector in Rust and Python#
The Rust SDK runs a Connector with run_connector: it registers, finds each consumer (waiting
while the reconciler has not created it), and calls deliver for each message, in order. A
message is acknowledged once deliver returns Ok; an error has it delivered again after
RETRY (5 s). The Python Connector does the same with one handler per kind:
use stellar_sdk::{Connector, Delivered, Options, run_connector};
use tokio_util::sync::CancellationToken;
struct Printer;
impl Connector for Printer {
fn software(&self) -> (String, String) {
("printer".to_owned(), "1.0.0".to_owned())
}
// The kinds it takes; all by default.
fn data(&self) -> Vec<String> {
vec!["measures".to_owned(), "tc_events".to_owned()]
}
async fn deliver(&self, message: Delivered) -> Result<(), String> {
// `message.id` is stable across deliveries: write it with the data to deduplicate.
println!("{} {} {}", message.id, message.subject, String::from_utf8_lossy(&message.payload));
Ok(())
}
}
#[tokio::main]
async fn main() -> anyhow::Result<()> {
let options = Options::from_env("csv-export"); // STELLAR_INSTANCE: the connector name
run_connector(std::sync::Arc::new(Printer), &options, None, CancellationToken::new()).await
}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() # target, component, instance, name, value, raw, time…
measures.append(message.id, {"time": m.time.isoformat(), "target": m.target,
"measure": 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"]})
if __name__ == "__main__":
connector.run()What a message carries#
Rust Delivered | Python Message | Content |
|---|---|---|
kind | kind | Kind of data (measures…) |
subject | subject, tokens | Subject it was published on |
id | id | Stable identifier: Nats-Msg-Id, else <stream>:<sequence> |
stream_sequence | stream_sequence | Sequence in its stream |
published | published | When the stream stored it |
headers | headers | Stellar-Config, Stellar-Ground-Time… |
payload | data, json() | Content, as published |
| — | target | Target, when the subject carries one |
| — | measure() | For measures: a Measure with target, component, instance, name, value, raw, time, ground_time, link, delivery |
In Python, a handler may be plain or async; an exception has the message delivered again a bit
later (NAK with delay). A connector needs at least one @connector.on(kind); an unknown kind
raises ValueError at declaration. The connector name is its instance name:
Connector(name, software=None, version="0.1.0").
The TimescaleDB connector#
stellar-timescale (examples/rust/connectors/timescale, crate stellar-connector-timescale) is
the reference connector: raw telemetry, measures, telecommand events, run events and alarm
transitions in TimescaleDB. It creates its tables at start-up and takes up where it stopped.
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-timescaleOne hypertable per kind of data, partitioned on time, primary key (time, id):
| Table | Columns |
|---|---|
measures | time, target, component, instance (null for a single instance), measure, value_num / value_bool / value_text, raw, ground_time, link, delivery |
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) |
Lag#
Every 15 s, the leader of the reconciler measures each consumer of a connector: messages not
delivered, messages not acknowledged, and the age of the oldest of them. It publishes it on
stellar.metrics.connector.<name> (a ConnectorLag: connector, kind, stream, consumer,
pending, unacked, oldest_unread, lag_seconds, retention_seconds, fill, active,
at_risk, reason, at) and as Prometheus gauges labelled by connector and kind:
stellar_connector_pending, stellar_connector_lag_seconds, stellar_connector_at_risk.
GET /v1/connectors gives the same, computed on demand (API).
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,
bound by max_bytes, is 90 % full while its oldest message is unread. An alert 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"Running a connector#
STELLAR_INSTANCE must be the name of the connector in the topology. The other variables are
those of every component (STELLAR_NATS_URL, STELLAR_NATS_CREDENTIALS, STELLAR_NATS_CA,
STELLAR_NATS_CERT, STELLAR_NATS_KEY): see Writing a Driver.
cd examples/python
uv run connector.py --output export # connector csv-export: export/measures.csv, export/tc_events.csv