Skip to main content

nros_node/executor/
node.rs

1//! Node — borrows the session to create typed entities.
2
3use core::marker::PhantomData;
4
5use nros_core::{RosAction, RosMessage, RosService};
6use nros_rmw::{ActionInfo, QosSettings, ServiceInfo, Session as _, TopicInfo, TransportError};
7
8use crate::{
9    rmw_type_registry::{MessageForRmw, register_type},
10    session,
11};
12
13use super::{
14    handles::{
15        ActionClient, ActionClientCallback, ActionServer, EmbeddedPublisher, EmbeddedServiceClient,
16        EmbeddedServiceServer, ServiceClientCallback, Subscription,
17    },
18    types::NodeError,
19};
20
21// ============================================================================
22// Node
23// ============================================================================
24
25/// Backend-agnostic node — borrows the session to create typed entities.
26pub struct NodeHandle<'a> {
27    name: heapless::String<64>,
28    namespace: heapless::String<64>,
29    session: &'a mut session::ConcreteSession,
30    domain_id: u32,
31    /// Phase 211.H — per-node QoS overrides lowered from the launch
32    /// `qos_overrides.<topic>.<role>.<policy>` params and baked into a
33    /// `&'static` table by the entry codegen. Folded into each entity's
34    /// `QosSettings` at `create_publisher`/`create_subscription` time
35    /// (setup-time, no alloc). Empty (`&[]`) by default → zero cost for
36    /// systems without overrides.
37    qos_overrides: &'static [nros_rmw::QosOverride],
38    /// RFC-0052 W3b.4 — baked contract-monitor table (mirror of
39    /// `qos_overrides`): `create_publisher` attaches the matching
40    /// endpoint's counter cell so `publish` can bump it lock-free.
41    monitors: &'static [crate::executor::monitor::MonitorSpec],
42    /// W3b.5 — baked subscriber age-contract table + epoch clock;
43    /// `create_subscription` attaches the matching endpoint's age cell.
44    age_monitors: &'static [crate::executor::monitor::AgeMonitorSpec],
45    epoch_us_fn: Option<fn() -> u64>,
46}
47
48impl<'a> NodeHandle<'a> {
49    /// Create a new node (called by Executor::create_node).
50    pub(crate) fn new(
51        name: heapless::String<64>,
52        namespace: heapless::String<64>,
53        session: &'a mut session::ConcreteSession,
54        domain_id: u32,
55    ) -> Self {
56        Self {
57            name,
58            namespace,
59            session,
60            domain_id,
61            qos_overrides: &[],
62            monitors: &[],
63            age_monitors: &[],
64            epoch_us_fn: None,
65        }
66    }
67
68    /// Phase 211.H — install the plan's QoS-override table on this node. Called
69    /// by the generated entry BEFORE the component constructs its entities, so
70    /// `create_publisher`/`create_subscription` fold the matching overrides in.
71    /// The table is `&'static` (codegen bakes it as a `static`), so there is no
72    /// lifetime to thread and no runtime allocation. Plan = authority: an
73    /// override for a topic the entity creates is applied transparently (the
74    /// user's `create_publisher(topic)` call is unchanged, matching rclcpp).
75    pub fn set_qos_overrides(&mut self, overrides: &'static [nros_rmw::QosOverride]) {
76        self.qos_overrides = overrides;
77    }
78
79    /// RFC-0052 W3b.4 — install the executor's monitor table on this node
80    /// (called by the entry glue / fixture alongside `set_qos_overrides`).
81    pub fn set_monitors(&mut self, monitors: &'static [crate::executor::monitor::MonitorSpec]) {
82        self.monitors = monitors;
83    }
84
85    /// W3b.5 — install the subscriber age-contract table + epoch clock
86    /// (auto-seeded from the executor's `set_age_table` / config epoch).
87    pub fn set_age_monitors(
88        &mut self,
89        table: &'static [crate::executor::monitor::AgeMonitorSpec],
90        epoch_us: Option<fn() -> u64>,
91    ) {
92        self.age_monitors = table;
93        self.epoch_us_fn = epoch_us;
94    }
95
96    /// The installed QoS-override table (empty unless the entry set one).
97    #[must_use]
98    pub fn qos_overrides(&self) -> &'static [nros_rmw::QosOverride] {
99        self.qos_overrides
100    }
101
102    /// Get the node name.
103    pub fn name(&self) -> &str {
104        &self.name
105    }
106
107    /// Phase 88.12 — return the [`nros_log::Logger`] keyed on the
108    /// node name.
109    ///
110    /// Loggers are interned in nros-log's bounded global table
111    /// ([`nros_log::MAX_LOGGERS`] slots). If the caller has
112    /// pre-registered a `'static Logger` whose name matches this
113    /// node's name (via [`nros_log::register_logger`]), this method
114    /// returns that exact reference — so subsequent `nros_*!` calls
115    /// share per-logger runtime threshold state with any other call
116    /// site that resolves the same name. Otherwise the call returns
117    /// [`nros_log::DEFAULT_LOGGER`], keeping the API total.
118    ///
119    /// ```ignore
120    /// // Pre-register if you want a dedicated threshold:
121    /// static MY_NODE_LOGGER: nros_log::Logger =
122    ///     nros_log::Logger::new("my_node");
123    /// nros_log::register_logger(&MY_NODE_LOGGER);
124    ///
125    /// // Inside any node-creating code:
126    /// let logger = node.logger();
127    /// nros_log::nros_info!(logger, "started; domain = {}", node.domain_id());
128    /// ```
129    #[must_use]
130    pub fn logger(&self) -> &'static nros_log::Logger {
131        nros_log::get_logger(self.name())
132    }
133
134    /// Get the domain ID.
135    pub fn domain_id(&self) -> u32 {
136        self.domain_id
137    }
138
139    /// Set the domain ID.
140    pub fn set_domain_id(&mut self, domain_id: u32) {
141        self.domain_id = domain_id;
142    }
143
144    /// Get a mutable reference to the underlying session.
145    pub fn session_mut(&mut self) -> &mut session::ConcreteSession {
146        self.session
147    }
148
149    // ------------------------------------------------------------------
150    // Routing-info builders (Phase 91.F)
151    //
152    // Every `create_*` below threads the same node identity (domain_id +
153    // name + namespace) into a TopicInfo / ServiceInfo / ActionInfo. The
154    // shape repeats verbatim ~12 times across this file. Centralised
155    // here so a future change to the routing-info shape (e.g. adding a
156    // `with_security_context`) updates one site instead of twelve, and
157    // so the per-`create_*` function bodies focus on the parts that
158    // actually differ between them.
159    // ------------------------------------------------------------------
160
161    // Associated fns (NOT `&self` methods) so the returned `*Info`
162    // value's borrow tracks only the explicit `&str` arguments, not
163    // the whole `Node`. A `&self` form would block the immediately-
164    // following `self.session.create_*(&info, …)` mut borrow on the
165    // `name` / `namespace` reborrow held inside the returned `*Info`,
166    // because going through a method call hides the field-disjoint
167    // path that lets `&self.name` + `&mut self.session` coexist.
168    fn topic_info<'b>(
169        domain_id: u32,
170        node_name: &'b str,
171        namespace: &'b str,
172        topic_name: &'b str,
173        type_name: &'b str,
174        type_hash: &'b str,
175    ) -> TopicInfo<'b> {
176        TopicInfo::new(topic_name, type_name, type_hash)
177            .with_domain(domain_id)
178            .with_node_name(node_name)
179            .with_namespace(namespace)
180    }
181
182    fn service_info<'b>(
183        domain_id: u32,
184        node_name: &'b str,
185        namespace: &'b str,
186        service_name: &'b str,
187        type_name: &'b str,
188        type_hash: &'b str,
189    ) -> ServiceInfo<'b> {
190        ServiceInfo::new(service_name, type_name, type_hash)
191            .with_domain(domain_id)
192            .with_node_name(node_name)
193            .with_namespace(namespace)
194    }
195
196    fn action_info<'b>(
197        domain_id: u32,
198        action_name: &'b str,
199        type_name: &'b str,
200        type_hash: &'b str,
201    ) -> ActionInfo<'b> {
202        // Action root only needs the domain — per-channel ServiceInfo /
203        // TopicInfo derived from action_info.{send_goal,cancel_goal,...}_key()
204        // carry the full node identity via service_info() / topic_info().
205        ActionInfo::new(action_name, type_name, type_hash).with_domain(domain_id)
206    }
207
208    // -- Publishers --
209
210    /// Create a publisher for the given topic.
211    pub fn create_publisher<M: MessageForRmw>(
212        &mut self,
213        topic_name: &str,
214    ) -> Result<EmbeddedPublisher<M>, NodeError> {
215        self.create_publisher_with_qos::<M>(topic_name, QosSettings::default())
216    }
217
218    /// Create a publisher with custom QoS settings.
219    pub fn create_publisher_with_qos<M: MessageForRmw>(
220        &mut self,
221        topic_name: &str,
222        qos: QosSettings,
223    ) -> Result<EmbeddedPublisher<M>, NodeError> {
224        // Phase 212.K.7.6.b — under `rmw-cyclonedds`, ensure the runtime
225        // type-descriptor exists before the cffi vtable creates the
226        // entity. No-op for other RMWs.
227        register_type::<M>()?;
228        // Phase 211.H — fold any plan qos_overrides for this topic+publisher
229        // into the profile (setup-time, no alloc) BEFORE validation, so an
230        // override the backend can't honour still errors loudly below.
231        let qos = qos.apply_overrides(
232            topic_name,
233            nros_rmw::QosOverrideRole::Publisher,
234            self.qos_overrides,
235        );
236        // Phase 108.B — synchronous QoS validation against backend's
237        // `supported_qos_policies()` mask. No silent downgrade.
238        qos.validate_against(nros_rmw::Session::supported_qos_policies(self.session))
239            .map_err(NodeError::Transport)?;
240        let topic = Self::topic_info(
241            self.domain_id,
242            &self.name,
243            &self.namespace,
244            topic_name,
245            <M as RosMessage>::TYPE_NAME,
246            <M as RosMessage>::TYPE_HASH,
247        );
248        let handle = self
249            .session
250            .create_publisher(&topic, qos)
251            .map_err(|_| NodeError::Transport(TransportError::PublisherCreationFailed))?;
252        // RFC-0052 W3b.4 — attach the contracted endpoint's counter cell
253        // (exact topic-name match against the baked table; None = free).
254        let monitor = self
255            .monitors
256            .iter()
257            .find(|m| m.topic == topic_name)
258            .map(|m| m.cell);
259        Ok(EmbeddedPublisher {
260            handle,
261            event_regs: crate::executor::handles::empty_event_regs(),
262            monitor,
263            _phantom: PhantomData,
264        })
265    }
266
267    /// Create a typeless publisher for non-ROS wire formats (e.g. PX4 uORB
268    /// raw POD bytes, custom binary protocols). The caller supplies the
269    /// `type_name` and `type_hash` strings used by backends that need them
270    /// for liveliness/discovery; backends that don't (uORB) can pass any
271    /// stable string.
272    pub fn create_publisher_raw(
273        &mut self,
274        topic_name: &str,
275        type_name: &str,
276        type_hash: &str,
277    ) -> Result<crate::executor::handles::EmbeddedRawPublisher, NodeError> {
278        self.create_publisher_raw_with_qos(topic_name, type_name, type_hash, QosSettings::default())
279    }
280
281    /// Typeless publisher with custom QoS.
282    pub fn create_publisher_raw_with_qos(
283        &mut self,
284        topic_name: &str,
285        type_name: &str,
286        type_hash: &str,
287        qos: QosSettings,
288    ) -> Result<crate::executor::handles::EmbeddedRawPublisher, NodeError> {
289        // Phase 211.H — apply plan qos_overrides (publisher side) before validate.
290        let qos = qos.apply_overrides(
291            topic_name,
292            nros_rmw::QosOverrideRole::Publisher,
293            self.qos_overrides,
294        );
295        qos.validate_against(nros_rmw::Session::supported_qos_policies(self.session))
296            .map_err(NodeError::Transport)?;
297        let topic = Self::topic_info(
298            self.domain_id,
299            &self.name,
300            &self.namespace,
301            topic_name,
302            type_name,
303            type_hash,
304        );
305        let handle = self
306            .session
307            .create_publisher(&topic, qos)
308            .map_err(|_| NodeError::Transport(TransportError::PublisherCreationFailed))?;
309        Ok(crate::executor::handles::EmbeddedRawPublisher {
310            handle,
311            arena: crate::executor::handles::TxArena::new(),
312            event_regs: crate::executor::handles::empty_event_regs(),
313        })
314    }
315
316    /// Phase 189.M1 — the customizable publisher **builder** (the `clone` tier;
317    /// see `docs/design/0022-entity-api-tiers.md`). Pick a mode with `.typed::<M>()`
318    /// or `.generic(type, hash)`, set knobs (`.qos`), then `.build()`. The
319    /// convenient `create_publisher` / `create_publisher_raw` are the `fork`
320    /// tier — sugar over this with defaults.
321    pub fn publisher<'t>(&mut self, topic: &'t str) -> PublisherBuilder<'_, 'a, 't> {
322        PublisherBuilder {
323            node: self,
324            topic,
325            qos: QosSettings::default(),
326        }
327    }
328
329    // -- Subscriptions --
330
331    /// Create a subscription for the given topic.
332    pub fn create_subscription<M: MessageForRmw>(
333        &mut self,
334        topic_name: &str,
335    ) -> Result<Subscription<M>, NodeError> {
336        self.create_subscription_sized::<M, { crate::config::DEFAULT_RX_BUF_SIZE }>(topic_name)
337    }
338
339    /// Create a subscription with custom buffer size.
340    pub fn create_subscription_sized<M: MessageForRmw, const RX_BUF: usize>(
341        &mut self,
342        topic_name: &str,
343    ) -> Result<Subscription<M, RX_BUF>, NodeError> {
344        self.create_subscription_with_qos::<M, RX_BUF>(topic_name, QosSettings::default())
345    }
346
347    /// Create a subscription with custom QoS and buffer size.
348    pub fn create_subscription_with_qos<M: MessageForRmw, const RX_BUF: usize>(
349        &mut self,
350        topic_name: &str,
351        qos: QosSettings,
352    ) -> Result<Subscription<M, RX_BUF>, NodeError> {
353        // Phase 212.K.7.6.b — see `create_publisher_with_qos`.
354        register_type::<M>()?;
355        // Phase 211.H — apply plan qos_overrides (subscription side) before validate.
356        let qos = qos.apply_overrides(
357            topic_name,
358            nros_rmw::QosOverrideRole::Subscription,
359            self.qos_overrides,
360        );
361        qos.validate_against(nros_rmw::Session::supported_qos_policies(self.session))
362            .map_err(NodeError::Transport)?;
363        let topic = Self::topic_info(
364            self.domain_id,
365            &self.name,
366            &self.namespace,
367            topic_name,
368            <M as RosMessage>::TYPE_NAME,
369            <M as RosMessage>::TYPE_HASH,
370        );
371        let handle = self
372            .session
373            .create_subscriber(&topic, qos)
374            .map_err(|_| NodeError::Transport(TransportError::SubscriberCreationFailed))?;
375        // W3b.5 — attach the contracted endpoint's age cell (stamped
376        // types only; needs an epoch source).
377        let age_mon = match (<M as RosMessage>::STAMP_OFFSET, self.epoch_us_fn) {
378            (Some(_), Some(epoch)) => self
379                .age_monitors
380                .iter()
381                .find(|a| a.topic == topic_name)
382                .map(|a| (a.cell, epoch)),
383            _ => None,
384        };
385        Ok(Subscription {
386            handle,
387            buffer: [0u8; RX_BUF],
388            event_regs: crate::executor::handles::empty_event_regs(),
389            age_mon,
390            _phantom: PhantomData,
391        })
392    }
393
394    /// Create a typeless subscription. Caller decodes raw bytes themselves.
395    pub fn create_subscription_raw(
396        &mut self,
397        topic_name: &str,
398        type_name: &str,
399        type_hash: &str,
400    ) -> Result<crate::executor::handles::RawSubscription, NodeError> {
401        self.create_subscription_raw_sized::<{ crate::config::DEFAULT_RX_BUF_SIZE }>(
402            topic_name, type_name, type_hash,
403        )
404    }
405
406    /// Typeless subscription with custom buffer size.
407    pub fn create_subscription_raw_sized<const RX_BUF: usize>(
408        &mut self,
409        topic_name: &str,
410        type_name: &str,
411        type_hash: &str,
412    ) -> Result<crate::executor::handles::RawSubscription<RX_BUF>, NodeError> {
413        // Phase 211.H — apply plan qos_overrides (subscription side) before
414        // validate, mirroring `create_publisher_raw_with_qos`. The raw entity
415        // paths honour node overrides exactly like the typed ones — an
416        // override the active RMW can't meet errors loudly, never silently.
417        let qos = QosSettings::default().apply_overrides(
418            topic_name,
419            nros_rmw::QosOverrideRole::Subscription,
420            self.qos_overrides,
421        );
422        qos.validate_against(nros_rmw::Session::supported_qos_policies(self.session))
423            .map_err(NodeError::Transport)?;
424        let topic = Self::topic_info(
425            self.domain_id,
426            &self.name,
427            &self.namespace,
428            topic_name,
429            type_name,
430            type_hash,
431        );
432        let handle = self
433            .session
434            .create_subscriber(&topic, qos)
435            .map_err(|_| NodeError::Transport(TransportError::SubscriberCreationFailed))?;
436        Ok(crate::executor::handles::RawSubscription {
437            handle,
438            buffer: [0u8; RX_BUF],
439            event_regs: crate::executor::handles::empty_event_regs(),
440        })
441    }
442
443    // -- Services --
444
445    /// Create a service server.
446    pub fn create_service<Svc: RosService>(
447        &mut self,
448        service_name: &str,
449    ) -> Result<EmbeddedServiceServer<Svc>, NodeError>
450    where
451        Svc::Request: MessageForRmw,
452        Svc::Reply: MessageForRmw,
453    {
454        self.create_service_sized::<Svc, { crate::config::DEFAULT_RX_BUF_SIZE }, { crate::config::DEFAULT_RX_BUF_SIZE }>(service_name, QosSettings::services_default())
455    }
456
457    /// Phase 193.2b — service server with an explicit QoS profile (applied to
458    /// both the request + reply endpoints; rclcpp's `create_service(name, qos)`).
459    pub fn create_service_with_qos<Svc: RosService>(
460        &mut self,
461        service_name: &str,
462        qos: QosSettings,
463    ) -> Result<EmbeddedServiceServer<Svc>, NodeError>
464    where
465        Svc::Request: MessageForRmw,
466        Svc::Reply: MessageForRmw,
467    {
468        self.create_service_sized::<Svc, { crate::config::DEFAULT_RX_BUF_SIZE }, { crate::config::DEFAULT_RX_BUF_SIZE }>(service_name, qos)
469    }
470
471    /// Create a service server with custom buffer sizes + QoS.
472    pub fn create_service_sized<Svc: RosService, const REQ_BUF: usize, const REPLY_BUF: usize>(
473        &mut self,
474        service_name: &str,
475        qos: QosSettings,
476    ) -> Result<EmbeddedServiceServer<Svc, REQ_BUF, REPLY_BUF>, NodeError>
477    where
478        Svc::Request: MessageForRmw,
479        Svc::Reply: MessageForRmw,
480    {
481        // Phase 212.K.7.6.b — register both halves of the service round-trip
482        // under cyclonedds. No-op for other RMWs.
483        register_type::<Svc::Request>()?;
484        register_type::<Svc::Reply>()?;
485        // Phase 193.5 — validate the service profile against the backend's
486        // supported policies (mirrors pub/sub); no silent downgrade. RELIABLE is
487        // effectively required for request/reply, so a backend that only honours
488        // a fixed profile rejects an incompatible request here.
489        qos.validate_against(nros_rmw::Session::supported_qos_policies(self.session))
490            .map_err(NodeError::Transport)?;
491        let info = Self::service_info(
492            self.domain_id,
493            &self.name,
494            &self.namespace,
495            service_name,
496            Svc::SERVICE_NAME,
497            Svc::SERVICE_HASH,
498        );
499        let handle = self
500            .session
501            .create_service_server(&info, qos)
502            .map_err(|_| NodeError::Transport(TransportError::ServiceServerCreationFailed))?;
503        Ok(EmbeddedServiceServer {
504            handle,
505            req_buffer: [0u8; REQ_BUF],
506            reply_buffer: [0u8; REPLY_BUF],
507            _phantom: PhantomData,
508        })
509    }
510
511    /// Create a service client.
512    pub fn create_client<Svc: RosService>(
513        &mut self,
514        service_name: &str,
515    ) -> Result<EmbeddedServiceClient<Svc>, NodeError>
516    where
517        Svc::Request: MessageForRmw,
518        Svc::Reply: MessageForRmw,
519    {
520        self.create_client_sized::<Svc, { crate::config::DEFAULT_RX_BUF_SIZE }, { crate::config::DEFAULT_RX_BUF_SIZE }>(service_name, QosSettings::services_default())
521    }
522
523    /// Phase 193.2b — service client with an explicit QoS profile.
524    pub fn create_client_with_qos<Svc: RosService>(
525        &mut self,
526        service_name: &str,
527        qos: QosSettings,
528    ) -> Result<EmbeddedServiceClient<Svc>, NodeError>
529    where
530        Svc::Request: MessageForRmw,
531        Svc::Reply: MessageForRmw,
532    {
533        self.create_client_sized::<Svc, { crate::config::DEFAULT_RX_BUF_SIZE }, { crate::config::DEFAULT_RX_BUF_SIZE }>(service_name, qos)
534    }
535
536    /// Create a service client with custom buffer sizes + QoS.
537    pub fn create_client_sized<Svc: RosService, const REQ_BUF: usize, const REPLY_BUF: usize>(
538        &mut self,
539        service_name: &str,
540        qos: QosSettings,
541    ) -> Result<EmbeddedServiceClient<Svc, REQ_BUF, REPLY_BUF>, NodeError>
542    where
543        Svc::Request: MessageForRmw,
544        Svc::Reply: MessageForRmw,
545    {
546        // Phase 212.K.7.6.b — see `create_service_sized`.
547        register_type::<Svc::Request>()?;
548        register_type::<Svc::Reply>()?;
549        // Phase 193.5 — validate against the backend's supported policies (no
550        // silent downgrade); request/reply effectively requires RELIABLE.
551        qos.validate_against(nros_rmw::Session::supported_qos_policies(self.session))
552            .map_err(NodeError::Transport)?;
553        let info = Self::service_info(
554            self.domain_id,
555            &self.name,
556            &self.namespace,
557            service_name,
558            Svc::SERVICE_NAME,
559            Svc::SERVICE_HASH,
560        );
561        let handle = self
562            .session
563            .create_service_client(&info, qos)
564            .map_err(|_| NodeError::Transport(TransportError::ServiceClientCreationFailed))?;
565        Ok(EmbeddedServiceClient {
566            handle,
567            req_buffer: [0u8; REQ_BUF],
568            reply_buffer: [0u8; REPLY_BUF],
569            in_flight: false,
570            _phantom: PhantomData,
571        })
572    }
573
574    /// Typeless service server. L1 counterpart of [`create_service`]
575    /// for the C / C++ FFI shims and callers that own their own
576    /// scheduler. Returns a [`crate::executor::handles::RawServiceServer`]
577    /// which polls request bytes directly.
578    pub fn create_service_raw(
579        &mut self,
580        service_name: &str,
581        type_name: &str,
582        type_hash: &str,
583    ) -> Result<crate::executor::handles::RawServiceServer, NodeError> {
584        self.create_service_raw_sized::<
585            { crate::config::DEFAULT_RX_BUF_SIZE },
586            { crate::config::DEFAULT_RX_BUF_SIZE },
587        >(service_name, type_name, type_hash)
588    }
589
590    /// Typeless service server with custom buffer sizes.
591    pub fn create_service_raw_sized<const REQ_BUF: usize, const RESP_BUF: usize>(
592        &mut self,
593        service_name: &str,
594        type_name: &str,
595        type_hash: &str,
596    ) -> Result<crate::executor::handles::RawServiceServer<REQ_BUF, RESP_BUF>, NodeError> {
597        let info = Self::service_info(
598            self.domain_id,
599            &self.name,
600            &self.namespace,
601            service_name,
602            type_name,
603            type_hash,
604        );
605        let handle = self
606            .session
607            .create_service_server(&info, QosSettings::services_default())
608            .map_err(|_| NodeError::Transport(TransportError::ServiceServerCreationFailed))?;
609        Ok(crate::executor::handles::RawServiceServer::new(handle))
610    }
611
612    /// Typeless service client. L1 counterpart of [`create_client`].
613    pub fn create_client_raw(
614        &mut self,
615        service_name: &str,
616        type_name: &str,
617        type_hash: &str,
618    ) -> Result<crate::executor::handles::RawServiceClient, NodeError> {
619        self.create_client_raw_sized::<
620            { crate::config::DEFAULT_RX_BUF_SIZE },
621            { crate::config::DEFAULT_RX_BUF_SIZE },
622        >(service_name, type_name, type_hash)
623    }
624
625    /// Typeless service client with custom buffer sizes.
626    pub fn create_client_raw_sized<const REQ_BUF: usize, const REPLY_BUF: usize>(
627        &mut self,
628        service_name: &str,
629        type_name: &str,
630        type_hash: &str,
631    ) -> Result<crate::executor::handles::RawServiceClient<REQ_BUF, REPLY_BUF>, NodeError> {
632        let info = Self::service_info(
633            self.domain_id,
634            &self.name,
635            &self.namespace,
636            service_name,
637            type_name,
638            type_hash,
639        );
640        let handle = self
641            .session
642            .create_service_client(&info, QosSettings::services_default())
643            .map_err(|_| NodeError::Transport(TransportError::ServiceClientCreationFailed))?;
644        Ok(crate::executor::handles::RawServiceClient::new(handle))
645    }
646
647    // -- Actions --
648
649    /// Phase 122.3.c.6 — typeless action server. Builds the 5
650    /// transport channels (`send_goal` / `cancel_goal` / `get_result`
651    /// services + `feedback` / `status` publishers) and returns the
652    /// raw `ActionServerCore` directly. Caller owns scheduling —
653    /// drives `try_recv_goal_request` / `publish_feedback_raw` /
654    /// `complete_goal_raw` / `try_handle_cancel` /
655    /// `try_handle_get_result_raw` on the returned core.
656    pub fn create_action_server_raw(
657        &mut self,
658        action_name: &str,
659        type_name: &str,
660        type_hash: &str,
661    ) -> Result<
662        super::action_core::ActionServerCore<
663            { crate::config::DEFAULT_RX_BUF_SIZE },
664            { crate::config::DEFAULT_RX_BUF_SIZE },
665            { crate::config::DEFAULT_RX_BUF_SIZE },
666            4,
667        >,
668        NodeError,
669    > {
670        self.create_action_server_raw_sized::<
671            { crate::config::DEFAULT_RX_BUF_SIZE },
672            { crate::config::DEFAULT_RX_BUF_SIZE },
673            { crate::config::DEFAULT_RX_BUF_SIZE },
674            4,
675        >(action_name, type_name, type_hash)
676    }
677
678    /// Typeless action server with custom buffer + goal-slot sizes.
679    pub fn create_action_server_raw_sized<
680        const GOAL_BUF: usize,
681        const RESULT_BUF: usize,
682        const FEEDBACK_BUF: usize,
683        const MAX_GOALS: usize,
684    >(
685        &mut self,
686        action_name: &str,
687        type_name: &str,
688        type_hash: &str,
689    ) -> Result<
690        super::action_core::ActionServerCore<GOAL_BUF, RESULT_BUF, FEEDBACK_BUF, MAX_GOALS>,
691        NodeError,
692    > {
693        let action_info = Self::action_info(self.domain_id, action_name, type_name, type_hash);
694
695        let send_goal_keyexpr: heapless::String<256> = action_info.send_goal_key();
696        let send_goal_info = Self::service_info(
697            self.domain_id,
698            &self.name,
699            &self.namespace,
700            &send_goal_keyexpr,
701            type_name,
702            type_hash,
703        );
704        let send_goal_server = self
705            .session
706            .create_service_server(&send_goal_info, QosSettings::services_default())
707            .map_err(|_| NodeError::ActionCreationFailed)?;
708
709        let cancel_goal_keyexpr: heapless::String<256> = action_info.cancel_goal_key();
710        let cancel_goal_info = Self::service_info(
711            self.domain_id,
712            &self.name,
713            &self.namespace,
714            &cancel_goal_keyexpr,
715            "action_msgs::srv::dds_::CancelGoal_",
716            type_hash,
717        );
718        let cancel_goal_server = self
719            .session
720            .create_service_server(&cancel_goal_info, QosSettings::services_default())
721            .map_err(|_| NodeError::ActionCreationFailed)?;
722
723        let get_result_keyexpr: heapless::String<256> = action_info.get_result_key();
724        let get_result_info = Self::service_info(
725            self.domain_id,
726            &self.name,
727            &self.namespace,
728            &get_result_keyexpr,
729            type_name,
730            type_hash,
731        );
732        let get_result_server = self
733            .session
734            .create_service_server(&get_result_info, QosSettings::services_default())
735            .map_err(|_| NodeError::ActionCreationFailed)?;
736
737        let feedback_keyexpr: heapless::String<256> = action_info.feedback_key();
738        let feedback_topic = Self::topic_info(
739            self.domain_id,
740            &self.name,
741            &self.namespace,
742            &feedback_keyexpr,
743            type_name,
744            type_hash,
745        );
746        let feedback_publisher = self
747            .session
748            .create_publisher(&feedback_topic, QosSettings::QOS_PROFILE_DEFAULT)
749            .map_err(|_| NodeError::ActionCreationFailed)?;
750
751        let status_keyexpr: heapless::String<256> = action_info.status_key();
752        let status_topic = Self::topic_info(
753            self.domain_id,
754            &self.name,
755            &self.namespace,
756            &status_keyexpr,
757            "action_msgs::msg::dds_::GoalStatusArray_",
758            type_hash,
759        );
760        let status_publisher = self
761            .session
762            .create_publisher(
763                &status_topic,
764                QosSettings::QOS_PROFILE_ACTION_STATUS_DEFAULT,
765            )
766            .map_err(|_| NodeError::ActionCreationFailed)?;
767
768        Ok(super::action_core::ActionServerCore {
769            send_goal_server,
770            cancel_goal_server,
771            get_result_server,
772            feedback_publisher,
773            status_publisher,
774            active_goals: heapless::Vec::new(),
775            completed_results: heapless::Vec::new(),
776            pending_get_results: heapless::Vec::new(),
777            result_slab: [0u8; RESULT_BUF],
778            result_slab_used: 0,
779            goal_buffer: [0u8; GOAL_BUF],
780            feedback_buffer: [0u8; FEEDBACK_BUF],
781            cancel_buffer: [0u8; 256],
782        })
783    }
784
785    /// Phase 122.3.c.6 — typeless action client. Same shape as
786    /// `create_action_server_raw` but builds the 3 service clients
787    /// + 1 feedback subscriber, returns the raw `ActionClientCore`.
788    pub fn create_action_client_raw(
789        &mut self,
790        action_name: &str,
791        type_name: &str,
792        type_hash: &str,
793    ) -> Result<
794        super::action_core::ActionClientCore<
795            { crate::config::DEFAULT_RX_BUF_SIZE },
796            { crate::config::DEFAULT_RX_BUF_SIZE },
797            { crate::config::DEFAULT_RX_BUF_SIZE },
798        >,
799        NodeError,
800    > {
801        self.create_action_client_raw_sized::<
802            { crate::config::DEFAULT_RX_BUF_SIZE },
803            { crate::config::DEFAULT_RX_BUF_SIZE },
804            { crate::config::DEFAULT_RX_BUF_SIZE },
805        >(action_name, type_name, type_hash)
806    }
807
808    /// Typeless action client with custom buffer sizes.
809    pub fn create_action_client_raw_sized<
810        const GOAL_BUF: usize,
811        const RESULT_BUF: usize,
812        const FEEDBACK_BUF: usize,
813    >(
814        &mut self,
815        action_name: &str,
816        type_name: &str,
817        type_hash: &str,
818    ) -> Result<super::action_core::ActionClientCore<GOAL_BUF, RESULT_BUF, FEEDBACK_BUF>, NodeError>
819    {
820        let action_info = Self::action_info(self.domain_id, action_name, type_name, type_hash);
821
822        let send_goal_keyexpr: heapless::String<256> = action_info.send_goal_key();
823        let send_goal_info = Self::service_info(
824            self.domain_id,
825            &self.name,
826            &self.namespace,
827            &send_goal_keyexpr,
828            type_name,
829            type_hash,
830        );
831        let send_goal_client = self
832            .session
833            .create_service_client(&send_goal_info, QosSettings::services_default())
834            .map_err(|_| NodeError::ActionCreationFailed)?;
835
836        let cancel_goal_keyexpr: heapless::String<256> = action_info.cancel_goal_key();
837        let cancel_goal_info = Self::service_info(
838            self.domain_id,
839            &self.name,
840            &self.namespace,
841            &cancel_goal_keyexpr,
842            "action_msgs::srv::dds_::CancelGoal_",
843            type_hash,
844        );
845        let cancel_goal_client = self
846            .session
847            .create_service_client(&cancel_goal_info, QosSettings::services_default())
848            .map_err(|_| NodeError::ActionCreationFailed)?;
849
850        let get_result_keyexpr: heapless::String<256> = action_info.get_result_key();
851        let get_result_info = Self::service_info(
852            self.domain_id,
853            &self.name,
854            &self.namespace,
855            &get_result_keyexpr,
856            type_name,
857            type_hash,
858        );
859        let get_result_client = self
860            .session
861            .create_service_client(&get_result_info, QosSettings::services_default())
862            .map_err(|_| NodeError::ActionCreationFailed)?;
863
864        let feedback_keyexpr: heapless::String<256> = action_info.feedback_key();
865        let feedback_topic = Self::topic_info(
866            self.domain_id,
867            &self.name,
868            &self.namespace,
869            &feedback_keyexpr,
870            type_name,
871            type_hash,
872        );
873        let feedback_subscriber = self
874            .session
875            .create_subscriber(&feedback_topic, QosSettings::BEST_EFFORT)
876            .map_err(|_| NodeError::ActionCreationFailed)?;
877
878        Ok(super::action_core::ActionClientCore::new(
879            send_goal_client,
880            cancel_goal_client,
881            get_result_client,
882            feedback_subscriber,
883        ))
884    }
885
886    /// Create an action server.
887    pub fn create_action_server<A: RosAction>(
888        &mut self,
889        action_name: &str,
890    ) -> Result<ActionServer<A>, NodeError>
891    where
892        A::Goal: MessageForRmw,
893        A::Result: MessageForRmw,
894        A::Feedback: MessageForRmw,
895        A::SendGoalRequest: MessageForRmw,
896        A::SendGoalResponse: MessageForRmw,
897        A::GetResultRequest: MessageForRmw,
898        A::GetResultResponse: MessageForRmw,
899        A::FeedbackMessage: MessageForRmw,
900    {
901        self.create_action_server_sized::<A, { crate::config::DEFAULT_RX_BUF_SIZE }, { crate::config::DEFAULT_RX_BUF_SIZE }, { crate::config::DEFAULT_RX_BUF_SIZE }, 4>(action_name)
902    }
903
904    /// Create an action server with custom buffer sizes.
905    pub fn create_action_server_sized<
906        A: RosAction,
907        const GOAL_BUF: usize,
908        const RESULT_BUF: usize,
909        const FEEDBACK_BUF: usize,
910        const MAX_GOALS: usize,
911    >(
912        &mut self,
913        action_name: &str,
914    ) -> Result<ActionServer<A, GOAL_BUF, RESULT_BUF, FEEDBACK_BUF, MAX_GOALS>, NodeError>
915    where
916        A::Goal: MessageForRmw,
917        A::Result: MessageForRmw,
918        A::Feedback: MessageForRmw,
919        A::SendGoalRequest: MessageForRmw,
920        A::SendGoalResponse: MessageForRmw,
921        A::GetResultRequest: MessageForRmw,
922        A::GetResultResponse: MessageForRmw,
923        A::FeedbackMessage: MessageForRmw,
924    {
925        // Phase 212.K.7.6.b + K.7.7.c — register the three user-facing
926        // message types AND the five action-protocol envelope types under
927        // cyclonedds. No-op for other RMWs. The envelopes are needed
928        // because the action service shapes (`*_SendGoal_Request`,
929        // `*_GetResult_Response`, …) are the actual on-wire CDR types,
930        // and the C++ Cyclone bridge auto-prepends a cdds_request_header_t
931        // for any TYPE_NAME ending `_Request`/`_Response`/`_Reply`.
932        register_type::<A::Goal>()?;
933        register_type::<A::Result>()?;
934        register_type::<A::Feedback>()?;
935        register_type::<A::SendGoalRequest>()?;
936        register_type::<A::SendGoalResponse>()?;
937        register_type::<A::GetResultRequest>()?;
938        register_type::<A::GetResultResponse>()?;
939        register_type::<A::FeedbackMessage>()?;
940        let action_info =
941            Self::action_info(self.domain_id, action_name, A::ACTION_NAME, A::ACTION_HASH);
942
943        // Each underlying ServiceInfo / TopicInfo also carries the
944        // node identity so the Zenoh shim declares a liveliness token
945        // for it. Without `with_node_name` the shim's
946        // `declare_entity_liveliness` short-circuits (`node_name.and_then`
947        // → None) and `wait_for_action_server` has nothing to find.
948        // Advertise the per-channel service / topic types ROS 2 matches on
949        // (`<Action>_SendGoal` / `<Action>_GetResult` / `<Action>_FeedbackMessage`),
950        // not the bare action type — see `action_core::action_service_base_type`.
951        let send_goal_type = super::action_core::action_service_base_type(
952            <A::SendGoalRequest as RosMessage>::TYPE_NAME,
953            A::ACTION_NAME,
954        );
955        let get_result_type = super::action_core::action_service_base_type(
956            <A::GetResultRequest as RosMessage>::TYPE_NAME,
957            A::ACTION_NAME,
958        );
959        let feedback_type = <A::FeedbackMessage as RosMessage>::TYPE_NAME;
960
961        let send_goal_keyexpr: heapless::String<256> = action_info.send_goal_key();
962        let send_goal_info = Self::service_info(
963            self.domain_id,
964            &self.name,
965            &self.namespace,
966            &send_goal_keyexpr,
967            send_goal_type,
968            A::ACTION_HASH,
969        );
970        let send_goal_server = self
971            .session
972            .create_service_server(&send_goal_info, QosSettings::services_default())
973            .map_err(|_| NodeError::ActionCreationFailed)?;
974
975        let cancel_goal_keyexpr: heapless::String<256> = action_info.cancel_goal_key();
976        let cancel_goal_info = Self::service_info(
977            self.domain_id,
978            &self.name,
979            &self.namespace,
980            &cancel_goal_keyexpr,
981            "action_msgs::srv::dds_::CancelGoal_",
982            A::ACTION_HASH,
983        );
984        let cancel_goal_server = self
985            .session
986            .create_service_server(&cancel_goal_info, QosSettings::services_default())
987            .map_err(|_| NodeError::ActionCreationFailed)?;
988
989        let get_result_keyexpr: heapless::String<256> = action_info.get_result_key();
990        let get_result_info = Self::service_info(
991            self.domain_id,
992            &self.name,
993            &self.namespace,
994            &get_result_keyexpr,
995            get_result_type,
996            A::ACTION_HASH,
997        );
998        let get_result_server = self
999            .session
1000            .create_service_server(&get_result_info, QosSettings::services_default())
1001            .map_err(|_| NodeError::ActionCreationFailed)?;
1002
1003        let feedback_keyexpr: heapless::String<256> = action_info.feedback_key();
1004        let feedback_topic = Self::topic_info(
1005            self.domain_id,
1006            &self.name,
1007            &self.namespace,
1008            &feedback_keyexpr,
1009            feedback_type,
1010            A::ACTION_HASH,
1011        );
1012        let feedback_publisher = self
1013            .session
1014            .create_publisher(&feedback_topic, QosSettings::QOS_PROFILE_DEFAULT)
1015            .map_err(|_| NodeError::ActionCreationFailed)?;
1016
1017        let status_keyexpr: heapless::String<256> = action_info.status_key();
1018        let status_topic = Self::topic_info(
1019            self.domain_id,
1020            &self.name,
1021            &self.namespace,
1022            &status_keyexpr,
1023            "action_msgs::msg::dds_::GoalStatusArray_",
1024            A::ACTION_HASH,
1025        );
1026        let status_publisher = self
1027            .session
1028            .create_publisher(
1029                &status_topic,
1030                QosSettings::QOS_PROFILE_ACTION_STATUS_DEFAULT,
1031            )
1032            .map_err(|_| NodeError::ActionCreationFailed)?;
1033
1034        Ok(ActionServer {
1035            core: super::action_core::ActionServerCore {
1036                send_goal_server,
1037                cancel_goal_server,
1038                get_result_server,
1039                feedback_publisher,
1040                status_publisher,
1041                active_goals: heapless::Vec::new(),
1042                completed_results: heapless::Vec::new(),
1043                pending_get_results: heapless::Vec::new(),
1044                result_slab: [0u8; RESULT_BUF],
1045                result_slab_used: 0,
1046                goal_buffer: [0u8; GOAL_BUF],
1047                feedback_buffer: [0u8; FEEDBACK_BUF],
1048                cancel_buffer: [0u8; 256],
1049            },
1050            typed_goals: heapless::Vec::new(),
1051            completed_goals: heapless::Vec::new(),
1052        })
1053    }
1054
1055    /// Create an action client.
1056    pub fn create_action_client<A: RosAction>(
1057        &mut self,
1058        action_name: &str,
1059    ) -> Result<ActionClient<A>, NodeError>
1060    where
1061        A::Goal: MessageForRmw,
1062        A::Result: MessageForRmw,
1063        A::Feedback: MessageForRmw,
1064        A::SendGoalRequest: MessageForRmw,
1065        A::SendGoalResponse: MessageForRmw,
1066        A::GetResultRequest: MessageForRmw,
1067        A::GetResultResponse: MessageForRmw,
1068        A::FeedbackMessage: MessageForRmw,
1069    {
1070        self.create_action_client_sized::<A, { crate::config::DEFAULT_RX_BUF_SIZE }, { crate::config::DEFAULT_RX_BUF_SIZE }, { crate::config::DEFAULT_RX_BUF_SIZE }>(action_name)
1071    }
1072
1073    /// Create an action client with custom buffer sizes.
1074    pub fn create_action_client_sized<
1075        A: RosAction,
1076        const GOAL_BUF: usize,
1077        const RESULT_BUF: usize,
1078        const FEEDBACK_BUF: usize,
1079    >(
1080        &mut self,
1081        action_name: &str,
1082    ) -> Result<ActionClient<A, GOAL_BUF, RESULT_BUF, FEEDBACK_BUF>, NodeError>
1083    where
1084        A::Goal: MessageForRmw,
1085        A::Result: MessageForRmw,
1086        A::Feedback: MessageForRmw,
1087        A::SendGoalRequest: MessageForRmw,
1088        A::SendGoalResponse: MessageForRmw,
1089        A::GetResultRequest: MessageForRmw,
1090        A::GetResultResponse: MessageForRmw,
1091        A::FeedbackMessage: MessageForRmw,
1092    {
1093        // Phase 212.K.7.6.b + K.7.7.c — see `create_action_server_sized`.
1094        register_type::<A::Goal>()?;
1095        register_type::<A::Result>()?;
1096        register_type::<A::Feedback>()?;
1097        register_type::<A::SendGoalRequest>()?;
1098        register_type::<A::SendGoalResponse>()?;
1099        register_type::<A::GetResultRequest>()?;
1100        register_type::<A::GetResultResponse>()?;
1101        register_type::<A::FeedbackMessage>()?;
1102        let action_info =
1103            Self::action_info(self.domain_id, action_name, A::ACTION_NAME, A::ACTION_HASH);
1104
1105        // Mirror `create_action_server_sized`: thread node identity through
1106        // each underlying ServiceInfo / TopicInfo so the Zenoh shim
1107        // declares the matching client-side liveliness tokens (and so the
1108        // discovery wildcard built from `send_goal_info` ends up in the
1109        // same domain as the server's tokens).
1110        // Same per-channel typing as the server side so the client's requesters
1111        // and feedback reader match a real ROS 2 action server over DDS.
1112        let send_goal_type = super::action_core::action_service_base_type(
1113            <A::SendGoalRequest as RosMessage>::TYPE_NAME,
1114            A::ACTION_NAME,
1115        );
1116        let get_result_type = super::action_core::action_service_base_type(
1117            <A::GetResultRequest as RosMessage>::TYPE_NAME,
1118            A::ACTION_NAME,
1119        );
1120        let feedback_type = <A::FeedbackMessage as RosMessage>::TYPE_NAME;
1121
1122        let send_goal_keyexpr: heapless::String<256> = action_info.send_goal_key();
1123        let send_goal_info = Self::service_info(
1124            self.domain_id,
1125            &self.name,
1126            &self.namespace,
1127            &send_goal_keyexpr,
1128            send_goal_type,
1129            A::ACTION_HASH,
1130        );
1131        let send_goal_client = self
1132            .session
1133            .create_service_client(&send_goal_info, QosSettings::services_default())
1134            .map_err(|_| NodeError::ActionCreationFailed)?;
1135
1136        let cancel_goal_keyexpr: heapless::String<256> = action_info.cancel_goal_key();
1137        let cancel_goal_info = Self::service_info(
1138            self.domain_id,
1139            &self.name,
1140            &self.namespace,
1141            &cancel_goal_keyexpr,
1142            "action_msgs::srv::dds_::CancelGoal_",
1143            A::ACTION_HASH,
1144        );
1145        let cancel_goal_client = self
1146            .session
1147            .create_service_client(&cancel_goal_info, QosSettings::services_default())
1148            .map_err(|_| NodeError::ActionCreationFailed)?;
1149
1150        let get_result_keyexpr: heapless::String<256> = action_info.get_result_key();
1151        let get_result_info = Self::service_info(
1152            self.domain_id,
1153            &self.name,
1154            &self.namespace,
1155            &get_result_keyexpr,
1156            get_result_type,
1157            A::ACTION_HASH,
1158        );
1159        let get_result_client = self
1160            .session
1161            .create_service_client(&get_result_info, QosSettings::services_default())
1162            .map_err(|_| NodeError::ActionCreationFailed)?;
1163
1164        let feedback_keyexpr: heapless::String<256> = action_info.feedback_key();
1165        let feedback_topic = Self::topic_info(
1166            self.domain_id,
1167            &self.name,
1168            &self.namespace,
1169            &feedback_keyexpr,
1170            feedback_type,
1171            A::ACTION_HASH,
1172        );
1173        let feedback_subscriber = self
1174            .session
1175            .create_subscriber(&feedback_topic, QosSettings::BEST_EFFORT)
1176            .map_err(|_| NodeError::ActionCreationFailed)?;
1177
1178        Ok(ActionClient {
1179            core: super::action_core::ActionClientCore {
1180                send_goal_client,
1181                cancel_goal_client,
1182                get_result_client,
1183                feedback_subscriber,
1184                goal_buffer: [0u8; GOAL_BUF],
1185                result_buffer: [0u8; RESULT_BUF],
1186                feedback_buffer: [0u8; FEEDBACK_BUF],
1187                goal_counter: 0,
1188                in_flight_send_goal: false,
1189                in_flight_cancel: false,
1190                in_flight_get_result: false,
1191            },
1192            _phantom: PhantomData,
1193        })
1194    }
1195}
1196
1197// ===================================================================
1198// Phase 189.M1 — entity builders (the `clone` tier)
1199// ===================================================================
1200
1201/// Publisher builder — `node.publisher(topic)`. Choose `.typed::<M>()` or
1202/// `.generic(type, hash)`, optionally `.qos(..)`, then `.build()`.
1203pub struct PublisherBuilder<'n, 'a, 't> {
1204    node: &'n mut NodeHandle<'a>,
1205    topic: &'t str,
1206    qos: QosSettings,
1207}
1208
1209impl<'n, 'a, 't> PublisherBuilder<'n, 'a, 't> {
1210    /// Set the QoS (also settable on the typed/generic builder).
1211    pub fn qos(mut self, qos: QosSettings) -> Self {
1212        self.qos = qos;
1213        self
1214    }
1215
1216    /// Phase 282 (#145) — mark this publisher "express": its samples bypass
1217    /// transport tx batching (sent immediately even when the batching knob is
1218    /// on). For control-tier / latency-sensitive topics.
1219    pub fn tx_express(mut self, express: bool) -> Self {
1220        self.qos.tx_express = express;
1221        self
1222    }
1223
1224    /// Typed publisher for a ROS message `M` (mirrors rclcpp/rclrs).
1225    pub fn typed<M: MessageForRmw>(self) -> TypedPublisherBuilder<'n, 'a, 't, M> {
1226        TypedPublisherBuilder {
1227            node: self.node,
1228            topic: self.topic,
1229            qos: self.qos,
1230            _phantom: PhantomData,
1231        }
1232    }
1233
1234    /// Generic (type-erased) publisher — the rclcpp `create_generic_publisher`
1235    /// form; raw CDR bytes via `publish_raw`.
1236    pub fn generic(
1237        self,
1238        type_name: &'t str,
1239        type_hash: &'t str,
1240    ) -> GenericPublisherBuilder<'n, 'a, 't> {
1241        GenericPublisherBuilder {
1242            node: self.node,
1243            topic: self.topic,
1244            type_name,
1245            type_hash,
1246            qos: self.qos,
1247        }
1248    }
1249}
1250
1251/// Typed publisher builder (`.typed::<M>()`).
1252pub struct TypedPublisherBuilder<'n, 'a, 't, M> {
1253    node: &'n mut NodeHandle<'a>,
1254    topic: &'t str,
1255    qos: QosSettings,
1256    _phantom: PhantomData<M>,
1257}
1258
1259impl<'n, 'a, 't, M: MessageForRmw> TypedPublisherBuilder<'n, 'a, 't, M> {
1260    pub fn qos(mut self, qos: QosSettings) -> Self {
1261        self.qos = qos;
1262        self
1263    }
1264
1265    /// Phase 282 (#145) — see [`PublisherBuilder::tx_express`].
1266    pub fn tx_express(mut self, express: bool) -> Self {
1267        self.qos.tx_express = express;
1268        self
1269    }
1270
1271    pub fn build(self) -> Result<EmbeddedPublisher<M>, NodeError> {
1272        self.node
1273            .create_publisher_with_qos::<M>(self.topic, self.qos)
1274    }
1275}
1276
1277/// Generic (type-erased) publisher builder (`.generic(type, hash)`).
1278pub struct GenericPublisherBuilder<'n, 'a, 't> {
1279    node: &'n mut NodeHandle<'a>,
1280    topic: &'t str,
1281    type_name: &'t str,
1282    type_hash: &'t str,
1283    qos: QosSettings,
1284}
1285
1286impl<'n, 'a, 't> GenericPublisherBuilder<'n, 'a, 't> {
1287    pub fn qos(mut self, qos: QosSettings) -> Self {
1288        self.qos = qos;
1289        self
1290    }
1291
1292    /// Phase 282 (#145) — see [`PublisherBuilder::tx_express`].
1293    pub fn tx_express(mut self, express: bool) -> Self {
1294        self.qos.tx_express = express;
1295        self
1296    }
1297
1298    pub fn build(self) -> Result<crate::executor::handles::EmbeddedRawPublisher, NodeError> {
1299        self.node.create_publisher_raw_with_qos(
1300            self.topic,
1301            self.type_name,
1302            self.type_hash,
1303            self.qos,
1304        )
1305    }
1306}
1307
1308// ============================================================================
1309// Phase 273 (RFC-0047) — CallbackGroup token (rclcpp/rclrs shape)
1310// ============================================================================
1311
1312/// A first-class callback group — a **name-only token** (rclcpp/rclrs shape).
1313///
1314/// Created via [`NodeCtx::create_callback_group`].  Passed to the `_in`
1315/// entity-create variants (`create_timer_in`, `create_subscription_in`,
1316/// `create_publisher_in`) to label entities with a group name.
1317///
1318/// The group is just a name — the actual `SchedContext` binding is seeded at
1319/// boot by `Executor::bind_group_sched` (from `system.toml group_tiers`, phase
1320/// 273 W2) and resolved at entity-registration time by
1321/// `apply_node_default_sched` (phase 273 W1).  No concurrency type (Mutually-
1322/// Exclusive vs Reentrant) is stored here — that is RFC-0047 OQ1 follow-up.
1323pub struct CallbackGroup {
1324    name: heapless::String<32>,
1325}
1326
1327impl CallbackGroup {
1328    /// The group name (e.g. `"ctrl"`, `"telem"`).
1329    pub fn name(&self) -> &str {
1330        &self.name
1331    }
1332}
1333
1334/// An executor-borrowing node handle — `exec.node(id)`. Hosts the
1335/// callback-registering entity builders (subscriptions register into the
1336/// executor's dispatch arena). It is a **short-lived `&mut Executor` borrow**:
1337/// create entities, then drop it before acquiring the next node handle; entity
1338/// handles (`HandleId`, publishers) are owned and outlive it (no `Arc` — see
1339/// `docs/design/0022-entity-api-tiers.md` §Borrow model).
1340pub struct NodeCtx<'e, 's> {
1341    executor: &'e mut super::spin::Executor<'s>,
1342    node_id: super::node_record::NodeId,
1343}
1344
1345impl<'e, 's> NodeCtx<'e, 's> {
1346    pub(crate) fn new(
1347        executor: &'e mut super::spin::Executor<'s>,
1348        node_id: super::node_record::NodeId,
1349    ) -> Self {
1350        Self { executor, node_id }
1351    }
1352
1353    /// Subscription builder (the `clone` tier). Pick a mode with `.typed::<M>()`
1354    /// or `.generic(type, hash)`, set knobs (`.qos`), then `.build(callback)`.
1355    pub fn subscription<'t>(&mut self, topic: &'t str) -> SubscriptionBuilder<'_, 'e, 't, 's> {
1356        SubscriptionBuilder {
1357            ctx: self,
1358            topic,
1359            qos: QosSettings::default(),
1360        }
1361    }
1362
1363    /// Publisher builder (the `clone` tier), symmetric with
1364    /// [`subscription`](Self::subscription). Pick `.typed::<M>()` or
1365    /// `.generic(type, hash)`, set `.qos()`, then `.build()`. The returned
1366    /// publisher handle is owned and outlives this `NodeCtx` — the bridge
1367    /// builds the dest publisher on one ctx, drops it, then registers the
1368    /// source subscription on another (see `0022-entity-api-tiers.md`).
1369    pub fn publisher<'t>(&mut self, topic: &'t str) -> CtxPublisherBuilder<'_, 'e, 't, 's> {
1370        CtxPublisherBuilder {
1371            ctx: self,
1372            topic,
1373            qos: QosSettings::default(),
1374        }
1375    }
1376
1377    /// Convenient typed publisher (the `fork` tier — rclcpp/rclrs shape).
1378    pub fn create_publisher<M: MessageForRmw>(
1379        &mut self,
1380        topic: &str,
1381    ) -> Result<EmbeddedPublisher<M>, NodeError> {
1382        self.executor
1383            .create_publisher_on::<M>(self.node_id, topic, QosSettings::default())
1384    }
1385
1386    /// Convenient generic (type-erased) publisher — rclcpp `create_generic_*`.
1387    pub fn create_generic_publisher(
1388        &mut self,
1389        topic: &str,
1390        type_name: &str,
1391        type_hash: &str,
1392    ) -> Result<crate::executor::handles::EmbeddedRawPublisher, NodeError> {
1393        self.executor.create_publisher_raw_on(
1394            self.node_id,
1395            topic,
1396            type_name,
1397            type_hash,
1398            QosSettings::default(),
1399        )
1400    }
1401
1402    /// Convenient typed subscription (the `fork` tier — rclcpp/rclrs shape).
1403    /// Sugar over the builder with default QoS + buffer.
1404    pub fn create_subscription<M, F>(
1405        &mut self,
1406        topic: &str,
1407        callback: F,
1408    ) -> Result<super::types::HandleId, NodeError>
1409    where
1410        M: MessageForRmw + 'static,
1411        F: FnMut(&M) + 'static,
1412    {
1413        self.executor
1414            .register_subscription_buffered_on::<M, F, { crate::config::DEFAULT_RX_BUF_SIZE }>(
1415                self.node_id,
1416                topic,
1417                QosSettings::default(),
1418                callback,
1419                None, // no group — node default
1420            )
1421    }
1422
1423    // -----------------------------------------------------------------------
1424    // Phase 273 (RFC-0047) — callback-group API (rclcpp/rclrs shape)
1425    // -----------------------------------------------------------------------
1426
1427    /// Create a named callback group — a thin token wrapping the group name.
1428    ///
1429    /// Pass the returned [`CallbackGroup`] to the `_in` create variants
1430    /// (`create_timer_in`, `create_subscription_in`, `create_publisher_in`)
1431    /// to label entities with this group. The executor binds the entity's
1432    /// callback to the `SchedContext` seeded via `bind_group_sched` for
1433    /// `(node_name, namespace, group_name)` (phase 273 W1/W2).
1434    ///
1435    /// ```ignore
1436    /// let ctrl = node.create_callback_group("ctrl");
1437    /// node.create_timer_in(&ctrl, TimerDuration::from_millis(10), || { /* … */ })?;
1438    /// ```
1439    ///
1440    /// Names longer than 32 bytes are silently truncated.
1441    pub fn create_callback_group(&self, name: &str) -> CallbackGroup {
1442        let mut s = heapless::String::<32>::new();
1443        // Silently truncate if the name is too long (defensive; caller
1444        // should use short ASCII group names like "ctrl"/"telem").
1445        for ch in name.chars() {
1446            if s.push(ch).is_err() {
1447                break;
1448            }
1449        }
1450        CallbackGroup { name: s }
1451    }
1452
1453    /// Create a repeating timer **in** a callback group (phase 273, rclcpp shape).
1454    ///
1455    /// The timer's callback is bound to the `SchedContext` associated with
1456    /// `group` in the executor's `group_sched_table` for this node. If no
1457    /// entry was seeded for the group, the node's default `SchedContext` applies
1458    /// (same as `register_timer`). `period` fires the callback repeatedly.
1459    pub fn create_timer_in<F>(
1460        &mut self,
1461        group: &CallbackGroup,
1462        period: crate::timer::TimerDuration,
1463        callback: F,
1464    ) -> Result<super::types::HandleId, NodeError>
1465    where
1466        F: FnMut() + 'static,
1467    {
1468        self.executor
1469            .register_timer_on(Some(self.node_id), period, callback, Some(group.name()))
1470    }
1471
1472    /// Create a typed subscription **in** a callback group (phase 273, rclcpp shape).
1473    ///
1474    /// The subscription's callback is bound to the `SchedContext` associated
1475    /// with `group` in the `group_sched_table` for this node. If no entry was
1476    /// seeded, the node default applies.
1477    pub fn create_subscription_in<M, F>(
1478        &mut self,
1479        group: &CallbackGroup,
1480        topic: &str,
1481        callback: F,
1482    ) -> Result<super::types::HandleId, NodeError>
1483    where
1484        M: MessageForRmw + 'static,
1485        F: FnMut(&M) + 'static,
1486    {
1487        self.executor
1488            .register_subscription_buffered_on::<M, F, { crate::config::DEFAULT_RX_BUF_SIZE }>(
1489                self.node_id,
1490                topic,
1491                QosSettings::default(),
1492                callback,
1493                Some(group.name()),
1494            )
1495    }
1496
1497    /// Create a typed publisher **in** a callback group (phase 273, rclcpp shape).
1498    ///
1499    /// Publishers do not have an executor-dispatched callback, so the group
1500    /// name has no scheduling effect today (publishers are explicitly driven by
1501    /// the user via `publish()`). The API is provided for symmetry and
1502    /// forward-compatibility (intra-process / loaned-message knobs may use
1503    /// it in the future).
1504    pub fn create_publisher_in<M: MessageForRmw>(
1505        &mut self,
1506        _group: &CallbackGroup,
1507        topic: &str,
1508    ) -> Result<EmbeddedPublisher<M>, NodeError> {
1509        // Publishers carry no executor callback slot; group is forward-compat.
1510        self.executor
1511            .create_publisher_on::<M>(self.node_id, topic, QosSettings::default())
1512    }
1513
1514    /// RFC-0041 / Phase 239.1 — callback-based service client (rclcpp
1515    /// `async_send_request(req, cb)` analogue). The reply is delivered to
1516    /// `callback` at `spin_once` (no `Promise` poll). Returns a
1517    /// [`ServiceClientCallback`] send handle; dual-mode — the `Promise`-based
1518    /// [`create_client`](Self::create_client) is unchanged.
1519    pub fn create_client_with_callback<Svc, F>(
1520        &mut self,
1521        service_name: &str,
1522        callback: F,
1523    ) -> Result<ServiceClientCallback<Svc>, NodeError>
1524    where
1525        Svc: RosService + 'static,
1526        Svc::Request: MessageForRmw,
1527        Svc::Reply: MessageForRmw,
1528        F: FnMut(&Svc::Reply) + 'static,
1529    {
1530        self.create_client_with_callback_sized::<
1531            Svc,
1532            F,
1533            { crate::config::DEFAULT_RX_BUF_SIZE },
1534            { crate::config::DEFAULT_RX_BUF_SIZE },
1535        >(service_name, callback)
1536    }
1537
1538    /// Callback-based service client with custom buffer sizes (Phase 239.1).
1539    pub fn create_client_with_callback_sized<Svc, F, const REQ_BUF: usize, const REPLY_BUF: usize>(
1540        &mut self,
1541        service_name: &str,
1542        callback: F,
1543    ) -> Result<ServiceClientCallback<Svc, REQ_BUF, REPLY_BUF>, NodeError>
1544    where
1545        Svc: RosService + 'static,
1546        Svc::Request: MessageForRmw,
1547        Svc::Reply: MessageForRmw,
1548        F: FnMut(&Svc::Reply) + 'static,
1549    {
1550        register_type::<Svc::Request>()?;
1551        register_type::<Svc::Reply>()?;
1552        let (_id, hdr) = self
1553            .executor
1554            .register_service_client_callback::<Svc, F, REPLY_BUF>(
1555                Some(self.node_id),
1556                service_name,
1557                Svc::SERVICE_NAME,
1558                Svc::SERVICE_HASH,
1559                QosSettings::services_default(),
1560                callback,
1561            )?;
1562        Ok(ServiceClientCallback::new(hdr))
1563    }
1564
1565    /// RFC-0041 / Phase 239.2 — callback-based action client (rclcpp
1566    /// `SendGoalOptions{goal_response_callback, feedback_callback,
1567    /// result_callback}` analogue). Goal-response / feedback / result are
1568    /// delivered to the closures at `spin_once`. Returns an
1569    /// [`ActionClientCallback`] send handle (`send_goal` / `get_result`);
1570    /// dual-mode — the `Promise`-based [`create_action_client`](Self::create_action_client)
1571    /// is unchanged.
1572    #[allow(clippy::type_complexity)]
1573    pub fn create_action_client_with_callbacks<A, GRespF, FbF, ResF>(
1574        &mut self,
1575        action_name: &str,
1576        on_goal_response: GRespF,
1577        on_feedback: FbF,
1578        on_result: ResF,
1579    ) -> Result<ActionClientCallback<A>, NodeError>
1580    where
1581        A: RosAction + 'static,
1582        A::Goal: MessageForRmw,
1583        A::Result: MessageForRmw,
1584        A::Feedback: MessageForRmw,
1585        GRespF: FnMut(&nros_core::GoalId, bool) + 'static,
1586        FbF: FnMut(&nros_core::GoalId, &A::Feedback) + 'static,
1587        ResF: FnMut(&nros_core::GoalId, nros_core::GoalStatus, &A::Result) + 'static,
1588    {
1589        self.create_action_client_with_callbacks_sized::<
1590            A,
1591            GRespF,
1592            FbF,
1593            ResF,
1594            { crate::config::DEFAULT_RX_BUF_SIZE },
1595            { crate::config::DEFAULT_RX_BUF_SIZE },
1596            { crate::config::DEFAULT_RX_BUF_SIZE },
1597        >(action_name, on_goal_response, on_feedback, on_result)
1598    }
1599
1600    /// Callback-based action client with custom buffer sizes (Phase 239.2).
1601    #[allow(clippy::type_complexity)]
1602    pub fn create_action_client_with_callbacks_sized<
1603        A,
1604        GRespF,
1605        FbF,
1606        ResF,
1607        const GOAL_BUF: usize,
1608        const RESULT_BUF: usize,
1609        const FEEDBACK_BUF: usize,
1610    >(
1611        &mut self,
1612        action_name: &str,
1613        on_goal_response: GRespF,
1614        on_feedback: FbF,
1615        on_result: ResF,
1616    ) -> Result<ActionClientCallback<A, GOAL_BUF, RESULT_BUF, FEEDBACK_BUF>, NodeError>
1617    where
1618        A: RosAction + 'static,
1619        A::Goal: MessageForRmw,
1620        A::Result: MessageForRmw,
1621        A::Feedback: MessageForRmw,
1622        GRespF: FnMut(&nros_core::GoalId, bool) + 'static,
1623        FbF: FnMut(&nros_core::GoalId, &A::Feedback) + 'static,
1624        ResF: FnMut(&nros_core::GoalId, nros_core::GoalStatus, &A::Result) + 'static,
1625    {
1626        register_type::<A::Goal>()?;
1627        register_type::<A::Result>()?;
1628        register_type::<A::Feedback>()?;
1629        let (_id, core) = self
1630            .executor
1631            .register_action_client_callback::<A, GRespF, FbF, ResF, GOAL_BUF, RESULT_BUF, FEEDBACK_BUF>(
1632                Some(self.node_id),
1633                action_name,
1634                A::ACTION_NAME,
1635                A::ACTION_HASH,
1636                // Feedback is a stream → buffer a short QoS-depth history (Phase
1637                // 239.5). Goal-response / result are single-outstanding (gated).
1638                8u16,
1639                on_goal_response,
1640                on_feedback,
1641                on_result,
1642            )?;
1643        Ok(ActionClientCallback::new(core))
1644    }
1645
1646    /// Convenient generic (type-erased) subscription — rclcpp `create_generic_*`.
1647    pub fn create_generic_subscription<F>(
1648        &mut self,
1649        topic: &str,
1650        type_name: &str,
1651        type_hash: &str,
1652        callback: F,
1653    ) -> Result<super::types::HandleId, NodeError>
1654    where
1655        F: FnMut(&[u8]) + 'static,
1656    {
1657        self.executor
1658            .register_subscription_buffered_raw_on::<F, { crate::config::DEFAULT_RX_BUF_SIZE }>(
1659                self.node_id,
1660                topic,
1661                type_name,
1662                type_hash,
1663                QosSettings::default(),
1664                callback,
1665            )
1666    }
1667
1668    /// Phase 250 (Wave 2) — generic (type-erased) subscription that surfaces E2E
1669    /// [`IntegrityStatus`](nros_rmw::IntegrityStatus) (CRC + sequence gap/dup) to
1670    /// the callback (`FnMut(&[u8], &IntegrityStatus)`). The declarative-`Node`
1671    /// analog of the typed `.typed::<M>().safety()` builder: the validator lives
1672    /// in the `RmwSubscriber`, so the raw bytes + status arrive together without
1673    /// a typed `M`. Wired by the declarative runtime's `.safety()` opt-in.
1674    #[cfg(feature = "safety-e2e")]
1675    pub fn create_generic_subscription_with_integrity<F>(
1676        &mut self,
1677        topic: &str,
1678        type_name: &str,
1679        type_hash: &str,
1680        callback: F,
1681    ) -> Result<super::types::HandleId, NodeError>
1682    where
1683        F: FnMut(&[u8], &nros_rmw::IntegrityStatus) + 'static,
1684    {
1685        self.executor
1686            .register_subscription_buffered_raw_safety_on::<F, { crate::config::DEFAULT_RX_BUF_SIZE }>(
1687                self.node_id,
1688                topic,
1689                type_name,
1690                type_hash,
1691                QosSettings::default(),
1692                callback,
1693            )
1694    }
1695
1696    /// Convenient borrowed (zero-copy) subscription (Phase 229.6, issue 0007 /
1697    /// RFC-0033 `borrowed` mode).
1698    ///
1699    /// `B` is the code-generated borrowed-message marker (e.g. `ImageBorrow`,
1700    /// emitted alongside the owned `Image` for a `.msg` with a `borrowed`-mode
1701    /// field). The callback receives `&B::View<'a>` — a lifetime-carrying
1702    /// message whose unbounded sequence/string fields borrow directly from the
1703    /// receive buffer (no `heapless::Vec` copy); the view is valid only for the
1704    /// callback's duration.
1705    ///
1706    /// Uses `KEEP_LAST(1)` QoS → triple buffer, as borrowed subscriptions
1707    /// require (a single well-defined slot for the callback's borrow). For an
1708    /// explicit deeper queue use the owned
1709    /// [`create_subscription`](Self::create_subscription); a borrowed
1710    /// subscription registered with `KEEP_LAST(N>1)` is rejected.
1711    pub fn create_subscription_borrowed<B, F>(
1712        &mut self,
1713        topic: &str,
1714        callback: F,
1715    ) -> Result<super::types::HandleId, NodeError>
1716    where
1717        B: nros_core::BorrowedMessage + 'static,
1718        F: for<'a> FnMut(&B::View<'a>) + 'static,
1719    {
1720        self.executor
1721            .register_subscription_buffered_borrowed_on::<B, F, { crate::config::DEFAULT_RX_BUF_SIZE }>(
1722                self.node_id,
1723                topic,
1724                QosSettings::default().keep_last(1),
1725                callback,
1726            )
1727    }
1728
1729    /// Service-server builder (the `clone` tier) — `node.service(name)`.
1730    /// Set `.qos()` (defaults to the services profile = RELIABLE+VOLATILE+
1731    /// KEEP_LAST(10)), then `.build::<Svc, _>(callback)` (Phase 193.2).
1732    pub fn service<'t>(&mut self, name: &'t str) -> CtxServiceBuilder<'_, 'e, 't, 's> {
1733        CtxServiceBuilder {
1734            ctx: self,
1735            name,
1736            qos: QosSettings::services_default(),
1737        }
1738    }
1739
1740    /// Convenient service server (the `fork` tier — rclrs/rclcpp shape), default
1741    /// services QoS. Mirror of `create_subscription`.
1742    pub fn create_service<Svc, F>(
1743        &mut self,
1744        name: &str,
1745        callback: F,
1746    ) -> Result<super::types::HandleId, NodeError>
1747    where
1748        Svc: RosService + 'static,
1749        Svc::Request: crate::rmw_type_registry::MessageForRmw,
1750        Svc::Reply: crate::rmw_type_registry::MessageForRmw,
1751        F: FnMut(&Svc::Request) -> Svc::Reply + 'static,
1752    {
1753        self.executor.register_service_sized_on::<
1754            Svc,
1755            F,
1756            { crate::config::DEFAULT_RX_BUF_SIZE },
1757            { crate::config::DEFAULT_RX_BUF_SIZE },
1758        >(self.node_id, name, QosSettings::services_default(), callback)
1759    }
1760}
1761
1762/// Service-server builder on a [`NodeCtx`] — `node.service(name)`.
1763pub struct CtxServiceBuilder<'c, 'e, 't, 's> {
1764    ctx: &'c mut NodeCtx<'e, 's>,
1765    name: &'t str,
1766    qos: QosSettings,
1767}
1768
1769impl<'c, 'e, 't, 's> CtxServiceBuilder<'c, 'e, 't, 's> {
1770    /// Service QoS (applies to both the request + reply endpoints). Defaults to
1771    /// `QosSettings::services_default()`.
1772    pub fn qos(mut self, qos: QosSettings) -> Self {
1773        self.qos = qos;
1774        self
1775    }
1776
1777    pub fn build<Svc, F>(self, callback: F) -> Result<super::types::HandleId, NodeError>
1778    where
1779        Svc: RosService + 'static,
1780        Svc::Request: crate::rmw_type_registry::MessageForRmw,
1781        Svc::Reply: crate::rmw_type_registry::MessageForRmw,
1782        F: FnMut(&Svc::Request) -> Svc::Reply + 'static,
1783    {
1784        self.ctx.executor.register_service_sized_on::<
1785            Svc,
1786            F,
1787            { crate::config::DEFAULT_RX_BUF_SIZE },
1788            { crate::config::DEFAULT_RX_BUF_SIZE },
1789        >(self.ctx.node_id, self.name, self.qos, callback)
1790    }
1791}
1792
1793/// Publisher builder on a [`NodeCtx`] — `node.publisher(topic)`.
1794pub struct CtxPublisherBuilder<'c, 'e, 't, 's> {
1795    ctx: &'c mut NodeCtx<'e, 's>,
1796    topic: &'t str,
1797    qos: QosSettings,
1798}
1799
1800impl<'c, 'e, 't, 's> CtxPublisherBuilder<'c, 'e, 't, 's> {
1801    pub fn qos(mut self, qos: QosSettings) -> Self {
1802        self.qos = qos;
1803        self
1804    }
1805
1806    /// Typed publisher for a ROS message `M`.
1807    pub fn typed<M: MessageForRmw>(self) -> CtxTypedPublisherBuilder<'c, 'e, 't, 's, M> {
1808        CtxTypedPublisherBuilder {
1809            ctx: self.ctx,
1810            topic: self.topic,
1811            qos: self.qos,
1812            _phantom: PhantomData,
1813        }
1814    }
1815
1816    /// Generic (type-erased) publisher.
1817    pub fn generic(
1818        self,
1819        type_name: &'t str,
1820        type_hash: &'t str,
1821    ) -> CtxGenericPublisherBuilder<'c, 'e, 't, 's> {
1822        CtxGenericPublisherBuilder {
1823            ctx: self.ctx,
1824            topic: self.topic,
1825            type_name,
1826            type_hash,
1827            qos: self.qos,
1828        }
1829    }
1830}
1831
1832/// Typed publisher builder on a `NodeCtx` (`.typed::<M>()`).
1833pub struct CtxTypedPublisherBuilder<'c, 'e, 't, 's, M> {
1834    ctx: &'c mut NodeCtx<'e, 's>,
1835    topic: &'t str,
1836    qos: QosSettings,
1837    _phantom: PhantomData<M>,
1838}
1839
1840impl<'c, 'e, 't, 's, M: MessageForRmw> CtxTypedPublisherBuilder<'c, 'e, 't, 's, M> {
1841    pub fn qos(mut self, qos: QosSettings) -> Self {
1842        self.qos = qos;
1843        self
1844    }
1845
1846    pub fn build(self) -> Result<EmbeddedPublisher<M>, NodeError> {
1847        self.ctx
1848            .executor
1849            .create_publisher_on::<M>(self.ctx.node_id, self.topic, self.qos)
1850    }
1851}
1852
1853/// Generic publisher builder on a `NodeCtx` (`.generic(type, hash)`).
1854pub struct CtxGenericPublisherBuilder<'c, 'e, 't, 's> {
1855    ctx: &'c mut NodeCtx<'e, 's>,
1856    topic: &'t str,
1857    type_name: &'t str,
1858    type_hash: &'t str,
1859    qos: QosSettings,
1860}
1861
1862impl<'c, 'e, 't, 's> CtxGenericPublisherBuilder<'c, 'e, 't, 's> {
1863    pub fn qos(mut self, qos: QosSettings) -> Self {
1864        self.qos = qos;
1865        self
1866    }
1867
1868    pub fn build(self) -> Result<crate::executor::handles::EmbeddedRawPublisher, NodeError> {
1869        self.ctx.executor.create_publisher_raw_on(
1870            self.ctx.node_id,
1871            self.topic,
1872            self.type_name,
1873            self.type_hash,
1874            self.qos,
1875        )
1876    }
1877}
1878
1879/// Subscription builder — `node.subscription(topic)`.
1880pub struct SubscriptionBuilder<'c, 'e, 't, 's> {
1881    ctx: &'c mut NodeCtx<'e, 's>,
1882    topic: &'t str,
1883    qos: QosSettings,
1884}
1885
1886impl<'c, 'e, 't, 's> SubscriptionBuilder<'c, 'e, 't, 's> {
1887    pub fn qos(mut self, qos: QosSettings) -> Self {
1888        self.qos = qos;
1889        self
1890    }
1891
1892    /// Typed subscription for a ROS message `M`.
1893    pub fn typed<M: MessageForRmw + 'static>(self) -> TypedSubscriptionBuilder<'c, 'e, 't, 's, M> {
1894        TypedSubscriptionBuilder {
1895            ctx: self.ctx,
1896            topic: self.topic,
1897            qos: self.qos,
1898            sched: None,
1899            _phantom: PhantomData,
1900        }
1901    }
1902
1903    /// Generic (type-erased) subscription — raw CDR bytes to the callback.
1904    pub fn generic(
1905        self,
1906        type_name: &'t str,
1907        type_hash: &'t str,
1908    ) -> GenericSubscriptionBuilder<'c, 'e, 't, 's> {
1909        GenericSubscriptionBuilder {
1910            ctx: self.ctx,
1911            topic: self.topic,
1912            type_name,
1913            type_hash,
1914            qos: self.qos,
1915            sched: None,
1916        }
1917    }
1918}
1919
1920/// Typed subscription builder (`.typed::<M>()`). `RX` is the staging-buffer
1921/// size, set via `.rx_buffer::<N>()` (defaults to `DEFAULT_RX_BUF_SIZE`).
1922pub struct TypedSubscriptionBuilder<
1923    'c,
1924    'e,
1925    't,
1926    's,
1927    M,
1928    const RX: usize = { crate::config::DEFAULT_RX_BUF_SIZE },
1929> {
1930    ctx: &'c mut NodeCtx<'e, 's>,
1931    topic: &'t str,
1932    qos: QosSettings,
1933    sched: Option<super::sched_context::SchedContextId>,
1934    _phantom: PhantomData<M>,
1935}
1936
1937impl<'c, 'e, 't, 's, M: MessageForRmw + 'static, const RX: usize>
1938    TypedSubscriptionBuilder<'c, 'e, 't, 's, M, RX>
1939{
1940    pub fn qos(mut self, qos: QosSettings) -> Self {
1941        self.qos = qos;
1942        self
1943    }
1944
1945    /// Bind the subscription's callback to a scheduling context.
1946    pub fn sched_context(mut self, sc: super::sched_context::SchedContextId) -> Self {
1947        self.sched = Some(sc);
1948        self
1949    }
1950
1951    /// Set the staging-buffer size (const-generic).
1952    pub fn rx_buffer<const N: usize>(self) -> TypedSubscriptionBuilder<'c, 'e, 't, 's, M, N> {
1953        TypedSubscriptionBuilder {
1954            ctx: self.ctx,
1955            topic: self.topic,
1956            qos: self.qos,
1957            sched: self.sched,
1958            _phantom: PhantomData,
1959        }
1960    }
1961
1962    /// Surface per-message [`MessageInfo`](nros_core::MessageInfo) (seq,
1963    /// publisher GID, timestamps) to the callback — `FnMut(&M, Option<&MessageInfo>)`,
1964    /// the rclrs shape. Distinct from the generic builder's `.message_info()`
1965    /// (which yields a `RawMessageInfo` with the wire attachment).
1966    pub fn message_info(self) -> TypedSubInfoBuilder<'c, 'e, 't, 's, M, RX> {
1967        TypedSubInfoBuilder {
1968            ctx: self.ctx,
1969            topic: self.topic,
1970            qos: self.qos,
1971            sched: self.sched,
1972            _phantom: PhantomData,
1973        }
1974    }
1975
1976    /// Surface E2E-safety validation (CRC + sequence gap/duplicate) to the
1977    /// callback — `FnMut(&M, &IntegrityStatus)`.
1978    #[cfg(feature = "safety-e2e")]
1979    pub fn safety(self) -> TypedSubSafetyBuilder<'c, 'e, 't, 's, M, RX> {
1980        TypedSubSafetyBuilder {
1981            ctx: self.ctx,
1982            topic: self.topic,
1983            qos: self.qos,
1984            sched: self.sched,
1985            _phantom: PhantomData,
1986        }
1987    }
1988
1989    pub fn build<F: FnMut(&M) + 'static>(
1990        self,
1991        callback: F,
1992    ) -> Result<super::types::HandleId, NodeError> {
1993        let handle = self
1994            .ctx
1995            .executor
1996            .register_subscription_buffered_on::<M, F, RX>(
1997                self.ctx.node_id,
1998                self.topic,
1999                self.qos,
2000                callback,
2001                None, // group threaded via create_subscription_in; builder uses sched override
2002            )?;
2003        if let Some(sc) = self.sched {
2004            self.ctx.executor.bind_handle_to_sched_context(handle, sc)?;
2005        }
2006        Ok(handle)
2007    }
2008}
2009
2010/// Typed subscription builder with `MessageInfo` (`.typed::<M>().message_info()`).
2011/// Callback is `FnMut(&M, Option<&MessageInfo>)`.
2012pub struct TypedSubInfoBuilder<
2013    'c,
2014    'e,
2015    't,
2016    's,
2017    M,
2018    const RX: usize = { crate::config::DEFAULT_RX_BUF_SIZE },
2019> {
2020    ctx: &'c mut NodeCtx<'e, 's>,
2021    topic: &'t str,
2022    qos: QosSettings,
2023    sched: Option<super::sched_context::SchedContextId>,
2024    _phantom: PhantomData<M>,
2025}
2026
2027impl<'c, 'e, 't, 's, M: MessageForRmw + 'static, const RX: usize>
2028    TypedSubInfoBuilder<'c, 'e, 't, 's, M, RX>
2029{
2030    pub fn qos(mut self, qos: QosSettings) -> Self {
2031        self.qos = qos;
2032        self
2033    }
2034
2035    pub fn sched_context(mut self, sc: super::sched_context::SchedContextId) -> Self {
2036        self.sched = Some(sc);
2037        self
2038    }
2039
2040    pub fn rx_buffer<const N: usize>(self) -> TypedSubInfoBuilder<'c, 'e, 't, 's, M, N> {
2041        TypedSubInfoBuilder {
2042            ctx: self.ctx,
2043            topic: self.topic,
2044            qos: self.qos,
2045            sched: self.sched,
2046            _phantom: PhantomData,
2047        }
2048    }
2049
2050    pub fn build<F: FnMut(&M, Option<&nros_core::MessageInfo>) + 'static>(
2051        self,
2052        callback: F,
2053    ) -> Result<super::types::HandleId, NodeError> {
2054        let handle = self
2055            .ctx
2056            .executor
2057            .register_subscription_with_info_sized_inner::<M, F, RX>(
2058                Some(self.ctx.node_id),
2059                self.topic,
2060                self.qos,
2061                callback,
2062            )?;
2063        if let Some(sc) = self.sched {
2064            self.ctx.executor.bind_handle_to_sched_context(handle, sc)?;
2065        }
2066        Ok(handle)
2067    }
2068}
2069
2070/// Typed subscription builder with E2E-safety validation
2071/// (`.typed::<M>().safety()`). Callback is `FnMut(&M, &IntegrityStatus)`.
2072#[cfg(feature = "safety-e2e")]
2073pub struct TypedSubSafetyBuilder<
2074    'c,
2075    'e,
2076    't,
2077    's,
2078    M,
2079    const RX: usize = { crate::config::DEFAULT_RX_BUF_SIZE },
2080> {
2081    ctx: &'c mut NodeCtx<'e, 's>,
2082    topic: &'t str,
2083    qos: QosSettings,
2084    sched: Option<super::sched_context::SchedContextId>,
2085    _phantom: PhantomData<M>,
2086}
2087
2088#[cfg(feature = "safety-e2e")]
2089impl<'c, 'e, 't, 's, M: MessageForRmw + 'static, const RX: usize>
2090    TypedSubSafetyBuilder<'c, 'e, 't, 's, M, RX>
2091{
2092    pub fn qos(mut self, qos: QosSettings) -> Self {
2093        self.qos = qos;
2094        self
2095    }
2096
2097    pub fn sched_context(mut self, sc: super::sched_context::SchedContextId) -> Self {
2098        self.sched = Some(sc);
2099        self
2100    }
2101
2102    pub fn rx_buffer<const N: usize>(self) -> TypedSubSafetyBuilder<'c, 'e, 't, 's, M, N> {
2103        TypedSubSafetyBuilder {
2104            ctx: self.ctx,
2105            topic: self.topic,
2106            qos: self.qos,
2107            sched: self.sched,
2108            _phantom: PhantomData,
2109        }
2110    }
2111
2112    pub fn build<F: FnMut(&M, &nros_rmw::IntegrityStatus) + 'static>(
2113        self,
2114        callback: F,
2115    ) -> Result<super::types::HandleId, NodeError> {
2116        let handle = self
2117            .ctx
2118            .executor
2119            .register_subscription_with_safety_sized_inner::<M, F, RX>(
2120                Some(self.ctx.node_id),
2121                self.topic,
2122                self.qos,
2123                callback,
2124            )?;
2125        if let Some(sc) = self.sched {
2126            self.ctx.executor.bind_handle_to_sched_context(handle, sc)?;
2127        }
2128        Ok(handle)
2129    }
2130}
2131
2132/// Generic (type-erased) subscription builder (`.generic(type, hash)`).
2133pub struct GenericSubscriptionBuilder<
2134    'c,
2135    'e,
2136    't,
2137    's,
2138    const RX: usize = { crate::config::DEFAULT_RX_BUF_SIZE },
2139> {
2140    ctx: &'c mut NodeCtx<'e, 's>,
2141    topic: &'t str,
2142    type_name: &'t str,
2143    type_hash: &'t str,
2144    qos: QosSettings,
2145    sched: Option<super::sched_context::SchedContextId>,
2146}
2147
2148impl<'c, 'e, 't, 's, const RX: usize> GenericSubscriptionBuilder<'c, 'e, 't, 's, RX> {
2149    pub fn qos(mut self, qos: QosSettings) -> Self {
2150        self.qos = qos;
2151        self
2152    }
2153
2154    pub fn sched_context(mut self, sc: super::sched_context::SchedContextId) -> Self {
2155        self.sched = Some(sc);
2156        self
2157    }
2158
2159    pub fn rx_buffer<const N: usize>(self) -> GenericSubscriptionBuilder<'c, 'e, 't, 's, N> {
2160        GenericSubscriptionBuilder {
2161            ctx: self.ctx,
2162            topic: self.topic,
2163            type_name: self.type_name,
2164            type_hash: self.type_hash,
2165            qos: self.qos,
2166            sched: self.sched,
2167        }
2168    }
2169
2170    /// Surface the sample's wire attachment + metadata to the callback
2171    /// (`FnMut(&[u8], &RawMessageInfo)`). The cross-RMW bridge reads the
2172    /// `bridge_origin` tag from `info.attachment()` for echo suppression.
2173    pub fn message_info(self) -> GenericSubInfoBuilder<'c, 'e, 't, 's, RX> {
2174        GenericSubInfoBuilder {
2175            ctx: self.ctx,
2176            topic: self.topic,
2177            type_name: self.type_name,
2178            type_hash: self.type_hash,
2179            qos: self.qos,
2180            sched: self.sched,
2181        }
2182    }
2183
2184    pub fn build<F: FnMut(&[u8]) + 'static>(
2185        self,
2186        callback: F,
2187    ) -> Result<super::types::HandleId, NodeError> {
2188        let handle = self
2189            .ctx
2190            .executor
2191            .register_subscription_buffered_raw_on::<F, RX>(
2192                self.ctx.node_id,
2193                self.topic,
2194                self.type_name,
2195                self.type_hash,
2196                self.qos,
2197                callback,
2198            )?;
2199        if let Some(sc) = self.sched {
2200            self.ctx.executor.bind_handle_to_sched_context(handle, sc)?;
2201        }
2202        Ok(handle)
2203    }
2204}
2205
2206/// Generic subscription builder with `MessageInfo` surfaced
2207/// (`.message_info()`). Callback is `FnMut(&[u8], &RawMessageInfo)`.
2208pub struct GenericSubInfoBuilder<
2209    'c,
2210    'e,
2211    't,
2212    's,
2213    const RX: usize = { crate::config::DEFAULT_RX_BUF_SIZE },
2214> {
2215    ctx: &'c mut NodeCtx<'e, 's>,
2216    topic: &'t str,
2217    type_name: &'t str,
2218    type_hash: &'t str,
2219    qos: QosSettings,
2220    sched: Option<super::sched_context::SchedContextId>,
2221}
2222
2223impl<'c, 'e, 't, 's, const RX: usize> GenericSubInfoBuilder<'c, 'e, 't, 's, RX> {
2224    pub fn qos(mut self, qos: QosSettings) -> Self {
2225        self.qos = qos;
2226        self
2227    }
2228
2229    pub fn sched_context(mut self, sc: super::sched_context::SchedContextId) -> Self {
2230        self.sched = Some(sc);
2231        self
2232    }
2233
2234    pub fn rx_buffer<const N: usize>(self) -> GenericSubInfoBuilder<'c, 'e, 't, 's, N> {
2235        GenericSubInfoBuilder {
2236            ctx: self.ctx,
2237            topic: self.topic,
2238            type_name: self.type_name,
2239            type_hash: self.type_hash,
2240            qos: self.qos,
2241            sched: self.sched,
2242        }
2243    }
2244
2245    pub fn build<F: FnMut(&[u8], &nros_core::RawMessageInfo) + 'static>(
2246        self,
2247        callback: F,
2248    ) -> Result<super::types::HandleId, NodeError> {
2249        let handle = self
2250            .ctx
2251            .executor
2252            .register_subscription_buffered_raw_info_on::<F, RX>(
2253                self.ctx.node_id,
2254                self.topic,
2255                self.type_name,
2256                self.type_hash,
2257                self.qos,
2258                callback,
2259            )?;
2260        if let Some(sc) = self.sched {
2261            self.ctx.executor.bind_handle_to_sched_context(handle, sc)?;
2262        }
2263        Ok(handle)
2264    }
2265}
2266
2267// `not(feature = "rmw-cffi")` — these tests use the `mock` backend
2268// (`crate::mock`, itself `cfg(all(test, not(rmw-cffi)))`); a workspace test
2269// build that unifies `rmw-cffi` on swaps `ConcreteSession` to the cffi session
2270// and drops `mock`, so the module must drop with it (matches the
2271// `mock_integration` gate in lifecycle_services.rs).
2272#[cfg(all(test, not(feature = "rmw-cffi")))]
2273mod builder_tests {
2274    use super::*;
2275    use crate::{executor::Executor, mock::MockSession};
2276    use nros_core::{CdrReader, CdrWriter, DeserError, Deserialize, SerError, Serialize};
2277
2278    struct TestMsg;
2279    impl RosMessage for TestMsg {
2280        const TYPE_NAME: &'static str = "test/msg/TestMsg";
2281        const TYPE_HASH: &'static str = "test_hash";
2282    }
2283    impl Serialize for TestMsg {
2284        fn serialize(&self, _w: &mut CdrWriter) -> Result<(), SerError> {
2285            Ok(())
2286        }
2287    }
2288    impl Deserialize for TestMsg {
2289        fn deserialize(_r: &mut CdrReader) -> Result<Self, DeserError> {
2290            Ok(Self)
2291        }
2292    }
2293    // Phase 212.K.7.6.b — minimal single-field `Message` impl so
2294    // `TypedPublisherBuilder::build` resolves under the cyclonedds-tightened
2295    // bound AND the runtime register call succeeds. `DescriptorBuilder`
2296    // rejects empty `FIELDS` with `BuildError::EmptySchema`; pretend
2297    // there's one byte so the bridge stub returns a non-NULL pointer.
2298    #[cfg(rmw_needs_type_descriptors)]
2299    impl nros_serdes::schema::Message for TestMsg {
2300        const TYPE_NAME: &'static str = "test/msg/TestMsg";
2301        const FIELDS: &'static [nros_serdes::schema::Field] = &[nros_serdes::schema::Field {
2302            name: "data",
2303            ty: nros_serdes::schema::FieldType::Uint8,
2304            offset: 0,
2305        }];
2306    }
2307
2308    fn s(v: &str) -> heapless::String<64> {
2309        heapless::String::try_from(v).unwrap()
2310    }
2311
2312    #[test]
2313    fn publisher_builder_typed_and_generic() {
2314        let mut session = MockSession::new();
2315        let mut node = NodeHandle::new(s("n"), s("/"), &mut session, 0);
2316
2317        // typed: node.publisher(t).typed::<M>().qos(..).build()
2318        let _typed = node
2319            .publisher("/chatter")
2320            .typed::<TestMsg>()
2321            .qos(QosSettings::default().keep_last(5))
2322            .build()
2323            .expect("typed publisher builds");
2324
2325        // generic: node.publisher(t).qos(..).generic(type, hash).build()
2326        let _generic = node
2327            .publisher("/chatter")
2328            .qos(QosSettings::default())
2329            .generic("std_msgs/msg/Int32", "hash")
2330            .build()
2331            .expect("generic publisher builds");
2332    }
2333
2334    #[test]
2335    fn subscription_builder_and_convenient() {
2336        let mut exec: Executor = Executor::from_session(MockSession::new());
2337        let id = exec.node_builder("n").build().expect("node");
2338
2339        // builder: typed
2340        let _h = exec
2341            .node_mut(id)
2342            .subscription("/chatter")
2343            .typed::<TestMsg>()
2344            .qos(QosSettings::default().keep_last(5))
2345            .build(|_m: &TestMsg| {})
2346            .expect("typed subscription builds");
2347
2348        // builder: generic (raw bytes)
2349        let _g = exec
2350            .node_mut(id)
2351            .subscription("/raw")
2352            .generic("std_msgs/msg/Int32", "hash")
2353            .build(|_b: &[u8]| {})
2354            .expect("generic subscription builds");
2355
2356        // builder: sized + sched-context (slice 3 knobs)
2357        let sc = exec.default_sched_context_id();
2358        let _s = exec
2359            .node_mut(id)
2360            .subscription("/sized")
2361            .typed::<TestMsg>()
2362            .rx_buffer::<64>()
2363            .sched_context(sc)
2364            .build(|_m: &TestMsg| {})
2365            .expect("sized + sched subscription builds");
2366
2367        // convenient (fork tier) — one node-ctx at a time, re-acquired
2368        let _c = exec
2369            .node_mut(id)
2370            .create_subscription::<TestMsg, _>("/conv", |_m: &TestMsg| {})
2371            .expect("convenient typed subscription builds");
2372    }
2373
2374    #[test]
2375    fn generic_message_info_builder() {
2376        // slice 3b — the bridge echo path: generic sub whose callback
2377        // receives the wire attachment via RawMessageInfo.
2378        let mut exec: Executor = Executor::from_session(MockSession::new());
2379        let id = exec.node_builder("n").build().expect("node");
2380
2381        let _i = exec
2382            .node_mut(id)
2383            .subscription("/info")
2384            .generic("std_msgs/msg/Int32", "hash")
2385            .message_info()
2386            .rx_buffer::<256>()
2387            .build(|_payload: &[u8], info: &nros_core::RawMessageInfo| {
2388                let _ = info.attachment();
2389            })
2390            .expect("generic + message_info subscription builds");
2391    }
2392
2393    #[test]
2394    fn typed_message_info_builder() {
2395        // M2.a — typed .message_info() (rclrs shape FnMut(&M, Option<&MessageInfo>)),
2396        // replacing register_subscription_with_info.
2397        let mut exec: Executor = Executor::from_session(MockSession::new());
2398        let id = exec.node_builder("n").build().expect("node");
2399        let _h = exec
2400            .node_mut(id)
2401            .subscription("/chatter")
2402            .typed::<TestMsg>()
2403            .qos(QosSettings::default().keep_last(5))
2404            .message_info()
2405            .build(|_m: &TestMsg, _info: Option<&nros_core::MessageInfo>| {})
2406            .expect("typed + message_info subscription builds");
2407    }
2408
2409    #[cfg(feature = "safety-e2e")]
2410    #[test]
2411    fn typed_safety_builder() {
2412        // M2.a — typed .safety(), replacing register_subscription_with_safety.
2413        let mut exec: Executor = Executor::from_session(MockSession::new());
2414        let id = exec.node_builder("n").build().expect("node");
2415        let _h = exec
2416            .node_mut(id)
2417            .subscription("/chatter")
2418            .typed::<TestMsg>()
2419            .safety()
2420            .build(|_m: &TestMsg, _status: &nros_rmw::IntegrityStatus| {})
2421            .expect("typed + safety subscription builds");
2422    }
2423
2424    #[test]
2425    fn generator_emitted_chain_compiles() {
2426        // Locks the exact builder chain the orchestration generator emits
2427        // for a subscriber (replaces register_subscription_raw_with_qos_sized_on).
2428        let mut exec: Executor = Executor::from_session(MockSession::new());
2429        let id = exec.node_builder("n").build().expect("node");
2430        let _h = exec
2431            .node_mut(id)
2432            .subscription("/topic")
2433            .generic("std_msgs/msg/Int32", "hash")
2434            .qos(QosSettings::default().keep_last(1))
2435            .rx_buffer::<1024>()
2436            .build(|_data: &[u8]| {})
2437            .expect("generator-shape subscription builds");
2438    }
2439
2440    #[test]
2441    fn nodectx_publisher_and_bridge_shape() {
2442        // NodeCtx publisher symmetry + the bridge two-ctx borrow pattern:
2443        // build the dest publisher on one NodeCtx (dropped), then register
2444        // the source subscription on another — the owned publisher outlives.
2445        let mut exec: Executor = Executor::from_session(MockSession::new());
2446        let id = exec.node_builder("n").build().expect("node");
2447
2448        // convenient + builder publisher on NodeCtx
2449        let _p = exec
2450            .node_mut(id)
2451            .create_publisher::<TestMsg>("/p")
2452            .expect("ctx convenient publisher");
2453        let dest_pub = exec
2454            .node_mut(id)
2455            .publisher("/fwd")
2456            .generic("std_msgs/msg/Int32", "hash")
2457            .build()
2458            .expect("ctx generic publisher builds"); // NodeCtx dropped here
2459
2460        // re-borrow exec for the source sub; closure owns dest_pub
2461        let _s = exec
2462            .node_mut(id)
2463            .subscription("/src")
2464            .generic("std_msgs/msg/Int32", "hash")
2465            .message_info()
2466            .build(move |payload: &[u8], _info: &nros_core::RawMessageInfo| {
2467                let _ = dest_pub.publish_raw(payload);
2468            })
2469            .expect("bridge-shape source subscription builds");
2470    }
2471}