From 0083dc8ad6a798b441df4e2a6ebbb5ea1a7daead Mon Sep 17 00:00:00 2001 From: ericek111 Date: Fri, 25 Sep 2026 10:46:08 +0000 Subject: [PATCH] Play incoming voice through TeamSpeak's adaptive jitter buffer Each speaker's stream held a fixed 60 ms and gave up after 200 ms of concealment, so a jittery connection concealed constantly and a stall ended the burst. It now runs through a port of the libspeex jitter buffer the official client links, configured and driven as it does: - the delay adapts to the measured jitter, aiming for at most 1% of packets arriving late - a speaker keeps one timeline, never re-anchored between talk bursts, so what the buffer learned carries over; only 5 s without packets drops it - the talker's stop marker never enters the buffer; playback ends when it reaches it, and its slot is played as silence, not concealed A pull now always hands out one 20 ms frame: what a longer packet or a delay increase produces beyond that is played on the following pulls, which Android's mixer used to cut off. Co-Authored-By: Claude Opus 5.5 --- .../com/ts3client/audio/JitterBuffer.java | 584 ++++++++++++++++++ .../java/com/ts3client/audio/VoiceStream.java | 327 ++++++---- .../java/com/ts3client/audio/FakeOpus.java | 7 +- .../com/ts3client/audio/JitterBufferTest.java | 353 +++++++++++ .../com/ts3client/audio/VoiceStreamTest.java | 84 ++- 5 files changed, 1226 insertions(+), 129 deletions(-) create mode 100644 ts3-client/core/src/main/java/com/ts3client/audio/JitterBuffer.java create mode 100644 ts3-client/core/src/test/java/com/ts3client/audio/JitterBufferTest.java 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"); + } }