Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
115 changes: 95 additions & 20 deletions Client.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,16 +4,30 @@
from PIL import Image, ImageTk
import io
import sys
import time
from CustomPacket import CustomPacket

MULTICAST_GROUP = '239.1.1.1'
MULTICAST_PORT = 5004
SOCKET_BUFFER_SIZE = 4 * 1024 * 1024
STALE_FRAME_SECONDS = 5
FRAGMENT_TIMEOUT_SECONDS = 0.02
MISSING_FRAGMENT_PREVIEW = 20

def get_default_interface_ip():
sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
try:
sock.connect((MULTICAST_GROUP, MULTICAST_PORT))
return sock.getsockname()[0]
finally:
sock.close()

class Client:
def __init__(self, master):
def __init__(self, master, interface_ip=None):
self.master = master
self.master.title("Multicast Video Client")
self.master.protocol("WM_DELETE_WINDOW", self.handler)
self.interface_ip = interface_ip or get_default_interface_ip()

# UI Elements
self.label = tk.Label(self.master, text="Waiting for multicast stream...", bg="black", fg="white", width=60, height=20)
Expand All @@ -26,7 +40,11 @@ def __init__(self, master):
self.running = True
self.expected_frame = 0
self.received_frames = 0
self.completed_frames = 0
self.expired_frames = 0
self.lost_frames = 0
self.fragment_buffers = {}
self.last_stats_update = 0

self.setup_socket()

Expand All @@ -37,6 +55,7 @@ def __init__(self, master):
def setup_socket(self):
# Create UDP socket
self.sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM, socket.IPPROTO_UDP)
self.sock.setsockopt(socket.SOL_SOCKET, socket.SO_RCVBUF, SOCKET_BUFFER_SIZE)

# Allow multiple clients on the same machine to bind to the same port
self.sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
Expand All @@ -50,36 +69,84 @@ def setup_socket(self):
self.sock.bind(('', MULTICAST_PORT))

# Join the multicast group
mreq = socket.inet_aton(MULTICAST_GROUP) + socket.inet_aton('0.0.0.0')
mreq = socket.inet_aton(MULTICAST_GROUP) + socket.inet_aton(self.interface_ip)
self.sock.setsockopt(socket.IPPROTO_IP, socket.IP_ADD_MEMBERSHIP, mreq)
print(f"Joined multicast group {MULTICAST_GROUP}:{MULTICAST_PORT}")
print(f"Joined multicast group {MULTICAST_GROUP}:{MULTICAST_PORT} via {self.interface_ip}")

def receive_loop(self):
while self.running:
try:
# Receive multicast packets (buffer size 65536 is large enough for our max frame size ~14KB)
data, _ = self.sock.recvfrom(65536)
data, _ = self.sock.recvfrom(2048)
if not data:
continue

# Decode received packets using our CustomPacket
frame_num, payload = CustomPacket.decode(data)
frame_num, fragment_index, fragment_count, payload = CustomPacket.decode(data)

if frame_num is not None:
# Loss detection
if self.expected_frame > 0 and frame_num > self.expected_frame:
self.lost_frames += (frame_num - self.expected_frame)

self.expected_frame = frame_num + 1
self.received_frames += 1

# Schedule display update on the main GUI thread
self.master.after(0, self.update_display, payload)
self.master.after(0, self.update_stats)

buffer = self.fragment_buffers.setdefault(frame_num, {
'count': fragment_count,
'fragments': {},
'updated_at': time.monotonic(),
})

if buffer['count'] == fragment_count:
buffer['fragments'][fragment_index] = payload
buffer['updated_at'] = time.monotonic()

if len(buffer['fragments']) == buffer['count']:
self.completed_frames += 1

# Loss detection is frame-based; display only complete frames.
if self.expected_frame > 0 and frame_num > self.expected_frame:
self.lost_frames += (frame_num - self.expected_frame)

