Skip to content

Commit f13dcda

Browse files
authored
Add Janus streaming example (#205)
1 parent fe1739a commit f13dcda

2 files changed

Lines changed: 187 additions & 0 deletions

File tree

‎examples/README.md‎

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -23,3 +23,11 @@ without worrying about copyrights.
2323
```
2424
OPENAI_API_KEY=sk-... uv run https://raw.githubusercontent.com/MarshalX/python-webrtc/main/examples/openai_live.py
2525
```
26+
27+
### [janus_streaming.py](janus_streaming.py)
28+
29+
**Watching a live stream in the terminal.** Plays a stream of the public [Janus](https://janus.conf.meetecho.com/) demo server: the video drawn with colored characters, the audio on your speakers. No account or key:
30+
31+
```
32+
uv run https://raw.githubusercontent.com/MarshalX/python-webrtc/main/examples/janus_streaming.py
33+
```

‎examples/janus_streaming.py‎

Lines changed: 179 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,179 @@
1+
#!/usr/bin/env -S uv run --script
2+
# /// script
3+
# requires-python = ">=3.9"
4+
# dependencies = ["wrtc>=0.0.0.dev10", "sounddevice", "httpx"]
5+
# ///
6+
"""Watches a stream of the public Janus demo server: the video in the terminal, the audio on your speakers.
7+
8+
uv run janus_streaming.py
9+
uv run janus_streaming.py --id 1
10+
11+
The video is drawn with colored half blocks, so the terminal needs 24-bit color. Press Ctrl+C to stop.
12+
"""
13+
14+
import argparse
15+
import asyncio
16+
import contextlib
17+
import shutil
18+
import sys
19+
import threading
20+
import uuid
21+
22+
import httpx
23+
import sounddevice
24+
25+
import webrtc
26+
27+
JANUS = 'https://janus.conf.meetecho.com/janus'
28+
29+
30+
class Janus:
31+
"""A session with the streaming plugin of a Janus server, over its HTTP API"""
32+
33+
def __init__(self):
34+
self.client = httpx.AsyncClient(timeout=60)
35+
36+
async def __aenter__(self):
37+
created = await self._post('', janus='create')
38+
self.session = f'/{created["data"]["id"]}'
39+
attached = await self._post(self.session, janus='attach', plugin='janus.plugin.streaming')
40+
self.handle = f'{self.session}/{attached["data"]["id"]}'
41+
return self
42+
43+
async def __aexit__(self, *exc_info):
44+
await self._post(self.session, janus='destroy')
45+
await self.client.aclose()
46+
47+
async def request(self, body, **extra):
48+
reply = await self._post(self.handle, janus='message', body=body, **extra)
49+
return reply.get('plugindata', {}).get('data')
50+
51+
async def event(self):
52+
"""Waits up to 30 s for an event: requests reply right away and send their results, like offers, as events"""
53+
return (await self.client.get(JANUS + self.session)).json()
54+
55+
async def _post(self, path, **message):
56+
reply = (await self.client.post(JANUS + path, json={'transaction': uuid.uuid4().hex, **message})).json()
57+
if reply['janus'] == 'error':
58+
raise RuntimeError(reply['error']['reason'])
59+
return reply
60+
61+
62+
class Speakers:
63+
"""Plays 16-bit audio from a buffer that keeps half a second at most"""
64+
65+
def __init__(self, rate, channels):
66+
self.buffer, self.lock, self.limit = bytearray(), threading.Lock(), rate * channels
67+
self.stream = sounddevice.RawOutputStream(rate, channels=channels, dtype='int16', callback=self._on_need)
68+
self.stream.start()
69+
70+
def play(self, samples):
71+
with self.lock:
72+
self.buffer += samples
73+
del self.buffer[: -self.limit]
74+
75+
def _on_need(self, out, *_):
76+
with self.lock:
77+
chunk = self.buffer[: len(out)]
78+
del self.buffer[: len(out)]
79+
out[:] = chunk.ljust(len(out), b'\0')
80+
81+
82+
def draw(rgbx, width, height):
83+
"""Draws a frame with ▀, whose foreground color is the upper pixel and background the lower one"""
84+
columns, rows = shutil.get_terminal_size()
85+
scale = max(width / columns, height / rows / 2)
86+
87+
def color(x, y):
88+
i = (int(y * scale) * width + int(x * scale)) * 4
89+
return '{};{};{}'.format(*rgbx[i : i + 3])
90+
91+
lines = (
92+
''.join(f'\033[38;2;{color(x, 2 * y)}m\033[48;2;{color(x, 2 * y + 1)}m▀' for x in range(int(width / scale)))
93+
for y in range(int(height / scale / 2))
94+
)
95+
sys.stdout.write('\033[H' + '\033[0m\033[K\n'.join(lines) + '\033[0m\033[J')
96+
sys.stdout.flush()
97+
98+
99+
@contextlib.contextmanager
100+
def fullscreen():
101+
"""Switches to the alternate screen, without the cursor"""
102+
sys.stdout.write('\033[?1049h\033[?25l')
103+
try:
104+
yield
105+
finally:
106+
sys.stdout.write('\033[?25h\033[?1049l')
107+
108+
109+
async def watch(track):
110+
# a buffer of one frame drops the frames the terminal is too slow for
111+
async for frame in webrtc.MediaStreamTrackProcessor(track, max_buffer_size=1).readable:
112+
with frame:
113+
rgbx = bytearray(frame.allocation_size({'format': 'RGBX'}))
114+
await frame.copy_to(rgbx, {'format': 'RGBX'})
115+
size = frame.visible_rect
116+
draw(rgbx, int(size.width), int(size.height))
117+
118+
119+
async def listen(track):
120+
speakers = None
121+
try:
122+
async for data in webrtc.MediaStreamTrackProcessor(track, max_buffer_size=50).readable:
123+
with data:
124+
options = {'plane_index': 0, 'format': 's16'}
125+
samples = bytearray(data.allocation_size(options))
126+
data.copy_to(samples, options)
127+
speakers = speakers or Speakers(int(data.sample_rate), data.number_of_channels)
128+
speakers.play(samples)
129+
finally:
130+
if speakers is not None:
131+
speakers.stream.close()
132+
133+
134+
async def answer(pc, offer):
135+
"""Answers with all the ICE candidates in the SDP, since there is no trickling"""
136+
gathered = asyncio.Event()
137+
pc.on('icegatheringstatechange', lambda _: pc.ice_gathering_state == 'complete' and gathered.set())
138+
await pc.set_remote_description(offer)
139+
await pc.set_local_description(await pc.create_answer())
140+
with contextlib.suppress(asyncio.TimeoutError):
141+
await asyncio.wait_for(gathered.wait(), 5)
142+
return {'type': 'answer', 'sdp': pc.local_description.sdp}
143+
144+
145+
async def main(stream_id):
146+
pc, tasks = webrtc.RTCPeerConnection(), []
147+
148+
@pc.on('track')
149+
def on_track(event):
150+
play = watch if event.track.kind == 'video' else listen
151+
tasks.append(asyncio.ensure_future(play(event.track)))
152+
153+
async with Janus() as janus:
154+
if stream_id is None:
155+
streams = (await janus.request({'request': 'list'}))['list']
156+
for stream in streams:
157+
print(f'{stream["id"]}: {stream.get("description")}')
158+
stream_id = streams[0]['id']
159+
160+
await janus.request({'request': 'watch', 'id': stream_id})
161+
event = {}
162+
while 'jsep' not in event:
163+
event = await janus.event()
164+
await janus.request({'request': 'start'}, jsep=await answer(pc, event['jsep']))
165+
166+
with fullscreen():
167+
try:
168+
while True: # polling for events keeps the session alive
169+
await janus.event()
170+
finally:
171+
pc.close() # ends the tracks, so watch and listen return
172+
await asyncio.gather(*tasks)
173+
174+
175+
if __name__ == '__main__':
176+
parser = argparse.ArgumentParser(description='Watch a stream of the Janus demo server in the terminal.')
177+
parser.add_argument('--id', type=int, help='the stream to watch, default: the first one')
178+
with contextlib.suppress(KeyboardInterrupt):
179+
asyncio.run(main(parser.parse_args().id))

0 commit comments

Comments
 (0)