-
Notifications
You must be signed in to change notification settings - Fork 1
Hub
The Hub (qitech_framework_hub) is the framework's async controller for building integrations. It holds the session with the Runtime and fans the data out to your own code: an HTTP or WebSocket API, a database writer, a cloud uplink, an automation script, and so on.
Where the TUI is a finished application for a person at a terminal, the Hub is a framework for your own controller. It handles the protocol for you, and you plug in two kinds of module:
| Module | Trait | Runs as | Use for |
|---|---|---|---|
| Listener | Listener |
synchronous callbacks inside the session task | receiving data: schemas, init events, reports, disconnects |
| Actor | Actor |
its own async task | doing things: serving an API, running automation, sending requests |
The examples/apps/debug app uses both: a listener (SocketIODispatcher) turns reports into shared state, and an actor (Server) serves that state over HTTP and Socket.IO on port 3001.
The Hub runs in the same process as the Runtime:
use qitech_framework::{run_with_hub, HubConfiguration};
#[tokio::main]
async fn main() {
let config_rt = RuntimeConfiguration::new()
.ethercat(EtherCATConfig::default())
.machine::<MyMachine>();
let config_hub = HubConfiguration::new()
.listener(MyListener::new())
.actor(MyApi::new());
run_with_hub(config_rt, config_hub).await.unwrap();
}run_with_hub does three things:
- It creates an in-process session with
session::mpsc(64). - It starts the Runtime on its own thread.
- It runs
qitech_framework_hub::runon the current tokio runtime.
It returns once every Hub task has finished. Actors usually run forever, so in practice it doesn't return.
If you only need the Hub with a different session provider, call it directly:
qitech_framework_hub::run(config_hub, provider).awaitAdd qitech_framework_hub to your dependencies for the Listener, Actor and ActorContext traits and types. HubConfiguration is also re-exported as qitech_framework::HubConfiguration.
pub trait Listener: Send {
fn on_schema_sync(&mut self, schema: &MachineSchema) {}
fn on_init_event_received(&mut self, event: &RuntimeInitEvent) {}
fn on_report_received(&mut self, report: &RuntimeReport) {}
fn on_runtime_disconnected(&mut self) {}
}All four methods have empty defaults, so you implement only the ones you need. They are called in protocol order:
-
on_schema_synconce per machine type, during the schema sync phase, -
on_init_event_receivedfor every init event (hardware discovery, machine builds, …), -
on_report_receivedfor every report while running, and -
on_runtime_disconnectedwhen the session ends while running.
Example: keeping the latest measurement values in shared state for an API to read:
struct LatestValues {
state: Arc<ArcSwap<HashMap<(MachineInstanceIdentification, String), Option<f64>>>>,
}
impl Listener for LatestValues {
fn on_report_received(&mut self, report: &RuntimeReport) {
let mut next = (**self.state.load()).clone();
for s in &report.machines.measurement_snapshots {
next.insert((s.machine, s.path.clone()), s.value);
}
self.state.store(Arc::new(next));
}
}Listeners run inside the session manager task, one after another, between receiving a report and receiving the next one. While a listener runs, no report is read from the Runtime.
A slow listener makes the Runtime's outgoing queue fill up, and a full queue either stalls or ends the Runtime (see Protocol). So:
- Do: update in-memory state, push into a channel, increment counters.
-
Don't: write to a database, make network calls, block on locks,
sleep.
For slow work, hand the data off to an actor through a channel:
impl Listener for DbFeeder {
fn on_report_received(&mut self, report: &RuntimeReport) {
// bounded channel; the actor on the other end does the slow inserts
if self.tx.try_send(report.clone()).is_err() {
tracing::error!("database writer is falling behind");
}
}
}Reports are deltas (see Journals). If you drop reports on the floor, your view of the machines diverges. Decide explicitly what happens when a consumer can't keep up.
pub trait Actor: Send + Sync {
fn run(self, ctx: ActorContext) -> impl Future<Output = ()> + Send + 'static;
}Every actor is spawned as its own tokio task when the Hub starts, and it runs until its future completes. It gets an ActorContext:
pub struct ActorContext {
pub schemas: Swappable<SchemaRegistry>, // machine types known from the current session
pub machines: Swappable<MachineRegistry>, // reserved, currently always empty
// + a private request sender
}
impl ActorContext {
pub async fn send_request(&self, request: RuntimeRequestKind)
-> Result<Result<(), RuntimeRequestError>, oneshot::error::RecvError>;
}-
ActorContextisClone. Hand it to every HTTP handler, background job and so on that needs it. The debug example passes it to axum as router state. -
schemasis anArc<ArcSwap<BTreeMap<MachineIdentification, MachineSchema>>>. Read it withctx.schemas.load(). It is replaced as a whole after each schema sync. - Actors don't receive reports. For live data, pair an actor with a listener that writes into state they share, as in the listener example above.
impl Actor for Automation {
async fn run(self, ctx: ActorContext) {
let result = ctx
.send_request(RuntimeRequestKind::SetConfigProperty {
target: LaserV1::IDENTIFICATION.unique(1),
path: "diameter.target".into(),
value: ScalarValue::Float(1.75),
})
.await;
match result {
Ok(Ok(())) => tracing::info!("accepted"),
Ok(Err(e)) => tracing::warn!(%e, "rejected by the runtime"),
Err(_) => tracing::error!("no runtime session"),
}
}
}send_request resolves when the Runtime's response arrives in a report, usually within one report interval (1/32 s by default). The result is nested:
| Result | Meaning |
|---|---|
Ok(Ok(())) |
the Runtime executed the request |
Ok(Err(RuntimeRequestError)) |
the Runtime rejected it, for example because the resource doesn't exist, isn't writable, a constraint was violated, or the command is disabled |
Err(RecvError) |
the request was never sent, because no session is running |
You don't choose request_ids. The Hub assigns and matches them for you.
This section is for people working on the Hub itself.
flowchart LR
Runtime(["Runtime"])
actors(["actors"])
subgraph run ["qitech_framework_hub::run"]
direction LR
sm["session_manager"]
listeners["listeners<br/>(sync callbacks)"]
rb["report broadcast<br/>(Arc#lt;RuntimeReport#gt;)"]
tm["transaction_manager<br/>request_id ↔ oneshot responder"]
end
Runtime <-- session --> sm
sm --> listeners
sm --> rb
rb --> tm
actors -- request channel --> tm
tm -- "per-session request channel<br/>(sent via request_dispatcher)" --> sm
classDef external fill:#334155,stroke:#1e293b,color:#f8fafc
classDef task fill:#dbeafe,stroke:#2563eb,color:#1e3a8a
classDef channel fill:#fef3c7,stroke:#d97706,color:#78350f
classDef callback fill:#f1f5f9,stroke:#64748b,color:#334155
class Runtime,actors external
class sm,tm task
class rb channel
class listeners callback
style run fill:transparent,stroke:#94a3b8,stroke-dasharray:5 5,color:#64748b
linkStyle default stroke:#64748b
run(config, provider) spawns everything into one JoinSet and waits for all tasks. A task that panics is logged with eprintln! and doesn't stop the others.
This task owns the controller session. In a loop, it:
- calls
provider.provide(). If that returnsTransportError::Disconnected, the provider is permanently gone and the task ends. On other errors it retries. - runs the handshake, then the schema sync. It calls
on_schema_syncfor each schema, rejects duplicates withSchemaSyncError::DuplicateItem, and publishes the result toschemas(ArcSwap::store). - runs the init phase, calling
on_init_event_receivedfor each event. - creates a fresh request channel for this session and hands its sender to the transaction manager (the "request dispatcher").
- while running, uses
select!(biased towards reports) to:- receive a report, wrap it in an
Arc, broadcast it, and callon_report_receivedon every listener, or - forward a request from the transaction manager to the Runtime.
- receive a report, wrap it in an
- on a receive error, calls
on_runtime_disconnectedand goes back to step 1.
This task matches requests with responses:
- It receives
(RuntimeRequestKind, oneshot::Sender)pairs from actors, gives each arequest_idfrom a counter that goes up by one, remembers the responder in apendingmap, and sends the request through the current session's channel. - If no session is running, it drops the responder straight away, so the caller gets
RecvError. - It subscribes to the report broadcast, and for every
RuntimeResponsein a report, sends the result to the matching responder.
Swappable<T> is Arc<ArcSwap<T>>: readers take a cheap, consistent snapshot with load(), and a writer replaces the whole value with store(). Actors, listeners and the session manager can share registries this way without locks.
| Channel | Capacity |
|---|---|
Runtime ↔ Hub session (run_with_hub) |
64 |
| report broadcast | 32 reports (~1 s at 32 Hz) |
| actor → transaction manager | 128 requests |
| transaction manager → session | 64 requests per session |
| request dispatcher | 1024 |