From: sbkelley Date: Tue, 2 Dec 2025 17:55:35 +0000 (-0500) Subject: Further reorganizing, graceful connection failures, connect/disconnect buttons X-Git-Url: https://skyeroc.xyz/gitweb/?a=commitdiff_plain;h=2c76dea4bf0aa8bcb581b01d1b11788b530fdf34;p=obs-ctl Further reorganizing, graceful connection failures, connect/disconnect buttons --- diff --git a/src/obs_ctl/__init__.py b/src/obs_ctl/__init__.py index 530c03f..205d6ca 100644 --- a/src/obs_ctl/__init__.py +++ b/src/obs_ctl/__init__.py @@ -2,18 +2,16 @@ import importlib.resources 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): @@ -32,81 +30,76 @@ class OBSCtlApp(App): 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(): diff --git a/src/obs_ctl/data/layout.tcss b/src/obs_ctl/data/layout.tcss index e0dbb18..bf3fa6d 100644 --- a/src/obs_ctl/data/layout.tcss +++ b/src/obs_ctl/data/layout.tcss @@ -38,6 +38,7 @@ RecordControls { } +/* all necessary */ BoolControl { background:$surface; border:blank; @@ -45,10 +46,12 @@ BoolControl { align:center top; } +StatsDisplay { + height:8; +} + TextReadout { - width:1fr; border:solid $primary; - background:$surface; .text_readout--label { height:1; align:center middle; diff --git a/src/obs_ctl/interface/__init__.py b/src/obs_ctl/interface/__init__.py index 9d6c3e8..fc2366a 100644 --- a/src/obs_ctl/interface/__init__.py +++ b/src/obs_ctl/interface/__init__.py @@ -1 +1,2 @@ -from .molecule import RecordControls, OBSSettingsPane \ No newline at end of file +"""Provides RecordControls, OBSSettingsPane, and StatsDisplay to obs-ctl""" +from .molecule import RecordControls, OBSSettingsPane, StatsDisplay \ No newline at end of file diff --git a/src/obs_ctl/interface/element/labeled_input.py b/src/obs_ctl/interface/element/labeled_input.py new file mode 100644 index 0000000..e69de29 diff --git a/src/obs_ctl/interface/element/text_readout.py b/src/obs_ctl/interface/element/text_readout.py index 55ede6b..d42f28d 100644 --- a/src/obs_ctl/interface/element/text_readout.py +++ b/src/obs_ctl/interface/element/text_readout.py @@ -3,6 +3,12 @@ from textual.reactive import reactive from textual.widgets import Label class TextReadout(Container): + DEFAULT_CSS = """ + TextReadout { + width:1fr; + background:$surface; + } + """ COMPONENT_CLASSES = { "text_readout--label", "text_readout--text", diff --git a/src/obs_ctl/interface/molecule/__init__.py b/src/obs_ctl/interface/molecule/__init__.py index 8271c2f..11f77d9 100644 --- a/src/obs_ctl/interface/molecule/__init__.py +++ b/src/obs_ctl/interface/molecule/__init__.py @@ -1,2 +1,3 @@ from .record_controls import RecordControls -from .settings_pane import OBSSettingsPane \ No newline at end of file +from .settings_pane import OBSSettingsPane +from .stats_display import StatsDisplay \ No newline at end of file diff --git a/src/obs_ctl/interface/molecule/stats_display.py b/src/obs_ctl/interface/molecule/stats_display.py new file mode 100644 index 0000000..8a9a938 --- /dev/null +++ b/src/obs_ctl/interface/molecule/stats_display.py @@ -0,0 +1,68 @@ +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 diff --git a/src/obs_ctl/interface/obs.py b/src/obs_ctl/interface/obs.py index 56de47a..870a504 100644 --- a/src/obs_ctl/interface/obs.py +++ b/src/obs_ctl/interface/obs.py @@ -4,6 +4,7 @@ from functools import partial import functools from pathlib import Path from time import ctime +import time import obsws_python as obs import obsws_python.error as obs_error @@ -18,29 +19,31 @@ MS_TO_HOURS = 1 / 1000 / 60 / 60 """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): @@ -63,14 +66,28 @@ class API(Static): 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" @@ -93,6 +110,7 @@ class API(Static): OUTPUT_PAUSED = 502 OUTPUT_NOT_PAUSED = 503 + # MESSAGES class APIOutputMSGBase(MSGBase): state = None def __init__(self, control, data): @@ -101,11 +119,19 @@ class API(Static): 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. @@ -113,37 +139,37 @@ class API(Static): 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)) @@ -153,42 +179,42 @@ class API(Static): 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: @@ -196,12 +222,19 @@ class API(Static): 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, @@ -223,4 +256,100 @@ class API(Static): 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