The Python SDK (stellar-mcs, in sdk/python) implements the contracts of Stellar Control for
drivers, transports, gateways and output connectors over NATS, and for pass feeders over the
HTTP API. A component only declares what it does, with decorators; the SDK handles
registration with the reconciler, key pair and JWT renewal, heartbeats, the bindings it
receives, JetStream consumers, telecommand events (ACK 1 and ACK 2), deadlines, fragments,
control verbs and reconnection.
The same package holds the client of the API: it launches runs of procedures, follows their events, answers their questions and controls them, from a script, a test bench or a notebook.
Install#
The package needs Python 3.11 or later; its dependencies are nats-py (with nkeys), httpx
and websockets.
uv add stellar-mcs # or: pip install stellar-mcsWithin the repository, examples/python uses the SDK from sdk/python directly
([tool.uv.sources] stellar-mcs = { path = "../../sdk/python", editable = true }); uv run
installs it on first use.
Public API#
Everything is imported from the package root:
from stellar_mcs import (
Driver, Transport, Gateway, Connector, PassFeeder, # components
SemanticTc, RawSample, FileChunk, LinkContext, Link, # messages
Downlink, Message, Measure, Pass, Replaced, # helpers
Bindings, BoundLink, Health, TcState, Timestamp, now,
Options, RegistrationRejected, PassRefused,
Client, Run, RunEvent, RunResult, RunStatus, Answer, Question, ApiError, # API client
)Components#
| Class | Constructor | Declarations |
|---|---|---|
Driver | Driver(software, version, *, catalogue, telecommands=(), measures=(), codec=None, output=None, params_schema=None, instance=None) | @driver.encode, @driver.decode, @driver.echo, @driver.chunks, @driver.decode_file |
Transport | Transport(software, version, *, input, output, params_schema=None, instance=None) | @transport.wrap, @transport.unwrap |
Gateway | Gateway(software, version, *, tags=(), uplink=None, downlink=None, instance=None) | @gateway.send, @gateway.receive, @gateway.link_available |
Connector | Connector(name, *, software=None, version="0.1.0") | @connector.on(kind) |
PassFeeder | PassFeeder(url, identity=None, *, source, token=None, timeout=10.0) | methods upsert, replace_window, passes |
instance is the default instance name, <software>-1 otherwise (the connector name for a
connector); STELLAR_INSTANCE overrides it.
Declarations#
| Decorator | Signature | Notes |
|---|---|---|
@driver.encode | (tc: SemanticTc) -> bytes or str, or (tc, context: LinkContext) | An exception gives ENCODE_FAILED, its message as the reason |
@driver.decode | (frame: bytes) -> Iterable[RawSample], or (frame, context) | A frame that is not telemetry yields nothing |
@driver.echo | (tc: SemanticTc, frame: bytes) -> bool or None | True conforming echo, False non-conforming, None not its echo |
@driver.chunks | (frame: bytes) -> Iterable[FileChunk] | File chunks carried by a frame |
@driver.decode_file | (file_type: str, content: bytes) -> Iterable[RawSample] or None | None for a type the driver does not decode; samples carry onboard_time |
@transport.wrap | (unit: bytes, context: LinkContext) -> bytes or Iterable[bytes] | Frames in order; an exception gives SEND_FAILED |
@transport.unwrap | (frame: bytes, context: LinkContext) -> Iterable[bytes] or None | Units the frame completes; an exception drops the frame |
@gateway.send | (target: str, frame: bytes) -> None, or (target, frame, context: LinkContext) | Returning means SENT; an exception SEND_FAILED. With a third parameter, the context of the frame: environment, mode of the target and its parameters |
@gateway.configure | (context: LinkContext) -> None | Optional: called when a link is bound to the gateway, and again each time its binding changes, such as a new mode of its target; runs in the background |
@gateway.receive | async (downlink: Downlink) -> None | Receives until the end; if it fails, the gateway is reported degraded |
@gateway.link_available | () -> bool | Link state of the heartbeats; always available without it |
@connector.on(kind) | (message: Message) -> None | measures, tm_raw, tc_events, run_events, alarms, files, streams; an exception has the message delivered again later |
@component.rpc(verb) | (request: dict) -> dict | Own control verb on stellar.ctl.rpc.<kind>.<instance>.<verb>; an exception replies {"error": …} |
Decorated functions may be plain or async.
Messages#
| Class | Fields and helpers |
|---|---|
SemanticTc | id, target, component, telecommand, instance, link, args, environment; tc["arg"], tc.get("arg", default), tc.bytes_arg("arg") (hexadecimal to bytes), tc.name (component.telecommand) |
RawSample | RawSample(component, measure, value, instance=None, onboard_time=None); bytes values are sent in hexadecimal |
FileChunk | FileChunk(file_id, generation, offset, data) |
LinkContext | target, link, environment, params, mode, parameters; context.param(name, default=None) reads a parameter of the link (params of the topology), context.parameter(name, default=None) a parameter of the target in its mode (rf.tx_baudrate, tcu[TCU1].anode_voltage; an instance without a value of its own takes the value given for every instance, tcu.anode_voltage) |
Link | Link(type, rate_bps): a gateway link in one direction, Link("rf", 64_000) |
Downlink | await publish(target, frame, *, ground_time=None, deferred=False), await publish_segment(target, stream, seq, segment, *, ground_time=None); the ground time is now by default |
Message | kind, subject, tokens, id (stable), stream_sequence, published, headers, data; json(), target, measure() |
Measure | target, component, instance, name, value, raw, time, ground_time, link, delivery |
Bindings, BoundLink | The links served, as last delivered (component.bindings); BoundLink.context(environment) resolves the parameters |
Health | HEALTHY, DEGRADED |
TcState | PENDING, REJECTED, ENCODED, ENCODE_FAILED, SENT, SEND_FAILED, VERIFIED, VERIFY_FAILED, VERIFY_TIMEOUT, COMPLETE, and ECHO (not a state) |
Timestamp, now() | An aware datetime, in UTC on the wire |
Pass, Replaced, PassRefused | See Writing a Pass Feeder |
Every component#
| Member | Role |
|---|---|
run(**options) | Runs until Ctrl-C or SIGTERM, settings from the environment and options; logs at STELLAR_LOG (INFO) |
async with running(options=None) | Runs in the background while in the block; returns once registered; registration errors are raised on entering |
await serve(options=None) | Runs until cancelled |
set_health(health, reason=None) | Health of the next heartbeats; degraded needs a reason |
bindings | The links served, as last delivered |
instance | The instance name |
RegistrationRejected is raised when the reconciler refuses the registration: it is final.
Examples#
from stellar_mcs import Driver, RawSample
driver = Driver(
"scpi-psu", "1.0.0",
catalogue="lab-psu@^1.0",
telecommands=["psu.set_voltage", "psu.output_on", "psu.output_off"],
measures=["psu.voltage", "psu.current", "psu.output_enabled"],
)
@driver.encode
def encode(tc):
return f"VOLT {tc['voltage']}"
@driver.decode
def decode(frame):
for field in frame.decode().split(";"):
key, value = field.split()
if key == "VOLT":
yield RawSample("psu", "voltage", float(value))
driver.run()examples/python carries a SCPI power supply into the MCS without hardware:
| File | Component |
|---|---|
driver.py | driver scpi-psu: telecommands of lab-psu to SCPI lines, status lines to measures |
gateway.py | gateway lab-eth: TCP link to the instrument, polling, link state, reconnection, a control verb |
connector.py | connector csv-export: measures and telecommand events to CSV files, idempotently |
pass_feeder.py | pass feeder: an FDS planning or provider bookings pushed as snapshots |
instrument.py | a fake SCPI power supply on TCP, in place of the hardware |
run_procedure.py | the API client: runs any procedure, prints its events and answers at the terminal |
hot_standby.py | the typed function of « Hot standby test » generated by stellar generate client |
lab_psu/ | the driver of lab-psu from the skeleton of stellar generate driver: _generated.py, driver.py, test_driver.py |
platform_v3_steps/ | the functions generated by stellar generate client for platform-v3-steps |
test_procedures_live.py | tests of the generated functions against a running MCS (STELLAR_TEST_API_URL) |
cd examples/python
uv run instrument.py --port 5025
uv run driver.py
STELLAR_INSTANCE=lab-eth-1 uv run gateway.py --target psu-lab-2 --instrument 127.0.0.1:5025
uv run connector.py --output export
stellar send psu-lab-2 psu.output_on
stellar send psu-lab-2 psu.set_voltage voltage="28 V" # ENCODED, SENT, then VERIFIED
stellar watch psu-lab-2The reconciler binds the driver and the gateway to psu-lab-2 of examples/config; the
connector csv-export fills export/measures.csv and export/tc_events.csv, and resumes
without duplicates after a restart.
A gateway that follows the mode of its target#
A gateway that must adapt the medium to the mode of its target (a data rate, a modulation)
reads the parameters of the target from the context. @gateway.configure applies them when the
link is bound and each time the binding changes, for instance after a configure of a
procedure or stellar mode; the send can also take the context of each frame:
from stellar_mcs import Gateway, Link, LinkContext
gateway = Gateway("uhf-gateway", "1.0.0", tags=["uhf"], uplink=Link("rf", 9600))
@gateway.configure
async def configure(context: LinkContext) -> None:
await modem.set_rate(context.parameter("rf.tx_baudrate"))
@gateway.send
async def send(target: str, frame: bytes, context: LinkContext) -> None:
await modem.transmit(frame, power=context.parameter("rf.tx_power"))See Modes and Parameters and Writing a Gateway.
API client#
Client talks to the HTTP API of a cell. The MCS executes a
run (leases, confirmations, log, validated versions); the client only launches it, follows its
events, answers its questions and decisions, and controls it. Everything is async.
from stellar_mcs import Answer, Client
async with Client("http://mcs:8080", "alice") as mcs: # or token="…"
result = await mcs.execute(
"Hot standby test",
environment="AIT",
targets={"sat": "flatsat-1", "psu": "psu-lab-2"},
inputs={"tcu": "TCU2", "bus_voltage": "28 V"},
on_ask=lambda question: True, # confirms every question
)
assert result.ok, result.reasonClient#
Client(url, identity=None, *, token=None, timeout=10.0, reconnections=5, transport=None)| Parameter | Meaning |
|---|---|
url | API of the MCS, http://mcs:8080 or https://…; WebSockets use ws:// or wss:// accordingly |
identity | Operator declared with X-Stellar-User, where the environment accepts declared identities; ignored when token is given |
token | OIDC token, sent as Authorization: Bearer |
timeout | Timeout of the HTTP requests and of the opening of a WebSocket, in seconds |
reconnections | Attempts to take up the events of a run after a lost connection |
transport | An httpx transport, for tests |
Use it as an async context (async with Client(...) as mcs), or call await mcs.close().
| Method | Returns | Does |
|---|---|---|
await launch(procedure, *, environment, targets, inputs=None, library=None) | Run | POST /v1/runs; returns once the MCS accepted the run. inputs are written as in a run request: "28 V", "TCU2", True, 3. library when several libraries have a procedure of that name |
await execute(procedure, *, environment, targets, inputs=None, library=None, on_ask=None, on_event=None) | RunResult | launch, then Run.wait |
run(id) | Run | A run already launched, by its identifier: to follow, answer or control it |
await upload(content: bytes) | str | POST /v1/uploads: stores a content and returns its hash, the value of a file input |
Run#
A run launched on the MCS: id, ir (hash of the resolved run) and warnings (for instance a
run that may outlast the pass in progress), as launch returned them.
| Method | Does |
|---|---|
async for event in run.events(after=0) | The events of the log from the start (those after position after), then as they come, up to Finished. A lost connection is taken up again, with a growing delay (0.4 s up to 5 s) and at most reconnections attempts in a row, without repeating an event; then ConnectionError |
await run.wait(on_ask=None, on_event=None) | Follows the run to its end and returns a RunResult. Each event goes to on_event; each question or decision to on_ask, whose reply is sent as the answer |
await run.answer(path, accepted=True, value=None) | Answers the question or decision of step path: value is a typed answer, or a decision ("replay", "skip", "abort") |
await run.suspend() | Suspends at the end of the current instruction; the leases are released |
await run.resume(choice="replay") | Resumes: replay the step where it stopped, skip it, or abort |
await run.abort() | Stops the run for good |
await run.status() | RunStatus: run, procedure, suspended, step, pending (a Question), finished, reason |
await run.report() | The HTML report of the run |
RunResult holds run, ok, reason and every events of the log;
result.failed_steps() lists the StepFinished attempts that failed.
Questions and answers#
on_ask receives a Question: run, path (the step asking), kind ("ask" or
"decision"), prompt (the question, or why a decision is needed) and role (the role whose
confirmation is awaited, for a hazardous telecommand). It may be plain or async, and returns:
| Reply | Sent as |
|---|---|
True / False | A confirmation, or a refusal (the step fails) |
Answer(accepted=True, value=…) | A typed answer (ask … as <type> into <name>) |
"replay", "skip", "abort" | A decision |
None | Nothing: someone else answers, from the web console, the CLI or another client |
Who may answer is decided by the MCS, as for any answer: see Questions, Decisions and Control and Hazardous Confirmations.
Events#
run.events() and on_event give typed events, from stellar_mcs.client. Every event has
run, seq (position in the log, from 1), at, event (its kind) and data (the event as
the API sent it).
| Class | Kind | Fields |
|---|---|---|
Started | started | procedure, library, environment, targets, ir, by, inputs |
ProcedureStarted | procedure_started | path |
ProcedureFinished | procedure_finished | path, ok |
StepStarted | step_started | path, attempt, inputs |
TelecommandSent | telecommand_sent | path, tc, target, telecommand (component[instance].telecommand), changes_state, retry_of |
TelecommandFinished | telecommand_finished | path, tc, state, detail, acks |
Checked | checked | path, statement, condition, ok, reason, samples |
Logged | logged | path, message, samples |
Asked | asked | path, prompt, role |
Answered | answered | path, accepted, by, value, role |
AnswerRefused | answer_refused | path, by, reason |
TakenOver | taken_over | by |
Suspended | suspended | path, by |
Continued | continued | path, by, choice |
DecisionRequired | decision_required | path, reason |
StepFinished | step_finished | path, ok, attempt, reason |
ModeChanged | mode_changed | path, target, mode |
Finished | finished | ok, reason |
ModeChanged is published when a configure sets the mode of a target:
from stellar_mcs.client import ModeChanged
if isinstance(event, ModeChanged):
print(f"{event.target} is now in {event.mode}")An event of a kind the SDK has no class for yet comes as a plain RunEvent, its fields in
data (event.event gives its kind).
See Running Procedures for the meaning of each event.
Errors#
A request the API refuses raises ApiError: status (HTTP), errors (the
[{code, message, help}] of the error body) and codes:
from stellar_mcs import ApiError
try:
run = await mcs.launch("Hot standby test", environment="in_orbit", targets={"sat": "sat1-fm"})
except ApiError as refused:
if "run::must-be-scheduled" in refused.codes:
...run.events() raises ApiError too when the WebSocket is refused (an unknown run).
A command-line runner#
examples/python/run_procedure.py runs any procedure, prints its events as they come and
answers at the terminal (--yes confirms every question). It exits with 0 when the run
succeeds, 1 when it fails, 2 when the API refuses it:
cd examples/python
uv run run_procedure.py "Hot standby test" --environment AIT \
--target sat=sim-1 --target psu=psu-sim-1 --input tcu=TCU2 --input "bus_voltage=28 V"--api defaults to STELLAR_API_URL (http://127.0.0.1:8080), --identity to STELLAR_USER
(the login name). Its core:
from stellar_mcs.client import Checked, Finished, StepStarted, TelecommandSent
def show(event: RunEvent) -> None:
match event:
case StepStarted(path=path, attempt=attempt):
print(f"▶ {path}" + (f" (attempt {attempt})" if attempt > 1 else ""))
case TelecommandSent(telecommand=telecommand, target=target):
print(f" → {telecommand} to {target}")
case Checked(condition=condition, ok=ok):
print(f" {'✓' if ok else '✗'} {condition}")
case Finished(ok=ok, reason=reason):
print("run succeeded" if ok else f"run failed: {reason}")
async with Client(args.api, args.identity) as mcs:
run = await mcs.launch(args.procedure, environment=args.environment,
targets=targets, inputs=inputs)
result = await run.wait(on_ask=ask, on_event=show)Typed functions per procedure#
stellar generate client writes, for a library, one pair of functions per procedure with its
roles, inputs and environments typed: <procedure>(client, …) follows the run to its end,
launch_<procedure>(client, …) returns the Run once accepted.
from platform_v3_steps import TcuId, hot_standby_test
async with Client("http://127.0.0.1:8080", "operator") as mcs:
result = await hot_standby_test(
mcs, sat="sim-1", psu="psu-sim-1", tcu=TcuId.TCU2, bus_voltage="28 V",
environment="AIT", # Literal["AIT", "IVV"]: where the procedure may run
on_ask=lambda question: True,
)See Code Generators. stellar generate driver
writes the typed skeleton of a driver in the same way: see
Writing a Driver.
Running#
component.run() reads its settings from the environment (Options.from_env()):
| Variable | Default | |
|---|---|---|
STELLAR_NATS_URL | nats://localhost:4222 | NATS server |
STELLAR_INSTANCE | <software>-1 | Instance name, unique per kind |
STELLAR_NATS_CREDENTIALS | none | Bootstrap credentials, limited to registration |
STELLAR_NATS_CA | none | CA certificates (PEM) the server is checked against: TLS only |
STELLAR_NATS_CERT, STELLAR_NATS_KEY | none | Client certificate and its key, for mTLS |
STELLAR_LOG | INFO | Log level of run() |
Within an application, build Options yourself and use the async context:
from stellar_mcs import Options
async with driver.running(Options(nats_url="nats://mcs:4222", instance="scpi-psu-2")):
...Options also has heartbeat_period (1 s), register_timeout (2 s) and metrics_period
(5 s, throughput reports of a gateway). Without CA nor client certificate, a tls:// URL checks
the server against the system roots.
Protocol helpers#
stellar_mcs.csp (CSP 1 and CSP 2 packets, CRC32), stellar_mcs.can (CAN frames on NATS) and
stellar_mcs.cfp (CSP over CAN, the CFP of libcsp): see
Writing a Transport.
What the Python SDK does not cover#
- the CFDP entity of drivers, and COP-1 in transports: they are in the Rust SDK only;
- a Prometheus endpoint.
A Python gateway still carries the frames of no telecommand (CFDP PDUs of a Rust driver) like the others, without event.
Development#
cd sdk/python
uv sync
uv run pytest # unit tests
STELLAR_TEST_NATS_URL=nats://127.0.0.1:4222 uv run pytest # and against a NATS server
uv run ruff check && uv run ruff format --check && uv run mypy src tests
uv build # wheel and sdist in dist/The integration tests (tests/test_components.py) play the reconciler themselves: a driver, a
gateway and a connector registered, bound, a telecommand to ACK 2 and its echo, a frame to
decoded values, a connector's message. They need a NATS server with JetStream
(nats-server -js) and nothing else of the MCS. tests/test_contract.py,
tests/test_passes.py and tests/test_tls.py cover the messages, the pass feeder and TLS;
tests/test_client.py covers the API client against a fake API, and tests/test_client_live.py
runs it against a running MCS when STELLAR_TEST_API_URL is set.