diff --git a/ts3-client/core/src/main/java/com/ts3client/audio/JitterBuffer.java b/ts3-client/core/src/main/java/com/ts3client/audio/JitterBuffer.java
new file mode 100644
index 0000000..2b24758
--- /dev/null
+++ b/ts3-client/core/src/main/java/com/ts3client/audio/JitterBuffer.java
@@ -0,0 +1,584 @@
+package com.ts3client.audio;
+
+/**
+ * Adaptive jitter buffer for incoming voice packets.
+ *
+ *
This is a port of the jitter buffer TeamSpeak itself uses: the official client links
+ * libspeex's {@code libspeex/jitter.c} verbatim (its diagnostics are still in the shipped
+ * binary) and drives it through the wrapper described in {@link #forTeamSpeakVoice}. The
+ * algorithm is what makes a voice stream survive a bad network: packets are reordered by
+ * timestamp, losses are reported so the codec can conceal them, and — most importantly —
+ * the playback delay continuously re-tunes itself to the measured jitter instead of being
+ * a fixed guess.
+ *
+ *
How the delay adapts
+ * Every packet's arrival is recorded as a timing: how early (positive) or late
+ * (negative) it was relative to the moment it was needed. Timings are kept in three
+ * rotating sub-windows, each holding only the {@value #MAX_TIMINGS} latest arrivals
+ * it saw, which makes the set a cheap running estimate of the distribution's late tail.
+ * {@link #computeOptimalDelay()} then walks the {@value #TOP_DELAY} latest timings and picks
+ * the delay minimising {@code delay + lateCost * (packets that would then be late)}. The
+ * result is applied a whole {@link #delayStep} at a time, by dropping a frame to shrink the
+ * delay or inserting a concealed one to grow it, so the correction is inaudible.
+ *
+ * Timestamp units
+ * Timestamps and spans are in caller-defined ticks; the buffer never converts them to time.
+ * All comparisons are wrap-safe, so a tick counter may overflow freely.
+ *
+ * Not thread-safe: every call must be externally serialised.
+ */
+public final class JitterBuffer {
+
+ /** Outcome of a {@link #get} call. */
+ public enum Status {
+ /** A packet was available; see {@link #payload()}. */
+ OK,
+ /** Nothing available — the caller should conceal {@link #span()} ticks of loss. */
+ MISSING,
+ /** The buffer is growing its delay; the caller should insert {@link #span()} ticks. */
+ INSERTION
+ }
+
+ /** How many packets may be held at once, awaiting playback or reordering. */
+ private static final int MAX_PACKETS = 200;
+
+ /** Timings kept per sub-window. Only the latest arrivals are retained. */
+ private static final int MAX_TIMINGS = 40;
+
+ /** Rotating sub-windows the timings are spread over. */
+ private static final int SUB_WINDOWS = 3;
+
+ /** How many of the latest timings the delay search considers. */
+ private static final int TOP_DELAY = 40;
+
+ /**
+ * The latest arrival timings seen within one sub-window, sorted ascending.
+ *
+ *
Only the {@value #MAX_TIMINGS} latest are kept: an arrival too early to displace
+ * any of them is counted but not stored, since the delay search only ever looks at the
+ * late tail.
+ */
+ private static final class TimingWindow {
+ final int[] timing = new int[MAX_TIMINGS];
+ /** Timings actually stored. */
+ int filled;
+ /** Arrivals seen, including those too early to store. */
+ int seen;
+
+ void clear() {
+ filled = 0;
+ seen = 0;
+ }
+
+ void add(int value) {
+ if (filled >= MAX_TIMINGS && value >= timing[filled - 1]) {
+ seen++;
+ return;
+ }
+
+ int pos = 0;
+ while (pos < filled && value >= timing[pos]) pos++;
+
+ if (pos < filled) {
+ int moveSize = filled - pos;
+ if (filled == MAX_TIMINGS) moveSize--;
+ System.arraycopy(timing, pos, timing, pos + 1, moveSize);
+ }
+ timing[pos] = value;
+
+ seen++;
+ if (filled < MAX_TIMINGS) filled++;
+ }
+ }
+
+ private final byte[][] data = new byte[MAX_PACKETS][];
+ private final int[] timestamp = new int[MAX_PACKETS];
+ private final int[] span = new int[MAX_PACKETS];
+ /** Value of {@link #nextStop} when the packet arrived, or 0 if that is unknown. */
+ private final int[] arrival = new int[MAX_PACKETS];
+
+ private final TimingWindow[] windows = new TimingWindow[SUB_WINDOWS];
+ /** {@link #windows} in age order, newest first; rotated in place. */
+ private final TimingWindow[] byAge = new TimingWindow[SUB_WINDOWS];
+
+ private final int delayStep;
+ private final int concealmentSize;
+
+ private int bufferMargin;
+ private int maxLateRate;
+ private int lateCost;
+ private int windowSize;
+ private int subWindowSize;
+
+ /** Playback position: the timestamp the next {@link #get} wants. */
+ private int pointerTimestamp;
+ /** Timestamp playback will have reached by the time the next packet is needed. */
+ private int nextStop;
+ /** Ticks handed out by the last {@link #get} beyond what was asked for. */
+ private int buffered;
+ /** Ticks of concealment {@link #get} has been told to insert before reading on. */
+ private int interpolationRequested;
+ /** Consecutive {@link #get} calls that found nothing. */
+ private int lostCount;
+ /** True until the first packet arrives and defines the timeline. */
+ private boolean resetState;
+ /** Self-tuned cost of a late packet, used when {@link #lateCost} is 0. */
+ private int autoTradeoff;
+
+ // Result of the most recent get(), kept here so the per-frame path allocates nothing.
+ private byte[] resultPayload;
+ private int resultTimestamp;
+ private int resultSpan;
+
+ /**
+ * @param stepSize granularity of delay changes, and the span of one concealed frame;
+ * normally the span of a single packet
+ * @param margin extra delay held on top of the measured optimum
+ * @param maxLateRate percentage of packets allowed to arrive too late, which also sets
+ * how long a window the delay is estimated over
+ * @param lateCost cost of one late packet relative to delay, or 0 to self-tune
+ */
+ public JitterBuffer(int stepSize, int margin, int maxLateRate, int lateCost) {
+ this.delayStep = stepSize;
+ this.concealmentSize = stepSize;
+ this.bufferMargin = margin;
+ this.lateCost = lateCost;
+ setMaxLateRate(maxLateRate);
+
+ for (int i = 0; i < SUB_WINDOWS; i++) {
+ windows[i] = new TimingWindow();
+ byAge[i] = windows[i];
+ }
+ reset();
+ }
+
+ /**
+ * A buffer configured the way the official TeamSpeak client configures its own.
+ *
+ *
Recovered from {@code ts3client_linux_amd64} (3.6.x): it calls
+ * {@code jitter_buffer_init(60)} and then sets margin 60 and max late rate 1, pulling
+ * one 60-tick frame per {@code jitter_buffer_get}. One tick is therefore a sixtieth of
+ * a 20 ms voice frame, so the delay moves in whole-frame steps, one frame of slack
+ * is held on top of the measured optimum, and the target is for no more than 1% of
+ * packets to arrive late — a deliberately conservative setting that estimates over a
+ * ~1.3 s window and is the reason TeamSpeak rides out network trouble smoothly.
+ *
+ * @param ticksPerFrame ticks spanned by one voice frame; see {@link #TEAMSPEAK_TICKS_PER_FRAME}
+ */
+ public static JitterBuffer forTeamSpeakVoice(int ticksPerFrame) {
+ return new JitterBuffer(ticksPerFrame, ticksPerFrame, 1, 0);
+ }
+
+ /** Ticks per voice frame in TeamSpeak's own timebase, i.e. a 20 ms frame in 1/3 ms units. */
+ public static final int TEAMSPEAK_TICKS_PER_FRAME = 60;
+
+ /**
+ * Sets the share of packets allowed to arrive too late, as a percentage.
+ *
+ *
This doubles as the estimator's time constant: a stricter rate needs a
+ * proportionally longer window to observe that many late packets.
+ */
+ public void setMaxLateRate(int rate) {
+ this.maxLateRate = Math.max(1, rate);
+ this.windowSize = 100 * TOP_DELAY / this.maxLateRate;
+ this.subWindowSize = this.windowSize / SUB_WINDOWS;
+ }
+
+ /** Sets the extra delay held on top of the measured optimum. */
+ public void setMargin(int margin) {
+ this.bufferMargin = margin;
+ }
+
+ public int margin() {
+ return bufferMargin;
+ }
+
+ /** Ticks currently held ahead of the playback position, i.e. the live buffering depth. */
+ public int bufferedSpan() {
+ int total = 0;
+ for (int i = 0; i < MAX_PACKETS; i++) {
+ if (data[i] != null && ge(timestamp[i], pointerTimestamp)) total += span[i];
+ }
+ return total;
+ }
+
+ /** True when no packet is held and playback has nothing left to drain. */
+ public boolean isEmpty() {
+ for (int i = 0; i < MAX_PACKETS; i++) {
+ if (data[i] != null) return false;
+ }
+ return true;
+ }
+
+ /**
+ * Drops every held packet and forgets both the timeline and the network estimate.
+ *
+ *
Reserved for the cases where a speaker's position genuinely stopped meaning
+ * anything — a new speaker, or a silence long enough that the sequence may have been
+ * recycled. Notably not something to do between talk bursts: the timeline stays
+ * valid across a pause because the sender stops numbering while the receiver stops
+ * playing, and the estimate needs far longer than one burst to be worth anything.
+ */
+ public void reset() {
+ for (int i = 0; i < MAX_PACKETS; i++) data[i] = null;
+ pointerTimestamp = 0;
+ nextStop = 0;
+ resetState = true;
+ lostCount = 0;
+ buffered = 0;
+ interpolationRequested = 0;
+ autoTradeoff = 32000;
+ for (TimingWindow w : windows) w.clear();
+ System.arraycopy(windows, 0, byAge, 0, SUB_WINDOWS);
+ }
+
+ /**
+ * Offers an arrived packet to the buffer.
+ *
+ * @param payload the packet, retained by reference until played or dropped
+ * @param timestamp the packet's position on the sender's timeline, in ticks
+ * @param span how many ticks the packet covers
+ */
+ public void put(byte[] payload, int timestamp, int span) {
+ // Drop anything playback has already moved past.
+ if (!resetState) {
+ for (int i = 0; i < MAX_PACKETS; i++) {
+ if (data[i] != null && le(this.timestamp[i] + this.span[i], pointerTimestamp)) {
+ data[i] = null;
+ }
+ }
+ }
+
+ // A packet that missed its slot still tells us how far behind the clock we are.
+ boolean late = !resetState && lt(timestamp, nextStop);
+ if (late) {
+ updateTimings(timestamp - nextStop - bufferMargin);
+ }
+
+ // The consumer has stopped fetching; don't sit on packets it will never take.
+ if (lostCount > 20) reset();
+
+ // Ignore packets so late that even concealment has moved past them.
+ if (!resetState && !ge(timestamp + span + delayStep, pointerTimestamp)) return;
+
+ int slot = -1;
+ for (int i = 0; i < MAX_PACKETS; i++) {
+ if (data[i] == null) {
+ slot = i;
+ break;
+ }
+ }
+ if (slot < 0) {
+ // Full: make room by discarding the oldest packet.
+ slot = 0;
+ int earliest = this.timestamp[0];
+ for (int i = 1; i < MAX_PACKETS; i++) {
+ if (lt(this.timestamp[i], earliest)) {
+ earliest = this.timestamp[i];
+ slot = i;
+ }
+ }
+ }
+
+ data[slot] = payload;
+ this.timestamp[slot] = timestamp;
+ this.span[slot] = span;
+ this.arrival[slot] = (resetState || late) ? 0 : nextStop;
+ }
+
+ /**
+ * Takes the next {@code desiredSpan} ticks of audio.
+ *
+ *
On {@link Status#OK} the packet is in {@link #payload()}. On {@link Status#MISSING}
+ * or {@link Status#INSERTION} the caller must produce {@link #span()} ticks of
+ * concealment instead — for Opus that means decoding a null packet.
+ */
+ public Status get(int desiredSpan) {
+ resultPayload = null;
+
+ // The first packet to arrive defines where the timeline starts.
+ if (resetState) {
+ int oldest = 0;
+ boolean found = false;
+ for (int i = 0; i < MAX_PACKETS; i++) {
+ if (data[i] != null && (!found || lt(timestamp[i], oldest))) {
+ oldest = timestamp[i];
+ found = true;
+ }
+ }
+ if (!found) {
+ resultTimestamp = 0;
+ resultSpan = interpolationRequested;
+ return Status.MISSING;
+ }
+ resetState = false;
+ pointerTimestamp = oldest;
+ nextStop = oldest;
+ }
+
+ // A delay increase decided during the last tick(): pad before reading on.
+ if (interpolationRequested != 0) {
+ resultTimestamp = pointerTimestamp;
+ resultSpan = interpolationRequested;
+ pointerTimestamp += interpolationRequested;
+ interpolationRequested = 0;
+ buffered = resultSpan - desiredSpan;
+ return Status.INSERTION;
+ }
+
+ int found = findPacket(desiredSpan);
+ if (found >= 0) {
+ lostCount = 0;
+
+ // An on-time arrival: record how much slack it had.
+ if (arrival[found] != 0) {
+ updateTimings(timestamp[found] - arrival[found] - bufferMargin);
+ }
+
+ resultPayload = data[found];
+ resultTimestamp = timestamp[found];
+ resultSpan = span[found];
+ data[found] = null;
+
+ pointerTimestamp = resultTimestamp + resultSpan;
+ buffered = resultSpan - desiredSpan;
+ return Status.OK;
+ }
+
+ lostCount++;
+
+ int optimal = computeOptimalDelay();
+ if (optimal < 0) {
+ // The measured jitter wants a deeper buffer: insert without advancing playback.
+ shiftTimings(-optimal);
+ resultTimestamp = pointerTimestamp;
+ resultSpan = -optimal;
+ buffered = resultSpan - desiredSpan;
+ return Status.INSERTION;
+ }
+
+ // Ordinary loss: conceal a whole number of frames and move on.
+ int concealed = roundDown(desiredSpan, concealmentSize);
+ resultTimestamp = pointerTimestamp;
+ resultSpan = concealed;
+ pointerTimestamp += concealed;
+ buffered = concealed - desiredSpan;
+ return Status.MISSING;
+ }
+
+ /**
+ * The held packet sitting at the playback position, without consuming it.
+ *
+ *
After a {@link Status#MISSING} the position has already stepped past the gap, so
+ * this is the packet immediately following the lost one — the one carrying the FEC copy
+ * that can rebuild it. {@code null} when that packet is missing too.
+ */
+ public byte[] nextPayload() {
+ for (int i = 0; i < MAX_PACKETS; i++) {
+ if (data[i] != null && timestamp[i] == pointerTimestamp) return data[i];
+ }
+ return null;
+ }
+
+ /** The timestamp the next {@link #get} will play from. */
+ public int position() {
+ return pointerTimestamp;
+ }
+
+ /** Whole frames currently held, given a frame spans {@code frameSpan} ticks. */
+ public int heldFrames(int frameSpan) {
+ int held = 0;
+ for (int i = 0; i < MAX_PACKETS; i++) {
+ if (data[i] != null) held += Math.max(1, span[i] / frameSpan);
+ }
+ return held;
+ }
+
+ /** The packet returned by the last {@link #get}, or {@code null} if it was not {@link Status#OK}. */
+ public byte[] payload() {
+ return resultPayload;
+ }
+
+ /** Ticks covered by the last {@link #get}, whether a real packet or concealment. */
+ public int span() {
+ return resultSpan;
+ }
+
+ /** Timestamp of the last {@link #get}. */
+ public int timestamp() {
+ return resultTimestamp;
+ }
+
+ /**
+ * Advances the playback clock by one frame, after the audio from {@link #get} has been
+ * handed to the output. This is also where the delay is re-tuned.
+ */
+ public void tick() {
+ updateDelay();
+ nextStop = pointerTimestamp - Math.max(0, buffered);
+ buffered = 0;
+ }
+
+ /**
+ * Advances the playback clock by one frame that the caller played from audio it still had
+ * from an earlier {@link #get}, instead of calling {@link #get} and {@link #tick}.
+ *
+ * @param remaining ticks of that audio still left after this frame
+ */
+ public void remainingSpan(int remaining) {
+ updateDelay();
+ nextStop = pointerTimestamp - remaining;
+ }
+
+ /** Picks the held packet that best covers the chunk starting at the playback position. */
+ private int findPacket(int desiredSpan) {
+ // Exact start, covering the whole chunk.
+ for (int i = 0; i < MAX_PACKETS; i++) {
+ if (data[i] != null && timestamp[i] == pointerTimestamp
+ && ge(timestamp[i] + span[i], pointerTimestamp + desiredSpan)) {
+ return i;
+ }
+ }
+ // An earlier packet still covering the whole chunk.
+ for (int i = 0; i < MAX_PACKETS; i++) {
+ if (data[i] != null && le(timestamp[i], pointerTimestamp)
+ && ge(timestamp[i] + span[i], pointerTimestamp + desiredSpan)) {
+ return i;
+ }
+ }
+ // An earlier packet covering part of it.
+ for (int i = 0; i < MAX_PACKETS; i++) {
+ if (data[i] != null && le(timestamp[i], pointerTimestamp)
+ && gt(timestamp[i] + span[i], pointerTimestamp)) {
+ return i;
+ }
+ }
+ // Anything starting inside the chunk: prefer the earliest, then the longest.
+ int best = -1;
+ for (int i = 0; i < MAX_PACKETS; i++) {
+ if (data[i] == null) continue;
+ if (!(lt(timestamp[i], pointerTimestamp + desiredSpan) && ge(timestamp[i], pointerTimestamp))) continue;
+ if (best < 0 || lt(timestamp[i], timestamp[best])
+ || (timestamp[i] == timestamp[best] && gt(span[i], span[best]))) {
+ best = i;
+ }
+ }
+ return best;
+ }
+
+ /**
+ * Applies the delay the timings call for, by scheduling an insertion to grow the buffer
+ * or skipping ahead to shrink it. The timings are shifted by the same amount so the
+ * next estimate is made relative to the new delay.
+ */
+ private void updateDelay() {
+ int optimal = computeOptimalDelay();
+ if (optimal == 0) return;
+
+ shiftTimings(-optimal);
+ pointerTimestamp += optimal;
+ if (optimal < 0) interpolationRequested = -optimal;
+ }
+
+ /**
+ * Finds the delay minimising {@code delay + lateCost * latePackets} over the observed
+ * timings, quantised to whole {@link #delayStep}s.
+ *
+ * @return ticks to shift playback by: negative to buffer more, positive to catch up
+ */
+ private int computeOptimalDelay() {
+ int totalSeen = 0;
+ for (TimingWindow w : windows) totalSeen += w.seen;
+ if (totalSeen == 0) return 0;
+
+ float lateFactor = lateCost != 0
+ ? lateCost * 100.0f / totalSeen
+ : (float) autoTradeoff * windowSize / totalSeen;
+
+ int[] pos = new int[SUB_WINDOWS];
+ int optimal = 0;
+ int bestCost = Integer.MAX_VALUE;
+ int late = 0;
+ boolean penaltyTaken = false;
+ int earliestSeen = 0;
+ int latestSeen = 0;
+
+ // Walk the timings from latest to earliest, costing each as a candidate delay.
+ for (int i = 0; i < TOP_DELAY; i++) {
+ int next = -1;
+ int timing = Integer.MAX_VALUE;
+ for (int j = 0; j < SUB_WINDOWS; j++) {
+ if (pos[j] < windows[j].filled && windows[j].timing[pos[j]] < timing) {
+ next = j;
+ timing = windows[j].timing[pos[j]];
+ }
+ }
+ if (next < 0) break;
+
+ if (i == 0) latestSeen = timing;
+ earliestSeen = timing;
+ timing = roundDown(timing, delayStep);
+ pos[next]++;
+
+ int cost = (int) (-timing + lateFactor * late);
+ if (cost < bestCost) {
+ bestCost = cost;
+ optimal = timing;
+ }
+
+ late++;
+ // Hysteresis: charge extra the first time we would start making packets late.
+ if (timing >= 0 && !penaltyTaken) {
+ penaltyTaken = true;
+ late += 4;
+ }
+ }
+
+ autoTradeoff = 1 + (earliestSeen - latestSeen) / TOP_DELAY;
+
+ // Too little evidence to justify shrinking the buffer.
+ if (totalSeen < TOP_DELAY && optimal > 0) return 0;
+ return optimal;
+ }
+
+ /** Records one arrival's slack, rotating the sub-windows when the newest fills up. */
+ private void updateTimings(int timing) {
+ timing = Math.max(-32767, Math.min(32767, timing));
+
+ if (byAge[0].seen >= subWindowSize) {
+ TimingWindow oldest = byAge[SUB_WINDOWS - 1];
+ System.arraycopy(byAge, 0, byAge, 1, SUB_WINDOWS - 1);
+ byAge[0] = oldest;
+ oldest.clear();
+ }
+ byAge[0].add(timing);
+ }
+
+ /** Re-bases every recorded timing after the delay moved. */
+ private void shiftTimings(int amount) {
+ for (TimingWindow w : windows) {
+ for (int i = 0; i < w.filled; i++) w.timing[i] += amount;
+ }
+ }
+
+ private static int roundDown(int value, int step) {
+ return value < 0 ? (value - step + 1) / step * step : value / step * step;
+ }
+
+ // Wrap-safe comparisons, so tick counters may overflow.
+ private static boolean lt(int a, int b) {
+ return a - b < 0;
+ }
+
+ private static boolean le(int a, int b) {
+ return a - b <= 0;
+ }
+
+ private static boolean gt(int a, int b) {
+ return a - b > 0;
+ }
+
+ private static boolean ge(int a, int b) {
+ return a - b >= 0;
+ }
+}
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
index a5abf94..9b272b8 100644
--- a/ts3-client/core/src/main/java/com/ts3client/audio/VoiceStream.java
+++ b/ts3-client/core/src/main/java/com/ts3client/audio/VoiceStream.java
@@ -4,13 +4,21 @@ 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.
+ * One speaker's incoming voice: packets go into an adaptive {@link JitterBuffer} as they
+ * arrive, and every {@link #pull} plays out the next 20 ms, concealing what went missing.
+ *
+ *
The buffer is driven the way the official client drives it. A speaker has one timeline,
+ * {@code unwrap16(packetId) * 60} ticks, which is never re-anchored between talk bursts: the
+ * sender stops numbering when it stops talking and playback stops at the same point, so the
+ * next burst lands where the last one left off and the delay the buffer has learned carries
+ * over. Only a long silence drops the timeline.
+ *
+ *
An empty packet is the talker letting go of the key. It never enters the buffer; its
+ * position is remembered, the burst ends once playback reaches it, and its slot is played as
+ * silence rather than concealed, since nothing was lost there.
*
*
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
@@ -23,36 +31,51 @@ final class VoiceStream {
/** 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) {
- }
+ private static final int FRAME = VoiceFormat.FRAME_SIZE;
+ private static final int FRAME_TICKS = JitterBuffer.TEAMSPEAK_TICKS_PER_FRAME;
+ /** Pulls in a row with nothing held before a burst whose end marker got lost is given up. */
+ private static final int EMPTY_PULLS_TO_END = 10;
+ /** How far back a talker's stop is remembered, as in the official client. */
+ private static final int STOP_MARK_WINDOW_TICKS = 3000;
+ /** Silence after which the timeline and the delay estimate are dropped, as in the official client. */
+ private static final long IDLE_RESET_NANOS = 5_000_000_000L;
+ private static final int EPOCH_UNSET = Integer.MIN_VALUE;
+
+ /** What a pull plays, decided under the lock and produced outside it. */
+ private enum Source { PACKET, CONCEAL, SILENCE, SURPLUS }
final int clientId;
private final OpusCodec opus;
private final BiConsumer talkListener;
private final Object lock = new Object();
- private final Map pending = new HashMap<>();
+ // Guarded by lock.
+ private final JitterBuffer jitter = JitterBuffer.forTeamSpeakVoice(FRAME_TICKS);
+ private int epoch = EPOCH_UNSET;
+ private int lastSeq;
+ private int[] stopMarks = new int[8];
+ private int stopMarkCount;
+ /** Channels of the latest packet: TeamSpeak streams {@code OPUS_MUSIC} in stereo. */
+ private int channels = 1;
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;
+ private int emptyPulls;
+ /** Frames of decoded audio left over from the last packet, still to be played. */
+ private int surplusFrames;
+ /** Frames of concealment still owed to a delay increase. */
+ private int concealOwed;
+ /** Set when the timeline was dropped; the decoder follows on the next pull. */
+ private boolean resetDecoder;
- /** Guards the decoder: a stop or close from another thread may reach it mid-pull. */
+ /** Guards the decoder: a 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];
+ /** Audio a packet decoded to beyond its first frame; only the pulling thread touches it. */
+ private final float[] surplus = new float[MAX_FRAME * VoiceFormat.MAX_CHANNELS];
+ private int surplusOffset;
+ private int surplusChannels;
private OpusDecoder decoder;
private int decoderChannels;
- /** Length of the last packet, which a lost one most likely had too. */
- private int lastFrameSize = VoiceFormat.FRAME_SIZE;
private boolean closed;
private volatile boolean talking;
@@ -64,36 +87,48 @@ final class VoiceStream {
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.
+ * Places an incoming packet on the speaker's timeline; 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);
+ long now = System.nanoTime();
+ boolean longIdle = lastPacketNanos != 0 && now - lastPacketNanos >= IDLE_RESET_NANOS;
+ lastPacketNanos = now;
+ boolean end = data == null || data.length == 0;
boolean started = false;
synchronized (lock) {
+ if (longIdle && !playing) {
+ jitter.reset();
+ epoch = EPOCH_UNSET;
+ stopMarkCount = 0;
+ resetDecoder = true;
+ }
+ boolean fresh = epoch == EPOCH_UNSET;
+ int ticks = unwrap(packetId & 0xFFFF) * FRAME_TICKS;
+ if (end) {
+ markStop(ticks);
+ return false;
+ }
+ // Playback already went past it; the buffer would only drop it.
+ if (!fresh && ticks - jitter.position() < 0) return false;
+
+ channels = channelsFor(codec);
+ jitter.put(data, ticks, FRAME_TICKS);
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
+ emptyPulls = 0;
}
- pending.putIfAbsent(frame.packetId(), frame);
- trimPending();
}
- if (!frame.end()) markTalking(true);
+ markTalking(true);
return started;
}
@@ -101,78 +136,139 @@ final class VoiceStream {
* 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
+ * @param out room for {@link VoiceFormat#FRAME_SIZE} samples per channel
+ * @return the samples written per channel, at most one frame; 0 if 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;
+ Source source;
+ byte[] packet = null;
+ int surplusTake = 0;
+ int packetChannels;
+ boolean reset;
synchronized (lock) {
if (!playing) return IDLE;
- if (warmup > 0) {
- warmup--;
- return 0;
- }
+ reset = resetDecoder;
+ resetDecoder = false;
+ packetChannels = channels;
- frame = pending.remove(expectedPacketId);
- boolean conceal = frame == null;
- if (conceal) {
- consecutivePlcFrames++;
+ if (surplusFrames > 0 || concealOwed > 0) {
+ // Still playing out what the last get() handed over.
+ if (surplusFrames > 0) {
+ source = Source.SURPLUS;
+ surplusTake = Math.min(FRAME, surplusFrames);
+ surplusFrames -= surplusTake;
+ } else {
+ source = Source.CONCEAL;
+ concealOwed--;
+ }
+ jitter.remainingSpan(remainingTicks());
} 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;
+ boolean empty = jitter.isEmpty();
+ emptyPulls = empty ? emptyPulls + 1 : 0;
+ if (empty && (stoppedAt(jitter.position()) || emptyPulls >= EMPTY_PULLS_TO_END)) {
+ stopLocked();
+ source = null;
+ } else {
+ JitterBuffer.Status status = jitter.get(FRAME_TICKS);
+ switch (status) {
+ case OK -> {
+ source = Source.PACKET;
+ packet = jitter.payload();
+ }
+ case MISSING -> source = isStopMark(jitter.timestamp()) ? Source.SILENCE : Source.CONCEAL;
+ default -> {
+ // The buffer grows its delay by inserting this many frames.
+ source = Source.CONCEAL;
+ concealOwed = Math.max(1, jitter.span() / FRAME_TICKS) - 1;
+ }
+ }
+ jitter.tick();
+ }
}
}
- if (channels == 0) {
- finish();
+ if (source == null) {
+ markTalking(false);
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 {
- OpusDecoder d = decoderFor(channels);
- if (data == null) {
- frames = d.conceal(pcm, lastFrameSize);
- } else {
- frames = d.decode(data, pcm);
- lastFrameSize = frames;
- }
+ if (reset && decoder != null) decoder.reset();
+ frames = produce(source, packet, surplusTake, packetChannels, out, outChannels, gain);
} catch (RuntimeException e) {
return 0;
}
}
- // Match the decoded stream to the output: duplicate mono across stereo, fold a
- // stereo (music) stream down onto mono.
+ return frames;
+ }
+
+ private int produce(Source source, byte[] packet, int surplusTake, int packetChannels, float[] out,
+ int outChannels, double gain) {
+ switch (source) {
+ case SURPLUS -> {
+ emit(surplus, surplusOffset, surplusTake, surplusChannels, out, outChannels, gain);
+ surplusOffset += surplusTake;
+ return surplusTake;
+ }
+ case SILENCE -> {
+ java.util.Arrays.fill(out, 0, FRAME * outChannels, 0f);
+ return FRAME;
+ }
+ case CONCEAL -> {
+ int frames = decoderFor(decoderChannels == 0 ? packetChannels : decoderChannels).conceal(pcm, FRAME);
+ emit(pcm, 0, frames, decoderChannels, out, outChannels, gain);
+ return frames;
+ }
+ default -> {
+ int frames = decoderFor(packetChannels).decode(packet, pcm);
+ int played = Math.min(frames, FRAME);
+ emit(pcm, 0, played, packetChannels, out, outChannels, gain);
+ if (frames > played) keepSurplus(played, frames - played, packetChannels);
+ return played;
+ }
+ }
+ }
+
+ /** Holds back what a packet decoded to beyond one frame, for the following pulls. */
+ private void keepSurplus(int from, int frames, int channels) {
+ System.arraycopy(pcm, from * channels, surplus, 0, frames * channels);
+ surplusOffset = 0;
+ surplusChannels = channels;
+ synchronized (lock) {
+ if (playing) surplusFrames = frames;
+ }
+ }
+
+ /** Ticks of audio already handed out but not yet played; guarded by the lock. */
+ private int remainingTicks() {
+ return (surplusFrames * FRAME_TICKS + FRAME - 1) / FRAME + concealOwed * FRAME_TICKS;
+ }
+
+ /**
+ * Matches decoded audio to the output: duplicates mono across stereo and folds a stereo
+ * (music) stream down onto mono.
+ */
+ private static void emit(float[] src, int fromFrame, int frames, int channels, float[] out, int outChannels,
+ double gain) {
+ int base = fromFrame * channels;
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];
+ v = src[base + i * channels + c];
} else if (channels == 1) {
- v = pcm[i];
+ v = src[base + i];
} else {
double sum = 0;
- for (int s = 0; s < channels; s++) sum += pcm[i * channels + s];
+ for (int s = 0; s < channels; s++) sum += src[base + i * channels + s];
v = sum / channels;
}
out[k] = (float) (v * gain);
}
}
- return frames;
}
/**
@@ -188,24 +284,66 @@ final class VoiceStream {
return decoder;
}
+ /**
+ * Widens the 16-bit packet id into a continuous sequence. A packet that arrives after the
+ * counter wrapped, but was sent before it, counts back into the previous epoch.
+ */
+ private int unwrap(int seq) {
+ if (epoch == EPOCH_UNSET) {
+ epoch = 0;
+ lastSeq = seq;
+ } else {
+ int delta = (short) (seq - lastSeq);
+ if (delta > 0) {
+ if (seq < lastSeq) epoch++;
+ lastSeq = seq;
+ } else if (seq > lastSeq) {
+ epoch--;
+ }
+ }
+ return (epoch << 16) | seq;
+ }
+
+ private void markStop(int ticks) {
+ int position = jitter.position();
+ int kept = 0;
+ for (int i = 0; i < stopMarkCount; i++) {
+ if (position - stopMarks[i] < STOP_MARK_WINDOW_TICKS) stopMarks[kept++] = stopMarks[i];
+ }
+ stopMarkCount = kept;
+ if (stopMarkCount == stopMarks.length) stopMarks = java.util.Arrays.copyOf(stopMarks, stopMarkCount * 2);
+ stopMarks[stopMarkCount++] = ticks;
+ }
+
+ private boolean isStopMark(int ticks) {
+ for (int i = 0; i < stopMarkCount; i++) {
+ if (stopMarks[i] == ticks) return true;
+ }
+ return false;
+ }
+
+ /** Whether the talker let go at or shortly before {@code position}, so the burst is spoken out. */
+ private boolean stoppedAt(int position) {
+ for (int i = 0; i < stopMarkCount; i++) {
+ int behind = position - stopMarks[i];
+ if (behind >= 0 && behind < STOP_MARK_WINDOW_TICKS) return true;
+ }
+ return false;
+ }
+
/** Ends the current talk burst, if any; the next packet starts a new one. */
void stop() {
synchronized (lock) {
stopLocked();
}
- finish();
+ markTalking(false);
}
private void stopLocked() {
playing = false;
- pending.clear();
- }
-
- private void finish() {
- synchronized (decoderLock) {
- if (decoder != null) decoder.reset();
- }
- markTalking(false);
+ surplusFrames = 0;
+ concealOwed = 0;
+ emptyPulls = 0;
}
boolean isPlaying() {
@@ -238,25 +376,4 @@ final class VoiceStream {
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/test/java/com/ts3client/audio/FakeOpus.java b/ts3-client/core/src/test/java/com/ts3client/audio/FakeOpus.java
index c91baf2..ab0139a 100644
--- a/ts3-client/core/src/test/java/com/ts3client/audio/FakeOpus.java
+++ b/ts3-client/core/src/test/java/com/ts3client/audio/FakeOpus.java
@@ -9,7 +9,7 @@ import java.util.Arrays;
import java.util.List;
/**
- * Stands in for libopus: a "packet" is one byte, decoded to a 20 ms frame holding that byte
+ * Stands in for libopus: a "packet" of n bytes decodes to n 20 ms frames holding the first byte
* divided by 100 in every sample. Concealment yields -1, so it is easy to spot.
*/
final class FakeOpus implements OpusCodec {
@@ -33,8 +33,9 @@ final class FakeOpus implements OpusCodec {
return new OpusDecoder() {
@Override
public int decode(byte[] packet, float[] out) {
- Arrays.fill(out, 0, VoiceFormat.FRAME_SIZE * channels, packet[0] / 100f);
- return VoiceFormat.FRAME_SIZE;
+ int frames = VoiceFormat.FRAME_SIZE * packet.length;
+ Arrays.fill(out, 0, frames * channels, packet[0] / 100f);
+ return frames;
}
@Override
diff --git a/ts3-client/core/src/test/java/com/ts3client/audio/JitterBufferTest.java b/ts3-client/core/src/test/java/com/ts3client/audio/JitterBufferTest.java
new file mode 100644
index 0000000..8d77387
--- /dev/null
+++ b/ts3-client/core/src/test/java/com/ts3client/audio/JitterBufferTest.java
@@ -0,0 +1,353 @@
+package com.ts3client.audio;
+
+import org.junit.jupiter.api.Test;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Random;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+class JitterBufferTest {
+
+ private static final int SPAN = JitterBuffer.TEAMSPEAK_TICKS_PER_FRAME;
+
+ private static JitterBuffer buffer() {
+ return JitterBuffer.forTeamSpeakVoice(SPAN);
+ }
+
+ /** Frame {@code n} of a stream, tagged so it can be identified on the way out. */
+ private static byte[] frame(int n) {
+ return new byte[]{(byte) n};
+ }
+
+ private static void put(JitterBuffer b, int n) {
+ b.put(frame(n), n * SPAN, SPAN);
+ }
+
+ /** Pulls one frame, returning its tag, or -1 for concealment. */
+ private static int pull(JitterBuffer b) {
+ JitterBuffer.Status status = b.get(SPAN);
+ int tag = status == JitterBuffer.Status.OK ? b.payload()[0] : -1;
+ b.tick();
+ return tag;
+ }
+
+ @Test
+ void playsAnUndisturbedStreamInOrder() {
+ JitterBuffer b = buffer();
+ List played = new ArrayList<>();
+ for (int i = 0; i < 20; i++) {
+ put(b, i);
+ played.add(pull(b));
+ }
+ // The buffer holds some delay back, so the tail is still queued; nothing is reordered.
+ List real = played.stream().filter(n -> n >= 0).toList();
+ assertEquals(real.stream().sorted().toList(), real);
+ assertTrue(real.size() >= 15, "most frames should have played: " + played);
+ }
+
+ @Test
+ void reordersPacketsThatArriveOutOfOrder() {
+ JitterBuffer b = buffer();
+ // Prime a delay so a swapped pair still lands in time.
+ for (int i = 0; i < 12; i++) {
+ put(b, i);
+ pull(b);
+ }
+
+ put(b, 13);
+ put(b, 12);
+ put(b, 14);
+
+ List played = new ArrayList<>();
+ for (int i = 0; i < 6; i++) played.add(pull(b));
+ List real = played.stream().filter(n -> n >= 0).toList();
+ assertEquals(real.stream().sorted().toList(), real, "played out of order: " + played);
+ assertTrue(real.contains(12) && real.contains(13) && real.contains(14), "lost a frame: " + played);
+ }
+
+ @Test
+ void reportsMissingFramesInsteadOfSkippingAhead() {
+ JitterBuffer b = buffer();
+ for (int i = 0; i < 30; i++) {
+ if (i != 20) put(b, i);
+ }
+
+ int concealed = 0;
+ int played = 0;
+ for (int i = 0; i < 30; i++) {
+ if (pull(b) < 0) concealed++;
+ else played++;
+ }
+ assertTrue(concealed >= 1, "the gap should have been concealed");
+ assertEquals(29, played, "every delivered frame should still be played");
+ }
+
+ @Test
+ void everyGetYieldsExactlyOneFrameOfAudio() {
+ JitterBuffer b = buffer();
+ Random random = new Random(7);
+ long ticks = 0;
+ for (int i = 0; i < 200; i++) {
+ if (random.nextInt(10) != 0) put(b, i);
+ b.get(SPAN);
+ assertEquals(0, b.span() % SPAN, "a partial frame cannot be rendered");
+ assertTrue(b.span() > 0);
+ ticks += b.span();
+ b.tick();
+ }
+ // Concealment for a growing delay may run long, but audio never stops or doubles up.
+ assertTrue(ticks >= 200L * SPAN, "output underran: " + ticks);
+ }
+
+ @Test
+ void growsTheDelayWhenArrivalsAreJittery() {
+ JitterBuffer b = buffer();
+ Random random = new Random(11);
+
+ // Deliver in bursts: frames pile up, then nothing arrives for a while.
+ int next = 0;
+ int concealedEarly = 0;
+ int concealedLate = 0;
+ for (int i = 0; i < 400; i++) {
+ if (i % 5 == 0) {
+ int burst = 3 + random.nextInt(4);
+ for (int j = 0; j < burst && next < 400; j++) put(b, next++);
+ }
+ boolean concealed = pull(b) < 0;
+ if (i < 100) concealedEarly += concealed ? 1 : 0;
+ else if (i >= 300) concealedLate += concealed ? 1 : 0;
+ }
+
+ assertTrue(concealedLate < concealedEarly,
+ "the buffer should settle: " + concealedEarly + " -> " + concealedLate);
+ }
+
+ @Test
+ void staysStableAcrossTimestampWraparound() {
+ JitterBuffer b = buffer();
+ // Start just below the point where the tick counter overflows.
+ int base = Integer.MAX_VALUE - 5 * SPAN;
+ for (int i = 0; i < 40; i++) {
+ b.put(frame(i), base + i * SPAN, SPAN);
+ }
+ int played = 0;
+ for (int i = 0; i < 40; i++) {
+ if (b.get(SPAN) == JitterBuffer.Status.OK) {
+ assertNotNull(b.payload());
+ played++;
+ }
+ b.tick();
+ }
+ assertTrue(played >= 35, "wraparound lost frames: " + played);
+ }
+
+ @Test
+ void resetForgetsTheTimelineSoTheNextBurstResyncs() {
+ JitterBuffer b = buffer();
+ for (int i = 0; i < 10; i++) {
+ put(b, i);
+ pull(b);
+ }
+ b.reset();
+ assertTrue(b.isEmpty());
+
+ // A new burst numbered far away must still play, not be treated as hopelessly late.
+ b.put(frame(1), 500 * SPAN, SPAN);
+ assertEquals(JitterBuffer.Status.OK, b.get(SPAN));
+ assertEquals(1, b.payload()[0]);
+ }
+
+ @Test
+ void dropsPacketsThatArriveHopelesslyLate() {
+ JitterBuffer b = buffer();
+ for (int i = 0; i < 40; i++) {
+ put(b, i);
+ pull(b);
+ }
+ assertFalse(b.isEmpty() && b.bufferedSpan() != 0);
+
+ // A frame from far in the past has nowhere to go and must not stall the stream.
+ b.put(frame(99), 0, SPAN);
+ for (int i = 0; i < 5; i++) {
+ assertTrue(pull(b) != 99, "a hopelessly late frame was played");
+ }
+ }
+
+ @Test
+ void teamSpeakConfigurationMatchesTheShippedClient() {
+ JitterBuffer b = JitterBuffer.forTeamSpeakVoice(JitterBuffer.TEAMSPEAK_TICKS_PER_FRAME);
+ assertEquals(60, JitterBuffer.TEAMSPEAK_TICKS_PER_FRAME);
+ assertEquals(60, b.margin());
+ }
+
+ /** Plays until nothing is held, leaving the position just past the last real frame. */
+ private static void drain(JitterBuffer b) {
+ for (int i = 0; i < 100 && !b.isEmpty(); i++) pull(b);
+ }
+
+ @Test
+ void offersTheFollowingPacketToRebuildALostOne() {
+ JitterBuffer b = buffer();
+ for (int i = 0; i < 12; i++) put(b, i);
+ drain(b);
+
+ // The frame due now never arrives, but the one after it is already here,
+ // carrying the FEC copy that can rebuild it.
+ int due = b.position();
+ b.put(frame(13), due + SPAN, SPAN);
+
+ assertEquals(JitterBuffer.Status.MISSING, b.get(SPAN));
+ assertNotNull(b.nextPayload(), "the successor should be offered for FEC recovery");
+ assertEquals(13, b.nextPayload()[0]);
+ b.tick();
+
+ // And it is still there to be played normally afterwards.
+ assertEquals(13, pull(b));
+ }
+
+ @Test
+ void offersNothingWhenTheSuccessorIsAlsoLost() {
+ JitterBuffer b = buffer();
+ for (int i = 0; i < 12; i++) put(b, i);
+ drain(b);
+
+ // Two frames in a row are gone; FEC only reaches back one.
+ int due = b.position();
+ b.put(frame(15), due + 2 * SPAN, SPAN);
+
+ assertEquals(JitterBuffer.Status.MISSING, b.get(SPAN));
+ assertNull(b.nextPayload(), "a double loss cannot be rebuilt from FEC");
+ }
+
+ @Test
+ void aPauseInSpeechCostsNothingWhenTheClockStopsToo() {
+ JitterBuffer b = buffer();
+ for (int i = 0; i < 40; i++) {
+ put(b, i);
+ pull(b);
+ }
+ drain(b);
+
+ // Silence. No packets arrive and nothing is pulled, so the position stays put —
+ // which is what lets the sender's numbering still line up when speech resumes.
+ int paused = b.position();
+ assertEquals(paused, b.position());
+
+ b.put(frame(7), paused, SPAN);
+ assertEquals(JitterBuffer.Status.OK, b.get(SPAN), "the next burst should play, not be judged late");
+ assertEquals(7, b.payload()[0]);
+ }
+
+ @Test
+ void keepsItsEstimateAcrossAPause() {
+ JitterBuffer b = buffer();
+ Random random = new Random(5);
+ int next = 0;
+ for (int i = 0; i < 300; i++) {
+ if (i % 4 == 0) {
+ for (int j = 0, burst = 2 + random.nextInt(4); j < burst; j++) put(b, next++);
+ }
+ pull(b);
+ }
+ drain(b);
+
+ JitterBuffer fresh = buffer();
+ assertTrue(concealedInBurst(b, b.position() / SPAN) <= concealedInBurst(fresh, 1000),
+ "a buffer that survived the pause should conceal no more than a cold one");
+ }
+
+ /** Plays a short jittery burst through a buffer and counts the concealed frames. */
+ private static int concealedInBurst(JitterBuffer b, int firstFrame) {
+ Random random = new Random(9);
+ int concealed = 0;
+ int next = firstFrame;
+ for (int i = 0; i < 60; i++) {
+ if (i % 4 == 0) {
+ for (int j = 0, burst = 2 + random.nextInt(4); j < burst; j++) {
+ b.put(frame(next), next * SPAN, SPAN);
+ next++;
+ }
+ }
+ if (pull(b) < 0) concealed++;
+ }
+ return concealed;
+ }
+
+ @Test
+ void resetForgetsTheNetworkEstimateToo() {
+ JitterBuffer b = buffer();
+ for (int i = 0; i < 100; i++) {
+ put(b, i);
+ pull(b);
+ }
+ b.reset();
+ // With no timings recorded there is nothing to justify moving the delay.
+ assertTrue(b.isEmpty());
+ assertEquals(0, b.heldFrames(SPAN));
+ }
+
+ @Test
+ void heldFramesCountsWhatIsWaiting() {
+ JitterBuffer b = buffer();
+ assertEquals(0, b.heldFrames(SPAN));
+ put(b, 0);
+ put(b, 1);
+ put(b, 2);
+ assertEquals(3, b.heldFrames(SPAN));
+ }
+
+ @Test
+ void aShortUtteranceComesOutWhole() {
+ JitterBuffer b = buffer();
+ int played = 0;
+ // Five frames, one per frame period — a single short word.
+ for (int i = 0; i < 5; i++) {
+ put(b, i);
+ if (pull(b) >= 0) played++;
+ }
+ // Then whatever the delay is still holding back.
+ for (int i = 0; i < 20 && !b.isEmpty(); i++) {
+ if (pull(b) >= 0) played++;
+ }
+ assertEquals(5, played, "every frame of a short utterance must be heard");
+ }
+
+ @Test
+ void pullingOnlyWhenAPacketIsWaitingStarvesTheStream() {
+ // The buffer adapts its delay in tick(), which only runs on a pull. A consumer that
+ // skips the pull whenever the buffer happens to be empty pins the delay at zero —
+ // and at zero delay the buffer is empty much of the time, so it never recovers.
+ int periods = 300;
+ assertEquals(periods, emitted(true), "a frame per period keeps the device fed");
+
+ int gated = emitted(false);
+ assertTrue(gated < periods * 9 / 10,
+ "gating the pull on a non-empty buffer should visibly starve the device: " + gated);
+ }
+
+ /** Runs a jittery arrival schedule for 300 periods, counting frames actually emitted. */
+ private static int emitted(boolean pullEveryPeriod) {
+ JitterBuffer b = buffer();
+ Random random = new Random(13);
+ int next = 0;
+ int frames = 0;
+ for (int i = 0; i < 300; i++) {
+ // Arrivals clump and stall, as they do on a real connection.
+ if (random.nextInt(4) != 0) {
+ for (int j = 0, burst = random.nextInt(3); j < burst; j++) put(b, next++);
+ }
+ if (pullEveryPeriod || !b.isEmpty()) {
+ b.get(SPAN);
+ b.tick();
+ frames++;
+ }
+ }
+ return frames;
+ }
+}
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
index 423f672..4d197bc 100644
--- a/ts3-client/core/src/test/java/com/ts3client/audio/VoiceStreamTest.java
+++ b/ts3-client/core/src/test/java/com/ts3client/audio/VoiceStreamTest.java
@@ -5,6 +5,7 @@ import org.junit.jupiter.api.Test;
import java.util.ArrayList;
import java.util.List;
+import java.util.Random;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
@@ -21,6 +22,10 @@ class VoiceStreamTest {
return stream.offer(packetId, CodecType.OPUS_VOICE, FakeOpus.packet(value));
}
+ private void offerEnd(int packetId) {
+ stream.offer(packetId, CodecType.OPUS_VOICE, new byte[0]);
+ }
+
/** 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);
@@ -28,13 +33,11 @@ class VoiceStreamTest {
}
@Test
- void holdsBackTheStartOfABurstThenPlaysInOrder() {
+ void playsPacketsInOrderWhateverOrderTheyArriveIn() {
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());
@@ -45,13 +48,10 @@ class VoiceStreamTest {
void concealsALostPacket() {
offer(0, 1);
offer(2, 3);
- pull();
- pull();
assertEquals(0.01f, pull());
assertEquals(FakeOpus.CONCEALED, pull());
assertEquals(0.03f, pull());
- assertEquals(List.of(VoiceFormat.FRAME_SIZE), opus.concealedFrameSizes,
- "a lost packet is concealed at the length of the one before, not the longest Opus frame");
+ assertEquals(List.of(VoiceFormat.FRAME_SIZE), opus.concealedFrameSizes);
}
@Test
@@ -60,38 +60,54 @@ class VoiceStreamTest {
offer(1, 2);
pull();
pull();
- pull();
- pull();
assertFalse(offer(0, 9));
- assertEquals(FakeOpus.CONCEALED, pull());
}
@Test
- void endPacketFinishesTheBurst() {
+ void theTalkersStopEndsTheBurst() {
offer(0, 1);
- stream.offer(1, CodecType.OPUS_VOICE, new byte[0]);
- pull();
- pull();
- pull();
+ offerEnd(1);
+ assertEquals(0.01f, 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() {
+ void theNextBurstCarriesOnWhereTheLastStopped() {
offer(0, 1);
+ offerEnd(1);
pull();
pull();
+
+ assertTrue(offer(2, 2), "the next packet starts a new burst");
+ assertEquals(0f, pull(), "the stop's slot is the talker's own silence, not a loss");
+ assertEquals(0.02f, pull());
+ assertEquals(0, opus.resets, "the decoder keeps its state between bursts");
+ assertEquals(List.of(), opus.concealedFrameSizes);
+ }
+
+ @Test
+ void givesUpOnABurstWhoseStopWasLost() {
+ offer(0, 1);
pull();
int concealed = 0;
while (stream.pull(out, 1, 1.0) != VoiceStream.IDLE) concealed++;
- assertEquals(10, concealed);
+ assertTrue(concealed > 0 && concealed < 10, "concealed " + concealed + " frames");
assertEquals(List.of(true, false), talk);
}
+ @Test
+ void playsALongPacketOneFrameAtATime() {
+ stream.offer(0, CodecType.OPUS_VOICE, new byte[]{40, 0});
+ offer(1, 50);
+ assertEquals(VoiceFormat.FRAME_SIZE, stream.pull(out, 1, 1.0));
+ assertEquals(0.4f, out[0]);
+ assertEquals(VoiceFormat.FRAME_SIZE, stream.pull(out, 1, 1.0));
+ assertEquals(0.4f, out[0]);
+ assertEquals(0.5f, pull());
+ }
+
@Test
void stopFromAnotherThreadEndsTheBurst() {
offer(0, 1);
@@ -103,11 +119,37 @@ class VoiceStreamTest {
@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]);
}
+
+ @Test
+ void neverHandsOutMoreThanOneFrame() {
+ for (int i = 0; i < 50; i++) offer(i, 1);
+ for (int i = 0; i < 50; i++) assertTrue(stream.pull(out, 1, 1.0) <= VoiceFormat.FRAME_SIZE);
+ }
+
+ /**
+ * Packets sent every 20 ms arrive up to 60 ms late. A fixed delay short of that conceals
+ * constantly; the adaptive one should learn the spread and stop concealing.
+ */
+ @Test
+ void learnsADelayThatRidesOutJitter() {
+ Random random = new Random(1);
+ int packets = 1500;
+ long[] arrival = new long[packets];
+ for (int i = 0; i < packets; i++) arrival[i] = i * 20L + random.nextInt(61);
+
+ int next = 0;
+ int concealedLate = 0;
+ for (long now = 0; now < packets * 20L; now += 20) {
+ for (; next < packets && arrival[next] <= now; next++) offer(next, 1);
+ // arrivals are not in order, but whatever is due by now has been offered
+ float sample = pull();
+ if (now >= packets * 10L && sample == FakeOpus.CONCEALED) concealedLate++;
+ }
+ assertTrue(concealedLate < packets / 2 / 50, "concealed " + concealedLate + " frames in the second half");
+ }
}