-
Notifications
You must be signed in to change notification settings - Fork 7
Expand file tree
/
Copy pathclient.py
More file actions
executable file
·162 lines (136 loc) · 5.81 KB
/
Copy pathclient.py
File metadata and controls
executable file
·162 lines (136 loc) · 5.81 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
#!/usr/bin/env python3
from __future__ import annotations
import argparse
import asyncio
import sys
import time
from contextlib import nullcontext, suppress
from pathlib import Path
import av
from livepeer_gateway.errors import LivepeerGatewayError
from livepeer_gateway.live_runner import stop_runner_session
from livepeer_gateway.media_output import MediaOutput
from livepeer_gateway.media_publish import MediaPublish
from livepeer_gateway.http import post_json
from livepeer_gateway.selection import reserve_session
DEFAULT_DISCOVERY = "http://localhost:8935/discovery"
ECHO_APP_ID = "livepeer-sample/echo"
DEFAULT_OUTPUT = "echo-out.ts"
BLUR_UPDATE_INTERVAL_S = 0.01
MAX_BLUR_RADIUS = 100
def _log(*args: object) -> None:
print(*args, file=sys.stderr)
def _parse_args() -> argparse.Namespace:
parser = argparse.ArgumentParser(description="Run the proxied echo Live Runner demo.")
parser.add_argument("input")
parser.add_argument("--discovery", default=DEFAULT_DISCOVERY)
parser.add_argument("--output", default=DEFAULT_OUTPUT)
parser.add_argument("--radius", type=int, default=75)
parser.add_argument("--max-frames", type=int, default=0, help="Stop after this many input video frames (0 = full file).")
parser.add_argument("--blur", action="store_true", help="Sweep blur radius while publishing the sample.")
return parser.parse_args()
def _channel_url(echo_response: dict[str, object], name: str) -> str:
url = echo_response.get(name)
if not isinstance(url, str) or not url:
raise LivepeerGatewayError(f"echo response missing {name!r} url")
return url
async def _publish_video(
input_path: Path,
publish_url: str,
*,
max_frames: int = 0,
app_url: str = "",
blur: bool = False,
) -> None:
input_ = av.open(str(input_path))
try:
if not input_.streams.video:
raise LivepeerGatewayError(f"No video stream found in input file: {input_path}")
publisher = MediaPublish(publish_url)
prev_pts_time: float | None = None
prev_wall: float | None = None
next_update_pts_time: float | None = None
blur_radius = 0
blur_direction = 1
try:
for index, frame in enumerate(input_.decode(video=0), start=1):
if max_frames > 0 and index > max_frames:
break
current_pts_time = None
if frame.pts is not None and frame.time_base is not None:
current_pts_time = float(frame.pts * frame.time_base)
if next_update_pts_time is None:
next_update_pts_time = current_pts_time
while (
blur
and app_url
and current_pts_time is not None
and next_update_pts_time is not None
and current_pts_time >= next_update_pts_time
):
await post_json(f"{app_url.rstrip('/')}/update", {"mode": "blur", "radius": blur_radius})
if blur_radius == MAX_BLUR_RADIUS:
blur_direction = -1
elif blur_radius == 0:
blur_direction = 1
blur_radius += blur_direction
next_update_pts_time += BLUR_UPDATE_INTERVAL_S
if (
prev_pts_time is not None
and prev_wall is not None
and current_pts_time is not None
):
delta_s = current_pts_time - prev_pts_time
elapsed_s = time.monotonic() - prev_wall
sleep_s = max(0.0, delta_s - elapsed_s)
if sleep_s > 0:
await asyncio.sleep(sleep_s)
if current_pts_time is not None:
prev_pts_time = current_pts_time
prev_wall = time.monotonic()
await publisher.write_frame(frame)
finally:
await publisher.close()
finally:
input_.close()
async def main() -> None:
args = _parse_args()
input_path = Path(args.input).expanduser()
output_stdout = args.output.strip().lower() in {"-", "stdout"}
output_path = None if output_stdout else Path(args.output).expanduser()
if not input_path.exists():
raise SystemExit(f"input file does not exist: {input_path}")
session = None
try:
session = await reserve_session(discovery_url=args.discovery, app=ECHO_APP_ID)
_log("runner_url:", session.runner.url if session.runner is not None else session.runner_url)
_log("session_id:", session.session_id)
_log("app_url:", session.app_url)
echo = await post_json(f"{session.app_url.rstrip('/')}/echo", {"radius": args.radius})
in_url = _channel_url(echo, "in")
out_url = _channel_url(echo, "out")
_log("in:", in_url)
_log("out:", out_url)
with nullcontext(sys.stdout.buffer) if output_stdout else output_path.open("wb") as fh:
def _write_chunk(chunk: bytes) -> None:
fh.write(chunk)
if output_stdout:
fh.flush()
async with MediaOutput(out_url, on_bytes=_write_chunk):
await _publish_video(
input_path,
in_url,
max_frames=max(0, args.max_frames),
app_url=session.app_url,
blur=args.blur,
)
_log("publish complete; waiting for output to drain...")
fh.flush()
except LivepeerGatewayError as exc:
raise SystemExit(f"ERROR: {exc}") from exc
finally:
if session is not None:
with suppress(Exception):
await stop_runner_session(session)
if __name__ == "__main__":
asyncio.run(main())