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 <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1,584 @@
|
||||
package com.ts3client.audio;
|
||||
|
||||
/**
|
||||
* Adaptive jitter buffer for incoming voice packets.
|
||||
*
|
||||
* <p>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.
|
||||
*
|
||||
* <h2>How the delay adapts</h2>
|
||||
* Every packet's arrival is recorded as a <em>timing</em>: 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} <em>latest</em> 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.
|
||||
*
|
||||
* <h2>Timestamp units</h2>
|
||||
* 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.
|
||||
*
|
||||
* <p>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.
|
||||
*
|
||||
* <p>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.
|
||||
*
|
||||
* <p>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.
|
||||
*
|
||||
* <p>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.
|
||||
*
|
||||
* <p>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 <em>not</em> 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.
|
||||
*
|
||||
* <p>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.
|
||||
*
|
||||
* <p>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;
|
||||
}
|
||||
}
|
||||
@@ -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.
|
||||
*
|
||||
* <p>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.
|
||||
*
|
||||
* <p>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.
|
||||
*
|
||||
* <p>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<Integer, Boolean> talkListener;
|
||||
|
||||
private final Object lock = new Object();
|
||||
private final Map<Integer, VoiceFrame> 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);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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<Integer> 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<Integer> 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<Integer> played = new ArrayList<>();
|
||||
for (int i = 0; i < 6; i++) played.add(pull(b));
|
||||
List<Integer> 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;
|
||||
}
|
||||
}
|
||||
@@ -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");
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user