205 lines
7.6 KiB
Python
205 lines
7.6 KiB
Python
#!/usr/bin/env python3
|
|
"""
|
|
WebRTC -> v4l2loopback bridge with stable output resolution (server-side reformat).
|
|
|
|
Usage:
|
|
python3 server.py --cert cert.pem --key key.pem --video-device /dev/video2 --host 0.0.0.0 --port 8443 --width 640 --height 480
|
|
|
|
Notes:
|
|
- If you omit --width/--height the server will use the first incoming frame's size.
|
|
- Requires the same dependencies as before: aiohttp, aiortc, av, numpy, ffmpeg installed system-wide.
|
|
"""
|
|
import argparse
|
|
import asyncio
|
|
import os
|
|
import ssl
|
|
import subprocess
|
|
from aiohttp import web
|
|
from aiortc import RTCPeerConnection, RTCSessionDescription
|
|
|
|
ROOT = os.path.join(os.path.dirname(__file__), "static")
|
|
pcs = set()
|
|
|
|
async def index(request):
|
|
return web.FileResponse(os.path.join(ROOT, "index.html"))
|
|
|
|
async def offer(request):
|
|
params = await request.json()
|
|
offer = RTCSessionDescription(sdp=params["sdp"], type=params["type"])
|
|
|
|
pc = RTCPeerConnection()
|
|
pcs.add(pc)
|
|
print("Created peer connection", pc)
|
|
|
|
# ffmpeg process holder per-track
|
|
ffmpeg_proc = {"proc": None}
|
|
|
|
@pc.on("iceconnectionstatechange")
|
|
def on_ice():
|
|
print("ICE:", pc.iceConnectionState)
|
|
if pc.iceConnectionState == "failed":
|
|
asyncio.ensure_future(pc.close())
|
|
|
|
@pc.on("track")
|
|
def on_track(track):
|
|
print("Track received:", track.kind)
|
|
|
|
if track.kind != "video":
|
|
return
|
|
|
|
async def run_video():
|
|
proc = None
|
|
target_w = request.app.get("force_width")
|
|
target_h = request.app.get("force_height")
|
|
|
|
try:
|
|
# get first frame (blocks until available)
|
|
first_frame = await track.recv()
|
|
# if client didn't force res, use first frame's dim
|
|
if not target_w or not target_h:
|
|
target_w = first_frame.width
|
|
target_h = first_frame.height
|
|
print(f"No forced size given; using first-frame size {target_w}x{target_h}")
|
|
else:
|
|
print(f"Using forced output size {target_w}x{target_h}")
|
|
|
|
# start ffmpeg once using target size
|
|
device = request.app["video_device"]
|
|
fps = request.app.get("force_fps", 30)
|
|
cmd = [
|
|
"ffmpeg",
|
|
"-hide_banner",
|
|
"-loglevel", "warning",
|
|
"-f", "rawvideo",
|
|
"-pix_fmt", "rgb24",
|
|
"-s", f"{target_w}x{target_h}",
|
|
"-r", str(fps),
|
|
"-i", "-",
|
|
"-f", "v4l2",
|
|
"-pix_fmt", "yuv420p",
|
|
device
|
|
]
|
|
print("Starting ffmpeg:", " ".join(cmd))
|
|
proc = subprocess.Popen(cmd, stdin=subprocess.PIPE, bufsize=0)
|
|
ffmpeg_proc["proc"] = proc
|
|
|
|
# reformat first frame to desired layout
|
|
try:
|
|
f_reformat = first_frame.reformat(width=target_w, height=target_h, format="rgb24")
|
|
except Exception:
|
|
# fallback: to_ndarray + numpy resize could go here, but reformat is preferred
|
|
f_reformat = first_frame
|
|
|
|
arr = f_reformat.to_ndarray(format="rgb24")
|
|
try:
|
|
proc.stdin.write(arr.tobytes())
|
|
proc.stdin.flush()
|
|
except BrokenPipeError:
|
|
print("BrokenPipeError writing first frame to ffmpeg")
|
|
except Exception as e:
|
|
print("Error writing first frame to ffmpeg:", e)
|
|
|
|
# loop and process subsequent frames: always reformat to target size/pixfmt
|
|
while True:
|
|
frame = await track.recv()
|
|
# reformat to target size/pixfmt; this creates a new VideoFrame
|
|
try:
|
|
out_frame = frame.reformat(width=target_w, height=target_h, format="rgb24")
|
|
except Exception:
|
|
# if reformat fails for any reason, skip the frame but continue
|
|
print("Warning: frame.reformat failed; skipping frame")
|
|
continue
|
|
|
|
arr = out_frame.to_ndarray(format="rgb24")
|
|
try:
|
|
proc.stdin.write(arr.tobytes())
|
|
except BrokenPipeError:
|
|
print("BrokenPipeError while writing frame to ffmpeg; terminating loop")
|
|
break
|
|
except Exception as e:
|
|
print("Error while writing frame to ffmpeg:", e)
|
|
break
|
|
|
|
except asyncio.CancelledError:
|
|
pass
|
|
except Exception as e:
|
|
print("Video worker error:", e)
|
|
finally:
|
|
print("Video worker ending")
|
|
p = ffmpeg_proc.get("proc")
|
|
if p:
|
|
try:
|
|
p.stdin.close()
|
|
except Exception:
|
|
pass
|
|
try:
|
|
p.terminate()
|
|
except Exception:
|
|
pass
|
|
try:
|
|
p.wait(timeout=1)
|
|
except Exception:
|
|
pass
|
|
ffmpeg_proc["proc"] = None
|
|
|
|
# schedule the video worker
|
|
asyncio.ensure_future(run_video())
|
|
|
|
@track.on("ended")
|
|
async def on_ended():
|
|
print("Track ended")
|
|
|
|
await pc.setRemoteDescription(offer)
|
|
answer = await pc.createAnswer()
|
|
await pc.setLocalDescription(answer)
|
|
|
|
return web.json_response(
|
|
{"sdp": pc.localDescription.sdp, "type": pc.localDescription.type}
|
|
)
|
|
|
|
async def on_shutdown(app):
|
|
coros = [pc.close() for pc in pcs]
|
|
await asyncio.gather(*coros)
|
|
pcs.clear()
|
|
|
|
def main():
|
|
parser = argparse.ArgumentParser()
|
|
parser.add_argument("--host", default="0.0.0.0")
|
|
parser.add_argument("--port", type=int, default=8443)
|
|
parser.add_argument("--cert", required=True, help="TLS cert (PEM)")
|
|
parser.add_argument("--key", required=True, help="TLS key (PEM)")
|
|
parser.add_argument("--video-device", default="/dev/video2", help="v4l2loopback device")
|
|
parser.add_argument("--width", type=int, dest="width", help="Force output width (optional)")
|
|
parser.add_argument("--height", type=int, dest="height", help="Force output height (optional)")
|
|
parser.add_argument("--fps", type=int, dest="fps", default=30, help="Output framerate (default 30)")
|
|
args = parser.parse_args()
|
|
|
|
if not os.path.exists(args.video_device):
|
|
print(f"Warning: video device {args.video_device} does not exist yet. Make sure v4l2loopback is loaded.")
|
|
|
|
app = web.Application()
|
|
app["video_device"] = args.video_device
|
|
app["force_width"] = args.width
|
|
app["force_height"] = args.height
|
|
app["force_fps"] = args.fps
|
|
|
|
app.router.add_get("/", index)
|
|
app.router.add_post("/offer", offer)
|
|
app.router.add_static("/static/", ROOT, show_index=False)
|
|
|
|
async def manifest_handler(request):
|
|
return web.FileResponse(os.path.join(ROOT, "manifest.json"))
|
|
async def sw_handler(request):
|
|
return web.FileResponse(os.path.join(ROOT, "sw.js"))
|
|
app.router.add_get("/manifest.json", manifest_handler)
|
|
app.router.add_get("/sw.js", sw_handler)
|
|
|
|
app.on_shutdown.append(on_shutdown)
|
|
|
|
ssl_context = ssl.create_default_context(ssl.Purpose.CLIENT_AUTH)
|
|
ssl_context.load_cert_chain(args.cert, args.key)
|
|
|
|
web.run_app(app, host=args.host, port=args.port, ssl_context=ssl_context)
|
|
|
|
if __name__ == "__main__":
|
|
main() |