diff --git a/ts3-client/core/src/main/java/com/ts3client/audio/AudioBackend.java b/ts3-client/core/src/main/java/com/ts3client/audio/AudioBackend.java index e0dec68..04f4a27 100644 --- a/ts3-client/core/src/main/java/com/ts3client/audio/AudioBackend.java +++ b/ts3-client/core/src/main/java/com/ts3client/audio/AudioBackend.java @@ -1,17 +1,30 @@ package com.ts3client.audio; +import com.ts3client.audio.opus.OpusCodec; import com.ts3client.config.Settings; import com.ts3client.sound.SoundPlayer; /** - * Factory for a platform's voice capture and playback. Injected into the - * connection layer so the core stays independent of any concrete audio stack. + * A platform's audio stack: its devices, its libopus and its sound player. Injected into + * the connection layer so the core stays independent of any concrete audio stack; the + * voice pipelines on top are the same everywhere. */ public interface AudioBackend { - VoiceInput createInput(Settings settings); + AudioIo io(); - VoiceOutput createOutput(Settings settings); + OpusCodec opus(); + + /** How incoming voice reaches this platform's devices. */ + StreamingVoiceOutput.Lines outputLines(); + + default VoiceInput createInput(Settings settings) { + return new CaptureVoiceInput(settings, io(), opus()); + } + + default VoiceOutput createOutput(Settings settings) { + return new StreamingVoiceOutput(io(), opus(), settings.outputDevice, outputLines()); + } /** Player for notification sounds (sound packs); shared by all connections. */ SoundPlayer createSoundPlayer(Settings settings); diff --git a/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/AudioCapture.java b/ts3-client/core/src/main/java/com/ts3client/audio/AudioCapture.java similarity index 71% rename from ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/AudioCapture.java rename to ts3-client/core/src/main/java/com/ts3client/audio/AudioCapture.java index 30441d4..827ed76 100644 --- a/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/AudioCapture.java +++ b/ts3-client/core/src/main/java/com/ts3client/audio/AudioCapture.java @@ -1,11 +1,11 @@ -package com.ts3client.audio.desktop; +package com.ts3client.audio; /** * An open microphone line, delivering 16-bit little-endian PCM at - * {@link AudioDevices#SAMPLE_RATE}. + * {@link VoiceFormat#SAMPLE_RATE}. * - *

Modelled on {@link javax.sound.sampled.TargetDataLine} so the capture pipeline does - * not care whether the audio comes from PipeWire or Java Sound. + *

Modelled on Java Sound's {@code TargetDataLine} so the capture pipeline does + * not care which platform audio API delivers it. */ public interface AudioCapture extends AutoCloseable { diff --git a/ts3-client/core/src/main/java/com/ts3client/audio/AudioDevice.java b/ts3-client/core/src/main/java/com/ts3client/audio/AudioDevice.java new file mode 100644 index 0000000..11145b7 --- /dev/null +++ b/ts3-client/core/src/main/java/com/ts3client/audio/AudioDevice.java @@ -0,0 +1,13 @@ +package com.ts3client.audio; + +/** A selectable audio device. {@code id} is what gets persisted in the settings. */ +public record AudioDevice(String id, String label) { + + /** The implicit entry that lets the platform choose. */ + public static final AudioDevice DEFAULT = new AudioDevice("", "(System default)"); + + @Override + public String toString() { + return label; + } +} diff --git a/ts3-client/core/src/main/java/com/ts3client/audio/AudioIo.java b/ts3-client/core/src/main/java/com/ts3client/audio/AudioIo.java new file mode 100644 index 0000000..22247de --- /dev/null +++ b/ts3-client/core/src/main/java/com/ts3client/audio/AudioIo.java @@ -0,0 +1,23 @@ +package com.ts3client.audio; + +import java.util.List; + +/** The platform's audio devices: what can be picked, and opening lines on them. */ +public interface AudioIo { + + /** Devices that can provide microphone lines, {@link AudioDevice#DEFAULT} first. */ + List inputDevices(); + + /** Devices that can provide speaker lines, {@link AudioDevice#DEFAULT} first. */ + List outputDevices(); + + /** + * Opens a capture line at {@link VoiceFormat#SAMPLE_RATE}, preferring + * {@code preferredChannels} but falling back to what the device takes; inspect + * {@link AudioCapture#channels()} for the result. An empty id means the default device. + */ + AudioCapture openCapture(String deviceId, int preferredChannels) throws Exception; + + /** Opens a playback line; see {@link #openCapture} for the channel negotiation. */ + AudioPlayback openPlayback(String deviceId, int preferredChannels) throws Exception; +} diff --git a/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/AudioPlayback.java b/ts3-client/core/src/main/java/com/ts3client/audio/AudioPlayback.java similarity index 70% rename from ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/AudioPlayback.java rename to ts3-client/core/src/main/java/com/ts3client/audio/AudioPlayback.java index 211e25f..3ac3346 100644 --- a/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/AudioPlayback.java +++ b/ts3-client/core/src/main/java/com/ts3client/audio/AudioPlayback.java @@ -1,11 +1,11 @@ -package com.ts3client.audio.desktop; +package com.ts3client.audio; /** * An open speaker line, accepting 16-bit little-endian PCM at - * {@link AudioDevices#SAMPLE_RATE}. + * {@link VoiceFormat#SAMPLE_RATE}. * - *

Modelled on {@link javax.sound.sampled.SourceDataLine} so the playback pipeline does - * not care whether the audio goes to PipeWire or Java Sound. + *

Modelled on Java Sound's {@code SourceDataLine} so the playback pipeline does + * not care which platform audio API plays it. */ public interface AudioPlayback extends AutoCloseable { diff --git a/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/DesktopVoiceInput.java b/ts3-client/core/src/main/java/com/ts3client/audio/CaptureVoiceInput.java similarity index 91% rename from ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/DesktopVoiceInput.java rename to ts3-client/core/src/main/java/com/ts3client/audio/CaptureVoiceInput.java index ce5330d..eb45ec9 100644 --- a/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/DesktopVoiceInput.java +++ b/ts3-client/core/src/main/java/com/ts3client/audio/CaptureVoiceInput.java @@ -1,11 +1,8 @@ -package com.ts3client.audio.desktop; +package com.ts3client.audio; import com.github.manevolent.ts3j.enums.CodecType; -import com.ts3client.audio.AudioFrameListener; -import com.ts3client.audio.InputLevel; -import com.ts3client.audio.OpusParameters; -import com.ts3client.audio.SpeechProbabilityDetector; -import com.ts3client.audio.VoiceInput; +import com.ts3client.audio.opus.OpusCodec; +import com.ts3client.audio.opus.OpusEncoder; import com.ts3client.audio.processing.AudioProcessor; import com.ts3client.audio.processing.ns.SuppressionLevel; import com.ts3client.audio.vad.RnnSpeechDetector; @@ -16,11 +13,11 @@ import java.util.concurrent.atomic.AtomicBoolean; import java.util.function.Consumer; /** - * Desktop {@link VoiceInput}: captures the microphone, applies voice-activation + * The {@link VoiceInput} every platform shares: captures the microphone, applies voice-activation * or push-to-talk gating, and Opus-encodes 20 ms frames. * - *

The capture line comes from {@link AudioDevices}, which picks PipeWire or Java - * Sound; it is opened in stereo when the device offers it. Voice + *

The capture line comes from the platform's {@link AudioIo}; it is opened in stereo + * when the device offers it. Voice * ({@code OPUS_VOICE}) is transmitted mono, from the downmix, so the pre-processing * chain and VAD see a single channel; the music codec ({@code OPUS_MUSIC}) transmits * the stereo capture as-is. @@ -30,7 +27,7 @@ import java.util.function.Consumer; * network sending. When the gate closes and the queue drains {@link #isReady()} * returns {@code false}, prompting ts3j to emit the terminating empty voice packet. */ -public final class DesktopVoiceInput implements VoiceInput { +public final class CaptureVoiceInput implements VoiceInput { /** * How long the gate stays open after the last active frame. TS3 counts 64 of its @@ -46,9 +43,9 @@ public final class DesktopVoiceInput implements VoiceInput { private static final int PREROLL_MS = 30; private static final int HANGOVER_FRAMES = - Math.max(1, HANGOVER_MS * AudioDevices.SAMPLE_RATE / 1000 / AudioDevices.FRAME_SIZE); + Math.max(1, HANGOVER_MS * VoiceFormat.SAMPLE_RATE / 1000 / VoiceFormat.FRAME_SIZE); private static final int PREROLL_FRAMES = - Math.max(1, PREROLL_MS * AudioDevices.SAMPLE_RATE / 1000 / AudioDevices.FRAME_SIZE); + Math.max(1, PREROLL_MS * VoiceFormat.SAMPLE_RATE / 1000 / VoiceFormat.FRAME_SIZE); private final ConcurrentLinkedQueue queue = new ConcurrentLinkedQueue<>(); private final AtomicBoolean muted = new AtomicBoolean(false); @@ -66,7 +63,7 @@ public final class DesktopVoiceInput implements VoiceInput { private volatile CodecType codec = CodecType.OPUS_VOICE; private final SpeechProbabilityDetector speechDetector = - new RnnSpeechDetector(AudioDevices.SAMPLE_RATE); + new RnnSpeechDetector(VoiceFormat.SAMPLE_RATE); private final AudioProcessor processor = new AudioProcessor(); private volatile Consumer levelListener; // input level, InputLevel scale @@ -75,12 +72,14 @@ public final class DesktopVoiceInput implements VoiceInput { private volatile AudioFrameListener monitorListener; // local monitoring of what is sent private boolean mutedTalking; + private final AudioIo io; + private final OpusCodec opus; private final String deviceName; private final Object encoderLock = new Object(); private volatile OpusParameters params; private OpusParameters appliedParams; - private int encoderApplication = -1; + private OpusCodec.Application encoderApplication; private volatile int encoderChannels = 1; private volatile int captureChannels = 1; @@ -96,7 +95,9 @@ public final class DesktopVoiceInput implements VoiceInput { private int prerollTail; private int prerollCount; - public DesktopVoiceInput(Settings settings) { + public CaptureVoiceInput(Settings settings, AudioIo io, OpusCodec opus) { + this.io = io; + this.opus = opus; this.deviceName = settings.inputDevice; this.mode = settings.inputMode; this.vadMode = settings.vadMode; @@ -173,8 +174,8 @@ public final class DesktopVoiceInput implements VoiceInput { } } - private static int applicationFor(OpusParameters p) { - return p.music ? Opus.OPUS_APPLICATION_AUDIO : Opus.OPUS_APPLICATION_VOIP; + private static OpusCodec.Application applicationFor(OpusParameters p) { + return p.music ? OpusCodec.Application.AUDIO : OpusCodec.Application.VOIP; } private static CodecType codecFor(OpusParameters p) { @@ -191,13 +192,13 @@ public final class DesktopVoiceInput implements VoiceInput { /** Creates or reconfigures the encoder to match {@code p}. Call under {@link #encoderLock}. */ private void applyParams(OpusParameters p) { - int application = applicationFor(p); + OpusCodec.Application application = applicationFor(p); int channels = channelsFor(p); // Application (VOIP vs AUDIO) and channel count are fixed at creation, so // switching voice<->music means building a new encoder. if (encoder == null || application != encoderApplication || channels != encoderChannels) { - OpusEncoder replacement = new OpusEncoder( - AudioDevices.SAMPLE_RATE, AudioDevices.FRAME_SIZE, channels, application); + OpusEncoder replacement = opus.createEncoder( + VoiceFormat.SAMPLE_RATE, VoiceFormat.FRAME_SIZE, channels, application); configureEncoder(replacement, p); OpusEncoder previous = encoder; encoder = replacement; @@ -217,7 +218,7 @@ public final class DesktopVoiceInput implements VoiceInput { enc.setVbr(p.vbr); enc.setInbandFec(p.fec); enc.setExpectedPacketLoss(p.expectedPacketLoss); - enc.setSignal(p.music ? Opus.OPUS_SIGNAL_MUSIC : Opus.OPUS_SIGNAL_VOICE); + enc.setSignal(p.music ? OpusEncoder.Signal.MUSIC : OpusEncoder.Signal.VOICE); } public void setLevelListener(Consumer l) { @@ -261,7 +262,7 @@ public final class DesktopVoiceInput implements VoiceInput { public synchronized void start() { if (running.get()) return; try { - line = AudioDevices.openCapture(deviceName, AudioDevices.MAX_CHANNELS); + line = io.openCapture(deviceName, VoiceFormat.MAX_CHANNELS); captureChannels = line.channels(); synchronized (encoderLock) { applyParams(params); @@ -301,14 +302,14 @@ public final class DesktopVoiceInput implements VoiceInput { } catch (Exception ignored) { } encoder = null; - encoderApplication = -1; + encoderApplication = null; appliedParams = null; } } } private void captureLoop() { - final int frameSamples = AudioDevices.FRAME_SIZE; + final int frameSamples = VoiceFormat.FRAME_SIZE; final int channels = captureChannels; final byte[] buf = new byte[frameSamples * 2 * channels]; final float[] pcm = new float[frameSamples * channels]; // interleaved capture @@ -409,7 +410,7 @@ public final class DesktopVoiceInput implements VoiceInput { /** Encodes and queues the buffered lead-in, oldest first. Call under {@link #encoderLock}. */ private void flushPreroll(boolean stereo) { - int expected = stereo ? AudioDevices.FRAME_SIZE * encoderChannels : AudioDevices.FRAME_SIZE; + int expected = stereo ? VoiceFormat.FRAME_SIZE * encoderChannels : VoiceFormat.FRAME_SIZE; for (int i = 0; i < prerollCount; i++) { int idx = (prerollTail - prerollCount + i + PREROLL_FRAMES) % PREROLL_FRAMES; float[] frame = preroll[idx]; diff --git a/ts3-client/core/src/main/java/com/ts3client/audio/MixedPlayout.java b/ts3-client/core/src/main/java/com/ts3client/audio/MixedPlayout.java new file mode 100644 index 0000000..fa2ebdb --- /dev/null +++ b/ts3-client/core/src/main/java/com/ts3client/audio/MixedPlayout.java @@ -0,0 +1,110 @@ +package com.ts3client.audio; + +import java.util.Arrays; +import java.util.Collection; +import java.util.function.Supplier; +import java.util.function.ToDoubleFunction; + +/** + * Mixes every speaker onto one playback line, as mobile audio APIs want: a single + * low-latency stream. The line's blocking writes pace the mix, so it runs on the device's + * clock; between talk bursts the mixer writes nothing and waits. + */ +final class MixedPlayout implements Playout { + + private static final int FRAME = VoiceFormat.FRAME_SIZE; + + private final AudioIo io; + private final Supplier outputDevice; + private final ToDoubleFunction gain; + private final Collection streams; + + private final Object lock = new Object(); + /** Guarded by {@link #lock}. */ + private boolean running = true; + private Thread mixer; + + MixedPlayout(AudioIo io, Supplier outputDevice, ToDoubleFunction gain, + Collection streams) { + this.io = io; + this.outputDevice = outputDevice; + this.gain = gain; + this.streams = streams; + } + + @Override + public void started(VoiceStream stream) { + synchronized (lock) { + if (!running) return; + if (mixer == null) { + mixer = new Thread(this::mixLoop, "ts3j-voice-mixer"); + mixer.setDaemon(true); + mixer.start(); + } + lock.notifyAll(); + } + } + + @Override + public void removed(VoiceStream stream) { + // The stream has already left the collection the mixer iterates. + } + + @Override + public void shutdown() { + synchronized (lock) { + running = false; + if (mixer != null) mixer.interrupt(); + lock.notifyAll(); + } + } + + private void mixLoop() { + AudioPlayback line = null; + try { + int channels = 0; + float[] mix = null; + float[] scratch = null; + byte[] out = null; + while (awaitSpeech()) { + if (line == null) { + line = io.openPlayback(outputDevice.get(), VoiceFormat.MAX_CHANNELS); + line.start(); + channels = line.channels(); + mix = new float[FRAME * channels]; + scratch = new float[VoiceStream.MAX_FRAME * channels]; + out = new byte[mix.length * 2]; + } + Arrays.fill(mix, 0f); + for (VoiceStream stream : streams) { + int frames = stream.pull(scratch, channels, gain.applyAsDouble(stream)); + // TeamSpeak clients send 20 ms packets; anything longer is cut to the frame. + int samples = Math.min(frames, FRAME) * channels; + for (int i = 0; i < samples; i++) mix[i] += scratch[i]; + } + Pcm16.encode(mix, mix.length, out); + line.write(out, 0, out.length); + } + } catch (Exception ignored) { + // No usable line: this mixer ends, and the next talk burst tries again. + } finally { + if (line != null) line.close(); + synchronized (lock) { + mixer = null; + } + } + } + + /** Waits until some stream is playing; false once shut down. */ + private boolean awaitSpeech() throws InterruptedException { + synchronized (lock) { + while (running) { + for (VoiceStream stream : streams) { + if (stream.isPlaying()) return true; + } + lock.wait(); + } + return false; + } + } +} diff --git a/ts3-client/core/src/main/java/com/ts3client/audio/Pcm16.java b/ts3-client/core/src/main/java/com/ts3client/audio/Pcm16.java new file mode 100644 index 0000000..ea38f07 --- /dev/null +++ b/ts3-client/core/src/main/java/com/ts3client/audio/Pcm16.java @@ -0,0 +1,20 @@ +package com.ts3client.audio; + +/** Conversion between float samples in [-1, 1] and 16-bit little-endian PCM. */ +public final class Pcm16 { + + private Pcm16() { + } + + /** Writes {@code count} samples to {@code out}, clipping anything out of range. */ + public static void encode(float[] samples, int count, byte[] out) { + for (int i = 0; i < count; i++) { + float f = samples[i]; + if (f > 1f) f = 1f; + else if (f < -1f) f = -1f; + int s = Math.round(f * 32767f); + out[2 * i] = (byte) (s & 0xFF); + out[2 * i + 1] = (byte) ((s >> 8) & 0xFF); + } + } +} diff --git a/ts3-client/core/src/main/java/com/ts3client/audio/PerSpeakerPlayout.java b/ts3-client/core/src/main/java/com/ts3client/audio/PerSpeakerPlayout.java new file mode 100644 index 0000000..dad2173 --- /dev/null +++ b/ts3-client/core/src/main/java/com/ts3client/audio/PerSpeakerPlayout.java @@ -0,0 +1,149 @@ +package com.ts3client.audio; + +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.ScheduledFuture; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.function.Supplier; +import java.util.function.ToDoubleFunction; + +/** + * Gives every speaker a playback line and worker thread of its own, so simultaneous + * speakers are mixed by the audio server and one slow decode never blocks another. On a + * PipeWire desktop every speaker then shows up as a separate stream in the mixer. + */ +final class PerSpeakerPlayout implements Playout { + + private static final long FRAME_NANOS = + TimeUnit.SECONDS.toNanos(VoiceFormat.FRAME_SIZE) / VoiceFormat.SAMPLE_RATE; + + private final class Speaker { + final VoiceStream stream; + final AudioPlayback line; + final int lineChannels; + final ExecutorService worker; + final AtomicBoolean tickQueued = new AtomicBoolean(); + /** Touched only by {@link #worker}. */ + final float[] pcm; + final byte[] out; + /** Guarded by this speaker. */ + ScheduledFuture ticks; + + Speaker(VoiceStream stream) throws Exception { + this.stream = stream; + this.line = io.openPlayback(outputDevice.get(), VoiceFormat.MAX_CHANNELS); + this.lineChannels = line.channels(); + this.line.start(); + this.pcm = new float[VoiceStream.MAX_FRAME * lineChannels]; + this.out = new byte[pcm.length * 2]; + this.worker = Executors.newSingleThreadExecutor(r -> { + Thread t = new Thread(r, "ts3j-play-" + stream.clientId); + t.setDaemon(true); + return t; + }); + } + + void close() { + synchronized (this) { + if (ticks != null) ticks.cancel(false); + } + worker.shutdownNow(); + line.close(); + } + } + + private final AudioIo io; + private final Supplier outputDevice; + private final ToDoubleFunction gain; + private final ScheduledExecutorService clock; + private final Map speakers = new ConcurrentHashMap<>(); + + PerSpeakerPlayout(AudioIo io, Supplier outputDevice, ToDoubleFunction gain, + ScheduledExecutorService clock) { + this.io = io; + this.outputDevice = outputDevice; + this.gain = gain; + this.clock = clock; + } + + @Override + public void started(VoiceStream stream) { + Speaker speaker = speakers.get(stream); + if (speaker == null) { + try { + speaker = new Speaker(stream); + } catch (Exception e) { + stream.stop(); // couldn't open a line; drop the burst + return; + } + Speaker existing = speakers.putIfAbsent(stream, speaker); + if (existing != null) { + speaker.close(); + speaker = existing; + } + } + synchronized (speaker) { + if (speaker.ticks == null) { + Speaker s = speaker; + speaker.ticks = clock.scheduleAtFixedRate( + () -> queueTick(s), FRAME_NANOS, FRAME_NANOS, TimeUnit.NANOSECONDS); + } + } + } + + private void queueTick(Speaker speaker) { + if (!speaker.tickQueued.compareAndSet(false, true)) return; + try { + speaker.worker.submit(() -> { + try { + tick(speaker); + } finally { + speaker.tickQueued.set(false); + } + }); + } catch (RuntimeException e) { + if (speaker.worker.isShutdown()) { + speaker.tickQueued.set(false); + } else { + throw e; + } + } + } + + private void tick(Speaker speaker) { + int frames = speaker.stream.pull(speaker.pcm, speaker.lineChannels, gain.applyAsDouble(speaker.stream)); + if (frames == VoiceStream.IDLE) { + synchronized (speaker) { + // A packet may have started the next burst since the pull; keep ticking then. + if (!speaker.stream.isPlaying() && speaker.ticks != null) { + speaker.ticks.cancel(false); + speaker.ticks = null; + } + } + return; + } + if (frames == 0) return; + int samples = frames * speaker.lineChannels; + Pcm16.encode(speaker.pcm, samples, speaker.out); + try { + speaker.line.write(speaker.out, 0, samples * 2); + } catch (RuntimeException ignored) { + } + } + + @Override + public void removed(VoiceStream stream) { + Speaker speaker = speakers.remove(stream); + if (speaker != null) speaker.close(); + } + + @Override + public void shutdown() { + for (Speaker speaker : speakers.values()) speaker.close(); + speakers.clear(); + } +} diff --git a/ts3-client/core/src/main/java/com/ts3client/audio/Playout.java b/ts3-client/core/src/main/java/com/ts3client/audio/Playout.java new file mode 100644 index 0000000..c691b3f --- /dev/null +++ b/ts3-client/core/src/main/java/com/ts3client/audio/Playout.java @@ -0,0 +1,13 @@ +package com.ts3client.audio; + +/** Drives the {@link VoiceStream}s of a {@link StreamingVoiceOutput} out to the speakers. */ +interface Playout { + + /** A stream began a talk burst and must now be pulled every 20 ms. */ + void started(VoiceStream stream); + + /** The stream's client left; it will not be pulled again. */ + void removed(VoiceStream stream); + + void shutdown(); +} diff --git a/ts3-client/core/src/main/java/com/ts3client/audio/StreamingVoiceOutput.java b/ts3-client/core/src/main/java/com/ts3client/audio/StreamingVoiceOutput.java new file mode 100644 index 0000000..bf69fe9 --- /dev/null +++ b/ts3-client/core/src/main/java/com/ts3client/audio/StreamingVoiceOutput.java @@ -0,0 +1,182 @@ +package com.ts3client.audio; + +import com.github.manevolent.ts3j.enums.CodecType; +import com.github.manevolent.ts3j.protocol.packet.PacketBody0Voice; +import com.github.manevolent.ts3j.protocol.packet.PacketBody1VoiceWhisper; +import com.ts3client.audio.opus.OpusCodec; + +import java.util.Map; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; +import java.util.function.BiConsumer; + +/** + * The {@link VoiceOutput} every platform shares: decodes incoming voice per speaker, each + * through its own jitter buffer, and plays it on the platform's {@link AudioIo}. + */ +public final class StreamingVoiceOutput implements VoiceOutput { + + /** How speakers reach the device. */ + public enum Lines { + /** A line per speaker, mixed by the platform's audio server. */ + PER_SPEAKER, + /** One line all speakers are mixed onto here. */ + SHARED + } + + /** + * A speaker is only meant to stop when the empty voice packet marking the end of a + * talk burst arrives — but that packet is UDP too, and a lost one otherwise leaves + * the talking indicator stuck until the speaker's next burst. Opus frames are 20 ms + * apart while someone is actually talking, so a gap several times that long, with no + * packet of either kind, is unambiguous: force the indicator off rather than trust + * the one packet that could go missing. + */ + private static final long TALK_TIMEOUT_NANOS = TimeUnit.MILLISECONDS.toNanos(400); + + private final OpusCodec opus; + private final Map streams = new ConcurrentHashMap<>(); + private final Set mutedClients = ConcurrentHashMap.newKeySet(); + /** Per-client gain on top of the master volume; absent means 1.0. */ + private final Map clientVolumes = new ConcurrentHashMap<>(); + + private volatile String outputDevice; + private volatile double masterVolume = 1.0; + private volatile boolean deafened = false; + + /** Notified (clientId, talking) on an audio thread when a speaker starts/stops. */ + private volatile BiConsumer talkListener; + + /** Catches talk bursts whose end packet never arrived, and paces per-speaker lines. */ + private final ScheduledExecutorService clock = Executors.newSingleThreadScheduledExecutor(r -> { + Thread t = new Thread(r, "ts3j-talk-watchdog"); + t.setDaemon(true); + return t; + }); + private final Playout playout; + + public StreamingVoiceOutput(AudioIo io, OpusCodec opus, String outputDevice, Lines lines) { + this.opus = opus; + this.outputDevice = outputDevice; + this.playout = switch (lines) { + case PER_SPEAKER -> new PerSpeakerPlayout(io, () -> this.outputDevice, this::gainFor, clock); + case SHARED -> new MixedPlayout(io, () -> this.outputDevice, this::gainFor, streams.values()); + }; + clock.scheduleWithFixedDelay(this::checkTalkTimeouts, 50, 50, TimeUnit.MILLISECONDS); + } + + private void checkTalkTimeouts() { + long now = System.nanoTime(); + for (VoiceStream s : streams.values()) { + if (s.isTalking() && now - s.lastPacketNanos() > TALK_TIMEOUT_NANOS) s.stop(); + } + } + + private double gainFor(VoiceStream stream) { + return masterVolume * clientVolumes.getOrDefault(stream.clientId, 1.0); + } + + private void fireTalk(int clientId, boolean talking) { + BiConsumer l = talkListener; + if (l != null) l.accept(clientId, talking); + } + + @Override + public void setTalkListener(BiConsumer l) { + this.talkListener = l; + } + + @Override + public void setMasterVolume(double v) { + this.masterVolume = Math.max(0, Math.min(2.0, v)); + } + + @Override + public void setDeafened(boolean d) { + this.deafened = d; + if (d) { + // Stop everyone talking immediately. + for (VoiceStream s : streams.values()) s.stop(); + } + } + + @Override + public boolean isDeafened() { + return deafened; + } + + @Override + public void setClientMuted(int clientId, boolean muted) { + if (muted) { + mutedClients.add(clientId); + VoiceStream s = streams.get(clientId); + if (s != null) s.stop(); + } else { + mutedClients.remove(clientId); + } + } + + @Override + public boolean isClientMuted(int clientId) { + return mutedClients.contains(clientId); + } + + @Override + public void setClientVolume(int clientId, double gain) { + if (gain == 1.0) clientVolumes.remove(clientId); + else clientVolumes.put(clientId, Math.max(0, gain)); + } + + @Override + public double getClientVolume(int clientId) { + return clientVolumes.getOrDefault(clientId, 1.0); + } + + @Override + public void setOutputDevice(String device) { + this.outputDevice = device; + } + + /** Entry point wired into {@code client.setVoiceHandler(...)}. */ + @Override + public void handleVoice(PacketBody0Voice voice) { + route(voice.getClientId(), voice.getPacketId(), voice.getCodecType(), voice.getCodecData()); + } + + /** Entry point wired into {@code client.setWhisperHandler(...)}. */ + @Override + public void handleWhisper(PacketBody1VoiceWhisper whisper) { + route(whisper.getClientId(), whisper.getPacketId(), whisper.getCodecType(), whisper.getCodecData()); + } + + private void route(int clientId, int packetId, CodecType codec, byte[] data) { + if (deafened) return; + if (mutedClients.contains(clientId)) return; + + VoiceStream stream = streams.computeIfAbsent(clientId, id -> new VoiceStream(id, opus, this::fireTalk)); + if (stream.offer(packetId, codec, data)) playout.started(stream); + } + + /** Forgets a client that left; the server may hand its id to somebody else. */ + @Override + public void removeClient(int clientId) { + VoiceStream s = streams.remove(clientId); + if (s != null) { + playout.removed(s); + s.close(); + } + mutedClients.remove(clientId); + clientVolumes.remove(clientId); + } + + @Override + public void shutdown() { + clock.shutdownNow(); + playout.shutdown(); + for (VoiceStream s : streams.values()) s.close(); + streams.clear(); + } +} diff --git a/ts3-client/core/src/main/java/com/ts3client/audio/VoiceFormat.java b/ts3-client/core/src/main/java/com/ts3client/audio/VoiceFormat.java new file mode 100644 index 0000000..cb4bd97 --- /dev/null +++ b/ts3-client/core/src/main/java/com/ts3client/audio/VoiceFormat.java @@ -0,0 +1,16 @@ +package com.ts3client.audio; + +/** + * The PCM format voice is captured, processed and played in. TeamSpeak's Opus codecs run at + * 48 kHz in 20 ms frames, mono for {@code OPUS_VOICE} and stereo for {@code OPUS_MUSIC}. + */ +public final class VoiceFormat { + + public static final int SAMPLE_RATE = 48_000; + /** Samples per channel in one 20 ms frame. */ + public static final int FRAME_SIZE = 960; + public static final int MAX_CHANNELS = 2; + + private VoiceFormat() { + } +} diff --git a/ts3-client/core/src/main/java/com/ts3client/audio/VoiceStream.java b/ts3-client/core/src/main/java/com/ts3client/audio/VoiceStream.java new file mode 100644 index 0000000..31ed5f3 --- /dev/null +++ b/ts3-client/core/src/main/java/com/ts3client/audio/VoiceStream.java @@ -0,0 +1,254 @@ +package com.ts3client.audio; + +import com.github.manevolent.ts3j.enums.CodecType; +import com.ts3client.audio.opus.OpusCodec; +import com.ts3client.audio.opus.OpusDecoder; + +import java.util.HashMap; +import java.util.Map; +import java.util.function.BiConsumer; + +/** + * One speaker's incoming voice: packets go into a small jitter buffer as they arrive, and + * every {@link #pull} plays out the next one in sequence, concealing any that went missing. + * + *

The stream keeps no clock of its own: whoever pulls sets the pace, one call per + * 20 ms frame. {@link #offer} and {@link #stop} may be called from any thread; pulls + * must come from one thread at a time. + */ +final class VoiceStream { + + /** Returned by {@link #pull} when no talk burst is playing. */ + static final int IDLE = -1; + + /** Longest Opus frame (120 ms @ 48 kHz) a packet may decode to, per channel. */ + static final int MAX_FRAME = 5760; + /** Initial playout delay: enough room for common UDP jitter without feeling laggy. */ + private static final int JITTER_DELAY_FRAMES = 3; + /** Cap PLC during a lost end-of-talk marker or a hard network stall. */ + private static final int MAX_CONSECUTIVE_PLC_FRAMES = 10; + /** Keep memory bounded when a talk burst runs far ahead of playback. */ + private static final int MAX_PENDING_PACKETS = 32; + + private record VoiceFrame(int packetId, CodecType codec, byte[] data, boolean end) { + } + + final int clientId; + private final OpusCodec opus; + private final BiConsumer talkListener; + + private final Object lock = new Object(); + private final Map pending = new HashMap<>(); + private boolean playing; + private int expectedPacketId; + private int activeChannels = 1; + private int consecutivePlcFrames; + /** Pulls still to sit out before the burst starts playing. */ + private int warmup; + + /** Guards the decoder: a stop or close from another thread may reach it mid-pull. */ + private final Object decoderLock = new Object(); + private final float[] pcm = new float[MAX_FRAME * VoiceFormat.MAX_CHANNELS]; + private OpusDecoder decoder; + private int decoderChannels; + private boolean closed; + + private volatile boolean talking; + private volatile long lastPacketNanos; + + VoiceStream(int clientId, OpusCodec opus, BiConsumer talkListener) { + this.clientId = clientId; + this.opus = opus; + this.talkListener = talkListener; + } + + /** TeamSpeak streams {@code OPUS_MUSIC} in stereo and everything else in mono. */ + private static int channelsFor(CodecType codec) { + return codec == CodecType.OPUS_MUSIC ? 2 : 1; + } + + /** + * Queues an incoming packet; an empty one marks the end of the talk burst. + * + * @return true if the packet started a new burst, so playout has to begin + */ + boolean offer(int packetId, CodecType codec, byte[] data) { + lastPacketNanos = System.nanoTime(); + VoiceFrame frame = new VoiceFrame(packetId & 0xFFFF, codec, data, data == null || data.length == 0); + boolean started = false; + + synchronized (lock) { + if (!playing) { + playing = true; + started = true; + expectedPacketId = frame.packetId(); + activeChannels = channelsFor(codec); + consecutivePlcFrames = 0; + warmup = JITTER_DELAY_FRAMES - 1; + } else if (packetDistance(expectedPacketId, frame.packetId()) < 0) { + return false; // too late for this burst; do not replay stale audio + } + pending.putIfAbsent(frame.packetId(), frame); + trimPending(); + } + if (!frame.end()) markTalking(true); + return started; + } + + /** + * Plays out the next frame into {@code out}, interleaved over {@code outChannels} and + * scaled by {@code gain} but not clipped, so it can still be mixed. + * + * @param out room for {@link #MAX_FRAME} samples per channel + * @return the samples written per channel; 0 while the burst is still buffering or a + * packet could not be decoded, or {@link #IDLE} once the burst is over + */ + int pull(float[] out, int outChannels, double gain) { + VoiceFrame frame; + int channels; + synchronized (lock) { + if (!playing) return IDLE; + if (warmup > 0) { + warmup--; + return 0; + } + + frame = pending.remove(expectedPacketId); + boolean conceal = frame == null; + if (conceal) { + consecutivePlcFrames++; + } else { + consecutivePlcFrames = 0; + } + if ((frame != null && frame.end()) || consecutivePlcFrames > MAX_CONSECUTIVE_PLC_FRAMES) { + stopLocked(); + frame = null; + channels = 0; + } else { + channels = frame != null ? channelsFor(frame.codec()) : activeChannels; + activeChannels = channels; + expectedPacketId = (expectedPacketId + 1) & 0xFFFF; + } + } + if (channels == 0) { + finish(); + return IDLE; + } + return decode(frame != null ? frame.data() : null, channels, out, outChannels, gain); + } + + private int decode(byte[] data, int channels, float[] out, int outChannels, double gain) { + int frames; + synchronized (decoderLock) { + if (closed) return 0; + try { + frames = decoderFor(channels).decode(data, pcm); + } catch (RuntimeException e) { + return 0; + } + } + // Match the decoded stream to the output: duplicate mono across stereo, fold a + // stereo (music) stream down onto mono. + for (int i = 0, k = 0; i < frames; i++) { + for (int c = 0; c < outChannels; c++, k++) { + double v; + if (channels == outChannels) { + v = pcm[i * channels + c]; + } else if (channels == 1) { + v = pcm[i]; + } else { + double sum = 0; + for (int s = 0; s < channels; s++) sum += pcm[i * channels + s]; + v = sum / channels; + } + out[k] = (float) (v * gain); + } + } + return frames; + } + + /** + * Returns a decoder matching the stream's channel count, rebuilding it when a speaker + * switches between the mono {@code OPUS_VOICE} and stereo {@code OPUS_MUSIC} codecs. + */ + private OpusDecoder decoderFor(int channels) { + if (decoder == null || decoderChannels != channels) { + if (decoder != null) decoder.close(); + decoder = opus.createDecoder(VoiceFormat.SAMPLE_RATE, MAX_FRAME, channels); + decoderChannels = channels; + } + return decoder; + } + + /** Ends the current talk burst, if any; the next packet starts a new one. */ + void stop() { + synchronized (lock) { + stopLocked(); + } + finish(); + } + + private void stopLocked() { + playing = false; + pending.clear(); + } + + private void finish() { + synchronized (decoderLock) { + if (decoder != null) decoder.reset(); + } + markTalking(false); + } + + boolean isPlaying() { + synchronized (lock) { + return playing; + } + } + + boolean isTalking() { + return talking; + } + + long lastPacketNanos() { + return lastPacketNanos; + } + + void close() { + synchronized (lock) { + stopLocked(); + } + synchronized (decoderLock) { + closed = true; + if (decoder != null) decoder.close(); + decoder = null; + } + } + + private void markTalking(boolean talking) { + if (this.talking == talking) return; + this.talking = talking; + talkListener.accept(clientId, talking); + } + + private static int packetDistance(int from, int to) { + int distance = (to - from) & 0xFFFF; + return distance >= 0x8000 ? distance - 0x10000 : distance; + } + + private void trimPending() { + while (pending.size() > MAX_PENDING_PACKETS) { + Integer furthest = null; + int furthestDistance = Integer.MIN_VALUE; + for (Integer packetId : pending.keySet()) { + int distance = packetDistance(expectedPacketId, packetId); + if (distance > furthestDistance) { + furthestDistance = distance; + furthest = packetId; + } + } + if (furthest == null) return; + pending.remove(furthest); + } + } +} diff --git a/ts3-client/core/src/main/java/com/ts3client/audio/opus/OpusCodec.java b/ts3-client/core/src/main/java/com/ts3client/audio/opus/OpusCodec.java new file mode 100644 index 0000000..15a5661 --- /dev/null +++ b/ts3-client/core/src/main/java/com/ts3client/audio/opus/OpusCodec.java @@ -0,0 +1,16 @@ +package com.ts3client.audio.opus; + +/** Creates Opus encoders and decoders; each platform binds its own libopus. */ +public interface OpusCodec { + + /** Opus' coding mode, fixed when an encoder is created. */ + enum Application { VOIP, AUDIO } + + OpusEncoder createEncoder(int sampleRate, int frameSize, int channels, Application application); + + /** @param maxFrameSize the longest frame, per channel, a packet may decode to */ + OpusDecoder createDecoder(int sampleRate, int maxFrameSize, int channels); + + /** libopus' version string, for display. */ + String version(); +} diff --git a/ts3-client/core/src/main/java/com/ts3client/audio/opus/OpusDecoder.java b/ts3-client/core/src/main/java/com/ts3client/audio/opus/OpusDecoder.java new file mode 100644 index 0000000..970ca58 --- /dev/null +++ b/ts3-client/core/src/main/java/com/ts3client/audio/opus/OpusDecoder.java @@ -0,0 +1,17 @@ +package com.ts3client.audio.opus; + +/** One Opus decoder, used from one thread at a time. */ +public interface OpusDecoder extends AutoCloseable { + + /** + * Decodes {@code packet} into interleaved samples, or conceals a lost packet when it is + * {@code null}. Returns the samples decoded per channel. + */ + int decode(byte[] packet, float[] out); + + /** Forgets the stream's state, as at the start of a new talk burst. */ + void reset(); + + @Override + void close(); +} diff --git a/ts3-client/core/src/main/java/com/ts3client/audio/opus/OpusEncoder.java b/ts3-client/core/src/main/java/com/ts3client/audio/opus/OpusEncoder.java new file mode 100644 index 0000000..2b15e49 --- /dev/null +++ b/ts3-client/core/src/main/java/com/ts3client/audio/opus/OpusEncoder.java @@ -0,0 +1,25 @@ +package com.ts3client.audio.opus; + +/** One Opus encoder. Its methods may be called from different threads, but not concurrently. */ +public interface OpusEncoder extends AutoCloseable { + + enum Signal { VOICE, MUSIC } + + void setBitrate(int bitsPerSecond); + + void setComplexity(int complexity); + + void setVbr(boolean vbr); + + void setInbandFec(boolean fec); + + void setExpectedPacketLoss(int percent); + + void setSignal(Signal signal); + + /** Encodes exactly one frame of interleaved samples. */ + byte[] encode(float[] pcm); + + @Override + void close(); +} diff --git a/ts3-client/core/src/test/java/com/ts3client/audio/FakeOpus.java b/ts3-client/core/src/test/java/com/ts3client/audio/FakeOpus.java new file mode 100644 index 0000000..7388f89 --- /dev/null +++ b/ts3-client/core/src/test/java/com/ts3client/audio/FakeOpus.java @@ -0,0 +1,53 @@ +package com.ts3client.audio; + +import com.ts3client.audio.opus.OpusCodec; +import com.ts3client.audio.opus.OpusDecoder; +import com.ts3client.audio.opus.OpusEncoder; + +import java.util.Arrays; + +/** + * Stands in for libopus: a "packet" is one byte, decoded to a 20 ms frame holding that byte + * divided by 100 in every sample. Concealment decodes to -1, so it is easy to spot. + */ +final class FakeOpus implements OpusCodec { + + static final float CONCEALED = -1f; + + int resets; + + static byte[] packet(int value) { + return new byte[]{(byte) value}; + } + + @Override + public OpusEncoder createEncoder(int sampleRate, int frameSize, int channels, Application application) { + throw new UnsupportedOperationException(); + } + + @Override + public OpusDecoder createDecoder(int sampleRate, int maxFrameSize, int channels) { + return new OpusDecoder() { + @Override + public int decode(byte[] packet, float[] out) { + float v = packet == null ? CONCEALED : packet[0] / 100f; + Arrays.fill(out, 0, VoiceFormat.FRAME_SIZE * channels, v); + return VoiceFormat.FRAME_SIZE; + } + + @Override + public void reset() { + resets++; + } + + @Override + public void close() { + } + }; + } + + @Override + public String version() { + return "fake"; + } +} diff --git a/ts3-client/core/src/test/java/com/ts3client/audio/MixedPlayoutTest.java b/ts3-client/core/src/test/java/com/ts3client/audio/MixedPlayoutTest.java new file mode 100644 index 0000000..ded30ea --- /dev/null +++ b/ts3-client/core/src/test/java/com/ts3client/audio/MixedPlayoutTest.java @@ -0,0 +1,102 @@ +package com.ts3client.audio; + +import com.github.manevolent.ts3j.enums.CodecType; +import com.github.manevolent.ts3j.protocol.ProtocolRole; +import com.github.manevolent.ts3j.protocol.packet.PacketBody0Voice; +import org.junit.jupiter.api.Test; + +import java.util.List; +import java.util.concurrent.BlockingQueue; +import java.util.concurrent.LinkedBlockingQueue; +import java.util.concurrent.Semaphore; +import java.util.concurrent.TimeUnit; + +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertTrue; + +class MixedPlayoutTest { + + /** A mono line that records every frame written and blocks until the test lets it go on. */ + private static final class RecordingLine implements AudioPlayback { + final BlockingQueue firstSamples = new LinkedBlockingQueue<>(); + final Semaphore permits = new Semaphore(0); + + @Override + public int channels() { + return 1; + } + + @Override + public void start() { + } + + @Override + public void write(byte[] buffer, int offset, int length) { + firstSamples.add((short) ((buffer[offset + 1] << 8) | (buffer[offset] & 0xFF))); + permits.acquireUninterruptibly(); + } + + @Override + public void drain() { + } + + @Override + public void close() { + } + } + + @Test + void mixesSimultaneousSpeakersOntoOneLine() throws Exception { + RecordingLine line = new RecordingLine(); + AudioIo io = new AudioIo() { + @Override + public List inputDevices() { + return List.of(); + } + + @Override + public List outputDevices() { + return List.of(); + } + + @Override + public AudioCapture openCapture(String deviceId, int preferredChannels) { + throw new UnsupportedOperationException(); + } + + @Override + public AudioPlayback openPlayback(String deviceId, int preferredChannels) { + return line; + } + }; + StreamingVoiceOutput output = new StreamingVoiceOutput(io, new FakeOpus(), "", StreamingVoiceOutput.Lines.SHARED); + try { + for (int i = 0; i < 10; i++) { + output.handleVoice(voice(1, i, 25)); + output.handleVoice(voice(2, i, 50)); + } + line.permits.release(100); + + short expected = (short) Math.round(0.75f * 32767f); + boolean mixed = false; + for (int i = 0; i < 20 && !mixed; i++) { + Short sample = line.firstSamples.poll(1, TimeUnit.SECONDS); + assertNotNull(sample, "the mixer stopped writing"); + mixed = sample == expected; + } + assertTrue(mixed, "no frame carried both speakers"); + } finally { + output.shutdown(); + line.permits.release(1000); + } + } + + private static PacketBody0Voice voice(int clientId, int packetId, int value) { + PacketBody0Voice voice = new PacketBody0Voice(ProtocolRole.SERVER); + voice.setClientId(clientId); + voice.setPacketId(packetId); + voice.setCodecType(CodecType.OPUS_VOICE); + voice.setCodecData(FakeOpus.packet(value)); + return voice; + } +} diff --git a/ts3-client/core/src/test/java/com/ts3client/audio/VoiceStreamTest.java b/ts3-client/core/src/test/java/com/ts3client/audio/VoiceStreamTest.java new file mode 100644 index 0000000..072699e --- /dev/null +++ b/ts3-client/core/src/test/java/com/ts3client/audio/VoiceStreamTest.java @@ -0,0 +1,111 @@ +package com.ts3client.audio; + +import com.github.manevolent.ts3j.enums.CodecType; +import org.junit.jupiter.api.Test; + +import java.util.ArrayList; +import java.util.List; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +class VoiceStreamTest { + + private final FakeOpus opus = new FakeOpus(); + private final List talk = new ArrayList<>(); + private final VoiceStream stream = new VoiceStream(7, opus, (id, talking) -> talk.add(talking)); + private final float[] out = new float[VoiceStream.MAX_FRAME * 2]; + + private boolean offer(int packetId, int value) { + return stream.offer(packetId, CodecType.OPUS_VOICE, FakeOpus.packet(value)); + } + + /** Pulls one frame and returns its first sample, or NaN if nothing was played. */ + private float pull() { + int frames = stream.pull(out, 1, 1.0); + return frames > 0 ? out[0] : Float.NaN; + } + + @Test + void holdsBackTheStartOfABurstThenPlaysInOrder() { + assertTrue(offer(10, 1)); + assertFalse(offer(12, 3)); + assertFalse(offer(11, 2)); + + assertEquals(Float.NaN, pull()); + assertEquals(Float.NaN, pull()); + assertEquals(0.01f, pull()); + assertEquals(0.02f, pull()); + assertEquals(0.03f, pull()); + assertEquals(List.of(true), talk); + } + + @Test + void concealsALostPacket() { + offer(0, 1); + offer(2, 3); + pull(); + pull(); + assertEquals(0.01f, pull()); + assertEquals(FakeOpus.CONCEALED, pull()); + assertEquals(0.03f, pull()); + } + + @Test + void dropsPacketsThatArriveTooLate() { + offer(0, 1); + offer(1, 2); + pull(); + pull(); + pull(); + pull(); + assertFalse(offer(0, 9)); + assertEquals(FakeOpus.CONCEALED, pull()); + } + + @Test + void endPacketFinishesTheBurst() { + offer(0, 1); + stream.offer(1, CodecType.OPUS_VOICE, new byte[0]); + pull(); + pull(); + pull(); + assertEquals(VoiceStream.IDLE, stream.pull(out, 1, 1.0)); + assertFalse(stream.isPlaying()); + assertEquals(List.of(true, false), talk); + assertEquals(1, opus.resets); + assertTrue(offer(2, 1), "the next packet starts a new burst"); + } + + @Test + void givesUpAfterTooMuchConcealment() { + offer(0, 1); + pull(); + pull(); + pull(); + int concealed = 0; + while (stream.pull(out, 1, 1.0) != VoiceStream.IDLE) concealed++; + assertEquals(10, concealed); + assertEquals(List.of(true, false), talk); + } + + @Test + void stopFromAnotherThreadEndsTheBurst() { + offer(0, 1); + stream.stop(); + assertEquals(VoiceStream.IDLE, stream.pull(out, 1, 1.0)); + assertEquals(List.of(true, false), talk); + } + + @Test + void spreadsMonoOverStereoWithGain() { + offer(0, 50); + stream.pull(out, 2, 0.5); + stream.pull(out, 2, 0.5); + assertEquals(VoiceFormat.FRAME_SIZE, stream.pull(out, 2, 0.5)); + assertEquals(0.25f, out[0]); + assertEquals(0.25f, out[1]); + assertEquals(0.25f, out[2 * VoiceFormat.FRAME_SIZE - 1]); + } +} diff --git a/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/AudioDevices.java b/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/AudioDevices.java index cc405d9..3ed2a7f 100644 --- a/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/AudioDevices.java +++ b/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/AudioDevices.java @@ -1,5 +1,10 @@ package com.ts3client.audio.desktop; +import com.ts3client.audio.AudioCapture; +import com.ts3client.audio.AudioDevice; +import com.ts3client.audio.AudioIo; +import com.ts3client.audio.AudioPlayback; +import com.ts3client.audio.VoiceFormat; import com.ts3client.audio.desktop.pipewire.PipeWire; import javax.sound.sampled.AudioFormat; @@ -36,11 +41,9 @@ import java.util.List; * {@link #openCapture} / {@link #openPlayback}, which fall back through the supported * channel counts. */ -public final class AudioDevices { +public final class AudioDevices implements AudioIo { - public static final int SAMPLE_RATE = 48_000; - public static final int FRAME_SIZE = 960; // 20 ms @ 48 kHz - public static final int MAX_CHANNELS = 2; + public static final AudioDevices INSTANCE = new AudioDevices(); /** Buffer size, in 20 ms frames, requested when opening a Java Sound line. */ private static final int BUFFER_FRAMES = 8; @@ -48,42 +51,31 @@ public final class AudioDevices { /** Marks a device id as a PipeWire node name rather than a Java Sound mixer name. */ private static final String PIPEWIRE_PREFIX = "pw:"; - /** A selectable audio device. {@code id} is what gets persisted in the settings. */ - public record Device(String id, String label) { - /** The implicit entry that lets the platform (PipeWire, Pulse, ALSA) choose. */ - public static final Device DEFAULT = new Device("", "(System default)"); - - @Override - public String toString() { - return label; - } - } - private AudioDevices() { } /** 48 kHz signed 16-bit little-endian PCM with the given channel count. */ public static AudioFormat format(int channels) { - return new AudioFormat(SAMPLE_RATE, 16, channels, true, false); + return new AudioFormat(VoiceFormat.SAMPLE_RATE, 16, channels, true, false); } - /** Devices that can provide microphone (capture) lines, system default first. */ - public static List inputDevices() { + @Override + public List inputDevices() { return devices(TargetDataLine.class, false); } - /** Devices that can provide speaker (playback) lines, system default first. */ - public static List outputDevices() { + @Override + public List outputDevices() { return devices(SourceDataLine.class, true); } - private static List devices(Class lineClass, boolean sinks) { - List devices = new ArrayList<>(); - devices.add(Device.DEFAULT); + private static List devices(Class lineClass, boolean sinks) { + List devices = new ArrayList<>(); + devices.add(AudioDevice.DEFAULT); for (PipeWire.Node node : PipeWire.nodes()) { if (node.sink() == sinks) { - devices.add(new Device(PIPEWIRE_PREFIX + node.name(), node.description())); + devices.add(new AudioDevice(PIPEWIRE_PREFIX + node.name(), node.description())); } } @@ -95,7 +87,7 @@ public final class AudioDevices { if (isDefaultMixer(name)) continue; if (!AudioSystem.getMixer(mi).isLineSupported(anyFormat)) continue; if (devices.stream().anyMatch(d -> d.id().equals(name))) continue; - devices.add(new Device(name, "ALSA: " + name)); + devices.add(new AudioDevice(name, "ALSA: " + name)); } return devices; } @@ -109,17 +101,13 @@ public final class AudioDevices { return mixerName.endsWith("[default]"); } - /** - * Opens a capture line, preferring {@code preferredChannels} and falling back to the - * other channel count if the device won't take it. Inspect {@link AudioCapture#channels()} - * for what was actually opened. - */ - public static AudioCapture openCapture(String deviceId, int preferredChannels) + @Override + public AudioCapture openCapture(String deviceId, int preferredChannels) throws LineUnavailableException { if (usePipeWire(deviceId)) { try { return PipeWire.openCapture("TS3J Microphone", pipeWireNode(deviceId), - SAMPLE_RATE, clampChannels(preferredChannels)); + VoiceFormat.SAMPLE_RATE, clampChannels(preferredChannels)); } catch (Exception e) { if (isPipeWireDevice(deviceId)) throw unavailable(true, deviceId, e); // The default device: Java Sound may still reach it through ALSA. @@ -128,13 +116,13 @@ public final class AudioDevices { return new JavaSoundCapture(open(TargetDataLine.class, deviceId, preferredChannels)); } - /** Opens a playback line; see {@link #openCapture} for the channel negotiation. */ - public static AudioPlayback openPlayback(String deviceId, int preferredChannels) + @Override + public AudioPlayback openPlayback(String deviceId, int preferredChannels) throws LineUnavailableException { if (usePipeWire(deviceId)) { try { return PipeWire.openPlayback("TS3J Playback", pipeWireNode(deviceId), - SAMPLE_RATE, clampChannels(preferredChannels)); + VoiceFormat.SAMPLE_RATE, clampChannels(preferredChannels)); } catch (Exception e) { if (isPipeWireDevice(deviceId)) throw unavailable(false, deviceId, e); } @@ -167,7 +155,7 @@ public final class AudioDevices { /** PipeWire converts freely, so any channel count works; keep it in the range we handle. */ private static int clampChannels(int channels) { - return Math.max(1, Math.min(MAX_CHANNELS, channels)); + return Math.max(1, Math.min(VoiceFormat.MAX_CHANNELS, channels)); } private static T open(Class lineClass, String deviceId, int preferredChannels) @@ -180,7 +168,7 @@ public final class AudioDevices { try { T line = lineClass.cast( (mixer != null) ? mixer.getLine(info) : AudioSystem.getLine(info)); - openLine(line, fmt, FRAME_SIZE * 2 * fmt.getChannels() * BUFFER_FRAMES); + openLine(line, fmt, VoiceFormat.FRAME_SIZE * 2 * fmt.getChannels() * BUFFER_FRAMES); return line; } catch (Exception e) { failure = e; diff --git a/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/DesktopAudioBackend.java b/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/DesktopAudioBackend.java index f01f494..71e6f8f 100644 --- a/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/DesktopAudioBackend.java +++ b/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/DesktopAudioBackend.java @@ -1,9 +1,10 @@ package com.ts3client.audio.desktop; import com.ts3client.audio.AudioBackend; -import com.ts3client.audio.VoiceInput; -import com.ts3client.audio.VoiceOutput; +import com.ts3client.audio.AudioIo; +import com.ts3client.audio.StreamingVoiceOutput; import com.ts3client.audio.desktop.pipewire.PipeWire; +import com.ts3client.audio.opus.OpusCodec; import com.ts3client.config.Settings; import com.ts3client.sound.SoundPlayer; @@ -13,19 +14,27 @@ import com.ts3client.sound.SoundPlayer; */ public final class DesktopAudioBackend implements AudioBackend { + private final OpusCodec opus = new NativeOpusCodec(); + @Override - public VoiceInput createInput(Settings settings) { - return new DesktopVoiceInput(settings); + public AudioIo io() { + return AudioDevices.INSTANCE; } @Override - public VoiceOutput createOutput(Settings settings) { - return new DesktopVoiceOutput(settings.outputDevice); + public OpusCodec opus() { + return opus; + } + + /** Every speaker gets a line of their own, so each shows up in the desktop's mixer. */ + @Override + public StreamingVoiceOutput.Lines outputLines() { + return StreamingVoiceOutput.Lines.PER_SPEAKER; } @Override public SoundPlayer createSoundPlayer(Settings settings) { - return new WavSoundPlayer(settings.outputDevice); + return new WavSoundPlayer(io(), settings.outputDevice); } @Override @@ -33,9 +42,9 @@ public final class DesktopAudioBackend implements AudioBackend { return codec() + ", " + audioSystem(); } - private static String codec() { + private String codec() { try { - return "Opus " + Opus.getVersionString(); + return "Opus " + opus.version(); } catch (Throwable t) { return "Opus (native library unavailable)"; } diff --git a/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/DesktopVoiceOutput.java b/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/DesktopVoiceOutput.java deleted file mode 100644 index 560df61..0000000 --- a/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/DesktopVoiceOutput.java +++ /dev/null @@ -1,404 +0,0 @@ -package com.ts3client.audio.desktop; - -import com.github.manevolent.ts3j.enums.CodecType; -import com.github.manevolent.ts3j.protocol.packet.PacketBody0Voice; -import com.github.manevolent.ts3j.protocol.packet.PacketBody1VoiceWhisper; -import com.ts3client.audio.VoiceOutput; - -import java.util.HashMap; -import java.util.Map; -import java.util.Set; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; -import java.util.concurrent.ScheduledFuture; -import java.util.concurrent.ScheduledExecutorService; -import java.util.concurrent.TimeUnit; -import java.util.concurrent.atomic.AtomicBoolean; -import java.util.function.BiConsumer; - -/** - * Desktop {@link VoiceOutput}: decodes and plays incoming voice per speaker. - * Each client gets its own Opus decoder, playback line and worker thread, so - * simultaneous speakers are mixed by the audio server and one slow decode never - * blocks another (or the network thread). Lines come from {@link AudioDevices}, - * so on a PipeWire desktop every speaker is a separate stream in the mixer. - */ -public final class DesktopVoiceOutput implements VoiceOutput { - - /** Longest Opus frame (120 ms @ 48 kHz) a packet may decode to, per channel. */ - private static final int MAX_FRAME = 5760; - /** TS3 voice packets normally carry one 20 ms Opus frame. */ - private static final long FRAME_NANOS = TimeUnit.MILLISECONDS.toNanos(20); - /** Initial playout delay: enough room for common UDP jitter without feeling laggy. */ - private static final long JITTER_DELAY_NANOS = FRAME_NANOS * 3; - /** Cap PLC during a lost end-of-talk marker or a hard network stall. */ - private static final int MAX_CONSECUTIVE_PLC_FRAMES = 10; - /** Keep memory bounded when a talk burst runs far ahead of playback. */ - private static final int MAX_PENDING_PACKETS = 32; - - /** - * A speaker is only meant to stop when the empty voice packet marking the end of a - * talk burst arrives — but that packet is UDP too, and a lost one otherwise leaves - * the talking indicator stuck until the speaker's next burst. Opus frames are 20 ms - * apart while someone is actually talking, so a gap several times that long, with no - * packet of either kind, is unambiguous: force the indicator off rather than trust - * the one packet that could go missing. - */ - private static final long TALK_TIMEOUT_NANOS = TimeUnit.MILLISECONDS.toNanos(400); - - /** One speaker's decode + playback pipeline. */ - private final class ClientStream { - final Object lock = new Object(); - final int clientId; - final AudioPlayback line; - final int lineChannels; - final ExecutorService worker; - final Map pending = new HashMap<>(); - final AtomicBoolean tickQueued = new AtomicBoolean(); - /** Decoder scratch and byte buffer, touched only by {@link #worker}. */ - final float[] pcm = new float[MAX_FRAME * AudioDevices.MAX_CHANNELS]; - byte[] out = new byte[0]; - OpusDecoder decoder; - int decoderChannels; - ScheduledFuture playout; - int expectedPacketId = -1; - int activeChannels = 1; - int consecutivePlcFrames; - volatile boolean talking; - volatile long lastPacketNanos; - - ClientStream(int clientId) throws Exception { - this.clientId = clientId; - this.line = AudioDevices.openPlayback(outputDevice, AudioDevices.MAX_CHANNELS); - this.lineChannels = line.channels(); - this.line.start(); - this.worker = Executors.newSingleThreadExecutor(r -> { - Thread t = new Thread(r, "ts3j-play-" + clientId); - t.setDaemon(true); - return t; - }); - } - - /** - * Returns a decoder matching the stream's channel count, rebuilding it when a - * speaker switches between the mono {@code OPUS_VOICE} and stereo - * {@code OPUS_MUSIC} codecs. - */ - OpusDecoder decoderFor(int channels) { - if (decoder == null || decoderChannels != channels) { - if (decoder != null) decoder.close(); - decoder = new OpusDecoder(AudioDevices.SAMPLE_RATE, MAX_FRAME, channels); - decoderChannels = channels; - } - return decoder; - } - - void close() { - synchronized (lock) { - if (playout != null) playout.cancel(false); - pending.clear(); - } - worker.shutdownNow(); - line.close(); - if (decoder != null) decoder.close(); - } - } - - private record VoiceFrame(int packetId, CodecType codec, byte[] data, boolean end) { - } - - private final Map streams = new ConcurrentHashMap<>(); - private final Set mutedClients = ConcurrentHashMap.newKeySet(); - /** Per-client gain on top of the master volume; absent means 1.0. */ - private final Map clientVolumes = new ConcurrentHashMap<>(); - - private volatile String outputDevice; - private volatile double masterVolume = 1.0; - private volatile boolean deafened = false; - - /** Notified (clientId, talking) on the EDT-agnostic worker thread when a speaker starts/stops. */ - private volatile BiConsumer talkListener; - - /** Catches a talk burst whose end packet never arrived; see {@link #TALK_TIMEOUT_NANOS}. */ - private final ScheduledExecutorService watchdog = Executors.newSingleThreadScheduledExecutor(r -> { - Thread t = new Thread(r, "ts3j-talk-watchdog"); - t.setDaemon(true); - return t; - }); - - public DesktopVoiceOutput(String outputDevice) { - this.outputDevice = outputDevice; - watchdog.scheduleWithFixedDelay(this::checkTalkTimeouts, 50, 50, TimeUnit.MILLISECONDS); - } - - private void checkTalkTimeouts() { - long now = System.nanoTime(); - for (ClientStream s : streams.values()) { - if (s.talking && now - s.lastPacketNanos > TALK_TIMEOUT_NANOS) { - s.worker.submit(() -> stopStream(s)); - } - } - } - - public void setTalkListener(BiConsumer l) { - this.talkListener = l; - } - - public void setMasterVolume(double v) { - this.masterVolume = Math.max(0, Math.min(2.0, v)); - } - - public void setDeafened(boolean d) { - this.deafened = d; - if (d) { - // Stop everyone talking immediately. - for (ClientStream s : streams.values()) { - s.worker.submit(() -> stopStream(s)); - } - } - } - - public boolean isDeafened() { - return deafened; - } - - public void setClientMuted(int clientId, boolean muted) { - if (muted) { - mutedClients.add(clientId); - ClientStream s = streams.get(clientId); - if (s != null) s.worker.submit(() -> stopStream(s)); - } else { - mutedClients.remove(clientId); - } - } - - public boolean isClientMuted(int clientId) { - return mutedClients.contains(clientId); - } - - public void setClientVolume(int clientId, double gain) { - if (gain == 1.0) clientVolumes.remove(clientId); - else clientVolumes.put(clientId, Math.max(0, gain)); - } - - public double getClientVolume(int clientId) { - return clientVolumes.getOrDefault(clientId, 1.0); - } - - public void setOutputDevice(String device) { - this.outputDevice = device; - } - - /** Entry point wired into {@code client.setVoiceHandler(...)}. */ - public void handleVoice(PacketBody0Voice voice) { - route(voice.getClientId(), voice.getPacketId(), voice.getCodecType(), voice.getCodecData()); - } - - /** Entry point wired into {@code client.setWhisperHandler(...)}. */ - public void handleWhisper(PacketBody1VoiceWhisper whisper) { - route(whisper.getClientId(), whisper.getPacketId(), whisper.getCodecType(), whisper.getCodecData()); - } - - /** TeamSpeak streams {@code OPUS_MUSIC} in stereo and everything else in mono. */ - private static int channelsFor(CodecType codec) { - return codec == CodecType.OPUS_MUSIC ? 2 : 1; - } - - private void route(int clientId, int packetId, CodecType codec, byte[] data) { - if (deafened) return; - if (mutedClients.contains(clientId)) return; - - ClientStream stream = streams.get(clientId); - if (stream == null) { - try { - stream = new ClientStream(clientId); - ClientStream existing = streams.putIfAbsent(clientId, stream); - if (existing != null) { - stream.close(); - stream = existing; - } - } catch (Exception e) { - return; // couldn't open a line; drop - } - } - - final ClientStream target = stream; - target.lastPacketNanos = System.nanoTime(); - VoiceFrame frame = new VoiceFrame(packetId & 0xFFFF, codec, data, data == null || data.length == 0); - - synchronized (target.lock) { - if (target.expectedPacketId < 0 || target.playout == null || target.playout.isDone()) { - target.expectedPacketId = frame.packetId(); - target.activeChannels = channelsFor(codec); - target.consecutivePlcFrames = 0; - target.playout = watchdog.scheduleAtFixedRate( - () -> queuePlayoutTick(target), - JITTER_DELAY_NANOS, - FRAME_NANOS, - TimeUnit.NANOSECONDS); - } else if (packetDistance(target.expectedPacketId, frame.packetId()) < 0) { - return; // too late for this burst; do not replay stale audio - } - - target.pending.putIfAbsent(frame.packetId(), frame); - trimPending(target); - } - if (!frame.end()) markTalking(target, true); - } - - private void queuePlayoutTick(ClientStream stream) { - if (!stream.tickQueued.compareAndSet(false, true)) return; - try { - stream.worker.submit(() -> { - try { - playNext(stream); - } finally { - stream.tickQueued.set(false); - } - }); - } catch (RuntimeException e) { - if (stream.worker.isShutdown()) { - stream.tickQueued.set(false); - } else { - throw e; - } - } - } - - private void playNext(ClientStream stream) { - VoiceFrame frame; - int channels; - boolean conceal; - - synchronized (stream.lock) { - if (stream.expectedPacketId < 0) return; - - frame = stream.pending.remove(stream.expectedPacketId); - if (frame != null && frame.end()) { - stopPlayoutLocked(stream); - return; - } - - conceal = frame == null; - channels = frame != null ? channelsFor(frame.codec()) : stream.activeChannels; - stream.activeChannels = channels; - - if (conceal) { - stream.consecutivePlcFrames++; - if (stream.consecutivePlcFrames > MAX_CONSECUTIVE_PLC_FRAMES) { - stopPlayoutLocked(stream); - return; - } - } else { - stream.consecutivePlcFrames = 0; - } - - stream.expectedPacketId = (stream.expectedPacketId + 1) & 0xFFFF; - } - - decodeAndPlay(stream, channels, conceal ? null : frame.data()); - } - - private void stopStream(ClientStream stream) { - synchronized (stream.lock) { - stopPlayoutLocked(stream); - } - } - - private void stopPlayoutLocked(ClientStream stream) { - stream.expectedPacketId = -1; - if (stream.playout != null) stream.playout.cancel(false); - stream.playout = null; - stream.pending.clear(); - finishStream(stream); - } - - private void finishStream(ClientStream stream) { - if (stream.decoder != null) stream.decoder.reset(); - markTalking(stream, false); - } - - private static int packetDistance(int from, int to) { - int distance = (to - from) & 0xFFFF; - return distance >= 0x8000 ? distance - 0x10000 : distance; - } - - private static void trimPending(ClientStream stream) { - while (stream.pending.size() > MAX_PENDING_PACKETS) { - Integer furthest = null; - int furthestDistance = Integer.MIN_VALUE; - for (Integer packetId : stream.pending.keySet()) { - int distance = packetDistance(stream.expectedPacketId, packetId); - if (distance > furthestDistance) { - furthestDistance = distance; - furthest = packetId; - } - } - if (furthest == null) return; - stream.pending.remove(furthest); - } - } - - private void decodeAndPlay(ClientStream stream, int channels, byte[] data) { - try { - float[] pcm = stream.pcm; - int frames = stream.decoderFor(channels).decode(data, pcm); - double vol = masterVolume * clientVolumes.getOrDefault(stream.clientId, 1.0); - - // Match the decoded stream to the line: duplicate mono across a stereo - // line, fold a stereo (music) stream down onto a mono-only line. - int lineChannels = stream.lineChannels; - int bytes = frames * lineChannels * 2; - if (stream.out.length < bytes) stream.out = new byte[bytes]; - byte[] out = stream.out; - - for (int i = 0, k = 0; i < frames; i++) { - for (int c = 0; c < lineChannels; c++, k += 2) { - double v; - if (channels == lineChannels) { - v = pcm[i * channels + c]; - } else if (channels == 1) { - v = pcm[i]; - } else { - double sum = 0; - for (int s = 0; s < channels; s++) sum += pcm[i * channels + s]; - v = sum / channels; - } - v *= vol; - if (v > 1.0) v = 1.0; - else if (v < -1.0) v = -1.0; - short s = (short) Math.round(v * 32767.0); - out[k] = (byte) (s & 0xFF); - out[k + 1] = (byte) ((s >> 8) & 0xFF); - } - } - stream.line.write(out, 0, bytes); - } catch (Exception ignored) { - } - } - - private void markTalking(ClientStream stream, boolean talking) { - if (stream.talking == talking) return; - stream.talking = talking; - BiConsumer l = talkListener; - if (l != null) l.accept(stream.clientId, talking); - } - - /** Drop a speaker's pipeline entirely (e.g. they left the server). */ - /** Forgets a client that left; the server may hand its id to somebody else. */ - public void removeClient(int clientId) { - ClientStream s = streams.remove(clientId); - if (s != null) s.close(); - mutedClients.remove(clientId); - clientVolumes.remove(clientId); - } - - public void shutdown() { - watchdog.shutdownNow(); - for (ClientStream s : streams.values()) { - s.close(); - } - streams.clear(); - } -} diff --git a/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/JavaSoundCapture.java b/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/JavaSoundCapture.java index a6187ed..0f733d7 100644 --- a/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/JavaSoundCapture.java +++ b/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/JavaSoundCapture.java @@ -1,5 +1,7 @@ package com.ts3client.audio.desktop; +import com.ts3client.audio.AudioCapture; + import javax.sound.sampled.TargetDataLine; /** {@link AudioCapture} backed by a Java Sound {@link TargetDataLine}. */ diff --git a/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/JavaSoundPlayback.java b/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/JavaSoundPlayback.java index a81ae2c..49eddf1 100644 --- a/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/JavaSoundPlayback.java +++ b/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/JavaSoundPlayback.java @@ -1,5 +1,7 @@ package com.ts3client.audio.desktop; +import com.ts3client.audio.AudioPlayback; + import javax.sound.sampled.SourceDataLine; /** {@link AudioPlayback} backed by a Java Sound {@link SourceDataLine}. */ diff --git a/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/NativeOpusCodec.java b/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/NativeOpusCodec.java new file mode 100644 index 0000000..143c851 --- /dev/null +++ b/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/NativeOpusCodec.java @@ -0,0 +1,25 @@ +package com.ts3client.audio.desktop; + +import com.ts3client.audio.opus.OpusCodec; +import com.ts3client.audio.opus.OpusDecoder; +import com.ts3client.audio.opus.OpusEncoder; + +/** The system's (or the bundled) libopus, reached through the FFM API. */ +public final class NativeOpusCodec implements OpusCodec { + + @Override + public OpusEncoder createEncoder(int sampleRate, int frameSize, int channels, Application application) { + int app = application == Application.AUDIO ? Opus.OPUS_APPLICATION_AUDIO : Opus.OPUS_APPLICATION_VOIP; + return new NativeOpusEncoder(sampleRate, frameSize, channels, app); + } + + @Override + public OpusDecoder createDecoder(int sampleRate, int maxFrameSize, int channels) { + return new NativeOpusDecoder(sampleRate, maxFrameSize, channels); + } + + @Override + public String version() { + return Opus.getVersionString(); + } +} diff --git a/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/OpusDecoder.java b/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/NativeOpusDecoder.java similarity index 93% rename from ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/OpusDecoder.java rename to ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/NativeOpusDecoder.java index a87b16e..55ba7db 100644 --- a/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/OpusDecoder.java +++ b/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/NativeOpusDecoder.java @@ -1,5 +1,7 @@ package com.ts3client.audio.desktop; +import com.ts3client.audio.opus.OpusDecoder; + import java.lang.foreign.Arena; import java.lang.foreign.MemorySegment; import java.lang.foreign.ValueLayout; @@ -11,7 +13,7 @@ import java.lang.foreign.ValueLayout; * (mono for {@code OPUS_VOICE}, stereo for {@code OPUS_MUSIC}); Opus will up-/down-mix * a mismatched stream, which costs the stereo image. */ -public final class OpusDecoder implements AutoCloseable { +final class NativeOpusDecoder implements OpusDecoder { private static final int MAX_PACKET_BYTES = 4096; @@ -23,7 +25,7 @@ public final class OpusDecoder implements AutoCloseable { private final int channels; private boolean closed; - public OpusDecoder(int sampleRate, int frameSize, int channels) { + NativeOpusDecoder(int sampleRate, int frameSize, int channels) { this.frameSize = frameSize; this.channels = channels; @@ -48,6 +50,7 @@ public final class OpusDecoder implements AutoCloseable { * @param out output buffer, at least {@code frameSize * channels} long * @return number of samples decoded per channel */ + @Override public int decode(byte[] packet, float[] out) { if (closed) throw new IllegalStateException("decoder closed"); @@ -70,15 +73,12 @@ public final class OpusDecoder implements AutoCloseable { return samples; } + @Override public void reset() { if (closed) return; Opus.decoderCtl(handle, Opus.OPUS_RESET_STATE, 0); } - public int getChannels() { - return channels; - } - @Override public void close() { if (closed) return; diff --git a/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/OpusEncoder.java b/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/NativeOpusEncoder.java similarity index 89% rename from ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/OpusEncoder.java rename to ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/NativeOpusEncoder.java index 530e4d6..d0fe712 100644 --- a/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/OpusEncoder.java +++ b/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/NativeOpusEncoder.java @@ -1,5 +1,7 @@ package com.ts3client.audio.desktop; +import com.ts3client.audio.opus.OpusEncoder; + import java.lang.foreign.Arena; import java.lang.foreign.MemorySegment; import java.lang.foreign.ValueLayout; @@ -14,7 +16,7 @@ import java.lang.foreign.ValueLayout; * are freed by {@link #close()}; access is serialised by the instance lock, so the * encoder may be driven from any thread. */ -public final class OpusEncoder implements AutoCloseable { +final class NativeOpusEncoder implements OpusEncoder { private static final int MAX_PACKET_BYTES = 4096; @@ -27,7 +29,7 @@ public final class OpusEncoder implements AutoCloseable { private final Object lock = new Object(); private boolean closed; - public OpusEncoder(int sampleRate, int frameSize, int channels, int application) { + NativeOpusEncoder(int sampleRate, int frameSize, int channels, int application) { this.frameSize = frameSize; this.channels = channels; @@ -44,28 +46,34 @@ public final class OpusEncoder implements AutoCloseable { this.packetBuffer = arena.allocate(MAX_PACKET_BYTES); } + @Override public void setBitrate(int bitsPerSecond) { ctl(Opus.OPUS_SET_BITRATE_REQUEST, bitsPerSecond); } + @Override public void setComplexity(int complexity) { ctl(Opus.OPUS_SET_COMPLEXITY_REQUEST, Math.max(0, Math.min(10, complexity))); } + @Override public void setVbr(boolean vbr) { ctl(Opus.OPUS_SET_VBR_REQUEST, vbr ? 1 : 0); } + @Override public void setInbandFec(boolean fec) { ctl(Opus.OPUS_SET_INBAND_FEC_REQUEST, fec ? 1 : 0); } + @Override public void setExpectedPacketLoss(int percent) { ctl(Opus.OPUS_SET_PACKET_LOSS_PERC_REQUEST, Math.max(0, Math.min(100, percent))); } - public void setSignal(int signal) { - ctl(Opus.OPUS_SET_SIGNAL_REQUEST, signal); + @Override + public void setSignal(Signal signal) { + ctl(Opus.OPUS_SET_SIGNAL_REQUEST, signal == Signal.MUSIC ? Opus.OPUS_SIGNAL_MUSIC : Opus.OPUS_SIGNAL_VOICE); } private void ctl(int request, int value) { @@ -85,6 +93,7 @@ public final class OpusEncoder implements AutoCloseable { * * @return a newly allocated byte array holding the encoded packet */ + @Override public byte[] encode(float[] pcm) { int expected = frameSize * channels; if (pcm.length != expected) { diff --git a/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/WavSoundPlayer.java b/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/WavSoundPlayer.java index 643cc58..8ca742b 100644 --- a/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/WavSoundPlayer.java +++ b/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/WavSoundPlayer.java @@ -1,5 +1,8 @@ package com.ts3client.audio.desktop; +import com.ts3client.audio.AudioIo; +import com.ts3client.audio.AudioPlayback; +import com.ts3client.audio.VoiceFormat; import com.ts3client.sound.SoundPlayer; import javax.sound.sampled.AudioFormat; @@ -28,7 +31,7 @@ public final class WavSoundPlayer implements SoundPlayer { * triggered, so it is kept around for a spell of quiet rather than per sound. */ private static final long IDLE_KEEP_OPEN_NANOS = 30_000_000_000L; - private static final long FRAME_NANOS = AudioDevices.FRAME_SIZE * 1_000_000_000L / AudioDevices.SAMPLE_RATE; + private static final long FRAME_NANOS = VoiceFormat.FRAME_SIZE * 1_000_000_000L / VoiceFormat.SAMPLE_RATE; /** * How far ahead of real time the mixer renders. The line would happily take a * few hundred milliseconds at once, but anything written is fixed: a sound that @@ -61,11 +64,13 @@ public final class WavSoundPlayer implements SoundPlayer { private final List voices = new ArrayList<>(); private final Object lock = new Object(); + private final AudioIo io; private volatile String outputDevice; private volatile boolean running = true; private Thread mixer; - public WavSoundPlayer(String outputDevice) { + public WavSoundPlayer(AudioIo io, String outputDevice) { + this.io = io; this.outputDevice = outputDevice; } @@ -121,11 +126,11 @@ public final class WavSoundPlayer implements SoundPlayer { private void mixLoop() { AudioPlayback line = null; try { - float[] mix = new float[AudioDevices.FRAME_SIZE]; + float[] mix = new float[VoiceFormat.FRAME_SIZE]; long due = 0; // when the next frame is due to leave the speaker while (awaitVoices()) { if (line == null) { - line = AudioDevices.openPlayback(outputDevice, AudioDevices.MAX_CHANNELS); + line = io.openPlayback(outputDevice, VoiceFormat.MAX_CHANNELS); line.start(); } long now = System.nanoTime(); @@ -236,7 +241,7 @@ public final class WavSoundPlayer implements SoundPlayer { } mono[i] = sum / channels; } - return resample(mono, source.getSampleRate(), AudioDevices.SAMPLE_RATE); + return resample(mono, source.getSampleRate(), VoiceFormat.SAMPLE_RATE); } } } diff --git a/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/pipewire/PipeWire.java b/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/pipewire/PipeWire.java index dc25deb..3210361 100644 --- a/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/pipewire/PipeWire.java +++ b/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/pipewire/PipeWire.java @@ -1,7 +1,7 @@ package com.ts3client.audio.desktop.pipewire; -import com.ts3client.audio.desktop.AudioCapture; -import com.ts3client.audio.desktop.AudioPlayback; +import com.ts3client.audio.AudioCapture; +import com.ts3client.audio.AudioPlayback; import java.util.List; diff --git a/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/pipewire/PipeWireCapture.java b/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/pipewire/PipeWireCapture.java index 88a5eaf..f79352d 100644 --- a/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/pipewire/PipeWireCapture.java +++ b/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/pipewire/PipeWireCapture.java @@ -1,6 +1,6 @@ package com.ts3client.audio.desktop.pipewire; -import com.ts3client.audio.desktop.AudioCapture; +import com.ts3client.audio.AudioCapture; import java.lang.foreign.MemorySegment; diff --git a/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/pipewire/PipeWirePlayback.java b/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/pipewire/PipeWirePlayback.java index d370a6a..eaaee05 100644 --- a/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/pipewire/PipeWirePlayback.java +++ b/ts3-client/desktop/src/main/java/com/ts3client/audio/desktop/pipewire/PipeWirePlayback.java @@ -1,6 +1,6 @@ package com.ts3client.audio.desktop.pipewire; -import com.ts3client.audio.desktop.AudioPlayback; +import com.ts3client.audio.AudioPlayback; import java.lang.foreign.MemorySegment; diff --git a/ts3-client/swing/src/main/java/com/ts3client/ui/DevicesPanel.java b/ts3-client/swing/src/main/java/com/ts3client/ui/DevicesPanel.java index f920c79..2cb1429 100644 --- a/ts3-client/swing/src/main/java/com/ts3client/ui/DevicesPanel.java +++ b/ts3-client/swing/src/main/java/com/ts3client/ui/DevicesPanel.java @@ -1,9 +1,10 @@ package com.ts3client.ui; import com.formdev.flatlaf.util.UIScale; +import com.ts3client.audio.AudioDevice; +import com.ts3client.audio.AudioIo; import com.ts3client.audio.VoiceInput; import com.ts3client.audio.VoiceOutput; -import com.ts3client.audio.desktop.AudioDevices; import com.ts3client.config.Settings; import javax.swing.Box; @@ -24,8 +25,8 @@ import java.util.function.Consumer; */ final class DevicesPanel extends FormPanel { - private final JComboBox inputCombo; - private final JComboBox outputCombo; + private final JComboBox inputCombo; + private final JComboBox outputCombo; private final JSlider inputGain; private final JSlider outputVol; private final JCheckBox denoiseCheck; @@ -34,16 +35,16 @@ final class DevicesPanel extends FormPanel { private final JCheckBox agcCheck; private final JCheckBox mutedWarningCheck; - DevicesPanel(Settings settings, VoiceOutput livePlayback, + DevicesPanel(Settings settings, AudioIo io, VoiceOutput livePlayback, Consumer> applyLive, Runnable onInputDeviceChanged, Consumer onOutputDeviceChanged) { GridBagConstraints c = gbc(); - List ins = AudioDevices.inputDevices(); - List outs = AudioDevices.outputDevices(); + List ins = io.inputDevices(); + List outs = io.outputDevices(); - inputCombo = new JComboBox<>(ins.toArray(new AudioDevices.Device[0])); - outputCombo = new JComboBox<>(outs.toArray(new AudioDevices.Device[0])); + inputCombo = new JComboBox<>(ins.toArray(new AudioDevice[0])); + outputCombo = new JComboBox<>(outs.toArray(new AudioDevice[0])); selectOrDefault(inputCombo, settings.inputDevice); selectOrDefault(outputCombo, settings.outputDevice); @@ -158,7 +159,7 @@ final class DevicesPanel extends FormPanel { target.mutedTalkWarning = mutedWarningCheck.isSelected(); } - private static void selectOrDefault(JComboBox combo, String deviceId) { + private static void selectOrDefault(JComboBox combo, String deviceId) { if (deviceId != null && !deviceId.isEmpty()) { for (int i = 0; i < combo.getItemCount(); i++) { if (deviceId.equals(combo.getItemAt(i).id())) { @@ -170,8 +171,8 @@ final class DevicesPanel extends FormPanel { combo.setSelectedIndex(0); } - private static String comboValue(JComboBox combo) { - AudioDevices.Device d = (AudioDevices.Device) combo.getSelectedItem(); + private static String comboValue(JComboBox combo) { + AudioDevice d = (AudioDevice) combo.getSelectedItem(); return d == null ? "" : d.id(); } } diff --git a/ts3-client/swing/src/main/java/com/ts3client/ui/MainFrame.java b/ts3-client/swing/src/main/java/com/ts3client/ui/MainFrame.java index 49785fd..673ec57 100644 --- a/ts3-client/swing/src/main/java/com/ts3client/ui/MainFrame.java +++ b/ts3-client/swing/src/main/java/com/ts3client/ui/MainFrame.java @@ -834,7 +834,7 @@ public final class MainFrame extends JFrame implements ServerTabPane.Listener { private void showSettings() { // Device and voice-gating changes apply to the capturing connection and the // playback of the visible one; the rest pick the new settings up on connect. - SettingsDialog dlg = new SettingsDialog(this, settings, + SettingsDialog dlg = new SettingsDialog(this, settings, audio, micTab == null ? null : micTab.connection().getMicrophone(), selected == null ? null : selected.connection().getPlayback(), sounds, hotkeys, diff --git a/ts3-client/swing/src/main/java/com/ts3client/ui/MicrophoneTest.java b/ts3-client/swing/src/main/java/com/ts3client/ui/MicrophoneTest.java index 930e76b..7aecf10 100644 --- a/ts3-client/swing/src/main/java/com/ts3client/ui/MicrophoneTest.java +++ b/ts3-client/swing/src/main/java/com/ts3client/ui/MicrophoneTest.java @@ -1,9 +1,8 @@ package com.ts3client.ui; +import com.ts3client.audio.AudioBackend; +import com.ts3client.audio.AudioPlayback; import com.ts3client.audio.VoiceInput; -import com.ts3client.audio.desktop.AudioDevices; -import com.ts3client.audio.desktop.AudioPlayback; -import com.ts3client.audio.desktop.DesktopVoiceInput; import com.ts3client.config.Settings; import javax.swing.SwingUtilities; @@ -13,7 +12,7 @@ import java.util.function.Consumer; /** * Drives the settings dialog's microphone test from a real capture chain. * - *

The dialog runs its own {@link DesktopVoiceInput} rather than borrowing the connected + *

The dialog runs its own {@link VoiceInput} rather than borrowing the connected * one, whose listeners belong to the connection. Because it is the same class that feeds * the server, the level and the gate shown here are exactly what would be transmitted — * pre-processing, voice detection, hangover and pre-roll included. @@ -29,7 +28,8 @@ final class MicrophoneTest { private final Consumer onLevel; private final Consumer onTransmitting; - private DesktopVoiceInput mic; + private final AudioBackend audio; + private VoiceInput mic; private final ArrayBlockingQueue loopbackQueue = new ArrayBlockingQueue<>(LOOPBACK_QUEUE_FRAMES); @@ -38,7 +38,8 @@ final class MicrophoneTest { private Thread loopbackThread; private String outputDevice = ""; - MicrophoneTest(Consumer onLevel, Consumer onTransmitting) { + MicrophoneTest(AudioBackend audio, Consumer onLevel, Consumer onTransmitting) { + this.audio = audio; this.onLevel = onLevel; this.onTransmitting = onTransmitting; } @@ -55,7 +56,7 @@ final class MicrophoneTest { */ boolean start(Settings settings) { stop(); - DesktopVoiceInput input = new DesktopVoiceInput(settings); + VoiceInput input = audio.createInput(settings); input.setLevelListener(db -> SwingUtilities.invokeLater(() -> onLevel.accept(db))); input.setTalkListener(talking -> SwingUtilities.invokeLater(() -> onTransmitting.accept(talking))); input.setMonitorListener(this::enqueueForLoopback); @@ -133,7 +134,7 @@ final class MicrophoneTest { try { // The monitored frames are mono for voice; ask for a matching line so no // channel juggling is needed, and up-mix only if the device insists on stereo. - line = AudioDevices.openPlayback(outputDevice, 1); + line = audio.io().openPlayback(outputDevice, 1); line.start(); int channels = line.channels(); while (loopbackRunning) { diff --git a/ts3-client/swing/src/main/java/com/ts3client/ui/SettingsDialog.java b/ts3-client/swing/src/main/java/com/ts3client/ui/SettingsDialog.java index 5370c06..b54203f 100644 --- a/ts3-client/swing/src/main/java/com/ts3client/ui/SettingsDialog.java +++ b/ts3-client/swing/src/main/java/com/ts3client/ui/SettingsDialog.java @@ -1,6 +1,7 @@ package com.ts3client.ui; import com.formdev.flatlaf.util.UIScale; +import com.ts3client.audio.AudioBackend; import com.ts3client.audio.OpusParameters; import com.ts3client.audio.VoiceInput; import com.ts3client.audio.VoiceOutput; @@ -39,7 +40,7 @@ public final class SettingsDialog extends JDialog { private final ClientVersionPanel clientVersionPanel; private final ContactsPanel contactsPanel; - public SettingsDialog(Frame owner, Settings settings, + public SettingsDialog(Frame owner, Settings settings, AudioBackend audio, VoiceInput liveMic, VoiceOutput livePlayback, SoundNotifier sounds, HotkeyService hotkeys, Runnable onApply) { super(owner, "Options", true); @@ -53,9 +54,9 @@ public final class SettingsDialog extends JDialog { hotkeysPanel = new HotkeysPanel(hotkeys); clientVersionPanel = new ClientVersionPanel(settings); contactsPanel = new ContactsPanel(settings); - devicesPanel = new DevicesPanel(settings, livePlayback, + devicesPanel = new DevicesPanel(settings, audio.io(), livePlayback, this::applyLive, this::restartTest, this::setTestOutputDevice); - voiceActivationPanel = new VoiceActivationPanel(settings, liveMic, hotkeys, + voiceActivationPanel = new VoiceActivationPanel(settings, audio, liveMic, hotkeys, hotkeysPanel, this::audioSnapshot); JTabbedPane tabs = new JTabbedPane(); diff --git a/ts3-client/swing/src/main/java/com/ts3client/ui/VoiceActivationPanel.java b/ts3-client/swing/src/main/java/com/ts3client/ui/VoiceActivationPanel.java index f0ef0d1..29fcd3e 100644 --- a/ts3-client/swing/src/main/java/com/ts3client/ui/VoiceActivationPanel.java +++ b/ts3-client/swing/src/main/java/com/ts3client/ui/VoiceActivationPanel.java @@ -1,6 +1,7 @@ package com.ts3client.ui; import com.formdev.flatlaf.util.UIScale; +import com.ts3client.audio.AudioBackend; import com.ts3client.audio.InputLevel; import com.ts3client.audio.OpusParameters; import com.ts3client.audio.VoiceInput; @@ -42,7 +43,7 @@ final class VoiceActivationPanel extends FormPanel { private final HotkeyService hotkeys; private final HotkeysPanel hotkeysPanel; private final Supplier audioSnapshot; - private final MicrophoneTest micTest = new MicrophoneTest(this::onTestLevel, this::onTestTalking); + private final MicrophoneTest micTest; private final JRadioButton vadRadio; private final JRadioButton pttRadio; @@ -62,10 +63,11 @@ final class VoiceActivationPanel extends FormPanel { private final JCheckBox fecCheck; private final JCheckBox musicCheck; - VoiceActivationPanel(Settings settings, VoiceInput liveMic, HotkeyService hotkeys, + VoiceActivationPanel(Settings settings, AudioBackend audio, VoiceInput liveMic, HotkeyService hotkeys, HotkeysPanel hotkeysPanel, Supplier audioSnapshot) { this.settings = settings; this.liveMic = liveMic; + this.micTest = new MicrophoneTest(audio, this::onTestLevel, this::onTestTalking); this.hotkeys = hotkeys; this.hotkeysPanel = hotkeysPanel; this.audioSnapshot = audioSnapshot;