self.expected_frame = frame_num + 1
frame = b''.join(buffer['fragments'][index] for index in range(buffer['count']))
del self.fragment_buffers[frame_num]

# Drop stale incomplete frames to avoid unbounded growth.
self.cleanup_stale_frames()

# Schedule display update on the main GUI thread
self.master.after(0, self.update_display, frame)
self.master.after(0, self.update_stats)

self.cleanup_stale_frames()
now = time.monotonic()
if now - self.last_stats_update >= 0.1:
self.last_stats_update = now
self.master.after(0, self.update_stats)
except Exception as e:
if self.running:
print(f"Error receiving packet: {e}")

def cleanup_stale_frames(self):
now = time.monotonic()
for stale_frame, stale_buffer in list(self.fragment_buffers.items()):
if stale_frame < self.expected_frame:
del self.fragment_buffers[stale_frame]
continue

timeout = max(STALE_FRAME_SECONDS, stale_buffer['count'] * FRAGMENT_TIMEOUT_SECONDS)
if now - stale_buffer['updated_at'] > timeout:
self.log_expired_frame(stale_frame, stale_buffer)
self.expired_frames += 1
del self.fragment_buffers[stale_frame]

def log_expired_frame(self, frame_num, buffer):
received = len(buffer['fragments'])
expected = buffer['count']
missing = [index for index in range(expected) if index not in buffer['fragments']]
preview = missing[:MISSING_FRAGMENT_PREVIEW]
suffix = "" if len(missing) <= MISSING_FRAGMENT_PREVIEW else f" ... +{len(missing) - MISSING_FRAGMENT_PREVIEW} more"
print(
f"Expired frame {frame_num}: received {received}/{expected} fragments, "
f"missing {len(missing)} [{', '.join(map(str, preview))}{suffix}]"
)

def update_display(self, payload):
# Display the video in real time
try:
Expand All @@ -91,16 +158,19 @@ def update_display(self, payload):
print(f"Error displaying frame: {e}")

def update_stats(self):
total = self.received_frames + self.lost_frames
loss_rate = (self.lost_frames / total * 100) if total > 0 else 0
self.stats_label.config(text=f"Packets Received: {self.received_frames} | Lost: {self.lost_frames} | Loss Rate: {loss_rate:.2f}%")
total_frames = self.completed_frames + self.lost_frames
loss_rate = (self.lost_frames / total_frames * 100) if total_frames > 0 else 0
incomplete_frames = len(self.fragment_buffers)
self.stats_label.config(
text=f"Packets Received: {self.received_frames} | Complete Frames: {self.completed_frames} | Incomplete Frames: {incomplete_frames} | Expired Frames: {self.expired_frames} | Lost Frames: {self.lost_frames} | Loss Rate: {loss_rate:.2f}%"
)

def handler(self):
"""Clean up when exiting."""
self.running = False
try:
# Leave the multicast group when exiting
mreq = socket.inet_aton(MULTICAST_GROUP) + socket.inet_aton('0.0.0.0')
mreq = socket.inet_aton(MULTICAST_GROUP) + socket.inet_aton(self.interface_ip)
self.sock.setsockopt(socket.IPPROTO_IP, socket.IP_DROP_MEMBERSHIP, mreq)
self.sock.close()
print("Left multicast group.")
Expand All @@ -110,6 +180,11 @@ def handler(self):
self.master.destroy()

if __name__ == "__main__":
if len(sys.argv) > 2:
print("Usage: python Client.py [interface_ip]")
sys.exit(1)

interface_ip = sys.argv[1] if len(sys.argv) == 2 else None
root = tk.Tk()
client = Client(root)
client = Client(root, interface_ip)
root.mainloop()
35 changes: 25 additions & 10 deletions CustomPacket.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,35 +2,50 @@

class CustomPacket:
MAGIC = 0x1234
HEADER_FORMAT = "!HII"
HEADER_FORMAT = "!HIIII"
HEADER_SIZE = struct.calcsize(HEADER_FORMAT)

