mirror of
https://github.com/sunnypilot/sunnypilot.git
synced 2026-08-21 15:43:45 +08:00
webrtc: cleanup api (#38447)
* add support for multi-track with wideRoad single track as default * remove stale debug mode * multi track support * make cameras list optional * remove multi camera support from athena, make it only local * update to teleop master
This commit is contained in:
@@ -589,7 +589,7 @@ def startStream(sdp: str, enabled: bool) -> dict:
|
||||
# wait for webrtcd end points to wake up
|
||||
wait_for_webrtcd()
|
||||
|
||||
return post_stream_request(StreamRequestBody(sdp, "wideRoad", enabled, bridge_services_in, ["carState", "deviceState"]))
|
||||
return post_stream_request(StreamRequestBody(sdp, ["wideRoad"], enabled, bridge_services_in, ["carState", "deviceState"]))
|
||||
|
||||
|
||||
def get_logs_to_send_sorted() -> list[str]:
|
||||
|
||||
@@ -8,7 +8,7 @@ WEBRTCD_PORT = 5001
|
||||
@dataclass
|
||||
class StreamRequestBody:
|
||||
sdp: str
|
||||
init_camera: str
|
||||
cameras: list[str]
|
||||
enabled: bool
|
||||
bridge_services_in: list[str] = field(default_factory=list)
|
||||
bridge_services_out: list[str] = field(default_factory=list)
|
||||
|
||||
@@ -221,7 +221,7 @@ class LivestreamBitrateController(AsyncTaskRunner):
|
||||
class StreamSession:
|
||||
shared_pub_master = DynamicPubMaster([])
|
||||
|
||||
def __init__(self, body: StreamRequestBody, debug_mode: bool = False):
|
||||
def __init__(self, body: StreamRequestBody):
|
||||
from openpilot.system.webrtc.device.video import LiveStreamVideoStreamTrack
|
||||
from teleoprtc.builder import WebRTCAnswerBuilder
|
||||
|
||||
@@ -230,8 +230,11 @@ class StreamSession:
|
||||
builder = WebRTCAnswerBuilder(body.sdp, bind_address=_default_route_ip())
|
||||
|
||||
self.enabled = body.enabled
|
||||
self.video_track = LiveStreamVideoStreamTrack(body.init_camera, self.enabled)
|
||||
builder.add_video_stream(body.init_camera, self.video_track)
|
||||
self.video_tracks = []
|
||||
for camera in body.cameras:
|
||||
track = LiveStreamVideoStreamTrack(camera, self.enabled)
|
||||
self.video_tracks.append(track)
|
||||
builder.add_video_stream(camera, track)
|
||||
self.stream = builder.stream()
|
||||
|
||||
self.is_body = "testJoystick" in body.bridge_services_in
|
||||
@@ -251,8 +254,8 @@ class StreamSession:
|
||||
self._cleanup_done = False
|
||||
self.logger = logging.getLogger("webrtcd")
|
||||
self.logger.info(
|
||||
"New stream session (%s), init camera %s, video enabled %s, incoming services %s, outgoing services %s",
|
||||
self.identifier, body.init_camera, body.enabled, body.bridge_services_in, body.bridge_services_out,
|
||||
"New stream session (%s), video cameras %s, video enabled %s, incoming services %s, outgoing services %s",
|
||||
self.identifier, [t.id for t in self.video_tracks], body.enabled, body.bridge_services_in, body.bridge_services_out,
|
||||
)
|
||||
|
||||
def start(self):
|
||||
@@ -277,14 +280,17 @@ class StreamSession:
|
||||
|
||||
match msg_type:
|
||||
case "livestreamCameraSwitch":
|
||||
self.video_track.switch_camera(payload["data"]["camera"])
|
||||
# only needed for 1 track stream
|
||||
if len(self.video_tracks) == 1:
|
||||
self.video_tracks[0].switch_camera(payload["data"]["camera"])
|
||||
case "livestreamSettings":
|
||||
if self.bitrate_controller is not None:
|
||||
self.bitrate_controller.set_quality(payload["data"]["quality"])
|
||||
case "livestreamVideoEnable":
|
||||
enabled = payload["data"]["enabled"]
|
||||
self.enabled = enabled
|
||||
self.video_track.enable(enabled)
|
||||
for track in self.video_tracks:
|
||||
track.enable(enabled)
|
||||
if self.outgoing_bridge is not None:
|
||||
self.outgoing_bridge.enable(enabled)
|
||||
if self.bitrate_controller is not None:
|
||||
@@ -297,8 +303,8 @@ class StreamSession:
|
||||
}})
|
||||
self.stream.get_messaging_channel().send(pong)
|
||||
case "enableTimingSei":
|
||||
if hasattr(self.video_track, 'timing_sei_enabled'):
|
||||
self.video_track.timing_sei_enabled = bool(payload["data"]["enabled"])
|
||||
for track in self.video_tracks:
|
||||
track.timing_sei_enabled = bool(payload["data"]["enabled"])
|
||||
case _:
|
||||
if msg_type not in self.incoming_bridge_services:
|
||||
return
|
||||
@@ -356,17 +362,16 @@ class StreamSession:
|
||||
await self.bitrate_controller.stop()
|
||||
if self.outgoing_bridge is not None:
|
||||
await self.outgoing_bridge.stop()
|
||||
if self.video_track is not None:
|
||||
self.video_track.stop()
|
||||
self.video_track = None
|
||||
for track in self.video_tracks:
|
||||
track.stop()
|
||||
self.video_tracks.clear()
|
||||
await self.stream.stop()
|
||||
|
||||
|
||||
class ServerState:
|
||||
def __init__(self, debug: bool):
|
||||
def __init__(self):
|
||||
self.streams: dict[str, StreamSession] = {}
|
||||
self.stream_lock = asyncio.Lock()
|
||||
self.debug = debug
|
||||
self.teardown: asyncio.TimerHandle | None = None
|
||||
|
||||
|
||||
@@ -391,7 +396,7 @@ def _text_response(text: str, status: int = 200) -> tuple[int, bytes, str]:
|
||||
|
||||
|
||||
async def handle_get_stream(state: ServerState, raw_body: bytes) -> tuple[int, bytes, str]:
|
||||
stream_dict, debug_mode = state.streams, state.debug
|
||||
stream_dict = state.streams
|
||||
body = StreamRequestBody(**json.loads(raw_body))
|
||||
|
||||
async with state.stream_lock:
|
||||
@@ -410,7 +415,7 @@ async def handle_get_stream(state: ServerState, raw_body: bytes) -> tuple[int, b
|
||||
await s.stop()
|
||||
stream_dict.pop(sid, None)
|
||||
|
||||
session = StreamSession(body, debug_mode)
|
||||
session = StreamSession(body)
|
||||
stream_dict[session.identifier] = session
|
||||
try:
|
||||
answer = await asyncio.wait_for(session.get_answer(), timeout=30)
|
||||
@@ -558,23 +563,23 @@ async def _shutdown(server: WebrtcdHTTPServer, state: ServerState, loop: asyncio
|
||||
loop.stop()
|
||||
|
||||
|
||||
def prewarm_stream_session_imports(debug_mode: bool = False) -> None:
|
||||
def prewarm_stream_session_imports() -> None:
|
||||
from openpilot.system.webrtc.device.video import LiveStreamVideoStreamTrack
|
||||
from teleoprtc.builder import WebRTCAnswerBuilder
|
||||
assert LiveStreamVideoStreamTrack
|
||||
assert WebRTCAnswerBuilder
|
||||
|
||||
|
||||
def webrtcd_thread(host: str, port: int, debug: bool):
|
||||
def webrtcd_thread(host: str, port: int):
|
||||
logging.basicConfig(level=logging.INFO, handlers=[logging.StreamHandler()])
|
||||
prewarm_start = time.monotonic()
|
||||
prewarm_stream_session_imports(debug)
|
||||
prewarm_stream_session_imports()
|
||||
prewarm_end = time.monotonic()
|
||||
logging.getLogger("webrtcd").info(f"webrtc prewarm finished in {(prewarm_end - prewarm_start) * 1000} ms")
|
||||
|
||||
loop = asyncio.new_event_loop()
|
||||
asyncio.set_event_loop(loop)
|
||||
state = ServerState(debug)
|
||||
state = ServerState()
|
||||
|
||||
server = WebrtcdHTTPServer((host, port), WebrtcdHandler)
|
||||
server.state = state
|
||||
@@ -608,10 +613,9 @@ def main():
|
||||
parser = argparse.ArgumentParser(description="WebRTC daemon")
|
||||
parser.add_argument("--host", type=str, default="0.0.0.0", help="Host to listen on")
|
||||
parser.add_argument("--port", type=int, default=5001, help="Port to listen on")
|
||||
parser.add_argument("--debug", action="store_true", help="Enable debug mode")
|
||||
args = parser.parse_args()
|
||||
|
||||
webrtcd_thread(args.host, args.port, args.debug)
|
||||
webrtcd_thread(args.host, args.port)
|
||||
|
||||
|
||||
if __name__=="__main__":
|
||||
|
||||
+1
-1
Submodule teleoprtc_repo updated: 3ac144dc9a...3f0b03a1d1
Reference in New Issue
Block a user