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
|
import asyncio
import asyncio.subprocess
from typing import List
from typing import Dict
from typing import Optional
from .logging import get_logger
from . import gpio
# =====
class Streamer: # pylint: disable=too-many-instance-attributes
def __init__(
self,
cap_power: int,
conv_power: int,
sync_delay: float,
init_delay: float,
init_restart_after: float,
quality: int,
cmd: List[str],
loop: asyncio.AbstractEventLoop,
) -> None:
self.__cap_power = (gpio.set_output(cap_power) if cap_power > 0 else cap_power)
self.__conv_power = (gpio.set_output(conv_power) if conv_power > 0 else conv_power)
self.__sync_delay = sync_delay
self.__init_delay = init_delay
self.__init_restart_after = init_restart_after
self.__quality = quality
self.__cmd = cmd
self.__loop = loop
self.__proc_task: Optional[asyncio.Task] = None
async def start(self, quality: int, no_init_restart: bool=False) -> None:
logger = get_logger()
logger.info("Starting streamer ...")
assert 1 <= quality <= 100
self.__quality = quality
await self.__inner_start()
if self.__init_restart_after > 0.0 and not no_init_restart:
logger.info("Stopping streamer to restart ...")
await self.__inner_stop()
logger.info("Starting again ...")
await self.__inner_start()
async def stop(self) -> None:
get_logger().info("Stopping streamer ...")
await self.__inner_stop()
def is_running(self) -> bool:
return bool(self.__proc_task)
def get_current_quality(self) -> int:
return self.__quality
def get_state(self) -> Dict:
return {
"is_running": self.is_running(),
"quality": self.__quality,
}
async def cleanup(self) -> None:
if self.is_running():
await self.stop()
async def __inner_start(self) -> None:
assert not self.__proc_task
await self.__set_hw_enabled(True)
self.__proc_task = self.__loop.create_task(self.__process())
async def __inner_stop(self) -> None:
assert self.__proc_task
self.__proc_task.cancel()
await asyncio.gather(self.__proc_task, return_exceptions=True)
await self.__set_hw_enabled(False)
self.__proc_task = None
async def __set_hw_enabled(self, enabled: bool) -> None:
# XXX: This sequence is very important to enable converter and cap board
if self.__cap_power > 0:
gpio.write(self.__cap_power, enabled)
if self.__conv_power > 0:
if enabled:
await asyncio.sleep(self.__sync_delay)
gpio.write(self.__conv_power, enabled)
if enabled:
await asyncio.sleep(self.__init_delay)
async def __process(self) -> None: # pylint: disable=too-many-branches
logger = get_logger(0)
while True: # pylint: disable=too-many-nested-blocks
proc: Optional[asyncio.subprocess.Process] = None # pylint: disable=no-member
try:
cmd = [part.format(quality=self.__quality) for part in self.__cmd]
proc = await asyncio.create_subprocess_exec(
*cmd,
stdout=asyncio.subprocess.PIPE,
stderr=asyncio.subprocess.STDOUT,
)
logger.info("Started streamer pid=%d: %s", proc.pid, cmd)
empty = 0
async for line_bytes in proc.stdout: # type: ignore
line = line_bytes.decode(errors="ignore").strip()
if line:
logger.info("Streamer: %s", line)
empty = 0
else:
empty += 1
if empty == 100: # asyncio bug
raise RuntimeError("Streamer/asyncio: too many empty lines")
raise RuntimeError("Streamer unexpectedly died")
except asyncio.CancelledError:
break
except Exception as err:
if proc:
logger.exception("Unexpected streamer error: pid=%d", proc.pid)
else:
logger.exception("Can't start streamer: %s", err)
await asyncio.sleep(1)
finally:
if proc and proc.returncode is None:
await self.__kill(proc)
async def __kill(self, proc: asyncio.subprocess.Process) -> None: # pylint: disable=no-member
try:
proc.terminate()
await asyncio.sleep(1)
if proc.returncode is None:
try:
proc.kill()
except Exception:
if proc.returncode is not None:
raise
await proc.wait()
get_logger().info("Streamer killed: pid=%d; retcode=%d", proc.pid, proc.returncode)
except Exception:
if proc.returncode is None:
get_logger().exception("Can't kill streamer pid=%d", proc.pid)
else:
get_logger().info("Streamer killed: pid=%d; retcode=%d", proc.pid, proc.returncode)
|