Implement capability event queue processing (#55)
This commit is contained in:
@@ -49,14 +49,6 @@ impl Clone for Packet {
|
||||
}
|
||||
}
|
||||
|
||||
impl Clone for Caps {
|
||||
fn clone(&self) -> Self {
|
||||
Self {
|
||||
simulator: self.simulator.clone(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn mutex<T>(value: &Mutex<T>) -> std::sync::MutexGuard<'_, T> {
|
||||
value
|
||||
.lock()
|
||||
@@ -236,6 +228,23 @@ pub struct SimConnectedEventArgs {
|
||||
simulator: Simulator,
|
||||
}
|
||||
|
||||
/// Event arguments raised when a simulator's `EventQueueGet` poller is live.
|
||||
#[derive(Clone)]
|
||||
pub struct EventQueueRunningEventArgs {
|
||||
simulator: Simulator,
|
||||
}
|
||||
|
||||
impl EventQueueRunningEventArgs {
|
||||
pub const fn new(simulator: Simulator) -> Result<Self, Error> {
|
||||
Ok(Self { simulator })
|
||||
}
|
||||
|
||||
#[must_use]
|
||||
pub fn simulator(&self) -> Simulator {
|
||||
self.simulator.clone()
|
||||
}
|
||||
}
|
||||
|
||||
impl SimConnectedEventArgs {
|
||||
pub fn new(simulator: Simulator) -> Result<Self, Error> {
|
||||
Ok(Self { simulator })
|
||||
@@ -849,6 +858,9 @@ pub struct SimulatorData {
|
||||
disconnect_candidate: AtomicBool,
|
||||
circuit_code: AtomicU32,
|
||||
seed_caps: RwLock<Option<Uri>>,
|
||||
caps_state: RwLock<Option<Arc<crate::caps::CapsInner>>>,
|
||||
event_queue_generation: Mutex<u64>,
|
||||
event_queue_wait: Condvar,
|
||||
manager: Mutex<Weak<NetworkManagerInner>>,
|
||||
transport: UDPBase,
|
||||
transport_thread: Mutex<Option<TransportThread>>,
|
||||
@@ -858,6 +870,14 @@ pub struct SimulatorData {
|
||||
|
||||
impl Drop for SimulatorData {
|
||||
fn drop(&mut self) {
|
||||
if let Some(caps) = self
|
||||
.caps_state
|
||||
.get_mut()
|
||||
.unwrap_or_else(std::sync::PoisonError::into_inner)
|
||||
.take()
|
||||
{
|
||||
caps.disconnect(true);
|
||||
}
|
||||
let runtime = self
|
||||
.transport_thread
|
||||
.get_mut()
|
||||
@@ -873,7 +893,6 @@ impl Drop for SimulatorData {
|
||||
}
|
||||
|
||||
/// Reference-identity simulator handle.
|
||||
#[derive(Clone)]
|
||||
pub struct Simulator {
|
||||
/// Capability client associated with this simulator.
|
||||
///
|
||||
@@ -883,6 +902,19 @@ pub struct Simulator {
|
||||
data: Arc<SimulatorData>,
|
||||
}
|
||||
|
||||
impl Clone for Simulator {
|
||||
fn clone(&self) -> Self {
|
||||
let mut clone = self.native_clone_without_caps();
|
||||
if let Some(inner) = read(&self.data.caps_state).clone() {
|
||||
clone.caps = Some(Box::new(Caps::native_from_inner(
|
||||
inner,
|
||||
clone.native_clone_without_caps(),
|
||||
)));
|
||||
}
|
||||
clone
|
||||
}
|
||||
}
|
||||
|
||||
impl std::ops::Deref for Simulator {
|
||||
type Target = SimulatorData;
|
||||
|
||||
@@ -987,6 +1019,9 @@ impl Simulator {
|
||||
disconnect_candidate: AtomicBool::new(false),
|
||||
circuit_code: AtomicU32::new(0),
|
||||
seed_caps: RwLock::new(None),
|
||||
caps_state: RwLock::new(None),
|
||||
event_queue_generation: Mutex::new(0),
|
||||
event_queue_wait: Condvar::new(),
|
||||
manager: Mutex::new(Weak::new()),
|
||||
transport,
|
||||
transport_thread: Mutex::new(None),
|
||||
@@ -997,6 +1032,29 @@ impl Simulator {
|
||||
Ok(Self { caps: None, data })
|
||||
}
|
||||
|
||||
pub(crate) fn native_clone_without_caps(&self) -> Self {
|
||||
Self {
|
||||
caps: None,
|
||||
data: Arc::clone(&self.data),
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) fn native_data_weak(&self) -> Weak<SimulatorData> {
|
||||
Arc::downgrade(&self.data)
|
||||
}
|
||||
|
||||
pub(crate) fn native_data_arc(&self) -> Arc<SimulatorData> {
|
||||
Arc::clone(&self.data)
|
||||
}
|
||||
|
||||
pub(crate) fn native_from_weak(data: &Weak<SimulatorData>) -> Option<Self> {
|
||||
data.upgrade().map(|data| Self { caps: None, data })
|
||||
}
|
||||
|
||||
pub(crate) fn native_from_data(data: Arc<SimulatorData>) -> Self {
|
||||
Self { caps: None, data }
|
||||
}
|
||||
|
||||
fn attach_manager(&self, manager: &Arc<NetworkManagerInner>, circuit_code: u32) {
|
||||
*mutex(&self.manager) = Arc::downgrade(manager);
|
||||
self.circuit_code.store(circuit_code, Ordering::Release);
|
||||
@@ -1084,6 +1142,10 @@ impl Simulator {
|
||||
}
|
||||
|
||||
pub fn native_disconnect(&self, send_close_circuit: bool) -> Result<(), Error> {
|
||||
if let Some(caps) = write(&self.caps_state).take() {
|
||||
caps.disconnect(true);
|
||||
}
|
||||
*write(&self.seed_caps) = None;
|
||||
self.disconnect_candidate.store(false, Ordering::Release);
|
||||
let was_connected = self.connected.swap(false, Ordering::AcqRel);
|
||||
self.handshake_complete.store(false, Ordering::Release);
|
||||
@@ -1198,14 +1260,114 @@ impl Simulator {
|
||||
seed_caps: Option<Uri>,
|
||||
changed_sim: Option<bool>,
|
||||
) -> Result<(), Error> {
|
||||
let mut current = write(&self.seed_caps);
|
||||
if !changed_sim.unwrap_or(false) && *current == seed_caps {
|
||||
if !changed_sim.unwrap_or(false) && *read(&self.seed_caps) == seed_caps {
|
||||
return Ok(());
|
||||
}
|
||||
*current = seed_caps;
|
||||
if let Some(previous) = write(&self.caps_state).take() {
|
||||
previous.disconnect(true);
|
||||
}
|
||||
write(&self.seed_caps).clone_from(&seed_caps);
|
||||
if let Some(seed_caps) = seed_caps {
|
||||
let caps = Caps::native_create(self, seed_caps)?;
|
||||
*write(&self.caps_state) = Some(caps);
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub fn native_is_event_queue_running(
|
||||
&self,
|
||||
block_until_running: Option<bool>,
|
||||
) -> Result<bool, Error> {
|
||||
if self.caps_running() {
|
||||
return Ok(true);
|
||||
}
|
||||
if !block_until_running.unwrap_or(false) {
|
||||
return Ok(false);
|
||||
}
|
||||
let mut generation = mutex(&self.event_queue_generation);
|
||||
for delay in [1_u64, 2, 4, 8] {
|
||||
if self.caps_running() || self.current_sim_caps_running() {
|
||||
return Ok(true);
|
||||
}
|
||||
self.start_current_event_queue();
|
||||
let observed = *generation;
|
||||
let (next, _) = self
|
||||
.event_queue_wait
|
||||
.wait_timeout_while(generation, Duration::from_secs(delay), |value| {
|
||||
*value == observed
|
||||
&& !self.client.cancellation_token().is_cancellation_requested()
|
||||
})
|
||||
.unwrap_or_else(std::sync::PoisonError::into_inner);
|
||||
generation = next;
|
||||
if *generation != observed {
|
||||
break;
|
||||
}
|
||||
}
|
||||
Ok(self.caps_running() || self.current_sim_caps_running())
|
||||
}
|
||||
|
||||
fn caps_running(&self) -> bool {
|
||||
read(&self.caps_state).as_ref().is_some_and(|inner| {
|
||||
Caps::native_from_inner(Arc::clone(inner), self.native_clone_without_caps())
|
||||
.is_event_queue_running()
|
||||
})
|
||||
}
|
||||
|
||||
fn current_sim_caps_running(&self) -> bool {
|
||||
mutex(&self.manager)
|
||||
.upgrade()
|
||||
.and_then(|manager| read(&manager.current_sim).clone())
|
||||
.is_some_and(|simulator| simulator.caps_running())
|
||||
}
|
||||
|
||||
fn start_current_event_queue(&self) {
|
||||
let simulator = mutex(&self.manager)
|
||||
.upgrade()
|
||||
.and_then(|manager| read(&manager.current_sim).clone())
|
||||
.unwrap_or_else(|| self.native_clone_without_caps());
|
||||
if let Some(inner) = read(&simulator.caps_state).clone() {
|
||||
let caps = Caps::native_from_inner(inner, simulator.native_clone_without_caps());
|
||||
if let Some(queue) = caps.event_queue() {
|
||||
let _ = queue.start();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) fn native_raise_event_queue_running(&self) {
|
||||
let mut generation = mutex(&self.event_queue_generation);
|
||||
*generation = generation.wrapping_add(1);
|
||||
drop(generation);
|
||||
self.event_queue_wait.notify_all();
|
||||
if let Some(manager) = mutex(&self.manager).upgrade() {
|
||||
let args = EventQueueRunningEventArgs {
|
||||
simulator: self.clone(),
|
||||
};
|
||||
manager
|
||||
.events
|
||||
.event_queue_running
|
||||
.emit_with(|| args.clone());
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) fn native_dispatch_caps_event(&self, event_name: &str, message: &dyn IMessage) {
|
||||
if let Some(manager) = mutex(&self.manager).upgrade() {
|
||||
manager
|
||||
.caps_events
|
||||
.invoke_raise_event(event_name, message, self.clone());
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) fn native_enqueue_caps_packet(&self, packet: Packet) -> Result<(), Error> {
|
||||
let manager = mutex(&self.manager)
|
||||
.upgrade()
|
||||
.ok_or(Error::InvalidOperation)?;
|
||||
manager.enqueue_incoming(NetworkManagerIncomingPacket {
|
||||
packet: Some(packet),
|
||||
simulator: Some(self.clone()),
|
||||
raw_data: None,
|
||||
})
|
||||
}
|
||||
|
||||
#[must_use]
|
||||
pub fn seed_capability(&self) -> Option<Uri> {
|
||||
read(&self.seed_caps).clone()
|
||||
@@ -1262,6 +1424,11 @@ impl Simulator {
|
||||
self.connected.load(Ordering::Acquire)
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
pub(crate) fn native_set_connected_for_tests(&self, connected: bool) {
|
||||
self.connected.store(connected, Ordering::Release);
|
||||
}
|
||||
|
||||
#[must_use]
|
||||
pub fn native_handshake_complete(&self) -> bool {
|
||||
self.handshake_complete.load(Ordering::Acquire)
|
||||
@@ -1397,6 +1564,7 @@ struct NetworkWorkers {
|
||||
#[derive(Default)]
|
||||
struct NetworkEvents {
|
||||
disconnected: EventRegistry<DisconnectedEventArgs>,
|
||||
event_queue_running: EventRegistry<EventQueueRunningEventArgs>,
|
||||
generic_streaming_message: EventRegistry<GenericStreamingMessageEventArgs>,
|
||||
packet_sent: EventRegistry<PacketSentEventArgs>,
|
||||
sim_changed: EventRegistry<SimChangedEventArgs>,
|
||||
@@ -1944,6 +2112,13 @@ impl NetworkManager {
|
||||
self.inner.events.disconnected.subscribe(handler)
|
||||
}
|
||||
|
||||
pub fn native_subscribe_event_queue_running(
|
||||
&self,
|
||||
handler: EventHandler<EventQueueRunningEventArgs>,
|
||||
) -> Subscription {
|
||||
self.inner.events.event_queue_running.subscribe(handler)
|
||||
}
|
||||
|
||||
pub fn native_subscribe_generic_streaming_message(
|
||||
&self,
|
||||
handler: EventHandler<GenericStreamingMessageEventArgs>,
|
||||
@@ -2254,6 +2429,8 @@ impl NetworkManager {
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use libremetaverse_structured_data::{OSD, OSDMap};
|
||||
use std::sync::mpsc;
|
||||
|
||||
#[test]
|
||||
fn manager_workers_and_callback_registries_do_not_retain_the_client() {
|
||||
@@ -2269,4 +2446,82 @@ mod tests {
|
||||
|
||||
assert!(probe.is_released());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn caps_packet_fallback_uses_the_bounded_udp_dispatch_pipeline() {
|
||||
let client = GridClient::new().expect("client");
|
||||
let manager = NetworkManager::new(client.clone()).expect("network manager");
|
||||
manager.inner.ensure_workers().expect("workers");
|
||||
let simulator =
|
||||
Simulator::new(client, "127.0.0.1:14000".parse().unwrap(), 10, None, None).unwrap();
|
||||
simulator.attach_manager(&manager.inner, 0);
|
||||
manager.simulators.push_if_missing(simulator.clone());
|
||||
let (sender, receiver) = mpsc::channel();
|
||||
let _subscription = manager.subscribe_packet(
|
||||
PacketType::PacketAck,
|
||||
Arc::new(move |args| {
|
||||
sender
|
||||
.send((args.packet().type_, args.simulator().handle))
|
||||
.unwrap();
|
||||
}),
|
||||
false,
|
||||
);
|
||||
let body = OSDMap::new_with_constructor().unwrap();
|
||||
body.add_with_string_osd("Packets".to_owned(), OSD::Array(Vec::new()))
|
||||
.unwrap();
|
||||
let packet = Packet::build_packet_with_string_osd_map("PacketAck".to_owned(), body)
|
||||
.unwrap()
|
||||
.expect("packet fallback");
|
||||
simulator.native_enqueue_caps_packet(packet).unwrap();
|
||||
assert_eq!(
|
||||
receiver.recv_timeout(Duration::from_secs(2)).unwrap(),
|
||||
(PacketType::PacketAck, 10)
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn event_queue_running_and_typed_caps_events_use_manager_registries() {
|
||||
let client = GridClient::new().expect("client");
|
||||
let manager = NetworkManager::new(client.clone()).expect("network manager");
|
||||
let simulator =
|
||||
Simulator::new(client, "127.0.0.1:14001".parse().unwrap(), 11, None, None).unwrap();
|
||||
simulator.attach_manager(&manager.inner, 0);
|
||||
|
||||
let (running_sender, running_receiver) = mpsc::channel();
|
||||
let _running = manager.native_subscribe_event_queue_running(Arc::new(move |args| {
|
||||
running_sender.send(args.simulator().handle).unwrap();
|
||||
}));
|
||||
simulator.native_raise_event_queue_running();
|
||||
assert_eq!(
|
||||
running_receiver
|
||||
.recv_timeout(Duration::from_secs(1))
|
||||
.unwrap(),
|
||||
11
|
||||
);
|
||||
|
||||
let (caps_sender, caps_receiver) = mpsc::channel();
|
||||
let _caps = manager.subscribe_caps(
|
||||
"UpdateAgentLanguage".to_owned(),
|
||||
CapsEventQueueCallback::from_handler(move |name, message, simulator| {
|
||||
let typed = (message as &dyn std::any::Any)
|
||||
.downcast_ref::<crate::messages::linden::UpdateAgentLanguageMessage>();
|
||||
caps_sender
|
||||
.send((name.to_owned(), typed.is_some(), simulator.handle))
|
||||
.unwrap();
|
||||
}),
|
||||
);
|
||||
let body = OSDMap::new_with_dictionary(HashMap::from([
|
||||
("language".to_owned(), OSD::String("en".to_owned())),
|
||||
("language_is_public".to_owned(), OSD::Boolean(true)),
|
||||
]))
|
||||
.unwrap();
|
||||
let message = crate::message_decoder::decode_event("UpdateAgentLanguage", &body)
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
simulator.native_dispatch_caps_event("UpdateAgentLanguage", message.as_ref());
|
||||
assert_eq!(
|
||||
caps_receiver.recv_timeout(Duration::from_secs(1)).unwrap(),
|
||||
("UpdateAgentLanguage".to_owned(), true, 11)
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user