From 758f82dbdddafa55a03d83ec16b36bb708ea3219 Mon Sep 17 00:00:00 2001 From: xmsama <36942744+xmsama@users.noreply.github.com> Date: Fri, 25 Sep 2026 05:55:45 +0800 Subject: [PATCH 1/2] feat(server): add opt-in latest-state output scheduling and stop barriers --- .../src/device/device_handle.rs | 71 ++- .../buttplug_server/src/device/device_task.rs | 6 +- .../src/device/latest_device_task.rs | 451 ++++++++++++++++++ crates/buttplug_server/src/device/mod.rs | 1 + crates/buttplug_server/src/device/protocol.rs | 17 + 5 files changed, 526 insertions(+), 20 deletions(-) create mode 100644 crates/buttplug_server/src/device/latest_device_task.rs diff --git a/crates/buttplug_server/src/device/device_handle.rs b/crates/buttplug_server/src/device/device_handle.rs index 693ba12ed..8b82419bf 100644 --- a/crates/buttplug_server/src/device/device_handle.rs +++ b/crates/buttplug_server/src/device/device_handle.rs @@ -309,20 +309,14 @@ impl DeviceHandle { // --- Private command handling methods --- - /// Run an output command through last-command deduplication and observation - /// emission, returning the protocol handler's hardware commands. Returns None - /// when the command equals the feature's last command and generates no work. - /// Shared by the normal output path and the stop path so both keep identical - /// dedupe-map and observation behaviour. - fn output_cmd_hardware_commands( - &self, - msg: &CheckedOutputCmdV4, - ) -> Option, ButtplugError>> { + /// Record an accepted output command and emit its observation. Returns false + /// when the command is identical to the feature's last command. + fn record_output_cmd(&self, msg: &CheckedOutputCmdV4) -> bool { if let Some(last_msg) = self.last_output_command.get(&msg.feature_id()) && *last_msg == *msg { trace!("No commands generated for incoming device packet, skipping and returning success."); - return None; + return false; } self .last_output_command @@ -339,14 +333,35 @@ impl DeviceHandle { }); } - Some(self.handler.handle_output_cmd(msg).map_err(|e| e.into())) + true } fn handle_outputcmd_v4(&self, msg: &CheckedOutputCmdV4) -> ButtplugServerResultFuture { - match self.output_cmd_hardware_commands(msg) { - None => future::ready(Ok(message::OkV0::default().into())).boxed(), - Some(Ok(commands)) => self.handle_hardware_commands(commands), - Some(Err(err)) => future::ready(Err(err)).boxed(), + if !self.record_output_cmd(msg) { + return future::ready(Ok(message::OkV0::default().into())).boxed(); + } + + match self.handler.handle_output_cmd(msg) { + Ok(commands) + if self.handler.use_latest_output_scheduler() + && msg.output_command().value() == 0 + && !matches!( + msg.output_command(), + message::OutputCommand::HwPositionWithDuration(_) | message::OutputCommand::Position(_) + ) => + { + let sender = self.internal_hw_msg_sender.clone(); + async move { + // Slider zero is an urgent barrier, not another replaceable sample. + let (message, ack) = DeviceTaskMessage::acknowledged(commands); + drop(ack); + let _ = sender.send(message).await; + Ok(message::OkV0::default().into()) + } + .boxed() + } + Ok(commands) => self.handle_hardware_commands(commands), + Err(err) => future::ready(Err(err.into())).boxed(), } } @@ -370,10 +385,28 @@ impl DeviceHandle { let mut hardware_commands: Vec = Vec::new(); if msg.outputs() { for stop_msg in self.stop_commands.iter() { - if let ButtplugDeviceCommandMessageUnionV4::OutputCmd(checked) = stop_msg - && let Some(Ok(cmds)) = self.output_cmd_hardware_commands(checked) - { - hardware_commands.extend(cmds); + if let ButtplugDeviceCommandMessageUnionV4::OutputCmd(checked) = stop_msg { + // Only opted-in state protocols force every explicit stop through deduplication. + if self.handler.use_latest_output_scheduler() { + // The source value may already have been recorded as zero while its + // hardware write is still pending. Force this one safety stop now. + self + .last_output_command + .insert(checked.feature_id(), checked.clone()); + if let Some(sender) = &self.output_observation_sender { + let _ = sender.send(OutputObservation { + device_index: self.definition.index(), + feature_index: checked.feature_index(), + output_type: checked.output_command().as_output_type().to_string(), + value: checked.output_command().value() as f64, + }); + } + } else if !self.record_output_cmd(checked) { + continue; + } + if let Ok(cmds) = self.handler.handle_stop_output_cmd(checked) { + hardware_commands.extend(cmds); + } } } } diff --git a/crates/buttplug_server/src/device/device_task.rs b/crates/buttplug_server/src/device/device_task.rs index 6723a2edf..89b03517a 100644 --- a/crates/buttplug_server/src/device/device_task.rs +++ b/crates/buttplug_server/src/device/device_task.rs @@ -84,10 +84,14 @@ pub struct DeviceTaskConfig { /// Run the device communication task under its device owner's task group. pub async fn run_owned_device_task( hardware: Arc, - _handler: Arc, + handler: Arc, config: DeviceTaskConfig, mut command_receiver: Receiver, ) { + if handler.use_latest_output_scheduler() && !config.requires_keepalive { + super::latest_device_task::run_latest_device_task(hardware, config, command_receiver).await; + return; + } run_device_task(hardware, config, &mut command_receiver).await; } diff --git a/crates/buttplug_server/src/device/latest_device_task.rs b/crates/buttplug_server/src/device/latest_device_task.rs new file mode 100644 index 000000000..fb7b0f39b --- /dev/null +++ b/crates/buttplug_server/src/device/latest_device_task.rs @@ -0,0 +1,451 @@ +// Buttplug Rust Source Code File - See https://buttplug.io for more info. +// Licensed under the BSD 3-Clause license. See LICENSE in the project root. + +//! Opt-in state scheduler: receive/coalesce while a single hardware write is in +//! flight. A stop is a barrier: newer motion cannot overwrite an unsent stop. +//! Never cancel a slow write and then submit another on the same connection; +//! after an error/timeout discard pending motion and disconnect instead. + +use std::{collections::VecDeque, sync::Arc, time::Duration}; + +use buttplug_core::{errors::ButtplugDeviceError, util::async_manager}; +use futures::{FutureExt, future::BoxFuture}; +use tokio::{ + select, + sync::{mpsc::Receiver, oneshot}, + time::Instant, +}; + +use super::{ + device_task::{DeviceTaskConfig, DeviceTaskMessage}, + hardware::{Hardware, HardwareCommand, HardwareEvent}, +}; + +const WRITE_TIMEOUT: Duration = Duration::from_secs(2); +const SLOW_WRITE: Duration = Duration::from_millis(200); + +#[derive(Default)] +struct Batch { + commands: VecDeque, + acks: Vec>, + queued_at: Option, +} + +impl Batch { + fn merge(&mut self, commands: impl IntoIterator) { + for command in commands { + self.queued_at.get_or_insert_with(Instant::now); + self.commands.retain(|old| !command.overlaps(old)); + self.commands.push_back(command); + } + } +} + +#[derive(Default)] +struct Pending { + normal: Batch, + urgent: Batch, +} + +impl Pending { + fn receive(&mut self, message: DeviceTaskMessage) { + if let Some(ack) = message.write_ack { + // Pending state before the stop can be replaced. State received AFTER + // it remains in normal and must not erase the stop barrier. + self.urgent.merge(self.normal.commands.drain(..)); + self.normal.queued_at = None; + self.urgent.merge(message.commands); + self.urgent.acks.push(ack); + } else { + self.normal.merge(message.commands); + } + } + + fn has_urgent(&self) -> bool { + !self.urgent.acks.is_empty() + } + + fn is_empty(&self) -> bool { + !self.has_urgent() && self.normal.commands.is_empty() + } + + fn take_next(&mut self) -> Batch { + if self.has_urgent() { + std::mem::take(&mut self.urgent) + } else { + std::mem::take(&mut self.normal) + } + } +} + +async fn write_batch( + hardware: Arc, + mut batch: Batch, + timeout: Duration, +) -> Result<(), ButtplugDeviceError> { + let queued_ms = batch + .queued_at + .map(|t| t.elapsed().as_millis()) + .unwrap_or(0); + debug!( + "Latest-state output dispatch: queued_ms={queued_ms}, packets={}, urgent={}", + batch.commands.len(), + !batch.acks.is_empty() + ); + while let Some(command) = batch.commands.pop_front() { + let started = Instant::now(); + let result = select! { + result = hardware.parse_message(&command) => result, + _ = async_manager::sleep(timeout) => { + Err(ButtplugDeviceError::DeviceCommunicationError( + format!("Latest-state output write timed out after {} ms", timeout.as_millis()))) + } + }; + let elapsed = started.elapsed(); + if elapsed >= SLOW_WRITE || result.is_err() { + warn!( + "Latest-state output write: elapsed_ms={}, success={}", + elapsed.as_millis(), + result.is_ok() + ); + } else { + debug!( + "Latest-state output write: elapsed_ms={}", + elapsed.as_millis() + ); + } + result?; + } + for ack in batch.acks { + let _ = ack.send(()); + } + Ok(()) +} + +pub(super) async fn run_latest_device_task( + hardware: Arc, + config: DeviceTaskConfig, + command_receiver: Receiver, +) { + run( + hardware, + config.message_gap.unwrap_or(Duration::ZERO), + command_receiver, + WRITE_TIMEOUT, + ) + .await; +} + +async fn run( + hardware: Arc, + gap: Duration, + mut receiver: Receiver, + write_timeout: Duration, +) { + let mut events = hardware.event_stream(); + let mut pending = Pending::default(); + let mut flight: Option>> = None; + let mut next_write = Instant::now(); + let mut closed = false; + info!( + "Latest-state scheduler enabled: gap_ms={}, timeout_ms={}", + gap.as_millis(), + write_timeout.as_millis() + ); + + loop { + if closed && flight.is_none() && pending.is_empty() { + return; + } + let can_send = flight.is_none() && !pending.is_empty(); + let delay = if pending.has_urgent() { + Duration::ZERO + } else { + next_write.saturating_duration_since(Instant::now()) + }; + select! { + biased; + event = events.recv() => { + match event { + Ok(HardwareEvent::Disconnected(_)) | Err(tokio::sync::broadcast::error::RecvError::Closed) => return, + _ => {} // Notifications/lag do not imply disconnection. + } + } + result = async { flight.as_mut().unwrap().await }, if flight.is_some() => { + flight = None; + if let Err(error) = result { + error!("Latest-state output failed; discarding pending motion and disconnecting: {error}"); + // Dropping the write future is not proof of OS-level cancellation. + // Do not retry or replay motion on this connection. + select! { + _ = hardware.disconnect() => {}, + _ = async_manager::sleep(Duration::from_secs(2)) => {}, + } + return; + } + } + // Timer precedes input so a continuous producer cannot starve output. + _ = async_manager::sleep(delay), if can_send => { + // Take a bounded snapshot of already queued messages before dispatch. + // Never drain indefinitely under an unbounded continuous producer. + let count = receiver.len(); + for _ in 0..count { + if let Ok(message) = receiver.try_recv() { + pending.receive(message); + } + } + let batch = pending.take_next(); + next_write = Instant::now() + gap; + flight = Some(write_batch(hardware.clone(), batch, write_timeout).boxed()); + } + message = receiver.recv(), if !closed => { + match message { + Some(message) => { + pending.receive(message); + } + None => { + closed = true; + // No owner remains to authorize replay of pending motion. + pending = Pending::default(); + }, + } + } + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::device::hardware::{ + HardwareInternal, + HardwareReadCmd, + HardwareReading, + HardwareSubscribeCmd, + HardwareUnsubscribeCmd, + HardwareWriteCmd, + }; + use buttplug_server_device_config::Endpoint; + use tokio::sync::{Semaphore, broadcast, mpsc}; + use uuid::Uuid; + + struct SlowHardware { + events: broadcast::Sender, + starts: mpsc::UnboundedSender, + permits: Arc, + fail: bool, + } + + impl HardwareInternal for SlowHardware { + fn event_stream(&self) -> broadcast::Receiver { + self.events.subscribe() + } + fn disconnect(&self) -> BoxFuture<'static, Result<(), ButtplugDeviceError>> { + let events = self.events.clone(); + async move { + let _ = events.send(HardwareEvent::Disconnected("test".into())); + Ok(()) + } + .boxed() + } + fn write_value( + &self, + msg: &HardwareWriteCmd, + ) -> BoxFuture<'static, Result<(), ButtplugDeviceError>> { + let starts = self.starts.clone(); + let value = msg.data()[0]; + let permits = self.permits.clone(); + let fail = self.fail; + async move { + starts.send(value).unwrap(); + permits.acquire().await.unwrap().forget(); + if fail { + Err(ButtplugDeviceError::DeviceCommunicationError( + "test failure".into(), + )) + } else { + Ok(()) + } + } + .boxed() + } + fn read_value( + &self, + msg: &HardwareReadCmd, + ) -> BoxFuture<'static, Result> { + let endpoint = msg.endpoint(); + async move { Ok(HardwareReading::new(endpoint, &[])) }.boxed() + } + fn subscribe( + &self, + _: &HardwareSubscribeCmd, + ) -> BoxFuture<'static, Result<(), ButtplugDeviceError>> { + async { Ok(()) }.boxed() + } + fn unsubscribe( + &self, + _: &HardwareUnsubscribeCmd, + ) -> BoxFuture<'static, Result<(), ButtplugDeviceError>> { + async { Ok(()) }.boxed() + } + } + + struct Harness { + sender: mpsc::Sender, + starts: mpsc::UnboundedReceiver, + permits: Arc, + events: broadcast::Sender, + task: tokio::task::JoinHandle<()>, + } + + impl Harness { + fn new(gap: Duration, timeout: Duration, fail: bool) -> Self { + let (events, _) = broadcast::channel(16); + let (starts, receiver) = mpsc::unbounded_channel(); + let permits = Arc::new(Semaphore::new(0)); + let hardware = Arc::new(Hardware::new( + "test", + "test", + &[Endpoint::Tx], + &None, + false, + Box::new(SlowHardware { + events: events.clone(), + starts, + permits: permits.clone(), + fail, + }), + )); + let (sender, commands) = mpsc::channel(32); + let task = tokio::spawn(run(hardware, gap, commands, timeout)); + Self { + sender, + starts: receiver, + permits, + events, + task, + } + } + async fn send(&self, value: u8) { + self + .sender + .send(DeviceTaskMessage::fire_and_forget(vec![command(value)])) + .await + .unwrap(); + } + async fn started(&mut self) -> u8 { + tokio::time::timeout(Duration::from_secs(1), self.starts.recv()) + .await + .unwrap() + .unwrap() + } + async fn finish(self) { + let _ = self.events.send(HardwareEvent::Disconnected("test".into())); + tokio::time::timeout(Duration::from_secs(1), self.task) + .await + .unwrap() + .unwrap(); + } + } + + fn command(value: u8) -> HardwareCommand { + HardwareWriteCmd::new(&[Uuid::nil()], Endpoint::Tx, vec![value], false).into() + } + + #[tokio::test] + async fn busy_write_keeps_only_latest_state() { + let mut h = Harness::new(Duration::from_millis(10), Duration::from_secs(2), false); + h.send(10).await; + assert_eq!(h.started().await, 10); + for value in 11..=100 { + h.send(value).await; + } + h.permits.add_permits(1); + assert_eq!(h.started().await, 100); + // No parallel writes while the last one is unresolved. + assert!( + tokio::time::timeout(Duration::from_millis(30), h.starts.recv()) + .await + .is_err() + ); + h.finish().await; + } + + #[tokio::test] + async fn stop_is_not_overwritten_by_new_motion() { + let mut h = Harness::new(Duration::from_millis(10), Duration::from_secs(2), false); + h.send(10).await; + assert_eq!(h.started().await, 10); + h.send(20).await; + let (stop, ack) = DeviceTaskMessage::acknowledged(vec![command(0)]); + h.sender.send(stop).await.unwrap(); + h.send(80).await; + h.permits.add_permits(1); + assert_eq!(h.started().await, 0); + h.permits.add_permits(1); + tokio::time::timeout(Duration::from_secs(1), ack) + .await + .unwrap() + .unwrap(); + assert_eq!(h.started().await, 80); + h.finish().await; + } + + #[tokio::test] + async fn first_command_and_stop_bypass_pacing_delay() { + let mut h = Harness::new(Duration::from_secs(5), Duration::from_secs(2), false); + h.send(20).await; + assert_eq!(h.started().await, 20); + h.permits.add_permits(1); + let (stop, _ack) = DeviceTaskMessage::acknowledged(vec![command(0)]); + h.sender.send(stop).await.unwrap(); + assert_eq!(h.started().await, 0); + h.finish().await; + } + + #[tokio::test] + async fn timeout_disconnects_without_replaying_queued_motion() { + let mut h = Harness::new(Duration::ZERO, Duration::from_millis(50), false); + let mut events = h.events.subscribe(); + h.send(20).await; + assert_eq!(h.started().await, 20); + h.send(80).await; + tokio::time::timeout(Duration::from_secs(1), &mut h.task) + .await + .unwrap() + .unwrap(); + assert!(matches!( + events.recv().await.unwrap(), + HardwareEvent::Disconnected(_) + )); + assert!(h.starts.try_recv().is_err()); + } + + #[tokio::test] + async fn failed_write_disconnects_without_replay() { + let mut h = Harness::new(Duration::ZERO, Duration::from_secs(2), true); + let mut events = h.events.subscribe(); + h.send(20).await; + assert_eq!(h.started().await, 20); + h.send(80).await; + h.permits.add_permits(1); + tokio::time::timeout(Duration::from_secs(1), &mut h.task) + .await + .unwrap() + .unwrap(); + assert!(matches!( + events.recv().await.unwrap(), + HardwareEvent::Disconnected(_) + )); + assert!(h.starts.try_recv().is_err()); + } + + #[tokio::test] + async fn disconnect_is_handled_while_write_is_blocked() { + let mut h = Harness::new(Duration::ZERO, Duration::from_secs(2), false); + h.send(20).await; + assert_eq!(h.started().await, 20); + h.send(80).await; + h.finish().await; + } +} diff --git a/crates/buttplug_server/src/device/mod.rs b/crates/buttplug_server/src/device/mod.rs index 26c48782f..4a5179c5d 100644 --- a/crates/buttplug_server/src/device/mod.rs +++ b/crates/buttplug_server/src/device/mod.rs @@ -98,6 +98,7 @@ mod device_handle; mod device_task; pub mod hardware; +mod latest_device_task; mod output_observation; pub mod protocol; pub mod protocol_impl; diff --git a/crates/buttplug_server/src/device/protocol.rs b/crates/buttplug_server/src/device/protocol.rs index e624b5706..526a7e344 100644 --- a/crates/buttplug_server/src/device/protocol.rs +++ b/crates/buttplug_server/src/device/protocol.rs @@ -283,6 +283,23 @@ pub trait ProtocolHandler: Sync + Send { } } + /// Opt into latest-state output scheduling. Only suitable for protocols whose + /// overlapping commands are replaceable state snapshots, without keepalives. + /// Other protocols retain the existing device task and stop behavior. + fn use_latest_output_scheduler(&self) -> bool { + false + } + + /// Handle an explicit device stop. The default is identical to an ordinary + /// output command; protocols with waveform conditioning can override it to + /// guarantee that safety stops are never filtered or delayed. + fn handle_stop_output_cmd( + &self, + cmd: &CheckedOutputCmdV4, + ) -> Result, ButtplugDeviceError> { + self.handle_output_cmd(cmd) + } + fn handle_output_vibrate_cmd( &self, _feature_index: u32, From f50d48a2573347b2f0daeb82524784c50d5870ab Mon Sep 17 00:00:00 2001 From: xmsama <36942744+xmsama@users.noreply.github.com> Date: Fri, 25 Sep 2026 06:02:47 +0800 Subject: [PATCH 2/2] feat(yiciyuan): correct FJB capabilities and add host rotation modes --- .../src/device/latest_device_task.rs | 305 +++++++++- crates/buttplug_server/src/device/mod.rs | 1 + .../src/device/protocol_impl/yiciyuan.rs | 538 +++++++++++++++++- .../src/device/yiciyuan_motion.rs | 499 ++++++++++++++++ .../buttplug-device-config-v5.json | 105 +++- .../buttplug-device-config-schema-v5.json | 4 + .../device-config/protocols/yiciyuan.yml | 72 ++- .../device-config/version.yaml | 2 +- .../tests/test_device_protocols.rs | 10 + .../util/device_test/client/client_v0/mod.rs | 3 +- .../util/device_test/client/client_v1/mod.rs | 3 +- .../util/device_test/client/client_v2/mod.rs | 3 +- .../util/device_test/client/client_v3/mod.rs | 3 +- .../util/device_test/client/client_v4/mod.rs | 3 +- .../test_yiciyuan_protocol_fjb02.yaml | 26 +- .../test_yiciyuan_protocol_fjb03.yaml | 152 +++++ .../tests/util/device_test/mod.rs | 3 + 17 files changed, 1673 insertions(+), 59 deletions(-) create mode 100644 crates/buttplug_server/src/device/yiciyuan_motion.rs create mode 100644 crates/buttplug_tests/tests/util/device_test/device_test_case/test_yiciyuan_protocol_fjb03.yaml diff --git a/crates/buttplug_server/src/device/latest_device_task.rs b/crates/buttplug_server/src/device/latest_device_task.rs index fb7b0f39b..c52534b11 100644 --- a/crates/buttplug_server/src/device/latest_device_task.rs +++ b/crates/buttplug_server/src/device/latest_device_task.rs @@ -9,6 +9,7 @@ use std::{collections::VecDeque, sync::Arc, time::Duration}; use buttplug_core::{errors::ButtplugDeviceError, util::async_manager}; +use buttplug_server_device_config::Endpoint; use futures::{FutureExt, future::BoxFuture}; use tokio::{ select, @@ -19,11 +20,40 @@ use tokio::{ use super::{ device_task::{DeviceTaskConfig, DeviceTaskMessage}, hardware::{Hardware, HardwareCommand, HardwareEvent}, + yiciyuan_motion::{MotionController, settings_for}, }; const WRITE_TIMEOUT: Duration = Duration::from_secs(2); const SLOW_WRITE: Duration = Duration::from_millis(200); +// Observability only: missing ticks never gate commands or trigger a reconnect. +#[derive(Default)] +struct TickMonitor { + last_tick: Option, + samples: u64, + max_gap: Duration, +} + +impl TickMonitor { + fn observe(&mut self, endpoint: Endpoint, data: &[u8], now: Instant) -> bool { + if endpoint != Endpoint::RxBLEBattery || data.len() < 3 || data[..2] != [0x35, 0x14] { + return false; + } + if let Some(last) = self.last_tick { + self.max_gap = self.max_gap.max(now.duration_since(last)); + } + self.last_tick = Some(now); + self.samples += 1; + true + } + + fn age_ms(&self, now: Instant) -> Option { + self + .last_tick + .map(|last| now.duration_since(last).as_millis()) + } +} + #[derive(Default)] struct Batch { commands: VecDeque, @@ -127,11 +157,23 @@ pub(super) async fn run_latest_device_task( config: DeviceTaskConfig, command_receiver: Receiver, ) { + let motion = (hardware.name() == "YCY-FJB-03") + .then(|| MotionController::new(settings_for(hardware.address()))); + // FJB-03: deliberately ignore legacy saved 30/250ms gap overrides. Only + // serialize writes and coalesce superseded state, with no artificial pacing. + let gap = if motion.is_some() { + Duration::ZERO + } else { + config.message_gap.unwrap_or(Duration::from_millis(75)) + }; + let updates = motion.as_ref().map(|_| super::yiciyuan_motion::subscribe()); run( hardware, - config.message_gap.unwrap_or(Duration::ZERO), + gap, command_receiver, WRITE_TIMEOUT, + motion, + updates, ) .await; } @@ -141,12 +183,22 @@ async fn run( gap: Duration, mut receiver: Receiver, write_timeout: Duration, + mut motion: Option, + mut updates: Option< + tokio::sync::watch::Receiver< + std::collections::HashMap, + >, + >, ) { let mut events = hardware.event_stream(); let mut pending = Pending::default(); let mut flight: Option>> = None; let mut next_write = Instant::now(); let mut closed = false; + let mut ticks = TickMonitor::default(); + let mut report_at = Instant::now() + Duration::from_secs(5); + let mut last_input = Instant::now(); + let device_name = hardware.name().clone(); info!( "Latest-state scheduler enabled: gap_ms={}, timeout_ms={}", gap.as_millis(), @@ -154,10 +206,21 @@ async fn run( ); loop { + if !closed && flight.is_none() { + if let (Some(controller), Some(updates)) = (&mut motion, &mut updates) { + let settings = updates + .borrow_and_update() + .get(hardware.address()) + .copied() + .unwrap_or_default(); + controller.reconfigure(settings, Instant::now()); + } + } if closed && flight.is_none() && pending.is_empty() { return; } let can_send = flight.is_none() && !pending.is_empty(); + let motion_deadline = motion.as_ref().and_then(|m| m.deadline()); let delay = if pending.has_urgent() { Duration::ZERO } else { @@ -168,9 +231,31 @@ async fn run( event = events.recv() => { match event { Ok(HardwareEvent::Disconnected(_)) | Err(tokio::sync::broadcast::error::RecvError::Closed) => return, + Ok(HardwareEvent::Notification(_, endpoint, data)) => { + let previous = ticks.last_tick; + if ticks.observe(endpoint, &data, Instant::now()) && previous.is_none() { + info!("{device_name} first device tick received"); + } + if endpoint == Endpoint::RxBLEBattery && data.len() >= 7 && data[..2] == [0x35, 0x10] { + info!("{device_name} device-info reply: model={}, firmware={}, motor_modes={}/{}/{}", + data[2], data[3], data[4], data[5], data[6]); + } + } _ => {} // Notifications/lag do not imply disconnection. } } + _ = async_manager::sleep(report_at.saturating_duration_since(Instant::now())) => { + let now = Instant::now(); + let age = ticks.age_ms(now); + info!("{device_name} tick summary: samples_5s={}, max_gap_ms={}, last_tick_age_ms={:?}, input_idle_ms={}", + ticks.samples, ticks.max_gap.as_millis(), age, now.duration_since(last_input).as_millis()); + if age.is_none_or(|age| age > 2000) { + warn!("{device_name} device ticks missing/stale; observation only, no automatic disconnect"); + } + ticks.samples = 0; + ticks.max_gap = Duration::ZERO; + report_at = now + Duration::from_secs(5); + } result = async { flight.as_mut().unwrap().await }, if flight.is_some() => { flight = None; if let Err(error) = result { @@ -183,7 +268,9 @@ async fn run( } return; } + if let Some(controller) = &mut motion { controller.written(Instant::now()); } } + _ = async { updates.as_mut().unwrap().changed().await }, if !closed && flight.is_none() && updates.is_some() => {} // Timer precedes input so a continuous producer cannot starve output. _ = async_manager::sleep(delay), if can_send => { // Take a bounded snapshot of already queued messages before dispatch. @@ -191,25 +278,46 @@ async fn run( let count = receiver.len(); for _ in 0..count { if let Ok(message) = receiver.try_recv() { + last_input = Instant::now(); pending.receive(message); } } - let batch = pending.take_next(); + let mut batch = pending.take_next(); + if let Some(controller) = &mut motion { + batch.commands = controller.transform(batch.commands, Instant::now()); + } + debug!("{device_name} dispatch tick_age_ms={:?}", ticks.age_ms(Instant::now())); next_write = Instant::now() + gap; flight = Some(write_batch(hardware.clone(), batch, write_timeout).boxed()); } message = receiver.recv(), if !closed => { match message { Some(message) => { + last_input = Instant::now(); pending.receive(message); } None => { closed = true; - // No owner remains to authorize replay of pending motion. + // Preserve the scheduler's no-replay policy for every opted-in model. pending = Pending::default(); + if let Some(controller) = &mut motion { + // The input owner is gone: cancel timers, discard pending motion, + // and make a best-effort zero batch after any in-flight write. + let (stop, ack) = DeviceTaskMessage::acknowledged(controller.stop_all().into_iter().collect()); + drop(ack); + pending.receive(stop); + } }, } } + _ = async_manager::sleep(motion_deadline.map(|t| t.saturating_duration_since(Instant::now())).unwrap_or(Duration::from_secs(3600))), + if !closed && flight.is_none() && pending.is_empty() && receiver.is_empty() && motion_deadline.is_some() => { + let commands = motion.as_mut().unwrap().tick(Instant::now()); + if !commands.is_empty() { + let batch = Batch { commands, acks: vec![], queued_at: Some(Instant::now()) }; + flight = Some(write_batch(hardware.clone(), batch, write_timeout).boxed()); + } + } } } } @@ -229,11 +337,31 @@ mod tests { use tokio::sync::{Semaphore, broadcast, mpsc}; use uuid::Uuid; + #[test] + fn tick_monitor_filters_packets_and_measures_idle_gap() { + let start = Instant::now(); + let mut ticks = TickMonitor::default(); + assert_eq!(ticks.age_ms(start), None); + assert!(!ticks.observe(Endpoint::Tx, &[0x35, 0x14, 0], start)); + assert!(!ticks.observe(Endpoint::RxBLEBattery, &[0x35, 0x13, 1, 80], start)); + assert!(!ticks.observe(Endpoint::RxBLEBattery, &[0x35, 0x14], start)); + assert!(ticks.observe(Endpoint::RxBLEBattery, &[0x35, 0x14, 0], start)); + assert!(ticks.observe( + Endpoint::RxBLEBattery, + &[0x35, 0x14, 1], + start + Duration::from_millis(100) + )); + assert_eq!(ticks.samples, 2); + assert_eq!(ticks.max_gap, Duration::from_millis(100)); + assert_eq!(ticks.age_ms(start + Duration::from_secs(3)), Some(2900)); + } + struct SlowHardware { events: broadcast::Sender, starts: mpsc::UnboundedSender, permits: Arc, fail: bool, + report_a: bool, } impl HardwareInternal for SlowHardware { @@ -253,7 +381,11 @@ mod tests { msg: &HardwareWriteCmd, ) -> BoxFuture<'static, Result<(), ButtplugDeviceError>> { let starts = self.starts.clone(); - let value = msg.data()[0]; + let value = if self.report_a && msg.data().len() == 6 { + msg.data()[2] + } else { + msg.data()[0] + }; let permits = self.permits.clone(); let fail = self.fail; async move { @@ -300,6 +432,27 @@ mod tests { impl Harness { fn new(gap: Duration, timeout: Duration, fail: bool) -> Self { + Self::with_motion(gap, timeout, fail, None) + } + fn with_motion( + gap: Duration, + timeout: Duration, + fail: bool, + motion: Option, + ) -> Self { + Self::with_updates(gap, timeout, fail, motion, None) + } + fn with_updates( + gap: Duration, + timeout: Duration, + fail: bool, + motion: Option, + updates: Option< + tokio::sync::watch::Receiver< + std::collections::HashMap, + >, + >, + ) -> Self { let (events, _) = broadcast::channel(16); let (starts, receiver) = mpsc::unbounded_channel(); let permits = Arc::new(Semaphore::new(0)); @@ -314,10 +467,11 @@ mod tests { starts, permits: permits.clone(), fail, + report_a: motion.is_some(), }), )); let (sender, commands) = mpsc::channel(32); - let task = tokio::spawn(run(hardware, gap, commands, timeout)); + let task = tokio::spawn(run(hardware, gap, commands, timeout, motion, updates)); Self { sender, starts: receiver, @@ -352,6 +506,147 @@ mod tests { HardwareWriteCmd::new(&[Uuid::nil()], Endpoint::Tx, vec![value], false).into() } + fn motion_command(value: u8) -> HardwareCommand { + HardwareWriteCmd::new( + &[Uuid::nil()], + Endpoint::Tx, + vec![0x35, 0x12, value, 0, 0, 0x47 + value], + false, + ) + .into() + } + fn motion_harness(timeout: Duration) -> Harness { + use super::super::yiciyuan_motion::{MotionMode, MotionSettings}; + Harness::with_motion( + Duration::ZERO, + timeout, + false, + Some(MotionController::new(MotionSettings { + mode: MotionMode::Alternate, + dwell_ms: 500, + pause_ms: 100, + })), + ) + } + #[tokio::test] + async fn live_mode_update_waits_for_inflight_write_and_stop_cancels_resume() { + use super::super::yiciyuan_motion::{MotionMode, MotionSettings}; + let (settings, updates) = tokio::sync::watch::channel(std::collections::HashMap::new()); + let mut h = Harness::with_updates( + Duration::ZERO, + Duration::from_secs(2), + false, + Some(MotionController::new(MotionSettings::default())), + Some(updates), + ); + h.sender + .send(DeviceTaskMessage::fire_and_forget(vec![motion_command(10)])) + .await + .unwrap(); + assert_eq!(h.started().await, 10); + settings.send_replace(std::collections::HashMap::from([( + "test".into(), + MotionSettings { + mode: MotionMode::Reverse, + ..Default::default() + }, + )])); + assert!( + tokio::time::timeout(Duration::from_millis(30), h.starts.recv()) + .await + .is_err() + ); + h.permits.add_permits(1); + assert_eq!(h.started().await, 0); + h.permits.add_permits(1); + assert_eq!(h.started().await, 30); + settings.send_replace(std::collections::HashMap::new()); + h.permits.add_permits(1); + assert_eq!(h.started().await, 0); + let (stop, ack) = DeviceTaskMessage::acknowledged(vec![motion_command(0)]); + h.sender.send(stop).await.unwrap(); + h.permits.add_permits(1); + assert_eq!(h.started().await, 0); + h.permits.add_permits(1); + ack.await.unwrap(); + settings.send_replace(std::collections::HashMap::from([( + "test".into(), + MotionSettings { + mode: MotionMode::Alternate, + ..Default::default() + }, + )])); + assert!( + tokio::time::timeout(Duration::from_millis(650), h.starts.recv()) + .await + .is_err() + ); + h.finish().await; + } + + #[tokio::test] + async fn motion_timer_runs_without_new_sender_updates_and_stops_on_zero() { + let mut h = motion_harness(Duration::from_secs(2)); + h.sender + .send(DeviceTaskMessage::fire_and_forget(vec![motion_command(10)])) + .await + .unwrap(); + assert_eq!(h.started().await, 10); + h.permits.add_permits(1); + assert_eq!(h.started().await, 0); + h.permits.add_permits(1); + assert_eq!(h.started().await, 30); + let (stop, ack) = DeviceTaskMessage::acknowledged(vec![motion_command(0)]); + h.sender.send(stop).await.unwrap(); + h.permits.add_permits(1); + assert_eq!(h.started().await, 0); + h.permits.add_permits(1); + ack.await.unwrap(); + assert!( + tokio::time::timeout(Duration::from_millis(650), h.starts.recv()) + .await + .is_err() + ); + h.finish().await; + } + #[tokio::test] + async fn stop_during_direction_pause_never_replays_pending_reverse() { + let mut h = motion_harness(Duration::from_secs(2)); + h.sender + .send(DeviceTaskMessage::fire_and_forget(vec![motion_command(10)])) + .await + .unwrap(); + assert_eq!(h.started().await, 10); + h.permits.add_permits(1); + assert_eq!(h.started().await, 0); // pause write still in flight + let (stop, ack) = DeviceTaskMessage::acknowledged(vec![motion_command(0)]); + h.sender.send(stop).await.unwrap(); + h.permits.add_permits(1); + assert_eq!(h.started().await, 0); + h.permits.add_permits(1); + ack.await.unwrap(); + assert!( + tokio::time::timeout(Duration::from_millis(650), h.starts.recv()) + .await + .is_err() + ); + h.finish().await; + } + #[tokio::test] + async fn motion_write_timeout_ends_timer_without_resuming() { + let mut h = motion_harness(Duration::from_millis(40)); + h.sender + .send(DeviceTaskMessage::fire_and_forget(vec![motion_command(10)])) + .await + .unwrap(); + assert_eq!(h.started().await, 10); + tokio::time::timeout(Duration::from_secs(1), &mut h.task) + .await + .unwrap() + .unwrap(); + assert!(h.starts.try_recv().is_err()); + } + #[tokio::test] async fn busy_write_keeps_only_latest_state() { let mut h = Harness::new(Duration::from_millis(10), Duration::from_secs(2), false); diff --git a/crates/buttplug_server/src/device/mod.rs b/crates/buttplug_server/src/device/mod.rs index 4a5179c5d..db5133164 100644 --- a/crates/buttplug_server/src/device/mod.rs +++ b/crates/buttplug_server/src/device/mod.rs @@ -104,6 +104,7 @@ pub mod protocol; pub mod protocol_impl; mod server_device_manager; mod server_device_manager_event_loop; +pub mod yiciyuan_motion; pub use device_handle::{DeviceCommand, DeviceEvent, DeviceHandle}; pub use output_observation::OutputObservation; diff --git a/crates/buttplug_server/src/device/protocol_impl/yiciyuan.rs b/crates/buttplug_server/src/device/protocol_impl/yiciyuan.rs index f3649667b..a6234a2e6 100644 --- a/crates/buttplug_server/src/device/protocol_impl/yiciyuan.rs +++ b/crates/buttplug_server/src/device/protocol_impl/yiciyuan.rs @@ -8,6 +8,7 @@ use async_trait::async_trait; use std::sync::Arc; use std::sync::atomic::{AtomicU8, Ordering}; +use std::time::Duration; use uuid::{Uuid, uuid}; use futures_util::future::BoxFuture; @@ -15,6 +16,7 @@ use futures_util::{FutureExt, future}; use buttplug_core::errors::ButtplugDeviceError; use buttplug_core::message::{InputReadingV4, InputType, InputTypeReading, InputValue}; +use buttplug_core::util::async_manager; use buttplug_server_device_config::Endpoint; use buttplug_server_device_config::{ @@ -39,16 +41,22 @@ use crate::device::{ generic_protocol_initializer_setup, }, }; +use crate::message::checked_output_cmd::CheckedOutputCmdV4; const YICIYUAN_PROTOCOL_UUID: Uuid = uuid!("d5987116-2fba-4c30-a7aa-ef567a3bf35d"); +// C-mode writes share the physical FF41 characteristic with A/B live writes, +// but are independent protocol state. A separate command identity prevents the +// device-task batcher from treating one packet family as a replacement for the other. +const YICIYUAN_C_MODE_COMMAND_UUID: Uuid = uuid!("b5ccbc68-d970-4e91-b0fa-7ebf74efbb91"); -// Device firmware accepts axes in the range 0..=0x14 (20). Buttplug v4 hands -// us 0..=100 per the YAML range; map by dividing by 5. -const DEVICE_MAX: u8 = 0x14; +// Strength remains unsigned. FJB-03's host motion controller adds the selected +// direction (20 + amplitude for reverse) immediately before the actual write. +const MOTOR_MAX_DEFAULT: u8 = 0x14; +const MOTOR_C_MODES_FJB03: u32 = 7; -// Output feature indices, matching the YAML order under `defaults.features`. +// Output feature indices, matching the YAML order for the three motor slots. const FEATURE_STROKE: u32 = 0; -const FEATURE_VIBE: u32 = 1; +const FEATURE_B: u32 = 1; const FEATURE_AXIS_C: u32 = 2; generic_protocol_initializer_setup!(Yiciyuan, "yiciyuan"); @@ -56,34 +64,100 @@ generic_protocol_initializer_setup!(Yiciyuan, "yiciyuan"); #[derive(Default)] pub struct YiciyuanInitializer {} +// App 6.8.7: settle after discovery, subscribe FF42, settle again, then +// query device info using the model-specific frame. No motor command here. +async fn initialize_fjb( + hardware: &Hardware, + settle: Duration, + timeout: Duration, +) -> Result<(), ButtplugDeviceError> { + let initialization = async { + async_manager::sleep(settle).await; + hardware + .subscribe(&HardwareSubscribeCmd::new( + YICIYUAN_PROTOCOL_UUID, + Endpoint::RxBLEBattery, + )) + .await?; + info!( + "{} initialization: notifications subscribed", + hardware.name() + ); + async_manager::sleep(settle).await; + let query = if hardware.name() == "YCY-FJB-03" { + vec![0x35, 0x10, 0x00, 0x00, 0x00, 0x45] + } else { + let mut query = vec![0u8; 16]; + query[0] = 0x35; + query[1] = 0x10; + query + }; + hardware + .write_value(&HardwareWriteCmd::new( + &[YICIYUAN_PROTOCOL_UUID], + Endpoint::Tx, + query, + false, + )) + .await?; + info!( + "{} initialization: device-info query submitted (not an execution acknowledgement)", + hardware.name() + ); + Ok(()) + }; + tokio::select! { + result = initialization => result, + _ = async_manager::sleep(timeout) => Err(ButtplugDeviceError::DeviceCommunicationError( + format!("{} notification/info initialization timed out", hardware.name()))), + } +} + #[async_trait] impl ProtocolInitializer for YiciyuanInitializer { async fn initialize( &mut self, - _hardware: Arc, + hardware: Arc, _def: &ServerDeviceDefinition, ) -> Result, ButtplugDeviceError> { - Ok(Arc::new(Yiciyuan::default())) + if matches!(hardware.name().as_str(), "YCY-FJB-02" | "YCY-FJB-03") { + initialize_fjb(&hardware, Duration::from_secs(1), Duration::from_secs(7)).await?; + } + Ok(Arc::new(Yiciyuan { + is_fjb03: hardware.name() == "YCY-FJB-03", + is_fjb02: hardware.name() == "YCY-FJB-02", + stroke: AtomicU8::new(0), + motor_b: AtomicU8::new(0), + axis_c: AtomicU8::new(0), + })) } } -/// Per-device state. The protocol sends all three axes in every packet, so -/// we keep the last commanded value for each axis here and rebuild the -/// packet on any axis change. -#[derive(Default)] +/// Per-device state. FJB-01/02 send their three motor slots in one packet. +/// FJB-03 uses the live A/B packet plus a separate fixed-mode C command. pub struct Yiciyuan { + is_fjb03: bool, + is_fjb02: bool, stroke: AtomicU8, - vibe: AtomicU8, + motor_b: AtomicU8, axis_c: AtomicU8, } impl Yiciyuan { fn store(&self, feature_index: u32, value: u32) -> Result<(), ButtplugDeviceError> { - // Map 0..=100 -> 0..=20 (DEVICE_MAX). Round half-up. - let level = ((value.min(100) as u16 * DEVICE_MAX as u16 + 50) / 100) as u8; + if self.is_fjb03 && feature_index == FEATURE_AXIS_C { + let mode = if value == 0 { + 0 + } else { + ((value.min(100) * MOTOR_C_MODES_FJB03 + 99) / 100) as u8 + }; + self.axis_c.store(mode, Ordering::Relaxed); + return Ok(()); + } + let level = ((value.min(100) as u16 * MOTOR_MAX_DEFAULT as u16 + 50) / 100) as u8; match feature_index { FEATURE_STROKE => self.stroke.store(level, Ordering::Relaxed), - FEATURE_VIBE => self.vibe.store(level, Ordering::Relaxed), + FEATURE_B => self.motor_b.store(level, Ordering::Relaxed), FEATURE_AXIS_C => self.axis_c.store(level, Ordering::Relaxed), _ => { return Err(ButtplugDeviceError::ProtocolSpecificError( @@ -96,37 +170,84 @@ impl Yiciyuan { } fn build_packet(&self) -> Vec { - // 16-byte motor-state frame: - // [0]=0x35 vendor magic, [1]=0x12 "set motor levels" sub-command, - // [2]=stroke, [3]=vibe, [4]=axis_c, [5..16]=reserved (zero). + let stroke = self.stroke.load(Ordering::Relaxed); + let motor_b = self.motor_b.load(Ordering::Relaxed); + if self.is_fjb03 { + let body = [0x35u8, 0x12, stroke, motor_b, 0x00]; + let checksum = body.iter().fold(0u16, |sum, byte| sum + *byte as u16) as u8; + let mut packet = Vec::from(body); + packet.push(checksum); + return packet; + } + // FJB-01/02: 16-byte motor-state frame with reserved bytes zero-padded. + // FJB-02 only exposes motor A, which mechanically couples two motions. let mut packet = vec![0u8; 16]; packet[0] = 0x35; packet[1] = 0x12; - packet[2] = self.stroke.load(Ordering::Relaxed); - packet[3] = self.vibe.load(Ordering::Relaxed); + packet[2] = stroke; + packet[3] = motor_b; packet[4] = self.axis_c.load(Ordering::Relaxed); packet } + fn build_fjb03_c_mode_packet(&self) -> Vec { + let mode = self.axis_c.load(Ordering::Relaxed); + let body = [0x35u8, 0x11, 0x04, mode]; + let checksum = body.iter().fold(0u16, |sum, byte| sum + *byte as u16) as u8; + let mut packet = Vec::from(body); + packet.push(checksum); + packet + } + + fn write_packet(&self, packet: Vec) -> HardwareCommand { + let command_id = if packet.get(1) == Some(&0x11) { + YICIYUAN_C_MODE_COMMAND_UUID + } else { + YICIYUAN_PROTOCOL_UUID + }; + HardwareWriteCmd::new(&[command_id], Endpoint::Tx, packet, false).into() + } + fn handle_axis_cmd( &self, feature_index: u32, value: u32, ) -> Result, ButtplugDeviceError> { self.store(feature_index, value)?; - Ok(vec![ - HardwareWriteCmd::new( - &[YICIYUAN_PROTOCOL_UUID], - Endpoint::Tx, - self.build_packet(), - false, - ) - .into(), - ]) + if self.is_fjb03 && feature_index == FEATURE_AXIS_C { + let mut commands = vec![self.write_packet(self.build_fjb03_c_mode_packet())]; + if value == 0 { + // Reapply the current A/B levels, whether this was a C-only stop or + // part of StopCmd. The C field in that packet is always zero. + commands.push(self.write_packet(self.build_packet())); + } + return Ok(commands); + } + + let mut commands = vec![self.write_packet(self.build_packet())]; + if self.is_fjb03 && self.axis_c.load(Ordering::Relaxed) > 0 { + // A live A/B frame may cancel the fixed C mode. Reapply it so moving + // stroke/suction controls does not silently turn vibration off. + commands.push(self.write_packet(self.build_fjb03_c_mode_packet())); + } + Ok(commands) } } impl ProtocolHandler for Yiciyuan { + fn use_latest_output_scheduler(&self) -> bool { + self.is_fjb02 || self.is_fjb03 + } + + fn handle_stop_output_cmd( + &self, + cmd: &CheckedOutputCmdV4, + ) -> Result, ButtplugDeviceError> { + // The DeviceHandle calls this only for explicit StopCmd and flushes the + // result urgently, so it must never apply waveform zero holding. + self.handle_output_cmd(cmd) + } + fn handle_output_oscillate_cmd( &self, feature_index: u32, @@ -145,6 +266,15 @@ impl ProtocolHandler for Yiciyuan { self.handle_axis_cmd(feature_index, speed) } + fn handle_output_constrict_cmd( + &self, + feature_index: u32, + _feature_id: Uuid, + level: u32, + ) -> Result, ButtplugDeviceError> { + self.handle_axis_cmd(feature_index, level) + } + fn handle_input_subscribe_cmd( &self, _device_index: u32, @@ -180,6 +310,11 @@ impl ProtocolHandler for Yiciyuan { feature_id: Uuid, sensor_type: InputType, ) -> BoxFuture<'_, Result<(), ButtplugDeviceError>> { + if (self.is_fjb02 || self.is_fjb03) && sensor_type == InputType::Battery { + // FF42 also carries protocol ticks. A client ending its battery + // subscription must not disable the connection-lifetime notification. + return future::ready(Ok(())).boxed(); + } match sensor_type { InputType::Battery => { async move { @@ -252,3 +387,348 @@ impl ProtocolHandler for Yiciyuan { .boxed() } } + +#[cfg(test)] +mod tests { + use super::*; + use crate::device::hardware::{HardwareInternal, HardwareReadCmd, HardwareReading}; + use buttplug_core::message::{OutputCommand, OutputValue}; + use std::sync::Mutex; + use tokio::sync::broadcast; + + struct InitHardware { + calls: Arc>>, + events: broadcast::Sender, + failure: u8, // 1: subscribe error, 2: write error, 3: subscribe stalls + } + + impl HardwareInternal for InitHardware { + fn disconnect(&self) -> BoxFuture<'static, Result<(), ButtplugDeviceError>> { + future::ready(Ok(())).boxed() + } + fn event_stream(&self) -> broadcast::Receiver { + self.events.subscribe() + } + fn read_value( + &self, + msg: &HardwareReadCmd, + ) -> BoxFuture<'static, Result> { + future::ready(Ok(HardwareReading::new(msg.endpoint(), &[]))).boxed() + } + fn write_value( + &self, + msg: &HardwareWriteCmd, + ) -> BoxFuture<'static, Result<(), ButtplugDeviceError>> { + self.calls.lock().unwrap().push(msg.clone().into()); + future::ready(if self.failure == 2 { + Err(ButtplugDeviceError::DeviceCommunicationError( + "test write failure".into(), + )) + } else { + Ok(()) + }) + .boxed() + } + fn subscribe( + &self, + msg: &HardwareSubscribeCmd, + ) -> BoxFuture<'static, Result<(), ButtplugDeviceError>> { + self.calls.lock().unwrap().push(msg.clone().into()); + if self.failure == 3 { + return future::pending().boxed(); + } + future::ready(if self.failure == 1 { + Err(ButtplugDeviceError::DeviceCommunicationError( + "test subscribe failure".into(), + )) + } else { + Ok(()) + }) + .boxed() + } + fn unsubscribe( + &self, + msg: &HardwareUnsubscribeCmd, + ) -> BoxFuture<'static, Result<(), ButtplugDeviceError>> { + self.calls.lock().unwrap().push(msg.clone().into()); + future::ready(Ok(())).boxed() + } + } + + fn init_hardware(name: &str, failure: u8) -> (Arc, Arc>>) { + let calls = Arc::new(Mutex::new(Vec::new())); + let (events, _) = broadcast::channel(16); + let hw = Hardware::new( + name, + "test", + &[Endpoint::Tx, Endpoint::RxBLEBattery], + &None, + false, + Box::new(InitHardware { + calls: calls.clone(), + events, + failure, + }), + ); + (Arc::new(hw), calls) + } + + #[tokio::test] + async fn fjb02_initialization_subscribes_then_queries_without_motor_output() { + let (hw, calls) = init_hardware("YCY-FJB-02", 0); + initialize_fjb(&hw, Duration::ZERO, Duration::from_secs(1)) + .await + .unwrap(); + let calls = calls.lock().unwrap(); + assert_eq!(calls.len(), 2); + assert!( + matches!(&calls[0], HardwareCommand::Subscribe(cmd) if cmd.endpoint() == Endpoint::RxBLEBattery) + ); + match &calls[1] { + HardwareCommand::Write(cmd) => { + assert_eq!(cmd.endpoint(), Endpoint::Tx); + assert_eq!( + cmd.data(), + &vec![0x35, 0x10, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0] + ); + assert!(!cmd.write_with_response()); + } + _ => panic!("Expected device-info query"), + } + } + + #[tokio::test] + async fn fjb03_initialization_uses_checksummed_query_without_motor_output() { + let (hw, calls) = init_hardware("YCY-FJB-03", 0); + initialize_fjb(&hw, Duration::ZERO, Duration::from_secs(1)) + .await + .unwrap(); + let calls = calls.lock().unwrap(); + assert_eq!(calls.len(), 2); + assert!( + matches!(&calls[0], HardwareCommand::Subscribe(cmd) if cmd.endpoint() == Endpoint::RxBLEBattery) + ); + match &calls[1] { + HardwareCommand::Write(cmd) => { + assert_eq!(cmd.endpoint(), Endpoint::Tx); + assert_eq!(cmd.data(), &vec![0x35, 0x10, 0, 0, 0, 0x45]); + assert!(!cmd.write_with_response()); + } + _ => panic!("Expected device-info query"), + } + } + + #[tokio::test] + async fn fjb_initialization_propagates_errors_and_times_out() { + for name in ["YCY-FJB-02", "YCY-FJB-03"] { + for failure in [1, 2, 3] { + let (hw, calls) = init_hardware(name, failure); + let result = initialize_fjb(&hw, Duration::ZERO, Duration::from_millis(50)).await; + assert!(result.is_err()); + assert_eq!( + calls.lock().unwrap().len(), + if failure == 2 { 2 } else { 1 } + ); + } + } + } + + #[tokio::test] + async fn other_models_do_not_run_fjb_initialization() { + use buttplug_server_device_config::ServerDeviceDefinitionBuilder; + let definition = ServerDeviceDefinitionBuilder::new("test", &Uuid::new_v4()).finish(); + for name in ["YCY-FJB-01", "Unknown"] { + let (hw, calls) = init_hardware(name, 0); + YiciyuanInitializer::default() + .initialize(hw, &definition) + .await + .unwrap(); + assert!(calls.lock().unwrap().is_empty()); + } + } + + #[tokio::test] + async fn fjb_battery_unsubscribe_keeps_protocol_notifications() { + for name in ["YCY-FJB-02", "YCY-FJB-03", "YCY-FJB-01"] { + let (hw, calls) = init_hardware(name, 0); + let mut device = fjb03(); + device.is_fjb03 = name == "YCY-FJB-03"; + device.is_fjb02 = name == "YCY-FJB-02"; + device + .handle_input_unsubscribe_cmd(hw, 1, Uuid::new_v4(), InputType::Battery) + .await + .unwrap(); + let calls = calls.lock().unwrap(); + if name == "YCY-FJB-01" { + assert!(matches!(&calls[..], [HardwareCommand::Unsubscribe(_)])); + } else { + assert!(calls.is_empty()); + } + } + } + + fn fjb03() -> Yiciyuan { + Yiciyuan { + is_fjb03: true, + is_fjb02: false, + stroke: AtomicU8::new(0), + motor_b: AtomicU8::new(0), + axis_c: AtomicU8::new(0), + } + } + + fn packet_data(commands: Vec) -> Vec> { + commands + .into_iter() + .map(|command| match command { + HardwareCommand::Write(write) => write.data().clone(), + other => panic!("Expected a write command, got {other:?}"), + }) + .collect() + } + + #[test] + fn fjb03_c_mode_uses_fixed_mode_command() { + let device = fjb03(); + assert_eq!( + packet_data(device.handle_axis_cmd(FEATURE_AXIS_C, 1).unwrap()), + vec![vec![0x35, 0x11, 0x04, 0x01, 0x4B]] + ); + assert_eq!( + packet_data(device.handle_axis_cmd(FEATURE_AXIS_C, 100).unwrap()), + vec![vec![0x35, 0x11, 0x04, 0x07, 0x51]] + ); + } + + #[test] + fn fjb03_reapplies_vibration_after_stroke_or_suction_changes() { + let device = fjb03(); + device.handle_axis_cmd(FEATURE_AXIS_C, 50).unwrap(); + assert_eq!( + packet_data(device.handle_axis_cmd(FEATURE_STROKE, 50).unwrap()), + vec![ + vec![0x35, 0x12, 0x0A, 0x00, 0x00, 0x51], + vec![0x35, 0x11, 0x04, 0x04, 0x4E], + ] + ); + assert_eq!( + packet_data(device.handle_axis_cmd(FEATURE_B, 50).unwrap()), + vec![ + vec![0x35, 0x12, 0x0A, 0x0A, 0x00, 0x5B], + vec![0x35, 0x11, 0x04, 0x04, 0x4E], + ] + ); + } + + #[test] + fn fjb03_stops_every_motor_in_any_feature_order() { + let stop_orders = [ + [FEATURE_STROKE, FEATURE_B, FEATURE_AXIS_C], + [FEATURE_STROKE, FEATURE_AXIS_C, FEATURE_B], + [FEATURE_B, FEATURE_STROKE, FEATURE_AXIS_C], + [FEATURE_B, FEATURE_AXIS_C, FEATURE_STROKE], + [FEATURE_AXIS_C, FEATURE_STROKE, FEATURE_B], + [FEATURE_AXIS_C, FEATURE_B, FEATURE_STROKE], + ]; + for stop_order in stop_orders { + let device = fjb03(); + for feature in [FEATURE_STROKE, FEATURE_B, FEATURE_AXIS_C] { + device.handle_axis_cmd(feature, 75).unwrap(); + } + let mut writes = Vec::new(); + for feature in stop_order { + writes.extend(packet_data(device.handle_axis_cmd(feature, 0).unwrap())); + } + assert_eq!(device.stroke.load(Ordering::Relaxed), 0); + assert_eq!(device.motor_b.load(Ordering::Relaxed), 0); + assert_eq!(device.axis_c.load(Ordering::Relaxed), 0); + assert!(writes.contains(&vec![0x35, 0x11, 0x04, 0x00, 0x4A])); + assert_eq!( + writes.last(), + Some(&vec![0x35, 0x12, 0x00, 0x00, 0x00, 0x47]) + ); + } + } + + #[test] + fn fjb03_stroke_maps_level_linearly_and_stops_immediately() { + let device = fjb03(); + assert_eq!( + packet_data(device.handle_axis_cmd(FEATURE_STROKE, 10).unwrap()), + vec![vec![0x35, 0x12, 0x02, 0x00, 0x00, 0x49]] + ); + assert_eq!( + packet_data(device.handle_axis_cmd(FEATURE_STROKE, 80).unwrap()), + vec![vec![0x35, 0x12, 0x10, 0x00, 0x00, 0x57]] + ); + assert_eq!( + packet_data(device.handle_axis_cmd(FEATURE_STROKE, 0).unwrap()), + vec![vec![0x35, 0x12, 0x00, 0x00, 0x00, 0x47]] + ); + assert_eq!( + packet_data(device.handle_axis_cmd(FEATURE_STROKE, 60).unwrap()), + vec![vec![0x35, 0x12, 0x0C, 0x00, 0x00, 0x53]] + ); + } + + #[test] + fn fjb03_protocol_emits_unsigned_amplitude_for_host_direction_controller() { + let device = fjb03(); + let mut previous = 0; + for value in 0..=100 { + let packets = packet_data(device.handle_axis_cmd(FEATURE_STROKE, value).unwrap()); + let packet = &packets[0]; + assert_eq!(packet.len(), 6); + assert!(packet[2] >= previous && packet[2] <= 20); + assert_eq!( + packet[5], + packet[..5].iter().copied().fold(0u8, u8::wrapping_add) + ); + previous = packet[2]; + } + assert_eq!(previous, 20); + for value in [101, u32::MAX] { + assert_eq!( + packet_data(device.handle_axis_cmd(FEATURE_STROKE, value).unwrap())[0], + vec![0x35, 0x12, 20, 0, 0, 0x5B] + ); + } + for (value, expected) in [(49, 10), (50, 10), (51, 10), (55, 11)] { + assert_eq!( + packet_data(device.handle_axis_cmd(FEATURE_STROKE, value).unwrap())[0][2], + expected + ); + } + } + + #[test] + fn fjb03_live_and_c_mode_packets_do_not_overlap_in_batching() { + let device = fjb03(); + let live = device + .handle_axis_cmd(FEATURE_STROKE, 50) + .unwrap() + .remove(0); + let c_mode = device + .handle_axis_cmd(FEATURE_AXIS_C, 50) + .unwrap() + .remove(0); + assert!(!live.overlaps(&c_mode)); + } + + #[test] + fn fjb03_never_holds_zero_and_explicit_stop_still_works() { + let device = fjb03(); + let feature_id = Uuid::new_v4(); + let zero = CheckedOutputCmdV4::new( + 1, + 0, + FEATURE_STROKE, + feature_id, + OutputCommand::Oscillate(OutputValue::new(0)), + ); + assert_eq!( + packet_data(device.handle_stop_output_cmd(&zero).unwrap()), + vec![vec![0x35, 0x12, 0x00, 0x00, 0x00, 0x47]] + ); + } +} diff --git a/crates/buttplug_server/src/device/yiciyuan_motion.rs b/crates/buttplug_server/src/device/yiciyuan_motion.rs new file mode 100644 index 000000000..25710aa69 --- /dev/null +++ b/crates/buttplug_server/src/device/yiciyuan_motion.rs @@ -0,0 +1,499 @@ +// Buttplug Rust Source Code File - See https://buttplug.io for more info. +// Licensed under the BSD 3-Clause license. See LICENSE in the project root. + +//! FJB-03 host-side direction modes. The client continues sending unsigned +//! Oscillate strength; direction is NEVER inferred from amplitude or slope. + +use super::hardware::{HardwareCommand, HardwareWriteCmd}; +use buttplug_server_device_config::Endpoint; +use serde::{Deserialize, Serialize}; +use std::{ + collections::{HashMap, VecDeque}, + sync::LazyLock, + time::Duration, +}; +use tokio::time::Instant; +use uuid::uuid; + +#[derive(Clone, Copy, Debug, Default, Deserialize, Serialize, PartialEq, Eq)] +#[serde(rename_all = "snake_case")] +pub enum MotionMode { + #[default] + Forward, + Reverse, + Alternate, +} + +#[derive(Clone, Copy, Debug, Deserialize, Serialize, PartialEq, Eq)] +#[serde(deny_unknown_fields)] +pub struct MotionSettings { + pub mode: MotionMode, + pub dwell_ms: u32, + pub pause_ms: u32, +} + +impl Default for MotionSettings { + fn default() -> Self { + Self { + mode: MotionMode::Forward, + dwell_ms: 1000, + pause_ms: 150, + } + } +} + +impl MotionSettings { + pub fn validate(&self) -> Result<(), String> { + if !(500..=10000).contains(&self.dwell_ms) || !(100..=1000).contains(&self.pause_ms) { + return Err("FJB-03: dwell_ms must be 500..10000; pause_ms must be 100..1000".into()); + } + Ok(()) + } +} + +static SETTINGS: LazyLock>> = + LazyLock::new(|| tokio::sync::watch::channel(HashMap::new()).0); + +/// Publish validated preferences without restarting connected devices. +pub fn configure(settings: HashMap) -> Result<(), String> { + for value in settings.values() { + value.validate()?; + } + SETTINGS.send_if_modified(|current| { + if *current == settings { + return false; + } + *current = settings; + true + }); + Ok(()) +} + +pub(super) fn settings_for(address: &str) -> MotionSettings { + SETTINGS.borrow().get(address).copied().unwrap_or_default() +} + +pub(super) fn subscribe() -> tokio::sync::watch::Receiver> { + SETTINGS.subscribe() +} + +#[derive(Clone, Copy, Debug)] +enum Phase { + Idle, + Running(Instant), + RunWrite, + Pausing(Instant), + PauseWrite, +} + +pub(super) struct MotionController { + settings: MotionSettings, + phase: Phase, + reverse: bool, + target_reverse: Option, + live: Option, + vibration: Option, +} + +impl MotionController { + pub fn new(settings: MotionSettings) -> Self { + Self { + settings, + phase: Phase::Idle, + reverse: settings.mode == MotionMode::Reverse, + target_reverse: None, + live: None, + vibration: None, + } + } + + pub fn deadline(&self) -> Option { + match self.phase { + Phase::Running(t) | Phase::Pausing(t) => Some(t), + _ => None, + } + } + + /// Called only between write batches. Never resurrect a stopped amplitude. + pub fn reconfigure(&mut self, settings: MotionSettings, now: Instant) { + if self.settings == settings { + return; + } + self.settings = settings; + self.target_reverse = None; + if self.live.as_ref().is_some_and(|cmd| cmd.data()[2] != 0) { + self.target_reverse = Some(settings.mode == MotionMode::Reverse); + // Expire the old cycle now; the next packet is zero, then a fresh pause. + self.phase = Phase::Running(now); + } else { + self.phase = Phase::Idle; + self.reverse = settings.mode == MotionMode::Reverse; + } + } + + fn advance(&mut self, now: Instant) { + match self.phase { + Phase::Running(t) if now >= t => self.phase = Phase::PauseWrite, + Phase::Pausing(t) if now >= t => { + self.reverse = self.target_reverse.take().unwrap_or(!self.reverse); + self.phase = Phase::RunWrite; + } + _ => {} + } + } + + /// Start timing only AFTER the whole write batch returns, especially the + /// zero-before-reversal. This is OS submission timing, not a motor ACK. + pub fn written(&mut self, now: Instant) { + match self.phase { + Phase::RunWrite => { + self.phase = if self.settings.mode == MotionMode::Alternate { + Phase::Running(now + Duration::from_millis(self.settings.dwell_ms as u64)) + } else { + Phase::Idle + }; + } + Phase::PauseWrite => { + self.phase = Phase::Pausing(now + Duration::from_millis(self.settings.pause_ms as u64)) + } + _ => {} + } + } + + fn render(&self, cmd: &HardwareWriteCmd) -> HardwareCommand { + let mut data = cmd.data().clone(); + let amplitude = data[2].min(20); + data[2] = if amplitude == 0 || matches!(self.phase, Phase::PauseWrite | Phase::Pausing(_)) { + 0 + } else if self.reverse { + 20 + amplitude + } else { + amplitude + }; + data[5] = data[..5].iter().copied().fold(0u8, u8::wrapping_add); + HardwareWriteCmd::new( + &cmd.command_id().iter().copied().collect::>(), + cmd.endpoint(), + data, + cmd.write_with_response(), + ) + .into() + } + + pub fn transform( + &mut self, + commands: VecDeque, + now: Instant, + ) -> VecDeque { + commands + .into_iter() + .map(|command| { + if let HardwareCommand::Write(cmd) = &command { + if cmd.endpoint() == Endpoint::Tx + && cmd.data().len() == 6 + && cmd.data()[..2] == [0x35, 0x12] + { + self.live = Some(cmd.clone()); + if cmd.data()[2] == 0 { + self.phase = Phase::Idle; + self.target_reverse = None; + self.reverse = self.settings.mode == MotionMode::Reverse; + } else if self.settings.mode == MotionMode::Alternate { + if matches!(self.phase, Phase::Idle) { + self.phase = Phase::RunWrite; + } + } + self.advance(now); + return self.render(cmd); + } + if cmd.endpoint() == Endpoint::Tx + && cmd.data().len() == 5 + && cmd.data()[..3] == [0x35, 0x11, 4] + { + self.vibration = (cmd.data()[3] != 0).then(|| cmd.clone()); + } + } + command + }) + .collect() + } + + pub fn tick(&mut self, now: Instant) -> VecDeque { + if !self.deadline().is_some_and(|t| now >= t) { + return VecDeque::new(); + } + self.advance(now); + let mut result = VecDeque::new(); + if let Some(live) = &self.live { + result.push_back(self.render(live)); + // Live A/B packets can cancel fixed C mode, including timer-generated + // packets. Keep suction unchanged and reapply active vibration. + if let Some(vibration) = &self.vibration { + result.push_back(vibration.clone().into()); + } + } + result + } + + pub fn stop_all(&mut self) -> VecDeque { + self.phase = Phase::Idle; + self.target_reverse = None; + self.live = None; + self.vibration = None; + VecDeque::from([ + HardwareWriteCmd::new( + &[uuid!("b5ccbc68-d970-4e91-b0fa-7ebf74efbb91")], + Endpoint::Tx, + vec![0x35, 0x11, 4, 0, 0x4a], + false, + ) + .into(), + HardwareWriteCmd::new( + &[uuid!("d5987116-2fba-4c30-a7aa-ef567a3bf35d")], + Endpoint::Tx, + vec![0x35, 0x12, 0, 0, 0, 0x47], + false, + ) + .into(), + ]) + } +} + +#[cfg(test)] +mod tests { + use super::*; + fn input(a: u8, b: u8) -> VecDeque { + VecDeque::from([HardwareWriteCmd::new( + &[], + Endpoint::Tx, + vec![0x35, 0x12, a, b, 0, 0x47 + a + b], + false, + ) + .into()]) + } + fn packet(commands: &VecDeque) -> Vec { + match &commands[0] { + HardwareCommand::Write(cmd) => cmd.data().clone(), + _ => panic!(), + } + } + fn alternate() -> MotionController { + MotionController::new(MotionSettings { + mode: MotionMode::Alternate, + ..Default::default() + }) + } + + #[test] + fn both_fixed_directions_use_equal_amplitude_and_zero_is_always_zero() { + for mode in [MotionMode::Forward, MotionMode::Reverse] { + let mut c = MotionController::new(MotionSettings { + mode, + ..Default::default() + }); + for a in 0..=20 { + let p = packet(&c.transform(input(a, 7), Instant::now())); + assert_eq!( + p[2], + if mode == MotionMode::Reverse && a > 0 { + a + 20 + } else { + a + } + ); + assert_eq!(p[3], 7); + assert_eq!(p[5], p[..5].iter().copied().fold(0u8, u8::wrapping_add)); + assert!(c.deadline().is_none()); + } + } + } + #[test] + fn fixed_input_alternates_without_more_client_commands() { + let mut c = alternate(); + let t = Instant::now(); + assert_eq!(packet(&c.transform(input(10, 8), t))[2], 10); + assert!(c.deadline().is_none()); + c.written(t); + assert!(c.tick(t + Duration::from_millis(999)).is_empty()); + assert_eq!(packet(&c.tick(t + Duration::from_secs(1)))[2], 0); + assert!(c.deadline().is_none()); + c.written(t + Duration::from_secs(2)); // slow zero submission: pause starts here + assert!(c.tick(t + Duration::from_millis(2149)).is_empty()); + let p = packet(&c.tick(t + Duration::from_millis(2150))); + assert_eq!((p[2], p[3]), (30, 8)); + c.written(t + Duration::from_millis(2150)); + assert_eq!(packet(&c.tick(t + Duration::from_millis(3150)))[2], 0); + c.written(t + Duration::from_millis(3150)); + assert_eq!(packet(&c.tick(t + Duration::from_millis(3300)))[2], 10); + } + #[test] + fn newest_amplitude_during_pause_is_used_without_shortening_pause() { + let mut c = alternate(); + let t = Instant::now(); + c.transform(input(10, 0), t); + c.written(t); + c.tick(t + Duration::from_secs(1)); + c.written(t + Duration::from_secs(1)); + assert_eq!( + packet(&c.transform(input(16, 9), t + Duration::from_millis(1050)))[2], + 0 + ); + c.written(t + Duration::from_millis(1050)); + assert_eq!(c.deadline(), Some(t + Duration::from_millis(1150))); + assert_eq!(packet(&c.tick(t + Duration::from_millis(1150)))[2], 36); + } + #[test] + fn zero_cancels_every_phase_and_never_restarts_from_a_timer() { + for phase in 0..3 { + let mut c = alternate(); + let t = Instant::now(); + c.transform(input(10, 0), t); + c.written(t); + if phase > 0 { + c.tick(t + Duration::from_secs(1)); + } + if phase > 1 { + c.written(t + Duration::from_secs(1)); + } + assert_eq!( + packet(&c.transform(input(0, 0), t + Duration::from_secs(2)))[2], + 0 + ); + c.written(t + Duration::from_secs(2)); + assert!(c.tick(t + Duration::from_secs(30)).is_empty()); + assert_eq!( + packet(&c.transform(input(8, 0), t + Duration::from_secs(31)))[2], + 8 + ); + } + } + #[test] + fn continuous_updates_do_not_reset_dwell_and_timer_preserves_vibration() { + let mut c = alternate(); + let t = Instant::now(); + c.transform(input(10, 5), t); + c.written(t); + c.transform( + VecDeque::from([HardwareWriteCmd::new( + &[], + Endpoint::Tx, + vec![0x35, 0x11, 4, 3, 0x4d], + false, + ) + .into()]), + t, + ); + c.transform(input(15, 5), t + Duration::from_millis(900)); + c.written(t + Duration::from_millis(900)); + assert_eq!(c.deadline(), Some(t + Duration::from_secs(1))); + let p = c.tick(t + Duration::from_secs(1)); + assert_eq!(p.len(), 2); + assert_eq!(packet(&p)[3], 5); + let stopped = c.stop_all(); + assert_eq!(stopped.len(), 2); + assert!(c.tick(t + Duration::from_secs(20)).is_empty()); + } + #[test] + fn live_switch_replaces_cycle_and_waits_for_zero_submission() { + let mut c = alternate(); + let t = Instant::now(); + c.transform(input(10, 8), t); + c.written(t); + c.reconfigure( + MotionSettings { + mode: MotionMode::Reverse, + ..Default::default() + }, + t, + ); + assert_eq!(packet(&c.tick(t))[2], 0); + assert!(c.deadline().is_none()); + c.written(t + Duration::from_secs(2)); + assert_eq!( + packet(&c.transform(input(15, 8), t + Duration::from_millis(2100)))[2], + 0 + ); + c.written(t + Duration::from_millis(2100)); + let p = packet(&c.tick(t + Duration::from_millis(2150))); + assert_eq!((p[2], p[3]), (35, 8)); + c.written(t + Duration::from_millis(2150)); + assert!(c.deadline().is_none()); + assert!(c.tick(t + Duration::from_secs(30)).is_empty()); + c.reconfigure(MotionSettings::default(), t + Duration::from_secs(31)); + assert_eq!(packet(&c.tick(t + Duration::from_secs(31)))[2], 0); + c.written(t + Duration::from_secs(31)); + assert_eq!(packet(&c.tick(t + Duration::from_millis(31150)))[2], 15); + } + + #[test] + fn live_settings_never_restart_zero_or_stop_and_latest_mode_wins() { + let mut c = alternate(); + let t = Instant::now(); + c.reconfigure(MotionSettings::default(), t); + assert!(c.tick(t).is_empty()); + c.transform(input(10, 4), t); + c.reconfigure( + MotionSettings { + mode: MotionMode::Reverse, + ..Default::default() + }, + t, + ); + assert_eq!(packet(&c.tick(t))[2], 0); + c.written(t); + c.reconfigure( + MotionSettings { + mode: MotionMode::Alternate, + ..Default::default() + }, + t, + ); + assert_eq!(packet(&c.tick(t))[2], 0); + c.written(t); + assert_eq!(packet(&c.tick(t + Duration::from_millis(150)))[2], 10); + c.written(t + Duration::from_millis(150)); + assert!(c.deadline().is_some()); + c.reconfigure(MotionSettings::default(), t); + c.transform(input(0, 4), t); + c.written(t); + assert!(c.tick(t + Duration::from_secs(30)).is_empty()); + c.stop_all(); + c.reconfigure( + MotionSettings { + mode: MotionMode::Reverse, + ..Default::default() + }, + t, + ); + assert!(c.tick(t + Duration::from_secs(30)).is_empty()); + } + + #[test] + fn invalid_settings_are_rejected_not_silently_applied() { + assert!(MotionSettings::default().validate().is_ok()); + for dwell_ms in [0, 499, 10001, u32::MAX] { + assert!( + MotionSettings { + dwell_ms, + ..Default::default() + } + .validate() + .is_err() + ); + } + for pause_ms in [0, 99, 1001] { + assert!( + MotionSettings { + pause_ms, + ..Default::default() + } + .validate() + .is_err() + ); + } + assert!( + serde_json::from_str::(r#"{"mode":"wrong","dwell_ms":1000,"pause_ms":150}"#) + .is_err() + ); + } +} diff --git a/crates/buttplug_server_device_config/build-config/buttplug-device-config-v5.json b/crates/buttplug_server_device_config/build-config/buttplug-device-config-v5.json index 4887c9826..c04405e0a 100644 --- a/crates/buttplug_server_device_config/build-config/buttplug-device-config-v5.json +++ b/crates/buttplug_server_device_config/build-config/buttplug-device-config-v5.json @@ -1,7 +1,7 @@ { "version": { "major": 5, - "minor": 56 + "minor": 58 }, "protocols": { "activejoy": { @@ -25708,7 +25708,8 @@ "btle": { "names": [ "YCY-FJB-01", - "YCY-FJB-02" + "YCY-FJB-02", + "YCY-FJB-03" ], "services": { "0000ff40-0000-1000-8000-00805f9b34fb": { @@ -25728,11 +25729,111 @@ "name": "Yiciyuan FJB-01" }, { + "features": [ + { + "description": "coupled stroke and rotation", + "id": "7d1de07f-0609-40de-aa1c-e6fe41b97d06", + "index": 0, + "output": { + "oscillate": { + "value": [ + 0, + 100 + ] + } + } + }, + { + "description": "battery level", + "id": "e284dbd3-d18d-409f-894c-bc5bbfeab060", + "index": 1, + "input": { + "battery": { + "command": [ + "Read" + ], + "value": [ + [ + 0, + 100 + ] + ] + } + } + } + ], "id": "48108f07-5871-445b-9f2a-10ceb1809b23", "identifier": [ "YCY-FJB-02" ], "name": "Yiciyuan FJB-02" + }, + { + "features": [ + { + "description": "coupled stroke and rotation", + "id": "74f218e9-204e-4600-baf9-43c942b5a6a0", + "index": 0, + "output": { + "oscillate": { + "value": [ + 0, + 100 + ] + } + } + }, + { + "description": "suction", + "id": "4bf007b8-e7df-4c4c-8fbe-ea112128a70f", + "index": 1, + "output": { + "constrict": { + "value": [ + 0, + 100 + ] + } + } + }, + { + "description": "vibration (7 fixed modes)", + "id": "f01acf53-731b-452e-be21-e912a65409c8", + "index": 2, + "output": { + "vibrate": { + "value": [ + 0, + 100 + ] + } + } + }, + { + "description": "battery level", + "id": "5fb5b0a4-aa3f-4e22-9aa3-32a3666d5141", + "index": 3, + "input": { + "battery": { + "command": [ + "Read" + ], + "value": [ + [ + 0, + 100 + ] + ] + } + } + } + ], + "id": "825357dd-9ed1-4f0f-ab74-7885a2d1bac2", + "identifier": [ + "YCY-FJB-03" + ], + "message_gap_ms": 0, + "name": "Yiciyuan FJB-03" } ], "defaults": { diff --git a/crates/buttplug_server_device_config/device-config/buttplug-device-config-schema-v5.json b/crates/buttplug_server_device_config/device-config/buttplug-device-config-schema-v5.json index 8cb05d6f9..b0be65bce 100644 --- a/crates/buttplug_server_device_config/device-config/buttplug-device-config-schema-v5.json +++ b/crates/buttplug_server_device_config/device-config/buttplug-device-config-schema-v5.json @@ -419,6 +419,10 @@ "protocol_variant": { "type": "string" }, + "message_gap_ms": { + "type": "integer", + "minimum": 0 + }, "features": { "$ref": "#/components/features" } diff --git a/crates/buttplug_server_device_config/device-config/protocols/yiciyuan.yml b/crates/buttplug_server_device_config/device-config/protocols/yiciyuan.yml index cf1b037df..36bc9d904 100644 --- a/crates/buttplug_server_device_config/device-config/protocols/yiciyuan.yml +++ b/crates/buttplug_server_device_config/device-config/protocols/yiciyuan.yml @@ -38,12 +38,11 @@ defaults: index: 3 id: d5987116-2fba-4c30-a7aa-ef567a3bf35d configurations: -# YCY-FJB-01 is the only model verified against physical hardware so far. -# YCY-FJB-02 is its successor in the same product line; the official app's -# code path for both models is identical (same vuex state, same hex-stringed -# 16-byte motor frame, same BLE service/characteristic UUIDs). Adding it -# here so the second model is recognised; flag confirmed device behaviour -# in a future PR once an FJB-02 owner can verify. +# BLE device-info replies confirm that FJB-02 has motor A (8 fixed modes) +# and no B/C motors. Its stroke and rotation move together under one control. +# FJB-03 physical testing identified A as stroke, B as suction and C as +# independent vibration. C responds to its seven fixed modes (35 11 04 N), +# not to real-time levels in the 35 12 packet. - identifier: - YCY-FJB-01 name: Yiciyuan FJB-01 @@ -52,11 +51,72 @@ configurations: - YCY-FJB-02 name: Yiciyuan FJB-02 id: 48108f07-5871-445b-9f2a-10ceb1809b23 + features: + - description: coupled stroke and rotation + id: 7d1de07f-0609-40de-aa1c-e6fe41b97d06 + output: + oscillate: + value: + - 0 + - 100 + index: 0 + - description: battery level + id: e284dbd3-d18d-409f-894c-bc5bbfeab060 + input: + battery: + value: + - - 0 + - 100 + command: + - Read + index: 1 +- identifier: + - YCY-FJB-03 + name: Yiciyuan FJB-03 + id: 825357dd-9ed1-4f0f-ab74-7885a2d1bac2 + # Host mode controller owns direction. No artificial pacing or zero hold. + message_gap_ms: 0 + features: + - description: coupled stroke and rotation + id: 74f218e9-204e-4600-baf9-43c942b5a6a0 + output: + oscillate: + value: + - 0 + - 100 + index: 0 + - description: suction + id: 4bf007b8-e7df-4c4c-8fbe-ea112128a70f + output: + constrict: + value: + - 0 + - 100 + index: 1 + - description: vibration (7 fixed modes) + id: f01acf53-731b-452e-be21-e912a65409c8 + output: + vibrate: + value: + - 0 + - 100 + index: 2 + - description: battery level + id: 5fb5b0a4-aa3f-4e22-9aa3-32a3666d5141 + input: + battery: + value: + - - 0 + - 100 + command: + - Read + index: 3 communication: - btle: names: - YCY-FJB-01 - YCY-FJB-02 + - YCY-FJB-03 services: 0000ff40-0000-1000-8000-00805f9b34fb: tx: 0000ff41-0000-1000-8000-00805f9b34fb diff --git a/crates/buttplug_server_device_config/device-config/version.yaml b/crates/buttplug_server_device_config/device-config/version.yaml index 373c81d3a..0b2b0cc11 100644 --- a/crates/buttplug_server_device_config/device-config/version.yaml +++ b/crates/buttplug_server_device_config/device-config/version.yaml @@ -1,3 +1,3 @@ version: major: 5 - minor: 56 + minor: 58 diff --git a/crates/buttplug_tests/tests/test_device_protocols.rs b/crates/buttplug_tests/tests/test_device_protocols.rs index 7c5043eba..0053c0b9c 100644 --- a/crates/buttplug_tests/tests/test_device_protocols.rs +++ b/crates/buttplug_tests/tests/test_device_protocols.rs @@ -270,6 +270,7 @@ async fn load_test_case(test_file: &str) -> DeviceTestCase { #[test_case("test_xuanhuan_protocol.yaml" ; "Xuanhuan Protocol")] #[test_case("test_yiciyuan_protocol.yaml" ; "Yiciyuan Protocol")] #[test_case("test_yiciyuan_protocol_fjb02.yaml" ; "Yiciyuan Protocol - FJB-02")] +#[test_case("test_yiciyuan_protocol_fjb03.yaml" ; "Yiciyuan Protocol - FJB-03")] #[test_case("test_sdl_gamepad_main_trigger.yaml" ; "SDL Gamepad Main And Triggers")] #[test_case("test_sdl_gamepad_triggers_only.yaml" ; "SDL Gamepad Triggers Only")] #[tokio::test] @@ -406,6 +407,7 @@ async fn test_device_protocols_embedded_v4(test_file: &str) { #[test_case("test_xuanhuan_protocol.yaml" ; "Xuanhuan Protocol")] #[test_case("test_yiciyuan_protocol.yaml" ; "Yiciyuan Protocol")] #[test_case("test_yiciyuan_protocol_fjb02.yaml" ; "Yiciyuan Protocol - FJB-02")] +#[test_case("test_yiciyuan_protocol_fjb03.yaml" ; "Yiciyuan Protocol - FJB-03")] #[tokio::test] async fn test_device_protocols_json_v4(test_file: &str) { //tracing_subscriber::fmt::init(); @@ -539,6 +541,7 @@ async fn test_device_protocols_json_v4(test_file: &str) { #[test_case("test_xuanhuan_protocol.yaml" ; "Xuanhuan Protocol")] #[test_case("test_yiciyuan_protocol.yaml" ; "Yiciyuan Protocol")] #[test_case("test_yiciyuan_protocol_fjb02.yaml" ; "Yiciyuan Protocol - FJB-02")] +#[test_case("test_yiciyuan_protocol_fjb03.yaml" ; "Yiciyuan Protocol - FJB-03")] #[test_case("test_sdl_gamepad_main_trigger.yaml" ; "SDL Gamepad Main And Triggers")] #[test_case("test_sdl_gamepad_triggers_only.yaml" ; "SDL Gamepad Triggers Only")] #[tokio::test] @@ -675,6 +678,7 @@ async fn test_device_protocols_embedded_v3(test_file: &str) { #[test_case("test_xuanhuan_protocol.yaml" ; "Xuanhuan Protocol")] #[test_case("test_yiciyuan_protocol.yaml" ; "Yiciyuan Protocol")] #[test_case("test_yiciyuan_protocol_fjb02.yaml" ; "Yiciyuan Protocol - FJB-02")] +#[test_case("test_yiciyuan_protocol_fjb03.yaml" ; "Yiciyuan Protocol - FJB-03")] #[tokio::test] async fn test_device_protocols_json_v3(test_file: &str) { //tracing_subscriber::fmt::init(); @@ -799,6 +803,7 @@ async fn test_device_protocols_json_v3(test_file: &str) { #[test_case("test_xuanhuan_protocol.yaml" ; "Xuanhuan Protocol")] #[test_case("test_yiciyuan_protocol.yaml" ; "Yiciyuan Protocol")] #[test_case("test_yiciyuan_protocol_fjb02.yaml" ; "Yiciyuan Protocol - FJB-02")] +#[test_case("test_yiciyuan_protocol_fjb03.yaml" ; "Yiciyuan Protocol - FJB-03")] #[tokio::test] async fn test_device_protocols_embedded_v2(test_file: &str) { //tracing_subscriber::fmt::init(); @@ -924,6 +929,7 @@ async fn test_device_protocols_embedded_v2(test_file: &str) { #[test_case("test_xuanhuan_protocol.yaml" ; "Xuanhuan Protocol")] #[test_case("test_yiciyuan_protocol.yaml" ; "Yiciyuan Protocol")] #[test_case("test_yiciyuan_protocol_fjb02.yaml" ; "Yiciyuan Protocol - FJB-02")] +#[test_case("test_yiciyuan_protocol_fjb03.yaml" ; "Yiciyuan Protocol - FJB-03")] #[tokio::test] async fn test_device_protocols_json_v2(test_file: &str) { util::device_test::client::client_v2::run_json_test_case(&load_test_case(test_file).await).await; @@ -1046,6 +1052,7 @@ async fn test_device_protocols_json_v2(test_file: &str) { #[test_case("test_xuanhuan_protocol.yaml" ; "Xuanhuan Protocol")] #[test_case("test_yiciyuan_protocol.yaml" ; "Yiciyuan Protocol")] #[test_case("test_yiciyuan_protocol_fjb02.yaml" ; "Yiciyuan Protocol - FJB-02")] +#[test_case("test_yiciyuan_protocol_fjb03.yaml" ; "Yiciyuan Protocol - FJB-03")] #[tokio::test] async fn test_device_protocols_embedded_v1(test_file: &str) { //tracing_subscriber::fmt::init(); @@ -1169,6 +1176,7 @@ async fn test_device_protocols_embedded_v1(test_file: &str) { #[test_case("test_xuanhuan_protocol.yaml" ; "Xuanhuan Protocol")] #[test_case("test_yiciyuan_protocol.yaml" ; "Yiciyuan Protocol")] #[test_case("test_yiciyuan_protocol_fjb02.yaml" ; "Yiciyuan Protocol - FJB-02")] +#[test_case("test_yiciyuan_protocol_fjb03.yaml" ; "Yiciyuan Protocol - FJB-03")] #[tokio::test] async fn test_device_protocols_json_v1(test_file: &str) { util::device_test::client::client_v1::run_json_test_case(&load_test_case(test_file).await).await; @@ -1248,6 +1256,7 @@ async fn test_device_protocols_json_v1(test_file: &str) { #[test_case("test_xuanhuan_protocol.yaml" ; "Xuanhuan Protocol")] #[test_case("test_yiciyuan_protocol.yaml" ; "Yiciyuan Protocol")] #[test_case("test_yiciyuan_protocol_fjb02.yaml" ; "Yiciyuan Protocol - FJB-02")] +#[test_case("test_yiciyuan_protocol_fjb03.yaml" ; "Yiciyuan Protocol - FJB-03")] #[tokio::test] async fn test_device_protocols_embedded_v0(test_file: &str) { //tracing_subscriber::fmt::init(); @@ -1319,6 +1328,7 @@ async fn test_device_protocols_embedded_v0(test_file: &str) { #[test_case("test_xuanhuan_protocol.yaml" ; "Xuanhuan Protocol")] #[test_case("test_yiciyuan_protocol.yaml" ; "Yiciyuan Protocol")] #[test_case("test_yiciyuan_protocol_fjb02.yaml" ; "Yiciyuan Protocol - FJB-02")] +#[test_case("test_yiciyuan_protocol_fjb03.yaml" ; "Yiciyuan Protocol - FJB-03")] #[tokio::test] async fn test_device_protocols_json_v0(test_file: &str) { util::device_test::client::client_v0::run_json_test_case(&load_test_case(test_file).await).await; diff --git a/crates/buttplug_tests/tests/util/device_test/client/client_v0/mod.rs b/crates/buttplug_tests/tests/util/device_test/client/client_v0/mod.rs index 264496f26..f45301eee 100644 --- a/crates/buttplug_tests/tests/util/device_test/client/client_v0/mod.rs +++ b/crates/buttplug_tests/tests/util/device_test/client/client_v0/mod.rs @@ -167,6 +167,7 @@ pub async fn run_test_case( .expect("Scanning should work."); if let Some(device_init) = &test_case.device_init { + let init_timeout = Duration::from_millis(test_case.device_init_timeout_ms.unwrap_or(500)); // Parse send message into client calls, receives into response checks for command in filter_commands(device_init, 0) { match command { @@ -183,7 +184,7 @@ pub async fn run_test_case( let device_receiver = &mut device_channels[*device_index as usize].receiver; for command in commands { tokio::select! { - _ = tokio::time::sleep(Duration::from_millis(500)) => { + _ = tokio::time::sleep(init_timeout) => { panic!("Timeout while waiting for device init output!") } event = device_receiver.recv() => { diff --git a/crates/buttplug_tests/tests/util/device_test/client/client_v1/mod.rs b/crates/buttplug_tests/tests/util/device_test/client/client_v1/mod.rs index af93d0394..f7fb50007 100644 --- a/crates/buttplug_tests/tests/util/device_test/client/client_v1/mod.rs +++ b/crates/buttplug_tests/tests/util/device_test/client/client_v1/mod.rs @@ -185,6 +185,7 @@ pub async fn run_test_case( .expect("Scanning should work."); if let Some(device_init) = &test_case.device_init { + let init_timeout = Duration::from_millis(test_case.device_init_timeout_ms.unwrap_or(500)); // Parse send message into client calls, receives into response checks for command in filter_commands(device_init, 1) { match command { @@ -201,7 +202,7 @@ pub async fn run_test_case( let device_receiver = &mut device_channels[*device_index as usize].receiver; for command in commands { tokio::select! { - _ = tokio::time::sleep(Duration::from_millis(500)) => { + _ = tokio::time::sleep(init_timeout) => { panic!("Timeout while waiting for device init output!") } event = device_receiver.recv() => { diff --git a/crates/buttplug_tests/tests/util/device_test/client/client_v2/mod.rs b/crates/buttplug_tests/tests/util/device_test/client/client_v2/mod.rs index 303bc506b..4229b652e 100644 --- a/crates/buttplug_tests/tests/util/device_test/client/client_v2/mod.rs +++ b/crates/buttplug_tests/tests/util/device_test/client/client_v2/mod.rs @@ -199,6 +199,7 @@ pub async fn run_test_case( .expect("Scanning should work."); if let Some(device_init) = &test_case.device_init { + let init_timeout = Duration::from_millis(test_case.device_init_timeout_ms.unwrap_or(500)); // Parse send message into client calls, receives into response checks for command in filter_commands(device_init, 2) { match command { @@ -215,7 +216,7 @@ pub async fn run_test_case( let device_receiver = &mut device_channels[*device_index as usize].receiver; for command in commands { tokio::select! { - _ = tokio::time::sleep(Duration::from_millis(500)) => { + _ = tokio::time::sleep(init_timeout) => { panic!("Timeout while waiting for device init output!") } event = device_receiver.recv() => { diff --git a/crates/buttplug_tests/tests/util/device_test/client/client_v3/mod.rs b/crates/buttplug_tests/tests/util/device_test/client/client_v3/mod.rs index 366760fb8..57b58239c 100644 --- a/crates/buttplug_tests/tests/util/device_test/client/client_v3/mod.rs +++ b/crates/buttplug_tests/tests/util/device_test/client/client_v3/mod.rs @@ -210,6 +210,7 @@ pub async fn run_test_case( .expect("Scanning should work."); if let Some(device_init) = &test_case.device_init { + let init_timeout = Duration::from_millis(test_case.device_init_timeout_ms.unwrap_or(500)); // Parse send message into client calls, receives into response checks for command in filter_commands(device_init, 3) { match command { @@ -226,7 +227,7 @@ pub async fn run_test_case( let device_receiver = &mut device_channels[*device_index as usize].receiver; for command in commands { tokio::select! { - _ = tokio::time::sleep(Duration::from_millis(500)) => { + _ = tokio::time::sleep(init_timeout) => { panic!("Timeout while waiting for device output!") } event = device_receiver.recv() => { diff --git a/crates/buttplug_tests/tests/util/device_test/client/client_v4/mod.rs b/crates/buttplug_tests/tests/util/device_test/client/client_v4/mod.rs index 63bab0826..df76e57b9 100644 --- a/crates/buttplug_tests/tests/util/device_test/client/client_v4/mod.rs +++ b/crates/buttplug_tests/tests/util/device_test/client/client_v4/mod.rs @@ -272,6 +272,7 @@ pub async fn run_test_case( .expect("Scanning should work."); if let Some(device_init) = &test_case.device_init { + let init_timeout = Duration::from_millis(test_case.device_init_timeout_ms.unwrap_or(500)); // Parse send message into client calls, receives into response checks for command in filter_commands(device_init, 4) { match command { @@ -288,7 +289,7 @@ pub async fn run_test_case( let device_receiver = &mut device_channels[*device_index as usize].receiver; for command in commands { tokio::select! { - _ = tokio::time::sleep(Duration::from_millis(500)) => { + _ = tokio::time::sleep(init_timeout) => { panic!("Timeout while waiting for device output!") } event = device_receiver.recv() => { diff --git a/crates/buttplug_tests/tests/util/device_test/device_test_case/test_yiciyuan_protocol_fjb02.yaml b/crates/buttplug_tests/tests/util/device_test/device_test_case/test_yiciyuan_protocol_fjb02.yaml index 77246517f..6bae7c6ca 100644 --- a/crates/buttplug_tests/tests/util/device_test/device_test_case/test_yiciyuan_protocol_fjb02.yaml +++ b/crates/buttplug_tests/tests/util/device_test/device_test_case/test_yiciyuan_protocol_fjb02.yaml @@ -2,11 +2,21 @@ devices: - identifier: name: "YCY-FJB-02" expected_name: "Yiciyuan FJB-02" + +device_init_timeout_ms: 3000 +device_init: + - !Commands + device_index: 0 + commands: + - !Subscribe + endpoint: rxblebattery + - !Write + endpoint: tx + data: [0x35, 0x10, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00] + write_with_response: false device_commands: - # Same protocol verification as test_yiciyuan_protocol.yaml — FJB-02 uses - # an identical control packet to FJB-01, so we re-run the multi-feature - # Scalar / Stop sequence to make sure the additional device identifier - # routes through the same handler without regression. + # FJB-02 reports motor counts 8/0/0. Its single control moves stroke and + # rotation together; B/C fields remain zero in the 16-byte frame. - !VersionGated min_spec_version: 3 commands: @@ -29,18 +39,12 @@ device_commands: - Index: 0 Scalar: 0.5 ActuatorType: Oscillate - - Index: 1 - Scalar: 0.75 - ActuatorType: Vibrate - - Index: 2 - Scalar: 0.25 - ActuatorType: Vibrate - !Commands device_index: 0 commands: - !Write endpoint: tx - data: [0x35, 0x12, 0x0A, 0x0F, 0x05, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00] + data: [0x35, 0x12, 0x0A, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00] write_with_response: false - !Messages diff --git a/crates/buttplug_tests/tests/util/device_test/device_test_case/test_yiciyuan_protocol_fjb03.yaml b/crates/buttplug_tests/tests/util/device_test/device_test_case/test_yiciyuan_protocol_fjb03.yaml new file mode 100644 index 000000000..8b5cb91a2 --- /dev/null +++ b/crates/buttplug_tests/tests/util/device_test/device_test_case/test_yiciyuan_protocol_fjb03.yaml @@ -0,0 +1,152 @@ +devices: + - identifier: + name: "YCY-FJB-03" + expected_name: "Yiciyuan FJB-03" +device_init_timeout_ms: 3000 +device_init: + - !Commands + device_index: 0 + commands: + - !Subscribe + endpoint: rxblebattery + - !Write + endpoint: tx + data: [0x35, 0x10, 0x00, 0x00, 0x00, 0x45] + write_with_response: false +device_commands: + # A couples stroke and rotation with proportional speed; B is suction, C vibration. + # Each message is separate so the expected command order is deterministic. + - !VersionGated + min_spec_version: 3 + commands: + - !Messages + device_index: 0 + messages: + - !Scalar + - Index: 0 + Scalar: 0.5 + ActuatorType: Oscillate + - !Commands + device_index: 0 + commands: + - !Write + endpoint: tx + data: [0x35, 0x12, 0x0A, 0x00, 0x00, 0x51] + write_with_response: false + + - !Messages + device_index: 0 + messages: + - !Scalar + - Index: 1 + Scalar: 0.75 + ActuatorType: Constrict + - !Commands + device_index: 0 + commands: + - !Write + endpoint: tx + data: [0x35, 0x12, 0x0A, 0x0F, 0x00, 0x60] + write_with_response: false + + - !Messages + device_index: 0 + messages: + - !Scalar + - Index: 2 + Scalar: 0.25 + ActuatorType: Vibrate + - !Commands + device_index: 0 + commands: + - !Write + endpoint: tx + data: [0x35, 0x11, 0x04, 0x02, 0x4C] + write_with_response: false + + - !Messages + device_index: 0 + messages: + - !Scalar + - Index: 0 + Scalar: 1.0 + ActuatorType: Oscillate + - !Commands + device_index: 0 + commands: + - !Write + endpoint: tx + data: [0x35, 0x12, 0x14, 0x0F, 0x00, 0x6A] + write_with_response: false + - !Write + endpoint: tx + data: [0x35, 0x11, 0x04, 0x02, 0x4C] + write_with_response: false + + - !Messages + device_index: 0 + messages: + - !Scalar + - Index: 1 + Scalar: 1.0 + ActuatorType: Constrict + - !Commands + device_index: 0 + commands: + - !Write + endpoint: tx + data: [0x35, 0x12, 0x14, 0x14, 0x00, 0x6F] + write_with_response: false + - !Write + endpoint: tx + data: [0x35, 0x11, 0x04, 0x02, 0x4C] + write_with_response: false + + - !Messages + device_index: 0 + messages: + - !Scalar + - Index: 2 + Scalar: 1.0 + ActuatorType: Vibrate + - !Commands + device_index: 0 + commands: + - !Write + endpoint: tx + data: [0x35, 0x11, 0x04, 0x07, 0x51] + write_with_response: false + + - !Messages + device_index: 0 + messages: + - !Scalar + - Index: 2 + Scalar: 0.0 + ActuatorType: Vibrate + - !Commands + device_index: 0 + commands: + - !Write + endpoint: tx + data: [0x35, 0x11, 0x04, 0x00, 0x4A] + write_with_response: false + - !Write + endpoint: tx + data: [0x35, 0x12, 0x14, 0x14, 0x00, 0x6F] + write_with_response: false + + - !Messages + device_index: 0 + messages: + - !Scalar + - Index: 1 + Scalar: 0.0 + ActuatorType: Constrict + - !Commands + device_index: 0 + commands: + - !Write + endpoint: tx + data: [0x35, 0x12, 0x14, 0x00, 0x00, 0x5B] + write_with_response: false diff --git a/crates/buttplug_tests/tests/util/device_test/mod.rs b/crates/buttplug_tests/tests/util/device_test/mod.rs index 4845b85f2..bb4547cc7 100644 --- a/crates/buttplug_tests/tests/util/device_test/mod.rs +++ b/crates/buttplug_tests/tests/util/device_test/mod.rs @@ -67,6 +67,9 @@ pub struct DeviceTestCase { device_config_file: Option, user_device_config_file: Option, device_init: Option>, + // Some protocols intentionally settle between initialization operations. + // Keep the historical 500ms default for all other fixtures. + device_init_timeout_ms: Option, device_commands: Vec, }