Skip to content

Commit f1ae455

Browse files
mondainclaude
andcommitted
Fix mediabunny stale state on republish, stream close cleanup, and mid-stream subscriber keyframe delivery
Replace computeIfAbsent with compute in MediaBunnyStreamRegistry.subscribe() to detect when a publisher has changed and detach the orphaned listener. Add onClosed() default method to IStreamListener and call it from ClientBroadcastStream.close() before clearing listeners. The mediabunny listener flushes remaining buffers and notifies the registry to remove the stale state and send a poison-pill to unblock waiting servlet threads. Cache the last video keyframe fragment on StreamState so new mid-stream subscribers receive init segment + keyframe fragment, preventing the VideoDecoder key chunk error on the player side. Co-Authored-By: Claude <noreply@anthropic.com>
1 parent 9e5f00c commit f1ae455

5 files changed

Lines changed: 83 additions & 5 deletions

File tree

common/src/main/java/org/red5/server/api/stream/IStreamListener.java

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,13 @@ public interface IStreamListener {
2323
* @param packet
2424
* the packet received
2525
*/
26-
public void packetReceived(IBroadcastStream stream, IStreamPacket packet);
26+
void packetReceived(IBroadcastStream stream, IStreamPacket packet);
27+
28+
/**
29+
* Stream is closed, notify the listener and release all resources.
30+
*/
31+
default void onClosed() {
32+
// default implementation does nothing
33+
}
2734

2835
}

