diff --git a/Client.py b/Client.py index e7c5799..24e6374 100755 --- a/Client.py +++ b/Client.py @@ -5,7 +5,9 @@ tkMessageBox = tkinter.messagebox -from RtpPacket import RtpPacket +from RtpPacket import RtpPacket +from Config import DEFAULT_MEDIA_FILE, RTP_MULTICAST_GROUP, RTP_MULTICAST_PORT, STATE_MULTICAST_GROUP, STATE_MULTICAST_PORT +from StatePacket import PAUSED as STREAM_PAUSED, PLAYING as STREAM_PLAYING, READY as STREAM_READY, STOPPED as STREAM_STOPPED, decode_state_packet CACHE_FILE_NAME = "cache-" CACHE_FILE_EXT = ".jpg" @@ -22,13 +24,13 @@ class Client: TEARDOWN = 3 # Initiation.. - def __init__(self, master, serveraddr, serverport, rtpport, filename): + def __init__(self, master, serveraddr, serverport, filename=DEFAULT_MEDIA_FILE): self.master = master self.master.protocol("WM_DELETE_WINDOW", self.handler) self.createWidgets() self.serverAddr = serveraddr self.serverPort = int(serverport) - self.rtpPort = int(rtpport) + self.rtpPort = RTP_MULTICAST_PORT self.fileName = filename self.rtspSeq = 0 self.sessionId = 0 @@ -36,8 +38,8 @@ def __init__(self, master, serveraddr, serverport, rtpport, filename): self.teardownAcked = 0 self.connectToServer() self.frameNbr = -1 - self.transport = "UDP" - self.quality = "SD" + self.transport = "UDP" + self.quality = "multicast" self.rtpConnection = None self.frameBuffer = [] # NOTE: Handle buffer access UI and network @@ -47,8 +49,15 @@ def __init__(self, master, serveraddr, serverport, rtpport, filename): # NOTE: Set up again request self.isDraining = False - self.pendingSetup = False - self.pendingPause = False + self.pendingSetup = False + self.pendingPause = False + self.pauseInProgress = False + self.setupInProgress = False + self.stateSocket = None + self.stateEvent = threading.Event() + self.lastStateVersion = 0 + self.serverStreamState = STREAM_STOPPED + self.startStateListener() def createWidgets(self): """Build GUI.""" @@ -81,9 +90,9 @@ def createWidgets(self): self.label.grid(row=0, column=0, columnspan=4, sticky=W+E+N+S, padx=5, pady=5) def setupMovie(self): - """Setup button handler - Đã thêm cấu hình xử lý riêng cho trạng thái PAUSE (READY)""" - if self.state == self.INIT: - self.chooseQuality() + """Setup button handler - Đã thêm cấu hình xử lý riêng cho trạng thái PAUSE (READY)""" + if self.state == self.INIT: + self.chooseMulticastStream() elif self.state == self.PLAYING: if tkMessageBox.askokcancel("Reset Video?", "Do you want to stop current playback and setup from the beginning?"): @@ -98,10 +107,10 @@ def setupMovie(self): self.frameBuffer.clear() # Clear buffer immediately self.state = self.INIT self.isBuffering = True - self.isDraining = False - self.pendingSetup = False - self.frameNbr = -1 - self.chooseQuality() + self.isDraining = False + self.pendingSetup = False + self.frameNbr = -1 + self.chooseMulticastStream() def startDrainingPipeline(self): # print("Stop network stream, draining remaining buffered frames to UI.") @@ -110,13 +119,13 @@ def startDrainingPipeline(self): self.isDraining = True self.state = self.PLAYING - def chooseQuality(self): - # Create a beautiful modal window to choose quality - self.quality_window = Toplevel(self.master) - self.quality_window.title("Video Quality") + def chooseMulticastStream(self): + # Multicast media uses UDP only; keep setup explicit without offering TCP. + self.quality_window = Toplevel(self.master) + self.quality_window.title("Multicast Stream") # Center the pop-up on the screen - w, h = 320, 160 + w, h = 360, 150 ws = self.master.winfo_screenwidth() hs = self.master.winfo_screenheight() x = (ws/2) - (w/2) @@ -129,51 +138,27 @@ def chooseQuality(self): title_label = Label( self.quality_window, - text="Choose Video Quality", + text="Start UDP Multicast Stream", fg="#FFFFFF", bg="#1E1E1E", font=("Helvetica", 12, "bold") ) title_label.pack(pady=12) - self.quality_var = StringVar(value="SD") - - frame = Frame(self.quality_window, bg="#1E1E1E") - frame.pack() - - rb_sd = Radiobutton( - frame, - text="SD (720p) - UDP", - variable=self.quality_var, - value="SD", - fg="#E0E0E0", - bg="#1E1E1E", - selectcolor="#2C2C2C", - activeforeground="#FFFFFF", - activebackground="#1E1E1E", - font=("Helvetica", 10) - ) - rb_sd.pack(anchor=W, pady=2) - - rb_hd = Radiobutton( - frame, - text="HD (1080p) - TCP", - variable=self.quality_var, - value="HD", - fg="#E0E0E0", - bg="#1E1E1E", - selectcolor="#2C2C2C", - activeforeground="#FFFFFF", - activebackground="#1E1E1E", - font=("Helvetica", 10) - ) - rb_hd.pack(anchor=W, pady=2) + info_label = Label( + self.quality_window, + text="RTP media is received from the shared UDP multicast group.", + fg="#E0E0E0", + bg="#1E1E1E", + font=("Helvetica", 10) + ) + info_label.pack(pady=8) btn_ok = Button( - self.quality_window, - text="OK", - width=12, - command=self.confirmQuality, + self.quality_window, + text="OK", + width=12, + command=self.confirmMulticastStream, fg="#FFFFFF", bg="#007ACC", activeforeground="#FFFFFF", @@ -188,59 +173,69 @@ def chooseQuality(self): self.quality_window.grab_set() self.master.wait_window(self.quality_window) - def confirmQuality(self): - self.quality = self.quality_var.get() - if self.quality == "HD": - self.transport = "TCP" - else: - self.transport = "UDP" - - self.quality_window.destroy() - self.sendRtspRequest(self.SETUP) + def confirmMulticastStream(self): + self.quality = "multicast" + self.transport = "UDP" + + self.quality_window.destroy() + self.setupInProgress = True + self.sendRtspRequest(self.SETUP) - def exitClient(self): - """Teardown button handler.""" - self.sendRtspRequest(self.TEARDOWN) - self.master.destroy() # Close the gui window + def exitClient(self): + """Teardown button handler.""" + self.closeStateListener() + self.sendRtspRequest(self.TEARDOWN) + self.master.destroy() # Close the gui window try: os.remove(CACHE_FILE_NAME + self.quality.lower() + "-" + str(self.sessionId) + CACHE_FILE_EXT) # Delete the cache image from video except: pass - def pauseMovie(self): - """Pause button handler.""" - if self.state == self.PLAYING: - self.pendingPause = True - self.isDraining = True - self.sendRtspRequest(self.PAUSE) + def pauseMovie(self): + """Pause button handler.""" + if self.state == self.PLAYING and not self.pauseInProgress: + self.pauseInProgress = True + self.pause.config(state=DISABLED) + self.sendRtspRequest(self.PAUSE) - def playMovie(self): - """Play button handler.""" - if self.state == self.READY: - # Kill the old listen thread before starting a new one. - if hasattr(self, 'playEvent'): - self.playEvent.set() # Signal old thread to stop - if hasattr(self, '_rtpThread') and self._rtpThread.is_alive(): self._rtpThread.join(timeout=1.0) # Wait for it to die - - # Now safe to create new event and thread - self.playEvent = threading.Event() - self.playEvent.clear() - - with self.bufferLock: - self.frameBuffer.clear() # Clear buffer when starting new playback - self.isBuffering = True - self.frameNbr = -1 # Reset frame number for new playback - - if self.transport == 'UDP': - self._rtpThread = threading.Thread(target=self.listenRtpWithUDP) - else: - self._rtpThread = threading.Thread(target=self.listenRtpWithTCP) - - self._rtpThread.start() - self.sendRtspRequest(self.PLAY) + def playMovie(self): + """Play button handler.""" + if self.state == self.READY: + self.startPlaybackPipeline(sendRtsp=True) + + def startPlaybackPipeline(self, sendRtsp=False): + """Start local RTP receive/render loops, optionally requesting PLAY first.""" + if not hasattr(self, 'rtpSocket') or not self.rtpSocket: + return + + if self.state == self.PLAYING and hasattr(self, '_rtpThread') and self._rtpThread.is_alive(): + return + + if self.state in (self.READY, self.PLAYING): + # Kill the old listen thread before starting a new one. + if hasattr(self, 'playEvent'): + self.playEvent.set() # Signal old thread to stop + if hasattr(self, '_rtpThread') and self._rtpThread.is_alive(): self._rtpThread.join(timeout=1.0) # Wait for it to die - # Start UI clock play loop to render frames from buffer - self.master.after(120, self.renderClientBufferLoop) + # Now safe to create new event and thread + self.playEvent = threading.Event() + self.playEvent.clear() + + with self.bufferLock: + hasBufferedFrames = len(self.frameBuffer) > 0 + self.isBuffering = not hasBufferedFrames + + self._rtpThread = threading.Thread(target=self.listenRtpWithUDP) + + self._rtpThread.start() + if sendRtsp: + self.sendRtspRequest(self.PLAY) + self.state = self.PLAYING + else: + self.state = self.PLAYING + + # Start UI clock play loop to render frames from buffer + self.master.after(120, self.renderClientBufferLoop) def renderClientBufferLoop(self): """ Render frames from buffer using clock """ @@ -259,20 +254,19 @@ def renderClientBufferLoop(self): frame_bytes = self.frameBuffer.pop(0) self.updateMovie(self.writeFrame(frame_bytes)) else: - if self.isDraining: - # print("Buffer completely drained. Stopping playback.") - self.isDraining = False - self.state = self.INIT - self.frameNbr = -1 - if self.pendingPause: - self.pendingPause = False - self.state = self.READY - self.sendRtspRequest(self.PAUSE) - - elif self.pendingSetup: - self.pendingSetup = False - self.state = self.INIT - self.chooseQuality() + if self.isDraining: + # print("Buffer completely drained. Stopping playback.") + self.isDraining = False + self.state = self.INIT + self.frameNbr = -1 + if self.pendingPause: + self.pendingPause = False + self.state = self.READY + + elif self.pendingSetup: + self.pendingSetup = False + self.state = self.INIT + self.chooseMulticastStream() return else: self.isBuffering = True @@ -324,80 +318,7 @@ def listenRtpWithUDP(self): pass break - def listenRtpWithTCP(self): - """Listen for video frames using TCP frame-by-frame""" - try: - # Accept connection from server with a timeout of 5 seconds - self.rtpSocket.settimeout(5.0) - self.rtpConnection, addr = self.rtpSocket.accept() - self.rtpConnection.settimeout(0.5) - print("RTP/TCP connection established with server") - except Exception as e: - print(f"RTP/TCP accept failed or timed out: {e}") - return - - while not self.playEvent.isSet(): - try: - # Read 5-byte length prefix (ASCII string) - length_bytes = self.recv_all(self.rtpConnection, 5) - if not length_bytes: - break - length = int(length_bytes.decode()) - data = self.recv_all(self.rtpConnection, length) - print("Received Frame over TCP") - print("Size of packet: " + str(len(data))) - - if not data: - break - - with self.bufferLock: - self.frameBuffer.append(data) - if len(self.frameBuffer) > 100: - self.frameBuffer.pop(0) # Discard oldest frame if buffer exceeds size - except: - # Stop listening upon requesting PAUSE or TEARDOWN - if self.playEvent.isSet(): - break - - # Upon receiving ACK for TEARDOWN request, - # close the RTP socket and connection - if self.teardownAcked == 1: - if hasattr(self, 'rtpConnection') and self.rtpConnection: - try: - self.rtpConnection.shutdown(socket.SHUT_RDWR) - self.rtpConnection.close() - except: - pass - try: - self.rtpSocket.shutdown(socket.SHUT_RDWR) - self.rtpSocket.close() - except: - pass - break - - if hasattr(self, 'rtpConnection') and self.rtpConnection: - try: - self.rtpConnection.close() - except: - pass - - def recv_all(self, sock, n): - """Helper to receive exactly n bytes from a TCP socket.""" - data = b'' - while len(data) < n: - try: - packet = sock.recv(n - len(data)) - if not packet: - return None - data += packet - except socket.timeout: - if len(data) > 0: - continue - else: - raise - return data - - def writeFrame(self, data): + def writeFrame(self, data): """Write the received frame to a temp image file. Return the image file.""" cachename = CACHE_FILE_NAME + self.quality.lower() + "-" + str(self.sessionId) + CACHE_FILE_EXT file = open(cachename, "wb") @@ -412,13 +333,129 @@ def updateMovie(self, imageFile): self.label.configure(image = photo, height=288) self.label.image = photo - def connectToServer(self): - """Connect to the Server. Start a new RTSP/TCP session.""" - self.rtspSocket = socket.socket(socket.AF_INET, socket.SOCK_STREAM) - try: - self.rtspSocket.connect((self.serverAddr, self.serverPort)) - except: - tkMessageBox.showwarning('Connection Failed', 'Connection to \'%s\' failed.' %self.serverAddr) + def connectToServer(self): + """Connect to the Server. Start a new RTSP/TCP session.""" + self.rtspSocket = socket.socket(socket.AF_INET, socket.SOCK_STREAM) + try: + self.rtspSocket.connect((self.serverAddr, self.serverPort)) + except: + tkMessageBox.showwarning('Connection Failed', 'Connection to \'%s\' failed.' %self.serverAddr) + + def startStateListener(self): + """Join the state multicast group and log official server state updates.""" + self.stateSocket = socket.socket(socket.AF_INET, socket.SOCK_DGRAM, socket.IPPROTO_UDP) + self.stateSocket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) + if hasattr(socket, 'SO_REUSEPORT'): + try: + self.stateSocket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEPORT, 1) + except OSError: + pass + + try: + self.stateSocket.bind(('', STATE_MULTICAST_PORT)) + mreq = socket.inet_aton(STATE_MULTICAST_GROUP) + socket.inet_aton('0.0.0.0') + self.stateSocket.setsockopt(socket.IPPROTO_IP, socket.IP_ADD_MEMBERSHIP, mreq) + self.stateSocket.settimeout(0.5) + except Exception as e: + print(f"Unable to join state multicast group {STATE_MULTICAST_GROUP}:{STATE_MULTICAST_PORT}: {e}") + try: + self.stateSocket.close() + except: + pass + self.stateSocket = None + return + + self._stateThread = threading.Thread(target=self.listenStateMulticast, daemon=True) + self._stateThread.start() + print(f"Listening for state multicast on {STATE_MULTICAST_GROUP}:{STATE_MULTICAST_PORT}") + + def listenStateMulticast(self): + """Receive official state multicast packets and apply them on the UI thread.""" + while not self.stateEvent.isSet(): + try: + data, address = self.stateSocket.recvfrom(1024) + packet = decode_state_packet(data) + if packet['is_server']: + print(f"State multicast received from {address[0]}:{address[1]}: {packet['state']} v{packet['version']}") + self.master.after(0, self.applyServerState, packet['state'], packet['version']) + else: + print(f"Ignoring non-server state multicast from {address[0]}:{address[1]}") + except socket.timeout: + continue + except OSError: + break + except Exception as e: + print(f"Ignoring invalid state multicast packet: {e}") + + def applyServerState(self, streamState, version): + """Update local playback from the server's official multicast state.""" + if version <= self.lastStateVersion: + return + self.lastStateVersion = version + self.serverStreamState = streamState + + if streamState == STREAM_READY: + if self.state == self.INIT: + if not self.setupInProgress: + print("Server stream is READY; auto-sending SETUP for this client") + self.setupInProgress = True + self.transport = "UDP" + self.sendRtspRequest(self.SETUP) + return + self.state = self.READY + self.isDraining = False + self.pendingPause = False + self.pauseInProgress = False + self.pause.config(state=NORMAL) + if hasattr(self, 'playEvent'): + self.playEvent.set() + + elif streamState == STREAM_PLAYING: + if self.state == self.INIT: + if not self.setupInProgress: + print("Server stream is PLAYING; auto-sending SETUP for this client") + self.setupInProgress = True + self.transport = "UDP" + self.sendRtspRequest(self.SETUP) + return + if self.state == self.READY: + self.startPlaybackPipeline(sendRtsp=False) + + elif streamState == STREAM_PAUSED: + if self.state == self.INIT: + return + self.state = self.READY + self.isDraining = False + self.pendingPause = False + self.pauseInProgress = False + self.isBuffering = False + if hasattr(self, 'playEvent'): + self.playEvent.set() + self.pause.config(state=NORMAL) + + elif streamState == STREAM_STOPPED: + self.state = self.INIT + self.isDraining = False + self.pendingSetup = False + self.pendingPause = False + self.pauseInProgress = False + self.setupInProgress = False + self.isBuffering = True + self.frameNbr = -1 + if hasattr(self, 'playEvent'): + self.playEvent.set() + with self.bufferLock: + self.frameBuffer.clear() + self.pause.config(state=NORMAL) + + def closeStateListener(self): + """Stop the state multicast listener.""" + self.stateEvent.set() + if self.stateSocket: + try: + self.stateSocket.close() + except: + pass def sendRtspRequest(self, requestCode): """Send RTSP request to the server.""" @@ -503,18 +540,29 @@ def parseRtspReply(self, data): # Process only if the session ID is the same if self.sessionId == session: if int(lines[0].split(' ')[1]) == 200: - if self.requestSent == self.SETUP: - # Update RTSP state. - self.state = self.READY - - # Open RTP port. - self.openRtpPort() + if self.requestSent == self.SETUP: + # Update RTSP state. + self.state = self.READY + self.setupInProgress = False + + with self.bufferLock: + self.frameBuffer.clear() + self.frameNbr = -1 + self.isBuffering = True + + # Open RTP port. + self.openRtpPort() + if self.serverStreamState == STREAM_PLAYING: + self.startPlaybackPipeline(sendRtsp=False) elif self.requestSent == self.PLAY: self.state = self.PLAYING - elif self.requestSent == self.PAUSE: - self.state = self.READY - # The play thread exits. A new thread is created on resume. - self.playEvent.set() + elif self.requestSent == self.PAUSE: + self.state = self.READY + self.pauseInProgress = False + self.pause.config(state=NORMAL) + # The play thread exits. A new thread is created on resume. + if hasattr(self, 'playEvent'): + self.playEvent.set() elif self.requestSent == self.TEARDOWN: self.state = self.INIT # Flag the teardownAcked to close the socket. @@ -529,22 +577,21 @@ def openRtpPort(self): except: pass - if self.transport == 'UDP': - self.rtpSocket = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) - self.rtpSocket.settimeout(0.5) - try: - self.rtpSocket.bind(('', self.rtpPort)) - except: - tkMessageBox.showwarning('Unable to Bind', 'Unable to bind PORT=%d' %self.rtpPort) - else: # TCP - self.rtpSocket = socket.socket(socket.AF_INET, socket.SOCK_STREAM) - self.rtpSocket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) - try: - self.rtpSocket.bind(('', self.rtpPort)) - self.rtpSocket.listen(1) - print(f"RTP/TCP socket listening on port {self.rtpPort}") - except: - tkMessageBox.showwarning('Unable to Bind', 'Unable to bind PORT=%d' %self.rtpPort) + self.rtpSocket = socket.socket(socket.AF_INET, socket.SOCK_DGRAM, socket.IPPROTO_UDP) + self.rtpSocket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) + if hasattr(socket, 'SO_REUSEPORT'): + try: + self.rtpSocket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEPORT, 1) + except OSError: + pass + self.rtpSocket.settimeout(0.5) + try: + self.rtpSocket.bind(('', RTP_MULTICAST_PORT)) + mreq = socket.inet_aton(RTP_MULTICAST_GROUP) + socket.inet_aton('0.0.0.0') + self.rtpSocket.setsockopt(socket.IPPROTO_IP, socket.IP_ADD_MEMBERSHIP, mreq) + print(f"Joined RTP multicast group {RTP_MULTICAST_GROUP}:{RTP_MULTICAST_PORT}") + except Exception as e: + tkMessageBox.showwarning('Unable to Join Multicast', f'Unable to join RTP multicast group {RTP_MULTICAST_GROUP}:{RTP_MULTICAST_PORT}: {e}') def handler(self): """Handler on explicitly closing the GUI window.""" diff --git a/ClientLauncher.py b/ClientLauncher.py index e7080e5..3fa3469 100755 --- a/ClientLauncher.py +++ b/ClientLauncher.py @@ -1,20 +1,21 @@ -import sys -from tkinter import Tk -from Client import Client - -if __name__ == "__main__": - try: - serverAddr = sys.argv[1] - serverPort = sys.argv[2] - rtpPort = sys.argv[3] - fileName = sys.argv[4] - except: - print("[Usage: ClientLauncher.py Server_name Server_port RTP_port Video_file]\n") +import sys +from tkinter import Tk +from Client import Client +from Config import DEFAULT_MEDIA_FILE + +if __name__ == "__main__": + try: + serverAddr = sys.argv[1] + serverPort = sys.argv[2] + fileName = sys.argv[3] if len(sys.argv) > 3 else DEFAULT_MEDIA_FILE + except: + print("[Usage: ClientLauncher.py Server_name Server_port [Server_media_file]]\n") + sys.exit(1) root = Tk() # Create a new client - app = Client(root, serverAddr, serverPort, rtpPort, fileName) + app = Client(root, serverAddr, serverPort, fileName) app.master.title("RTPClient") root.mainloop() diff --git a/Config.py b/Config.py new file mode 100644 index 0000000..391c08f --- /dev/null +++ b/Config.py @@ -0,0 +1,9 @@ +RTP_MULTICAST_GROUP = "239.10.10.1" +RTP_MULTICAST_PORT = 5004 + +STATE_MULTICAST_GROUP = "239.10.10.2" +STATE_MULTICAST_PORT = 7000 + +MULTICAST_TTL = 1 + +DEFAULT_MEDIA_FILE = "movie.Mjpeg" diff --git a/MulticastStreamManager.py b/MulticastStreamManager.py new file mode 100644 index 0000000..ff48924 --- /dev/null +++ b/MulticastStreamManager.py @@ -0,0 +1,211 @@ +import socket +import threading + +from Config import RTP_MULTICAST_GROUP, RTP_MULTICAST_PORT, STATE_MULTICAST_GROUP, STATE_MULTICAST_PORT, MULTICAST_TTL +from RtpPacket import RtpPacket +from StatePacket import READY, PLAYING, PAUSED, STOPPED, encode_state_packet +from VideoStream import VideoStream + + +class MulticastStreamManager: + MAX_RTP_PAYLOAD_SIZE = 1400 + + def __init__(self): + self.lock = threading.Lock() + self.videoStream = None + self.filename = None + self.rtpSocket = None + self.worker = None + self.stopEvent = threading.Event() + self.state = STOPPED + self.stateVersion = 0 + self.rtpSeq = 0 + self.clientSessions = set() + + def register_client(self, session): + """Track a client that joined the shared multicast stream.""" + with self.lock: + self.clientSessions.add(session) + clientCount = len(self.clientSessions) + + print(f"Multicast client registered: {session} ({clientCount} active)") + + def unregister_client(self, session): + """Remove one client and stop the stream only when no clients remain.""" + should_stop = False + with self.lock: + if session not in self.clientSessions: + return + + self.clientSessions.remove(session) + clientCount = len(self.clientSessions) + should_stop = clientCount == 0 + + if should_stop: + self._stop_locked(close_socket=True) + self.state = STOPPED + + print(f"Multicast client unregistered: {session} ({clientCount} active)") + + if should_stop: + self._send_state_multicast(STOPPED) + + def setup(self, filename): + """Prepare one shared video stream for multicast delivery.""" + with self.lock: + if self.videoStream is not None and self.filename == filename and self.state != STOPPED: + currentState = self.state + should_announce_ready = False + else: + self._stop_locked(close_socket=True) + self.videoStream = VideoStream(filename) + self.filename = filename + self.rtpSeq = 0 + self.state = READY + currentState = READY + should_announce_ready = True + + if should_announce_ready: + self._send_state_multicast(READY) + elif currentState != STOPPED: + self._send_state_multicast(currentState) + + def play(self): + """Start or resume the shared multicast sender.""" + with self.lock: + if self.videoStream is None: + raise IOError + + self.stopEvent.clear() + self._open_rtp_socket_locked() + should_start = self.worker is None or not self.worker.is_alive() + self.state = PLAYING + if should_start: + self.worker = threading.Thread(target=self._send_rtp_loop, daemon=True) + self.worker.start() + + self._send_state_multicast(PLAYING) + + def pause(self): + """Pause the shared multicast sender without resetting the video.""" + with self.lock: + if self.state == PAUSED: + should_announce_paused = True + worker = None + elif self.state != PLAYING: + return + else: + should_announce_paused = True + self.state = PAUSED + self.stopEvent.set() + worker = self.worker + + if worker: + worker.join(timeout=0.5) + + with self.lock: + if self.worker is worker: + self.worker = None + + if should_announce_paused: + self._send_state_multicast(PAUSED) + + def stop(self): + """Stop the shared stream and release multicast resources.""" + with self.lock: + self.clientSessions.clear() + self._stop_locked(close_socket=True) + self.state = STOPPED + + self._send_state_multicast(STOPPED) + + def _stop_locked(self, close_socket): + self.stopEvent.set() + worker = self.worker + self.worker = None + + if worker and worker.is_alive() and worker is not threading.current_thread(): + self.lock.release() + try: + worker.join(timeout=0.5) + finally: + self.lock.acquire() + + if close_socket and self.rtpSocket: + try: + self.rtpSocket.close() + except: + pass + self.rtpSocket = None + + self.videoStream = None + self.filename = None + + def _open_rtp_socket_locked(self): + if self.rtpSocket: + return + + self.rtpSocket = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) + self.rtpSocket.setsockopt(socket.IPPROTO_IP, socket.IP_MULTICAST_TTL, MULTICAST_TTL) + print(f"RTP multicast sender ready: {RTP_MULTICAST_GROUP}:{RTP_MULTICAST_PORT} ttl={MULTICAST_TTL}") + + def _send_rtp_loop(self): + while not self.stopEvent.wait(0.05): + with self.lock: + videoStream = self.videoStream + rtpSocket = self.rtpSocket + + if videoStream is None or rtpSocket is None: + break + + data = videoStream.nextFrame() + if not data: + with self.lock: + self.state = STOPPED + self.worker = None + self.stopEvent.set() + self._send_state_multicast(STOPPED) + break + + self._send_frame(data, rtpSocket) + + def _send_frame(self, data, rtpSocket): + frameSize = len(data) + bytesSent = 0 + + while bytesSent < frameSize and not self.stopEvent.isSet(): + chunkSize = min(self.MAX_RTP_PAYLOAD_SIZE, frameSize - bytesSent) + chunkData = data[bytesSent:bytesSent + chunkSize] + markerBit = (bytesSent + chunkSize) == frameSize + + with self.lock: + packet = self._make_rtp(chunkData, self.rtpSeq, markerBit) + self.rtpSeq += 1 + + try: + rtpSocket.sendto(packet, (RTP_MULTICAST_GROUP, RTP_MULTICAST_PORT)) + bytesSent += chunkSize + except Exception as e: + print(f"Sending multicast RTP error: {e}") + break + + def _make_rtp(self, payload, seqnum, markerBit): + rtpPacket = RtpPacket() + rtpPacket.encode(2, 0, 0, 0, seqnum, markerBit, 26, 0, payload) + return rtpPacket.getPacket() + + def _send_state_multicast(self, state): + with self.lock: + self.stateVersion += 1 + version = self.stateVersion + + packet = encode_state_packet(state, version) + stateSocket = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) + try: + stateSocket.setsockopt(socket.IPPROTO_IP, socket.IP_MULTICAST_TTL, MULTICAST_TTL) + stateSocket.sendto(packet, (STATE_MULTICAST_GROUP, STATE_MULTICAST_PORT)) + print(f"State multicast sent: {state} v{version}") + except Exception as e: + print(f"Sending state multicast error: {e}") + finally: + stateSocket.close() diff --git a/README.md b/README.md index ba5c248..01cbd60 100644 --- a/README.md +++ b/README.md @@ -13,15 +13,13 @@ Designed and implemented as a professional university project for the **Computer ## 🌟 Key Features * **⚡ Non-Blocking Multiplexed Server:** Built with the Python `selectors` module to implement I/O multiplexing. The server manages multiple client RTSP signaling channels concurrently on a single thread without blocking, ensuring optimal resource utilization. -* **🔄 Hybrid Dual-Transport Architecture (UDP & TCP):** - * **SD (Standard Definition) Streaming:** Streamed over **UDP** for low-overhead, real-time packet delivery. - * **HD (High Definition) Streaming:** Streamed over **TCP** with custom length-prefixed framing (5-byte headers) to ensure 100% reliable packet reassembly and eliminate packet-loss artifacts in high-bitrate media. +* **📡 UDP Multicast Media Architecture:** Video frames are packetized as RTP and sent once to a shared UDP multicast group. RTSP remains a per-client TCP control channel for SETUP, PLAY, PAUSE, and TEARDOWN. * **📶 Client-Side Jitter Buffer (Double-Buffered State Machine):** * Implements a thread-safe, lock-protected queue that acts as a client caching layer. * Utilizes a **Low-Water / High-Water Mark state machine** (`minBufferSize = 15` frames) to dynamically pause/resume visual rendering, absorbing network jitter and packet arrival fluctuations. -* **🎬 Seamless Dynamic Quality Switching:** Clients can switch stream resolutions on-the-fly. The client triggers a pipeline draining phase that safely plays out cached buffer frames before establishing a new RTP session, preserving session state and preventing UI flickering. +* **🎬 Shared Live Stream Control:** Clients join the same multicast media group and follow the server's official stream state announcements. * **📦 UDP Fragmentation & RTP Reassembly:** Features robust fragmentation for UDP payloads exceeding MTU limits (1400 bytes). Employs RTP Marker bits to detect frame boundaries, facilitating perfect client-side JPEG reassembly. -* **🎨 Premium Responsive Tkinter GUI:** A dark-themed, sleek user interface displaying live video playback, buffer state visualizations, and intuitive playback controls (Setup, Play, Pause, Quality Selection, and Teardown). +* **🎨 Premium Responsive Tkinter GUI:** A dark-themed, sleek user interface displaying live video playback, buffer state visualizations, and intuitive playback controls (Setup, Play, Pause, and Teardown). --- @@ -36,9 +34,9 @@ sequenceDiagram participant Client as RTP/RTSP Client participant Server as Multiplexed RTSP Server - User->>Client: Clicks "Setup" & Selects Quality (SD/HD) - Client->>Server: RTSP SETUP (Transport: UDP/TCP) - Server-->>Client: RTSP 200 OK (Session ID, Server Ports) + User->>Client: Clicks "Setup" + Client->>Server: RTSP SETUP (Transport: UDP) + Server-->>Client: RTSP 200 OK (Session ID) User->>Client: Clicks "Play" Client->>Server: RTSP PLAY @@ -48,21 +46,17 @@ sequenceDiagram rect rgb(20, 20, 30) Note over Server, Client: Media Streaming Loop Server->>Server: Read MJPEG Frame - alt Quality == SD (UDP) - Server->>Server: Fragment Frame (MTU = 1400 bytes) - Server->>Client: Send RTP Packets over UDP - else Quality == HD (TCP) - Server->>Client: Send RTP Packet over TCP (5-byte length-prefix) - end + Server->>Server: Fragment Frame (MTU = 1400 bytes) + Server->>Client: Send RTP Packets to UDP multicast group Client->>Client: Reassemble & Queue in Jitter Buffer Client->>Client: Render from Buffer (Min Buffer: 15 frames) end - User->>Client: Clicks "Setup" (Switch Quality) + User->>Client: Clicks "Setup" (Reset Stream) Note over Client: Enters isDraining Pipeline Mode Note over Client: Plays out cached frames to avoid abrupt cut - Client->>Server: RTSP SETUP (New Quality Selected) - Server->>Server: Close Old RTP Channel & Re-initialize Video Stream + Client->>Server: RTSP SETUP + Server->>Server: Re-initialize Shared Video Stream Server-->>Client: RTSP 200 OK Client->>Server: RTSP PLAY Server-->>Client: RTSP 200 OK @@ -80,7 +74,7 @@ The client state engine integrates network operations and buffer rendering trans | INIT | +--------------------------------+ | - | SETUP (Selects SD/HD) + | SETUP (Join Multicast) v +--------------------------------+ | READY | @@ -135,12 +129,62 @@ The client state engine integrates network operations and buffer rendering trans ``` 2. **Launch the Video Client:** ```bash - python ClientLauncher.py + python ClientLauncher.py [server_media_file] # Example (running client locally connecting to server): - python ClientLauncher.py 127.0.0.1 8554 25000 movie.Mjpeg + python ClientLauncher.py 127.0.0.1 8554 ``` + The optional `server_media_file` is the media resource requested from the server, not a file read by the client. If omitted, the client requests `movie.Mjpeg` from the server. + +### 📡 Multicast Demo + +The RTSP server still listens on the port you pass to `Server.py`, but the media and state channels use fixed multicast groups from `Config.py`: + +```text +RTP media multicast: 239.10.10.1:5004 +State multicast: 239.10.10.2:7000 +Multicast TTL: 1 +``` + +Run one server: + +```bash +python Server.py 8554 +``` + +Then launch two or more clients in separate terminals: + +```bash +python ClientLauncher.py 127.0.0.1 8554 +python ClientLauncher.py 127.0.0.1 8554 +``` + +The clients do not read `movie.Mjpeg` locally. They request the server-side media resource over RTSP, then receive RTP frames from the shared multicast port `5004`. + +Expected demo flow: + +```text +1. Start the server. +2. Start Client A and click Setup. +3. Start Client B and click Setup. +4. Click Play from either client. +5. Both clients should follow the same server stream state and receive the same RTP multicast media. +6. Click Pause from either client to pause the shared stream for all clients. +7. Close one client; the other client should remain connected unless it is the last active client. +``` + +Useful console logs during the demo: + +```text +Multicast client registered: ( active) +RTP multicast sender ready: 239.10.10.1:5004 ttl=1 +State multicast sent: PLAYING v +Joined RTP multicast group 239.10.10.1:5004 +Listening for state multicast on 239.10.10.2:7000 +State multicast received from :: PLAYING v +``` + --- ## 📁 Repository Structure @@ -149,9 +193,12 @@ The client state engine integrates network operations and buffer rendering trans . ├── Client.py # RTSP Client core logic & Jitter Buffer controls ├── ClientLauncher.py # GUI and Client initialization entry point +├── Config.py # Shared multicast group, port, and TTL settings +├── MulticastStreamManager.py # Shared RTP multicast sender and server state owner ├── RtpPacket.py # RTP Packet encapsulation & header parsing (12-byte headers) ├── Server.py # Non-blocking I/O multiplexing RTSP Server entry point ├── ServerWorker.py # RTSP State Machine, TCP/UDP Media Streamer, Frame Fragmenter +├── StatePacket.py # UDP multicast state announcement encoder/decoder ├── VideoStream.py # MJPEG parser and frame extraction engine ├── movie.Mjpeg # Sample MJPEG video file ├── doc/ diff --git a/Server.py b/Server.py index 54bf57b..0b58252 100755 --- a/Server.py +++ b/Server.py @@ -1,12 +1,14 @@ import sys, socket import selectors -from ServerWorker import ServerWorker +from ServerWorker import ServerWorker +from MulticastStreamManager import MulticastStreamManager class Server: - def __init__(self): - self.serverWorkers = {} - self.selectors = selectors.DefaultSelector() + def __init__(self): + self.serverWorkers = {} + self.selectors = selectors.DefaultSelector() + self.streamManager = MulticastStreamManager() def accept_client(self, rtspSocket): try: @@ -14,7 +16,7 @@ def accept_client(self, rtspSocket): clientSocket.setblocking(False) clientInfo = {'rtspSocket': (clientSocket, clientAddress)} - worker = ServerWorker(clientInfo) + worker = ServerWorker(clientInfo, self.streamManager) self.serverWorkers[clientSocket] = worker self.selectors.register(clientSocket, selectors.EVENT_READ, self.handle_client_request) except Exception as e: @@ -26,11 +28,15 @@ def handle_client_request(self, clientSocket): try: worker.recvRtspRequest() - except Exception as e: - print(f"Cleaning up disconnected client: {e}") - try: - self.selectors.unregister(clientSocket) - except: + except Exception as e: + print(f"Cleaning up disconnected client: {e}") + try: + worker.close() + except: + pass + try: + self.selectors.unregister(clientSocket) + except: pass clientSocket.close() if clientSocket in self.serverWorkers: diff --git a/ServerWorker.py b/ServerWorker.py index b866572..a3ffc34 100755 --- a/ServerWorker.py +++ b/ServerWorker.py @@ -4,7 +4,6 @@ from PIL import Image from VideoStream import VideoStream -from RtpPacket import RtpPacket class ServerWorker: SETUP = 'SETUP' @@ -22,9 +21,9 @@ class ServerWorker: CON_ERR_500 = 2 clientInfo = {} - rtpSeq = 0 - def __init__(self, clientInfo): + def __init__(self, clientInfo, streamManager): self.clientInfo = clientInfo + self.streamManager = streamManager def run(self): threading.Thread(target=self.recvRtspRequest).start() @@ -60,90 +59,141 @@ def processRtspRequest(self, data): # Process SETUP request if requestType == self.SETUP: print("processing SETUP\n") - if 'event' in self.clientInfo and self.clientInfo['event']: - try: self.clientInfo['event'].set() - except: pass - - if 'worker' in self.clientInfo and self.clientInfo['worker']: - try: self.clientInfo['worker'].join(timeout=0.2) - except: pass - - if 'rtpSocket' in self.clientInfo and self.clientInfo['rtpSocket']: - try: self.clientInfo['rtpSocket'].close() - except: pass - # print("Cleaned up previous session resources") - - try: - self.clientInfo['videoStream'] = VideoStream(filename) - self.state = self.READY - except IOError: - self.replyRtsp(self.FILE_NOT_FOUND_404, seq[1]) - self.rtpSeq = 0 - # Generate a randomized RTSP session ID - if self.clientInfo.get('session') is None: - self.clientInfo['session'] = randint(100000, 999999) - - # Send RTSP reply - self.replyRtsp(self.OK_200, seq[1]) - - # Parse Transport header + previousTransport = self.clientInfo.get('transport') transport_line = request[2] self.clientInfo['transport'] = 'UDP' if 'TCP' in transport_line: self.clientInfo['transport'] = 'TCP' - - # Extract client_port robustly + parts = transport_line.split(';') for part in parts: if 'client_port' in part: self.clientInfo['rtpPort'] = part.split('=')[1].strip() - # print("Successfully parsed client_port: " + self.clientInfo['rtpPort']) break - - # Process PLAY request - elif requestType == self.PLAY: - if self.state == self.READY: - print("processing PLAY\n") - self.state = self.PLAYING - - # Create a new socket for RTP + + if previousTransport == 'UDP' and self.clientInfo.get('transport') == 'TCP': + self._unregisterMulticastClient() + + self._cleanupClientRtpResources() + + try: if self.clientInfo.get('transport', 'UDP') == 'TCP': - self.clientInfo["rtpSocket"] = socket.socket(socket.AF_INET, socket.SOCK_STREAM) - address = self.clientInfo['rtspSocket'][1][0] - port = int(self.clientInfo['rtpPort']) - print(f"Connecting RTP/TCP socket to {address}:{port}") - self.clientInfo["rtpSocket"].connect((address, port)) + self.clientInfo['videoStream'] = VideoStream(filename) + self.clientInfo['rtpSeq'] = 0 else: - self.clientInfo["rtpSocket"] = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) - + self.streamManager.setup(filename) + self.state = self.READY + setupSucceeded = True + except IOError: + setupSucceeded = False + self.replyRtsp(self.FILE_NOT_FOUND_404, seq[1]) + + # Generate a randomized RTSP session ID + if self.clientInfo.get('session') is None: + self.clientInfo['session'] = randint(100000, 999999) + + # Send RTSP reply + if setupSucceeded: + if self.clientInfo.get('transport', 'UDP') != 'TCP': + self._registerMulticastClient() self.replyRtsp(self.OK_200, seq[1]) - - # Create a new thread and start sending RTP packets - self.clientInfo['event'] = threading.Event() - self.clientInfo['worker']= threading.Thread(target=self.sendRtp) - self.clientInfo['worker'].start() - + + # Process PLAY request + elif requestType == self.PLAY: + if self.state == self.READY or self.clientInfo.get('multicastRegistered'): + print("processing PLAY\n") + try: + if self.clientInfo.get('transport', 'UDP') == 'TCP': + if self.state != self.READY: + return + self.clientInfo["rtpSocket"] = socket.socket(socket.AF_INET, socket.SOCK_STREAM) + address = self.clientInfo['rtspSocket'][1][0] + port = int(self.clientInfo['rtpPort']) + print(f"Connecting RTP/TCP socket to {address}:{port}") + self.clientInfo["rtpSocket"].connect((address, port)) + self.clientInfo['event'] = threading.Event() + self.clientInfo['worker']= threading.Thread(target=self.sendRtpWithTCP) + self.clientInfo['worker'].start() + else: + self.streamManager.play() + + self.state = self.PLAYING + self.replyRtsp(self.OK_200, seq[1]) + except Exception: + self.replyRtsp(self.CON_ERR_500, seq[1]) + # Process PAUSE request elif requestType == self.PAUSE: - if self.state == self.PLAYING: + if self.state == self.PLAYING or self.clientInfo.get('multicastRegistered'): print("processing PAUSE\n") self.state = self.READY - - self.clientInfo['event'].set() - + + if self.clientInfo.get('transport', 'UDP') == 'TCP': + if 'event' not in self.clientInfo or not self.clientInfo['event']: + return + self.clientInfo['event'].set() + else: + self.streamManager.pause() + self.replyRtsp(self.OK_200, seq[1]) - + # Process TEARDOWN request elif requestType == self.TEARDOWN: print("processing TEARDOWN\n") - self.clientInfo['event'].set() - + if self.clientInfo.get('transport', 'UDP') == 'TCP': + if 'event' in self.clientInfo and self.clientInfo['event']: + self.clientInfo['event'].set() + else: + self._unregisterMulticastClient() + self.replyRtsp(self.OK_200, seq[1]) - + self.state = self.INIT + # Close the RTP socket - self.clientInfo['rtpSocket'].close() + if 'rtpSocket' in self.clientInfo and self.clientInfo['rtpSocket']: + self.clientInfo['rtpSocket'].close() + + def close(self): + """Release resources owned by this RTSP client session.""" + self._cleanupClientRtpResources() + self._unregisterMulticastClient() + + def _registerMulticastClient(self): + if self.clientInfo.get('multicastRegistered'): + return + + session = self.clientInfo.get('session') + if session is None: + return + + self.streamManager.register_client(session) + self.clientInfo['multicastRegistered'] = True + + def _unregisterMulticastClient(self): + if not self.clientInfo.get('multicastRegistered'): + return + + session = self.clientInfo.get('session') + if session is not None: + self.streamManager.unregister_client(session) + + self.clientInfo['multicastRegistered'] = False + + def _cleanupClientRtpResources(self): + if 'event' in self.clientInfo and self.clientInfo['event']: + try: self.clientInfo['event'].set() + except: pass + + if 'worker' in self.clientInfo and self.clientInfo['worker']: + try: self.clientInfo['worker'].join(timeout=0.2) + except: pass + + if 'rtpSocket' in self.clientInfo and self.clientInfo['rtpSocket']: + try: self.clientInfo['rtpSocket'].close() + except: pass + self.clientInfo['rtpSocket'] = None # NOTE: Implement fragmentation for frames exceeding the MTU def sendRtpWithTCP(self): @@ -167,60 +217,6 @@ def sendRtpWithTCP(self): print(f"Sending video frame with TCP error: {e}") break - def sendRtpWithUDP(self): - """Send RTP packets using UDP""" - MAX_RTP_PAYLOAD_SIZE = 1400 # MTU - RTP header size - while True: - self.clientInfo['event'].wait(0.05) - - # Stop sending if request is PAUSE or TEARDOWN - if self.clientInfo['event'].isSet(): - break - - data = self.clientInfo['videoStream'].nextFrame() - if data: - frameSize = len(data) - bytesSent = 0 - while bytesSent < frameSize: - chunkSize = min(MAX_RTP_PAYLOAD_SIZE, frameSize - bytesSent) - chunkData = data[bytesSent:bytesSent + chunkSize] - markerBit = (bytesSent + chunkSize) == frameSize - try: - address = self.clientInfo['rtspSocket'][1][0] - port = int(self.clientInfo['rtpPort']) - packet = self.makeRtp(chunkData, self.rtpSeq, markerBit) - self.clientInfo['rtpSocket'].sendto(packet, (address, port)) - - self.rtpSeq += 1 - bytesSent += chunkSize - except Exception as e: - print(f"Sending RTP with UDP error: {e}") - break - - def sendRtp(self): - """Send RTP packets.""" - if self.clientInfo.get('transport', 'UDP') == 'TCP': - return self.sendRtpWithTCP() - else: - return self.sendRtpWithUDP() - - def makeRtp(self, payload, frameNbr, markerBit): - """RTP-packetize the video data.""" - version = 2 - padding = 0 - extension = 0 - cc = 0 - marker = markerBit - pt = 26 # MJPEG type - seqnum = frameNbr - ssrc = 0 - - rtpPacket = RtpPacket() - - rtpPacket.encode(version, padding, extension, cc, seqnum, marker, pt, ssrc, payload) - - return rtpPacket.getPacket() - def replyRtsp(self, code, seq): """Send RTSP reply to the client.""" if code == self.OK_200: diff --git a/StatePacket.py b/StatePacket.py new file mode 100644 index 0000000..0318f1f --- /dev/null +++ b/StatePacket.py @@ -0,0 +1,45 @@ +import json + + +MAGIC = "VS_STATE_V1" +STATE_UPDATE = "STATE_UPDATE" + +READY = "READY" +PLAYING = "PLAYING" +PAUSED = "PAUSED" +STOPPED = "STOPPED" + +VALID_STATES = {READY, PLAYING, PAUSED, STOPPED} + + +def encode_state_packet(state, version, is_server=True): + """Encode a stream state update for UDP multicast transport.""" + if state not in VALID_STATES: + raise ValueError("Invalid stream state: " + str(state)) + + packet = { + "magic": MAGIC, + "is_server": bool(is_server), + "type": STATE_UPDATE, + "state": state, + "version": int(version), + } + return json.dumps(packet, separators=(",", ":")).encode("utf-8") + + +def decode_state_packet(data): + """Decode and validate a stream state update packet.""" + packet = json.loads(data.decode("utf-8")) + + if packet.get("magic") != MAGIC: + raise ValueError("Invalid state packet magic") + if packet.get("type") != STATE_UPDATE: + raise ValueError("Invalid state packet type") + if packet.get("state") not in VALID_STATES: + raise ValueError("Invalid stream state: " + str(packet.get("state"))) + + return { + "is_server": bool(packet.get("is_server")), + "state": packet["state"], + "version": int(packet.get("version", 0)), + } diff --git a/doc/multicast_design.md b/doc/multicast_design.md new file mode 100644 index 0000000..2e4049d --- /dev/null +++ b/doc/multicast_design.md @@ -0,0 +1,134 @@ +# Multicast Design Checkpoint + +This document records the current behavior and the target multicast design before code changes are made. It is intended as a safe checkpoint for review and rollback. + +## Current Behavior + +The current project uses RTSP over TCP for control and sends media per client. + +Control path: + +```text +Client -> Server +RTSP over TCP +SETUP / PLAY / PAUSE / TEARDOWN +``` + +Media path: + +```text +ServerWorker -> one client +RTP over UDP, or custom frame delivery over TCP +``` + +For UDP streaming, the server sends each RTP packet directly to the requesting client's IP address and RTP port. Each client has its own `ServerWorker`, its own `VideoStream`, and its own media sending loop. + +Current UDP media model: + +```text +Client A SETUP/PLAY -> ServerWorker A -> RTP packets to Client A +Client B SETUP/PLAY -> ServerWorker B -> RTP packets to Client B +``` + +This is multi-client unicast, not multicast. Multiple clients can connect, but the server still sends separate media streams to each client. + +## Target Multicast Behavior + +The target design keeps RTSP as the unicast control protocol and changes RTP media delivery to UDP multicast. + +Target control path: + +```text +Client -> Server +RTSP over TCP +SETUP / PLAY / PAUSE / TEARDOWN +``` + +Target media path: + +```text +Server -> RTP multicast group +UDP multicast +All joined clients receive the same RTP stream +``` + +Target media model: + +```text +Server -> 239.10.10.1:5004 + -> Client A, if joined + -> Client B, if joined +``` + +The server sends each RTP packet once to the multicast group. Clients receive video by joining the multicast group and listening on the multicast RTP port. + +## Target State Model + +The server owns the official shared stream state. + +Proposed states: + +```text +STOPPED +READY +PLAYING +PAUSED +``` + +The target design may use a separate state announcement multicast channel so all clients can observe the official state. + +State channel: + +```text +Server -> 239.10.10.2:7000 +UDP multicast state announcements +``` + +Example flow: + +```text +Client A sends RTSP PAUSE to server +Server changes shared stream state to PAUSED +Server multicasts STATE_PAUSED +All clients receive STATE_PAUSED and update local UI/playback state +``` + +Clients do not directly control each other. Clients request control from the server, and the server announces the official state. + +## Assumptions + +- RTSP remains TCP unicast per client. +- RTP media is multicast over UDP. +- The multicast stream is shared by all clients. +- PLAY and PAUSE are global stream controls in the multicast design. +- TEARDOWN should disconnect the requesting client; it should not necessarily stop the shared stream for all clients. +- Late-joining clients start from the current live stream position, not from the beginning of the video. +- A simple LAN/demo environment is assumed. +- Security and authentication for state packets are out of scope. + +## Out Of Scope For Initial Multicast Work + +- RTCP implementation. +- H.264 encoding changes. +- Authentication or anti-spoofing for multicast state messages. +- WAN multicast routing support. +- Production-grade stream discovery. + +## Implementation Direction + +The main architectural change is moving from per-client media streams to one shared multicast stream. + +Current ownership: + +```text +ServerWorker owns VideoStream and RTP sending loop +``` + +Target ownership: + +```text +Shared multicast stream manager owns VideoStream and RTP sending loop +ServerWorker handles RTSP control requests and delegates shared stream control +``` + +This keeps the existing RTSP client/server structure while replacing the media delivery path with multicast.