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