use std::{
collections::HashMap,
marker::PhantomData,
sync::{
Arc,
atomic::{AtomicBool, Ordering},
},
time::{Duration, Instant},
};
use tokio_util::sync::CancellationToken;
use zenoh::{Result, Wait};
use super::{
GoalId, GoalInfo, GoalStatus, ZAction,
messages::*,
state::{SafeGoalManager, ServerGoalState},
};
use crate::{
Builder, attachment::Attachment, entity::TypeInfo, msg::ZMessage,
topic_name::qualify_topic_name,
};
pub(crate) struct CancelDispatcher {
routes: parking_lot::Mutex<HashMap<GoalId, flume::Sender<zenoh::query::Query>>>,
}
impl CancelDispatcher {
pub(crate) fn new() -> Self {
Self {
routes: parking_lot::Mutex::new(HashMap::new()),
}
}
pub(crate) fn register(&self, goal_id: GoalId) -> flume::Receiver<zenoh::query::Query> {
let (tx, rx) = flume::bounded(4);
self.routes.lock().insert(goal_id, tx);
rx
}
pub(crate) fn deregister(&self, goal_id: GoalId) {
self.routes.lock().remove(&goal_id);
}
pub(crate) fn drain(&self, queue: &Arc<crate::queue::BoundedQueue<zenoh::query::Query>>) {
while let Some(query) = queue.try_recv() {
let Some(payload) = query.payload() else {
tracing::warn!("CancelDispatcher: cancel query has no payload");
continue;
};
let goal_id =
match <CancelGoalServiceRequest as ZMessage>::deserialize(&payload.to_bytes()) {
Ok(r) => r.goal_info.goal_id,
Err(e) => {
tracing::warn!("CancelDispatcher: failed to parse cancel request: {}", e);
continue;
}
};
let routes = self.routes.lock();
if let Some(tx) = routes.get(&goal_id) {
if tx.try_send(query).is_err() {
tracing::warn!(
"CancelDispatcher: per-goal channel full for goal {:?}",
goal_id
);
}
} else {
tracing::warn!(
"CancelDispatcher: no handle registered for goal {:?}",
goal_id
);
}
}
}
}
pub(crate) struct InnerServer<A: ZAction> {
pub(crate) goal_server: Arc<crate::service::ZServer<GoalService<A>>>,
pub(crate) result_server: Arc<crate::service::ZServer<ResultService<A>>>,
pub(crate) cancel_server: Arc<crate::service::ZServer<CancelService<A>>>,
pub(crate) feedback_pub:
Arc<crate::pubsub::ZPub<FeedbackMessage<A>, <FeedbackMessage<A> as ZMessage>::Serdes>>,
pub(crate) status_pub:
Arc<crate::pubsub::ZPub<StatusMessage, <StatusMessage as ZMessage>::Serdes>>,
pub(crate) goal_manager: Arc<SafeGoalManager<A>>,
pub(crate) result_handler_token: CancellationToken,
pub(crate) cancel_dispatcher: Arc<CancelDispatcher>,
}
pub(crate) struct ShutdownGuard {
pub(crate) token: CancellationToken,
}
impl Drop for ShutdownGuard {
fn drop(&mut self) {
tracing::debug!("ZActionServer handle dropped, triggering shutdown");
self.token.cancel();
}
}
pub struct ZActionServerBuilder<'a, A: ZAction> {
pub action_name: String,
pub node: &'a crate::node::ZNode,
pub result_timeout: Duration,
pub goal_timeout: Option<Duration>,
pub goal_service_qos: Option<crate::qos::QosProfile>,
pub result_service_qos: Option<crate::qos::QosProfile>,
pub cancel_service_qos: Option<crate::qos::QosProfile>,
pub feedback_topic_qos: Option<crate::qos::QosProfile>,
pub status_topic_qos: Option<crate::qos::QosProfile>,
pub goal_type_info: Option<TypeInfo>,
pub result_type_info: Option<TypeInfo>,
pub feedback_type_info: Option<TypeInfo>,
pub _phantom: std::marker::PhantomData<A>,
}
impl<'a, A: ZAction> ZActionServerBuilder<'a, A> {
pub fn with_result_timeout(mut self, timeout: Duration) -> Self {
self.result_timeout = timeout;
self
}
pub fn with_goal_timeout(mut self, timeout: Duration) -> Self {
self.goal_timeout = Some(timeout);
self
}
pub fn with_goal_service_qos(mut self, qos: crate::qos::QosProfile) -> Self {
self.goal_service_qos = Some(qos);
self
}
pub fn with_result_service_qos(mut self, qos: crate::qos::QosProfile) -> Self {
self.result_service_qos = Some(qos);
self
}
pub fn with_cancel_service_qos(mut self, qos: crate::qos::QosProfile) -> Self {
self.cancel_service_qos = Some(qos);
self
}
pub fn with_feedback_topic_qos(mut self, qos: crate::qos::QosProfile) -> Self {
self.feedback_topic_qos = Some(qos);
self
}
pub fn with_status_topic_qos(mut self, qos: crate::qos::QosProfile) -> Self {
self.status_topic_qos = Some(qos);
self
}
pub fn with_goal_type_info(mut self, info: TypeInfo) -> Self {
self.goal_type_info = Some(info);
self
}
pub fn with_result_type_info(mut self, info: TypeInfo) -> Self {
self.result_type_info = Some(info);
self
}
pub fn with_feedback_type_info(mut self, info: TypeInfo) -> Self {
self.feedback_type_info = Some(info);
self
}
}
impl<'a, A: ZAction> ZActionServerBuilder<'a, A> {
pub fn new(action_name: &str, node: &'a crate::node::ZNode) -> Self {
Self {
action_name: action_name.to_string(),
node,
result_timeout: Duration::from_secs(10),
goal_timeout: None,
goal_service_qos: None,
result_service_qos: None,
cancel_service_qos: None,
feedback_topic_qos: None,
status_topic_qos: None,
goal_type_info: None,
result_type_info: None,
feedback_type_info: None,
_phantom: std::marker::PhantomData,
}
}
}
fn reply_result<A: ZAction>(query: zenoh::query::Query, result: A::Result, status: GoalStatus) {
let response = GetResultResponse::<A> {
status: status as i8,
result,
};
let response_bytes = <GetResultResponse<A> as ZMessage>::serialize(&response);
let attachment: Attachment = query.attachment().unwrap().try_into().unwrap();
let _ = query
.reply(query.key_expr().clone(), response_bytes)
.attachment(attachment)
.wait();
tracing::debug!("Sent result response");
}
async fn handle_result_requests_legacy_inner<A: ZAction>(
inner: &InnerServer<A>,
query: zenoh::query::Query,
) {
tracing::debug!("Received result request");
let payload = query.payload().unwrap().to_bytes();
let request = match <GetResultRequest as ZMessage>::deserialize(&payload) {
Ok(r) => r,
Err(e) => {
tracing::error!("Failed to deserialize result request: {}", e);
return;
}
};
let goal_id = request.goal_id;
let (result_data, maybe_rx) = inner.goal_manager.modify(|manager| {
if let Some(ServerGoalState::Terminated { result, status, .. }) =
manager.goals.get(&goal_id)
{
(Some((result.clone(), *status)), None)
} else {
let (tx, rx) = tokio::sync::oneshot::channel();
manager.result_futures.entry(goal_id).or_default().push(tx);
(None, Some(rx))
}
});
if let Some((result, status)) = result_data {
tracing::debug!("Goal {:?} already terminated ({:?})", goal_id, status);
reply_result::<A>(query, result, status);
} else if let Some(rx) = maybe_rx {
tokio::spawn(async move {
match rx.await {
Ok((result, status)) => {
tracing::debug!(
"Goal {:?} terminated ({:?}), sending result",
goal_id,
status
);
reply_result::<A>(query, result, status);
}
Err(_) => {
tracing::warn!("Result future dropped for goal {:?}", goal_id);
}
}
});
}
}
impl<'a, A: ZAction> Builder for ZActionServerBuilder<'a, A> {
type Output = ZActionServer<A>;
fn build(self) -> Result<Self::Output> {
let action_name = self.node.remap_rules.apply(&self.action_name);
if action_name.is_empty() {
return Err(zenoh::Error::from("Action name cannot be empty"));
}
let qualified_action_name = qualify_topic_name(
&action_name,
&self.node.entity.namespace,
&self.node.entity.name,
)?;
tracing::debug!(
"Action name: '{}', namespace: '{}', qualified: '{}'",
action_name,
self.node.entity.namespace,
qualified_action_name
);
let goal_service_name = format!("{}/_action/send_goal", qualified_action_name);
let result_service_name = format!("{}/_action/get_result", qualified_action_name);
let cancel_service_name = format!("{}/_action/cancel_goal", qualified_action_name);
let feedback_topic_name = format!("{}/_action/feedback", qualified_action_name);
let status_topic_name = format!("{}/_action/status", qualified_action_name);
let goal_type_info = Some(self.goal_type_info.unwrap_or_else(A::send_goal_type_info));
let mut goal_server_builder = self
.node
.create_service_impl::<GoalService<A>>(&goal_service_name, goal_type_info);
if let Some(qos) = self.goal_service_qos {
goal_server_builder.entity.qos = qos.to_protocol_qos();
}
let goal_server = goal_server_builder.build()?;
let result_type_info = Some(
self.result_type_info
.unwrap_or_else(A::get_result_type_info),
);
let mut result_server_builder = self
.node
.create_service_impl::<ResultService<A>>(&result_service_name, result_type_info);
if let Some(qos) = self.result_service_qos {
result_server_builder.entity.qos = qos.to_protocol_qos();
}
let result_server = result_server_builder.build()?;
tracing::debug!("Created result server for: {}", result_service_name);
let cancel_type_info = Some(A::cancel_goal_type_info());
let mut cancel_server_builder = self
.node
.create_service_impl::<CancelService<A>>(&cancel_service_name, cancel_type_info);
if let Some(qos) = self.cancel_service_qos {
cancel_server_builder.entity.qos = qos.to_protocol_qos();
}
let cancel_server = cancel_server_builder.build()?;
let feedback_type_info = Some(
self.feedback_type_info
.unwrap_or_else(A::feedback_type_info),
);
let mut feedback_pub_builder = self
.node
.create_pub_impl::<FeedbackMessage<A>>(&feedback_topic_name, feedback_type_info);
if let Some(qos) = self.feedback_topic_qos {
feedback_pub_builder.entity.qos = qos.to_protocol_qos();
}
let feedback_pub = feedback_pub_builder.build()?;
let status_type_info = Some(A::status_type_info());
let mut status_pub_builder = self
.node
.create_pub_impl::<StatusMessage>(&status_topic_name, status_type_info);
if let Some(qos) = self.status_topic_qos {
status_pub_builder.entity.qos = qos.to_protocol_qos();
}
let status_pub = status_pub_builder.build()?;
let goal_manager = Arc::new(SafeGoalManager::new(self.result_timeout, self.goal_timeout));
let cancellation_token = CancellationToken::new();
let result_handler_token = CancellationToken::new();
let inner = Arc::new(InnerServer {
goal_server: Arc::new(goal_server),
result_server: Arc::new(result_server),
cancel_server: Arc::new(cancel_server),
feedback_pub: Arc::new(feedback_pub),
status_pub: Arc::new(status_pub),
goal_manager,
result_handler_token: result_handler_token.clone(),
cancel_dispatcher: Arc::new(CancelDispatcher::new()),
});
let weak_inner = Arc::downgrade(&inner);
let global_shutdown = cancellation_token.clone();
let handler_token = result_handler_token.clone();
tokio::spawn(async move {
tokio::select! {
_ = global_shutdown.cancelled() => {
tracing::debug!("Result handler stopping due to global shutdown");
},
_ = handler_token.cancelled() => {
tracing::debug!("Result handler stopping - switching to full driver mode");
},
_ = async {
while let Some(inner) = weak_inner.upgrade() {
let query = inner.result_server.queue().recv_async().await;
handle_result_requests_legacy_inner(&inner, query).await;
}
} => {},
}
});
Ok(ZActionServer {
inner,
_shutdown: Arc::new(ShutdownGuard {
token: cancellation_token,
}),
})
}
}
pub struct ZActionServer<A: ZAction> {
inner: Arc<InnerServer<A>>,
_shutdown: Arc<ShutdownGuard>,
}
impl<A: ZAction> std::fmt::Debug for ZActionServer<A> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("ZActionServer")
.field("goal_server", &self.inner.goal_server)
.finish_non_exhaustive()
}
}
impl<A: ZAction> Clone for ZActionServer<A> {
fn clone(&self) -> Self {
Self {
inner: self.inner.clone(),
_shutdown: self._shutdown.clone(),
}
}
}
impl<A: ZAction> ZActionServer<A> {
pub(crate) fn from_inner(inner: Arc<InnerServer<A>>) -> Self {
let dummy_token = CancellationToken::new();
Self {
inner,
_shutdown: Arc::new(ShutdownGuard { token: dummy_token }),
}
}
}
impl<A: ZAction> ZActionServer<A> {
fn goal_server(&self) -> &Arc<crate::service::ZServer<GoalService<A>>> {
&self.inner.goal_server
}
fn result_server(&self) -> &Arc<crate::service::ZServer<ResultService<A>>> {
&self.inner.result_server
}
pub(crate) fn cancel_server(&self) -> &Arc<crate::service::ZServer<CancelService<A>>> {
&self.inner.cancel_server
}
pub(crate) fn cancel_dispatcher(&self) -> &Arc<CancelDispatcher> {
&self.inner.cancel_dispatcher
}
fn feedback_pub(
&self,
) -> &Arc<crate::pubsub::ZPub<FeedbackMessage<A>, <FeedbackMessage<A> as ZMessage>::Serdes>>
{
&self.inner.feedback_pub
}
fn status_pub(
&self,
) -> &Arc<crate::pubsub::ZPub<StatusMessage, <StatusMessage as ZMessage>::Serdes>> {
&self.inner.status_pub
}
pub fn goal_manager(&self) -> &Arc<SafeGoalManager<A>> {
&self.inner.goal_manager
}
fn result_handler_token(&self) -> &CancellationToken {
&self.inner.result_handler_token
}
}
impl<A: ZAction> ZActionServer<A> {
fn publish_status(&self) {
let status_list: Vec<GoalStatusInfo> = self.goal_manager().read(|manager| {
manager
.goals
.iter()
.map(|(goal_id, state)| {
let status = match state {
ServerGoalState::Accepted { .. } => GoalStatus::Accepted,
ServerGoalState::Executing { .. } => GoalStatus::Executing,
ServerGoalState::Canceling { .. } => GoalStatus::Canceling,
ServerGoalState::Terminated { status, .. } => *status,
};
GoalStatusInfo {
goal_info: GoalInfo::new(*goal_id),
status,
}
})
.collect()
});
let msg = StatusMessage { status_list };
let _ = self.status_pub().publish(&msg);
}
pub async fn recv_goal(&self) -> Result<GoalHandle<A, Requested>> {
let query = self.goal_server().queue().recv_async().await;
let payload = query.payload().unwrap().to_bytes();
let request = <SendGoalRequest<A> as ZMessage>::deserialize(&payload)
.map_err(|e| zenoh::Error::from(e.to_string()))?;
Ok(GoalHandle {
goal: request.goal,
info: GoalInfo::new(request.goal_id),
server: self.clone(),
query: Some(query),
cancel_flag: None,
cancel_rx: None,
_state: PhantomData,
})
}
pub async fn recv_cancel(&self) -> Result<(CancelGoalServiceRequest, zenoh::query::Query)> {
let query = self.cancel_server().queue().recv_async().await;
let payload = query.payload().unwrap().to_bytes();
let request = <CancelGoalServiceRequest as ZMessage>::deserialize(&payload)
.map_err(|e| zenoh::Error::from(e.to_string()))?;
Ok((request, query))
}
pub fn is_cancel_request_ready(&self) -> bool {
!self.cancel_server().queue().is_empty()
}
pub fn request_cancel(&self, goal_id: GoalId) -> bool {
self.goal_manager().read(|manager| {
if let Some(ServerGoalState::Executing { cancel_flag, .. }) =
manager.goals.get(&goal_id)
{
cancel_flag.store(true, Ordering::Relaxed);
true
} else {
false
}
})
}
pub async fn recv_result_request(&self) -> Result<(GoalId, zenoh::query::Query)> {
let query = self.result_server().queue().recv_async().await;
let payload = query.payload().unwrap().to_bytes();
let request = <ResultRequest as ZMessage>::deserialize(&payload)
.map_err(|e| zenoh::Error::from(e.to_string()))?;
Ok((request.goal_id, query))
}
pub fn send_goal_response_low(
&self,
query: &zenoh::query::Query,
response: &GoalResponse,
) -> Result<()> {
let response_bytes = <GoalResponse as ZMessage>::serialize(response);
let attachment: Attachment = query.attachment().unwrap().try_into().unwrap();
let _ = query
.reply(query.key_expr().clone(), response_bytes)
.attachment(attachment)
.wait();
Ok(())
}
pub async fn recv_cancel_request_low(
&self,
) -> Result<(CancelGoalServiceRequest, zenoh::query::Query)> {
let query = self.cancel_server().queue().recv_async().await;
let payload = query.payload().unwrap().to_bytes();
let request = <CancelGoalServiceRequest as ZMessage>::deserialize(&payload)
.map_err(|e| zenoh::Error::from(e.to_string()))?;
Ok((request, query))
}
pub fn send_cancel_response_low(
&self,
query: &zenoh::query::Query,
response: &CancelGoalServiceResponse,
) -> Result<()> {
let response_bytes = <CancelGoalServiceResponse as ZMessage>::serialize(response);
let attachment: Attachment = query.attachment().unwrap().try_into().unwrap();
let _ = query
.reply(query.key_expr().clone(), response_bytes)
.attachment(attachment)
.wait();
Ok(())
}
pub fn send_result_response_low(
&self,
query: &zenoh::query::Query,
response: &GetResultResponse<A>,
) -> Result<()> {
let response_bytes = <GetResultResponse<A> as ZMessage>::serialize(response);
let attachment: Attachment = query.attachment().unwrap().try_into().unwrap();
let _ = query
.reply(query.key_expr().clone(), response_bytes)
.attachment(attachment)
.wait();
Ok(())
}
pub fn with_handler<F, Fut>(self, handler: F) -> Self
where
F: Fn(GoalHandle<A, Executing>) -> Fut + Send + Sync + 'static,
Fut: std::future::Future<Output = ()> + Send + 'static,
{
tracing::debug!("Cancelling default result handler to switch to full driver mode");
self.result_handler_token().cancel();
let weak_inner = Arc::downgrade(&self.inner);
let shutdown_token = self._shutdown.token.clone();
tokio::spawn(async move {
crate::action::driver::run_driver_loop(weak_inner, shutdown_token, handler).await;
});
self
}
pub fn expire_goals(&self) -> Vec<GoalId> {
let expired = self.goal_manager().modify(|manager| {
let now = Instant::now();
let mut expired = Vec::new();
manager.goals.retain(|goal_id, state| {
let should_expire = match state {
ServerGoalState::Accepted { expires_at, .. }
| ServerGoalState::Executing { expires_at, .. }
| ServerGoalState::Terminated { expires_at, .. } => {
expires_at.is_some_and(|exp| now >= exp)
}
ServerGoalState::Canceling { .. } => false,
};
if should_expire {
expired.push(*goal_id);
false } else {
true }
});
expired
});
if !expired.is_empty() {
self.publish_status();
}
expired
}
pub fn set_result_timeout(&self, timeout: Duration) {
self.goal_manager().modify(|manager| {
manager.result_timeout = timeout;
});
}
pub fn result_timeout(&self) -> Duration {
self.goal_manager().read(|manager| manager.result_timeout)
}
}
pub struct Requested;
pub struct Accepted;
pub struct Executing;
pub type RequestedGoal<A> = GoalHandle<A, Requested>;
pub type AcceptedGoal<A> = GoalHandle<A, Accepted>;
pub type ExecutingGoal<A> = GoalHandle<A, Executing>;
pub struct GoalHandle<A: ZAction, State> {
pub goal: A::Goal,
pub info: GoalInfo,
pub(crate) server: ZActionServer<A>,
pub(crate) query: Option<zenoh::query::Query>,
pub(crate) cancel_flag: Option<Arc<AtomicBool>>,
pub(crate) cancel_rx: Option<flume::Receiver<zenoh::query::Query>>,
pub(crate) _state: PhantomData<State>,
}
impl<A: ZAction> GoalHandle<A, Requested> {
pub fn goal(&self) -> &A::Goal {
&self.goal
}
pub fn info(&self) -> &GoalInfo {
&self.info
}
pub fn accept(mut self) -> GoalHandle<A, Accepted> {
let response = SendGoalResponse {
accepted: true,
stamp_sec: self.info.stamp.sec,
stamp_nanosec: self.info.stamp.nanosec,
};
let response_bytes = <SendGoalResponse as ZMessage>::serialize(&response);
if let Some(query) = self.query.take() {
let attachment: Attachment = query.attachment().unwrap().try_into().unwrap();
let _ = query
.reply(query.key_expr().clone(), response_bytes)
.attachment(attachment)
.wait();
}
self.server.goal_manager().modify(|manager| {
let expires_at = manager.goal_timeout.map(|timeout| Instant::now() + timeout);
manager.goals.insert(
self.info.goal_id,
ServerGoalState::Accepted {
goal: self.goal.clone(),
timestamp: Instant::now(),
expires_at,
},
);
});
self.server.publish_status();
GoalHandle {
goal: self.goal,
info: self.info,
server: self.server,
query: None,
cancel_flag: None,
cancel_rx: None,
_state: PhantomData,
}
}
pub fn reject(mut self) -> Result<()> {
let response = GoalResponse {
accepted: false,
stamp_sec: 0,
stamp_nanosec: 0,
};
let response_bytes = <GoalResponse as ZMessage>::serialize(&response);
if let Some(query) = self.query.take() {
let attachment: Attachment = query.attachment().unwrap().try_into().unwrap();
let _ = query
.reply(query.key_expr().clone(), response_bytes)
.attachment(attachment)
.wait();
}
Ok(())
}
}
impl<A: ZAction> GoalHandle<A, Accepted> {
pub fn goal(&self) -> &A::Goal {
&self.goal
}
pub fn info(&self) -> &GoalInfo {
&self.info
}
pub fn execute(self) -> GoalHandle<A, Executing> {
let cancel_flag = Arc::new(AtomicBool::new(false));
let cancel_rx = self.server.cancel_dispatcher().register(self.info.goal_id);
self.server.goal_manager().modify(|manager| {
let expires_at = manager.goal_timeout.map(|timeout| Instant::now() + timeout);
manager.goals.insert(
self.info.goal_id,
ServerGoalState::Executing {
goal: self.goal.clone(),
cancel_flag: cancel_flag.clone(),
expires_at,
},
);
});
self.server.publish_status();
GoalHandle {
goal: self.goal,
info: self.info,
server: self.server,
query: None,
cancel_flag: Some(cancel_flag),
cancel_rx: Some(cancel_rx),
_state: PhantomData,
}
}
}
impl<A: ZAction> GoalHandle<A, Executing> {
pub fn goal(&self) -> &A::Goal {
&self.goal
}
pub fn info(&self) -> &GoalInfo {
&self.info
}
pub async fn wait_for_feedback_subscriber(
&self,
count: usize,
timeout: std::time::Duration,
) -> bool {
self.server
.feedback_pub()
.wait_for_subscription(count, timeout)
.await
}
pub fn publish_feedback(&self, feedback: A::Feedback) -> Result<()> {
let msg = FeedbackMessage {
goal_id: self.info.goal_id,
feedback,
};
self.server.feedback_pub().publish(&msg)
}
pub fn is_cancel_requested(&self) -> bool {
self.cancel_flag
.as_ref()
.map(|flag| flag.load(Ordering::Relaxed))
.unwrap_or(false)
}
pub fn try_process_cancel(&self) -> bool {
if self.is_cancel_requested() {
return true;
}
self.server
.cancel_dispatcher()
.drain(self.server.cancel_server().queue());
let Some(cancel_rx) = &self.cancel_rx else {
return false;
};
if let Ok(query) = cancel_rx.try_recv() {
let payload = match query.payload() {
Some(p) => p.to_bytes(),
None => return false,
};
let request = match <CancelGoalServiceRequest as ZMessage>::deserialize(&payload) {
Ok(r) => r,
Err(e) => {
tracing::error!("try_process_cancel: deserialize error: {}", e);
return false;
}
};
self.server.request_cancel(self.info.goal_id);
let response = CancelGoalServiceResponse {
return_code: 1,
goals_canceling: vec![request.goal_info],
};
let response_bytes = <CancelGoalServiceResponse as ZMessage>::serialize(&response);
if let Some(raw_attachment) = query.attachment()
&& let Ok(attachment) = Attachment::try_from(raw_attachment)
{
let _ = query
.reply(query.key_expr().clone(), response_bytes)
.attachment(attachment)
.wait();
}
return true;
}
false
}
pub fn succeed(self, result: A::Result) -> Result<()> {
self.terminate(result, GoalStatus::Succeeded)
}
pub fn abort(self, result: A::Result) -> Result<()> {
self.terminate(result, GoalStatus::Aborted)
}
pub fn canceled(self, result: A::Result) -> Result<()> {
self.terminate(result, GoalStatus::Canceled)
}
fn terminate(self, result: A::Result, status: GoalStatus) -> Result<()> {
self.server
.cancel_dispatcher()
.deregister(self.info.goal_id);
let futures_to_notify = self.server.goal_manager().modify(|manager| {
let now = Instant::now();
let expires_at = Some(now + manager.result_timeout);
manager.goals.insert(
self.info.goal_id,
ServerGoalState::Terminated {
result: result.clone(),
status,
timestamp: now,
expires_at,
},
);
manager
.result_futures
.remove(&self.info.goal_id)
.unwrap_or_default()
});
for tx in futures_to_notify {
let _ = tx.send((result.clone(), status));
}
self.server.publish_status();
Ok(())
}
}