Stellar ControlMission control · by Stellar Systems v0.1.0

SDKs and Integration

Writing an Output Connector

An output connector consumes the data of the MCS to store it or relay it to another system.

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:

YAML
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.

KindStreamSubjects of target TPayload
tm_rawTM_RAWstellar.tm.raw.*.Traw frame, header 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.> (all runs: their subjects carry no target)run event (JSON)
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, 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 (+ACK on 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-Id header when it has one (telecommand events, run events, alarm transitions), else <stream>:<stream sequence>. Both SDKs give it as the id of 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_files object store, when it takes files.

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:

Rust
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
}
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()   # 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 DeliveredPython MessageContent
kindkindKind of data (measures…)
subjectsubject, tokensSubject it was published on
ididStable identifier: Nats-Msg-Id, else <stream>:<sequence>
stream_sequencestream_sequenceSequence in its stream
publishedpublishedWhen the stream stored it
headersheadersStellar-Config, Stellar-Ground-Time…
payloaddata, json()Content, as published
—targetTarget, 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.

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

One hypertable per kind of data, partitioned on time, primary key (time, id):

TableColumns
measurestime, target, component, instance (null for a single instance), measure, value_num / value_bool / value_text, raw, ground_time, link, delivery
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)

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:

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"

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.

Shell
cd examples/python
uv run connector.py --output export     # connector csv-export: export/measures.csv, export/tc_events.csv

Stellar Control · v0.1.0

↑↓ to moveEnter to open