@staticmethod
def encode(frame_num, payload):
def encode(frame_num, payload, fragment_index=0, fragment_count=1):
"""
Encode the frame into a custom packet.
Encode a frame fragment into a custom packet.
Header:
- 2 bytes: Magic number (0x1234)
- 4 bytes: Frame number
- 4 bytes: Fragment index
- 4 bytes: Fragment count
- 4 bytes: Payload length
"""
header = struct.pack(CustomPacket.HEADER_FORMAT, CustomPacket.MAGIC, frame_num, len(payload))
header = struct.pack(
CustomPacket.HEADER_FORMAT,
CustomPacket.MAGIC,
frame_num,
fragment_index,
fragment_count,
len(payload),
)
return header + payload

@staticmethod
def decode(data):
"""
Decode the custom packet.
Returns (frame_num, payload) or (None, None) if invalid.
Returns (frame_num, fragment_index, fragment_count, payload) or
(None, None, None, None) if invalid.
"""
if len(data) < CustomPacket.HEADER_SIZE:
return None, None
return None, None, None, None

header = data[:CustomPacket.HEADER_SIZE]
magic, frame_num, length = struct.unpack(CustomPacket.HEADER_FORMAT, header)
magic, frame_num, fragment_index, fragment_count, length = struct.unpack(CustomPacket.HEADER_FORMAT, header)

if magic != CustomPacket.MAGIC:
return None, None

return None, None, None, None
if fragment_count == 0 or fragment_index >= fragment_count:
return None, None, None, None

payload = data[CustomPacket.HEADER_SIZE:CustomPacket.HEADER_SIZE+length]
return frame_num, payload
if len(payload) != length:
return None, None, None, None

return frame_num, fragment_index, fragment_count, payload
106 changes: 84 additions & 22 deletions Server.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,17 @@

MULTICAST_GROUP = '239.1.1.1'
MULTICAST_PORT = 5004
UDP_MTU = 1400
SOCKET_BUFFER_SIZE = 4 * 1024 * 1024
FRAGMENT_SEND_INTERVAL = 0.005

def get_default_interface_ip():
sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
try:
sock.connect((MULTICAST_GROUP, MULTICAST_PORT))
return sock.getsockname()[0]
finally:
sock.close()

class VideoStream:
def __init__(self, filename):
Expand All @@ -15,38 +26,80 @@ def __init__(self, filename):
print(f"Error: Could not open {filename}")
sys.exit(1)
self.frameNum = 0
self.lengthPrefixed = self._is_length_prefixed()

def _is_length_prefixed(self):
"""Detect the original sample format: 5 ASCII digits before each frame."""
prefix = self.file.read(5)
self.file.seek(0)
return len(prefix) == 5 and prefix.isdigit()

def nextFrame(self):
# The first 5 bytes represent the frame length
if not self.lengthPrefixed:
return self._next_jpeg_frame()

data = self.file.read(5)
if data:
try:
framelength = int(data)
frame = self.file.read(framelength)
self.frameNum += 1
return frame
except ValueError:
if not data:
return None

framelength = int(data)
frame = self.file.read(framelength)
if len(frame) != framelength:
return None

self.frameNum += 1
return frame

def _next_jpeg_frame(self):
"""Read one JPEG image from a standard concatenated MJPEG stream."""
frame = bytearray()
prev = None

while True:
byte = self.file.read(1)
if not byte:
return None
return None

value = byte[0]
if prev == 0xFF and value == 0xD8:
frame.extend((0xFF, 0xD8))
break
prev = value

prev = None
while True:
byte = self.file.read(1)
if not byte:
return None

value = byte[0]
frame.append(value)
if prev == 0xFF and value == 0xD9:
self.frameNum += 1
return bytes(frame)
prev = value

def reset(self):
self.file.seek(0)
self.frameNum = 0

