from textual import on
from textual.app import App
-from textual.containers import Vertical, Horizontal, Grid, ScrollableContainer
+from textual.containers import Vertical, Horizontal, ScrollableContainer
from textual.widgets import Footer, Header, Placeholder, TabbedContent, TabPane
from . import log, data
-from .interface.element.text_readout import TextReadout
-from .interface import RecordControls, OBSSettingsPane, obs
+from .interface import RecordControls, StatsDisplay, OBSSettingsPane, obs
logger = log.getLogger(__name__)
class OBSCtlApp(App):
- logger.debug(importlib.resources.files(data))
CSS_PATH = importlib.resources.files(data).joinpath("layout.tcss")
def compose(self):
yield Placeholder("CONNECTION BLOCK", id="conection-plc")
with TabPane("obs-ctl Settings", id="clients-tab"):
yield Placeholder("CLIENT SETTINGS BLOCK", id="client-plc")
- with Horizontal(id="blocks_h_container2", classes="blocks_h_containers"):
- with Grid(id="machine_container", classes="blocks gridblocks"):
- yield TextReadout("CPU Usage", "-.-%", id="cpu_usage", classes="stats")
- yield TextReadout("RAM Usage", "-b", id="memory_usage", classes="stats")
- yield TextReadout("Time 'til Full", "N/A", id="disk_full", classes="stats")
- with Grid(id="obs_container", classes="blocks gridblocks"):
- yield TextReadout("FPS", "--.--", id="fps", classes="stats")
- yield TextReadout("Netwrk frames dropped", "-/- (-.-%)", id="stream_dropped_frames", classes="stats")
- yield TextReadout("Render frames dropped", "-/- (-.-%)", id="rendering_lag", classes="stats")
- yield TextReadout("Encdng frames dropped", "-/- (-.-%)", id="encoding_lag", classes="stats")
+ yield StatsDisplay(id="stats_display")
@on(obs.API.Report, "#api")
def on_api_report(self, message: obs.API.Report):
self.update_monitors(message.data)
-
+
+
@on(RecordControls.RecordStarted, "#record_controls")
def on_record_started(self, message: RecordControls.RecordStarted):
- api: obs.API = self.query_exactly_one("#api")
- api.start_record()
+ self._query_api("start_record")
+
@on(RecordControls.RecordStopped, "#record_controls")
def on_record_stopped(self, message: RecordControls.RecordStopped):
- api: obs.API = self.query_exactly_one("#api")
- api.stop_record()
+ self._query_api("stop_record")
+
@on(RecordControls.StreamStarted, "#record_controls")
def on_stream_started(self, message: RecordControls.StreamStarted):
- api: obs.API = self.query_exactly_one("#api")
- api.start_stream()
+ self._query_api("start_stream")
+
@on(RecordControls.StreamStopped, "#record_controls")
def on_stream_stopped(self, message: RecordControls.StreamStopped):
- api: obs.API = self.query_exactly_one("#api")
- api.stop_stream()
+ self._query_api("stop_stream")
+
@on(RecordControls.RecordPaused, "#record_controls")
def on_record_paused(self, message: RecordControls.RecordPaused):
- api: obs.API = self.query_exactly_one("#api")
- api.start_pause()
+ self._query_api("start_pause")
+
@on(RecordControls.RecordUnpaused, "#record_controls")
def on_record_unpaused(self, message: RecordControls.RecordUnpaused):
+ self._query_api("stop_pause")
+
+
+ @on(RecordControls.OBSWSConnected, "#record_controls")
+ def on_obs_ws_connected(self, message: RecordControls.OBSWSConnected):
+ self._query_api("connect")
+
+
+ @on(RecordControls.OBSWSDisconnected, "#record_controls")
+ def on_obs_ws_connected(self, message: RecordControls.OBSWSDisconnected):
+ self._query_api("disconnect")
+
+
+ def _query_api(self, query: str, *args, **kwargs):
api: obs.API = self.query_exactly_one("#api")
- api.stop_pause()
+ getattr(api, query)(*args, **kwargs)
+
@on(OBSSettingsPane.RecordPathChanged, "#settings-content")
def on_settings_record_path_changed(self, message: OBSSettingsPane.RecordPathChanged):
api: obs.API = self.query_exactly_one("#api")
api.set_record_path(message.record_path)
+
def update_monitors(self, obs_stats: obs.OBSStats):
record_controls: RecordControls = self.query_exactly_one("#record_controls")
- cpu_usage: TextReadout = self.query_exactly_one("#cpu_usage")
- memory_usage: TextReadout = self.query_exactly_one("#memory_usage")
- active_fp: TextReadout = self.query_exactly_one("#fps")
- disk_full: TextReadout = self.query_exactly_one("#disk_full")
- stream_dropped_frames: TextReadout = self.query_exactly_one("#stream_dropped_frames")
- rendering_lag: TextReadout = self.query_exactly_one("#rendering_lag")
- encoding_lag: TextReadout = self.query_exactly_one("#encoding_lag")
obs_settings: OBSSettingsPane = self.query_exactly_one("#settings-content")
+ stats_display: StatsDisplay = self.query_exactly_one("#stats_display")
+ record_controls.obsws_connected = obs_stats.obsws_connection_active
record_controls.streaming = obs_stats.stream_output_active
record_controls.recording = obs_stats.record_output_active
record_controls.paused = obs_stats.record_output_paused
-
- cpu_usage.text = f"{obs_stats.cpu_usage_percent:0.5}%"
- memory_usage.text = f"{obs_stats.memory_usage_mb:0.5}MB"
- active_fp.text = f"{obs_stats.active_fps}"
- disk_full.text = f"{obs_stats.estimated_hours_remaining:0.5} hours"
- stream_dropped_frames.text = f"{obs_stats.stream_output_skipped_frames}/{obs_stats.stream_output_total_frames}"
- rendering_lag.text = f"{obs_stats.stats_render_skipped_frames}/{obs_stats.stats_render_total_frames}"
- encoding_lag.text = f"{obs_stats.stats_output_skipped_frames}/{obs_stats.stats_output_total_frames}"
obs_settings.record_path = obs_stats.record_directory
-
+ stats_display.update_text_readouts(obs_stats)
def main():
--- /dev/null
+from textual.containers import Grid, Horizontal, Container
+from textual.reactive import reactive
+
+from ..element import TextReadout
+from .. import obs
+
+
+class StatsDisplay(Container):
+ DEFAULT_CSS="""
+ StatsDisplay {
+ Grid {
+ grid-size:4;
+ }
+ }
+ """
+ COMPONENT_CLASSES = {
+ "stats_display--cpu",
+ "stats_display--ram",
+ "stats_display--hdd",
+ "stats_display--fps",
+ "stats_display--dropped_net",
+ "stats_display--droped_ren",
+ "stats-display--dropped_enc"
+ }
+
+
+ def compose(self):
+ with Grid():
+ yield TextReadout("CPU Usage", "-.-%", id="cpu", classes="stats_display--cpu")
+ yield TextReadout("RAM Usage", "-b", id="ram", classes="stats_display--ram")
+ yield TextReadout("Time 'til Full", "N/A", id="hdd", classes="stats_display--hdd")
+ yield TextReadout("FPS", "--.--", id="fps", classes="stats_display-fps")
+ yield TextReadout("Netwrk frames dropped", "-/- (-.-%)", id="dropped_net", classes="stats_display--dropped_net")
+ yield TextReadout("Render frames dropped", "-/- (-.-%)", id="dropped_ren", classes="stats_display--dropped_ren")
+ yield TextReadout("Encdng frames dropped", "-/- (-.-%)", id="dropped_enc", classes="stats_display--dropped_enc")
+
+
+ def update_text_readouts(self, obs_stats: obs.OBSStats):
+ def _try_catch_div_by_zero(drp: float, tot: float) -> float:
+ if tot > 0:
+ return float(drp) / float(tot)
+ else:
+ return 0.0
+
+ cpu_usage: TextReadout = self.query_exactly_one("#cpu")
+ memory_usage: TextReadout = self.query_exactly_one("#ram")
+ active_fp: TextReadout = self.query_exactly_one("#fps")
+ disk_full: TextReadout = self.query_exactly_one("#hdd")
+ stream_dropped_frames: TextReadout = self.query_exactly_one("#dropped_net")
+ rendering_lag: TextReadout = self.query_exactly_one("#dropped_ren")
+ encoding_lag: TextReadout = self.query_exactly_one("#dropped_enc")
+ netwrk_drp = obs_stats.stream_output_skipped_frames
+ netwrk_tot = obs_stats.stream_output_total_frames
+ render_drp = obs_stats.stats_render_skipped_frames
+ render_tot = obs_stats.stats_render_total_frames
+ encdng_drp = obs_stats.stats_output_skipped_frames
+ encdng_tot = obs_stats.stats_output_total_frames
+ netwrk_drp_pct: float = _try_catch_div_by_zero(netwrk_drp, netwrk_tot)
+ render_drp_pct: float = _try_catch_div_by_zero(render_drp, render_tot)
+ encdng_drp_pct: float = _try_catch_div_by_zero(encdng_drp, encdng_tot)
+
+ cpu_usage.text = f"{obs_stats.cpu_usage_percent:0.5}%"
+ memory_usage.text = f"{obs_stats.memory_usage_mb:0.5}MB"
+ active_fp.text = f"{obs_stats.active_fps}"
+ disk_full.text = f"{obs_stats.estimated_hours_remaining:0.5} hours"
+ stream_dropped_frames.text = f"{netwrk_drp}/{netwrk_tot} ({netwrk_drp_pct:0.4}%)"
+ rendering_lag.text = f"{render_drp}/{render_tot} ({render_drp_pct:0.4}%)"
+ encoding_lag.text = f"{encdng_drp}/{encdng_tot} ({encdng_drp_pct:0.4}%)"
\ No newline at end of file
import functools
from pathlib import Path
from time import ctime
+import time
import obsws_python as obs
import obsws_python.error as obs_error
"""Multiply milliseconds by this reciprocal to get hours out."""
MB_TO_BYTES = 1_048_576
"""Multiply megabytes by this number to get bytes. MB, not MiB"""
-
+TRY_AGAIN = 3
+"""Time to wait before trying again if OBS WebSocket isn't ready for requests."""
logger = getLogger(__name__)
@dataclass(frozen=True)
class OBSStats:
"""A dataclass for transmitting OBS stats."""
- record_output_active: bool
- record_output_paused: bool
- record_output_bytes: float
+ obsws_connection_active: bool = False
+ record_output_active: bool = False
+ record_output_paused: bool = False
+ record_output_bytes: float = 0.0
"""Total bytes to disk."""
- record_output_duration: float
- stream_output_active: bool
- stream_output_skipped_frames: float
- stream_output_total_frames: float
- stats_render_skipped_frames: float
- stats_render_total_frames: float
- stats_output_skipped_frames: float
- stats_output_total_frames: float
- cpu_usage_percent: float
- memory_usage_mb: float
- active_fps: float
- available_disk_space: float
- record_directory: str
+ record_output_duration: float = 0.0
+ stream_output_active: bool = False
+ stream_output_skipped_frames: float = 0.0
+ stream_output_total_frames: float = 0.0
+ stats_render_skipped_frames: float = 0.0
+ stats_render_total_frames: float = 0.0
+ stats_output_skipped_frames: float = 0.0
+ stats_output_total_frames: float = 0.0
+ cpu_usage_percent: float = 0.0
+ memory_usage_mb: float = 0.0
+ active_fps: float = 0.0
+ available_disk_space: float = 0.0
+ record_directory: str = ""
@property
def estimated_hours_remaining(self):
It stores no states. It reads, writes, and outputs states.
- Messages: StreamStateUpdated, RecordStateUpdated"""
+ Messages: StreamStateUpdated, RecordStateUpdated, Report"""
event_client: obs.EventClient = None
request_client: obs.ReqClient = None
+ is_connected: reactive[bool] = False
+ _trying_request: bool = False
_timer: Timer = None
interval: reactive[float] = reactive(1.0)
+ ws_host: reactive[str] = reactive("localhost")
+ ws_port: reactive[int] = reactive(4455)
+ ws_password: reactive[str] = reactive("")
+ ws_timeout: (reactive[float] | reactive[int] | reactive[None]) = reactive(None)
+ # ENUMS
class OutputStates(StrEnum):
+ """Named output states from OBS WebSocket
+
+ UKNOWN should never occur. Otherwise, they are relatively-self-explanatory.
+
+ UNKNOWN, STARTED, STARTING, PAUSED, RESUMED, STOPPED, STOPPING,
+ RECONNECTING, RECONNECTED.
+ """
UNKNOWN = "OBS_WEBSOCKET_OUTPUT_UNKNOWN"
STARTED = "OBS_WEBSOCKET_OUTPUT_STARTED"
STARTING = "OBS_WEBSOCKET_OUTPUT_STARTING"
OUTPUT_PAUSED = 502
OUTPUT_NOT_PAUSED = 503
+ # MESSAGES
class APIOutputMSGBase(MSGBase):
state = None
def __init__(self, control, data):
logger.debug(f"{self} is an obs.API Message; state: {_state.name}")
self.state = _state
- class StreamStateUpdated(APIOutputMSGBase): ...
+ class StreamStateUpdated(APIOutputMSGBase):
+ """Posted when OBS WebSocket reports a change to the Stream state.
+
+ state should be something from API.OutputStates."""
- class RecordStateUpdated(APIOutputMSGBase): ...
+ class RecordStateUpdated(APIOutputMSGBase):
+ """Posted when OBS WebSocket reports a change to the Record state.
+
+ state should be something from API.OutputStates."""
+ pass
class Report(MSGBase):
+ """Posted ever self.interval second. data is an OBSStats object."""
data: OBSStats = None
def __init__(self, control, data: OBSStats):
# this message is sent too often to log it.
self.data = data
- def __init__(self, interval, content = "", *, expand = False, shrink = False, markup = True, name = None, id = None, classes = None, disabled = False):
- super().__init__(content, expand=expand, shrink=shrink, markup=markup, name=name, id=id, classes=classes, disabled=disabled)
+ def __init__(
+ self,
+ interval,
+ ws_host: str = "localhost",
+ ws_port: int = 4455,
+ ws_password: str = "",
+ ws_timeout: int | float | None = None,
+ *,
+ content = "",
+ expand = False,
+ shrink = False,
+ markup = True,
+ name = None,
+ id = None,
+ classes = None,
+ disabled = False
+ ):
logger.debug(f"Initializing obs.API widget: {self}")
- self.request_client = obs.ReqClient()
- self.event_client = obs.EventClient()
- on_stream_state_changed = self.on_stream_state_changed
- on_record_state_changed = self.on_record_state_changed
- self.event_client.callback.register(on_stream_state_changed)
- self.event_client.callback.register(on_record_state_changed)
+ super().__init__(content, expand=expand, shrink=shrink, markup=markup, name=name, id=id, classes=classes, disabled=disabled)
self.set_reactive(API.interval, interval)
+ self.set_reactive(API.ws_host, ws_host)
+ self.set_reactive(API.ws_port, ws_port)
+ self.set_reactive(API.ws_password, ws_password)
+ self.set_reactive(API.ws_timeout, ws_timeout)
+ self.connect(host=ws_host, port=ws_port, password=ws_password, timeout=ws_timeout)
def on_mount(self):
- self.reset_timer(self.interval)
-
-
- def _try_request(self, request:functools.partial):
- if type(request) is functools.partial:
- name = request.func.__name__
- else:
- name = request.__name__
- logger.info(f"Sending request: {name}")
- try:
- request()
- except obs_error.OBSSDKRequestError as err:
- if err.code in API.KnownReqErr:
- logger.debug(f"Error handled: {API.KnownReqErr(err.code)}")
- else:
- raise
-
-
+ self._reset_timer(self.interval)
+
+
def on_stream_state_changed(self: "API", data):
self.post_message(API.StreamStateUpdated(self, data))
def watch_interval(self, value):
- self.reset_timer(value)
-
-
- def reset_timer(self, interval):
- self._timer: Timer = self.set_interval(interval, self.send_report)
-
-
- def send_report(self):
- self.post_message(API.Report(self, self.get_obs_stats()))
+ self._reset_timer(value)
def stop_stream(self):
+ """Requests to stop streaming."""
self._try_request(self.request_client.stop_stream)
def start_stream(self):
+ """Requests to start streaming."""
self._try_request(self.request_client.start_stream)
def start_record(self):
+ """Requests to start recording."""
self._try_request(self.request_client.start_record)
def stop_record(self):
+ """Requests to stop recording."""
self._try_request(self.request_client.stop_record)
def start_pause(self):
+ """Requests to pause recording."""
self._try_request(self.request_client.pause_record)
def stop_pause(self):
+ """Requests to resume recording."""
self._try_request(self.request_client.resume_record)
+ # LOGGER
def set_record_path(self, path: str):
+ """Sets OBS's record path to path, if the path exists."""
if Path(path).exists():
self._try_request(partial(self.request_client.set_record_directory, path))
else:
def get_obs_stats(self) -> OBSStats:
- record_status = self.request_client.get_record_status()
- stream_status = self.request_client.get_stream_status()
- obs_stats = self.request_client.get_stats()
- record_directory = self.request_client.get_record_directory()
-
+ if self.is_connected:
+ try:
+ record_status = self._try_request(self.request_client.get_record_status)
+ stream_status = self._try_request(self.request_client.get_stream_status)
+ obs_stats = self._try_request(self.request_client.get_stats)
+ record_directory = self._try_request(self.request_client.get_record_directory)
+ except AttributeError:
+ return OBSStats()
+ else:
+ return OBSStats()
+
output = OBSStats(
+ obsws_connection_active = True, # being here implies that it's true.
record_output_active=record_status.output_active,
record_output_paused=record_status.output_paused,
record_output_bytes=record_status.output_bytes,
record_directory=record_directory.record_directory,
)
- return output
\ No newline at end of file
+ return output
+
+
+ # LOGGER
+ def connect(self, host=None, port=None, password=None, timeout=None):
+ host = self.ws_host if host is None else host
+ port = self.ws_port if port is None else port
+ password = self.ws_password if password is None else password
+ timeout = self.ws_timeout if timeout is None else timeout
+ try:
+ logger.info(f"Connecting to OBS WebSocket at {host}:{port}")
+ self.request_client = obs.ReqClient(
+ host=host,
+ port=port,
+ password=password,
+ timeout=timeout
+ )
+ self.event_client = obs.EventClient(
+ host=host,
+ port=port,
+ password=password,
+ timeout=timeout
+ )
+ except ConnectionRefusedError as err:
+ logger.debug(err)
+ else:
+ self.is_connected = True
+ self.event_client.callback.register(self.on_stream_state_changed)
+ self.event_client.callback.register(self.on_record_state_changed)
+
+
+ def disconnect(self):
+ logger.info("Disconnecting from OBS WebSocket")
+ self.request_client = None
+ self.event_client = None
+ self.is_connected = False
+
+
+ # Regarding the pattern in _try_request --
+ # number = 5
+ # while number > 0
+ # try:
+ # raise Error
+ # except Error:
+ # continue
+ # finally:
+ # number -= 1
+ #
+ # Tested; finally is reached. No infinite while loop.
+ # LOGGER
+ def _try_request(self, request:functools.partial):
+ if type(request) is functools.partial:
+ name = request.func.__name__
+ else:
+ name = request.__name__
+ logger.info(f"Sending request: {name}")
+
+ output = None
+ _sentry = 3
+ handled: bool = False
+ while not handled and _sentry >= 0:
+ try:
+ output = request()
+ handled = True
+ except obs_error.OBSSDKRequestError as err:
+ if err.code == API.KnownReqErr.NOT_READY:
+ logger.info(f"OBS WebSocket not ready; trying again in 3 seconds.")
+ logger.debug(f"At most {_sentry} more tries.")
+
+ time.sleep(TRY_AGAIN)
+ continue
+ elif err.code in API.KnownReqErr:
+ logger.debug(f"Error handled: {API.KnownReqErr(err.code)}")
+
+ handled = True
+ else:
+ raise
+ except Exception as err:
+ logger.error(err)
+
+ self.disconnect()
+ handled = True
+ finally:
+ _sentry -= 1
+ if _sentry < 0:
+ logger.info(f"Maximum tries attempted ({TRY_AGAIN}). Aborting request {name}.")
+ return output
+
+
+ # untested
+ def _reset_timer(self, interval):
+ self._timer: Timer = self.set_interval(interval, self._send_report)
+
+
+ def _send_report(self):
+ data = self.get_obs_stats()
+ self.post_message(API.Report(self, data))
\ No newline at end of file