Smooth incoming voice under jitter
This commit is contained in:
@@ -5,13 +5,16 @@ import com.github.manevolent.ts3j.protocol.packet.PacketBody0Voice;
|
|||||||
import com.github.manevolent.ts3j.protocol.packet.PacketBody1VoiceWhisper;
|
import com.github.manevolent.ts3j.protocol.packet.PacketBody1VoiceWhisper;
|
||||||
import com.ts3client.audio.VoiceOutput;
|
import com.ts3client.audio.VoiceOutput;
|
||||||
|
|
||||||
|
import java.util.HashMap;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
import java.util.Set;
|
import java.util.Set;
|
||||||
import java.util.concurrent.ConcurrentHashMap;
|
import java.util.concurrent.ConcurrentHashMap;
|
||||||
import java.util.concurrent.ExecutorService;
|
import java.util.concurrent.ExecutorService;
|
||||||
import java.util.concurrent.Executors;
|
import java.util.concurrent.Executors;
|
||||||
|
import java.util.concurrent.ScheduledFuture;
|
||||||
import java.util.concurrent.ScheduledExecutorService;
|
import java.util.concurrent.ScheduledExecutorService;
|
||||||
import java.util.concurrent.TimeUnit;
|
import java.util.concurrent.TimeUnit;
|
||||||
|
import java.util.concurrent.atomic.AtomicBoolean;
|
||||||
import java.util.function.BiConsumer;
|
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. */
|
/** Longest Opus frame (120 ms @ 48 kHz) a packet may decode to, per channel. */
|
||||||
private static final int MAX_FRAME = 5760;
|
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
|
* 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
|
* packet of either kind, is unambiguous: force the indicator off rather than trust
|
||||||
* the one packet that could go missing.
|
* 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. */
|
/** One speaker's decode + playback pipeline. */
|
||||||
private final class ClientStream {
|
private final class ClientStream {
|
||||||
|
final Object lock = new Object();
|
||||||
final int clientId;
|
final int clientId;
|
||||||
final AudioPlayback line;
|
final AudioPlayback line;
|
||||||
final int lineChannels;
|
final int lineChannels;
|
||||||
final ExecutorService worker;
|
final ExecutorService worker;
|
||||||
|
final Map<Integer, VoiceFrame> pending = new HashMap<>();
|
||||||
|
final AtomicBoolean tickQueued = new AtomicBoolean();
|
||||||
/** Decoder scratch and byte buffer, touched only by {@link #worker}. */
|
/** Decoder scratch and byte buffer, touched only by {@link #worker}. */
|
||||||
final float[] pcm = new float[MAX_FRAME * AudioDevices.MAX_CHANNELS];
|
final float[] pcm = new float[MAX_FRAME * AudioDevices.MAX_CHANNELS];
|
||||||
byte[] out = new byte[0];
|
byte[] out = new byte[0];
|
||||||
OpusDecoder decoder;
|
OpusDecoder decoder;
|
||||||
int decoderChannels;
|
int decoderChannels;
|
||||||
|
ScheduledFuture<?> playout;
|
||||||
|
int expectedPacketId = -1;
|
||||||
|
int activeChannels = 1;
|
||||||
|
int consecutivePlcFrames;
|
||||||
volatile boolean talking;
|
volatile boolean talking;
|
||||||
volatile long lastPacketNanos;
|
volatile long lastPacketNanos;
|
||||||
|
|
||||||
@@ -77,12 +95,19 @@ public final class DesktopVoiceOutput implements VoiceOutput {
|
|||||||
}
|
}
|
||||||
|
|
||||||
void close() {
|
void close() {
|
||||||
|
synchronized (lock) {
|
||||||
|
if (playout != null) playout.cancel(false);
|
||||||
|
pending.clear();
|
||||||
|
}
|
||||||
worker.shutdownNow();
|
worker.shutdownNow();
|
||||||
line.close();
|
line.close();
|
||||||
if (decoder != null) decoder.close();
|
if (decoder != null) decoder.close();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private record VoiceFrame(int packetId, CodecType codec, byte[] data, boolean end) {
|
||||||
|
}
|
||||||
|
|
||||||
private final Map<Integer, ClientStream> streams = new ConcurrentHashMap<>();
|
private final Map<Integer, ClientStream> streams = new ConcurrentHashMap<>();
|
||||||
private final Set<Integer> mutedClients = ConcurrentHashMap.newKeySet();
|
private final Set<Integer> mutedClients = ConcurrentHashMap.newKeySet();
|
||||||
|
|
||||||
@@ -109,14 +134,7 @@ public final class DesktopVoiceOutput implements VoiceOutput {
|
|||||||
long now = System.nanoTime();
|
long now = System.nanoTime();
|
||||||
for (ClientStream s : streams.values()) {
|
for (ClientStream s : streams.values()) {
|
||||||
if (s.talking && now - s.lastPacketNanos > TALK_TIMEOUT_NANOS) {
|
if (s.talking && now - s.lastPacketNanos > TALK_TIMEOUT_NANOS) {
|
||||||
s.worker.submit(() -> {
|
s.worker.submit(() -> stopStream(s));
|
||||||
try {
|
|
||||||
s.line.drain();
|
|
||||||
} catch (Exception ignored) {
|
|
||||||
}
|
|
||||||
if (s.decoder != null) s.decoder.reset();
|
|
||||||
markTalking(s, false);
|
|
||||||
});
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -134,7 +152,7 @@ public final class DesktopVoiceOutput implements VoiceOutput {
|
|||||||
if (d) {
|
if (d) {
|
||||||
// Stop everyone talking immediately.
|
// Stop everyone talking immediately.
|
||||||
for (ClientStream s : streams.values()) {
|
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) {
|
public void setClientMuted(int clientId, boolean muted) {
|
||||||
if (muted) mutedClients.add(clientId);
|
if (muted) {
|
||||||
else mutedClients.remove(clientId);
|
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) {
|
public boolean isClientMuted(int clientId) {
|
||||||
@@ -158,12 +181,12 @@ public final class DesktopVoiceOutput implements VoiceOutput {
|
|||||||
|
|
||||||
/** Entry point wired into {@code client.setVoiceHandler(...)}. */
|
/** Entry point wired into {@code client.setVoiceHandler(...)}. */
|
||||||
public void handleVoice(PacketBody0Voice voice) {
|
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(...)}. */
|
/** Entry point wired into {@code client.setWhisperHandler(...)}. */
|
||||||
public void handleWhisper(PacketBody1VoiceWhisper whisper) {
|
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. */
|
/** 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;
|
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 (deafened) return;
|
||||||
if (mutedClients.contains(clientId)) return;
|
if (mutedClients.contains(clientId)) return;
|
||||||
|
|
||||||
@@ -191,23 +214,119 @@ public final class DesktopVoiceOutput implements VoiceOutput {
|
|||||||
|
|
||||||
final ClientStream target = stream;
|
final ClientStream target = stream;
|
||||||
target.lastPacketNanos = System.nanoTime();
|
target.lastPacketNanos = System.nanoTime();
|
||||||
|
VoiceFrame frame = new VoiceFrame(packetId & 0xFFFF, codec, data, data == null || data.length == 0);
|
||||||
|
|
||||||
if (data == null || data.length == 0) {
|
synchronized (target.lock) {
|
||||||
// End of a talk burst: flush and reset the decoder, mark silent.
|
if (target.expectedPacketId < 0 || target.playout == null || target.playout.isDone()) {
|
||||||
target.worker.submit(() -> {
|
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 {
|
try {
|
||||||
target.line.drain();
|
playNext(stream);
|
||||||
} catch (Exception ignored) {
|
} 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);
|
decodeAndPlay(stream, channels, conceal ? null : frame.data());
|
||||||
final int channels = channelsFor(codec);
|
}
|
||||||
target.worker.submit(() -> decodeAndPlay(target, channels, 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) {
|
private void decodeAndPlay(ClientStream stream, int channels, byte[] data) {
|
||||||
|
|||||||
Reference in New Issue
Block a user