common/src/main/java/org/red5/server/stream/ClientBroadcastStream.java

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -236,8 +236,15 @@ public void close() {
236236
if (recordingListener != null) {
237237
recordingListener.clear();
238238
}
239-
// clear listeners
239+
// notify and clear listeners
240240
if (!listeners.isEmpty()) {
241+
for (IStreamListener listener : listeners) {
242+
try {
243+
listener.onClosed();
244+
} catch (Exception e) {
245+
log.warn("Error notifying listener onClosed: {}", listener, e);
246+
}
247+
}
241248
listeners.clear();
242249
}
243250
// deregister with jmx

server/src/main/java/org/red5/server/net/mediabunny/MediaBunnyServlet.java

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -117,6 +117,10 @@ private void streamQueue(AsyncContext asyncContext, MediaBunnyStreamRegistry.Str
117117
BlockingQueue<byte[]> queue = subscription.getQueue();
118118
while (true) {
119119
byte[] chunk = queue.take();
120+
if (chunk.length == 0) {
121+
log.debug("MediaBunny received end-of-stream signal");
122+
break;
123+
}
120124
chunkCount++;
121125
if (!loggedFirstChunk) {
122126
loggedFirstChunk = true;

server/src/main/java/org/red5/server/net/mediabunny/MediaBunnyStreamListener.java

Lines changed: 19 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -155,6 +155,14 @@ public void packetReceived(IBroadcastStream stream, IStreamPacket packet) {
155155
}
156156
}
157157

158+
@Override
159+
public void onClosed() {
160+
log.info("Closing MediaBunnyStreamListener for stream {}", streamKey);
161+
flushVideoBuffer();
162+
flushAudioBuffer();
163+
registry.onStreamClosed(streamKey);
164+
}
165+
158166
private void handleVideo(VideoData video) {
159167
if (video.getCodecId() != VideoCodec.AVC.getId()) {
160168
if (log.isDebugEnabled()) {
@@ -220,6 +228,7 @@ private void handleVideo(VideoData video) {
220228
} else if (!videoBuffer.samples.isEmpty()) {
221229
flushVideoBuffer();
222230
}
231+
videoBuffer.startsWithKeyframe = true;
223232
} else if (!videoKeyframeSeen) {
224233
return;
225234
}
@@ -403,16 +412,21 @@ private void flushBuffer(TrackBuffer buffer, long sequence, long trackId, long g
403412
if (buffer.samples.isEmpty()) {
404413
return;
405414
}
415+
boolean keyframe = buffer.startsWithKeyframe;
406416
Fmp4FragmentBuilder.FragmentConfig config = new Fmp4FragmentBuilder.FragmentConfig().setSequenceNumber(sequence).setTrackId(trackId).setBaseDecodeTime(buffer.fragmentStartDecodeTime).setGroupId(groupId).setMediaType(mediaType).setMediaData(buffer.media.toByteArray()).setSamples(new ArrayList<>(buffer.samples));
407417
byte[] fragment;
408418
try {
409419
fragment = fragmentBuilder.buildFragment(config).serialize();
410420
patchBrandBox(fragment, "styp", "iso6");
411421
fragment = stripLeadingBox(fragment, "styp");
412422
if (log.isDebugEnabled()) {
413-
log.debug("Built {} fragment for stream {} ({} bytes)", mediaType, streamKey, fragment.length);
423+
log.debug("Built {} fragment for stream {} ({} bytes, keyframe={})", mediaType, streamKey, fragment.length, keyframe);
424+
}
425+
if (keyframe && mediaType == CmafFragment.MediaType.VIDEO) {
426+
registry.onKeyframeFragment(streamKey, fragment);
427+
} else {
428+
registry.onFragment(streamKey, fragment);
414429
}
415-
registry.onFragment(streamKey, fragment);
416430
} catch (IOException e) {
417431
log.warn("Failed to build fragment for stream: {}", streamKey, e);
418432
}
@@ -888,11 +902,14 @@ private static class TrackBuffer {
888902

889903
private long fragmentStartDecodeTime;
890904

905+
private boolean startsWithKeyframe;
906+
891907
private void reset() {
892908
media.reset();
893909
samples.clear();
894910
bufferedDuration = 0;
895911
fragmentStartDecodeTime = 0;
912+
startsWithKeyframe = false;
896913
}
897914
}
898915
}

server/src/main/java/org/red5/server/net/mediabunny/MediaBunnyStreamRegistry.java

Lines changed: 44 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -35,7 +35,18 @@ public StreamSubscription subscribe(IScope scope, String streamName) {
3535
log.debug("Subscribing to stream: {}", streamName);
3636
String key = buildKey(scope, streamName);
3737
log.debug("Subscriber key: {}", key);
38-
StreamState state = streams.computeIfAbsent(key, id -> createState(scope, streamName));
38+
StreamState state = streams.compute(key, (k, existing) -> {
39+
IBroadcastScope bs = scope.getBroadcastScope(streamName);
40+
IClientBroadcastStream currentStream = (bs != null) ? bs.getClientBroadcastStream() : null;
41+
if (existing != null) {
42+
if (currentStream != null && currentStream == existing.stream) {
43+
return existing; // same stream, reuse
44+
}
45+
log.info("Stream changed for {}, detaching old listener", k);
46+
existing.detach();
47+
}
48+
return createState(scope, streamName);
49+
});
3950
if (state == null) {
4051
throw new IllegalStateException("Stream not found: " + streamName);
4152
}
@@ -50,6 +61,10 @@ public StreamSubscription subscribe(IScope scope, String streamName) {
5061
if (initSegment != null) {
5162
queue.offer(initSegment);
5263
}
64+
byte[] keyframe = state.keyframeFragment;
65+
if (keyframe != null) {
66+
queue.offer(keyframe);
67+
}
5368
return new StreamSubscription(key, queue, this);
5469
}
5570

@@ -66,6 +81,18 @@ public void unsubscribe(String key, BlockingQueue<byte[]> queue) {
6681
}
6782
}
6883

84+
void onStreamClosed(String key) {
85+
StreamState state = streams.remove(key);
86+
if (state != null) {
87+
log.info("Stream closed for {}, removed state and notifying {} subscribers", key, state.subscribers.size());
88+
// poison-pill empty array to unblock waiting subscribers
89+
for (BlockingQueue<byte[]> queue : state.subscribers) {
90+
queue.offer(new byte[0]);
91+
}
92+
}
93+
pendingInitSegments.remove(key);
94+
}
95+
6996
void onInitSegment(String key, byte[] initSegment) {
7097
StreamState state = streams.get(key);
7198
if (state == null) {
@@ -79,6 +106,20 @@ void onInitSegment(String key, byte[] initSegment) {
79106
}
80107
}
81108

109+
void onKeyframeFragment(String key, byte[] fragment) {
110+
StreamState state = streams.get(key);
111+
if (state == null) {
112+
return;
113+
}
114+
state.keyframeFragment = fragment;
115+
if (log.isDebugEnabled()) {
116+
log.debug("Dispatching keyframe fragment for {} to {} subscribers ({} bytes)", key, state.subscribers.size(), fragment.length);
117+
}
118+
for (BlockingQueue<byte[]> queue : state.subscribers) {
119+
queue.offer(fragment);
120+
}
121+
}
122+
82123
void onFragment(String key, byte[] fragment) {
83124
StreamState state = streams.get(key);
84125
if (state == null) {
@@ -132,6 +173,8 @@ static class StreamState {
132173

133174
private volatile byte[] initSegment;
134175

176+
private volatile byte[] keyframeFragment;
177+
135178
StreamState(IClientBroadcastStream stream, MediaBunnyStreamListener listener) {
136179
this.stream = stream;
137180
this.listener = listener;

0 commit comments

Comments
 (0)