1use super::*;
2#[cfg(target_os = "linux")]
3use crate::hw::alsa::MidiHub;
4#[cfg(target_os = "macos")]
5use crate::hw::coreaudio::{HwDriver, HwOptions, MidiHub};
6#[cfg(target_os = "windows")]
7use crate::hw::options::HwOptions;
8#[cfg(target_os = "freebsd")]
9use crate::hw::oss::MidiHub;
10#[cfg(target_os = "openbsd")]
11use crate::hw::sndio::{HwDriver, HwOptions, MidiHub};
12#[cfg(target_os = "windows")]
13use crate::hw::wasapi::{self, HwDriver, MidiHub};
14#[cfg(target_os = "macos")]
15use crate::workers::coreaudio_worker::HwWorker;
16#[cfg(target_os = "openbsd")]
17use crate::workers::sndio_worker::HwWorker;
18#[cfg(target_os = "windows")]
19use crate::workers::wasapi_worker::HwWorker;
20use crate::{
21 history::{History, UndoEntry},
22 message::{Action, HwMidiEvent, Message, ProcessTask, SessionSlotState},
23 midi::io::MidiEvent,
24 osc::{OscArg, OscServer, build_error_packet, build_osc_packet},
25 state::State,
26 workers::worker::Worker,
27};
28use std::{
29 collections::{HashMap, VecDeque},
30 net::{SocketAddr, UdpSocket},
31 path::{Path, PathBuf},
32 sync::{Arc, atomic::Ordering},
33 time::{Duration, Instant},
34};
35use tokio::sync::mpsc::{Receiver, Sender, channel};
36use tracing::error;
37
38impl Engine {
39 pub fn state(&self) -> Arc<State> {
40 self.state.clone()
41 }
42
43 pub(crate) fn timing_at_sample(&self, sample: usize) -> (f64, u16, u16) {
44 let bpm = self
45 .tempo_points
46 .iter()
47 .filter(|p| p.sample <= sample)
48 .max_by_key(|p| p.sample)
49 .map(|p| p.bpm)
50 .unwrap_or(self.tempo_bpm)
51 .max(1.0);
52 let (num, den) = self
53 .time_signature_points
54 .iter()
55 .filter(|p| p.sample <= sample)
56 .max_by_key(|p| p.sample)
57 .map(|p| (p.numerator.max(1), p.denominator.max(1)))
58 .unwrap_or((self.tsig_num.max(1), self.tsig_denom.max(1)));
59 (bpm, num, den)
60 }
61
62 pub(crate) fn update_global_tempo_from_map(&mut self) {
63 let (bpm, num, den) = self.timing_at_sample(0);
64 self.tempo_bpm = bpm;
65 self.tsig_num = num;
66 self.tsig_denom = den;
67 }
68
69 pub(crate) fn meter_linear_to_db(peak: f32) -> f32 {
70 if peak <= 1.0e-6 {
71 -90.0
72 } else {
73 (20.0 * peak.log10()).clamp(-90.0, 20.0)
74 }
75 }
76
77 pub(crate) const METER_PUBLISH_INTERVAL: Duration = Duration::from_millis(50);
78 pub(crate) const SESSION_RUNTIME_REPORT_INTERVAL: Duration = Duration::from_millis(50);
79 pub(crate) const TRACK_PROCESS_TIMEOUT: Duration = Duration::from_millis(250);
80 #[cfg(any(target_os = "freebsd", target_os = "linux", target_os = "openbsd"))]
81 pub(crate) const HW_OUT_METER_LINEAR_EPSILON: f32 = 0.0025;
82
83 #[cfg(all(unix, not(target_os = "macos")))]
84 pub(crate) fn session_plugins_dir(&self) -> Option<PathBuf> {
85 self.session_dir.as_ref().map(|d| d.join("plugins"))
86 }
87
88 pub(crate) fn session_audio_dir(&self) -> Option<PathBuf> {
89 self.session_dir.as_ref().map(|d| d.join("audio"))
90 }
91
92 pub(crate) fn session_midi_dir(&self) -> Option<PathBuf> {
93 self.session_dir.as_ref().map(|d| d.join("midi"))
94 }
95
96 pub(crate) fn session_peaks_dir(&self) -> Option<PathBuf> {
97 self.session_dir.as_ref().map(|d| d.join("peaks"))
98 }
99
100 pub(crate) fn ensure_session_subdirs(&self) {
101 if let Some(root) = &self.session_dir {
102 let _ = std::fs::create_dir_all(root.join("plugins"));
103 let _ = std::fs::create_dir_all(root.join("audio"));
104 let _ = std::fs::create_dir_all(root.join("midi"));
105 let _ = std::fs::create_dir_all(root.join("peaks"));
106 }
107 }
108
109 pub fn new(rx: Receiver<Message>, tx: Sender<Message>) -> Self {
110 let (meter_snapshot_producer, _) =
111 crate::triple_buffer::triple_buffer(crate::meter::MeterSnapshot::default());
112 let (transport_snapshot_producer, _) =
113 crate::triple_buffer::triple_buffer(crate::meter::TransportSnapshot::default());
114 let (session_runtime_snapshot_producer, _) =
115 crate::triple_buffer::triple_buffer(crate::meter::SessionRuntimeSnapshot::default());
116 Self::new_with_snapshots(
117 rx,
118 tx,
119 meter_snapshot_producer,
120 transport_snapshot_producer,
121 session_runtime_snapshot_producer,
122 )
123 }
124
125 pub fn new_with_snapshots(
126 rx: Receiver<Message>,
127 tx: Sender<Message>,
128 meter_snapshot_producer: crate::triple_buffer::TripleBufferProducer<
129 crate::meter::MeterSnapshot,
130 >,
131 transport_snapshot_producer: crate::triple_buffer::TripleBufferProducer<
132 crate::meter::TransportSnapshot,
133 >,
134 session_runtime_snapshot_producer: crate::triple_buffer::TripleBufferProducer<
135 crate::meter::SessionRuntimeSnapshot,
136 >,
137 ) -> Self {
138 let state = Arc::new(State::default());
139 let initial_state_snapshot = state.lock().snapshot();
140 let state_snapshot = Arc::new(crate::state::StateSlot::from_pointee(
141 initial_state_snapshot.clone(),
142 ));
143 let collector = basedrop::Collector::new();
146 let hw_ports = Arc::new(arc_swap::ArcSwap::from_pointee(
147 crate::plan_builder::HwPorts {
148 buffer_size: 1024,
149 ..Default::default()
150 },
151 ));
152 let initial_plan =
153 { crate::render_plan::RenderPlan::compile(&initial_state_snapshot, &[], &[], 1024) };
154 let plan_slot = Arc::new(crate::render_plan::PlanSlot::from_pointee(
155 basedrop::Owned::new(&collector.handle(), initial_plan),
156 ));
157 let plan_builder = crate::plan_builder::PlanBuilder::spawn(
158 state_snapshot.clone(),
159 hw_ports.clone(),
160 plan_slot.clone(),
161 collector,
162 );
163 let executor = crate::executor::CycleExecutor::new(plan_slot.clone());
164 Self {
165 rx,
166 tx,
167 clients: vec![],
168 state,
169 state_snapshot,
170 workers: vec![],
171 hw_driver: None,
172 hw_driver_info: None,
173 hw_input_ports: Vec::new(),
174 hw_output_ports: Vec::new(),
175 #[cfg(unix)]
176 jack_runtime: None,
177 midi_hub: Some(MidiHub::default()),
178 hw_worker: None,
179 osc_server: None,
180 osc_reply_socket: None,
181 osc_reply_target: None,
182 pending_hw_midi_events: vec![],
183 pending_hw_midi_events_by_device: HashMap::new(),
184 pending_hw_midi_out_events: vec![],
185 pending_hw_midi_out_events_by_device: vec![],
186 active_hw_notes_by_track: HashMap::new(),
187 active_hw_notes_cycle_start: HashMap::new(),
188 midi_hw_in_routes: vec![],
189 midi_hw_out_routes: vec![],
190 midi_hw_thru_routes: vec![],
191 ready_workers: vec![],
192 pending_requests: VecDeque::new(),
193 awaiting_hwfinished: false,
194 handling_hwfinished: false,
195 transport_panic_flush_pending: false,
196 transport_restart_pending: false,
197 notified_loop_wrap_sample: None,
198 transport_sample: 0,
199 hw_input_latency_frames: 0,
200 hw_output_latency_frames: 0,
201 loop_enabled: false,
202 loop_range_samples: None,
203 metronome_enabled: false,
204 tempo_bpm: 120.0,
205 tsig_num: 4,
206 tsig_denom: 4,
207 tempo_points: vec![crate::message::TempoPoint {
208 sample: 0,
209 bpm: 120.0,
210 }],
211 time_signature_points: vec![crate::message::TimeSignaturePoint {
212 sample: 0,
213 numerator: 4,
214 denominator: 4,
215 }],
216 punch_enabled: false,
217 punch_range_samples: None,
218 audio_recordings: std::collections::HashMap::new(),
219 midi_recordings: std::collections::HashMap::new(),
220 completed_audio_recordings: Vec::new(),
221 completed_midi_recordings: Vec::new(),
222 playing: false,
223 transport_running: false,
224 clip_playback_enabled: true,
225 session_clip_playback_enabled: false,
226 session_transport_sample: 0,
227 session_scene_queue: None,
228 session_scene_queue_length_samples: 0,
229 session_current_scene: None,
230 session_current_scene_previous_scene: None,
231 session_current_scene_start_sample: 0,
232 session_current_scene_length_samples: 0,
233 session_completed_clip_passes: Vec::new(),
234 session_reported_clip_passes: std::collections::HashSet::new(),
235 record_enabled: false,
236 step_recording_enabled: false,
237 session_dir: None,
238 hw_out_level_db: 0.0,
239 hw_out_balance: 0.0,
240 hw_out_muted: false,
241 last_hw_out_meter_publish: None,
242 #[cfg(any(target_os = "freebsd", target_os = "linux", target_os = "openbsd"))]
243 last_hw_out_meter_linear: vec![],
244 hw_out_peak_hold_linear: vec![],
245 #[cfg(any(target_os = "freebsd", target_os = "linux", target_os = "openbsd"))]
246 hw_out_meter_publish_phase: false,
247 last_track_meter_publish: None,
248 last_meter_snapshot_publish: None,
249 last_session_report_publish: None,
250 track_meter_linear_by_track: HashMap::new(),
251 meter_snapshot_producer,
252 transport_snapshot_producer,
253 session_runtime_snapshot_producer,
254 executor,
255 plan_builder,
256 plan_slot,
257 hw_ports,
258 pending_node_jobs: VecDeque::new(),
259 latest_hw_out_meter_db: Arc::new(Vec::new()),
260 latest_track_meter_snapshot: Arc::new(Vec::new()),
261 history: History::default(),
262 history_group: None,
263 history_suspended: false,
264 offline_bounce_jobs: HashMap::new(),
265 pending_bounce_starts: Vec::new(),
266 bounce_worker_tracks: HashMap::new(),
267 pending_midi_learn: None,
268 pending_global_midi_learn: None,
269 pending_session_midi_learn: None,
270 global_midi_learn_play_pause: None,
271 global_midi_learn_stop: None,
272 global_midi_learn_record_toggle: None,
273 session_midi_learn_slots: HashMap::new(),
274 session_midi_learn_scenes: HashMap::new(),
275 session_midi_learn_stop_track: HashMap::new(),
276 session_midi_learn_stop_all: None,
277 midi_cc_gate: HashMap::new(),
278 modulators: Vec::new(),
279 modulator_values: None,
280 }
281 }
282
283 pub(crate) fn publish_state_snapshot(&self) {
284 let snapshot = self.state.lock().snapshot();
285 self.state_snapshot.store(Arc::new(snapshot));
286 }
287
288 pub(crate) fn hw_driver_cycle_samples(&self) -> Option<usize> {
289 self.hw_driver_info.map(|info| info.cycle_samples)
290 }
291
292 #[cfg(unix)]
293 pub(crate) fn jack_cycle_samples(&self) -> Option<usize> {
294 self.jack_runtime.as_ref().map(|j| j.buffer_size)
295 }
296
297 #[cfg(not(unix))]
298 pub(crate) fn jack_cycle_samples(&self) -> Option<usize> {
299 None
300 }
301
302 pub(crate) fn current_cycle_samples(&self) -> usize {
303 self.hw_driver_cycle_samples()
304 .or_else(|| self.jack_cycle_samples())
305 .unwrap_or(0)
306 }
307
308 pub(crate) fn sample_rate(&self) -> f64 {
309 if let Some(info) = self.hw_driver_info {
310 info.sample_rate as f64
311 } else {
312 #[cfg(unix)]
313 {
314 self.jack_runtime
315 .as_ref()
316 .map(|j| j.sample_rate as f64)
317 .unwrap_or(48_000.0)
318 }
319 #[cfg(not(unix))]
320 {
321 48_000.0
322 }
323 }
324 }
325
326 pub(crate) async fn set_hw_playing(&mut self, playing: bool) {
327 if let Some(worker) = &self.hw_worker {
328 let _ = worker.tx.send(Message::HWSetPlaying(playing)).await;
331 } else if let Some(driver) = self.hw_driver.as_mut() {
332 driver.set_playing(playing);
333 }
334 }
335
336 pub(crate) fn active_transport_sample(&self) -> usize {
337 if self.session_clip_playback_enabled && self.playing {
338 self.session_transport_sample
339 } else {
340 self.transport_sample
341 }
342 }
343
344 pub(crate) fn compute_modulator_values(
345 &self,
346 sample: usize,
347 ) -> Arc<std::collections::HashMap<usize, f32>> {
348 let sample_rate = self.sample_rate();
349 let (bpm, tsig_num, tsig_denom) = self.timing_at_sample(sample);
350 let values: std::collections::HashMap<usize, f32> = self
351 .modulators
352 .iter()
353 .filter(|m| m.enabled)
354 .map(|m| {
355 (
356 m.id,
357 m.value_at(sample, sample_rate, bpm, tsig_num, tsig_denom),
358 )
359 })
360 .collect();
361 Arc::new(values)
362 }
363
364 pub(crate) fn apply_modulators(&mut self, sample: usize) -> Vec<Action> {
365 use crate::modulator::ModulatorTarget;
366 let values = self.compute_modulator_values(sample);
367 self.modulator_values = Some(values.clone());
368 let mut echoes = Vec::new();
369 let mut per_track: HashMap<String, (Option<f32>, Option<f32>)> = HashMap::new();
370 let mut clap_params: HashMap<(String, usize, u32), f64> = HashMap::new();
371 let mut vst3_params: HashMap<(String, usize, u32), f32> = HashMap::new();
372 #[cfg(all(unix, not(target_os = "macos")))]
373 let mut lv2_params: HashMap<(String, usize, u32), f32> = HashMap::new();
374 let mut midi_cc_events: HashMap<String, Vec<MidiEvent>> = HashMap::new();
375
376 let map_f32 = |value: f32, min: f32, max: f32| -> f32 {
377 crate::modulator::map_value(value, min, max)
378 };
379 let map_f64 = |value: f32, min: f64, max: f64| -> f64 {
380 crate::modulator::map_value_f64(value, min, max)
381 };
382
383 for m in &self.modulators {
384 if !m.enabled {
385 continue;
386 }
387 let Some(&value) = values.get(&m.id) else {
388 continue;
389 };
390 for target in &m.targets {
391 match target {
392 ModulatorTarget::TrackVolume {
393 track_name,
394 min,
395 max,
396 } => {
397 let clamped = map_f32(value, *min, *max);
398 per_track.entry(track_name.clone()).or_default().0 = Some(clamped);
399 }
400 ModulatorTarget::TrackBalance {
401 track_name,
402 min,
403 max,
404 } => {
405 let clamped = map_f32(value, *min, *max);
406 per_track.entry(track_name.clone()).or_default().1 = Some(clamped);
407 }
408 ModulatorTarget::HwOutVolume { min, max } => {
409 let clamped = map_f32(value, *min, *max);
410 if (self.hw_out_level_db - clamped).abs() > f32::EPSILON {
411 self.hw_out_level_db = clamped;
412 echoes
413 .push(Action::TrackAutomationLevel("hw:out".to_string(), clamped));
414 }
415 }
416 ModulatorTarget::HwOutBalance { min, max } => {
417 let next = map_f32(value, *min, *max).clamp(-1.0, 1.0);
418 if (self.hw_out_balance - next).abs() > f32::EPSILON {
419 self.hw_out_balance = next;
420 echoes.push(Action::TrackAutomationBalance("hw:out".to_string(), next));
421 }
422 }
423 ModulatorTarget::ClapParameter {
424 track_name,
425 instance_id,
426 param_id,
427 min,
428 max,
429 } => {
430 let param_value = map_f64(value, *min, *max);
431 clap_params
432 .insert((track_name.clone(), *instance_id, *param_id), param_value);
433 }
434 ModulatorTarget::Vst3Parameter {
435 track_name,
436 instance_id,
437 param_id,
438 min,
439 max,
440 } => {
441 let param_value = map_f32(value, *min, *max);
442 vst3_params
443 .insert((track_name.clone(), *instance_id, *param_id), param_value);
444 }
445 #[cfg(all(unix, not(target_os = "macos")))]
446 ModulatorTarget::Lv2Parameter {
447 track_name,
448 instance_id,
449 index,
450 min,
451 max,
452 } => {
453 let param_value = map_f32(value, *min, *max);
454 lv2_params.insert((track_name.clone(), *instance_id, *index), param_value);
455 }
456 ModulatorTarget::MidiCc {
457 track_name,
458 channel,
459 cc,
460 } => {
461 let cc_value = (value * 127.0).round() as u8;
462 midi_cc_events
463 .entry(track_name.clone())
464 .or_default()
465 .push(MidiEvent::new(
466 0,
467 vec![0xB0 | (*channel).min(15), (*cc).min(127), cc_value],
468 ));
469 }
470 }
471 }
472 }
473 let state = self.state_snapshot.load_full();
474 for (track_name, (level, balance)) in per_track {
475 if let Some(level) = level
476 && let Some(track) = state.tracks.get(&track_name).cloned()
477 {
478 let t = track.lock();
479 if (t.level() - level).abs() > f32::EPSILON {
480 t.set_level(level);
481 echoes.push(Action::TrackAutomationLevel(track_name.clone(), level));
482 }
483 }
484 if let Some(balance) = balance
485 && let Some(track) = state.tracks.get(&track_name).cloned()
486 {
487 let t = track.lock();
488 let next = balance.clamp(-1.0, 1.0);
489 if (t.balance() - next).abs() > f32::EPSILON {
490 t.set_balance(next);
491 echoes.push(Action::TrackAutomationBalance(track_name.clone(), next));
492 }
493 }
494 }
495
496 for (track_name, events) in midi_cc_events {
497 if let Some(track) = state.tracks.get(&track_name).cloned() {
498 track.lock().rt.pending_modulator_midi_events.extend(events);
499 }
500 }
501
502 for ((track_name, instance_id, param_id), value) in clap_params {
503 if let Some(track) = state.tracks.get(&track_name).cloned()
504 && track
505 .lock()
506 .set_clap_parameter(instance_id, param_id, value)
507 .is_ok()
508 {
509 echoes.push(Action::TrackSetClapParameter {
510 track_name,
511 instance_id,
512 param_id,
513 value,
514 });
515 }
516 }
517 for ((track_name, instance_id, param_id), value) in vst3_params {
518 if let Some(track) = state.tracks.get(&track_name).cloned()
519 && track
520 .lock()
521 .set_vst3_parameter(instance_id, param_id, value)
522 .is_ok()
523 {
524 echoes.push(Action::TrackSetVst3Parameter {
525 track_name,
526 instance_id,
527 param_id,
528 value,
529 });
530 }
531 }
532 #[cfg(all(unix, not(target_os = "macos")))]
533 for ((track_name, instance_id, index), value) in lv2_params {
534 if let Some(track) = state.tracks.get(&track_name).cloned()
535 && track
536 .lock()
537 .set_lv2_control_value(instance_id, index as usize, f64::from(value))
538 .is_ok()
539 {
540 echoes.push(Action::TrackSetLv2ControlValue {
541 track_name,
542 instance_id,
543 index,
544 value,
545 });
546 }
547 }
548
549 echoes
550 }
551
552 pub(crate) fn session_end_sample(&self) -> usize {
553 self.state
554 .lock()
555 .tracks
556 .values()
557 .map(|track| {
558 let track = track.lock();
559 let audio_end = track
560 .audio
561 .clips()
562 .iter()
563 .map(|clip| clip.end)
564 .max()
565 .unwrap_or(0);
566 let midi_end = track
567 .midi
568 .clips()
569 .iter()
570 .map(|clip| clip.end)
571 .max()
572 .unwrap_or(0);
573 audio_end.max(midi_end)
574 })
575 .max()
576 .unwrap_or(0)
577 }
578
579 pub(crate) fn normalize_transport_sample(&self, sample: usize) -> usize {
580 if self.loop_enabled
581 && let Some((loop_start, loop_end)) = self.loop_range_samples
582 && loop_end > loop_start
583 && sample >= loop_end
584 {
585 let loop_len = loop_end - loop_start;
586 return loop_start + (sample - loop_start) % loop_len;
587 }
588 sample
589 }
590
591 pub(crate) fn scheduled_loop_wrap_for_next_cycle(&self) -> Option<(usize, usize, usize)> {
592 if !self.playing || !self.loop_enabled {
593 return None;
594 }
595 let (loop_start, loop_end) = self.loop_range_samples?;
596 if loop_end <= loop_start || self.transport_sample >= loop_end {
597 return None;
598 }
599 let cycle_samples = self.current_cycle_samples();
600 if cycle_samples == 0 {
601 return None;
602 }
603 let next = self.transport_sample.saturating_add(cycle_samples);
604 if next < loop_end {
605 return None;
606 }
607 let after_frames = loop_end.saturating_sub(self.transport_sample);
608 Some((
609 after_frames,
610 loop_start,
611 self.normalize_transport_sample(next),
612 ))
613 }
614
615 pub(crate) fn cycle_segments(&self, frames: usize) -> Vec<(usize, usize, usize)> {
616 if frames == 0 {
617 return vec![];
618 }
619 if !self.loop_enabled {
620 return vec![(
621 self.transport_sample,
622 self.transport_sample.saturating_add(frames),
623 0,
624 )];
625 }
626 let Some((loop_start, loop_end)) = self.loop_range_samples else {
627 return vec![(
628 self.transport_sample,
629 self.transport_sample.saturating_add(frames),
630 0,
631 )];
632 };
633 if loop_end <= loop_start {
634 return vec![(
635 self.transport_sample,
636 self.transport_sample.saturating_add(frames),
637 0,
638 )];
639 }
640 let mut segments = Vec::new();
641 let mut remaining = frames;
642 let mut out_offset = 0usize;
643 let mut current = self.transport_sample;
644 while remaining > 0 {
645 let take = loop_end.saturating_sub(current).min(remaining);
646 if take == 0 {
647 current = loop_start;
648 continue;
649 }
650 segments.push((current, current.saturating_add(take), out_offset));
651 out_offset = out_offset.saturating_add(take);
652 remaining -= take;
653 current = if remaining > 0 {
654 loop_start
655 } else {
656 current.saturating_add(take)
657 };
658 }
659 segments
660 }
661
662 pub(crate) fn recording_segments_for_cycle(&self, frames: usize) -> Vec<(usize, usize, usize)> {
663 let segments = self.cycle_segments(frames);
664 let comp = self.hw_input_latency_frames;
665 let segments: Vec<_> = if comp > 0 {
666 segments
667 .into_iter()
668 .map(|(start, end, offset)| {
669 (start.saturating_sub(comp), end.saturating_sub(comp), offset)
670 })
671 .collect()
672 } else {
673 segments
674 };
675 if !self.punch_enabled {
676 return segments;
677 }
678 let Some((punch_start, punch_end)) = self.punch_range_samples else {
679 return vec![];
680 };
681 if punch_end <= punch_start {
682 return vec![];
683 }
684 let mut clipped = Vec::new();
685 for (segment_start, segment_end, frame_offset) in segments {
686 let start = segment_start.max(punch_start);
687 let end = segment_end.min(punch_end);
688 if end <= start {
689 continue;
690 }
691 let clipped_offset = frame_offset.saturating_add(start.saturating_sub(segment_start));
692 clipped.push((start, end, clipped_offset));
693 }
694 clipped
695 }
696
697 pub async fn init(&mut self) {
698 let max_threads = num_cpus::get();
699 for id in 0..max_threads {
700 let (tx, rx) = channel::<Message>(32);
701 let tx_thread = self.tx.clone();
702 let handler = tokio::spawn(async move {
703 let wrk = Worker::new(id, rx, tx_thread, 8);
704 wrk.await.work().await;
705 });
706 let (node_job_tx, mut node_job_rx) = rtrb::RingBuffer::new(64);
707 let (mut node_result_tx, node_result_rx) = rtrb::RingBuffer::new(64);
708 let node_quit = Arc::new(AtomicBool::new(false));
709 let node_quit_thread = node_quit.clone();
710 let node_thread_handle = std::thread::Builder::new()
711 .name(format!("maolan-node-worker-{id}"))
712 .spawn(move || {
713 crate::enable_flush_denormals_to_zero();
714 if let Err(e) = Worker::try_enable_realtime(8) {
715 tracing::warn!(
716 "Node worker {} realtime priority {} not enabled: {}",
717 id,
718 8,
719 e
720 );
721 }
722 while !node_quit_thread.load(std::sync::atomic::Ordering::Acquire) {
723 match node_job_rx.pop() {
724 Ok(job) => {
725 let mut result = Worker::process_node_job_result(id, job);
726 loop {
727 match node_result_tx.push(result) {
728 Ok(()) => break,
729 Err(rtrb::PushError::Full(returned)) => {
730 if node_quit_thread
731 .load(std::sync::atomic::Ordering::Acquire)
732 {
733 break;
734 }
735 result = returned;
736 std::thread::yield_now();
737 }
738 }
739 }
740 }
741 Err(rtrb::PopError::Empty) => {
742 std::thread::park();
743 }
744 }
745 }
746 })
747 .expect("failed to spawn node worker thread");
748 let node_thread = node_thread_handle.thread().clone();
749 std::mem::forget(node_thread_handle);
750 self.workers.push(WorkerData::with_node_mailbox(
751 tx.clone(),
752 handler,
753 node_job_tx,
754 node_result_rx,
755 node_thread,
756 node_quit,
757 ));
758 }
759 }
760
761 pub(crate) async fn notify_clients(&mut self, action: Result<Action, String>) {
762 self.clients.retain(|client| !client.is_closed());
763 for client in self.clients.iter() {
764 if client
765 .send(Message::Response(action.clone()))
766 .await
767 .is_err()
768 {}
769 }
770 if let Some(reply_to) = self.osc_reply_target {
771 match &action {
772 Err(reason) => {
773 self.send_osc_reply(reply_to, &build_error_packet(reason));
774 }
775 Ok(Action::TrackList(names)) => {
776 let args: Vec<OscArg> = names
777 .iter()
778 .map(|name| OscArg::String(name.clone()))
779 .collect();
780 self.send_osc_reply(
781 reply_to,
782 &build_osc_packet("/response/tracks", &"s".repeat(names.len()), &args),
783 );
784 }
785 Ok(Action::TransportState {
786 sample,
787 tempo_bpm,
788 playing,
789 paused: _,
790 tsig_num,
791 tsig_denom,
792 }) => {
793 self.send_osc_reply(
794 reply_to,
795 &build_osc_packet(
796 "/response/transport",
797 "idffii",
798 &[
799 OscArg::Int(*sample as i32),
800 OscArg::Int(if *playing { 1 } else { 0 }),
801 OscArg::Float(*tempo_bpm as f32),
802 OscArg::Float(0.0), OscArg::Int(*tsig_num as i32),
804 OscArg::Int(*tsig_denom as i32),
805 ],
806 ),
807 );
808 }
809 Ok(Action::MeterSnapshot {
810 hw_out_db,
811 track_meters,
812 }) => {
813 let mut args: Vec<OscArg> = Vec::new();
814 args.push(OscArg::Int(hw_out_db.len() as i32));
815 for db in hw_out_db.iter() {
816 args.push(OscArg::Float(*db));
817 }
818 args.push(OscArg::Int(track_meters.len() as i32));
819 for (name, channels) in track_meters.iter() {
820 args.push(OscArg::String(name.clone()));
821 args.push(OscArg::Int(channels.len() as i32));
822 for db in channels.iter() {
823 args.push(OscArg::Float(*db));
824 }
825 }
826 let types = args
827 .iter()
828 .map(|a| match a {
829 OscArg::String(_) => 's',
830 OscArg::Int(_) => 'i',
831 OscArg::Float(_) => 'f',
832 })
833 .collect::<String>();
834 self.send_osc_reply(
835 reply_to,
836 &build_osc_packet("/response/meters", &types, &args),
837 );
838 }
839 Ok(Action::TrackPluginGraph {
840 track_name,
841 plugins,
842 connections: _,
843 connectable_connections: _,
844 }) => {
845 let mut args: Vec<OscArg> = vec![OscArg::String(track_name.clone())];
846 args.push(OscArg::Int(plugins.len() as i32));
847 for plugin in plugins.iter() {
848 args.push(OscArg::Int(plugin.instance_id as i32));
849 args.push(OscArg::String(plugin.format.clone()));
850 args.push(OscArg::String(plugin.uri.clone()));
851 args.push(OscArg::String(plugin.name.clone()));
852 args.push(OscArg::Int(plugin.bypassed as i32));
853 }
854 let types = args
855 .iter()
856 .map(|a| match a {
857 OscArg::String(_) => 's',
858 OscArg::Int(_) => 'i',
859 OscArg::Float(_) => 'f',
860 })
861 .collect::<String>();
862 self.send_osc_reply(
863 reply_to,
864 &build_osc_packet("/response/plugins", &types, &args),
865 );
866 }
867 Ok(Action::ClapPlugins(plugins)) => {
868 let args: Vec<OscArg> = plugins
869 .iter()
870 .map(|p| OscArg::String(format!("{}|{}", p.path, p.name)))
871 .collect();
872 let types = "s".repeat(args.len());
873 self.send_osc_reply(
874 reply_to,
875 &build_osc_packet("/response/clap_plugins", &types, &args),
876 );
877 }
878 Ok(Action::Vst3Plugins(plugins)) => {
879 let args: Vec<OscArg> = plugins
880 .iter()
881 .map(|p| OscArg::String(format!("{}|{}", p.id, p.name)))
882 .collect();
883 let types = "s".repeat(args.len());
884 self.send_osc_reply(
885 reply_to,
886 &build_osc_packet("/response/vst3_plugins", &types, &args),
887 );
888 }
889 #[cfg(all(unix, not(target_os = "macos")))]
890 Ok(Action::Lv2Plugins(plugins)) => {
891 let args: Vec<OscArg> = plugins
892 .iter()
893 .map(|p| OscArg::String(format!("{}|{}", p.uri, p.name)))
894 .collect();
895 let types = "s".repeat(args.len());
896 self.send_osc_reply(
897 reply_to,
898 &build_osc_packet("/response/lv2_plugins", &types, &args),
899 );
900 }
901 Ok(Action::ClapPluginsUnavailable { error })
902 | Ok(Action::Vst3PluginsUnavailable { error }) => {
903 self.send_osc_reply(reply_to, &build_error_packet(error));
904 }
905 #[cfg(all(unix, not(target_os = "macos")))]
906 Ok(Action::Lv2PluginsUnavailable { error }) => {
907 self.send_osc_reply(reply_to, &build_error_packet(error));
908 }
909 Ok(Action::TrackClapParameters {
910 track_name,
911 instance_id,
912 parameters,
913 }) => {
914 let json = serde_json::json!(
915 parameters
916 .iter()
917 .map(|p| serde_json::json!({
918 "id": p.id,
919 "name": p.name,
920 "module": p.module,
921 "min_value": p.min_value,
922 "max_value": p.max_value,
923 "default_value": p.default_value,
924 }))
925 .collect::<Vec<_>>()
926 )
927 .to_string();
928 self.send_osc_reply(
929 reply_to,
930 &build_osc_packet(
931 "/response/plugin_parameters",
932 "siss",
933 &[
934 OscArg::String(track_name.clone()),
935 OscArg::Int(*instance_id as i32),
936 OscArg::String("clap".to_string()),
937 OscArg::String(json),
938 ],
939 ),
940 );
941 }
942 Ok(Action::TrackVst3Parameters {
943 track_name,
944 instance_id,
945 parameters,
946 }) => {
947 let json = serde_json::to_string(parameters).unwrap_or_default();
948 self.send_osc_reply(
949 reply_to,
950 &build_osc_packet(
951 "/response/plugin_parameters",
952 "siss",
953 &[
954 OscArg::String(track_name.clone()),
955 OscArg::Int(*instance_id as i32),
956 OscArg::String("vst3".to_string()),
957 OscArg::String(json),
958 ],
959 ),
960 );
961 }
962 #[cfg(all(unix, not(target_os = "macos")))]
963 Ok(Action::TrackLv2PluginControls {
964 track_name,
965 instance_id,
966 controls,
967 instance_access_handle: _,
968 }) => {
969 let json = serde_json::json!(
970 controls
971 .iter()
972 .map(|c| serde_json::json!({
973 "index": c.index,
974 "name": c.name,
975 "min": c.min,
976 "max": c.max,
977 "value": c.value,
978 }))
979 .collect::<Vec<_>>()
980 )
981 .to_string();
982 self.send_osc_reply(
983 reply_to,
984 &build_osc_packet(
985 "/response/plugin_parameters",
986 "siss",
987 &[
988 OscArg::String(track_name.clone()),
989 OscArg::Int(*instance_id as i32),
990 OscArg::String("lv2".to_string()),
991 OscArg::String(json),
992 ],
993 ),
994 );
995 }
996 Ok(Action::TrackClapNoteNames {
997 track_name,
998 note_names,
999 }) => {
1000 let json = serde_json::to_string(note_names).unwrap_or_default();
1001 self.send_osc_reply(
1002 reply_to,
1003 &build_osc_packet(
1004 "/response/clap_note_names",
1005 "ss",
1006 &[OscArg::String(track_name.clone()), OscArg::String(json)],
1007 ),
1008 );
1009 }
1010 #[cfg(all(unix, not(target_os = "macos")))]
1011 Ok(Action::TrackLv2Midnam {
1012 track_name,
1013 note_names,
1014 }) => {
1015 let json = serde_json::to_string(note_names).unwrap_or_default();
1016 self.send_osc_reply(
1017 reply_to,
1018 &build_osc_packet(
1019 "/response/lv2_midnam",
1020 "ss",
1021 &[OscArg::String(track_name.clone()), OscArg::String(json)],
1022 ),
1023 );
1024 }
1025 _ => {}
1026 }
1027 }
1028 }
1029
1030 pub(crate) fn send_osc_reply(&mut self, reply_to: SocketAddr, packet: &[u8]) {
1031 if self.osc_reply_socket.is_none() {
1032 self.osc_reply_socket = UdpSocket::bind("0.0.0.0:0").ok();
1033 }
1034 if let Some(socket) = self.osc_reply_socket.as_ref() {
1035 let _ = socket.send_to(packet, reply_to);
1036 }
1037 }
1038
1039 pub(crate) async fn dispatch_request(&mut self, a: Action) {
1040 match a {
1041 Action::TrackOfflineBounceCancel { track_name } => {
1042 if let Some(job) = self.offline_bounce_jobs.get(&track_name) {
1043 job.cancel.store(true, Ordering::Relaxed);
1044 }
1045 }
1046 Action::TrackOfflineBounceCancelAll => {
1047 for job in self.offline_bounce_jobs.values() {
1048 job.cancel.store(true, Ordering::Relaxed);
1049 }
1050 }
1051 _ if !self.offline_bounce_jobs.is_empty() => {
1052 self.pending_requests.push_back(a);
1053 }
1054 Action::OpenAudioDevice { .. }
1055 | Action::OpenMidiInputDevice(_)
1056 | Action::OpenMidiOutputDevice(_)
1057 | Action::RequestMeterSnapshot
1058 | Action::RequestTrackList
1059 | Action::RequestTransportState
1060 | Action::Quit
1061 | Action::Log { .. }
1062 | Action::Play
1063 | Action::Pause
1064 | Action::Stop
1065 | Action::TransportPosition(_)
1066 | Action::JumpToEnd
1067 | Action::SetLoopEnabled(_)
1068 | Action::SetLoopRange(_)
1069 | Action::SetPunchEnabled(_)
1070 | Action::SetPunchRange(_)
1071 | Action::SetMetronomeEnabled(_)
1072 | Action::SetTempo(_)
1073 | Action::SetTimeSignature { .. }
1074 | Action::SetTempoMap { .. }
1075 | Action::SetOscEnabled(_)
1076 | Action::SetClipPlaybackEnabled(_)
1077 | Action::SetRecordEnabled(_)
1078 | Action::SetStepRecording(_)
1079 | Action::StepRecordMidiNote { .. }
1080 | Action::SetClipIdentity { .. }
1081 | Action::SetSessionPath(_)
1082 | Action::ClearHistory
1083 | Action::BeginSessionRestore
1084 | Action::PianoKey { .. }
1085 | Action::ModifyMidiNotes { .. }
1086 | Action::ModifyMidiControllers { .. }
1087 | Action::DeleteMidiControllers { .. }
1088 | Action::InsertMidiControllers { .. }
1089 | Action::DeleteMidiNotes { .. }
1090 | Action::InsertMidiNotes { .. }
1091 | Action::SetMidiSysExEvents { .. }
1092 | Action::Session(_) => {
1093 self.handle_request(a).await;
1094 }
1095 #[cfg(all(unix, not(target_os = "macos")))]
1096 Action::ListLv2Plugins => {
1097 self.handle_request(a).await;
1098 }
1099 Action::ListVst3Plugins => {
1100 self.handle_request(a).await;
1101 }
1102 Action::ListClapPlugins => {
1103 self.handle_request(a).await;
1104 }
1105 Action::ListClapPluginsWithCapabilities => {
1106 self.handle_request(a).await;
1107 }
1108 _ => {
1109 self.pending_requests.push_back(a);
1110 if self.can_schedule_hw_cycle() {
1111 self.request_hw_cycle().await;
1112 } else {
1113 while let Some(next) = self.pending_requests.pop_front() {
1114 self.handle_request(next).await;
1115 }
1116 }
1117 }
1118 };
1119 self.publish_clap_state_dirty().await;
1120 }
1121
1122 pub(crate) fn spawn_plugin_host_stderr_reader(
1123 &self,
1124 stderr: std::process::ChildStderr,
1125 source: String,
1126 ) {
1127 let tx = self.tx.clone();
1128 std::thread::spawn(move || {
1129 use std::io::{BufRead, BufReader};
1130 let reader = BufReader::new(stderr);
1131 for line in reader.lines() {
1132 if let Ok(line) = line
1133 && !line.is_empty()
1134 {
1135 let _ = tx.blocking_send(Message::Request(Action::Log {
1136 source: source.clone(),
1137 message: line,
1138 }));
1139 }
1140 }
1141 });
1142 }
1143
1144 pub(crate) fn set_osc_enabled_with<F>(
1145 &mut self,
1146 enabled: bool,
1147 start_server: F,
1148 ) -> Result<(), String>
1149 where
1150 F: FnOnce(Sender<Message>) -> Result<OscServer, String>,
1151 {
1152 if enabled {
1153 if self.osc_server.is_none() {
1154 self.osc_server = Some(start_server(self.tx.clone())?);
1155 }
1156 } else if let Some(mut server) = self.osc_server.take() {
1157 server.stop();
1158 }
1159 Ok(())
1160 }
1161
1162 pub(crate) async fn request_hw_cycle(&mut self) {
1163 if self.awaiting_hwfinished {
1164 tracing::debug!("request_hw_cycle skipped (already awaiting)");
1165 return;
1166 }
1167 tracing::debug!("request_hw_cycle sending TracksFinished");
1168 self.apply_hw_out_gain_and_meter().await;
1169 self.publish_meter_snapshot_if_due();
1170 if let Some((after_frames, loop_start, cycle_end_sample)) =
1171 self.scheduled_loop_wrap_for_next_cycle()
1172 {
1173 self.notified_loop_wrap_sample = Some(cycle_end_sample);
1174 self.notify_clients(Ok(Action::TransportPositionAt {
1175 sample: loop_start,
1176 after_frames,
1177 }))
1178 .await;
1179 } else {
1180 self.notified_loop_wrap_sample = None;
1181 }
1182 if let Some(worker) = &self.hw_worker {
1183 if !self.pending_hw_midi_out_events_by_device.is_empty() {
1184 let out_events = std::mem::take(&mut self.pending_hw_midi_out_events_by_device);
1185 if let Err(e) = worker.tx.send(Message::HWMidiOutEvents(out_events)).await {
1186 error!("Error sending HWMidiOutEvents {e}");
1187 }
1188 }
1189 match worker.tx.send(Message::TracksFinished).await {
1190 Ok(_) => {
1191 self.awaiting_hwfinished = true;
1192 }
1193 Err(e) => {
1194 error!("Error sending TracksFinished {e}");
1195 }
1196 }
1197 }
1198 }
1199
1200 pub(crate) fn should_publish_hw_out_meters(&mut self) -> bool {
1201 let now = Instant::now();
1202 match self.last_hw_out_meter_publish {
1203 Some(last) if now.duration_since(last) < Self::METER_PUBLISH_INTERVAL => false,
1204 _ => {
1205 self.last_hw_out_meter_publish = Some(now);
1206 true
1207 }
1208 }
1209 }
1210
1211 pub(crate) fn should_publish_track_meters(&mut self) -> bool {
1212 let now = Instant::now();
1213 match self.last_track_meter_publish {
1214 Some(last) if now.duration_since(last) < Self::METER_PUBLISH_INTERVAL => false,
1215 _ => {
1216 self.last_track_meter_publish = Some(now);
1217 true
1218 }
1219 }
1220 }
1221
1222 pub(crate) fn should_publish_hw_out_linear(&mut self, peaks_linear: &[f32]) -> bool {
1223 #[cfg(any(target_os = "freebsd", target_os = "linux", target_os = "openbsd"))]
1224 {
1225 self.hw_out_meter_publish_phase = !self.hw_out_meter_publish_phase;
1226 if !self.hw_out_meter_publish_phase {
1227 return false;
1228 }
1229 let changed = if self.last_hw_out_meter_linear.len() != peaks_linear.len() {
1230 true
1231 } else {
1232 self.last_hw_out_meter_linear
1233 .iter()
1234 .zip(peaks_linear.iter())
1235 .any(|(prev, next)| (prev - next).abs() >= Self::HW_OUT_METER_LINEAR_EPSILON)
1236 };
1237 if !changed {
1238 return false;
1239 }
1240 self.last_hw_out_meter_linear.clear();
1241 self.last_hw_out_meter_linear
1242 .extend_from_slice(peaks_linear);
1243 true
1244 }
1245 #[cfg(not(any(target_os = "freebsd", target_os = "linux", target_os = "openbsd")))]
1246 {
1247 let _ = peaks_linear;
1248 false
1249 }
1250 }
1251
1252 pub(crate) async fn maybe_notify_hw_out_meter(&mut self, _meter_db: Vec<f32>) {
1253 {}
1254 }
1255
1256 pub(crate) async fn apply_hw_out_gain_and_meter(&mut self) {
1257 let gain = if self.hw_out_muted {
1258 0.0
1259 } else {
1260 10.0_f32.powf(self.hw_out_level_db / 20.0)
1261 };
1262 let should_notify_interval = self.should_publish_hw_out_meters();
1263 if let Some(worker) = &self.hw_worker {
1264 let _ = worker
1265 .tx
1266 .send(Message::HWSetOutputGainBalance {
1267 gain,
1268 balance: self.hw_out_balance,
1269 })
1270 .await;
1271 if !should_notify_interval {
1272 return;
1273 }
1274 } else {
1275 #[cfg(unix)]
1276 {
1277 if let Some(jack) = self.jack_runtime.as_ref() {
1278 jack.set_output_gain_linear(gain);
1279 jack.set_output_balance(self.hw_out_balance);
1280 if !should_notify_interval {
1281 return;
1282 }
1283 } else {
1284 return;
1285 }
1286 }
1287 #[cfg(not(unix))]
1288 {
1289 return;
1290 }
1291 }
1292 let peaks_linear = if self.hw_worker.is_some() {
1293 let plan = self.executor.plan().clone();
1294 crate::hw::common::output_meter_linear_from_plan(&plan, gain, self.hw_out_balance)
1295 } else {
1296 #[cfg(unix)]
1297 {
1298 if self.jack_runtime.is_none() {
1299 return;
1300 }
1301 let plan = self.executor.plan().clone();
1304 crate::hw::common::output_meter_linear_from_plan(&plan, gain, self.hw_out_balance)
1305 }
1306 #[cfg(not(unix))]
1307 {
1308 return;
1309 }
1310 };
1311 if self.hw_out_peak_hold_linear.len() != peaks_linear.len() {
1312 self.hw_out_peak_hold_linear.resize(peaks_linear.len(), 0.0);
1313 }
1314 let mut held_peaks = Vec::with_capacity(peaks_linear.len());
1315 for (idx, peak_now) in peaks_linear.iter().copied().enumerate() {
1316 let held = self.hw_out_peak_hold_linear[idx] * 0.92;
1317 let next = peak_now.max(held);
1318 self.hw_out_peak_hold_linear[idx] = next;
1319 held_peaks.push(next);
1320 }
1321 let should_notify =
1322 should_notify_interval && self.should_publish_hw_out_linear(&held_peaks);
1323 let meter_db: Vec<f32> = held_peaks
1324 .into_iter()
1325 .map(Self::meter_linear_to_db)
1326 .collect();
1327 self.latest_hw_out_meter_db = Arc::new(meter_db.clone());
1328 if should_notify {
1329 self.maybe_notify_hw_out_meter(meter_db).await;
1330 }
1331 }
1332
1333 pub(crate) fn preload_track_clips_spawn(&self) {
1334 let tracks: Vec<_> = self
1335 .state_snapshot
1336 .load_full()
1337 .tracks
1338 .values()
1339 .cloned()
1340 .collect();
1341 for track in tracks {
1342 tokio::task::spawn_blocking(move || {
1343 track.lock().preload_clips();
1344 });
1345 }
1346 }
1347
1348 pub(crate) fn preload_track_clips(
1357 &self,
1358 ) -> impl std::future::Future<Output = ()> + Send + 'static {
1359 let tracks: Vec<_> = self
1360 .state_snapshot
1361 .load_full()
1362 .tracks
1363 .values()
1364 .cloned()
1365 .collect();
1366 Self::preload_track_handles(tracks)
1367 }
1368
1369 async fn preload_track_handles(tracks: Vec<crate::state::TrackHandle>) {
1370 if tracks.is_empty() {
1371 return;
1372 }
1373 let mut handles = Vec::with_capacity(tracks.len());
1374 for track in tracks {
1375 handles.push(tokio::task::spawn_blocking(move || {
1376 track.lock().preload_clips();
1377 }));
1378 }
1379 for handle in handles {
1380 if let Err(e) = handle.await {
1381 tracing::warn!("Clip preload task panicked: {e}");
1382 }
1383 }
1384 }
1385
1386 pub(crate) fn task_track_name(task: &ProcessTask) -> String {
1387 match task {
1388 ProcessTask::Track(t) | ProcessTask::FolderInput(t) | ProcessTask::FolderOutput(t) => {
1389 t.lock().name.clone()
1390 }
1391 ProcessTask::Plugin { track, .. } => track.lock().name.clone(),
1392 }
1393 }
1394
1395 pub(crate) fn prepare_task_track(&self, task: &ProcessTask) {
1396 let track = match task {
1397 ProcessTask::Track(t) | ProcessTask::FolderInput(t) | ProcessTask::FolderOutput(t) => t,
1398 ProcessTask::Plugin { track, .. } => track,
1399 };
1400 let mut t = track.lock();
1401 let transport_sample = if self.session_clip_playback_enabled && self.playing {
1402 self.session_transport_sample
1403 } else {
1404 self.transport_sample
1405 };
1406 t.set_transport_sample(transport_sample);
1407 t.set_loop_config(self.loop_enabled, self.loop_range_samples);
1408 t.set_transport_timing(self.tempo_bpm, self.tsig_num, self.tsig_denom);
1409 t.set_clip_playback_enabled(self.clip_playback_enabled && self.playing);
1410 t.set_session_clip_playback_enabled(self.session_clip_playback_enabled && self.playing);
1411 t.set_record_tap_enabled(self.playing && self.record_enabled);
1412 t.audio.set_processing(true);
1413 }
1414
1415 pub(crate) async fn dispatch_node_jobs(&mut self, jobs: Vec<crate::executor::NodeJob>) -> bool {
1419 self.pending_node_jobs.extend(jobs);
1420 let mut cycle_complete = false;
1421 while !self.pending_node_jobs.is_empty() {
1422 let Some(worker_index) = self.take_ready_worker_index() else {
1423 break;
1424 };
1425 let Some(job) = self.pending_node_jobs.pop_front() else {
1426 break;
1427 };
1428 if let Some(crate::render_plan::Op::Task { task, .. }) =
1429 job.plan.nodes.get(job.node as usize)
1430 {
1431 self.prepare_task_track(task);
1432 }
1433 let worker = &mut self.workers[worker_index];
1434 if let Some(node_job_tx) = worker.node_job_tx.as_mut() {
1435 match node_job_tx.push(job) {
1436 Ok(()) => {
1437 if let Some(thread) = &worker.node_thread {
1438 thread.unpark();
1439 }
1440 }
1441 Err(rtrb::PushError::Full(job)) => {
1442 self.pending_node_jobs.push_front(job);
1443 self.push_ready_worker(worker_index);
1444 break;
1445 }
1446 }
1447 } else {
1448 let node = job.node;
1449 error!("Worker {worker_index} has no node-job mailbox");
1450 let outcome = self.executor.abandon_node(node, Instant::now());
1451 self.log_silenced_nodes(&outcome.silenced);
1452 cycle_complete |= outcome.cycle_complete;
1453 self.pending_node_jobs.extend(outcome.jobs);
1454 }
1455 }
1456 cycle_complete
1457 }
1458
1459 fn ensure_metronome_wiring(&mut self) {
1467 let Some(track) = self
1468 .state_snapshot
1469 .load_full()
1470 .tracks
1471 .get(Self::METRONOME_TRACK)
1472 .cloned()
1473 else {
1474 return;
1475 };
1476 let frames = self.current_cycle_samples();
1477 let (_, changed) = track.lock().ensure_metronome_source(frames);
1478 if changed {
1479 self.plan_builder.mark_dirty();
1480 }
1481 }
1482
1483 pub(crate) async fn start_plan_cycle(&mut self) -> bool {
1484 if !self.playing || !self.executor.cycle_complete() || !self.offline_bounce_jobs.is_empty()
1488 {
1489 return false;
1490 }
1491 self.refresh_realtime_infection();
1492 self.ensure_metronome_wiring();
1493 let jobs = self.executor.start_cycle(Instant::now());
1494 if self.dispatch_node_jobs(jobs).await {
1495 self.on_all_tracks_finished().await;
1496 return true;
1497 }
1498 false
1499 }
1500
1501 pub(crate) async fn on_node_done(
1505 &mut self,
1506 worker_id: usize,
1507 epoch: u64,
1508 node: u32,
1509 output_linear: Vec<f32>,
1510 parameter_updates: Vec<Action>,
1511 latency_changed: bool,
1512 ) {
1513 self.push_ready_worker(worker_id);
1514 let mut complete = self.dispatch_node_jobs(Vec::new()).await;
1515 if epoch != self.executor.epoch() {
1516 tracing::debug!(
1517 "dropping stale NodeDone (epoch {} vs {}) for node {}",
1518 epoch,
1519 self.executor.epoch(),
1520 node
1521 );
1522 return;
1523 }
1524 if latency_changed {
1525 self.publish_state_snapshot();
1526 self.plan_builder.mark_dirty();
1527 }
1528 let plan = self.executor.plan().clone();
1529 if let Some(crate::render_plan::Op::Task { task, .. }) = plan.nodes.get(node as usize) {
1530 let track_name = Self::task_track_name(task);
1531 self.track_meter_linear_by_track
1532 .insert(track_name, output_linear);
1533 }
1534 for action in parameter_updates {
1535 self.notify_clients(Ok(action)).await;
1536 }
1537 let now = Instant::now();
1538 let (jobs, done) = self.executor.on_node_done(epoch, node, now);
1539 complete |= done;
1540 complete |= self.dispatch_node_jobs(jobs).await;
1541 let outcome = self
1542 .executor
1543 .force_timeouts(now, Self::TRACK_PROCESS_TIMEOUT);
1544 self.log_silenced_nodes(&outcome.silenced);
1545 complete |= self.dispatch_node_jobs(outcome.jobs).await;
1546 complete |= outcome.cycle_complete;
1547 if complete {
1548 self.on_all_tracks_finished().await;
1549 }
1550 }
1551
1552 pub(crate) async fn poll_node_worker_results(&mut self) {
1553 let mut results = Vec::new();
1554 for worker in &mut self.workers {
1555 if let Some(rx) = worker.node_result_rx.as_mut() {
1556 while let Ok(result) = rx.pop() {
1557 results.push(result);
1558 }
1559 }
1560 }
1561 for result in results {
1562 self.on_node_done(
1563 result.worker_id,
1564 result.epoch,
1565 result.node,
1566 result.output_linear,
1567 result.parameter_updates,
1568 result.latency_changed,
1569 )
1570 .await;
1571 }
1572 }
1573
1574 pub(crate) async fn poll_jack_hw_finished(&mut self) {
1575 #[cfg(unix)]
1576 {
1577 let finished = self
1578 .jack_runtime
1579 .as_ref()
1580 .map(|jack| jack.take_hw_finished_count())
1581 .unwrap_or(0);
1582 if finished > 0 {
1583 self.handle_hw_finished().await;
1584 }
1585 }
1586 }
1587
1588 pub(crate) async fn handle_hw_finished(&mut self) {
1589 if !self.awaiting_hwfinished {
1590 tracing::debug!("HWFinished ignored (not awaiting)");
1591 return;
1592 }
1593 tracing::debug!("HWFinished handling; playing={}", self.playing);
1594 self.handling_hwfinished = true;
1595 self.awaiting_hwfinished = false;
1596 #[cfg(unix)]
1597 {
1598 if let Some(jack) = self.jack_runtime.as_mut() {
1599 if !self.pending_hw_midi_out_events.is_empty() {
1600 let out_events = std::mem::take(&mut self.pending_hw_midi_out_events);
1601 jack.write_events(&out_events);
1602 }
1603 let mut in_events = vec![];
1604 jack.read_events_into(&mut in_events);
1605 if !in_events.is_empty() {
1606 self.pending_hw_midi_events.extend(in_events);
1607 }
1608 let dropped = jack.take_midi_events_dropped();
1609 if dropped > 0 {
1610 tracing::warn!(
1611 "JACK MIDI ring full; {dropped} events dropped since last cycle"
1612 );
1613 }
1614 }
1615 }
1616 #[cfg(unix)]
1617 if self.jack_runtime.is_some() {
1618 self.sync_from_jack_transport().await;
1619 }
1620 while let Some(a) = self.pending_requests.pop_front() {
1621 self.handle_request(a).await;
1622 }
1623 self.apply_mute_solo_policy();
1624 self.append_recorded_cycle();
1625 self.flush_completed_recordings().await;
1626 let hw_in_routes = self.midi_hw_in_routes.clone();
1627 let pending_hw_in_by_device = self.pending_hw_midi_events_by_device.clone();
1628 let mut reconfigured_tracks = Vec::new();
1629 let state = self.state_snapshot.load_full();
1630 for (track_name, track) in state.tracks.iter() {
1631 let mut track_lock = track.lock();
1632 if self.jack_runtime_is_some() {
1633 if !self.pending_hw_midi_events.is_empty() {
1634 track_lock.push_hw_midi_events(&self.pending_hw_midi_events);
1635 }
1636 } else {
1637 for route in hw_in_routes.iter().filter(|r| &r.to_track == track_name) {
1638 if let Some(events) = pending_hw_in_by_device.get(&route.device) {
1639 track_lock.push_hw_midi_events_to_port(route.to_port, events);
1640 }
1641 }
1642 }
1643 if track_lock.setup() {
1644 reconfigured_tracks.push(track_name.clone());
1645 }
1646 }
1647 self.publish_track_meters();
1648 self.publish_session_runtime_reports().await;
1649 self.publish_clap_state_dirty().await;
1650 for track_name in reconfigured_tracks {
1651 let track = state.tracks.get(&track_name).cloned();
1652 if let Some(track) = track {
1653 let (plugins, connections, connectable_connections) = {
1654 let track_lock = track.lock();
1655 (
1656 track_lock.plugin_graph_plugins(false),
1657 track_lock.plugin_graph_connections(),
1658 track_lock.connectable_connections(),
1659 )
1660 };
1661 self.notify_clients(Ok(Action::TrackPluginGraph {
1662 track_name: track_name.clone(),
1663 plugins,
1664 connections,
1665 connectable_connections,
1666 }))
1667 .await;
1668 }
1669 }
1670 self.pending_hw_midi_events.clear();
1671 self.pending_hw_midi_events_by_device.clear();
1672 if self.transport_running {
1673 if self.transport_panic_flush_pending {
1674 self.transport_panic_flush_pending = false;
1675 } else if self.transport_restart_pending {
1676 self.transport_restart_pending = false;
1677 } else {
1678 let next = self
1679 .transport_sample
1680 .saturating_add(self.current_cycle_samples());
1681 let normalized = self.normalize_transport_sample(next);
1682 let wrapped = normalized != next;
1683 self.transport_sample = normalized;
1684 self.publish_transport_snapshot();
1685 if wrapped {
1686 if self.notified_loop_wrap_sample == Some(self.transport_sample) {
1687 self.notified_loop_wrap_sample = None;
1688 } else {
1689 self.notify_clients(Ok(Action::TransportPosition(self.transport_sample)))
1690 .await;
1691 }
1692 }
1693 }
1694 }
1695 if self.session_clip_playback_enabled && self.playing {
1696 self.session_transport_sample = self
1697 .session_transport_sample
1698 .saturating_add(self.current_cycle_samples());
1699 }
1700 {
1701 let echoes = self.apply_modulators(self.active_transport_sample());
1702 for action in echoes {
1703 self.notify_clients(Ok(action)).await;
1704 }
1705 }
1706 self.start_plan_cycle().await;
1707 #[cfg(unix)]
1708 {
1709 if self.jack_runtime.is_some() {
1710 self.awaiting_hwfinished = true;
1711 }
1712 }
1713 self.handling_hwfinished = false;
1714 }
1715
1716 pub(crate) async fn on_executor_tick(&mut self) {
1719 if self.executor.cycle_complete() {
1720 return;
1721 }
1722 let outcome = self
1723 .executor
1724 .force_timeouts(Instant::now(), Self::TRACK_PROCESS_TIMEOUT);
1725 self.log_silenced_nodes(&outcome.silenced);
1726 let mut complete = self.dispatch_node_jobs(outcome.jobs).await;
1727 complete |= outcome.cycle_complete;
1728 if complete {
1729 self.on_all_tracks_finished().await;
1730 }
1731 }
1732
1733 pub(crate) fn log_silenced_nodes(&self, silenced: &[u32]) {
1734 for &node in silenced {
1735 let plan = self.executor.plan();
1736 let name = match plan.nodes.get(node as usize) {
1737 Some(crate::render_plan::Op::Task { task, .. }) => Self::task_track_name(task),
1738 _ => format!("node {node}"),
1739 };
1740 tracing::warn!(
1741 "Node {} ('{}') exceeded process timeout ({} ms); forced silent completion for cycle",
1742 node,
1743 name,
1744 Self::TRACK_PROCESS_TIMEOUT.as_millis()
1745 );
1746 }
1747 }
1748
1749 pub(crate) async fn send_bounce_job(
1752 &mut self,
1753 worker_index: usize,
1754 job: crate::message::OfflineBounceWork,
1755 ) {
1756 let track_name = job.track_name.clone();
1757 self.bounce_worker_tracks
1758 .insert(worker_index, track_name.clone());
1759 let worker = &self.workers[worker_index];
1760 if let Err(e) = worker.tx.send(Message::ProcessOfflineBounce(job)).await {
1761 self.bounce_worker_tracks.remove(&worker_index);
1762 self.offline_bounce_jobs.remove(&track_name);
1763 self.push_ready_worker(worker_index);
1764 self.notify_clients(Err(format!("Failed to schedule offline bounce: {e}")))
1765 .await;
1766 }
1767 }
1768
1769 async fn drain_pending_requests_if_idle(&mut self) {
1772 if self.offline_bounce_jobs.is_empty() {
1773 while let Some(next) = self.pending_requests.pop_front() {
1774 self.handle_request(next).await;
1775 }
1776 }
1777 }
1778
1779 pub(crate) async fn on_all_tracks_finished(&mut self) {
1780 let pending = std::mem::take(&mut self.pending_bounce_starts);
1783 for (worker_index, job) in pending {
1784 self.send_bounce_job(worker_index, job).await;
1785 }
1786 if self.transport_restart_pending {
1787 let state = self.state_snapshot.load_full();
1788 for track in state.tracks.values() {
1789 track.lock().take_hw_midi_out_events();
1790 }
1791 } else if self.hw_worker.is_some() {
1792 self.active_hw_notes_cycle_start = self.active_hw_notes_by_track.clone();
1793 let mut out_events = self.collect_hw_midi_output_events_by_device();
1794 if self.loop_enabled
1795 && let Some((_, loop_end)) = self.loop_range_samples
1796 {
1797 let cycle_end = self
1798 .transport_sample
1799 .saturating_add(self.current_cycle_samples());
1800 if self.transport_sample < loop_end && cycle_end >= loop_end {
1801 let wrap_frame = loop_end
1802 .saturating_sub(self.transport_sample)
1803 .min(self.current_cycle_samples())
1804 as u32;
1805 out_events.extend(self.note_off_events_for_active_snapshot(
1806 &self.active_hw_notes_cycle_start,
1807 wrap_frame,
1808 ));
1809 out_events.sort_by(|a, b| {
1810 a.event
1811 .frame
1812 .cmp(&b.event.frame)
1813 .then_with(|| a.device.cmp(&b.device))
1814 });
1815 }
1816 }
1817 self.pending_hw_midi_out_events_by_device.extend(out_events);
1818 } else {
1819 self.pending_hw_midi_out_events = self.collect_hw_midi_output_events();
1820 }
1821 self.request_hw_cycle().await;
1822 }
1823
1824 pub(crate) fn take_ready_worker_index(&mut self) -> Option<usize> {
1825 while !self.ready_workers.is_empty() {
1826 let worker_index = self.ready_workers.remove(0);
1827 if worker_index < self.workers.len() {
1828 return Some(worker_index);
1829 }
1830 }
1831 None
1832 }
1833
1834 pub(crate) fn push_ready_worker(&mut self, worker_index: usize) {
1835 self.ready_workers.push(worker_index);
1836 }
1837
1838 pub(crate) fn publish_track_meters(&mut self) {
1839 if !self.should_publish_track_meters() {
1840 return;
1841 }
1842 let tracks: Vec<(String, crate::state::TrackHandle)> = self
1843 .state_snapshot
1844 .load_full()
1845 .tracks
1846 .iter()
1847 .map(|(name, track)| (name.clone(), track.clone()))
1848 .collect();
1849 let mut snapshot = Vec::with_capacity(tracks.len());
1850 for (name, track) in &tracks {
1851 let linear = self
1852 .track_meter_linear_by_track
1853 .get(name)
1854 .cloned()
1855 .unwrap_or_else(|| track.lock().output_meter_linear());
1856 let output_db = linear
1857 .iter()
1858 .copied()
1859 .map(Self::meter_linear_to_db)
1860 .collect::<Vec<_>>();
1861 snapshot.push((name.clone(), output_db));
1862 }
1863 self.latest_track_meter_snapshot = Arc::new(snapshot);
1864 }
1865
1866 fn record_session_completed_clip_pass(
1867 &mut self,
1868 track_name: String,
1869 scene_index: usize,
1870 clip_id: String,
1871 pass_index: usize,
1872 start_sample: usize,
1873 length_samples: usize,
1874 ) {
1875 let key = (
1876 track_name.clone(),
1877 scene_index,
1878 clip_id.clone(),
1879 pass_index,
1880 start_sample,
1881 );
1882 if self.session_reported_clip_passes.insert(key) {
1883 self.session_completed_clip_passes
1884 .push(crate::meter::SessionCompletedClipPass {
1885 track_name,
1886 scene_index,
1887 clip_id,
1888 pass_index,
1889 start_sample,
1890 length_samples,
1891 });
1892 }
1893 }
1894
1895 fn record_completed_session_scene_span(
1896 &mut self,
1897 tracks: &HashMap<String, crate::state::TrackHandle>,
1898 scene_index: usize,
1899 previous_scene: Option<usize>,
1900 scene_start_sample: usize,
1901 elapsed_samples: usize,
1902 ) {
1903 for (track_name, track) in tracks {
1904 let track = track.lock();
1905 let slot = track.rt.session_slots.get(&scene_index);
1906 let prev_slot = previous_scene.and_then(|scene| track.rt.session_slots.get(&scene));
1907 let playing_clip_id = track
1908 .rt
1909 .playing_session_clips
1910 .last()
1911 .map(|clip| clip.clip_id.as_str());
1912 let clip_id = match (slot, prev_slot) {
1913 (Some(slot), _) if slot.play_enabled => Some(slot.clip_id.as_str()),
1914 (Some(slot), _) if slot.stop_enabled => None,
1915 (_, Some(prev_slot)) if prev_slot.play_enabled => Some(prev_slot.clip_id.as_str()),
1916 (_, Some(prev_slot)) if prev_slot.stop_enabled => None,
1917 _ => playing_clip_id,
1918 }
1919 .filter(|clip_id| !clip_id.is_empty());
1920 let Some(clip_id) = clip_id else { continue };
1921 let clip_length = track
1922 .session_clip_length(clip_id, crate::kind::Kind::Audio)
1923 .or_else(|| track.session_clip_length(clip_id, crate::kind::Kind::MIDI))
1924 .unwrap_or(0);
1925 if clip_length == 0 {
1926 continue;
1927 }
1928 let completed_passes = elapsed_samples / clip_length;
1929 for pass_index in 0..completed_passes {
1930 self.record_session_completed_clip_pass(
1931 track_name.clone(),
1932 scene_index,
1933 clip_id.to_string(),
1934 pass_index,
1935 scene_start_sample.saturating_add(pass_index * clip_length),
1936 clip_length,
1937 );
1938 }
1939 }
1940 }
1941
1942 pub(crate) async fn publish_session_runtime_reports(&mut self) {
1943 let mut current = HashMap::<(String, usize), (SessionSlotState, usize, usize)>::new();
1944 {
1945 let state = self.state_snapshot.load_full();
1946 if let Some((queued_scene, queued_launch_at)) = self.session_scene_queue {
1947 let still_pending = state.tracks.values().any(|track| {
1955 let track = track.lock();
1956 track.rt.pending_session_launches.iter().any(|launch| {
1957 launch.scene_index == queued_scene
1958 && launch.launch_at_sample == queued_launch_at
1959 }) || track
1960 .rt
1961 .playing_session_clips
1962 .iter()
1963 .any(|clip| clip.stop_at_sample == Some(queued_launch_at))
1964 });
1965 if !still_pending && self.session_transport_sample >= queued_launch_at {
1966 if let Some(current_scene) = self.session_current_scene {
1967 let scene_start = self.session_current_scene_start_sample;
1968 let elapsed = queued_launch_at.saturating_sub(scene_start);
1969 self.record_completed_session_scene_span(
1970 &state.tracks,
1971 current_scene,
1972 self.session_current_scene_previous_scene,
1973 scene_start,
1974 elapsed,
1975 );
1976 }
1977 self.session_scene_queue = None;
1978 self.session_current_scene_previous_scene = self.session_current_scene;
1979 self.session_current_scene = Some(queued_scene);
1980 self.session_current_scene_start_sample = queued_launch_at;
1981 self.session_current_scene_length_samples =
1982 self.session_scene_queue_length_samples;
1983 self.session_scene_queue_length_samples = 0;
1984 }
1985 }
1986 if let Some(current_scene) = self.session_current_scene
1987 && self.session_current_scene_length_samples > 0
1988 {
1989 self.record_completed_session_scene_span(
1990 &state.tracks,
1991 current_scene,
1992 self.session_current_scene_previous_scene,
1993 self.session_current_scene_start_sample,
1994 self.session_transport_sample
1995 .saturating_sub(self.session_current_scene_start_sample),
1996 );
1997 }
1998 for (track_name, track) in &state.tracks {
1999 let track = track.lock();
2000 for launch in &track.rt.pending_session_launches {
2001 current.insert(
2002 (track_name.clone(), launch.scene_index),
2003 (SessionSlotState::Queued, 0, 0),
2004 );
2005 }
2006 for clip in &track.rt.playing_session_clips {
2007 if self.session_current_scene_length_samples == 0
2008 && let Some(clip_length) =
2009 track.session_clip_length(&clip.clip_id, clip.kind)
2010 {
2011 let scene_report = self
2012 .session_current_scene
2013 .map(|_| self.session_current_scene_length_samples)
2014 .filter(|length| *length > 0)
2015 .map(|length| {
2016 (
2017 self.session_current_scene.unwrap_or(clip.scene_index),
2018 self.session_current_scene_start_sample,
2019 self.session_transport_sample
2020 .saturating_sub(self.session_current_scene_start_sample),
2021 length,
2022 )
2023 });
2024 let launch_sample = self
2025 .session_transport_sample
2026 .saturating_sub(clip.elapsed_samples);
2027 let (
2028 report_scene_index,
2029 report_start_sample,
2030 report_elapsed,
2031 report_length,
2032 ) = scene_report.unwrap_or((
2033 clip.scene_index,
2034 launch_sample,
2035 clip.elapsed_samples,
2036 clip_length,
2037 ));
2038 if report_length == 0 {
2039 continue;
2040 }
2041 let completed_passes = report_elapsed / report_length;
2042 for pass_index in 0..completed_passes {
2043 self.record_session_completed_clip_pass(
2044 track_name.clone(),
2045 report_scene_index,
2046 clip.clip_id.clone(),
2047 pass_index,
2048 report_start_sample.saturating_add(pass_index * report_length),
2049 report_length,
2050 );
2051 }
2052 }
2053 current.insert(
2054 (track_name.clone(), clip.scene_index),
2055 (
2056 SessionSlotState::Playing,
2057 clip.play_position_samples,
2058 clip.elapsed_samples,
2059 ),
2060 );
2061 }
2062 }
2063 }
2064
2065 if self
2066 .last_session_report_publish
2067 .is_some_and(|t| t.elapsed() < Self::SESSION_RUNTIME_REPORT_INTERVAL)
2068 {
2069 return;
2070 }
2071
2072 let snapshot = self.session_runtime_snapshot_producer.write_buffer();
2073 snapshot.session_sample = self.session_transport_sample;
2074 snapshot.slots.clear();
2075 snapshot.slots.extend(current.iter().map(
2076 |((track_name, scene_index), (state, play_position_samples, elapsed_samples))| {
2077 crate::meter::SessionRuntimeSlotSnapshot {
2078 track_name: track_name.clone(),
2079 scene_index: *scene_index,
2080 state: *state,
2081 play_position_samples: *play_position_samples,
2082 elapsed_samples: *elapsed_samples,
2083 }
2084 },
2085 ));
2086 snapshot.completed_clip_passes = self.session_completed_clip_passes.clone();
2087 snapshot.current_scene = self.session_current_scene;
2088 self.session_runtime_snapshot_producer.publish();
2089 self.last_session_report_publish = Some(Instant::now());
2090 }
2091
2092 pub(crate) async fn publish_clap_state_dirty(&mut self) {
2093 let tracks: Vec<(String, crate::state::TrackHandle)> = self
2094 .state_snapshot
2095 .load_full()
2096 .tracks
2097 .iter()
2098 .map(|(name, track)| (name.clone(), track.clone()))
2099 .collect();
2100 for (track_name, track) in &tracks {
2101 let dirty = track.lock().take_dirty_clap_instances();
2102 for instance_id in dirty {
2103 self.notify_clients(Ok(Action::TrackClapStateDirty {
2104 track_name: track_name.clone(),
2105 instance_id,
2106 }))
2107 .await;
2108 }
2109 }
2110 }
2111
2112 pub(crate) fn reset_meters_after_stop(&mut self) {
2113 self.last_hw_out_meter_publish = None;
2114 self.last_track_meter_publish = None;
2115 self.last_meter_snapshot_publish = None;
2116 self.hw_out_peak_hold_linear.fill(0.0);
2117 #[cfg(any(target_os = "freebsd", target_os = "linux", target_os = "openbsd"))]
2118 {
2119 self.last_hw_out_meter_linear.clear();
2120 }
2121 let hw_channels = self.latest_hw_out_meter_db.len();
2122 self.latest_hw_out_meter_db = Arc::new(vec![-90.0; hw_channels]);
2123
2124 let tracks: Vec<(String, crate::state::TrackHandle)> = self
2125 .state_snapshot
2126 .load_full()
2127 .tracks
2128 .iter()
2129 .map(|(name, track)| (name.clone(), track.clone()))
2130 .collect();
2131 self.track_meter_linear_by_track.clear();
2132 let mut snapshot = Vec::with_capacity(tracks.len());
2133 for (name, track) in tracks {
2134 let mut t = track.lock();
2135 t.clear_output_meters();
2136 let width = t.output_meter_linear().len();
2137 let zero_linear = vec![0.0; width];
2138 self.track_meter_linear_by_track
2139 .insert(name.clone(), zero_linear);
2140 snapshot.push((name, vec![-90.0; width]));
2141 }
2142 self.latest_track_meter_snapshot = Arc::new(snapshot);
2143 self.publish_meter_snapshot();
2144 }
2145
2146 pub(crate) fn publish_meter_snapshot_if_due(&mut self) {
2147 let now = Instant::now();
2148 if self
2149 .last_meter_snapshot_publish
2150 .is_some_and(|last| now.duration_since(last) < Self::METER_PUBLISH_INTERVAL)
2151 {
2152 return;
2153 }
2154 self.last_meter_snapshot_publish = Some(now);
2155 self.publish_meter_snapshot();
2156 }
2157
2158 pub(crate) fn publish_meter_snapshot(&mut self) {
2159 let snapshot = self.meter_snapshot_producer.write_buffer();
2160 snapshot.hw_out_db.clear();
2161 snapshot
2162 .hw_out_db
2163 .extend(self.latest_hw_out_meter_db.iter().copied());
2164 snapshot.track_meters.clear();
2165 snapshot
2166 .track_meters
2167 .extend(self.latest_track_meter_snapshot.iter().cloned());
2168 self.meter_snapshot_producer.publish();
2169 }
2170
2171 pub(crate) fn publish_transport_snapshot(&mut self) {
2172 let snapshot = self.transport_snapshot_producer.write_buffer();
2173 snapshot.sample = self.transport_sample;
2174 snapshot.tempo_bpm = self.tempo_bpm;
2175 snapshot.playing = self.playing;
2176 snapshot.transport_running = self.transport_running;
2177 snapshot.tsig_num = self.tsig_num;
2178 snapshot.tsig_denom = self.tsig_denom;
2179 self.transport_snapshot_producer.publish();
2180 }
2181
2182 pub(crate) async fn handle_request(&mut self, a: Action) {
2183 match a {
2184 Action::Log { source, message } => {
2185 self.notify_clients(Ok(Action::Log { source, message }))
2186 .await;
2187 }
2188 Action::Undo => {
2189 let actions = match self.history.undo() {
2190 Some(actions) => actions,
2191 None => {
2192 self.notify_clients(Ok(Action::Undo)).await;
2193 self.notify_clients(Ok(Action::HistoryState {
2194 dirty: self.history.is_dirty(),
2195 }))
2196 .await;
2197 return;
2198 }
2199 };
2200
2201 let was_suspended = self.history_suspended;
2202 self.history_suspended = true;
2203 for action in actions {
2204 self.handle_request_inner(action, false).await;
2205 }
2206 self.history_suspended = was_suspended;
2207 self.notify_clients(Ok(Action::Undo)).await;
2208 self.notify_clients(Ok(Action::HistoryState {
2209 dirty: self.history.is_dirty(),
2210 }))
2211 .await;
2212 }
2213 Action::Redo => {
2214 let actions = match self.history.redo() {
2215 Some(actions) => actions,
2216 None => {
2217 self.notify_clients(Ok(Action::Redo)).await;
2218 self.notify_clients(Ok(Action::HistoryState {
2219 dirty: self.history.is_dirty(),
2220 }))
2221 .await;
2222 return;
2223 }
2224 };
2225
2226 let was_suspended = self.history_suspended;
2227 self.history_suspended = true;
2228 for action in actions {
2229 self.handle_request_inner(action, false).await;
2230 }
2231 self.history_suspended = was_suspended;
2232 self.notify_clients(Ok(Action::Redo)).await;
2233 self.notify_clients(Ok(Action::HistoryState {
2234 dirty: self.history.is_dirty(),
2235 }))
2236 .await;
2237 }
2238 Action::ApplyGroupedActions(actions) => {
2239 self.handle_request_inner(Action::BeginHistoryGroup, true)
2240 .await;
2241 for action in actions {
2242 self.handle_request_inner(action, true).await;
2243 }
2244 self.handle_request_inner(Action::EndHistoryGroup, true)
2245 .await;
2246 }
2247 Action::Session(_) => {
2248 self.handle_request_inner(a, false).await;
2249 }
2250 other => {
2251 self.handle_request_inner(other, true).await;
2252 }
2253 }
2254 self.publish_state_snapshot();
2255 }
2256
2257 pub(crate) async fn handle_quit(&mut self, a: Action) {
2258 self.flush_recordings().await;
2259 if let Some(mut worker) = self.hw_worker.take() {
2267 let panic_events = self.panic_events_for_all_hw_midi_outputs();
2270 if !panic_events.is_empty() {
2271 let _ = worker.tx.send(Message::HWMidiOutEvents(panic_events)).await;
2272 }
2273 if let Err(e) = worker.tx.send(Message::Request(a.clone())).await {
2276 error!("Error sending quit message to HW worker: {e}");
2277 }
2278 if let Some(handle) = worker.handle.take() {
2279 handle
2280 .await
2281 .unwrap_or_else(|e| error!("Error waiting for HW worker to quit: {e}"));
2282 }
2283 }
2284 if let Some(hw) = self.hw_driver.as_mut() {
2290 hw.close_fds();
2291 }
2292 if let Some(midi_hub) = self.midi_hub.as_mut() {
2293 midi_hub.close_all();
2294 }
2295 self.hw_driver = None;
2296 self.hw_driver_info = None;
2297 self.hw_input_ports.clear();
2298 self.hw_output_ports.clear();
2299 self.notify_clients(Ok(Action::Quit)).await;
2300 self.ready_workers.clear();
2301 while !self.workers.is_empty() {
2302 let mut worker = self.workers.remove(0);
2303 if let Err(e) = worker.tx.send(Message::Request(a.clone())).await {
2304 error!("Error sending quit message to worker: {e}");
2305 }
2306 if let Some(handle) = worker.handle.take() {
2307 handle
2308 .await
2309 .unwrap_or_else(|e| error!("Error waiting for worker to quit: {e}"));
2310 }
2311 }
2312 #[cfg(unix)]
2313 {
2314 self.jack_runtime = None;
2315 }
2316 self.osc_server = None;
2317 }
2318
2319 #[inline]
2320 pub(crate) fn box_bool<'a>(
2321 fut: impl std::future::Future<Output = bool> + Send + 'a,
2322 ) -> std::pin::Pin<Box<dyn std::future::Future<Output = bool> + Send + 'a>> {
2323 Box::pin(fut)
2324 }
2325
2326 pub(crate) async fn handle_request_inner(
2327 &mut self,
2328 mut action_to_process: Action,
2329 record_history: bool,
2330 ) {
2331 let a = action_to_process.clone();
2332 let suppress_timing_history = self.playing
2333 && matches!(
2334 &action_to_process,
2335 Action::SetTempo(_) | Action::SetTimeSignature { .. } | Action::SetTempoMap { .. }
2336 );
2337 let mut inverse_actions = self.prepare_inverse_actions(
2338 &action_to_process,
2339 record_history,
2340 suppress_timing_history,
2341 );
2342
2343 match action_to_process {
2344 Action::Play => {
2345 if Self::box_bool(self.handle_play(a.clone())).await {
2346 return;
2347 }
2348 }
2349 Action::Pause => {
2350 if Self::box_bool(self.handle_pause(a.clone())).await {
2351 return;
2352 }
2353 }
2354 Action::Stop => {
2355 if Self::box_bool(self.handle_stop(a.clone())).await {
2356 return;
2357 }
2358 }
2359 Action::SessionPlay => {
2360 if Self::box_bool(self.handle_session_play(a.clone())).await {
2361 return;
2362 }
2363 }
2364 Action::JumpToEnd => {
2365 self.transport_sample = self.normalize_transport_sample(self.session_end_sample());
2366 self.publish_transport_snapshot();
2367 self.notify_clients(Ok(Action::TransportPosition(self.transport_sample)))
2368 .await;
2369 }
2370 Action::Panic => {
2371 if Self::box_bool(self.handle_panic(a.clone())).await {
2372 return;
2373 }
2374 }
2375 Action::Session(ref session_action) => {
2376 self.handle_session_action(session_action.clone()).await;
2377 }
2378 Action::SessionRuntimeReport { .. } => {}
2379 Action::SessionMidiLearnTriggered { .. } => {}
2380 Action::SetClipPlaybackEnabled(enabled) => {
2381 self.clip_playback_enabled = enabled;
2382 for track in self.state_snapshot.load_full().tracks.values() {
2383 track.lock().set_clip_playback_enabled(enabled);
2384 }
2385 }
2386 Action::SetSessionClipPlaybackEnabled(enabled) => {
2387 self.session_clip_playback_enabled = enabled;
2388 for track in self.state_snapshot.load_full().tracks.values() {
2389 track.lock().set_session_clip_playback_enabled(enabled);
2390 }
2391 }
2392 Action::TransportPosition(..) => {
2393 if Self::box_bool(self.handle_transport_position(a.clone())).await {
2394 return;
2395 }
2396 }
2397 Action::SetLoopEnabled(enabled) => {
2398 self.loop_enabled = enabled && self.loop_range_samples.is_some();
2399 self.notified_loop_wrap_sample = None;
2400 }
2401 Action::SetLoopRange(..) => {
2402 if Self::box_bool(self.handle_set_loop_range(a.clone())).await {
2403 return;
2404 }
2405 }
2406 Action::SetPunchEnabled(enabled) => {
2407 self.punch_enabled = enabled && self.punch_range_samples.is_some();
2408 }
2409 Action::SetPunchRange(range) => {
2410 self.punch_range_samples = range.and_then(|(start, end)| {
2411 if end > start {
2412 Some((start, end))
2413 } else {
2414 None
2415 }
2416 });
2417 self.punch_enabled = self.punch_range_samples.is_some();
2418 }
2419 Action::SetMetronomeEnabled(enabled) => {
2420 self.metronome_enabled = enabled;
2421 if enabled {
2422 self.ensure_metronome_track().await;
2423 }
2424 if let Some(track) = self
2425 .state_snapshot
2426 .load_full()
2427 .tracks
2428 .get(Self::METRONOME_TRACK)
2429 .cloned()
2430 {
2431 track.lock().set_metronome_enabled(enabled);
2432 }
2433 }
2434 Action::SetTempo(bpm) => {
2435 self.tempo_bpm = bpm.max(1.0);
2436 self.publish_transport_snapshot();
2437 }
2438 Action::SetTimeSignature {
2439 numerator,
2440 denominator,
2441 } => {
2442 self.tsig_num = numerator.max(1);
2443 self.tsig_denom = denominator.max(1);
2444 self.publish_transport_snapshot();
2445 }
2446 Action::SetTempoMap {
2447 ref tempo_points,
2448 ref time_signature_points,
2449 } => {
2450 self.tempo_points = tempo_points.clone();
2451 self.time_signature_points = time_signature_points.clone();
2452 self.update_global_tempo_from_map();
2453 self.publish_transport_snapshot();
2454 }
2455 Action::SetOscEnabled(enabled) => {
2456 if let Err(err) = self.set_osc_enabled_with(enabled, OscServer::start) {
2457 self.notify_clients(Err(err)).await;
2458 }
2459 }
2460 Action::SetRecordEnabled(..) => {
2461 if Self::box_bool(self.handle_set_record_enabled(a.clone())).await {
2462 return;
2463 }
2464 }
2465 Action::SetModulators(ref modulators) => {
2466 self.modulators = modulators.clone();
2467 let echoes = self.apply_modulators(self.active_transport_sample());
2468 for action in echoes {
2469 self.notify_clients(Ok(action)).await;
2470 }
2471 }
2472 Action::SetTrackAutomationLanes {
2473 ref track_name,
2474 ref lanes,
2475 mode,
2476 } => {
2477 if let Some(track) = self.state_snapshot.load_full().tracks.get(track_name) {
2478 let mut track = track.lock();
2479 track.automation_lanes = lanes.clone();
2480 track.set_automation_mode(mode);
2481 }
2482 }
2483 Action::TrackAutomationToggleLane { .. } => {
2484 if Self::box_bool(self.handle_track_automation_toggle_lane(a.clone())).await {
2485 return;
2486 }
2487 }
2488 Action::TrackAutomationInsertPoint { .. } => {
2489 if Self::box_bool(self.handle_track_automation_insert_point(a.clone())).await {
2490 return;
2491 }
2492 }
2493 Action::TrackAutomationDeletePoint { .. } => {
2494 if Self::box_bool(self.handle_track_automation_delete_point(a.clone())).await {
2495 return;
2496 }
2497 }
2498 Action::TrackAutomationSetMode {
2499 ref track_name,
2500 mode,
2501 } => {
2502 if let Some(track) = self
2503 .state_snapshot
2504 .load_full()
2505 .tracks
2506 .get(track_name)
2507 .cloned()
2508 {
2509 track.lock().set_automation_mode(mode);
2510 }
2511 }
2512 Action::RequestTrackList => {
2513 let names: Vec<String> = self
2514 .state_snapshot
2515 .load_full()
2516 .tracks
2517 .keys()
2518 .cloned()
2519 .collect();
2520 self.notify_clients(Ok(Action::TrackList(names))).await;
2521 }
2522 Action::TrackList(_) => {}
2523 Action::RequestTransportState => {
2524 self.notify_clients(Ok(Action::TransportState {
2525 sample: self.transport_sample,
2526 tempo_bpm: self.tempo_bpm,
2527 playing: self.playing,
2528 paused: !self.transport_running && self.playing,
2529 tsig_num: self.tsig_num,
2530 tsig_denom: self.tsig_denom,
2531 }))
2532 .await;
2533 }
2534 Action::TransportState { .. } => {}
2535 Action::SetStepRecording(enabled) => {
2536 self.step_recording_enabled = enabled;
2537 }
2538 Action::BeginHistoryGroup if self.history_group.is_none() => {
2539 self.history_group = Some(UndoEntry {
2540 forward_actions: vec![],
2541 inverse_actions: vec![],
2542 });
2543 }
2544 Action::EndHistoryGroup => {
2545 if let Some(mut group) = self.history_group.take()
2546 && !group.forward_actions.is_empty()
2547 && !group.inverse_actions.is_empty()
2548 {
2549 let mut add_tracks = Vec::new();
2550 let mut connections = Vec::new();
2551 let mut rest = Vec::new();
2552 for action in group.inverse_actions {
2553 if matches!(action, Action::AddTrack { .. }) {
2554 add_tracks.push(action);
2555 } else if matches!(action, Action::Connect { .. }) {
2556 connections.push(action);
2557 } else {
2558 rest.push(action);
2559 }
2560 }
2561 group.inverse_actions = add_tracks;
2562 group.inverse_actions.extend(rest);
2563 group.inverse_actions.extend(connections);
2564 self.history.record(group);
2565 }
2566 }
2567 Action::SetSessionPath(ref path) => {
2568 self.session_dir = Some(Path::new(path).to_path_buf());
2569 self.ensure_session_subdirs();
2570 #[cfg(all(unix, not(target_os = "macos")))]
2571 let _lv2_dir = self.session_plugins_dir();
2572 for track in self.state_snapshot.load_full().tracks.values() {
2573 track.lock().set_session_base_dir(self.session_dir.clone());
2574 }
2575 }
2576 Action::MarkHistorySavePoint => {
2577 self.history.mark_save_point();
2578 self.notify_clients(Ok(Action::HistoryState {
2579 dirty: self.history.is_dirty(),
2580 }))
2581 .await;
2582 }
2583 Action::ClearHistory => {
2584 self.history.clear();
2585 self.history.mark_save_point();
2586 }
2587 Action::BeginSessionRestore => {
2588 self.history_suspended = true;
2589 self.history.clear();
2590 }
2591 Action::EndSessionRestore => {
2592 self.history.clear();
2593 self.history_suspended = false;
2594 self.preload_track_clips_spawn();
2595 }
2596 Action::Quit => {
2597 self.handle_quit(a.clone()).await;
2598 return;
2599 }
2600 Action::AddTrack {
2601 ref name,
2602 audio_ins,
2603 midi_ins,
2604 audio_outs,
2605 midi_outs,
2606 folder,
2607 } => {
2608 self.handle_add_track(
2609 name.clone(),
2610 audio_ins,
2611 midi_ins,
2612 audio_outs,
2613 midi_outs,
2614 folder,
2615 )
2616 .await;
2617 }
2618 Action::TrackAddAudioInput(..) => {
2619 if Self::box_bool(self.handle_track_add_audio_input(a.clone())).await {
2620 return;
2621 }
2622 }
2623 Action::TrackAddAudioOutput(..) => {
2624 if Self::box_bool(self.handle_track_add_audio_output(a.clone())).await {
2625 return;
2626 }
2627 }
2628 Action::TrackRemoveAudioInput(..) => {
2629 if Self::box_bool(self.handle_track_remove_audio_input(a.clone())).await {
2630 return;
2631 }
2632 }
2633 Action::TrackRemoveAudioOutput(..) => {
2634 if Self::box_bool(self.handle_track_remove_audio_output(a.clone())).await {
2635 return;
2636 }
2637 }
2638 Action::RenameTrack { .. } => {
2639 if Self::box_bool(self.handle_rename_track(a.clone())).await {
2640 return;
2641 }
2642 }
2643 Action::RemoveTrack(ref name) => {
2644 self.handle_remove_track(name.clone(), record_history).await;
2645 inverse_actions = None;
2646 }
2647 Action::TrackLevel(ref name, level) => {
2648 if name == "hw:out" {
2649 self.hw_out_level_db = level;
2650 } else if let Some(track) = self.state_snapshot.load_full().tracks.get(name) {
2651 track.lock().set_level(level);
2652 }
2653 }
2654 Action::TrackBalance(ref name, balance) => {
2655 if name == "hw:out" {
2656 self.hw_out_balance = balance.clamp(-1.0, 1.0);
2657 } else if let Some(track) = self.state_snapshot.load_full().tracks.get(name) {
2658 track.lock().set_balance(balance);
2659 }
2660 }
2661 Action::TrackAutomationLevel(ref name, level) => {
2662 tracing::debug!(%name, level, "engine received TrackAutomationLevel");
2663 if name == "hw:out" {
2664 self.hw_out_level_db = level;
2665 } else if let Some(track) = self.state_snapshot.load_full().tracks.get(name) {
2666 track.lock().set_level(level);
2667 }
2668 }
2669 Action::TrackAutomationBalance(ref name, balance) => {
2670 if name == "hw:out" {
2671 self.hw_out_balance = balance.clamp(-1.0, 1.0);
2672 } else if let Some(track) = self.state_snapshot.load_full().tracks.get(name) {
2673 track.lock().set_balance(balance);
2674 }
2675 }
2676 Action::TrackMidiCc { .. } => {
2677 if Self::box_bool(self.handle_track_midi_cc(a.clone())).await {
2678 return;
2679 }
2680 }
2681 Action::RequestMeterSnapshot => {
2682 self.notify_clients(Ok(Action::MeterSnapshot {
2683 hw_out_db: self.latest_hw_out_meter_db.clone(),
2684 track_meters: self.latest_track_meter_snapshot.clone(),
2685 }))
2686 .await;
2687 return;
2688 }
2689 Action::TrackMeters { .. } => {}
2690 Action::MeterSnapshot { .. } => {}
2691 Action::TrackToggleArm(..) => {
2692 if Self::box_bool(self.handle_track_toggle_arm(a.clone())).await {
2693 return;
2694 }
2695 }
2696 Action::TrackToggleMute(ref name) => {
2697 if name == "hw:out" {
2698 self.hw_out_muted = !self.hw_out_muted;
2699 } else if let Some(track) = self.state_snapshot.load_full().tracks.get(name) {
2700 track.lock().mute();
2701 }
2702 }
2703 Action::TrackTogglePhase(ref name) => {
2704 if let Some(track) = self.state_snapshot.load_full().tracks.get(name) {
2705 track.lock().invert_phase();
2706 }
2707 }
2708 Action::TrackToggleSolo(ref name) => {
2709 if name == "hw:out" {
2710 return;
2711 }
2712 if let Some(track) = self.state_snapshot.load_full().tracks.get(name) {
2713 track.lock().solo();
2714 }
2715 }
2716 Action::TrackToggleMaster(ref name) => {
2717 if let Some(track) = self.state_snapshot.load_full().tracks.get(name) {
2718 track.lock().toggle_master();
2719 }
2720 }
2721 Action::TrackToggleInputMonitor {
2722 ref track_name,
2723 lane,
2724 } => {
2725 if let Some(track) = self.state_snapshot.load_full().tracks.get(track_name) {
2726 track.lock().toggle_input_monitor(lane);
2727 }
2728 }
2729 Action::TrackToggleDiskMonitor {
2730 ref track_name,
2731 lane,
2732 } => {
2733 if let Some(track) = self.state_snapshot.load_full().tracks.get(track_name) {
2734 track.lock().toggle_disk_monitor(lane);
2735 }
2736 }
2737 Action::TrackToggleMidiInputMonitor {
2738 ref track_name,
2739 lane,
2740 } => {
2741 if let Some(track) = self.state_snapshot.load_full().tracks.get(track_name) {
2742 track.lock().toggle_midi_input_monitor(lane);
2743 }
2744 }
2745 Action::TrackToggleMidiDiskMonitor {
2746 ref track_name,
2747 lane,
2748 } => {
2749 if let Some(track) = self.state_snapshot.load_full().tracks.get(track_name) {
2750 track.lock().toggle_midi_disk_monitor(lane);
2751 }
2752 }
2753 Action::TrackSetColor {
2754 ref track_name,
2755 color,
2756 } => {
2757 if let Some(track) = self.state_snapshot.load_full().tracks.get(track_name) {
2758 track.lock().color = color;
2759 }
2760 }
2761 Action::TrackArmMidiLearn {
2762 ref track_name,
2763 target,
2764 } => {
2765 if let Err(e) = self.track_handle_or_err(track_name) {
2766 self.notify_clients(Err(e)).await;
2767 return;
2768 }
2769 self.pending_midi_learn = Some((track_name.clone(), target, None));
2770 }
2771 Action::GlobalArmMidiLearn { target } => {
2772 self.pending_global_midi_learn = Some(target);
2773 }
2774 Action::SessionArmMidiLearn { ref target } => {
2775 self.pending_session_midi_learn = Some(target.clone());
2776 }
2777 Action::TrackSetMidiLearnBinding { .. } => {
2778 if Self::box_bool(self.handle_track_set_midi_learn_binding(a.clone())).await {
2779 return;
2780 }
2781 }
2782 Action::SetGlobalMidiLearnBinding { .. } => {
2783 if Self::box_bool(self.handle_set_global_midi_learn_binding(a.clone())).await {
2784 return;
2785 }
2786 }
2787 Action::SetSessionMidiLearnBinding { .. } => {
2788 if Self::box_bool(self.handle_set_session_midi_learn_binding(a.clone())).await {
2789 return;
2790 }
2791 }
2792 Action::TrackSetFolder { .. } => {
2793 if Self::box_bool(self.handle_track_set_folder(a.clone())).await {
2794 return;
2795 }
2796 }
2797 Action::TrackSetParent {
2798 ref track_name,
2799 ref parent_name,
2800 } => {
2801 self.handle_track_set_parent(track_name.as_str(), parent_name.as_deref())
2802 .await;
2803 }
2804 Action::TrackToggleFolder { .. } => {
2805 if Self::box_bool(self.handle_track_toggle_folder(a.clone())).await {
2806 return;
2807 }
2808 }
2809 Action::TrackSetMidiLaneChannel { .. } => {
2810 if Self::box_bool(self.handle_track_set_midi_lane_channel(a.clone())).await {
2811 return;
2812 }
2813 }
2814 Action::TrackSetFrozen { .. } => {
2815 if Self::box_bool(self.handle_track_set_frozen(a.clone())).await {
2816 return;
2817 }
2818 }
2819 Action::TrackSetSessionSlot { .. } => {
2820 if Self::box_bool(self.handle_track_set_session_slot(a.clone())).await {
2821 return;
2822 }
2823 }
2824 Action::TrackSetSessionSlotPlayEnabled { .. } => {
2825 if self
2826 .handle_track_set_session_slot_play_enabled(a.clone())
2827 .await
2828 {
2829 return;
2830 }
2831 }
2832 Action::TrackSetSessionSlotStopEnabled { .. } => {
2833 if self
2834 .handle_track_set_session_slot_stop_enabled(a.clone())
2835 .await
2836 {
2837 return;
2838 }
2839 }
2840 Action::TrackOfflineBounce { .. } => {
2841 self.handle_track_offline_bounce(action_to_process).await;
2842 return;
2843 }
2844 Action::TrackOfflineBounceCancel { .. } => {}
2845 Action::TrackOfflineBounceCancelAll => {}
2846 Action::TrackOfflineBounceCanceled { .. } => {}
2847 Action::TrackOfflineBounceProgress { .. } => {}
2848 Action::PianoKey {
2849 ref track_name,
2850 note,
2851 velocity,
2852 on,
2853 } => {
2854 if let Some(track) = self.state_snapshot.load_full().tracks.get(track_name) {
2855 let status = if on { 0x90 } else { 0x80 };
2856 let event = MidiEvent::new(0, vec![status, note.min(127), velocity.min(127)]);
2857 track.lock().push_hw_midi_events(&[event]);
2858 }
2859 }
2860 Action::ModifyMidiNotes { .. }
2861 | Action::ModifyMidiControllers { .. }
2862 | Action::DeleteMidiControllers { .. }
2863 | Action::InsertMidiControllers { .. }
2864 | Action::DeleteMidiNotes { .. }
2865 | Action::InsertMidiNotes { .. } => {
2866 if let Err(e) = self.apply_midi_edit_action(&action_to_process) {
2867 self.notify_clients(Err(e)).await;
2868 return;
2869 }
2870 }
2871 Action::SetMidiSysExEvents { .. } => {
2872 if let Err(e) = self.apply_midi_edit_action(&action_to_process) {
2873 self.notify_clients(Err(e)).await;
2874 return;
2875 }
2876 }
2877 Action::TrackClearDefaultPassthrough { .. } => {
2878 if Self::box_bool(self.handle_track_clear_default_passthrough(a.clone())).await {
2879 return;
2880 }
2881 }
2882 Action::TrackClearPlugins { .. } => {
2883 if Self::box_bool(self.handle_track_clear_plugins(a.clone())).await {
2884 return;
2885 }
2886 }
2887 #[cfg(all(unix, not(target_os = "macos")))]
2888 Action::ClipSetLv2PluginState { ref track_name, .. } => {
2889 self.notify_clients(Err(format!(
2890 "Track '{}': clip LV2 plugin state changes are not supported",
2891 track_name
2892 )))
2893 .await;
2894 }
2895 Action::TrackGetClapNoteNames { .. } => {
2896 if Self::box_bool(self.handle_track_get_clap_note_names(a.clone())).await {
2897 return;
2898 }
2899 }
2900 #[cfg(all(unix, not(target_os = "macos")))]
2901 Action::TrackGetLv2Midnam { .. } => {
2902 if Self::box_bool(self.handle_track_get_lv2_midnam(a.clone())).await {
2903 return;
2904 }
2905 }
2906 Action::TrackGetPluginGraph { .. } => {
2907 if Self::box_bool(self.handle_track_get_plugin_graph(a.clone())).await {
2908 return;
2909 }
2910 }
2911 Action::TrackPluginGraph { .. } => {}
2912 Action::TrackConnectPluginAudio { .. } => {
2913 if Self::box_bool(self.handle_track_connect_plugin_audio(a.clone())).await {
2914 return;
2915 }
2916 }
2917 Action::TrackConnectPluginMidi { .. } => {
2918 if Self::box_bool(self.handle_track_connect_plugin_midi(a.clone())).await {
2919 return;
2920 }
2921 }
2922 Action::TrackDisconnectPluginAudio { .. } => {
2923 if Self::box_bool(self.handle_track_disconnect_plugin_audio(a.clone())).await {
2924 return;
2925 }
2926 }
2927 Action::TrackDisconnectPluginMidi { .. } => {
2928 if Self::box_bool(self.handle_track_disconnect_plugin_midi(a.clone())).await {
2929 return;
2930 }
2931 }
2932 Action::TrackConnectAudio { .. } => {
2933 if Self::box_bool(self.handle_track_connect_audio(a.clone())).await {
2934 return;
2935 }
2936 }
2937 Action::TrackDisconnectAudio { .. } => {
2938 if Self::box_bool(self.handle_track_disconnect_audio(a.clone())).await {
2939 return;
2940 }
2941 }
2942 Action::TrackConnectMidi { .. } => {
2943 if Self::box_bool(self.handle_track_connect_midi(a.clone())).await {
2944 return;
2945 }
2946 }
2947 Action::TrackDisconnectMidi { .. } => {
2948 if Self::box_bool(self.handle_track_disconnect_midi(a.clone())).await {
2949 return;
2950 }
2951 }
2952 #[cfg(all(unix, not(target_os = "macos")))]
2953 Action::ListLv2Plugins => {
2954 if Self::box_bool(self.handle_list_lv2_plugins(a.clone())).await {
2955 return;
2956 }
2957 }
2958 #[cfg(all(unix, not(target_os = "macos")))]
2959 Action::Lv2Plugins(_) => {}
2960 #[cfg(all(unix, not(target_os = "macos")))]
2961 Action::Lv2PluginsUnavailable { .. } => {}
2962 Action::ListVst3Plugins => {
2963 if Self::box_bool(self.handle_list_vst3_plugins(a.clone())).await {
2964 return;
2965 }
2966 }
2967 Action::Vst3Plugins(_) => {}
2968 Action::Vst3PluginsUnavailable { .. } => {}
2969 Action::ListClapPlugins => {
2970 if Self::box_bool(self.handle_list_clap_plugins(a.clone())).await {
2971 return;
2972 }
2973 }
2974 Action::ListClapPluginsWithCapabilities => {
2975 if self
2976 .handle_list_clap_plugins_with_capabilities(a.clone())
2977 .await
2978 {
2979 return;
2980 }
2981 }
2982 Action::ClapPlugins(_) => {}
2983 Action::ClapPluginsUnavailable { .. } => {}
2984 Action::TrackLoadClapPlugin {
2985 ref track_name,
2986 ref plugin_id,
2987 instance_id,
2988 } => {
2989 if self
2990 .handle_track_load_clap_plugin(
2991 track_name.as_str(),
2992 plugin_id.as_str(),
2993 instance_id,
2994 )
2995 .await
2996 {
2997 return;
2998 }
2999 }
3000 Action::TrackUnloadClapPlugin {
3001 ref track_name,
3002 ref plugin_id,
3003 } => {
3004 if self
3005 .handle_track_unload_clap_plugin(track_name.as_str(), plugin_id.as_str())
3006 .await
3007 {
3008 return;
3009 }
3010 }
3011 Action::TrackUnloadClapPluginInstance {
3012 ref track_name,
3013 instance_id,
3014 } => {
3015 if self
3016 .handle_track_unload_clap_plugin_instance(track_name.as_str(), instance_id)
3017 .await
3018 {
3019 return;
3020 }
3021 }
3022 Action::TrackShowClapGui { .. } => {
3023 if Self::box_bool(self.handle_track_show_clap_gui(a.clone())).await {
3024 return;
3025 }
3026 }
3027 Action::TrackLoadVst3Plugin {
3028 ref track_name,
3029 ref plugin_id,
3030 instance_id,
3031 } => {
3032 if self
3033 .handle_track_load_vst3_plugin(
3034 track_name.as_str(),
3035 plugin_id.as_str(),
3036 instance_id,
3037 )
3038 .await
3039 {
3040 return;
3041 }
3042 }
3043 Action::TrackUnloadVst3Plugin {
3044 ref track_name,
3045 ref plugin_id,
3046 } => {
3047 if self
3048 .handle_track_unload_vst3_plugin(track_name.as_str(), plugin_id.as_str())
3049 .await
3050 {
3051 return;
3052 }
3053 }
3054 Action::TrackUnloadVst3PluginInstance {
3055 ref track_name,
3056 instance_id,
3057 } => {
3058 if self
3059 .handle_track_unload_vst3_plugin_instance(track_name.as_str(), instance_id)
3060 .await
3061 {
3062 return;
3063 }
3064 }
3065 Action::TrackShowVst3Gui { .. } => {
3066 if Self::box_bool(self.handle_track_show_vst3_gui(a.clone())).await {
3067 return;
3068 }
3069 }
3070 #[cfg(all(unix, not(target_os = "macos")))]
3071 Action::TrackLoadLv2Plugin {
3072 ref track_name,
3073 ref plugin_uri,
3074 instance_id,
3075 } => {
3076 if self
3077 .handle_track_load_lv2_plugin(
3078 track_name.as_str(),
3079 plugin_uri.as_str(),
3080 instance_id,
3081 )
3082 .await
3083 {
3084 return;
3085 }
3086 }
3087 #[cfg(all(unix, not(target_os = "macos")))]
3088 Action::TrackUnloadLv2Plugin {
3089 ref track_name,
3090 ref plugin_uri,
3091 } => {
3092 if self
3093 .handle_track_unload_lv2_plugin(track_name.as_str(), plugin_uri.as_str())
3094 .await
3095 {
3096 return;
3097 }
3098 }
3099 #[cfg(all(unix, not(target_os = "macos")))]
3100 Action::TrackUnloadLv2PluginInstance {
3101 ref track_name,
3102 instance_id,
3103 } => {
3104 if self
3105 .handle_track_unload_lv2_plugin_instance(track_name.as_str(), instance_id)
3106 .await
3107 {
3108 return;
3109 }
3110 }
3111 #[cfg(all(unix, not(target_os = "macos")))]
3112 Action::TrackShowLv2Gui { .. } => {
3113 if Self::box_bool(self.handle_track_show_lv2_gui(a.clone())).await {
3114 return;
3115 }
3116 }
3117 Action::TrackSetPluginResourceDir { .. } => {
3118 if Self::box_bool(self.handle_track_set_plugin_resource_dir(a.clone())).await {
3119 return;
3120 }
3121 }
3122 Action::TrackClapFileReferences { .. } => {
3123 if Self::box_bool(self.handle_track_clap_file_references(a.clone())).await {
3124 return;
3125 }
3126 }
3127 Action::TrackUpdateClapFileReference { .. } => {
3128 if self
3129 .handle_track_update_clap_file_reference(a.clone())
3130 .await
3131 {
3132 return;
3133 }
3134 }
3135 Action::ClipSetPluginResourceDir { .. } => {
3136 if Self::box_bool(self.handle_clip_set_plugin_resource_dir(a.clone())).await {
3137 return;
3138 }
3139 }
3140 Action::ClipClapFileReferences { .. } => {
3141 if Self::box_bool(self.handle_clip_clap_file_references(a.clone())).await {
3142 return;
3143 }
3144 }
3145 Action::ClipUpdateClapFileReference { .. } => {
3146 if Self::box_bool(self.handle_clip_update_clap_file_reference(a.clone())).await {
3147 return;
3148 }
3149 }
3150 Action::TrackSetClapParameter { .. } => {
3151 if Self::box_bool(self.handle_track_set_clap_parameter(a.clone())).await {
3152 return;
3153 }
3154 }
3155 Action::ClipSetClapParameter { .. } => {
3156 if Self::box_bool(self.handle_clip_set_clap_parameter(a.clone())).await {
3157 return;
3158 }
3159 }
3160 Action::TrackSetClapParameterAt { .. } => {
3161 if Self::box_bool(self.handle_track_set_clap_parameter_at(a.clone())).await {
3162 return;
3163 }
3164 }
3165 Action::TrackBeginClapParameterEdit { .. } => {
3166 if Self::box_bool(self.handle_track_begin_clap_parameter_edit(a.clone())).await {
3167 return;
3168 }
3169 }
3170 Action::TrackEndClapParameterEdit { .. } => {
3171 if Self::box_bool(self.handle_track_end_clap_parameter_edit(a.clone())).await {
3172 return;
3173 }
3174 }
3175 Action::TrackGetClapParameters { .. } => {
3176 if Self::box_bool(self.handle_track_get_clap_parameters(a.clone())).await {
3177 return;
3178 }
3179 }
3180 Action::TrackClapParameters { .. } => {}
3181 Action::TrackClapSnapshotState { .. } => {
3182 if Self::box_bool(self.handle_track_clap_snapshot_state(a.clone())).await {
3183 return;
3184 }
3185 }
3186 Action::ClipClapSnapshotState { .. } => {
3187 if Self::box_bool(self.handle_clip_clap_snapshot_state(a.clone())).await {
3188 return;
3189 }
3190 }
3191 Action::TrackClapStateSnapshot { .. } => {}
3192 Action::ClipClapStateSnapshot { .. } => {}
3193 Action::TrackClapStateDirty { .. } => {}
3194 Action::ClipClapStateDirty { .. } => {}
3195 Action::TrackClapRestoreState { .. } => {
3196 if Self::box_bool(self.handle_track_clap_restore_state(a.clone())).await {
3197 return;
3198 }
3199 }
3200 Action::ClipClapRestoreState { .. } => {
3201 if Self::box_bool(self.handle_clip_clap_restore_state(a.clone())).await {
3202 return;
3203 }
3204 }
3205 Action::TrackSnapshotAllClapStates { .. } => {
3206 if Self::box_bool(self.handle_track_snapshot_all_clap_states(a.clone())).await {
3207 return;
3208 }
3209 }
3210 Action::TrackSnapshotAllClapStatesDone { .. } => {}
3211 Action::TrackGetVst3Graph { .. } => {
3212 if Self::box_bool(self.handle_track_get_vst3_graph(a.clone())).await {
3213 return;
3214 }
3215 }
3216 Action::TrackVst3Graph { .. } => {}
3217 Action::TrackSetVst3Parameter { .. } => {
3218 if Self::box_bool(self.handle_track_set_vst3_parameter(a.clone())).await {
3219 return;
3220 }
3221 }
3222 Action::TrackSetPluginBypassed { .. } => {
3223 if Self::box_bool(self.handle_track_set_plugin_bypassed(a.clone())).await {
3224 return;
3225 }
3226 }
3227 Action::TrackGetVst3Parameters { .. } => {
3228 if Self::box_bool(self.handle_track_get_vst3_parameters(a.clone())).await {
3229 return;
3230 }
3231 }
3232 Action::TrackVst3Parameters { .. } => {}
3233 #[cfg(all(unix, not(target_os = "macos")))]
3234 Action::TrackGetLv2PluginControls { .. } => {
3235 if Self::box_bool(self.handle_track_get_lv2_plugin_controls(a.clone())).await {
3236 return;
3237 }
3238 }
3239 #[cfg(all(unix, not(target_os = "macos")))]
3240 Action::TrackLv2SnapshotState { .. } => {
3241 if Self::box_bool(self.handle_track_lv2_snapshot_state(a.clone())).await {
3242 return;
3243 }
3244 }
3245 #[cfg(all(unix, not(target_os = "macos")))]
3246 Action::ClipLv2SnapshotState { .. } => {
3247 if Self::box_bool(self.handle_clip_lv2_snapshot_state(a.clone())).await {
3248 return;
3249 }
3250 }
3251 Action::TrackVst3SnapshotState { .. } => {
3252 if Self::box_bool(self.handle_track_vst3_snapshot_state(a.clone())).await {
3253 return;
3254 }
3255 }
3256 Action::ClipVst3SnapshotState { .. } => {
3257 if Self::box_bool(self.handle_clip_vst3_snapshot_state(a.clone())).await {
3258 return;
3259 }
3260 }
3261 Action::TrackVst3StateSnapshot { .. } => {}
3262 Action::ClipVst3StateSnapshot { .. } => {}
3263 Action::TrackVst3RestoreState { .. } => {
3264 if Self::box_bool(self.handle_track_vst3_restore_state(a.clone())).await {
3265 return;
3266 }
3267 }
3268 Action::TrackConnectVst3Audio { .. } => {
3269 if Self::box_bool(self.handle_track_connect_vst3_audio(a.clone())).await {
3270 return;
3271 }
3272 }
3273 Action::TrackDisconnectVst3Audio { .. } => {
3274 if Self::box_bool(self.handle_track_disconnect_vst3_audio(a.clone())).await {
3275 return;
3276 }
3277 }
3278 Action::ClipMove { .. } => {
3279 self.handle_clip_move(a.clone()).await;
3280 }
3281 Action::AddClip { .. } => {
3282 if Self::box_bool(self.handle_add_clip(a.clone())).await {
3283 return;
3284 }
3285 }
3286 Action::AddGroupedClip { .. } => {
3287 if Self::box_bool(self.handle_add_grouped_clip(a.clone())).await {
3288 return;
3289 }
3290 }
3291 Action::RemoveClip {
3292 ref track_name,
3293 kind,
3294 ref clip_indices,
3295 } => {
3296 self.remove_clips_from_track(track_name, kind, clip_indices);
3297 }
3298 Action::MoveClipToUnused {
3299 ref track_name,
3300 kind,
3301 ref clip_indices,
3302 } => {
3303 self.move_clips_to_unused(track_name, kind, clip_indices);
3304 }
3305 Action::DeleteUnusedClips { ref clip_ids } => {
3306 self.delete_unused_clips(clip_ids);
3307 }
3308 Action::SetUnusedClips {
3309 ref audio,
3310 ref midi,
3311 } => {
3312 self.set_unused_clips(audio.clone(), midi.clone());
3313 }
3314 Action::RenameClip {
3315 ref track_name,
3316 kind,
3317 clip_index,
3318 ref new_name,
3319 } => {
3320 self.rename_clip_references(track_name, kind, clip_index, new_name);
3321 }
3322 Action::SetClipIdentity {
3323 ref track_name,
3324 kind,
3325 clip_index,
3326 ref new_id,
3327 ref new_name,
3328 } => {
3329 self.set_clip_identity(track_name, kind, clip_index, new_id, new_name);
3330 }
3331 Action::SetClipSourceName {
3332 ref track_name,
3333 kind,
3334 clip_index,
3335 ref name,
3336 } => {
3337 self.set_clip_source_name(track_name, clip_index, kind, name.clone());
3338 }
3339 Action::SetClipFade { .. } => {
3340 if Self::box_bool(self.handle_set_clip_fade(a.clone())).await {
3341 return;
3342 }
3343 }
3344 Action::SetClipBounds {
3345 ref track_name,
3346 clip_index,
3347 kind,
3348 start,
3349 length,
3350 offset,
3351 } => {
3352 self.set_clip_bounds(track_name, clip_index, kind, start, length, offset);
3353 }
3354 Action::SyncClipBounds {
3355 ref track_name,
3356 clip_index,
3357 kind,
3358 start,
3359 length,
3360 offset,
3361 } => {
3362 self.set_clip_bounds(track_name, clip_index, kind, start, length, offset);
3363 }
3364 Action::SetClipMuted {
3365 ref track_name,
3366 clip_index,
3367 kind,
3368 muted,
3369 } => {
3370 self.set_clip_muted(track_name, clip_index, kind, muted);
3371 }
3372 Action::SetClipPluginGraphJson {
3373 ref track_name,
3374 clip_index,
3375 ref plugin_graph_json,
3376 } => {
3377 self.set_clip_plugin_graph_json(track_name, clip_index, plugin_graph_json.clone());
3378 }
3379 Action::SetClipPitchCorrection { .. } => {
3380 if Self::box_bool(self.handle_set_clip_pitch_correction(a.clone())).await {
3381 return;
3382 }
3383 }
3384 Action::Connect {
3385 ref from_track,
3386 from_port,
3387 ref to_track,
3388 to_port,
3389 kind,
3390 } => {
3391 self.handle_connect(
3392 from_track.as_str(),
3393 from_port,
3394 to_track.as_str(),
3395 to_port,
3396 kind,
3397 )
3398 .await;
3399 }
3400 Action::Disconnect { .. } => {
3401 self.handle_disconnect(a.clone()).await;
3402 }
3403 Action::OpenAudioDevice { .. } => {
3404 let (done, updated) = self.handle_open_audio_device(a.clone()).await;
3405 if done {
3406 return;
3407 }
3408 if let Some(action) = updated {
3409 action_to_process = action;
3410 }
3411 }
3412 Action::JackAddAudioInputPort => {
3413 if Self::box_bool(self.handle_jack_add_audio_input_port(a.clone())).await {
3414 return;
3415 }
3416 }
3417 Action::JackRemoveAudioInputPort(_removed_port) => {
3418 if self
3419 .handle_jack_remove_audio_input_port(_removed_port, a.clone())
3420 .await
3421 {
3422 return;
3423 }
3424 }
3425 Action::JackAddAudioOutputPort => {
3426 if Self::box_bool(self.handle_jack_add_audio_output_port(a.clone())).await {
3427 return;
3428 }
3429 }
3430 Action::JackRemoveAudioOutputPort(_removed_port) => {
3431 if self
3432 .handle_jack_remove_audio_output_port(_removed_port, a.clone())
3433 .await
3434 {
3435 return;
3436 }
3437 }
3438 Action::JackGetGraph => {
3439 #[cfg(unix)]
3440 {
3441 match self
3442 .jack_runtime
3443 .as_ref()
3444 .ok_or(
3445 "JACK runtime is not active; open the JACK backend first".to_string(),
3446 )
3447 .and_then(|jack| jack.graph_info())
3448 {
3449 Ok(graph) => self.notify_clients(Ok(Action::JackGraph(graph))).await,
3450 Err(e) => self.notify_clients(Err(e)).await,
3451 }
3452 }
3453 #[cfg(not(unix))]
3454 {
3455 self.notify_clients(Err(
3456 "JACK backend is not available on this platform build".to_string(),
3457 ))
3458 .await;
3459 }
3460 return;
3461 }
3462 Action::JackConnect {
3463 ref source,
3464 ref destination,
3465 } => {
3466 #[cfg(unix)]
3467 {
3468 match self
3469 .jack_runtime
3470 .as_ref()
3471 .ok_or(
3472 "JACK runtime is not active; open the JACK backend first".to_string(),
3473 )
3474 .and_then(|jack| jack.connect_ports_by_name(source, destination))
3475 .and_then(|_| {
3476 self.jack_runtime
3477 .as_ref()
3478 .expect("JACK runtime was checked")
3479 .graph_info()
3480 }) {
3481 Ok(graph) => {
3482 self.notify_clients(Ok(a.clone())).await;
3483 self.notify_clients(Ok(Action::JackGraph(graph))).await;
3484 }
3485 Err(e) => self.notify_clients(Err(e)).await,
3486 }
3487 }
3488 #[cfg(not(unix))]
3489 {
3490 self.notify_clients(Err(
3491 "JACK backend is not available on this platform build".to_string(),
3492 ))
3493 .await;
3494 }
3495 return;
3496 }
3497 Action::JackDisconnect {
3498 ref source,
3499 ref destination,
3500 } => {
3501 #[cfg(unix)]
3502 {
3503 match self
3504 .jack_runtime
3505 .as_ref()
3506 .ok_or(
3507 "JACK runtime is not active; open the JACK backend first".to_string(),
3508 )
3509 .and_then(|jack| jack.disconnect_ports_by_name(source, destination))
3510 .and_then(|_| {
3511 self.jack_runtime
3512 .as_ref()
3513 .expect("JACK runtime was checked")
3514 .graph_info()
3515 }) {
3516 Ok(graph) => {
3517 self.notify_clients(Ok(a.clone())).await;
3518 self.notify_clients(Ok(Action::JackGraph(graph))).await;
3519 }
3520 Err(e) => self.notify_clients(Err(e)).await,
3521 }
3522 }
3523 #[cfg(not(unix))]
3524 {
3525 self.notify_clients(Err(
3526 "JACK backend is not available on this platform build".to_string(),
3527 ))
3528 .await;
3529 }
3530 return;
3531 }
3532 Action::OpenMidiInputDevice(ref device) => {
3533 if let Some(worker) = &self.hw_worker {
3534 if let Err(e) = worker
3535 .tx
3536 .send(Message::HWOpenMidiInputDevice(device.clone()))
3537 .await
3538 {
3539 self.notify_clients(Err(format!("Failed to send MIDI input open: {e}")))
3540 .await;
3541 }
3542 return;
3543 }
3544 let Some(midi_hub) = self.midi_hub.as_mut() else {
3545 self.notify_clients(Err("Hardware MIDI hub is not available".to_string()))
3546 .await;
3547 return;
3548 };
3549 if let Err(e) = midi_hub.open_input(device) {
3550 self.notify_clients(Err(e)).await;
3551 return;
3552 }
3553 }
3554 Action::OpenMidiOutputDevice(ref device) => {
3555 if let Some(worker) = &self.hw_worker {
3556 if let Err(e) = worker
3557 .tx
3558 .send(Message::HWOpenMidiOutputDevice(device.clone()))
3559 .await
3560 {
3561 self.notify_clients(Err(format!("Failed to send MIDI output open: {e}")))
3562 .await;
3563 }
3564 return;
3565 }
3566 let Some(midi_hub) = self.midi_hub.as_mut() else {
3567 self.notify_clients(Err("Hardware MIDI hub is not available".to_string()))
3568 .await;
3569 return;
3570 };
3571 if let Err(e) = midi_hub.open_output(device) {
3572 self.notify_clients(Err(e)).await;
3573 return;
3574 }
3575 }
3576 Action::RequestSessionDiagnostics => {
3577 self.handle_request_session_diagnostics().await;
3578 }
3579 Action::RequestMidiLearnMappingsReport => {
3580 self.handle_request_midi_learn_mappings_report().await;
3581 }
3582 Action::ClearAllMidiLearnBindings => {
3583 if Self::box_bool(self.handle_clear_all_midi_learn_bindings(a.clone())).await {
3584 return;
3585 }
3586 }
3587 #[cfg(all(unix, not(target_os = "macos")))]
3588 Action::TrackLv2PluginControls { .. } => {}
3589 #[cfg(all(unix, not(target_os = "macos")))]
3590 Action::ClipLv2PluginControls { .. } => {}
3591 #[cfg(all(unix, not(target_os = "macos")))]
3592 Action::TrackLv2StateSnapshot { .. } => {}
3593 #[cfg(all(unix, not(target_os = "macos")))]
3594 Action::ClipLv2StateSnapshot { .. } => {}
3595 #[cfg(all(unix, not(target_os = "macos")))]
3596 Action::TrackLv2Midnam { .. } => {}
3597 Action::TrackClapNoteNames { .. } => {}
3598 Action::SessionDiagnosticsReport { .. } => {}
3599 Action::MidiLearnMappingsReport { .. } => {}
3600 Action::HWInfo { .. } => {}
3601 Action::HistoryState { .. } => {}
3602 Action::Undo => {}
3603 Action::Redo => {}
3604 Action::ApplyGroupedActions(_) => {}
3605 _ => {}
3606 }
3607
3608 if let Some(inverse) = inverse_actions {
3609 if let Some(group) = self.history_group.as_mut() {
3610 group.forward_actions.push(action_to_process.clone());
3611 group.inverse_actions.splice(0..0, inverse);
3612 } else {
3613 self.history.record(UndoEntry {
3614 forward_actions: vec![action_to_process.clone()],
3615 inverse_actions: inverse,
3616 });
3617 }
3618 }
3619
3620 self.notify_clients(Ok(action_to_process)).await;
3621 }
3622 pub async fn work(&mut self) {
3623 let mut tick = tokio::time::interval(Duration::from_millis(1));
3626 tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
3627 loop {
3628 let message = tokio::select! {
3629 message = self.rx.recv() => {
3630 let Some(message) = message else {
3631 break;
3632 };
3633 message
3634 }
3635 _ = tick.tick() => {
3636 self.poll_node_worker_results().await;
3637 self.poll_jack_hw_finished().await;
3638 self.on_executor_tick().await;
3639 continue;
3640 }
3641 };
3642 self.poll_node_worker_results().await;
3643 self.poll_jack_hw_finished().await;
3644 match message {
3645 Message::Ready(id) => {
3646 if let Some(track_name) = self.bounce_worker_tracks.remove(&id) {
3650 self.offline_bounce_jobs.remove(&track_name);
3651 }
3652 self.push_ready_worker(id);
3653 if self.dispatch_node_jobs(Vec::new()).await {
3654 self.on_all_tracks_finished().await;
3655 }
3656 self.drain_pending_requests_if_idle().await;
3657 }
3658 Message::NodeDone {
3659 worker_id,
3660 epoch,
3661 node,
3662 output_linear,
3663 parameter_updates,
3664 latency_changed,
3665 } => {
3666 self.on_node_done(
3667 worker_id,
3668 epoch,
3669 node,
3670 output_linear,
3671 parameter_updates,
3672 latency_changed,
3673 )
3674 .await;
3675 }
3676 Message::Channel(s) => {
3677 self.clients.push(s);
3678 }
3679 Message::Response(result) => {
3680 self.notify_clients(result).await;
3681 }
3682
3683 Message::Request(a) => {
3684 self.dispatch_request(a).await;
3685 self.plan_builder.mark_dirty();
3688 }
3689 Message::OscRequest { action, reply_to } => {
3690 self.osc_reply_target = Some(reply_to);
3691 self.dispatch_request(action).await;
3692 self.osc_reply_target = None;
3693 self.plan_builder.mark_dirty();
3694 }
3695 Message::OfflineBounceFinished { result } => {
3696 if let Ok(Action::TrackOfflineBounce { track_name, .. })
3697 | Ok(Action::TrackOfflineBounceCanceled { track_name, .. }) = &result
3698 {
3699 self.offline_bounce_jobs.remove(track_name);
3700 }
3701 self.notify_clients(result).await;
3702 self.drain_pending_requests_if_idle().await;
3703 }
3704 Message::HWFinished => {
3705 if !self.awaiting_hwfinished {
3706 tracing::debug!("HWFinished ignored (not awaiting)");
3707 continue;
3708 }
3709 tracing::debug!("HWFinished handling; playing={}", self.playing);
3710 self.handling_hwfinished = true;
3711 self.awaiting_hwfinished = false;
3712 #[cfg(unix)]
3713 {
3714 if let Some(jack) = self.jack_runtime.as_mut() {
3715 if !self.pending_hw_midi_out_events.is_empty() {
3716 let out_events =
3717 std::mem::take(&mut self.pending_hw_midi_out_events);
3718 jack.write_events(&out_events);
3719 }
3720 let mut in_events = vec![];
3721 jack.read_events_into(&mut in_events);
3722 if !in_events.is_empty() {
3723 self.pending_hw_midi_events.extend(in_events);
3724 }
3725 let dropped = jack.take_midi_events_dropped();
3726 if dropped > 0 {
3727 tracing::warn!(
3728 "JACK MIDI ring full; {dropped} events dropped since last cycle"
3729 );
3730 }
3731 }
3732 }
3733 #[cfg(unix)]
3734 if self.jack_runtime.is_some() {
3735 self.sync_from_jack_transport().await;
3736 }
3737 while let Some(a) = self.pending_requests.pop_front() {
3738 self.handle_request(a).await;
3739 }
3740 self.apply_mute_solo_policy();
3741 self.append_recorded_cycle();
3742 self.flush_completed_recordings().await;
3743 let hw_in_routes = self.midi_hw_in_routes.clone();
3744 let pending_hw_in_by_device = self.pending_hw_midi_events_by_device.clone();
3745 let mut reconfigured_tracks = Vec::new();
3746 let state = self.state_snapshot.load_full();
3747 for (track_name, track) in state.tracks.iter() {
3748 let mut track_lock = track.lock();
3749 if self.jack_runtime_is_some() {
3750 if !self.pending_hw_midi_events.is_empty() {
3751 track_lock.push_hw_midi_events(&self.pending_hw_midi_events);
3752 }
3753 } else {
3754 for route in hw_in_routes.iter().filter(|r| &r.to_track == track_name) {
3755 if let Some(events) = pending_hw_in_by_device.get(&route.device) {
3756 track_lock.push_hw_midi_events_to_port(route.to_port, events);
3757 }
3758 }
3759 }
3760 if track_lock.setup() {
3761 reconfigured_tracks.push(track_name.clone());
3762 }
3763 }
3764 self.publish_track_meters();
3765 self.publish_session_runtime_reports().await;
3766 self.publish_clap_state_dirty().await;
3767 for track_name in reconfigured_tracks {
3768 let track = state.tracks.get(&track_name).cloned();
3769 if let Some(track) = track {
3770 let (plugins, connections, connectable_connections) = {
3771 let track_lock = track.lock();
3772 (
3773 track_lock.plugin_graph_plugins(false),
3774 track_lock.plugin_graph_connections(),
3775 track_lock.connectable_connections(),
3776 )
3777 };
3778 self.notify_clients(Ok(Action::TrackPluginGraph {
3779 track_name: track_name.clone(),
3780 plugins,
3781 connections,
3782 connectable_connections,
3783 }))
3784 .await;
3785 }
3786 }
3787 self.pending_hw_midi_events.clear();
3788 self.pending_hw_midi_events_by_device.clear();
3789 if self.transport_running {
3790 if self.transport_panic_flush_pending {
3791 self.transport_panic_flush_pending = false;
3792 } else if self.transport_restart_pending {
3793 self.transport_restart_pending = false;
3794 } else {
3795 let next = self
3796 .transport_sample
3797 .saturating_add(self.current_cycle_samples());
3798 let normalized = self.normalize_transport_sample(next);
3799 let wrapped = normalized != next;
3800 self.transport_sample = normalized;
3801 self.publish_transport_snapshot();
3802 if wrapped {
3803 if self.notified_loop_wrap_sample == Some(self.transport_sample) {
3804 self.notified_loop_wrap_sample = None;
3805 } else {
3806 self.notify_clients(Ok(Action::TransportPosition(
3807 self.transport_sample,
3808 )))
3809 .await;
3810 }
3811 }
3812 }
3813 }
3814 if self.session_clip_playback_enabled && self.playing {
3815 self.session_transport_sample = self
3816 .session_transport_sample
3817 .saturating_add(self.current_cycle_samples());
3818 }
3819 {
3820 let echoes = self.apply_modulators(self.active_transport_sample());
3821 for action in echoes {
3822 self.notify_clients(Ok(action)).await;
3823 }
3824 }
3825 self.start_plan_cycle().await;
3826 #[cfg(unix)]
3827 {
3828 if self.jack_runtime.is_some() {
3829 self.awaiting_hwfinished = true;
3830 }
3831 }
3832 self.handling_hwfinished = false;
3833 }
3834 Message::HWMidiEvents(events) => {
3835 for hw_event in events {
3836 let thru_targets: Vec<String> = self
3837 .midi_hw_thru_routes
3838 .iter()
3839 .filter(|route| route.from_device == hw_event.device)
3840 .map(|route| route.to_device.clone())
3841 .collect();
3842 for device in thru_targets {
3843 self.pending_hw_midi_out_events_by_device.push(HwMidiEvent {
3844 device,
3845 event: hw_event.event.clone(),
3846 });
3847 }
3848 if hw_event.event.data.len() >= 3 {
3849 let status = hw_event.event.data[0];
3850 if status & 0xF0 == 0xB0 {
3851 let channel = status & 0x0F;
3852 let cc = hw_event.event.data[1];
3853 let value = hw_event.event.data[2];
3854 self.handle_incoming_hw_cc(&hw_event.device, channel, cc, value)
3855 .await;
3856 }
3857 if self.step_recording_enabled && status & 0xF0 == 0x90 {
3858 let channel = status & 0x0F;
3859 let pitch = hw_event.event.data[1];
3860 let velocity = hw_event.event.data[2];
3861 if velocity > 0 {
3862 self.notify_clients(Ok(Action::StepRecordMidiNote {
3863 device: hw_event.device.clone(),
3864 channel,
3865 pitch,
3866 velocity,
3867 }))
3868 .await;
3869 }
3870 }
3871 }
3872 self.pending_hw_midi_events_by_device
3873 .entry(hw_event.device)
3874 .or_default()
3875 .push(hw_event.event);
3876 }
3877 }
3878 _ => {}
3879 }
3880 }
3881 }
3882
3883 pub(crate) fn collect_hw_midi_output_events(&self) -> Vec<MidiEvent> {
3884 let mut events = vec![];
3885 for track in self.state_snapshot.load_full().tracks.values() {
3886 events.extend(
3887 track
3888 .lock()
3889 .take_hw_midi_out_events()
3890 .into_iter()
3891 .map(|evt| evt.event),
3892 );
3893 }
3894 events.sort_by_key(|a| a.frame);
3895 events
3896 }
3897
3898 pub(crate) fn collect_hw_midi_output_events_by_device(&mut self) -> Vec<HwMidiEvent> {
3899 let mut events = Vec::<HwMidiEvent>::new();
3900 let routes = self.midi_hw_out_routes.clone();
3901 let mut events_by_track = HashMap::<String, Vec<crate::track::HwMidiOutEvent>>::new();
3902 {
3903 let state = self.state_snapshot.load_full();
3904 for route in &routes {
3905 if events_by_track.contains_key(&route.from_track) {
3906 continue;
3907 }
3908 let Some(track) = state.tracks.get(&route.from_track) else {
3909 continue;
3910 };
3911 events_by_track.insert(
3912 route.from_track.clone(),
3913 track.lock().take_hw_midi_out_events(),
3914 );
3915 }
3916 }
3917
3918 for route in routes {
3919 let Some(track_events) = events_by_track.get(&route.from_track) else {
3920 continue;
3921 };
3922 for hw_event in track_events
3923 .iter()
3924 .filter(|evt| evt.port == route.from_port)
3925 {
3926 self.update_active_hw_notes_for_track(
3927 &route.from_track,
3928 &route.device,
3929 &hw_event.event.data,
3930 );
3931 events.push(HwMidiEvent {
3932 device: route.device.clone(),
3933 event: hw_event.event.clone(),
3934 });
3935 }
3936 }
3937 events.sort_by(|a, b| {
3938 a.event
3939 .frame
3940 .cmp(&b.event.frame)
3941 .then_with(|| a.device.cmp(&b.device))
3942 });
3943 events
3944 }
3945}