def main():
if len(sys.argv) != 2:
print("Usage: python Server.py <file MJPEG>")
if len(sys.argv) not in (2, 3):
print("Usage: python Server.py <file MJPEG> [interface_ip]")
sys.exit(1)

filename = sys.argv[1]
interface_ip = sys.argv[2] if len(sys.argv) == 3 else get_default_interface_ip()

# Create UDP socket for multicast
sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM, socket.IPPROTO_UDP)
sock.setsockopt(socket.SOL_SOCKET, socket.SO_SNDBUF, SOCKET_BUFFER_SIZE)
sock.setsockopt(socket.IPPROTO_IP, socket.IP_MULTICAST_TTL, 2)
sock.setsockopt(socket.IPPROTO_IP, socket.IP_MULTICAST_IF, socket.inet_aton(interface_ip))

video_stream = VideoStream(filename)

print(f"Starting multicast streaming to {MULTICAST_GROUP}:{MULTICAST_PORT}...")
print(f"Starting multicast streaming to {MULTICAST_GROUP}:{MULTICAST_PORT} via {interface_ip}...")

while True:
frame = video_stream.nextFrame()
Expand All @@ -55,16 +108,25 @@ def main():
video_stream.reset()
continue

# Packetize the frame using our custom format
packet = CustomPacket.encode(video_stream.frameNum, frame)

# Send every frame to the multicast IP address
try:
sock.sendto(packet, (MULTICAST_GROUP, MULTICAST_PORT))
except Exception as e:
print(f"Failed to send packet: {e}")
# Split large frames so each UDP datagram stays under the target MTU.
max_payload_size = UDP_MTU - CustomPacket.HEADER_SIZE
fragment_count = (len(frame) + max_payload_size - 1) // max_payload_size

for fragment_index in range(fragment_count):
start = fragment_index * max_payload_size
payload = frame[start:start + max_payload_size]
packet = CustomPacket.encode(video_stream.frameNum, payload, fragment_index, fragment_count)

try:
sock.sendto(packet, (MULTICAST_GROUP, MULTICAST_PORT))
except Exception as e:
print(f"Failed to send packet: {e}")

# Avoid dropping large FHD frames by blasting hundreds of UDP packets at once.
if fragment_count > 1:
time.sleep(FRAGMENT_SEND_INTERVAL)

# Broadcast frames at approximately 20 FPS (50 ms/frame)
# Keep a baseline frame interval; large fragmented frames may run slower due to pacing.
time.sleep(0.05)

if __name__ == "__main__":
Expand Down
Binary file removed doc/Socket_Requirement.pdf
Binary file not shown.
Binary file removed doc/assets/Client_Init.png
Binary file not shown.
Binary file removed doc/assets/Client_running_2.png
Binary file not shown.
Binary file removed doc/assets/Clients_running.png
Binary file not shown.
Binary file removed doc/assets/Server_Init.png
Binary file not shown.
Binary file removed doc/assets/Statistics.png
Binary file not shown.
Binary file removed doc/assets/client_waiting_stream.png
Binary file not shown.
Binary file removed doc/assets/pause_packet_client.png
Binary file not shown.
Binary file removed doc/assets/pause_packet_server.png
Binary file not shown.
Binary file removed doc/assets/play_packet_client_udp.png
Binary file not shown.
Binary file removed doc/assets/play_packet_server.png
Binary file not shown.
Binary file removed doc/assets/play_packet_tcp.png
Binary file not shown.
Binary file removed doc/assets/setup_packet_client.png
Binary file not shown.
Binary file removed doc/assets/setup_packet_server.png
Binary file not shown.
Binary file removed doc/assets/setup_ui_udp.png
Binary file not shown.
Binary file removed doc/assets/tcp_setup.png
Binary file not shown.
Binary file removed doc/assets/teardown_packet_client.png
Binary file not shown.
Binary file removed doc/assets/teardown_packet_server.png
Binary file not shown.
Binary file removed doc/logohcmus.jpg
Binary file not shown.
Loading