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 index 79bdfce..b57a573 100644 --- 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 @@ -5,13 +5,16 @@ 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; /** @@ -25,6 +28,14 @@ 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 @@ -34,19 +45,26 @@ public final class DesktopVoiceOutput implements VoiceOutput { * 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(200); + 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; @@ -77,12 +95,19 @@ public final class DesktopVoiceOutput implements VoiceOutput { } 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(); @@ -109,14 +134,7 @@ public final class DesktopVoiceOutput implements VoiceOutput { long now = System.nanoTime(); for (ClientStream s : streams.values()) { if (s.talking && now - s.lastPacketNanos > TALK_TIMEOUT_NANOS) { - s.worker.submit(() -> { - try { - s.line.drain(); - } catch (Exception ignored) { - } - if (s.decoder != null) s.decoder.reset(); - markTalking(s, false); - }); + s.worker.submit(() -> stopStream(s)); } } } @@ -134,7 +152,7 @@ public final class DesktopVoiceOutput implements VoiceOutput { if (d) { // Stop everyone talking immediately. for (ClientStream s : streams.values()) { - markTalking(s, false); + s.worker.submit(() -> stopStream(s)); } } } @@ -144,8 +162,13 @@ public final class DesktopVoiceOutput implements VoiceOutput { } public void setClientMuted(int clientId, boolean muted) { - if (muted) mutedClients.add(clientId); - else mutedClients.remove(clientId); + 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) { @@ -158,12 +181,12 @@ public final class DesktopVoiceOutput implements VoiceOutput { /** Entry point wired into {@code client.setVoiceHandler(...)}. */ public void handleVoice(PacketBody0Voice voice) { - route(voice.getClientId(), voice.getCodecType(), voice.getCodecData()); + 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.getCodecType(), whisper.getCodecData()); + route(whisper.getClientId(), whisper.getPacketId(), whisper.getCodecType(), whisper.getCodecData()); } /** TeamSpeak streams {@code OPUS_MUSIC} in stereo and everything else in mono. */ @@ -171,7 +194,7 @@ public final class DesktopVoiceOutput implements VoiceOutput { return codec == CodecType.OPUS_MUSIC ? 2 : 1; } - private void route(int clientId, CodecType codec, byte[] data) { + private void route(int clientId, int packetId, CodecType codec, byte[] data) { if (deafened) return; if (mutedClients.contains(clientId)) return; @@ -191,23 +214,119 @@ public final class DesktopVoiceOutput implements VoiceOutput { final ClientStream target = stream; target.lastPacketNanos = System.nanoTime(); + VoiceFrame frame = new VoiceFrame(packetId & 0xFFFF, codec, data, data == null || data.length == 0); - if (data == null || data.length == 0) { - // End of a talk burst: flush and reset the decoder, mark silent. - target.worker.submit(() -> { + 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 { - target.line.drain(); - } catch (Exception ignored) { + playNext(stream); + } finally { + stream.tickQueued.set(false); } - if (target.decoder != null) target.decoder.reset(); - markTalking(target, false); }); - return; + } 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; } - markTalking(target, true); - final int channels = channelsFor(codec); - target.worker.submit(() -> decodeAndPlay(target, channels, data)); + 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) {