diff --git a/ts3-client/core/src/main/java/com/ts3client/audio/PerSpeakerPlayout.java b/ts3-client/core/src/main/java/com/ts3client/audio/PerSpeakerPlayout.java index dad2173..82d9b40 100644 --- a/ts3-client/core/src/main/java/com/ts3client/audio/PerSpeakerPlayout.java +++ b/ts3-client/core/src/main/java/com/ts3client/audio/PerSpeakerPlayout.java @@ -20,6 +20,8 @@ final class PerSpeakerPlayout implements Playout { private static final long FRAME_NANOS = TimeUnit.SECONDS.toNanos(VoiceFormat.FRAME_SIZE) / VoiceFormat.SAMPLE_RATE; + /** Silence after which a speaker's line and thread are given back; the next burst reopens them. */ + private static final long IDLE_RELEASE_NANOS = TimeUnit.SECONDS.toNanos(30); private final class Speaker { final VoiceStream stream; @@ -32,6 +34,8 @@ final class PerSpeakerPlayout implements Playout { final byte[] out; /** Guarded by this speaker. */ ScheduledFuture ticks; + /** Guarded by this speaker; a closed speaker is never ticked again. */ + boolean closed; Speaker(VoiceStream stream) throws Exception { this.stream = stream; @@ -49,6 +53,7 @@ final class PerSpeakerPlayout implements Playout { void close() { synchronized (this) { + closed = true; if (ticks != null) ticks.cancel(false); } worker.shutdownNow(); @@ -68,30 +73,52 @@ final class PerSpeakerPlayout implements Playout { this.outputDevice = outputDevice; this.gain = gain; this.clock = clock; + clock.scheduleWithFixedDelay(this::releaseIdle, 1, 1, TimeUnit.SECONDS); } @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 + while (true) { + 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.closed) { + // Released a moment ago; this burst gets a fresh line. + speakers.remove(stream, speaker); + continue; + } + if (speaker.ticks == null) { + Speaker s = speaker; + speaker.ticks = clock.scheduleAtFixedRate( + () -> queueTick(s), FRAME_NANOS, FRAME_NANOS, TimeUnit.NANOSECONDS); + } 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); + } + + /** Gives back the line and thread of every speaker that has been silent for a while. */ + private void releaseIdle() { + long now = System.nanoTime(); + for (Speaker speaker : speakers.values()) { + synchronized (speaker) { + if (speaker.ticks != null || now - speaker.stream.lastPacketNanos() < IDLE_RELEASE_NANOS) continue; + speaker.closed = true; } + speakers.remove(speaker.stream, speaker); + speaker.close(); } }