Files
MetaCrate/crates/libremetaverse/src/event_queue.rs
Chili Palmer aa269f1785
Some checks failed
Native code generation / deterministic (push) Failing after 1m15s
Imaging and meshing gate / native (push) Successful in 3m58s
Native Rust workspace compile / compile (push) Successful in 4m1s
Implement capability event queue processing (#55)
2026-08-09 16:49:20 +00:00

1121 lines
36 KiB
Rust

//! Bounded native implementation of the long-polling capabilities event queue.
#![allow(clippy::missing_errors_doc)] // Public Result shapes are fixed by the compatibility map.
#![allow(clippy::needless_pass_by_value)] // Mapped APIs preserve owned CLR argument shapes.
#![allow(clippy::type_complexity)] // Delegate signatures are fixed by the compatibility map.
use crate::interfaces::IMessage;
use crate::{Error, HttpCapsClient, Simulator};
use libremetaverse_structured_data::{OSD, OSDMap, OSDParser};
use libremetaverse_types::compat::{CancellationToken, CancellationTokenSource, Object, Uri};
use std::collections::HashMap;
use std::fmt;
use std::panic::{AssertUnwindSafe, catch_unwind};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex, Weak};
use std::thread::{self, JoinHandle};
use std::time::Duration;
const LLSD_XML: &str = "application/llsd+xml";
fn mutex<T>(value: &Mutex<T>) -> std::sync::MutexGuard<'_, T> {
value
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
}
/// Explicit event-queue resource and retry policy.
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct EventQueuePolicy {
pub max_events_per_response: usize,
pub max_event_nodes: usize,
pub max_event_binary_bytes: usize,
pub initial_retry_delay: Duration,
pub maximum_retry_delay: Duration,
pub graceful_shutdown_timeout: Duration,
}
impl Default for EventQueuePolicy {
fn default() -> Self {
Self {
max_events_per_response: 1_024,
max_event_nodes: 100_000,
max_event_binary_bytes: 8 * 1024 * 1024,
initial_retry_delay: Duration::from_secs(1),
maximum_retry_delay: Duration::from_secs(30),
graceful_shutdown_timeout: Duration::from_secs(2),
}
}
}
impl EventQueuePolicy {
fn validate(&self) -> Result<(), Error> {
if self.max_events_per_response == 0
|| self.max_event_nodes == 0
|| self.max_event_binary_bytes == 0
|| self.initial_retry_delay.is_zero()
|| self.maximum_retry_delay < self.initial_retry_delay
|| self.graceful_shutdown_timeout.is_zero()
{
return Err(Error::Argument);
}
Ok(())
}
}
type ConnectedHandler = dyn Fn() + Send + Sync;
/// Native representation of the C# connected delegate.
#[derive(Clone)]
pub struct EventQueueClientConnectedCallback {
handler: Arc<ConnectedHandler>,
}
impl EventQueueClientConnectedCallback {
#[must_use]
pub fn from_handler(handler: impl Fn() + Send + Sync + 'static) -> Self {
Self {
handler: Arc::new(handler),
}
}
pub fn new(_object: Object, _method: isize) -> Result<Self, Error> {
Err(Error::InvalidOperation)
}
pub fn begin_invoke(
&self,
callback: Box<dyn Fn(&dyn std::any::Any) + Send + Sync>,
object: Object,
) -> Result<Box<dyn std::any::Any + Send + Sync>, Error> {
self.invoke()?;
callback(&object);
Ok(Box::new(object))
}
pub fn end_invoke(&self, _result: Box<dyn std::any::Any + Send + Sync>) -> Result<(), Error> {
Ok(())
}
pub fn invoke(&self) -> Result<(), Error> {
(self.handler)();
Ok(())
}
fn invoke_safely(&self) {
let _ = catch_unwind(AssertUnwindSafe(|| (self.handler)()));
}
}
impl fmt::Debug for EventQueueClientConnectedCallback {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.write_str("EventQueueClientConnectedCallback(<handler>)")
}
}
type EventHandler = dyn Fn(String, OSDMap) + Send + Sync;
/// Native representation of the C# raw event delegate.
#[derive(Clone)]
pub struct EventQueueClientEventCallback {
handler: Arc<EventHandler>,
}
impl EventQueueClientEventCallback {
#[must_use]
pub fn from_handler(handler: impl Fn(String, OSDMap) + Send + Sync + 'static) -> Self {
Self {
handler: Arc::new(handler),
}
}
pub fn new(_object: Object, _method: isize) -> Result<Self, Error> {
Err(Error::InvalidOperation)
}
pub fn begin_invoke(
&self,
event_name: String,
body: OSDMap,
callback: Box<dyn Fn(&dyn std::any::Any) + Send + Sync>,
object: Object,
) -> Result<Box<dyn std::any::Any + Send + Sync>, Error> {
self.invoke(event_name, body)?;
callback(&object);
Ok(Box::new(object))
}
pub fn end_invoke(&self, _result: Box<dyn std::any::Any + Send + Sync>) -> Result<(), Error> {
Ok(())
}
pub fn invoke(&self, event_name: String, body: OSDMap) -> Result<(), Error> {
(self.handler)(event_name, body);
Ok(())
}
fn invoke_safely(&self, event_name: String, body: OSDMap) {
let _ = catch_unwind(AssertUnwindSafe(|| (self.handler)(event_name, body)));
}
}
impl fmt::Debug for EventQueueClientEventCallback {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.write_str("EventQueueClientEventCallback(<handler>)")
}
}
#[derive(Clone, Default)]
struct QueueCallbacks {
connected: Option<EventQueueClientConnectedCallback>,
event: Option<EventQueueClientEventCallback>,
}
struct QueueTask {
cancellation: CancellationTokenSource,
graceful: Arc<AtomicBool>,
handle: JoinHandle<()>,
}
enum SimulatorReference {
Strong(Arc<crate::network_manager::SimulatorData>),
Weak(Weak<crate::network_manager::SimulatorData>),
}
impl SimulatorReference {
fn simulator(&self) -> Option<Simulator> {
match self {
Self::Strong(data) => Some(Simulator::native_from_data(Arc::clone(data))),
Self::Weak(data) => Simulator::native_from_weak(data),
}
}
}
struct EventQueueClientInner {
address: Uri,
simulator: SimulatorReference,
http: HttpCapsClient,
policy: EventQueuePolicy,
callbacks: Mutex<QueueCallbacks>,
task: Mutex<Option<QueueTask>>,
running: AtomicBool,
disposed: AtomicBool,
}
impl EventQueueClientInner {
fn simulator(&self) -> Option<Simulator> {
self.simulator.simulator()
}
fn stop_task(&self, immediate: bool) {
let task = mutex(&self.task).take();
let Some(task) = task else {
self.running.store(false, Ordering::Release);
return;
};
task.graceful.store(!immediate, Ordering::Release);
task.cancellation.cancel();
if task.handle.thread().id() != thread::current().id() {
let _ = task.handle.join();
}
self.running.store(false, Ordering::Release);
}
}
impl Drop for EventQueueClientInner {
fn drop(&mut self) {
self.disposed.store(true, Ordering::Release);
if let Some(task) = self
.task
.get_mut()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.take()
{
task.cancellation.cancel();
if task.handle.thread().id() != thread::current().id() {
let _ = task.handle.join();
}
}
}
}
/// Long-polling `EventQueueGet` client with bounded parsing and deterministic ordering.
pub struct EventQueueClient {
pub on_connected: Option<EventQueueClientConnectedCallback>,
pub on_event: Option<EventQueueClientEventCallback>,
inner: Arc<EventQueueClientInner>,
}
impl Clone for EventQueueClient {
fn clone(&self) -> Self {
Self {
on_connected: self.on_connected.clone(),
on_event: self.on_event.clone(),
inner: Arc::clone(&self.inner),
}
}
}
impl EventQueueClient {
pub fn new(event_queue_location: Uri, sim: Simulator) -> Result<Self, Error> {
Self::with_policy(event_queue_location, sim, EventQueuePolicy::default())
}
/// Rust extension allowing deterministic retry timing in offline tests and embedders.
pub fn with_policy(
event_queue_location: Uri,
sim: Simulator,
policy: EventQueuePolicy,
) -> Result<Self, Error> {
Self::with_policy_and_reference(
event_queue_location,
&sim,
policy,
SimulatorReference::Strong(sim.native_data_arc()),
)
}
pub(crate) fn new_for_caps(event_queue_location: Uri, sim: Simulator) -> Result<Self, Error> {
Self::with_policy_and_reference(
event_queue_location,
&sim,
EventQueuePolicy::default(),
SimulatorReference::Weak(sim.native_data_weak()),
)
}
fn with_policy_and_reference(
event_queue_location: Uri,
sim: &Simulator,
policy: EventQueuePolicy,
simulator: SimulatorReference,
) -> Result<Self, Error> {
policy.validate()?;
if !is_http_uri(&event_queue_location) {
return Err(Error::Argument);
}
let http = sim.client.native_http_caps_client();
Ok(Self {
on_connected: None,
on_event: None,
inner: Arc::new(EventQueueClientInner {
address: event_queue_location,
simulator,
http,
policy,
callbacks: Mutex::new(QueueCallbacks::default()),
task: Mutex::new(None),
running: AtomicBool::new(false),
disposed: AtomicBool::new(false),
}),
})
}
pub fn dispose(&self) -> Result<(), Error> {
if self.inner.disposed.swap(true, Ordering::AcqRel) {
return Ok(());
}
self.inner.stop_task(true);
Ok(())
}
pub fn start(&self) -> Result<(), Error> {
if self.inner.disposed.load(Ordering::Acquire) {
return Err(Error::InvalidOperation);
}
let mut slot = mutex(&self.inner.task);
if self.inner.running.load(Ordering::Acquire) {
return Ok(());
}
if let Some(previous) = slot.take()
&& previous.handle.thread().id() != thread::current().id()
{
let _ = previous.handle.join();
}
*mutex(&self.inner.callbacks) = QueueCallbacks {
connected: self.on_connected.clone(),
event: self.on_event.clone(),
};
let cancellation = CancellationTokenSource::new();
let token = cancellation.token();
let graceful = Arc::new(AtomicBool::new(false));
let graceful_worker = Arc::clone(&graceful);
let weak = Arc::downgrade(&self.inner);
let (ready_sender, ready_receiver) = std::sync::mpsc::sync_channel(1);
self.inner.running.store(true, Ordering::Release);
let handle = thread::Builder::new()
.name("libremetaverse-event-queue".to_owned())
.spawn(move || {
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build();
let Ok(runtime) = runtime else {
if let Some(inner) = weak.upgrade() {
inner.running.store(false, Ordering::Release);
}
let _ = ready_sender.send(false);
return;
};
let _ = ready_sender.send(true);
runtime.block_on(run_event_queue(weak.clone(), token, graceful_worker));
if let Some(inner) = weak.upgrade() {
inner.running.store(false, Ordering::Release);
}
})
.map_err(|_| {
self.inner.running.store(false, Ordering::Release);
Error::InvalidOperation
})?;
if ready_receiver.recv_timeout(Duration::from_secs(2)) != Ok(true) {
cancellation.cancel();
let _ = handle.join();
self.inner.running.store(false, Ordering::Release);
return Err(Error::InvalidOperation);
}
*slot = Some(QueueTask {
cancellation,
graceful,
handle,
});
Ok(())
}
pub fn stop(&self, immediate: bool) -> Result<(), Error> {
self.inner.stop_task(immediate);
Ok(())
}
#[must_use]
pub fn running(&self) -> bool {
self.inner.running.load(Ordering::Acquire) && !self.inner.disposed.load(Ordering::Acquire)
}
}
impl fmt::Debug for EventQueueClient {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("EventQueueClient")
.field("running", &self.running())
.field("disposed", &self.inner.disposed.load(Ordering::Acquire))
.finish_non_exhaustive()
}
}
async fn run_event_queue(
weak: Weak<EventQueueClientInner>,
cancellation: CancellationToken,
graceful: Arc<AtomicBool>,
) {
let mut ack = 0_i32;
let mut retry = 0_u32;
loop {
let Some(inner) = weak.upgrade() else {
return;
};
if cancellation.is_cancellation_requested() {
if graceful.load(Ordering::Acquire) {
send_done(&inner, ack).await;
}
return;
}
let Some(simulator) = inner.simulator() else {
return;
};
if !simulator.native_is_connected() {
return;
}
let request = EventQueueAck {
ack_id: ack,
done: false,
};
let Ok(payload) = request
.serialize()
.and_then(|map| OSDParser::serialize_llsd_xml_bytes(OSD::Map(map.snapshot())))
else {
return;
};
let response = inner
.http
.post_with_uri_string_bytes_cancellation_token_i_progress(
inner.address.clone(),
LLSD_XML.to_owned(),
payload,
cancellation.clone(),
None,
)
.await;
drop(simulator);
match response {
Ok((response, data)) if response.is_success_status_code() => {
if let Some(callback) = mutex(&inner.callbacks).connected.clone() {
callback.invoke_safely();
}
match parse_response(response.content_type.as_deref(), data, &inner.policy) {
Ok(parsed) => {
ack = parsed.sequence;
retry = 0;
let callback = mutex(&inner.callbacks).event.clone();
if let Some(callback) = callback {
for event in parsed.events {
callback.invoke_safely(event.name, event.body);
}
}
}
Err(_) => retry = retry.saturating_add(1),
}
}
Ok((response, _)) if matches!(response.status_code, 404 | 410 | 499) => return,
Err(Error::Cancelled) if cancellation.is_cancellation_requested() => {
if graceful.load(Ordering::Acquire) {
send_done(&inner, ack).await;
}
return;
}
Ok(_) | Err(_) => retry = retry.saturating_add(1),
}
if retry > 0 {
let delay = retry_delay(&inner, retry);
tokio::select! {
() = tokio::time::sleep(delay) => {}
() = cancellation.cancelled() => {
if graceful.load(Ordering::Acquire) {
send_done(&inner, ack).await;
}
return;
}
}
}
}
}
async fn send_done(inner: &EventQueueClientInner, ack: i32) {
let Ok(map) = (EventQueueAck {
ack_id: ack,
done: true,
})
.serialize() else {
return;
};
let Ok(payload) = OSDParser::serialize_llsd_xml_bytes(OSD::Map(map.snapshot())) else {
return;
};
let timeout = CancellationTokenSource::new();
let timeout_worker = timeout.clone();
let duration = inner.policy.graceful_shutdown_timeout;
tokio::spawn(async move {
tokio::time::sleep(duration).await;
timeout_worker.cancel();
});
let _ = inner
.http
.post_with_uri_string_bytes_cancellation_token_i_progress(
inner.address.clone(),
LLSD_XML.to_owned(),
payload,
timeout.token(),
None,
)
.await;
}
fn retry_delay(inner: &EventQueueClientInner, retry: u32) -> Duration {
let multiplier = 1_u32
.checked_shl(retry.saturating_sub(1).min(30))
.unwrap_or(u32::MAX);
let base = inner
.policy
.initial_retry_delay
.saturating_mul(multiplier)
.min(inner.policy.maximum_retry_delay);
let max_jitter = base / 8;
if max_jitter.is_zero() {
return base;
}
let mut hash = 0xcbf2_9ce4_8422_2325_u64 ^ u64::from(retry);
for byte in inner.address.0.as_bytes() {
hash ^= u64::from(*byte);
hash = hash.wrapping_mul(0x100_0000_01b3);
}
let jitter_nanos = u64::try_from(max_jitter.as_nanos())
.unwrap_or(u64::MAX)
.saturating_add(1);
base.saturating_add(Duration::from_nanos(hash % jitter_nanos))
}
struct ParsedEvent {
name: String,
body: OSDMap,
}
struct ParsedResponse {
sequence: i32,
events: Vec<ParsedEvent>,
}
fn parse_response(
content_type: Option<&str>,
data: Vec<u8>,
policy: &EventQueuePolicy,
) -> Result<ParsedResponse, Error> {
if data.is_empty() || !is_likely_llsd(content_type, &data) {
return Err(Error::Parse {
position: 0,
context: "event queue response is not LLSD/XML",
});
}
let OSD::Map(mut map) = OSDParser::deserialize_llsd_xml_with_bytes(data)? else {
return Err(Error::Parse {
position: 0,
context: "event queue response is not a map",
});
};
let Some(OSD::Integer(sequence)) = map.remove("id") else {
return Err(Error::Parse {
position: 0,
context: "event queue response has no integer id",
});
};
let Some(OSD::Array(raw_events)) = map.remove("events") else {
return Err(Error::Parse {
position: 0,
context: "event queue response has no event array",
});
};
if raw_events.len() > policy.max_events_per_response {
return Err(Error::Argument);
}
let mut events = Vec::with_capacity(raw_events.len());
for raw_event in raw_events {
let OSD::Map(mut event) = raw_event else {
return Err(Error::Parse {
position: 0,
context: "event queue entry is not a map",
});
};
let name = event
.remove("message")
.ok_or(Error::Parse {
position: 0,
context: "event queue entry has no message",
})?
.as_string()?;
let Some(OSD::Map(body)) = event.remove("body") else {
return Err(Error::Parse {
position: 0,
context: "event queue entry body is not a map",
});
};
OSD::Map(body.clone()).validate_limits(
OSD::DEFAULT_MAX_DEPTH,
policy.max_event_nodes,
policy.max_event_binary_bytes,
)?;
events.push(ParsedEvent {
name,
body: OSDMap::new_with_dictionary(body)?,
});
}
Ok(ParsedResponse { sequence, events })
}
fn is_likely_llsd(content_type: Option<&str>, data: &[u8]) -> bool {
if let Some(content_type) = content_type {
let content_type = content_type.to_ascii_lowercase();
return content_type.contains("xml") || content_type.contains("llsd");
}
let prefix = String::from_utf8_lossy(&data[..data.len().min(256)]);
let prefix = prefix.trim_start();
if prefix.to_ascii_lowercase().starts_with("<!doctype")
|| prefix.to_ascii_lowercase().starts_with("<html")
{
return false;
}
prefix.starts_with("<? LLSD/")
|| prefix.to_ascii_lowercase().starts_with("<?xml")
|| prefix.to_ascii_lowercase().starts_with("<llsd")
|| prefix.starts_with('<')
}
fn is_http_uri(uri: &Uri) -> bool {
reqwest::Url::parse(&uri.0)
.ok()
.is_some_and(|uri| matches!(uri.scheme(), "http" | "https") && uri.host().is_some())
}
/// Request or response body carried by `EventQueueGet`.
pub enum EventMessageBlock {
None,
Ack(EventQueueAck),
Events(EventQueueEvent),
}
impl EventMessageBlock {
pub fn deserialize(&mut self, map: OSDMap) -> Result<(), Error> {
if map.get("ack").is_some() {
let mut value = EventQueueAck::new()?;
value.deserialize(map)?;
*self = Self::Ack(value);
Ok(())
} else if map.get("events").is_some() {
let mut value = EventQueueEvent::new()?;
value.deserialize(map)?;
*self = Self::Events(value);
Ok(())
} else {
Err(Error::Parse {
position: 0,
context: "EventQueueGet has no message block",
})
}
}
pub fn serialize(&self) -> Result<OSDMap, Error> {
match self {
Self::Ack(value) => value.serialize(),
Self::Events(value) => value.serialize(),
Self::None => Err(Error::InvalidOperation),
}
}
}
pub struct EventQueueAck {
pub ack_id: i32,
pub done: bool,
}
impl EventQueueAck {
pub const fn new() -> Result<Self, Error> {
Ok(Self {
ack_id: 0,
done: false,
})
}
pub fn deserialize(&mut self, map: OSDMap) -> Result<(), Error> {
self.ack_id = map.get("ack").unwrap_or_default().as_integer()?;
self.done = map.get("done").unwrap_or_default().as_boolean()?;
Ok(())
}
pub fn serialize(&self) -> Result<OSDMap, Error> {
OSDMap::new_with_dictionary(HashMap::from([
("ack".to_owned(), OSD::Integer(self.ack_id)),
("done".to_owned(), OSD::Boolean(self.done)),
]))
}
}
pub struct EventQueueEventQueueEvent {
pub event_message: Option<Box<dyn IMessage>>,
pub message_key: String,
}
impl EventQueueEventQueueEvent {
pub fn new() -> Result<Self, Error> {
Ok(Self {
event_message: None,
message_key: String::new(),
})
}
}
pub struct EventQueueEvent {
pub message_events: Vec<EventQueueEventQueueEvent>,
pub sequence: i32,
}
impl EventQueueEvent {
pub const MAX_MESSAGE_EVENTS: usize = 1_024;
pub const fn new() -> Result<Self, Error> {
Ok(Self {
message_events: Vec::new(),
sequence: 0,
})
}
pub fn deserialize(&mut self, map: OSDMap) -> Result<(), Error> {
self.sequence = map.get("id").unwrap_or_default().as_integer()?;
let Some(OSD::Array(events)) = map.get("events") else {
return Err(Error::Parse {
position: 0,
context: "EventQueueEvent has no event array",
});
};
if events.len() > Self::MAX_MESSAGE_EVENTS {
return Err(Error::Argument);
}
let mut decoded = Vec::with_capacity(events.len());
for event in events {
let OSD::Map(mut event) = event else {
return Err(Error::Parse {
position: 0,
context: "EventQueueEvent entry is not a map",
});
};
let key = event.remove("message").unwrap_or_default().as_string()?;
let body = match event.remove("body") {
Some(OSD::Map(body)) => OSDMap::new_with_dictionary(body)?,
_ => OSDMap::new_with_constructor()?,
};
let event_message = crate::message_decoder::decode_event(&key, &body)?;
decoded.push(EventQueueEventQueueEvent {
event_message,
message_key: key,
});
}
self.message_events = decoded;
Ok(())
}
pub fn serialize(&self) -> Result<OSDMap, Error> {
if self.message_events.len() > Self::MAX_MESSAGE_EVENTS {
return Err(Error::Argument);
}
let mut events = Vec::with_capacity(self.message_events.len());
for event in &self.message_events {
let body = event
.event_message
.as_ref()
.map_or_else(OSDMap::new_with_constructor, |message| message.serialize())?;
events.push(OSD::Map(HashMap::from([
("body".to_owned(), OSD::Map(body.snapshot())),
("message".to_owned(), OSD::String(event.message_key.clone())),
])));
}
OSDMap::new_with_dictionary(HashMap::from([
("events".to_owned(), OSD::Array(events)),
("id".to_owned(), OSD::Integer(self.sequence)),
]))
}
}
pub struct EventQueueGetMessage {
pub messages: EventMessageBlock,
}
impl EventQueueGetMessage {
pub const fn new() -> Result<Self, Error> {
Ok(Self {
messages: EventMessageBlock::None,
})
}
pub fn deserialize(&mut self, map: OSDMap) -> Result<(), Error> {
self.messages.deserialize(map)
}
pub fn serialize(&self) -> Result<OSDMap, Error> {
self.messages.serialize()
}
}
impl IMessage for EventQueueGetMessage {
fn deserialize(&mut self, map: OSDMap) -> Result<(), Error> {
Self::deserialize(self, map)
}
fn serialize(&self) -> Result<OSDMap, Error> {
Self::serialize(self)
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::GridClient;
use libremetaverse_types::compat::{HttpMessageHandler, HttpRequest, HttpResponse};
use std::collections::BTreeMap;
use std::sync::atomic::AtomicU64;
use std::sync::mpsc;
use std::time::Instant;
fn lock<T>(value: &Mutex<T>) -> std::sync::MutexGuard<'_, T> {
value
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
}
fn response(status_code: u16, body: Vec<u8>) -> HttpResponse {
HttpResponse {
status_code,
headers: BTreeMap::new(),
content_type: Some(LLSD_XML.to_owned()),
body,
}
}
fn encoded_response(id: i32, names: &[&str]) -> Vec<u8> {
let events = names
.iter()
.map(|name| {
OSD::Map(HashMap::from([
("message".to_owned(), OSD::String((*name).to_owned())),
(
"body".to_owned(),
OSD::Map(HashMap::from([(
"ordinal".to_owned(),
OSD::Integer(if *name == "first" { 1 } else { 2 }),
)])),
),
]))
})
.collect();
OSDParser::serialize_llsd_xml_bytes(OSD::Map(HashMap::from([
("id".to_owned(), OSD::Integer(id)),
("events".to_owned(), OSD::Array(events)),
])))
.expect("fixture")
}
fn request_ack(request: &HttpRequest) -> (i32, bool) {
let OSD::Map(map) =
OSDParser::deserialize_llsd_xml_with_bytes(request.body.clone()).expect("request LLSD")
else {
panic!("request map")
};
(
map.get("ack")
.unwrap_or(&OSD::Undefined)
.as_integer()
.unwrap(),
map.get("done")
.unwrap_or(&OSD::Undefined)
.as_boolean()
.unwrap(),
)
}
fn simulator_with_handler(handler: HttpMessageHandler) -> Simulator {
let mut client = GridClient::new().expect("client");
client.set_http_caps_client(HttpCapsClient::new(handler).expect("HTTP client"));
let simulator = Simulator::new(
client,
"127.0.0.1:13000".parse().expect("endpoint"),
1,
None,
None,
)
.expect("simulator");
simulator.native_set_connected_for_tests(true);
simulator
}
#[test]
fn polling_preserves_ack_shape_event_order_and_prompt_cancellation() {
let requests = Arc::new(Mutex::new(Vec::<HttpRequest>::new()));
let request_log = Arc::clone(&requests);
let calls = Arc::new(AtomicU64::new(0));
let call_count = Arc::clone(&calls);
let handler = HttpMessageHandler::new(move |request, cancellation| {
let request_log = Arc::clone(&request_log);
let call = call_count.fetch_add(1, Ordering::AcqRel);
async move {
lock(&request_log).push(request);
if call == 0 {
response(200, encoded_response(7, &["first", "unknown-second"]))
} else {
cancellation.cancelled().await;
response(499, Vec::new())
}
}
});
let simulator = simulator_with_handler(handler);
let mut queue = EventQueueClient::new(
Uri("https://caps.example.test/event-queue?token=secret".to_owned()),
simulator,
)
.expect("queue");
let connected = Arc::new(AtomicU64::new(0));
let connected_count = Arc::clone(&connected);
queue.on_connected = Some(EventQueueClientConnectedCallback::from_handler(move || {
connected_count.fetch_add(1, Ordering::AcqRel);
}));
let (event_sender, event_receiver) = mpsc::channel();
queue.on_event = Some(EventQueueClientEventCallback::from_handler(
move |name, body| {
event_sender
.send((
name,
body.get("ordinal")
.unwrap_or_default()
.as_integer()
.unwrap(),
))
.expect("event receiver");
},
));
queue.start().expect("start");
assert_eq!(
event_receiver.recv_timeout(Duration::from_secs(2)).unwrap(),
("first".to_owned(), 1)
);
assert_eq!(
event_receiver.recv_timeout(Duration::from_secs(2)).unwrap(),
("unknown-second".to_owned(), 2)
);
let deadline = Instant::now() + Duration::from_secs(2);
while lock(&requests).len() < 2 && Instant::now() < deadline {
thread::sleep(Duration::from_millis(5));
}
let started = Instant::now();
queue.stop(true).expect("stop");
assert!(started.elapsed() < Duration::from_millis(500));
assert!(!queue.running());
assert_eq!(connected.load(Ordering::Acquire), 1);
let requests = lock(&requests);
assert_eq!(request_ack(&requests[0]), (0, false));
assert_eq!(request_ack(&requests[1]), (7, false));
}
#[test]
fn graceful_stop_posts_done_with_last_ack() {
let requests = Arc::new(Mutex::new(Vec::<HttpRequest>::new()));
let request_log = Arc::clone(&requests);
let calls = Arc::new(AtomicU64::new(0));
let call_count = Arc::clone(&calls);
let handler = HttpMessageHandler::new(move |request, cancellation| {
let request_log = Arc::clone(&request_log);
let call = call_count.fetch_add(1, Ordering::AcqRel);
async move {
lock(&request_log).push(request);
if call == 0 {
response(200, encoded_response(19, &[]))
} else if call == 1 {
cancellation.cancelled().await;
response(499, Vec::new())
} else {
response(200, encoded_response(19, &[]))
}
}
});
let simulator = simulator_with_handler(handler);
let queue = EventQueueClient::new(
Uri("https://caps.example.test/event-queue".to_owned()),
simulator,
)
.expect("queue");
queue.start().expect("start");
let deadline = Instant::now() + Duration::from_secs(2);
while lock(&requests).len() < 2 && Instant::now() < deadline {
thread::sleep(Duration::from_millis(5));
}
queue.stop(false).expect("graceful stop");
let requests = lock(&requests);
assert_eq!(requests.len(), 3);
assert_eq!(request_ack(&requests[2]), (19, true));
}
#[test]
fn transient_status_retries_and_terminal_status_stops() {
let calls = Arc::new(AtomicU64::new(0));
let call_count = Arc::clone(&calls);
let (sender, receiver) = mpsc::channel();
let handler = HttpMessageHandler::new(move |_request, _cancellation| {
let call = call_count.fetch_add(1, Ordering::AcqRel);
let sender = sender.clone();
async move {
let _ = sender.send(call);
match call {
0 => response(502, Vec::new()),
1 => response(200, encoded_response(3, &[])),
_ => response(404, Vec::new()),
}
}
});
let simulator = simulator_with_handler(handler);
let policy = EventQueuePolicy {
initial_retry_delay: Duration::from_millis(5),
maximum_retry_delay: Duration::from_millis(10),
..EventQueuePolicy::default()
};
let queue = EventQueueClient::with_policy(
Uri("https://caps.example.test/event-queue".to_owned()),
simulator,
policy,
)
.expect("queue");
queue.start().expect("start");
assert_eq!(receiver.recv_timeout(Duration::from_secs(2)).unwrap(), 0);
assert_eq!(receiver.recv_timeout(Duration::from_secs(2)).unwrap(), 1);
assert_eq!(receiver.recv_timeout(Duration::from_secs(2)).unwrap(), 2);
let deadline = Instant::now() + Duration::from_secs(2);
while queue.running() && Instant::now() < deadline {
thread::sleep(Duration::from_millis(5));
}
assert!(!queue.running());
queue.stop(true).expect("observe task");
}
#[test]
fn parser_rejects_excess_fanout_and_non_llsd_content() {
let too_many = vec![OSD::Map(HashMap::new()); 3];
let body = OSDParser::serialize_llsd_xml_bytes(OSD::Map(HashMap::from([
("id".to_owned(), OSD::Integer(1)),
("events".to_owned(), OSD::Array(too_many)),
])))
.unwrap();
let policy = EventQueuePolicy {
max_events_per_response: 2,
..EventQueuePolicy::default()
};
assert!(matches!(
parse_response(Some(LLSD_XML), body, &policy),
Err(Error::Argument)
));
assert!(
parse_response(
Some("text/html"),
b"<html>proxy</html>".to_vec(),
&EventQueuePolicy::default(),
)
.is_err()
);
}
#[test]
fn queue_message_blocks_round_trip_ack_and_known_events() {
let ack = EventQueueAck {
ack_id: 42,
done: true,
};
let mut get = EventQueueGetMessage::new().unwrap();
get.deserialize(ack.serialize().unwrap()).unwrap();
let EventMessageBlock::Ack(decoded) = get.messages else {
panic!("ack block")
};
assert_eq!((decoded.ack_id, decoded.done), (42, true));
let fixture = OSDMap::new_with_dictionary(HashMap::from([
("id".to_owned(), OSD::Integer(9)),
(
"events".to_owned(),
OSD::Array(vec![OSD::Map(HashMap::from([
(
"message".to_owned(),
OSD::String("UpdateAgentLanguage".to_owned()),
),
(
"body".to_owned(),
OSD::Map(HashMap::from([
("language".to_owned(), OSD::String("en".to_owned())),
("language_is_public".to_owned(), OSD::Boolean(true)),
])),
),
]))]),
),
]))
.unwrap();
let mut event = EventQueueEvent::new().unwrap();
event.deserialize(fixture).unwrap();
assert_eq!(event.sequence, 9);
assert_eq!(event.message_events.len(), 1);
assert!(event.message_events[0].event_message.is_some());
assert_eq!(event.serialize().unwrap().get("id"), Some(OSD::Integer(9)));
}
}