pub mod abi;
pub mod contract;
pub mod error;
pub mod handle;
pub mod lease;
pub mod liveliness;
mod lock;
pub mod metadata;
pub mod query;
pub mod runtime_metrics;
pub mod server;
pub mod session;
pub mod time;
pub mod topic;
pub mod tree;
mod outbound;
#[cfg(test)]
mod test_support;
pub use crate::identity::{ExecutionId, ParticipantId, ProducerId, TimelineId};
pub use abi::{Codec, CodecError, CodecId, EncodingError, EncodingMetadata, MessagePack};
pub use contract::{
DeliveryFamily, Direction, Endpoint, EndpointKind, EndpointSemantics, Event, Family, In, Out,
Payload, Query, QueryEndpoint, Robot, RobotEndpoint, Runtime, Sample, Setpoint, State, Stream,
StreamDelivered, Supervisor, WorldClock,
};
pub use error::{BusError, KeyProblem, MetadataProblem, OutboundBound, Result, SessionIdRole};
pub use handle::publisher::{
EventPublisher, SamplePublisher, SetpointPublisher, StatePublisher, StreamPublisher,
};
pub use handle::querier::{DEFAULT_QUERY_TIMEOUT, Querier};
pub use handle::stamp::{StepStamp, StepToken, WorldStepToken};
pub use handle::subscriber::{
EventReceiver, MAX_SETPOINT_SOURCES, MAX_STREAM_SOURCES, Observed, ReceiveTerminal,
SampleReceiver, SetpointReceiver, StateView, StreamEvent, StreamReceiver, TimelineRetention,
};
pub use lease::{
ExclusiveProducerLease, FixedSourceAdmission, FixedSourceLease, LEASE_TRACE_TARGET,
LeaseDecision, LeaseRejection, MAX_READY_PRODUCERS,
};
pub use liveliness::{
KeyLivelinessObserver, KeyLivelinessToken, LivelinessStatus, ParticipantReadyEvent,
ParticipantReadyEvents, ParticipantReadyObserver, ParticipantReadyStatus,
ParticipantReadyToken,
};
pub use metadata::{
BusMetadata, ParticipantSourceIdentity, SourceAttribution, SourceLabel, SourceLabelError,
StreamPosition,
};
pub use query::{QueryCode, QueryError, QueryFailure, QueryResult};
pub use runtime_metrics::{
RuntimeBufferKind, RuntimeDirection, RuntimeMetricKey, RuntimeMetricSnapshot,
};
pub use server::{IncomingQuery, ServerQueryable};
pub use session::{BusCloseReport, BusCloseTimeout, BusFault, BusHandle, BusHealth, BusTerminal};
pub(crate) use session::{BusConfig, BusOwner};
pub use time::{
CaptureStamp, LocalInstant, RetiredTimelines, RobotInstant, RobotTimeError, TimeWindow, Timed,
TimelineMismatch, WallTimestamp,
};
pub use topic::{
AskQuery, KeySegment, KeySegmentError, Publish, ServeQuery, Subscribe, Topic, TopicKind,
WildcardPublish,
};
pub use tree::{BoundEndpoint, TopicSegment};
#[doc(hidden)]
pub mod __compat {
use crate::__compat::surface::{ContractRecord, ContractSurface};
use crate::__compat::wire::DescribeWire;
use crate::bus::abi::CodecId;
use crate::bus::liveliness::PARTICIPANT_LIVELINESS_PREFIX;
use crate::bus::metadata::BusMetadata;
use crate::bus::query::QueryFailure;
use crate::bus::session::{BUS_KEY_PREFIX, ZENOH_WIRE_PROTOCOL_VERSION};
#[must_use]
pub fn contract_surface() -> String {
let mut records = Vec::new();
contract_records(&mut records);
ContractSurface::new(records).canonical_json()
}
pub(crate) fn contract_records(out: &mut Vec<ContractRecord>) {
out.extend([
ContractRecord::envelope("BusMetadata", BusMetadata::wire_schema()),
ContractRecord::envelope("QueryFailure", QueryFailure::wire_schema()),
ContractRecord::identifier(
"bus-key-root",
format!("{BUS_KEY_PREFIX}/{{execution}}"),
),
ContractRecord::identifier(
"bus-key-composition",
format!("{BUS_KEY_PREFIX}/{{execution}}/{{topic}}"),
),
ContractRecord::identifier(
"participant-ready-key",
format!(
"{BUS_KEY_PREFIX}/{{execution}}/{PARTICIPANT_LIVELINESS_PREFIX}/{{participant}}/{{producer}}"
),
),
ContractRecord::identifier("encoding", CodecId::MessagePack.encoding_string()),
ContractRecord::identifier(
"zenoh-wire-protocol",
ZENOH_WIRE_PROTOCOL_VERSION.to_string(),
),
]);
}
#[cfg(test)]
mod tests {
use super::contract_surface;
#[test]
fn the_surface_names_every_bus_owned_wire_fact_deterministically() {
let rendered = contract_surface();
serde_json::from_str::<serde_json::Value>(&rendered).expect("the surface is JSON");
assert_eq!(contract_surface(), rendered);
for expected in [
r#""name":"BusMetadata""#,
r#""name":"QueryFailure""#,
r#""value":"phoxal/{execution}""#,
r#""value":"phoxal/{execution}/{topic}""#,
r#""value":"phoxal/{execution}/liveliness/participants/{participant}/{producer}""#,
r#""value":"phoxal/v0;codec=1""#,
r#""name":"zenoh-wire-protocol","record":"identifier","value":"9""#,
r#""name":"produced_at""#,
] {
assert!(
rendered.contains(expected),
"{expected} missing: {rendered}"
);
}
}
#[test]
fn the_declared_envelope_shapes_are_the_shapes_the_bus_writes() {
use crate::__compat::wire::DescribeWire;
use crate::identity::{ParticipantId, TimelineId};
use crate::bus::abi::CodecId;
use crate::bus::metadata::{
BusMetadata, ParticipantSourceIdentity, SourceAttribution, StreamPosition,
};
use crate::bus::query::{QueryCode, QueryFailure};
use crate::bus::test_support::producer;
use crate::bus::time::{RobotInstant, TimeWindow};
let metadata = BusMetadata {
codec: CodecId::MessagePack.as_u8(),
sequence: 3,
stream_position: Some(StreamPosition { sequence: 1 }),
produced_at: Some(TimeWindow::exact(RobotInstant::new(TimelineId::mint(), 8))),
source: SourceAttribution::Participant(ParticipantSourceIdentity::new(
ParticipantId::new("drive").expect("a participant id"),
producer(1),
)),
};
let json = serde_json::to_value(&metadata).expect("the attachment serializes");
assert_eq!(BusMetadata::wire_schema().conforms(&json), Ok(()));
let failure = QueryFailure::new(QueryCode::NotFound, "no such entity");
let json = serde_json::to_value(&failure).expect("a query failure serializes");
assert_eq!(QueryFailure::wire_schema().conforms(&json), Ok(()));
}
}
}