Move the voice pipelines into core
Capture processing, the jitter buffer and decoding were desktop classes, so any other platform would have had to copy them. They now live in core and reach the platform only through two interfaces: - AudioIo lists devices and opens capture/playback lines; the desktop's AudioDevices implements it over PipeWire and Java Sound. - OpusCodec creates encoders and decoders; the desktop binds libopus through the FFM API as before. DesktopVoiceInput becomes CaptureVoiceInput unchanged in behaviour. DesktopVoiceOutput splits into VoiceStream, one speaker's jitter buffer and decoder, paced by whoever pulls it, and StreamingVoiceOutput around it. Speakers reach the device either on a line each (the desktop, so each shows up in the PipeWire mixer) or mixed onto one shared line, as mobile audio APIs want. The settings dialog's device lists and microphone test now go through the AudioBackend instead of desktop classes. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
@@ -1,17 +1,30 @@
|
||||
package com.ts3client.audio;
|
||||
|
||||
import com.ts3client.audio.opus.OpusCodec;
|
||||
import com.ts3client.config.Settings;
|
||||
import com.ts3client.sound.SoundPlayer;
|
||||
|
||||
/**
|
||||
* Factory for a platform's voice capture and playback. Injected into the
|
||||
* connection layer so the core stays independent of any concrete audio stack.
|
||||
* A platform's audio stack: its devices, its libopus and its sound player. Injected into
|
||||
* the connection layer so the core stays independent of any concrete audio stack; the
|
||||
* voice pipelines on top are the same everywhere.
|
||||
*/
|
||||
public interface AudioBackend {
|
||||
|
||||
VoiceInput createInput(Settings settings);
|
||||
AudioIo io();
|
||||
|
||||
VoiceOutput createOutput(Settings settings);
|
||||
OpusCodec opus();
|
||||
|
||||
/** How incoming voice reaches this platform's devices. */
|
||||
StreamingVoiceOutput.Lines outputLines();
|
||||
|
||||
default VoiceInput createInput(Settings settings) {
|
||||
return new CaptureVoiceInput(settings, io(), opus());
|
||||
}
|
||||
|
||||
default VoiceOutput createOutput(Settings settings) {
|
||||
return new StreamingVoiceOutput(io(), opus(), settings.outputDevice, outputLines());
|
||||
}
|
||||
|
||||
/** Player for notification sounds (sound packs); shared by all connections. */
|
||||
SoundPlayer createSoundPlayer(Settings settings);
|
||||
|
||||
@@ -0,0 +1,26 @@
|
||||
package com.ts3client.audio;
|
||||
|
||||
/**
|
||||
* An open microphone line, delivering 16-bit little-endian PCM at
|
||||
* {@link VoiceFormat#SAMPLE_RATE}.
|
||||
*
|
||||
* <p>Modelled on Java Sound's {@code TargetDataLine} so the capture pipeline does
|
||||
* not care which platform audio API delivers it.
|
||||
*/
|
||||
public interface AudioCapture extends AutoCloseable {
|
||||
|
||||
/** Channel count actually negotiated with the device. */
|
||||
int channels();
|
||||
|
||||
/** Begins capturing; audio read before this call is not delivered. */
|
||||
void start();
|
||||
|
||||
/**
|
||||
* Blocks until {@code length} bytes are captured. Returns fewer bytes only when the
|
||||
* line is closing or the calling thread was interrupted.
|
||||
*/
|
||||
int read(byte[] buffer, int offset, int length);
|
||||
|
||||
@Override
|
||||
void close();
|
||||
}
|
||||
@@ -0,0 +1,13 @@
|
||||
package com.ts3client.audio;
|
||||
|
||||
/** A selectable audio device. {@code id} is what gets persisted in the settings. */
|
||||
public record AudioDevice(String id, String label) {
|
||||
|
||||
/** The implicit entry that lets the platform choose. */
|
||||
public static final AudioDevice DEFAULT = new AudioDevice("", "(System default)");
|
||||
|
||||
@Override
|
||||
public String toString() {
|
||||
return label;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,23 @@
|
||||
package com.ts3client.audio;
|
||||
|
||||
import java.util.List;
|
||||
|
||||
/** The platform's audio devices: what can be picked, and opening lines on them. */
|
||||
public interface AudioIo {
|
||||
|
||||
/** Devices that can provide microphone lines, {@link AudioDevice#DEFAULT} first. */
|
||||
List<AudioDevice> inputDevices();
|
||||
|
||||
/** Devices that can provide speaker lines, {@link AudioDevice#DEFAULT} first. */
|
||||
List<AudioDevice> outputDevices();
|
||||
|
||||
/**
|
||||
* Opens a capture line at {@link VoiceFormat#SAMPLE_RATE}, preferring
|
||||
* {@code preferredChannels} but falling back to what the device takes; inspect
|
||||
* {@link AudioCapture#channels()} for the result. An empty id means the default device.
|
||||
*/
|
||||
AudioCapture openCapture(String deviceId, int preferredChannels) throws Exception;
|
||||
|
||||
/** Opens a playback line; see {@link #openCapture} for the channel negotiation. */
|
||||
AudioPlayback openPlayback(String deviceId, int preferredChannels) throws Exception;
|
||||
}
|
||||
@@ -0,0 +1,26 @@
|
||||
package com.ts3client.audio;
|
||||
|
||||
/**
|
||||
* An open speaker line, accepting 16-bit little-endian PCM at
|
||||
* {@link VoiceFormat#SAMPLE_RATE}.
|
||||
*
|
||||
* <p>Modelled on Java Sound's {@code SourceDataLine} so the playback pipeline does
|
||||
* not care which platform audio API plays it.
|
||||
*/
|
||||
public interface AudioPlayback extends AutoCloseable {
|
||||
|
||||
/** Channel count actually negotiated with the device. */
|
||||
int channels();
|
||||
|
||||
/** Begins playing whatever is written from now on. */
|
||||
void start();
|
||||
|
||||
/** Blocks until the audio has been queued for playback. */
|
||||
void write(byte[] buffer, int offset, int length);
|
||||
|
||||
/** Blocks until everything already written has been played out. */
|
||||
void drain();
|
||||
|
||||
@Override
|
||||
void close();
|
||||
}
|
||||
@@ -0,0 +1,530 @@
|
||||
package com.ts3client.audio;
|
||||
|
||||
import com.github.manevolent.ts3j.enums.CodecType;
|
||||
import com.ts3client.audio.opus.OpusCodec;
|
||||
import com.ts3client.audio.opus.OpusEncoder;
|
||||
import com.ts3client.audio.processing.AudioProcessor;
|
||||
import com.ts3client.audio.processing.ns.SuppressionLevel;
|
||||
import com.ts3client.audio.vad.RnnSpeechDetector;
|
||||
import com.ts3client.config.Settings;
|
||||
|
||||
import java.util.concurrent.ConcurrentLinkedQueue;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.function.Consumer;
|
||||
|
||||
/**
|
||||
* The {@link VoiceInput} every platform shares: captures the microphone, applies voice-activation
|
||||
* or push-to-talk gating, and Opus-encodes 20 ms frames.
|
||||
*
|
||||
* <p>The capture line comes from the platform's {@link AudioIo}; it is opened in stereo
|
||||
* when the device offers it. Voice
|
||||
* ({@code OPUS_VOICE}) is transmitted mono, from the downmix, so the pre-processing
|
||||
* chain and VAD see a single channel; the music codec ({@code OPUS_MUSIC}) transmits
|
||||
* the stereo capture as-is.
|
||||
*
|
||||
* <p>A dedicated capture thread fills a small packet queue while ts3j polls
|
||||
* {@link #isReady()} / {@link #provide()} on its own timer, decoupling capture from
|
||||
* network sending. When the gate closes and the queue drains {@link #isReady()}
|
||||
* returns {@code false}, prompting ts3j to emit the terminating empty voice packet.
|
||||
*/
|
||||
public final class CaptureVoiceInput implements VoiceInput {
|
||||
|
||||
/**
|
||||
* How long the gate stays open after the last active frame. TS3 counts 64 of its
|
||||
* 10 ms preprocessor frames before closing, so the tail of a word is never clipped.
|
||||
*/
|
||||
private static final int HANGOVER_MS = 640;
|
||||
|
||||
/**
|
||||
* Audio replayed when the gate opens, so a word's onset is not swallowed by the frame
|
||||
* that detected it. TS3 keeps up to three 10 ms buffers ({@code vad_extrabuffersize}
|
||||
* defaults to 2, plus one).
|
||||
*/
|
||||
private static final int PREROLL_MS = 30;
|
||||
|
||||
private static final int HANGOVER_FRAMES =
|
||||
Math.max(1, HANGOVER_MS * VoiceFormat.SAMPLE_RATE / 1000 / VoiceFormat.FRAME_SIZE);
|
||||
private static final int PREROLL_FRAMES =
|
||||
Math.max(1, PREROLL_MS * VoiceFormat.SAMPLE_RATE / 1000 / VoiceFormat.FRAME_SIZE);
|
||||
|
||||
private final ConcurrentLinkedQueue<byte[]> queue = new ConcurrentLinkedQueue<>();
|
||||
private final AtomicBoolean muted = new AtomicBoolean(false);
|
||||
private final AtomicBoolean localMuted = new AtomicBoolean(false);
|
||||
private final AtomicBoolean transmitting = new AtomicBoolean(false);
|
||||
private final AtomicBoolean pttDown = new AtomicBoolean(false);
|
||||
private final AtomicBoolean running = new AtomicBoolean(false);
|
||||
|
||||
private volatile Settings.InputMode mode;
|
||||
private volatile Settings.VadMode vadMode;
|
||||
private volatile double thresholdDb;
|
||||
private volatile double speechThreshold;
|
||||
private volatile boolean vadOverPtt;
|
||||
private volatile double inputGain;
|
||||
private volatile CodecType codec = CodecType.OPUS_VOICE;
|
||||
|
||||
private final SpeechProbabilityDetector speechDetector =
|
||||
new RnnSpeechDetector(VoiceFormat.SAMPLE_RATE);
|
||||
private final AudioProcessor processor = new AudioProcessor();
|
||||
|
||||
private volatile Consumer<Double> levelListener; // input level, InputLevel scale
|
||||
private volatile Consumer<Boolean> talkListener; // local talk-state changes
|
||||
private volatile Runnable mutedTalkListener; // speech detected while muted
|
||||
private volatile AudioFrameListener monitorListener; // local monitoring of what is sent
|
||||
private boolean mutedTalking;
|
||||
|
||||
private final AudioIo io;
|
||||
private final OpusCodec opus;
|
||||
private final String deviceName;
|
||||
|
||||
private final Object encoderLock = new Object();
|
||||
private volatile OpusParameters params;
|
||||
private OpusParameters appliedParams;
|
||||
private OpusCodec.Application encoderApplication;
|
||||
private volatile int encoderChannels = 1;
|
||||
private volatile int captureChannels = 1;
|
||||
|
||||
private Thread captureThread;
|
||||
private AudioCapture line;
|
||||
private OpusEncoder encoder;
|
||||
|
||||
private int hangover;
|
||||
private boolean lastTransmitting;
|
||||
|
||||
// Ring of recent frames captured while the gate was shut.
|
||||
private final float[][] preroll = new float[PREROLL_FRAMES][];
|
||||
private int prerollTail;
|
||||
private int prerollCount;
|
||||
|
||||
public CaptureVoiceInput(Settings settings, AudioIo io, OpusCodec opus) {
|
||||
this.io = io;
|
||||
this.opus = opus;
|
||||
this.deviceName = settings.inputDevice;
|
||||
this.mode = settings.inputMode;
|
||||
this.vadMode = settings.vadMode;
|
||||
this.thresholdDb = settings.vadThresholdDb;
|
||||
this.speechThreshold = settings.speechThreshold;
|
||||
this.vadOverPtt = settings.vadOverPtt;
|
||||
this.inputGain = settings.inputVolume;
|
||||
this.params = OpusParameters.from(settings);
|
||||
this.codec = codecFor(params);
|
||||
processor.setNoiseSuppression(settings.denoise);
|
||||
processor.setSuppressionLevel(SuppressionLevel.fromDenoiserLevel(settings.denoiserLevel));
|
||||
processor.setTransientSuppression(settings.typingAttenuation);
|
||||
processor.setGainControl(settings.agc);
|
||||
}
|
||||
|
||||
// ---- live configuration (safe to call from the UI thread) ----
|
||||
|
||||
public void setMode(Settings.InputMode mode) {
|
||||
this.mode = mode;
|
||||
}
|
||||
|
||||
public void setVadMode(Settings.VadMode mode) {
|
||||
this.vadMode = mode;
|
||||
speechDetector.reset();
|
||||
}
|
||||
|
||||
public void setThresholdDb(double db) {
|
||||
this.thresholdDb = db;
|
||||
}
|
||||
|
||||
public void setSpeechThreshold(double threshold) {
|
||||
this.speechThreshold = threshold;
|
||||
}
|
||||
|
||||
public void setVadOverPtt(boolean enabled) {
|
||||
this.vadOverPtt = enabled;
|
||||
}
|
||||
|
||||
public void setInputGain(double gain) {
|
||||
this.inputGain = gain;
|
||||
}
|
||||
|
||||
public void setNoiseSuppression(boolean enabled) {
|
||||
processor.setNoiseSuppression(enabled);
|
||||
}
|
||||
|
||||
public void setDenoiserLevel(int level) {
|
||||
processor.setSuppressionLevel(SuppressionLevel.fromDenoiserLevel(level));
|
||||
}
|
||||
|
||||
public void setTypingAttenuation(boolean enabled) {
|
||||
processor.setTransientSuppression(enabled);
|
||||
}
|
||||
|
||||
public void setAgc(boolean enabled) {
|
||||
processor.setGainControl(enabled);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void keyPressed() {
|
||||
processor.keyPressed();
|
||||
}
|
||||
|
||||
/**
|
||||
* Stages new encoder settings; the capture thread picks them up on the next frame
|
||||
* (or {@link #start()} applies them directly when idle).
|
||||
*/
|
||||
public void setOpusParameters(OpusParameters p) {
|
||||
this.params = p;
|
||||
synchronized (encoderLock) {
|
||||
// While capturing, the codec flag flips together with the encoder swap so
|
||||
// no packet is ever tagged with a codec it wasn't encoded for.
|
||||
if (encoder == null) this.codec = codecFor(p);
|
||||
}
|
||||
}
|
||||
|
||||
private static OpusCodec.Application applicationFor(OpusParameters p) {
|
||||
return p.music ? OpusCodec.Application.AUDIO : OpusCodec.Application.VOIP;
|
||||
}
|
||||
|
||||
private static CodecType codecFor(OpusParameters p) {
|
||||
return p.music ? CodecType.OPUS_MUSIC : CodecType.OPUS_VOICE;
|
||||
}
|
||||
|
||||
/**
|
||||
* Channels to transmit: TeamSpeak's {@code OPUS_MUSIC} stream is stereo,
|
||||
* {@code OPUS_VOICE} is mono. A mono-only capture device caps this at one.
|
||||
*/
|
||||
private int channelsFor(OpusParameters p) {
|
||||
return (p.music && captureChannels >= 2) ? 2 : 1;
|
||||
}
|
||||
|
||||
/** Creates or reconfigures the encoder to match {@code p}. Call under {@link #encoderLock}. */
|
||||
private void applyParams(OpusParameters p) {
|
||||
OpusCodec.Application application = applicationFor(p);
|
||||
int channels = channelsFor(p);
|
||||
// Application (VOIP vs AUDIO) and channel count are fixed at creation, so
|
||||
// switching voice<->music means building a new encoder.
|
||||
if (encoder == null || application != encoderApplication || channels != encoderChannels) {
|
||||
OpusEncoder replacement = opus.createEncoder(
|
||||
VoiceFormat.SAMPLE_RATE, VoiceFormat.FRAME_SIZE, channels, application);
|
||||
configureEncoder(replacement, p);
|
||||
OpusEncoder previous = encoder;
|
||||
encoder = replacement;
|
||||
encoderApplication = application;
|
||||
encoderChannels = channels;
|
||||
if (previous != null) previous.close();
|
||||
} else {
|
||||
configureEncoder(encoder, p);
|
||||
}
|
||||
codec = codecFor(p);
|
||||
appliedParams = p;
|
||||
}
|
||||
|
||||
private static void configureEncoder(OpusEncoder enc, OpusParameters p) {
|
||||
enc.setBitrate(p.bitrate);
|
||||
enc.setComplexity(p.complexity);
|
||||
enc.setVbr(p.vbr);
|
||||
enc.setInbandFec(p.fec);
|
||||
enc.setExpectedPacketLoss(p.expectedPacketLoss);
|
||||
enc.setSignal(p.music ? OpusEncoder.Signal.MUSIC : OpusEncoder.Signal.VOICE);
|
||||
}
|
||||
|
||||
public void setLevelListener(Consumer<Double> l) {
|
||||
this.levelListener = l;
|
||||
}
|
||||
|
||||
public void setTalkListener(Consumer<Boolean> l) {
|
||||
this.talkListener = l;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setMonitorListener(AudioFrameListener l) {
|
||||
this.monitorListener = l;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setMutedTalkListener(Runnable l) {
|
||||
this.mutedTalkListener = l;
|
||||
}
|
||||
|
||||
public void setPushToTalk(boolean down) {
|
||||
this.pttDown.set(down);
|
||||
}
|
||||
|
||||
public void setMuted(boolean m) {
|
||||
this.muted.set(m);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setLocalMuted(boolean m) {
|
||||
this.localMuted.set(m);
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean isLocalMuted() {
|
||||
return localMuted.get();
|
||||
}
|
||||
|
||||
// ---- lifecycle ----
|
||||
|
||||
public synchronized void start() {
|
||||
if (running.get()) return;
|
||||
try {
|
||||
line = io.openCapture(deviceName, VoiceFormat.MAX_CHANNELS);
|
||||
captureChannels = line.channels();
|
||||
synchronized (encoderLock) {
|
||||
applyParams(params);
|
||||
}
|
||||
} catch (Throwable t) {
|
||||
cleanup();
|
||||
throw new RuntimeException("Could not start microphone: " + t.getMessage(), t);
|
||||
}
|
||||
speechDetector.reset();
|
||||
processor.reset();
|
||||
running.set(true);
|
||||
captureThread = new Thread(this::captureLoop, "ts3j-mic-capture");
|
||||
captureThread.setDaemon(true);
|
||||
captureThread.start();
|
||||
}
|
||||
|
||||
public synchronized void stop() {
|
||||
running.set(false);
|
||||
if (captureThread != null) {
|
||||
captureThread.interrupt();
|
||||
captureThread = null;
|
||||
}
|
||||
cleanup();
|
||||
queue.clear();
|
||||
setTransmitting(false);
|
||||
}
|
||||
|
||||
private void cleanup() {
|
||||
if (line != null) {
|
||||
line.close();
|
||||
line = null;
|
||||
}
|
||||
synchronized (encoderLock) {
|
||||
if (encoder != null) {
|
||||
try {
|
||||
encoder.close();
|
||||
} catch (Exception ignored) {
|
||||
}
|
||||
encoder = null;
|
||||
encoderApplication = null;
|
||||
appliedParams = null;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private void captureLoop() {
|
||||
final int frameSamples = VoiceFormat.FRAME_SIZE;
|
||||
final int channels = captureChannels;
|
||||
final byte[] buf = new byte[frameSamples * 2 * channels];
|
||||
final float[] pcm = new float[frameSamples * channels]; // interleaved capture
|
||||
// Analysis (level, VAD) always runs on a mono downmix; with a mono device that
|
||||
// is the capture buffer itself, so nothing is copied.
|
||||
final float[] mono = (channels == 1) ? pcm : new float[frameSamples];
|
||||
|
||||
// Held locally so a concurrent stop() closing the line cannot null it mid-loop.
|
||||
final AudioCapture capture = line;
|
||||
capture.start();
|
||||
|
||||
while (running.get()) {
|
||||
// A short read only happens once the line is closing, or on interruption.
|
||||
if (capture.read(buf, 0, buf.length) < buf.length) break;
|
||||
|
||||
// 16-bit LE -> float, with input gain
|
||||
for (int i = 0; i < pcm.length; i++) {
|
||||
int lo = buf[2 * i] & 0xFF;
|
||||
int hi = buf[2 * i + 1];
|
||||
short s = (short) ((hi << 8) | lo);
|
||||
float f = (float) (s / 32768.0 * inputGain);
|
||||
if (f > 1f) f = 1f;
|
||||
else if (f < -1f) f = -1f;
|
||||
pcm[i] = f;
|
||||
}
|
||||
if (channels > 1) {
|
||||
for (int i = 0; i < frameSamples; i++) {
|
||||
float sum = 0;
|
||||
for (int c = 0; c < channels; c++) sum += pcm[i * channels + c];
|
||||
mono[i] = sum / channels;
|
||||
}
|
||||
}
|
||||
|
||||
byte[] packet = null;
|
||||
synchronized (encoderLock) {
|
||||
OpusParameters p = params;
|
||||
if (p != appliedParams && encoder != null) applyParams(p);
|
||||
boolean stereo = encoderChannels == 2;
|
||||
|
||||
// As in TS3, the speech detector judges the raw microphone signal, while
|
||||
// the level meter and volume gate see the processed one.
|
||||
double probability = usesSpeechDetector()
|
||||
? speechDetector.process(mono, frameSamples) : 0.0;
|
||||
|
||||
// Processing is a voice-chain stage on the mono path. A stereo (music)
|
||||
// stream bypasses it and is transmitted as captured.
|
||||
if (!stereo) processor.process(mono, frameSamples);
|
||||
|
||||
double db = InputLevel.toDb(mono, frameSamples);
|
||||
Consumer<Double> ll = levelListener;
|
||||
if (ll != null) ll.accept(db);
|
||||
|
||||
boolean wasOpen = transmitting.get();
|
||||
boolean open = decideGate(db, probability);
|
||||
setTransmitting(open);
|
||||
|
||||
if (open && !muted.get() && encoder != null) {
|
||||
try {
|
||||
// On the opening edge, send the buffered lead-in first so the
|
||||
// word's onset isn't lost to the frame that detected it.
|
||||
if (!wasOpen) {
|
||||
flushPreroll(stereo);
|
||||
}
|
||||
float[] sent = stereo ? pcm : mono;
|
||||
packet = encoder.encode(sent);
|
||||
monitor(sent, stereo ? encoderChannels : 1);
|
||||
} catch (Exception ignored) {
|
||||
}
|
||||
}
|
||||
if (!open) {
|
||||
rememberForPreroll(stereo ? pcm : mono);
|
||||
}
|
||||
}
|
||||
|
||||
if (packet != null && packet.length > 0) {
|
||||
offer(packet);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private void offer(byte[] packet) {
|
||||
queue.offer(packet);
|
||||
// Guard against unbounded growth if the network stalls.
|
||||
while (queue.size() > 10) queue.poll();
|
||||
}
|
||||
|
||||
/** Keeps the most recent frames while the gate is shut, for {@link #flushPreroll}. */
|
||||
private void rememberForPreroll(float[] frame) {
|
||||
float[] slot = preroll[prerollTail];
|
||||
if (slot == null || slot.length != frame.length) {
|
||||
slot = new float[frame.length];
|
||||
preroll[prerollTail] = slot;
|
||||
}
|
||||
System.arraycopy(frame, 0, slot, 0, frame.length);
|
||||
prerollTail = (prerollTail + 1) % PREROLL_FRAMES;
|
||||
if (prerollCount < PREROLL_FRAMES) prerollCount++;
|
||||
}
|
||||
|
||||
/** Encodes and queues the buffered lead-in, oldest first. Call under {@link #encoderLock}. */
|
||||
private void flushPreroll(boolean stereo) {
|
||||
int expected = stereo ? VoiceFormat.FRAME_SIZE * encoderChannels : VoiceFormat.FRAME_SIZE;
|
||||
for (int i = 0; i < prerollCount; i++) {
|
||||
int idx = (prerollTail - prerollCount + i + PREROLL_FRAMES) % PREROLL_FRAMES;
|
||||
float[] frame = preroll[idx];
|
||||
// A channel-count change between capture and flush invalidates the buffer.
|
||||
if (frame == null || frame.length != expected) continue;
|
||||
try {
|
||||
byte[] p = encoder.encode(frame);
|
||||
if (p != null && p.length > 0) offer(p);
|
||||
monitor(frame, stereo ? encoderChannels : 1);
|
||||
} catch (Exception ignored) {
|
||||
}
|
||||
}
|
||||
prerollCount = 0;
|
||||
}
|
||||
|
||||
/** Whether the gate may consult the speech detector, so it must follow the signal. */
|
||||
private boolean usesSpeechDetector() {
|
||||
if (vadMode == Settings.VadMode.VOLUME_GATE) return false;
|
||||
Settings.InputMode m = mode;
|
||||
return m == Settings.InputMode.VOICE_ACTIVATION || (m == Settings.InputMode.PUSH_TO_TALK && vadOverPtt);
|
||||
}
|
||||
|
||||
private boolean decideGate(double db, double probability) {
|
||||
if (muted.get() || localMuted.get()) {
|
||||
hangover = 0;
|
||||
detectMutedSpeech(db);
|
||||
return false;
|
||||
}
|
||||
mutedTalking = false;
|
||||
switch (mode) {
|
||||
case CONTINUOUS:
|
||||
return true;
|
||||
case PUSH_TO_TALK:
|
||||
if (pttDown.get()) return true;
|
||||
return vadOverPtt && voiceActivated(db, probability);
|
||||
case VOICE_ACTIVATION:
|
||||
default:
|
||||
return voiceActivated(db, probability);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Applies the selected VAD mode with a hangover so trailing syllables aren't clipped.
|
||||
*
|
||||
* <p>The modes match the TS3 client's: the volume gate alone, the speech detector
|
||||
* alone, or both together. Volume Gate skips the detector entirely, as TS3 does.
|
||||
*/
|
||||
private boolean voiceActivated(double db, double probability) {
|
||||
boolean detected = switch (vadMode) {
|
||||
case VOLUME_GATE -> db >= thresholdDb;
|
||||
case AUTOMATIC -> probability >= speechThreshold;
|
||||
case HYBRID -> db >= thresholdDb && probability >= speechThreshold;
|
||||
};
|
||||
if (detected) {
|
||||
hangover = HANGOVER_FRAMES;
|
||||
return true;
|
||||
}
|
||||
if (hangover > 0) {
|
||||
hangover--;
|
||||
return true;
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
/**
|
||||
* Reports the start of a talk burst that the mute is swallowing. The volume gate
|
||||
* alone decides here: the speech detector's state is kept for real transmission.
|
||||
*/
|
||||
private void detectMutedSpeech(double db) {
|
||||
boolean talking = db >= thresholdDb;
|
||||
if (talking && !mutedTalking) {
|
||||
Runnable listener = mutedTalkListener;
|
||||
if (listener != null) listener.run();
|
||||
}
|
||||
mutedTalking = talking;
|
||||
}
|
||||
|
||||
/** Hands a transmitted frame to the monitor, if one is attached. */
|
||||
private void monitor(float[] frame, int channels) {
|
||||
AudioFrameListener l = monitorListener;
|
||||
if (l != null) l.onFrame(frame, channels);
|
||||
}
|
||||
|
||||
private void setTransmitting(boolean t) {
|
||||
transmitting.set(t);
|
||||
if (t != lastTransmitting) {
|
||||
lastTransmitting = t;
|
||||
Consumer<Boolean> tl = talkListener;
|
||||
if (tl != null) tl.accept(t);
|
||||
}
|
||||
}
|
||||
|
||||
// ---- ts3j Microphone contract ----
|
||||
|
||||
@Override
|
||||
public boolean isMuted() {
|
||||
return muted.get();
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean isReady() {
|
||||
// Keep the ts3j sender "active" while we're transmitting or still have
|
||||
// buffered packets. When both are false ts3j sends the terminating packet.
|
||||
return transmitting.get() || !queue.isEmpty();
|
||||
}
|
||||
|
||||
@Override
|
||||
public CodecType getCodec() {
|
||||
return codec;
|
||||
}
|
||||
|
||||
@Override
|
||||
public byte[] provide() {
|
||||
byte[] p = queue.poll();
|
||||
return p != null ? p : new byte[0];
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,110 @@
|
||||
package com.ts3client.audio;
|
||||
|
||||
import java.util.Arrays;
|
||||
import java.util.Collection;
|
||||
import java.util.function.Supplier;
|
||||
import java.util.function.ToDoubleFunction;
|
||||
|
||||
/**
|
||||
* Mixes every speaker onto one playback line, as mobile audio APIs want: a single
|
||||
* low-latency stream. The line's blocking writes pace the mix, so it runs on the device's
|
||||
* clock; between talk bursts the mixer writes nothing and waits.
|
||||
*/
|
||||
final class MixedPlayout implements Playout {
|
||||
|
||||
private static final int FRAME = VoiceFormat.FRAME_SIZE;
|
||||
|
||||
private final AudioIo io;
|
||||
private final Supplier<String> outputDevice;
|
||||
private final ToDoubleFunction<VoiceStream> gain;
|
||||
private final Collection<VoiceStream> streams;
|
||||
|
||||
private final Object lock = new Object();
|
||||
/** Guarded by {@link #lock}. */
|
||||
private boolean running = true;
|
||||
private Thread mixer;
|
||||
|
||||
MixedPlayout(AudioIo io, Supplier<String> outputDevice, ToDoubleFunction<VoiceStream> gain,
|
||||
Collection<VoiceStream> streams) {
|
||||
this.io = io;
|
||||
this.outputDevice = outputDevice;
|
||||
this.gain = gain;
|
||||
this.streams = streams;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void started(VoiceStream stream) {
|
||||
synchronized (lock) {
|
||||
if (!running) return;
|
||||
if (mixer == null) {
|
||||
mixer = new Thread(this::mixLoop, "ts3j-voice-mixer");
|
||||
mixer.setDaemon(true);
|
||||
mixer.start();
|
||||
}
|
||||
lock.notifyAll();
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void removed(VoiceStream stream) {
|
||||
// The stream has already left the collection the mixer iterates.
|
||||
}
|
||||
|
||||
@Override
|
||||
public void shutdown() {
|
||||
synchronized (lock) {
|
||||
running = false;
|
||||
if (mixer != null) mixer.interrupt();
|
||||
lock.notifyAll();
|
||||
}
|
||||
}
|
||||
|
||||
private void mixLoop() {
|
||||
AudioPlayback line = null;
|
||||
try {
|
||||
int channels = 0;
|
||||
float[] mix = null;
|
||||
float[] scratch = null;
|
||||
byte[] out = null;
|
||||
while (awaitSpeech()) {
|
||||
if (line == null) {
|
||||
line = io.openPlayback(outputDevice.get(), VoiceFormat.MAX_CHANNELS);
|
||||
line.start();
|
||||
channels = line.channels();
|
||||
mix = new float[FRAME * channels];
|
||||
scratch = new float[VoiceStream.MAX_FRAME * channels];
|
||||
out = new byte[mix.length * 2];
|
||||
}
|
||||
Arrays.fill(mix, 0f);
|
||||
for (VoiceStream stream : streams) {
|
||||
int frames = stream.pull(scratch, channels, gain.applyAsDouble(stream));
|
||||
// TeamSpeak clients send 20 ms packets; anything longer is cut to the frame.
|
||||
int samples = Math.min(frames, FRAME) * channels;
|
||||
for (int i = 0; i < samples; i++) mix[i] += scratch[i];
|
||||
}
|
||||
Pcm16.encode(mix, mix.length, out);
|
||||
line.write(out, 0, out.length);
|
||||
}
|
||||
} catch (Exception ignored) {
|
||||
// No usable line: this mixer ends, and the next talk burst tries again.
|
||||
} finally {
|
||||
if (line != null) line.close();
|
||||
synchronized (lock) {
|
||||
mixer = null;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/** Waits until some stream is playing; false once shut down. */
|
||||
private boolean awaitSpeech() throws InterruptedException {
|
||||
synchronized (lock) {
|
||||
while (running) {
|
||||
for (VoiceStream stream : streams) {
|
||||
if (stream.isPlaying()) return true;
|
||||
}
|
||||
lock.wait();
|
||||
}
|
||||
return false;
|
||||
}
|
||||
}
|
||||
}
|
||||
20
ts3-client/core/src/main/java/com/ts3client/audio/Pcm16.java
Normal file
20
ts3-client/core/src/main/java/com/ts3client/audio/Pcm16.java
Normal file
@@ -0,0 +1,20 @@
|
||||
package com.ts3client.audio;
|
||||
|
||||
/** Conversion between float samples in [-1, 1] and 16-bit little-endian PCM. */
|
||||
public final class Pcm16 {
|
||||
|
||||
private Pcm16() {
|
||||
}
|
||||
|
||||
/** Writes {@code count} samples to {@code out}, clipping anything out of range. */
|
||||
public static void encode(float[] samples, int count, byte[] out) {
|
||||
for (int i = 0; i < count; i++) {
|
||||
float f = samples[i];
|
||||
if (f > 1f) f = 1f;
|
||||
else if (f < -1f) f = -1f;
|
||||
int s = Math.round(f * 32767f);
|
||||
out[2 * i] = (byte) (s & 0xFF);
|
||||
out[2 * i + 1] = (byte) ((s >> 8) & 0xFF);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,149 @@
|
||||
package com.ts3client.audio;
|
||||
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.ExecutorService;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.ScheduledFuture;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.function.Supplier;
|
||||
import java.util.function.ToDoubleFunction;
|
||||
|
||||
/**
|
||||
* Gives every speaker a playback line and worker thread of its own, so simultaneous
|
||||
* speakers are mixed by the audio server and one slow decode never blocks another. On a
|
||||
* PipeWire desktop every speaker then shows up as a separate stream in the mixer.
|
||||
*/
|
||||
final class PerSpeakerPlayout implements Playout {
|
||||
|
||||
private static final long FRAME_NANOS =
|
||||
TimeUnit.SECONDS.toNanos(VoiceFormat.FRAME_SIZE) / VoiceFormat.SAMPLE_RATE;
|
||||
|
||||
private final class Speaker {
|
||||
final VoiceStream stream;
|
||||
final AudioPlayback line;
|
||||
final int lineChannels;
|
||||
final ExecutorService worker;
|
||||
final AtomicBoolean tickQueued = new AtomicBoolean();
|
||||
/** Touched only by {@link #worker}. */
|
||||
final float[] pcm;
|
||||
final byte[] out;
|
||||
/** Guarded by this speaker. */
|
||||
ScheduledFuture<?> ticks;
|
||||
|
||||
Speaker(VoiceStream stream) throws Exception {
|
||||
this.stream = stream;
|
||||
this.line = io.openPlayback(outputDevice.get(), VoiceFormat.MAX_CHANNELS);
|
||||
this.lineChannels = line.channels();
|
||||
this.line.start();
|
||||
this.pcm = new float[VoiceStream.MAX_FRAME * lineChannels];
|
||||
this.out = new byte[pcm.length * 2];
|
||||
this.worker = Executors.newSingleThreadExecutor(r -> {
|
||||
Thread t = new Thread(r, "ts3j-play-" + stream.clientId);
|
||||
t.setDaemon(true);
|
||||
return t;
|
||||
});
|
||||
}
|
||||
|
||||
void close() {
|
||||
synchronized (this) {
|
||||
if (ticks != null) ticks.cancel(false);
|
||||
}
|
||||
worker.shutdownNow();
|
||||
line.close();
|
||||
}
|
||||
}
|
||||
|
||||
private final AudioIo io;
|
||||
private final Supplier<String> outputDevice;
|
||||
private final ToDoubleFunction<VoiceStream> gain;
|
||||
private final ScheduledExecutorService clock;
|
||||
private final Map<VoiceStream, Speaker> speakers = new ConcurrentHashMap<>();
|
||||
|
||||
PerSpeakerPlayout(AudioIo io, Supplier<String> outputDevice, ToDoubleFunction<VoiceStream> gain,
|
||||
ScheduledExecutorService clock) {
|
||||
this.io = io;
|
||||
this.outputDevice = outputDevice;
|
||||
this.gain = gain;
|
||||
this.clock = clock;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void started(VoiceStream stream) {
|
||||
Speaker speaker = speakers.get(stream);
|
||||
if (speaker == null) {
|
||||
try {
|
||||
speaker = new Speaker(stream);
|
||||
} catch (Exception e) {
|
||||
stream.stop(); // couldn't open a line; drop the burst
|
||||
return;
|
||||
}
|
||||
Speaker existing = speakers.putIfAbsent(stream, speaker);
|
||||
if (existing != null) {
|
||||
speaker.close();
|
||||
speaker = existing;
|
||||
}
|
||||
}
|
||||
synchronized (speaker) {
|
||||
if (speaker.ticks == null) {
|
||||
Speaker s = speaker;
|
||||
speaker.ticks = clock.scheduleAtFixedRate(
|
||||
() -> queueTick(s), FRAME_NANOS, FRAME_NANOS, TimeUnit.NANOSECONDS);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private void queueTick(Speaker speaker) {
|
||||
if (!speaker.tickQueued.compareAndSet(false, true)) return;
|
||||
try {
|
||||
speaker.worker.submit(() -> {
|
||||
try {
|
||||
tick(speaker);
|
||||
} finally {
|
||||
speaker.tickQueued.set(false);
|
||||
}
|
||||
});
|
||||
} catch (RuntimeException e) {
|
||||
if (speaker.worker.isShutdown()) {
|
||||
speaker.tickQueued.set(false);
|
||||
} else {
|
||||
throw e;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private void tick(Speaker speaker) {
|
||||
int frames = speaker.stream.pull(speaker.pcm, speaker.lineChannels, gain.applyAsDouble(speaker.stream));
|
||||
if (frames == VoiceStream.IDLE) {
|
||||
synchronized (speaker) {
|
||||
// A packet may have started the next burst since the pull; keep ticking then.
|
||||
if (!speaker.stream.isPlaying() && speaker.ticks != null) {
|
||||
speaker.ticks.cancel(false);
|
||||
speaker.ticks = null;
|
||||
}
|
||||
}
|
||||
return;
|
||||
}
|
||||
if (frames == 0) return;
|
||||
int samples = frames * speaker.lineChannels;
|
||||
Pcm16.encode(speaker.pcm, samples, speaker.out);
|
||||
try {
|
||||
speaker.line.write(speaker.out, 0, samples * 2);
|
||||
} catch (RuntimeException ignored) {
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void removed(VoiceStream stream) {
|
||||
Speaker speaker = speakers.remove(stream);
|
||||
if (speaker != null) speaker.close();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void shutdown() {
|
||||
for (Speaker speaker : speakers.values()) speaker.close();
|
||||
speakers.clear();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,13 @@
|
||||
package com.ts3client.audio;
|
||||
|
||||
/** Drives the {@link VoiceStream}s of a {@link StreamingVoiceOutput} out to the speakers. */
|
||||
interface Playout {
|
||||
|
||||
/** A stream began a talk burst and must now be pulled every 20 ms. */
|
||||
void started(VoiceStream stream);
|
||||
|
||||
/** The stream's client left; it will not be pulled again. */
|
||||
void removed(VoiceStream stream);
|
||||
|
||||
void shutdown();
|
||||
}
|
||||
@@ -0,0 +1,182 @@
|
||||
package com.ts3client.audio;
|
||||
|
||||
import com.github.manevolent.ts3j.enums.CodecType;
|
||||
import com.github.manevolent.ts3j.protocol.packet.PacketBody0Voice;
|
||||
import com.github.manevolent.ts3j.protocol.packet.PacketBody1VoiceWhisper;
|
||||
import com.ts3client.audio.opus.OpusCodec;
|
||||
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.function.BiConsumer;
|
||||
|
||||
/**
|
||||
* The {@link VoiceOutput} every platform shares: decodes incoming voice per speaker, each
|
||||
* through its own jitter buffer, and plays it on the platform's {@link AudioIo}.
|
||||
*/
|
||||
public final class StreamingVoiceOutput implements VoiceOutput {
|
||||
|
||||
/** How speakers reach the device. */
|
||||
public enum Lines {
|
||||
/** A line per speaker, mixed by the platform's audio server. */
|
||||
PER_SPEAKER,
|
||||
/** One line all speakers are mixed onto here. */
|
||||
SHARED
|
||||
}
|
||||
|
||||
/**
|
||||
* A speaker is only meant to stop when the empty voice packet marking the end of a
|
||||
* talk burst arrives — but that packet is UDP too, and a lost one otherwise leaves
|
||||
* the talking indicator stuck until the speaker's next burst. Opus frames are 20 ms
|
||||
* apart while someone is actually talking, so a gap several times that long, with no
|
||||
* packet of either kind, is unambiguous: force the indicator off rather than trust
|
||||
* the one packet that could go missing.
|
||||
*/
|
||||
private static final long TALK_TIMEOUT_NANOS = TimeUnit.MILLISECONDS.toNanos(400);
|
||||
|
||||
private final OpusCodec opus;
|
||||
private final Map<Integer, VoiceStream> streams = new ConcurrentHashMap<>();
|
||||
private final Set<Integer> mutedClients = ConcurrentHashMap.newKeySet();
|
||||
/** Per-client gain on top of the master volume; absent means 1.0. */
|
||||
private final Map<Integer, Double> clientVolumes = new ConcurrentHashMap<>();
|
||||
|
||||
private volatile String outputDevice;
|
||||
private volatile double masterVolume = 1.0;
|
||||
private volatile boolean deafened = false;
|
||||
|
||||
/** Notified (clientId, talking) on an audio thread when a speaker starts/stops. */
|
||||
private volatile BiConsumer<Integer, Boolean> talkListener;
|
||||
|
||||
/** Catches talk bursts whose end packet never arrived, and paces per-speaker lines. */
|
||||
private final ScheduledExecutorService clock = Executors.newSingleThreadScheduledExecutor(r -> {
|
||||
Thread t = new Thread(r, "ts3j-talk-watchdog");
|
||||
t.setDaemon(true);
|
||||
return t;
|
||||
});
|
||||
private final Playout playout;
|
||||
|
||||
public StreamingVoiceOutput(AudioIo io, OpusCodec opus, String outputDevice, Lines lines) {
|
||||
this.opus = opus;
|
||||
this.outputDevice = outputDevice;
|
||||
this.playout = switch (lines) {
|
||||
case PER_SPEAKER -> new PerSpeakerPlayout(io, () -> this.outputDevice, this::gainFor, clock);
|
||||
case SHARED -> new MixedPlayout(io, () -> this.outputDevice, this::gainFor, streams.values());
|
||||
};
|
||||
clock.scheduleWithFixedDelay(this::checkTalkTimeouts, 50, 50, TimeUnit.MILLISECONDS);
|
||||
}
|
||||
|
||||
private void checkTalkTimeouts() {
|
||||
long now = System.nanoTime();
|
||||
for (VoiceStream s : streams.values()) {
|
||||
if (s.isTalking() && now - s.lastPacketNanos() > TALK_TIMEOUT_NANOS) s.stop();
|
||||
}
|
||||
}
|
||||
|
||||
private double gainFor(VoiceStream stream) {
|
||||
return masterVolume * clientVolumes.getOrDefault(stream.clientId, 1.0);
|
||||
}
|
||||
|
||||
private void fireTalk(int clientId, boolean talking) {
|
||||
BiConsumer<Integer, Boolean> l = talkListener;
|
||||
if (l != null) l.accept(clientId, talking);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setTalkListener(BiConsumer<Integer, Boolean> l) {
|
||||
this.talkListener = l;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setMasterVolume(double v) {
|
||||
this.masterVolume = Math.max(0, Math.min(2.0, v));
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setDeafened(boolean d) {
|
||||
this.deafened = d;
|
||||
if (d) {
|
||||
// Stop everyone talking immediately.
|
||||
for (VoiceStream s : streams.values()) s.stop();
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean isDeafened() {
|
||||
return deafened;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setClientMuted(int clientId, boolean muted) {
|
||||
if (muted) {
|
||||
mutedClients.add(clientId);
|
||||
VoiceStream s = streams.get(clientId);
|
||||
if (s != null) s.stop();
|
||||
} else {
|
||||
mutedClients.remove(clientId);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean isClientMuted(int clientId) {
|
||||
return mutedClients.contains(clientId);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setClientVolume(int clientId, double gain) {
|
||||
if (gain == 1.0) clientVolumes.remove(clientId);
|
||||
else clientVolumes.put(clientId, Math.max(0, gain));
|
||||
}
|
||||
|
||||
@Override
|
||||
public double getClientVolume(int clientId) {
|
||||
return clientVolumes.getOrDefault(clientId, 1.0);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setOutputDevice(String device) {
|
||||
this.outputDevice = device;
|
||||
}
|
||||
|
||||
/** Entry point wired into {@code client.setVoiceHandler(...)}. */
|
||||
@Override
|
||||
public void handleVoice(PacketBody0Voice voice) {
|
||||
route(voice.getClientId(), voice.getPacketId(), voice.getCodecType(), voice.getCodecData());
|
||||
}
|
||||
|
||||
/** Entry point wired into {@code client.setWhisperHandler(...)}. */
|
||||
@Override
|
||||
public void handleWhisper(PacketBody1VoiceWhisper whisper) {
|
||||
route(whisper.getClientId(), whisper.getPacketId(), whisper.getCodecType(), whisper.getCodecData());
|
||||
}
|
||||
|
||||
private void route(int clientId, int packetId, CodecType codec, byte[] data) {
|
||||
if (deafened) return;
|
||||
if (mutedClients.contains(clientId)) return;
|
||||
|
||||
VoiceStream stream = streams.computeIfAbsent(clientId, id -> new VoiceStream(id, opus, this::fireTalk));
|
||||
if (stream.offer(packetId, codec, data)) playout.started(stream);
|
||||
}
|
||||
|
||||
/** Forgets a client that left; the server may hand its id to somebody else. */
|
||||
@Override
|
||||
public void removeClient(int clientId) {
|
||||
VoiceStream s = streams.remove(clientId);
|
||||
if (s != null) {
|
||||
playout.removed(s);
|
||||
s.close();
|
||||
}
|
||||
mutedClients.remove(clientId);
|
||||
clientVolumes.remove(clientId);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void shutdown() {
|
||||
clock.shutdownNow();
|
||||
playout.shutdown();
|
||||
for (VoiceStream s : streams.values()) s.close();
|
||||
streams.clear();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,16 @@
|
||||
package com.ts3client.audio;
|
||||
|
||||
/**
|
||||
* The PCM format voice is captured, processed and played in. TeamSpeak's Opus codecs run at
|
||||
* 48 kHz in 20 ms frames, mono for {@code OPUS_VOICE} and stereo for {@code OPUS_MUSIC}.
|
||||
*/
|
||||
public final class VoiceFormat {
|
||||
|
||||
public static final int SAMPLE_RATE = 48_000;
|
||||
/** Samples per channel in one 20 ms frame. */
|
||||
public static final int FRAME_SIZE = 960;
|
||||
public static final int MAX_CHANNELS = 2;
|
||||
|
||||
private VoiceFormat() {
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,254 @@
|
||||
package com.ts3client.audio;
|
||||
|
||||
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.
|
||||
*
|
||||
* <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
|
||||
* must come from one thread at a time.
|
||||
*/
|
||||
final class VoiceStream {
|
||||
|
||||
/** Returned by {@link #pull} when no talk burst is playing. */
|
||||
static final int IDLE = -1;
|
||||
|
||||
/** 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) {
|
||||
}
|
||||
|
||||
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<>();
|
||||
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;
|
||||
|
||||
/** Guards the decoder: a stop or 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];
|
||||
private OpusDecoder decoder;
|
||||
private int decoderChannels;
|
||||
private boolean closed;
|
||||
|
||||
private volatile boolean talking;
|
||||
private volatile long lastPacketNanos;
|
||||
|
||||
VoiceStream(int clientId, OpusCodec opus, BiConsumer<Integer, Boolean> talkListener) {
|
||||
this.clientId = clientId;
|
||||
this.opus = opus;
|
||||
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.
|
||||
*
|
||||
* @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);
|
||||
boolean started = false;
|
||||
|
||||
synchronized (lock) {
|
||||
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
|
||||
}
|
||||
pending.putIfAbsent(frame.packetId(), frame);
|
||||
trimPending();
|
||||
}
|
||||
if (!frame.end()) markTalking(true);
|
||||
return started;
|
||||
}
|
||||
|
||||
/**
|
||||
* 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
|
||||
*/
|
||||
int pull(float[] out, int outChannels, double gain) {
|
||||
VoiceFrame frame;
|
||||
int channels;
|
||||
synchronized (lock) {
|
||||
if (!playing) return IDLE;
|
||||
if (warmup > 0) {
|
||||
warmup--;
|
||||
return 0;
|
||||
}
|
||||
|
||||
frame = pending.remove(expectedPacketId);
|
||||
boolean conceal = frame == null;
|
||||
if (conceal) {
|
||||
consecutivePlcFrames++;
|
||||
} 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;
|
||||
}
|
||||
}
|
||||
if (channels == 0) {
|
||||
finish();
|
||||
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 {
|
||||
frames = decoderFor(channels).decode(data, pcm);
|
||||
} catch (RuntimeException e) {
|
||||
return 0;
|
||||
}
|
||||
}
|
||||
// Match the decoded stream to the output: duplicate mono across stereo, fold a
|
||||
// stereo (music) stream down onto mono.
|
||||
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];
|
||||
} else if (channels == 1) {
|
||||
v = pcm[i];
|
||||
} else {
|
||||
double sum = 0;
|
||||
for (int s = 0; s < channels; s++) sum += pcm[i * channels + s];
|
||||
v = sum / channels;
|
||||
}
|
||||
out[k] = (float) (v * gain);
|
||||
}
|
||||
}
|
||||
return frames;
|
||||
}
|
||||
|
||||
/**
|
||||
* Returns a decoder matching the stream's channel count, rebuilding it when a speaker
|
||||
* switches between the mono {@code OPUS_VOICE} and stereo {@code OPUS_MUSIC} codecs.
|
||||
*/
|
||||
private OpusDecoder decoderFor(int channels) {
|
||||
if (decoder == null || decoderChannels != channels) {
|
||||
if (decoder != null) decoder.close();
|
||||
decoder = opus.createDecoder(VoiceFormat.SAMPLE_RATE, MAX_FRAME, channels);
|
||||
decoderChannels = channels;
|
||||
}
|
||||
return decoder;
|
||||
}
|
||||
|
||||
/** Ends the current talk burst, if any; the next packet starts a new one. */
|
||||
void stop() {
|
||||
synchronized (lock) {
|
||||
stopLocked();
|
||||
}
|
||||
finish();
|
||||
}
|
||||
|
||||
private void stopLocked() {
|
||||
playing = false;
|
||||
pending.clear();
|
||||
}
|
||||
|
||||
private void finish() {
|
||||
synchronized (decoderLock) {
|
||||
if (decoder != null) decoder.reset();
|
||||
}
|
||||
markTalking(false);
|
||||
}
|
||||
|
||||
boolean isPlaying() {
|
||||
synchronized (lock) {
|
||||
return playing;
|
||||
}
|
||||
}
|
||||
|
||||
boolean isTalking() {
|
||||
return talking;
|
||||
}
|
||||
|
||||
long lastPacketNanos() {
|
||||
return lastPacketNanos;
|
||||
}
|
||||
|
||||
void close() {
|
||||
synchronized (lock) {
|
||||
stopLocked();
|
||||
}
|
||||
synchronized (decoderLock) {
|
||||
closed = true;
|
||||
if (decoder != null) decoder.close();
|
||||
decoder = null;
|
||||
}
|
||||
}
|
||||
|
||||
private void markTalking(boolean talking) {
|
||||
if (this.talking == talking) return;
|
||||
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);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,16 @@
|
||||
package com.ts3client.audio.opus;
|
||||
|
||||
/** Creates Opus encoders and decoders; each platform binds its own libopus. */
|
||||
public interface OpusCodec {
|
||||
|
||||
/** Opus' coding mode, fixed when an encoder is created. */
|
||||
enum Application { VOIP, AUDIO }
|
||||
|
||||
OpusEncoder createEncoder(int sampleRate, int frameSize, int channels, Application application);
|
||||
|
||||
/** @param maxFrameSize the longest frame, per channel, a packet may decode to */
|
||||
OpusDecoder createDecoder(int sampleRate, int maxFrameSize, int channels);
|
||||
|
||||
/** libopus' version string, for display. */
|
||||
String version();
|
||||
}
|
||||
@@ -0,0 +1,17 @@
|
||||
package com.ts3client.audio.opus;
|
||||
|
||||
/** One Opus decoder, used from one thread at a time. */
|
||||
public interface OpusDecoder extends AutoCloseable {
|
||||
|
||||
/**
|
||||
* Decodes {@code packet} into interleaved samples, or conceals a lost packet when it is
|
||||
* {@code null}. Returns the samples decoded per channel.
|
||||
*/
|
||||
int decode(byte[] packet, float[] out);
|
||||
|
||||
/** Forgets the stream's state, as at the start of a new talk burst. */
|
||||
void reset();
|
||||
|
||||
@Override
|
||||
void close();
|
||||
}
|
||||
@@ -0,0 +1,25 @@
|
||||
package com.ts3client.audio.opus;
|
||||
|
||||
/** One Opus encoder. Its methods may be called from different threads, but not concurrently. */
|
||||
public interface OpusEncoder extends AutoCloseable {
|
||||
|
||||
enum Signal { VOICE, MUSIC }
|
||||
|
||||
void setBitrate(int bitsPerSecond);
|
||||
|
||||
void setComplexity(int complexity);
|
||||
|
||||
void setVbr(boolean vbr);
|
||||
|
||||
void setInbandFec(boolean fec);
|
||||
|
||||
void setExpectedPacketLoss(int percent);
|
||||
|
||||
void setSignal(Signal signal);
|
||||
|
||||
/** Encodes exactly one frame of interleaved samples. */
|
||||
byte[] encode(float[] pcm);
|
||||
|
||||
@Override
|
||||
void close();
|
||||
}
|
||||
@@ -0,0 +1,53 @@
|
||||
package com.ts3client.audio;
|
||||
|
||||
import com.ts3client.audio.opus.OpusCodec;
|
||||
import com.ts3client.audio.opus.OpusDecoder;
|
||||
import com.ts3client.audio.opus.OpusEncoder;
|
||||
|
||||
import java.util.Arrays;
|
||||
|
||||
/**
|
||||
* Stands in for libopus: a "packet" is one byte, decoded to a 20 ms frame holding that byte
|
||||
* divided by 100 in every sample. Concealment decodes to -1, so it is easy to spot.
|
||||
*/
|
||||
final class FakeOpus implements OpusCodec {
|
||||
|
||||
static final float CONCEALED = -1f;
|
||||
|
||||
int resets;
|
||||
|
||||
static byte[] packet(int value) {
|
||||
return new byte[]{(byte) value};
|
||||
}
|
||||
|
||||
@Override
|
||||
public OpusEncoder createEncoder(int sampleRate, int frameSize, int channels, Application application) {
|
||||
throw new UnsupportedOperationException();
|
||||
}
|
||||
|
||||
@Override
|
||||
public OpusDecoder createDecoder(int sampleRate, int maxFrameSize, int channels) {
|
||||
return new OpusDecoder() {
|
||||
@Override
|
||||
public int decode(byte[] packet, float[] out) {
|
||||
float v = packet == null ? CONCEALED : packet[0] / 100f;
|
||||
Arrays.fill(out, 0, VoiceFormat.FRAME_SIZE * channels, v);
|
||||
return VoiceFormat.FRAME_SIZE;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void reset() {
|
||||
resets++;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
public String version() {
|
||||
return "fake";
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,102 @@
|
||||
package com.ts3client.audio;
|
||||
|
||||
import com.github.manevolent.ts3j.enums.CodecType;
|
||||
import com.github.manevolent.ts3j.protocol.ProtocolRole;
|
||||
import com.github.manevolent.ts3j.protocol.packet.PacketBody0Voice;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.concurrent.BlockingQueue;
|
||||
import java.util.concurrent.LinkedBlockingQueue;
|
||||
import java.util.concurrent.Semaphore;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertNotNull;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
class MixedPlayoutTest {
|
||||
|
||||
/** A mono line that records every frame written and blocks until the test lets it go on. */
|
||||
private static final class RecordingLine implements AudioPlayback {
|
||||
final BlockingQueue<Short> firstSamples = new LinkedBlockingQueue<>();
|
||||
final Semaphore permits = new Semaphore(0);
|
||||
|
||||
@Override
|
||||
public int channels() {
|
||||
return 1;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void start() {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void write(byte[] buffer, int offset, int length) {
|
||||
firstSamples.add((short) ((buffer[offset + 1] << 8) | (buffer[offset] & 0xFF)));
|
||||
permits.acquireUninterruptibly();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void drain() {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void mixesSimultaneousSpeakersOntoOneLine() throws Exception {
|
||||
RecordingLine line = new RecordingLine();
|
||||
AudioIo io = new AudioIo() {
|
||||
@Override
|
||||
public List<AudioDevice> inputDevices() {
|
||||
return List.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<AudioDevice> outputDevices() {
|
||||
return List.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public AudioCapture openCapture(String deviceId, int preferredChannels) {
|
||||
throw new UnsupportedOperationException();
|
||||
}
|
||||
|
||||
@Override
|
||||
public AudioPlayback openPlayback(String deviceId, int preferredChannels) {
|
||||
return line;
|
||||
}
|
||||
};
|
||||
StreamingVoiceOutput output = new StreamingVoiceOutput(io, new FakeOpus(), "", StreamingVoiceOutput.Lines.SHARED);
|
||||
try {
|
||||
for (int i = 0; i < 10; i++) {
|
||||
output.handleVoice(voice(1, i, 25));
|
||||
output.handleVoice(voice(2, i, 50));
|
||||
}
|
||||
line.permits.release(100);
|
||||
|
||||
short expected = (short) Math.round(0.75f * 32767f);
|
||||
boolean mixed = false;
|
||||
for (int i = 0; i < 20 && !mixed; i++) {
|
||||
Short sample = line.firstSamples.poll(1, TimeUnit.SECONDS);
|
||||
assertNotNull(sample, "the mixer stopped writing");
|
||||
mixed = sample == expected;
|
||||
}
|
||||
assertTrue(mixed, "no frame carried both speakers");
|
||||
} finally {
|
||||
output.shutdown();
|
||||
line.permits.release(1000);
|
||||
}
|
||||
}
|
||||
|
||||
private static PacketBody0Voice voice(int clientId, int packetId, int value) {
|
||||
PacketBody0Voice voice = new PacketBody0Voice(ProtocolRole.SERVER);
|
||||
voice.setClientId(clientId);
|
||||
voice.setPacketId(packetId);
|
||||
voice.setCodecType(CodecType.OPUS_VOICE);
|
||||
voice.setCodecData(FakeOpus.packet(value));
|
||||
return voice;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,111 @@
|
||||
package com.ts3client.audio;
|
||||
|
||||
import com.github.manevolent.ts3j.enums.CodecType;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
class VoiceStreamTest {
|
||||
|
||||
private final FakeOpus opus = new FakeOpus();
|
||||
private final List<Boolean> talk = new ArrayList<>();
|
||||
private final VoiceStream stream = new VoiceStream(7, opus, (id, talking) -> talk.add(talking));
|
||||
private final float[] out = new float[VoiceStream.MAX_FRAME * 2];
|
||||
|
||||
private boolean offer(int packetId, int value) {
|
||||
return stream.offer(packetId, CodecType.OPUS_VOICE, FakeOpus.packet(value));
|
||||
}
|
||||
|
||||
/** 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);
|
||||
return frames > 0 ? out[0] : Float.NaN;
|
||||
}
|
||||
|
||||
@Test
|
||||
void holdsBackTheStartOfABurstThenPlaysInOrder() {
|
||||
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());
|
||||
assertEquals(List.of(true), talk);
|
||||
}
|
||||
|
||||
@Test
|
||||
void concealsALostPacket() {
|
||||
offer(0, 1);
|
||||
offer(2, 3);
|
||||
pull();
|
||||
pull();
|
||||
assertEquals(0.01f, pull());
|
||||
assertEquals(FakeOpus.CONCEALED, pull());
|
||||
assertEquals(0.03f, pull());
|
||||
}
|
||||
|
||||
@Test
|
||||
void dropsPacketsThatArriveTooLate() {
|
||||
offer(0, 1);
|
||||
offer(1, 2);
|
||||
pull();
|
||||
pull();
|
||||
pull();
|
||||
pull();
|
||||
assertFalse(offer(0, 9));
|
||||
assertEquals(FakeOpus.CONCEALED, pull());
|
||||
}
|
||||
|
||||
@Test
|
||||
void endPacketFinishesTheBurst() {
|
||||
offer(0, 1);
|
||||
stream.offer(1, CodecType.OPUS_VOICE, new byte[0]);
|
||||
pull();
|
||||
pull();
|
||||
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() {
|
||||
offer(0, 1);
|
||||
pull();
|
||||
pull();
|
||||
pull();
|
||||
int concealed = 0;
|
||||
while (stream.pull(out, 1, 1.0) != VoiceStream.IDLE) concealed++;
|
||||
assertEquals(10, concealed);
|
||||
assertEquals(List.of(true, false), talk);
|
||||
}
|
||||
|
||||
@Test
|
||||
void stopFromAnotherThreadEndsTheBurst() {
|
||||
offer(0, 1);
|
||||
stream.stop();
|
||||
assertEquals(VoiceStream.IDLE, stream.pull(out, 1, 1.0));
|
||||
assertEquals(List.of(true, false), talk);
|
||||
}
|
||||
|
||||
@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]);
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user