From 69e4c3e6d0b098f5dcff595f1cb315a9f7a609c6 Mon Sep 17 00:00:00 2001 From: RioPlay Date: Sun, 13 Sep 2026 09:37:02 -0500 Subject: [PATCH] Add explicit OBS session arming and lifecycle guards --- docs/desktop-roadmap.md | 22 +- docs/plans/active/obs-audio-implementation.md | 24 +- docs/plans/active/obs-session-arm.md | 160 +++++ docs/plans/active/windows-obs-audio-pipe.md | 14 +- docs/workspaces.md | 4 +- native/obs-plugin/README.md | 10 +- native/obs-plugin/src/admission.c | 181 +++++ native/obs-plugin/src/admission.h | 20 + native/obs-plugin/src/bridge.c | 84 ++- native/obs-plugin/src/plugin_state.c | 267 ++++++- native/obs-plugin/src/plugin_state.h | 21 +- native/obs-plugin/src/session_protocol.c | 45 ++ native/obs-plugin/src/session_protocol.h | 22 + .../tests/admission_io_fault_test.c | 189 +++++ native/obs-plugin/tests/admission_io_test.c | 672 ++++++++++++++++++ native/obs-plugin/tests/bridge_test.c | 132 +++- native/obs-plugin/tests/plugin_state_test.c | 443 +++++++++++- .../obs-plugin/tests/session_protocol_test.c | 49 ++ native/obs-plugin/tools/build.py | 2 +- native/obs-plugin/tools/test_native.py | 24 +- tests/test_obs_audio_arm.py | 163 +++++ tests/test_obs_audio_pipe.py | 40 +- tests/windows_pipe_server.py | 9 + utterleaf/obs_audio_pipe.py | 37 + 24 files changed, 2583 insertions(+), 51 deletions(-) create mode 100644 docs/plans/active/obs-session-arm.md create mode 100644 native/obs-plugin/src/session_protocol.c create mode 100644 native/obs-plugin/src/session_protocol.h create mode 100644 native/obs-plugin/tests/admission_io_fault_test.c create mode 100644 native/obs-plugin/tests/admission_io_test.c create mode 100644 native/obs-plugin/tests/session_protocol_test.c create mode 100644 tests/test_obs_audio_arm.py diff --git a/docs/desktop-roadmap.md b/docs/desktop-roadmap.md index d0d68120..05e621ad 100644 --- a/docs/desktop-roadmap.md +++ b/docs/desktop-roadmap.md @@ -150,9 +150,9 @@ binary is included in desktop or Android releases. The [control dependency record](desktop-obs-control-resource.md) records the pinned library, reviewed full license and development-wheel provenance. -Continue on `feat/obs-enrollment-flow` with frontend acceptance and the desktop -connection/controller after the pairing setup increment. Arm/PCM and live -recognition remain later gates. No live-capture app entry point, audio endpoint, plugin binary +Pairing setup remains in draft PR #36 on `feat/obs-enrollment-flow`; the separate +Arm increment continues on `feat/obs-session-arm`. The desktop controller, PCM and +live recognition remain later gates. No live-capture app entry point, audio endpoint, plugin binary publication or OBS capture integration exists yet. The linked native pairing Tools flow, exclusive per-user owner and strict @@ -178,10 +178,22 @@ pass after the final status-card adjustment. Native TaskDialog activation, Escape, modal cleanup and zero-mutation checks also pass in an isolated fixture; all 29 native driver commands pass with `--ui`. Normal/compact/enlarged desktop renders were inspected. Independent final desktop setup review is clear, with -29 focused and 187 related tests rerun on the final source. Its own CI remains -pending; no existing release changes. Earlier linked checkpoint `15e1c77` passed +29 focused and 187 related tests rerun on the final source. All five exact-source +CI jobs pass at `e71b627` in [34760760046](https://github.com/RioPlay/utterleaf/actions/runs/34760760046). +No existing release changes. Earlier linked checkpoint `15e1c77` passed all five desktop CI jobs in [34758907354](https://github.com/RioPlay/utterleaf/actions/runs/34758907354). +The separate [Arm integration](plans/active/obs-session-arm.md) is active on +`feat/obs-session-arm`. It adds a fixed authenticated-pipe command, guarded +frontend consent, bounded native I/O and an explicit desktop Arm operation. +The 199-test affected desktop bundle and all 34 native build/test commands pass, +along with the linked build and headless refusal smoke. Independent source review +is clear. Startup/started/stopping states remain busy even when public activity +flags are false; a failed start without STOPPED requires a later actual STOPPED +or OBS restart before fresh Arm. Visible failure/recovery controls and real OBS +acceptance remain pending with the controller/PCM work. This is source work, +not live OBS support or a new release. + Current integrated desktop regression: **1,432 passed, 13 skipped in 39.47 seconds**, including the reviewed OBS control, native identity, pipe and audio handshake. The combined focused OBS/pipe/privacy/configuration/boundary/packaging bundle diff --git a/docs/plans/active/obs-audio-implementation.md b/docs/plans/active/obs-audio-implementation.md index 4a080edc..f43b3359 100644 --- a/docs/plans/active/obs-audio-implementation.md +++ b/docs/plans/active/obs-audio-implementation.md @@ -11,19 +11,25 @@ in source; its physical microphone and release gates remain separate. ## Current increment -The [enrollment plan](obs-native-enrollment.md) now records the integrated private -pairing stores and the typed desktop Issue/Prepare adapter on -`feat/obs-enrollment-flow`. Its focused regression passes 341 tests with no skips, -including 66 enrollment cases. Native owner, Tools and vendor dispatch are still -pending; no app entry point or capture activation is exposed. +The [enrollment plan](obs-native-enrollment.md) records the integrated private +pairing stores, native owner/Tools/vendor flow and typed desktop Issue/Prepare +adapter. The Windows [pairing setup](obs-desktop-pairing-ui.md) at `e71b627` passed +all five desktop CI jobs in run 34760760046 and remains in draft PR #36. +The separate [explicit Arm increment](obs-session-arm.md) on +`feat/obs-session-arm` passes 199 desktop tests, 34 native build/test commands, +the linked build and headless smoke, with independent source review. Its +conservative frontend guard refuses fresh Arm after a failed start without +STOPPED until a later actual STOPPED or OBS restart. Controller guidance and real +frontend recovery remain acceptance gates. No live-capture app entry or audio +activation is exposed. The bounded audio protocol, consent/session receiver and read-only loopback WebSocket control client are reviewed. The Windows TCP peer identity gate is now implemented and integrated, as recorded below. The separate original C module now -has a verified inert build and `libobs` load/unload prerequisite; this does not -establish the native server, vendor registration, PCM callback, Arm, or real OBS -application workflow. Continue with the original native plugin/server and its -authenticated audio endpoint. Root owns the integrated source, +has a verified build and headless-refusal prerequisite, linked native enrollment +and authenticated pipe admission. These do not establish actual OBS frontend +load, the PCM callback or the complete application workflow. Continue the +original mix callback and visible controller. Root owns the integrated source, tests and this plan after delegated handoff; assign each next component one writer and a separate reviewer under the repository's ownership rules. No runtime/UI entry is exposed until the authenticated transport can satisfy the diff --git a/docs/plans/active/obs-session-arm.md b/docs/plans/active/obs-session-arm.md new file mode 100644 index 00000000..0a253f9d --- /dev/null +++ b/docs/plans/active/obs-session-arm.md @@ -0,0 +1,160 @@ +# Explicit OBS session arming + +September 13, 2026. Continue on `feat/obs-session-arm`, stacked on reviewed +pairing setup `e71b627` in draft PR #36. That exact setup source passed all five +desktop jobs in [CI 34760760046](https://github.com/RioPlay/utterleaf/actions/runs/34760760046). +This plan implements the next part of the full [live audio contract](obs-audio-design.md). +The published previews remain unchanged. + +## Goal and area + +After explicit enrollment and pipe authentication, require a separate one-use +Arm transaction. Admit it only while OBS is idle. Only a later ordered frontend +STARTING then STARTED may authorize that session's future audio delivery. Keep +the authenticated connection alive while armed, detect client departure, and +retire it on disarm, stop, cancellation, revocation or exit. Connecting or finding +an existing stream must never grant capture consent. + +Area: native authenticated I/O, bounded session command codec, plugin worker and +frontend state coordination; desktop `ObsAudioPipe` explicit Arm operation; +focused native/desktop interop and race tests. The connection controller and PCM +callback/recognition remain required follow-ups, not replacements for this scope. + +## Constraints + +- Keep enrollment proof and Hello/ACK unchanged. Arm travels only inside the + already authenticated duplex pipe, never vendor JSON. Keep the prepared session + ID and requested additional mix mask until worker teardown. +- Keep unauthenticated and authenticated-but-unarmed deadlines bounded at 15 + seconds each. Armed waiting has no arbitrary recording-duration countdown. + Retain one joined worker, one request slot and bounded wire/I/O buffers. +- Use the existing retained process/pipe identity and cancellation/drain logic + before and after native I/O. Never allocate from untrusted wire lengths. +- Serialize Arm's OBS idle observation with frontend state transitions. Do not + block the frontend on a worker that needs frontend work. Queued callbacks must + survive late invocation after EXIT without touching retired heap state or OBS. +- Success acknowledges a committed Arm, not successful capture. Lost replies and + errors are terminal; no automatic retry/reconnect/rearm. Session IDs are routing + metadata and do not independently authenticate anything. +- No OBS installation, consumer profile or microphone changes during fixtures. + Original native GPL source remains separate from Apache desktop code and Android. + +## Acceptance and verification + +The fixed Arm record is 28 bytes, `<4sBBH16sBBH`: `ULAC`, version 1, kind +1=request or 2=reply, zero header reserved, the prepared 16-byte session ID, +prepared additional mix mask (0–63), status, zero trailing reserved. Request +status is 0; reply status is 1=committed or 2=refused. Only the exact prepared +session/mask may be used. A reply precedes any future ULAP frame on that pipe. +An unrecognized or malformed command terminates the attempt. + +- Native exact I/O: authenticated roundtrip, fragmentation, wrong/preauth peer, + timeout, cancellation, EOF and failed-read wiping using actual child processes. +- Command framing: exact lengths, magic/version/kind/session/mask/reserved/status + checks; malformed, duplicate and replayed commands terminate. +- Consent ordering: refuse already-active/starting streams; STARTING before Arm + commit refuses, Arm before STARTING then STARTED authorizes only once; no Start + or PCM before that transition. Failed OBS start must not manufacture STARTED. +- Lifetime: client close, pending UI task, stop, replacement, forget and EXIT + races cancel/drain/join; late callbacks cannot access OBS or freed runtime. +- Desktop Arm sends once, checks its bound reply, preserves coalesced future PCM, + and closes on error/cancel; concurrent close must unblock it. +- Run focused native fixture drivers and affected Python OBS/pipe tests, then + native integration build/test checks and independent source/evidence review. + Record exact commands and results here as implemented. Fixture results do not + establish real OBS frontend/audio load or physical-device usability. + +## Non-goals and stop + +Do not add an alternate recorder, change OBS routing, silently attach to existing +streams, publish a plugin, or treat this as completed live transcription. Stop +editing the Arm increment when its observable checks and independent review pass; +continue the actual mix callback, controller, live recognition and acceptance +gates required by the desktop roadmap. + +## OBS lifecycle boundary + +The pinned OBS 32.2.2 +[streaming frontend](https://raw.githubusercontent.com/obsproject/obs-studio/ba2f32bdf791005443988a4955e963663e16b1ed/frontend/widgets/OBSBasic_Streaming.cpp) +emits STARTING before calling the output start operation and STARTED later. Its +synchronous start-error path emits a private Qt stopped signal without necessarily +a public frontend STOPPED event. Once STARTING is observed, the bridge keeps its +frontend busy guard set through STARTED and STOPPING, and clears it only on +STOPPED. The Arm check also requires both current activity queries to be false. + +`obs_output_active()` covers active/reconnecting outputs, so false alone does +not rule out an asynchronous connection attempt. The pinned +[libobs output implementation](https://raw.githubusercontent.com/obsproject/obs-studio/ba2f32bdf791005443988a4955e963663e16b1ed/libobs/obs-output.c) +does not provide a public completion event for a synchronously rejected start. +A later queued UI callback is not proof that the start call has returned: nested +frontend event processing can run queued work early. An absent signal or elapsed +grace period therefore cannot safely establish idle. A 30-second deadline +for STARTING to reach STARTED terminally disarms an orphaned/slow startup; it +does not establish that OBS became idle. Waiting while armed and recording after +STARTED have no duration cutoff. + +Known limit: after a failed start that omits STOPPED, a fresh Arm remains refused +until OBS emits STOPPED in a later actual lifecycle or OBS restarts. Utterleaf +does not initiate that lifecycle or change stream settings. Clear user-facing +explanation and real OBS failure/recovery acceptance remain requirements for the +pending desktop controller. These component fixtures do not establish that +experience. If EXIT is missed, unload closes native gates without calling OBS. + +## Component verification + +The final native driver passes **34 build/test commands**, including actual +child-process exact I/O, injected API failures, fixed commands, Windows runtime +thread/event ordering, queued frontend lifecycle and real-libobs vendor dispatch. +Exact-I/O cases cover fragmentation, bounds, partial-read wiping, EOF, typed +cancellation and blocked-write draining. Injected event-creation, zero/over-count +transfer and failed probe cases cannot report success; cancellation or deadline +expiry during the final peer check wipes the read and wins over success. + +The runtime fixture covers late generations, accepted/refused Arm, STARTING before +the queued callback, startup expiry, late STARTED, and cancellation through +revocation/close. Bridge fixtures refuse startup/started/stopping states even +when activity flags are false, and prove stale/duplicate/closed UI tasks cannot +touch retired state. The explicit desktop Arm and affected +protocol/session/pipe/privacy/boundary bundle passes **199 tests in 5.50 seconds**, +with no skips. Separate source review cleared the final codec/runtime and +conservative bridge gate; root reviewed the delegated I/O and desktop changes. + +The linked x64 DLL builds against the pinned headers and installed OBS runtime. +Headless smoke opens it, refuses initialization before any pairing-store startup, +returns from shutdown, and leaves fixed-store metadata unchanged. The DLL SHA-256 +is `9ce5f225fef16af2ccdcdbf5e30ea637df5a306863e0285d8329e33f1e7925a3`. +Three existing linker warnings remain confined to the test heap shim. The native +dialog fixture was not repeated for this increment; its earlier setup evidence +remains separate. No OBS application, source or audio device ran in these checks. +Exact-source CI for this Arm checkpoint remains pending. + +From the owning worktree, the desktop command was: + +```powershell +$env:PYTHONPATH = (Get-Location).Path +& "C:\Users\unknown\Projects\Mindict\.venv\Scripts\python.exe" -m pytest tests/test_obs_audio_pipe.py tests/test_obs_audio_arm.py tests/test_obs_session.py tests/test_obs_protocol.py tests/test_windows_pipe.py tests/test_privacy.py tests/test_repo_boundaries.py -o addopts='' -q +``` + +Local native evidence is under `.grok/obs-session-arm/`; the separate exact-I/O +development receipts under `.grok/obs-enrollment-flow/` are historical. The +authoritative final receipts are `build/build-receipt.json`, +`build/smoke-receipt.json`, and `verification/test-receipt.json` under that Arm +directory. Commands from the owning worktree: + +```powershell +$py = "C:\Users\unknown\Projects\Mindict\.venv\Scripts\python.exe" +$toolchain = "C:\Users\unknown\.local\llvm-mingw-20260616-ucrt-x86_64" +$obsBin = "C:\Program Files\obs-studio\bin\64bit" +$headers = "C:\Users\unknown\Projects\Mindict\.grok\obs-native-build\headers" +$build = "C:\Users\unknown\Projects\Mindict\.grok\obs-session-arm\build" +$verification = "C:\Users\unknown\Projects\Mindict\.grok\obs-session-arm\verification" +& $py native/obs-plugin/tools/build.py --toolchain $toolchain --obs-bin $obsBin --headers $headers --output $build +& $py native/obs-plugin/tools/smoke.py --build $build +& $py native/obs-plugin/tools/test_native.py --toolchain $toolchain --output $verification --build $build --headers $headers +``` + +Continue the visible controller and actual PCM path. Future Start/PCM writes must +remain serialized after the complete successful Arm reply, even if native state +already reached STARTING/STARTED while that reply was in flight. This increment +emits no frames. Real frontend recovery, audio load, redistribution/toolchain +review and release acceptance remain open. diff --git a/docs/plans/active/windows-obs-audio-pipe.md b/docs/plans/active/windows-obs-audio-pipe.md index 71207a51..9cf79dc9 100644 --- a/docs/plans/active/windows-obs-audio-pipe.md +++ b/docs/plans/active/windows-obs-audio-pipe.md @@ -6,6 +6,16 @@ Follows the reviewed [TCP peer identity gate](windows-obs-peer-identity.md) and [OBS audio design](obs-audio-design.md). +September 13 follow-up: [explicit session arming](obs-session-arm.md) adds a +separate `ObsAudioPipe.arm()` transaction after authentication. The caller supplies +the prepared additional mix mask and deadline; the exact matching native reply +must commit Arm before `read_frames()` accepts ULAP. Calling Arm twice, reading +before Arm, refusal, malformed replies, cancellation or identity loss closes the +connection. Coalesced bytes following the fixed reply remain available to the +frame reader. This is internal source behavior; the application controller and +live OBS audio workflow are still pending. Evidence below describes the earlier +handshake/receiver increment unless explicitly linked to that follow-up. + ## Goal and area Implement the local client side of the original OBS audio bridge: bind one named @@ -13,7 +23,7 @@ pipe to the already verified OBS process before exchanging session credentials, then receive bounded binary audio for the existing protocol/session receiver. The audio connection must retain its own process identity so a degraded control connection cannot silently replace or invalidate an otherwise healthy active -audio session. This component does not itself arm or record. +audio session. Authentication alone does not arm or record. The source is `windows_pipe.py`, `obs_audio_pipe.py` and the process-lease additions to `windows_peer_identity.py`, `obs_websocket.py` and `obs_control.py`, all under @@ -85,7 +95,7 @@ its first read and consumes a session once. The client never reconnects a failed session. Mutable secret/Hello construction buffers are cleared best-effort; Python/HMAC copies cannot be guaranteed erased. -After ACK, each bounded read rechecks server PID and the independent process +After ACK and the explicit Arm reply, each bounded read rechecks server PID and the independent process before and after I/O. Bytes feed the existing ULAP decoder; every frame must carry the same session ID. End must be the last frame, with no trailing partial frame. Consent, post-arm start ordering, bus selection and temporary stores remain the diff --git a/docs/workspaces.md b/docs/workspaces.md index 1a4ecdcd..05af2e87 100644 --- a/docs/workspaces.md +++ b/docs/workspaces.md @@ -9,7 +9,7 @@ tree, one desktop OBS tree, and published/unpublished release receipts. | Stream | Checkout | Next PR | | --- | --- | --- | | Android keyboard | `android-keyboard-hardening` on `main` | Signed alpha18 candidate. Phone/TalkBack/landscape stay Ernest. No new keyboard features in this slice. | -| Desktop OBS | `desktop-obs-bridge` | Land draft stack from the base: 36 enrollment → 37 arm → 38 PCM → 39 disarm → 40 controller → 41 provenance → 42 timelines. Rebase each onto current `main`/parent before undrafting. | +| Desktop OBS | `desktop-obs-bridge` | PR #36 merged. Next: draft [PR #37](https://github.com/RioPlay/utterleaf/pull/37) arm, then 38 PCM → 39 disarm → 40 controller → 41 provenance → 42 timelines. | | Android CI split | same Android tree, later | Follow-up only: stop downloading speech models and running the full emulator suite on every keyboard PR. | Do not mix these in one branch. Desktop CI already skips Android-only paths. @@ -17,7 +17,7 @@ Do not mix these in one branch. Desktop CI already skips Android-only paths. | Branch | Purpose | Next work | | --- | --- | --- | | `main` | Current Android keyboard checkout (`android-keyboard-hardening` worktree) | [Alpha17 polish](plans/active/android-alpha17-polish.md) is merged; signed candidate and phone acceptance remain | -| `feat/obs-enrollment-flow` | Desktop OBS stack base (`desktop-obs-bridge` worktree); draft [PR #36](https://github.com/RioPlay/utterleaf/pull/36) | Rebase onto current `main`, then pairing UI/vendor requests. Later stack PRs 37–42 stay parked. | +| `feat/obs-session-arm` | Desktop OBS stack (`desktop-obs-bridge` worktree); draft [PR #37](https://github.com/RioPlay/utterleaf/pull/37) | Arm/lifecycle guards on current `main`. PRs 38–42 stay parked until this lands. | | `checkpoint/mixed-work-20260912` | Preserved mixed development snapshot; not a release or PR | Recovery/reference only; leave the original source environment intact | | `release/desktop-0.4.6rc2` | Published Windows x64 CPU prerelease | Preserve the immutable RC2 tag and release evidence | | `release/desktop-0.4.6rc1` | Unpublished RC1 candidate retained for audit | Reference only; tag and downloaded artifact unchanged | diff --git a/native/obs-plugin/README.md b/native/obs-plugin/README.md index d4031788..502ff49d 100644 --- a/native/obs-plugin/README.md +++ b/native/obs-plugin/README.md @@ -17,8 +17,14 @@ storage, cancellation and gap handling, visible control, local recognition, and live OBS acceptance. Separate native Hello/ACK and retained-process pipe admission components now have an explicit test build, described below, and the bounded admission worker is linked into the module. Frontend/device behavior -and real OBS acceptance remain unverified; stream-following, PCM, controller, -Arm and live capture remain unimplemented. +and real OBS acceptance remain unverified. The separate +[Arm increment](../../docs/plans/active/obs-session-arm.md) now implements fixed +authenticated-pipe commands, bounded exact I/O and native/desktop Arm state in +reviewed development source. Its 199 desktop tests, 34 native build/test commands, +linked build and headless smoke pass. Startup remains busy until a real STOPPED +event: a synchronous OBS failure lacking that event can require a later completed +stream lifecycle or OBS restart before fresh Arm. PCM, the visible capture +controller, live recognition and actual frontend acceptance remain open. ## Inputs and legal boundary diff --git a/native/obs-plugin/src/admission.c b/native/obs-plugin/src/admission.c index 6a3e46df..62e57b28 100644 --- a/native/obs-plugin/src/admission.c +++ b/native/obs-plugin/src/admission.c @@ -15,6 +15,33 @@ #define UL_IO_POLL_MS 50u #define UL_PIPE_BUFFER_BYTES 4096u +#ifndef UL_ADMISSION_IO_EVENT_CREATE +#define UL_ADMISSION_IO_EVENT_CREATE() CreateEventW(NULL, TRUE, FALSE, NULL) +#endif + +#ifndef UL_ADMISSION_PEER_EXPECTED +#define UL_ADMISSION_PEER_EXPECTED(admission) peer_is_expected(admission) +#endif + +#ifndef UL_ADMISSION_READ_PIPE +#define UL_ADMISSION_READ_PIPE(admission, overlapped, buffer, size, started, \ + timeout_ms, transferred) \ + read_pipe(admission, overlapped, buffer, size, started, timeout_ms, \ + transferred) +#endif + +#ifndef UL_ADMISSION_WRITE_PIPE +#define UL_ADMISSION_WRITE_PIPE(admission, overlapped, buffer, size, started, \ + timeout_ms, transferred) \ + write_pipe(admission, overlapped, buffer, size, started, timeout_ms, \ + transferred) +#endif + +#ifndef UL_ADMISSION_PEEK_PIPE +#define UL_ADMISSION_PEEK_PIPE(pipe, available) \ + PeekNamedPipe(pipe, NULL, 0, NULL, available, NULL) +#endif + typedef struct ul_identity { PSID user_sid; DWORD user_sid_size; @@ -520,6 +547,122 @@ static int operation_to_public(ul_operation_result result) return UL_ADMISSION_IO_ERROR; } +static int admission_gate_result(const ul_admission *admission) +{ + if (admission == NULL || admission->pipe == INVALID_HANDLE_VALUE) + return UL_ADMISSION_IO_ERROR; + if (InterlockedCompareExchange( + (volatile LONG *)&admission->cancelled, 0, 0) != 0) + return UL_ADMISSION_CANCELLED; + if (InterlockedCompareExchange( + (volatile LONG *)&admission->authenticated, 0, 0) == 0) + return UL_ADMISSION_REJECTED; + return UL_ADMISSION_AUTH_OK; +} + +static int terminal_failure(ul_admission *admission, int result) +{ + if (admission != NULL) { + ul_admission_cancel(admission); + if (admission->pipe != INVALID_HANDLE_VALUE) + DisconnectNamedPipe(admission->pipe); + } + return result; +} + +static int admission_peer_result(ul_admission *admission, ULONGLONG started, + DWORD timeout_ms) +{ + DWORD remaining = 0; + BOOL expected; + ul_operation_result operation = + progress_result(admission, started, timeout_ms, &remaining); + + if (operation != UL_OPERATION_OK) + return operation_to_public(operation); + expected = UL_ADMISSION_PEER_EXPECTED(admission); + operation = progress_result(admission, started, timeout_ms, &remaining); + if (operation != UL_OPERATION_OK) + return operation_to_public(operation); + return expected ? UL_ADMISSION_AUTH_OK : UL_ADMISSION_REJECTED; +} + +static int admission_exact_io(ul_admission *admission, void *buffer, + DWORD size, DWORD timeout_ms, BOOL writing) +{ + HANDLE io_event = NULL; + OVERLAPPED overlapped; + ULONGLONG started; + DWORD offset = 0; + int result = UL_ADMISSION_IO_ERROR; + + if (buffer == NULL || size == 0u || + size > UL_ADMISSION_MAX_IO_BYTES || timeout_ms == 0u || + timeout_ms > 30000u) + goto cleanup; + started = GetTickCount64(); + result = admission_gate_result(admission); + if (result != UL_ADMISSION_AUTH_OK) + goto cleanup; + result = admission_peer_result(admission, started, timeout_ms); + if (result != UL_ADMISSION_AUTH_OK) + goto cleanup; + result = UL_ADMISSION_IO_ERROR; + io_event = UL_ADMISSION_IO_EVENT_CREATE(); + if (io_event == NULL) + goto cleanup; + SecureZeroMemory(&overlapped, sizeof(overlapped)); + overlapped.hEvent = io_event; + + while (offset < size) { + DWORD transferred = 0; + ul_operation_result operation; + int peer_result = + admission_peer_result(admission, started, timeout_ms); + + if (peer_result != UL_ADMISSION_AUTH_OK) { + result = peer_result; + goto cleanup; + } + + if (writing) + operation = UL_ADMISSION_WRITE_PIPE( + admission, &overlapped, (const uint8_t *)buffer + offset, + size - offset, started, timeout_ms, &transferred); + else + operation = UL_ADMISSION_READ_PIPE( + admission, &overlapped, (uint8_t *)buffer + offset, + size - offset, started, timeout_ms, &transferred); + if (operation != UL_OPERATION_OK) { + result = operation_to_public(operation); + goto cleanup; + } + peer_result = admission_peer_result(admission, started, timeout_ms); + if (peer_result != UL_ADMISSION_AUTH_OK) { + result = peer_result; + goto cleanup; + } + if (transferred == 0u || transferred > size - offset) + goto cleanup; + offset += transferred; + } + result = admission_peer_result(admission, started, timeout_ms); + if (result != UL_ADMISSION_AUTH_OK) + goto cleanup; + +cleanup: + SecureZeroMemory(&overlapped, sizeof(overlapped)); + if (io_event != NULL) + CloseHandle(io_event); + if (result != UL_ADMISSION_AUTH_OK) { + if (!writing && buffer != NULL && size > 0u && + size <= UL_ADMISSION_MAX_IO_BYTES) + SecureZeroMemory(buffer, size); + return terminal_failure(admission, result); + } + return result; +} + ul_admission *ul_admission_create(DWORD expected_pid, const uint8_t session[16]) { @@ -716,6 +859,44 @@ int ul_admission_authenticate(ul_admission *admission, DWORD timeout_ms) return result; } +int ul_admission_read_exact(ul_admission *admission, void *buffer, DWORD size, + DWORD timeout_ms) +{ + return admission_exact_io(admission, buffer, size, timeout_ms, FALSE); +} + +int ul_admission_write_all(ul_admission *admission, const void *buffer, + DWORD size, DWORD timeout_ms) +{ + return admission_exact_io(admission, (void *)buffer, size, timeout_ms, + TRUE); +} + +int ul_admission_probe(ul_admission *admission, DWORD *available) +{ + DWORD observed = 0; + int result = UL_ADMISSION_IO_ERROR; + + if (available == NULL) + return terminal_failure(admission, result); + *available = 0; + result = admission_gate_result(admission); + if (result != UL_ADMISSION_AUTH_OK) + return terminal_failure(admission, result); + if (!UL_ADMISSION_PEER_EXPECTED(admission)) + return terminal_failure(admission, UL_ADMISSION_REJECTED); + result = UL_ADMISSION_IO_ERROR; + if (!UL_ADMISSION_PEEK_PIPE(admission->pipe, &observed)) + return terminal_failure(admission, result); + if (!UL_ADMISSION_PEER_EXPECTED(admission)) + return terminal_failure(admission, UL_ADMISSION_REJECTED); + result = admission_gate_result(admission); + if (result != UL_ADMISSION_AUTH_OK) + return terminal_failure(admission, result); + *available = observed; + return UL_ADMISSION_AUTH_OK; +} + void ul_admission_cancel(ul_admission *admission) { if (admission == NULL) diff --git a/native/obs-plugin/src/admission.h b/native/obs-plugin/src/admission.h index 45f8c5de..f5bda073 100644 --- a/native/obs-plugin/src/admission.h +++ b/native/obs-plugin/src/admission.h @@ -20,6 +20,8 @@ enum ul_admission_result { UL_ADMISSION_IO_ERROR = 4, }; +#define UL_ADMISSION_MAX_IO_BYTES 65585u + /* * Create a once-only pending admission for an expected process. This must only * be called from an already-authenticated, bounded vendor-request context. @@ -36,6 +38,24 @@ ul_admission *ul_admission_create(DWORD expected_pid, */ int ul_admission_authenticate(ul_admission *admission, DWORD timeout_ms); +/* + * Exact, bounded I/O for the single admission worker after authentication. + * Calls must be serialized by that worker. size must be from 1 through + * UL_ADMISSION_MAX_IO_BYTES and timeout_ms from 1 through 30000. Success is + * UL_ADMISSION_AUTH_OK. Any failure terminally cancels and disconnects the + * admission; a failed read also wipes its complete caller-supplied buffer. + */ +int ul_admission_read_exact(ul_admission *admission, void *buffer, DWORD size, + DWORD timeout_ms); +int ul_admission_write_all(ul_admission *admission, const void *buffer, + DWORD size, DWORD timeout_ms); + +/* + * Non-consuming liveness check for the same serialized worker. Zero available + * bytes is a successful live result. Errors terminally cancel and disconnect. + */ +int ul_admission_probe(ul_admission *admission, DWORD *available); + /* * May run concurrently with authenticate until the caller joins that worker. * Cancellation prevents future pipe borrowing but cannot revoke a handle that diff --git a/native/obs-plugin/src/bridge.c b/native/obs-plugin/src/bridge.c index 34beda9e..b2a89a5b 100644 --- a/native/obs-plugin/src/bridge.c +++ b/native/obs-plugin/src/bridge.c @@ -16,6 +16,59 @@ OBS_DECLARE_MODULE() static obs_websocket_vendor vendor; static bool issue_registered; static bool prepare_registered; +static SRWLOCK frontend_gate = SRWLOCK_INIT; +static bool frontend_open; +static uintptr_t queued_arm; +/* Public active queries can stay false during startup. Once STARTING arrives, + * only STOPPED proves idle; a synchronous failure emits no public STOPPED and + * therefore remains fail closed. */ +static bool stream_busy; + +/* OBS 32.2.2 queues worker UI tasks through Qt. The only task parameter is a + * numeric generation, never a borrowed runtime/event/admission pointer. */ +static void check_arm_on_frontend(void *parameter) +{ + uintptr_t generation = (uintptr_t)parameter; + obs_output_t *output; + bool idle; + AcquireSRWLockExclusive(&frontend_gate); + if (!frontend_open || generation == 0u || queued_arm != generation) { + ReleaseSRWLockExclusive(&frontend_gate); + return; + } + queued_arm = 0u; + idle = !stream_busy && !obs_frontend_streaming_active(); + output = obs_frontend_get_streaming_output(); /* Already a new reference. */ + if (output != NULL) { + idle = idle && !obs_output_active(output); + obs_output_release(output); + } + ul_plugin_arm_checked(generation, idle); + ReleaseSRWLockExclusive(&frontend_gate); +} + +static bool queue_arm(uintptr_t generation) +{ + bool queued = false; + AcquireSRWLockExclusive(&frontend_gate); + if (frontend_open && generation != 0u && queued_arm == 0u) { + queued_arm = generation; + obs_queue_task(OBS_TASK_UI, check_arm_on_frontend, + (void *)generation, false); + queued = true; + } + ReleaseSRWLockExclusive(&frontend_gate); + return queued; +} + +static void close_frontend(void) +{ + AcquireSRWLockExclusive(&frontend_gate); + frontend_open = false; + queued_arm = 0u; + stream_busy = true; + ReleaseSRWLockExclusive(&frontend_gate); +} static void pairing_menu(void *private_data) { @@ -27,8 +80,29 @@ static void pairing_menu(void *private_data) static void frontend_event(enum obs_frontend_event event, void *private_data) { (void)private_data; - if (event != OBS_FRONTEND_EVENT_EXIT) + if (event != OBS_FRONTEND_EVENT_EXIT) { + AcquireSRWLockExclusive(&frontend_gate); + if (frontend_open) { + switch (event) { + case OBS_FRONTEND_EVENT_STREAMING_STARTING: + stream_busy = true; + ul_plugin_stream_event(UL_STREAM_STARTING); break; + case OBS_FRONTEND_EVENT_STREAMING_STARTED: + stream_busy = true; + ul_plugin_stream_event(UL_STREAM_STARTED); break; + case OBS_FRONTEND_EVENT_STREAMING_STOPPING: + stream_busy = true; + ul_plugin_stream_event(UL_STREAM_STOPPING); break; + case OBS_FRONTEND_EVENT_STREAMING_STOPPED: + stream_busy = false; + ul_plugin_stream_event(UL_STREAM_STOPPED); break; + default: break; + } + } + ReleaseSRWLockExclusive(&frontend_gate); return; + } + close_frontend(); ul_vendor_set_enabled(false); ul_plugin_stop_accepting(); if (prepare_registered) @@ -48,6 +122,13 @@ bool obs_module_load(void) return false; if (!ul_plugin_start()) return false; + if (!ul_plugin_set_arm_scheduler(queue_arm)) { + ul_plugin_close(); + return false; + } + AcquireSRWLockExclusive(&frontend_gate); + frontend_open = true; + ReleaseSRWLockExclusive(&frontend_gate); obs_frontend_add_event_callback(frontend_event, NULL); obs_frontend_add_tools_menu_item("Utterleaf pairing...", pairing_menu, NULL); blog(LOG_INFO, "[Utterleaf OBS bridge] pairing controls loaded; recording remains off"); @@ -82,6 +163,7 @@ void obs_module_post_load(void) void obs_module_unload(void) { /* Native-only fallback: frontend/websocket teardown order is not assumed. */ + close_frontend(); ul_vendor_set_enabled(false); ul_plugin_close(); } diff --git a/native/obs-plugin/src/plugin_state.c b/native/obs-plugin/src/plugin_state.c index 19df1b3d..dd72271f 100644 --- a/native/obs-plugin/src/plugin_state.c +++ b/native/obs-plugin/src/plugin_state.c @@ -1,5 +1,6 @@ // SPDX-License-Identifier: GPL-2.0-or-later #include "plugin_state.h" +#include "session_protocol.h" #include @@ -16,8 +17,13 @@ #define UL_PLUGIN_READY_TIMEOUT_MS 15000u #endif +#ifndef UL_PLUGIN_START_TIMEOUT_MS +#define UL_PLUGIN_START_TIMEOUT_MS 30000u +#endif + typedef struct ul_plugin_runtime { SRWLOCK operation_lock; + SRWLOCK session_lock; HANDLE leases_zero; volatile LONG leases; volatile LONG closing; @@ -27,7 +33,16 @@ typedef struct ul_plugin_runtime { ul_authorizer *authorizer; HANDLE worker; HANDLE worker_cancel; + HANDLE arm_result; ul_admission *worker_admission; + ul_prepare_options session_options; + uintptr_t session_generation; + uint64_t stream_epoch; + uint64_t pending_epoch; + ULONGLONG ready_started; + ULONGLONG stream_starting_at; + ul_session_phase session_phase; + ul_arm_scheduler schedule_arm; volatile LONG worker_active; ul_plugin_status status; ul_pairing_result storage_result; @@ -96,6 +111,129 @@ static void mutation_end(ul_plugin_runtime *runtime) InterlockedExchange(&runtime->mutation_active, 0); } +/* session_lock is held; never joins or acquires operation_lock. Admission is + * removed under this same lock before the authorizer can destroy it. */ +static void session_cancel_locked(ul_plugin_runtime *runtime) +{ + runtime->session_phase = UL_SESSION_TERMINAL; + SetEvent(runtime->worker_cancel); + SetEvent(runtime->arm_result); + if (runtime->worker_admission != NULL) + ul_admission_cancel(runtime->worker_admission); +} + +static void session_cancel(ul_plugin_runtime *runtime) +{ + AcquireSRWLockExclusive(&runtime->session_lock); + session_cancel_locked(runtime); + ReleaseSRWLockExclusive(&runtime->session_lock); +} + +static void session_release(ul_plugin_runtime *runtime) +{ + AcquireSRWLockExclusive(&runtime->session_lock); + runtime->worker_admission = NULL; + runtime->session_phase = UL_SESSION_NONE; + SecureZeroMemory(&runtime->session_options, sizeof(runtime->session_options)); + ReleaseSRWLockExclusive(&runtime->session_lock); +} + +static DWORD ready_remaining(ULONGLONG started) +{ + ULONGLONG elapsed = GetTickCount64() - started; + return elapsed >= UL_PLUGIN_READY_TIMEOUT_MS ? 0u : + UL_PLUGIN_READY_TIMEOUT_MS - (DWORD)elapsed; +} + +static int session_worker(ul_plugin_runtime *runtime, ul_admission *admission) +{ + uint8_t command[UL_SESSION_COMMAND_BYTES]; + ul_arm_scheduler scheduler; + uintptr_t generation; + HANDLE waits[2] = {runtime->worker_cancel, runtime->arm_result}; + ULONGLONG started = GetTickCount64(); + DWORD remaining, waited; + bool accepted; + int result; + + AcquireSRWLockExclusive(&runtime->session_lock); + runtime->ready_started = started; + ReleaseSRWLockExclusive(&runtime->session_lock); + result = ul_admission_read_exact(admission, command, sizeof(command), + UL_PLUGIN_READY_TIMEOUT_MS); + if (result != UL_ADMISSION_AUTH_OK) + return result; + if (!ul_session_arm_request(command, sizeof(command), + runtime->session_options.session, + runtime->session_options.additional_mix_mask)) + return UL_ADMISSION_REJECTED; + + AcquireSRWLockExclusive(&runtime->session_lock); + accepted = runtime->session_phase == UL_SESSION_READY && + InterlockedCompareExchange(&runtime->closing, 0, 0) == 0; + scheduler = runtime->schedule_arm; + generation = runtime->session_generation; + if (accepted) { + runtime->pending_epoch = runtime->stream_epoch; + runtime->session_phase = UL_SESSION_ARM_PENDING; + } + ReleaseSRWLockExclusive(&runtime->session_lock); + if (!accepted) + return UL_ADMISSION_REJECTED; + if (scheduler == NULL || !scheduler(generation)) { + AcquireSRWLockExclusive(&runtime->session_lock); + if (runtime->session_phase == UL_SESSION_ARM_PENDING) { + runtime->session_phase = UL_SESSION_TERMINAL; + SetEvent(runtime->arm_result); + } + ReleaseSRWLockExclusive(&runtime->session_lock); + } + remaining = ready_remaining(started); + if (remaining == 0u) + return UL_ADMISSION_TIMEOUT; + waited = WaitForMultipleObjects(2, waits, FALSE, remaining); + if (waited != WAIT_OBJECT_0 + 1u) + return waited == WAIT_TIMEOUT ? UL_ADMISSION_TIMEOUT : UL_ADMISSION_CANCELLED; + AcquireSRWLockShared(&runtime->session_lock); + accepted = runtime->session_phase == UL_SESSION_ARMED || + runtime->session_phase == UL_SESSION_STARTING || + runtime->session_phase == UL_SESSION_STARTED; + ReleaseSRWLockShared(&runtime->session_lock); + remaining = ready_remaining(started); + if (remaining == 0u) + return UL_ADMISSION_TIMEOUT; + if (!ul_session_arm_reply(runtime->session_options.session, + runtime->session_options.additional_mix_mask, accepted, command)) + return UL_ADMISSION_REJECTED; + result = ul_admission_write_all(admission, command, sizeof(command), remaining); + if (result != UL_ADMISSION_AUTH_OK || !accepted) + return result == UL_ADMISSION_AUTH_OK ? UL_ADMISSION_REJECTED : result; + + /* Armed waiting has no duration cutoff. Until the PCM worker is integrated, + * retain consent only; never manufacture a Start descriptor or audio. */ + for (;;) { + DWORD available = 0; + bool start_expired = false; + waited = WaitForSingleObject(runtime->worker_cancel, 50u); + if (waited != WAIT_TIMEOUT) + return UL_ADMISSION_CANCELLED; + AcquireSRWLockExclusive(&runtime->session_lock); + if (runtime->session_phase == UL_SESSION_STARTING && + GetTickCount64() - runtime->stream_starting_at >= UL_PLUGIN_START_TIMEOUT_MS) { + session_cancel_locked(runtime); + start_expired = true; + } + ReleaseSRWLockExclusive(&runtime->session_lock); + if (start_expired) + return UL_ADMISSION_TIMEOUT; + result = ul_admission_probe(admission, &available); + if (result != UL_ADMISSION_AUTH_OK) + return result; + if (available != 0u) /* No second Arm or other command is accepted. */ + return UL_ADMISSION_REJECTED; + } +} + static DWORD WINAPI admission_worker(void *context) { ul_plugin_runtime *runtime = (ul_plugin_runtime *)context; @@ -106,10 +244,9 @@ static DWORD WINAPI admission_worker(void *context) result = ul_admission_authenticate(admission, UL_PLUGIN_AUTH_TIMEOUT_MS); if (result == UL_ADMISSION_AUTH_OK) { pipe = ul_admission_pipe(admission); - (void)WaitForSingleObject(runtime->worker_cancel, - UL_PLUGIN_READY_TIMEOUT_MS); + result = session_worker(runtime, admission); } - ul_admission_cancel(admission); + session_cancel(runtime); if (pipe != NULL && pipe != INVALID_HANDLE_VALUE) (void)DisconnectNamedPipe(pipe); InterlockedExchange(&runtime->worker_active, 0); @@ -131,7 +268,7 @@ static bool reap_worker(ul_plugin_runtime *runtime) return false; CloseHandle(runtime->worker); runtime->worker = NULL; - runtime->worker_admission = NULL; + session_release(runtime); ul_authorizer_release(runtime->authorizer); ResetEvent(runtime->worker_cancel); return true; @@ -142,14 +279,14 @@ static void retire_authorizer(ul_plugin_runtime *runtime) { if (runtime->authorizer == NULL) return; - SetEvent(runtime->worker_cancel); + session_cancel(runtime); ul_authorizer_revoke(runtime->authorizer); if (runtime->worker != NULL) { (void)WaitForSingleObject(runtime->worker, INFINITE); CloseHandle(runtime->worker); runtime->worker = NULL; - runtime->worker_admission = NULL; } + session_release(runtime); InterlockedExchange(&runtime->worker_active, 0); ul_authorizer_destroy(runtime->authorizer); runtime->authorizer = NULL; @@ -175,9 +312,12 @@ bool ul_plugin_start(void) if (runtime == NULL) goto permanent_failure; InitializeSRWLock(&runtime->operation_lock); + InitializeSRWLock(&runtime->session_lock); runtime->leases_zero = CreateEventW(NULL, TRUE, TRUE, NULL); runtime->worker_cancel = CreateEventW(NULL, TRUE, FALSE, NULL); - if (runtime->leases_zero == NULL || runtime->worker_cancel == NULL) + runtime->arm_result = CreateEventW(NULL, TRUE, FALSE, NULL); + if (runtime->leases_zero == NULL || runtime->worker_cancel == NULL || + runtime->arm_result == NULL) goto fail; result = ul_pairing_store_open(&runtime->store); if (result != UL_PAIRING_OK) { @@ -224,6 +364,7 @@ bool ul_plugin_start(void) ul_pairing_store_destroy(runtime->store); ReleaseSRWLockExclusive(&runtime->operation_lock); CloseHandle(runtime->worker_cancel); + CloseHandle(runtime->arm_result); CloseHandle(runtime->leases_zero); SecureZeroMemory(runtime, sizeof(*runtime)); HeapFree(GetProcessHeap(), 0, runtime); @@ -239,6 +380,8 @@ bool ul_plugin_start(void) fail: SecureZeroMemory(key, sizeof(key)); + if (runtime->arm_result != NULL) + CloseHandle(runtime->arm_result); if (runtime->worker_cancel != NULL) CloseHandle(runtime->worker_cancel); if (runtime->leases_zero != NULL) @@ -261,14 +404,13 @@ void ul_plugin_stop_accepting(void) plugin_global.snapshot.storage_result = UL_PAIRING_CANCELLED; plugin_global.snapshot.owns_store = false; plugin_global.snapshot.admission_pending = false; - if (runtime != NULL) + if (runtime != NULL) { InterlockedExchange(&runtime->closing, 1); + session_cancel(runtime); + if (InterlockedCompareExchange(&runtime->mutation_active, 0, 0) != 0) + ul_pairing_cancel_request(&runtime->mutation_cancel); + } ReleaseSRWLockExclusive(&plugin_global.lock); - if (runtime == NULL) - return; - if (InterlockedCompareExchange(&runtime->mutation_active, 0, 0) != 0) - ul_pairing_cancel_request(&runtime->mutation_cancel); - SetEvent(runtime->worker_cancel); } void ul_plugin_close(void) @@ -293,6 +435,7 @@ void ul_plugin_close(void) runtime->store = NULL; ReleaseSRWLockExclusive(&runtime->operation_lock); CloseHandle(runtime->worker_cancel); + CloseHandle(runtime->arm_result); CloseHandle(runtime->leases_zero); SecureZeroMemory(runtime, sizeof(*runtime)); HeapFree(GetProcessHeap(), 0, runtime); @@ -516,22 +659,27 @@ bool ul_plugin_prepare(const uint8_t *challenge, size_t challenge_size, if (InterlockedCompareExchange(&runtime->closing, 0, 0) != 0 || runtime->status != UL_PLUGIN_PAIRED || !runtime->owns_store || runtime->authorizer == NULL || !reap_worker(runtime) || - runtime->worker != NULL) + runtime->worker != NULL || runtime->session_generation == UINTPTR_MAX) goto cleanup; admission = ul_authorizer_prepare(runtime->authorizer, challenge, challenge_size, proof, proof_size, &options); - SecureZeroMemory(&options, sizeof(options)); if (admission == NULL || InterlockedCompareExchange(&runtime->closing, 0, 0) != 0) goto cleanup; ResetEvent(runtime->worker_cancel); + ResetEvent(runtime->arm_result); + AcquireSRWLockExclusive(&runtime->session_lock); runtime->worker_admission = admission; + runtime->session_options = options; + runtime->session_generation++; + runtime->session_phase = UL_SESSION_READY; + ReleaseSRWLockExclusive(&runtime->session_lock); InterlockedExchange(&runtime->worker_active, 1); runtime->worker = CreateThread(NULL, 0, admission_worker, runtime, 0, NULL); if (runtime->worker == NULL) { InterlockedExchange(&runtime->worker_active, 0); - runtime->worker_admission = NULL; + session_release(runtime); ul_authorizer_release(runtime->authorizer); goto cleanup; } @@ -544,3 +692,90 @@ bool ul_plugin_prepare(const uint8_t *challenge, size_t challenge_size, runtime_release(runtime); return success; } + +bool ul_plugin_set_arm_scheduler(ul_arm_scheduler scheduler) +{ + ul_plugin_runtime *runtime = runtime_acquire(); + bool accepted = false; + if (runtime == NULL || scheduler == NULL) { + if (runtime != NULL) + runtime_release(runtime); + return false; + } + AcquireSRWLockExclusive(&runtime->operation_lock); + AcquireSRWLockExclusive(&runtime->session_lock); + if (runtime->schedule_arm == NULL && runtime->worker == NULL && + InterlockedCompareExchange(&runtime->closing, 0, 0) == 0) { + runtime->schedule_arm = scheduler; + accepted = true; + } + ReleaseSRWLockExclusive(&runtime->session_lock); + ReleaseSRWLockExclusive(&runtime->operation_lock); + runtime_release(runtime); + return accepted; +} + +void ul_plugin_arm_checked(uintptr_t generation, bool idle) +{ + ul_plugin_runtime *runtime = runtime_acquire(); + if (runtime == NULL) + return; + AcquireSRWLockExclusive(&runtime->session_lock); + if (generation == runtime->session_generation && + runtime->session_phase == UL_SESSION_ARM_PENDING) { + runtime->session_phase = idle && + runtime->stream_epoch == runtime->pending_epoch && + runtime->stream_epoch != UINT64_MAX && + ready_remaining(runtime->ready_started) != 0u && + InterlockedCompareExchange(&runtime->closing, 0, 0) == 0 ? + UL_SESSION_ARMED : UL_SESSION_TERMINAL; + SetEvent(runtime->arm_result); + } + ReleaseSRWLockExclusive(&runtime->session_lock); + runtime_release(runtime); +} + +void ul_plugin_stream_event(ul_stream_event event) +{ + ul_plugin_runtime *runtime = runtime_acquire(); + if (runtime == NULL) + return; + AcquireSRWLockExclusive(&runtime->session_lock); + if (event < UL_STREAM_STARTING || event > UL_STREAM_STOPPED || + runtime->stream_epoch == UINT64_MAX) { + session_cancel_locked(runtime); + } else { + runtime->stream_epoch++; + if (runtime->session_phase != UL_SESSION_NONE && + runtime->session_phase != UL_SESSION_TERMINAL) { + if (event == UL_STREAM_STARTING && runtime->session_phase == UL_SESSION_ARMED) { + runtime->session_phase = UL_SESSION_STARTING; + runtime->stream_starting_at = GetTickCount64(); + } else if (event == UL_STREAM_STARTED && runtime->session_phase == UL_SESSION_STARTING && + GetTickCount64() - runtime->stream_starting_at < UL_PLUGIN_START_TIMEOUT_MS) + runtime->session_phase = UL_SESSION_STARTED; + else if (runtime->session_phase == UL_SESSION_ARM_PENDING) { + /* A valid Arm is refused while its verified pipe is still + * usable. Lifecycle teardown retains the cancel-first path. */ + runtime->session_phase = UL_SESSION_TERMINAL; + SetEvent(runtime->arm_result); + } else + session_cancel_locked(runtime); + } + } + ReleaseSRWLockExclusive(&runtime->session_lock); + runtime_release(runtime); +} + +ul_session_phase ul_plugin_session_status(void) +{ + ul_plugin_runtime *runtime = runtime_acquire(); + ul_session_phase phase = UL_SESSION_NONE; + if (runtime != NULL) { + AcquireSRWLockShared(&runtime->session_lock); + phase = runtime->session_phase; + ReleaseSRWLockShared(&runtime->session_lock); + runtime_release(runtime); + } + return phase; +} diff --git a/native/obs-plugin/src/plugin_state.h b/native/obs-plugin/src/plugin_state.h index 35ed6b8c..9332637e 100644 --- a/native/obs-plugin/src/plugin_state.h +++ b/native/obs-plugin/src/plugin_state.h @@ -27,10 +27,29 @@ typedef struct ul_plugin_snapshot { ul_plugin_status status; ul_pairing_result storage_result; bool owns_store; - /* False means no active authenticate/READY worker. Reaping is deferred. */ + /* False means no active authentication/READY/armed worker. Reaping is deferred. */ bool admission_pending; } ul_plugin_snapshot; +typedef enum ul_session_phase { + UL_SESSION_NONE = 0, UL_SESSION_READY, UL_SESSION_ARM_PENDING, + UL_SESSION_ARMED, UL_SESSION_STARTING, UL_SESSION_STARTED, + UL_SESSION_TERMINAL +} ul_session_phase; + +typedef enum ul_stream_event { + UL_STREAM_STARTING = 1, UL_STREAM_STARTED, UL_STREAM_STOPPING, + UL_STREAM_STOPPED +} ul_stream_event; + +/* Pinned callback; queues nonblocking frontend work with numeric generation + * only. It must refuse after frontend shutdown, without calling OBS. */ +typedef bool (*ul_arm_scheduler)(uintptr_t generation); +bool ul_plugin_set_arm_scheduler(ul_arm_scheduler scheduler); +void ul_plugin_arm_checked(uintptr_t generation, bool idle); +void ul_plugin_stream_event(ul_stream_event event); +ul_session_phase ul_plugin_session_status(void); + /* * Permanently pins this DLL generation before callbacks may be registered. * A process may start this component once; close is final and cannot be reset. diff --git a/native/obs-plugin/src/session_protocol.c b/native/obs-plugin/src/session_protocol.c new file mode 100644 index 00000000..2983a100 --- /dev/null +++ b/native/obs-plugin/src/session_protocol.c @@ -0,0 +1,45 @@ +// SPDX-License-Identifier: GPL-2.0-or-later +#include "session_protocol.h" + +#include + +static bool valid_session(const uint8_t session[16]) +{ + uint8_t combined = 0; + size_t index; + if (session == NULL) + return false; + for (index = 0; index < 16u; ++index) + combined |= session[index]; + return combined != 0; +} + +bool ul_session_arm_request(const uint8_t *data, size_t size, + const uint8_t session[16], uint8_t mix_mask) +{ + return data != NULL && size == UL_SESSION_COMMAND_BYTES && + valid_session(session) && mix_mask <= 63u && + memcmp(data, "ULAC", 4u) == 0 && + data[4] == 1u && data[5] == 1u && + data[6] == 0u && data[7] == 0u && + memcmp(data + 8u, session, 16u) == 0 && + data[24] == mix_mask && data[25] == 0u && + data[26] == 0u && data[27] == 0u; +} + +bool ul_session_arm_reply(const uint8_t session[16], uint8_t mix_mask, + bool accepted, uint8_t out[UL_SESSION_COMMAND_BYTES]) +{ + if (out == NULL) + return false; + memset(out, 0, UL_SESSION_COMMAND_BYTES); + if (!valid_session(session) || mix_mask > 63u) + return false; + memcpy(out, "ULAC", 4u); + out[4] = 1u; + out[5] = 2u; + memcpy(out + 8u, session, 16u); + out[24] = mix_mask; + out[25] = accepted ? 1u : 2u; + return true; +} diff --git a/native/obs-plugin/src/session_protocol.h b/native/obs-plugin/src/session_protocol.h new file mode 100644 index 00000000..221de3fd --- /dev/null +++ b/native/obs-plugin/src/session_protocol.h @@ -0,0 +1,22 @@ +// SPDX-License-Identifier: GPL-2.0-or-later +#ifndef UTTERLEAF_OBS_SESSION_PROTOCOL_H +#define UTTERLEAF_OBS_SESSION_PROTOCOL_H + +#include +#include +#include + +#define UL_SESSION_COMMAND_BYTES 28u + +/* Fixed metadata only, inside the authenticated pipe after Hello/ACK. + * The prepared session and mix mask remain authoritative. A command does not + * authenticate a caller or establish an idle OBS state on its own. */ +bool ul_session_arm_request(const uint8_t *data, size_t size, + const uint8_t session[16], uint8_t mix_mask); + +/* Reply status 1 acknowledges committed Arm; status 2 refuses it. Both are + * terminal for this one Arm attempt. Clears output on invalid local arguments. */ +bool ul_session_arm_reply(const uint8_t session[16], uint8_t mix_mask, + bool accepted, uint8_t out[UL_SESSION_COMMAND_BYTES]); + +#endif diff --git a/native/obs-plugin/tests/admission_io_fault_test.c b/native/obs-plugin/tests/admission_io_fault_test.c new file mode 100644 index 00000000..4f83dbf3 --- /dev/null +++ b/native/obs-plugin/tests/admission_io_fault_test.c @@ -0,0 +1,189 @@ +// SPDX-License-Identifier: GPL-2.0-or-later +#include + +#include +#include +#include + +static int fault_mode; +static LONG peer_checks; +static LONG cancel_on_peer_check; +static LONG delay_on_peer_check; +static volatile LONG *fault_cancelled; +static HANDLE fault_cancel_event; + +static HANDLE fault_event_create(void) +{ + if (fault_mode == 1) + return NULL; + return CreateEventW(NULL, TRUE, FALSE, NULL); +} + +static BOOL fault_peer_expected(const void *admission) +{ + LONG current; + + (void)admission; + current = InterlockedIncrement(&peer_checks); + if (current == InterlockedCompareExchange(&cancel_on_peer_check, 0, 0)) { + InterlockedExchange(fault_cancelled, 1); + SetEvent(fault_cancel_event); + } + if (current == InterlockedCompareExchange(&delay_on_peer_check, 0, 0)) + Sleep(30); + return TRUE; +} + +static void fault_read(void *buffer, DWORD size, DWORD *transferred) +{ + if (fault_mode == 2) { + *transferred = 0; + } else if (fault_mode == 3) { + *transferred = size + 1u; + } else { + memset(buffer, 0x3c, 1); + *transferred = 1; + } +} + +static BOOL fault_peek(HANDLE pipe, DWORD *available) +{ + (void)pipe; + *available = 0; + return fault_mode != 5; +} + +#define UL_ADMISSION_IO_EVENT_CREATE() fault_event_create() +#define UL_ADMISSION_PEER_EXPECTED(admission) fault_peer_expected(admission) +#define UL_ADMISSION_READ_PIPE(admission, overlapped, buffer, size, started, \ + timeout_ms, transferred) \ + (fault_read(buffer, size, transferred), UL_OPERATION_OK) +#define UL_ADMISSION_PEEK_PIPE(pipe, available) fault_peek(pipe, available) +#include "../src/admission.c" + +static int check(BOOL condition, const char *message) +{ + if (!condition) { + fprintf(stderr, "FAIL: %s\n", message); + return 1; + } + return 0; +} + +static BOOL all_zero(const uint8_t *buffer, DWORD size) +{ + DWORD index; + for (index = 0; index < size; ++index) { + if (buffer[index] != 0u) + return FALSE; + } + return TRUE; +} + +static void initialize_admission(ul_admission *admission) +{ + SecureZeroMemory(admission, sizeof(*admission)); + admission->pipe = (HANDLE)(uintptr_t)1; + admission->cancel_event = CreateEventW(NULL, TRUE, FALSE, NULL); + admission->authenticated = 1; + fault_cancelled = &admission->cancelled; + fault_cancel_event = admission->cancel_event; +} + +static void finish_admission(ul_admission *admission) +{ + if (admission->cancel_event != NULL) + CloseHandle(admission->cancel_event); + fault_cancelled = NULL; + fault_cancel_event = NULL; + SecureZeroMemory(admission, sizeof(*admission)); +} + +static int failure_case(int mode, const char *message) +{ + ul_admission admission; + uint8_t buffer[8]; + int failures = 0; + + fault_mode = mode; + InterlockedExchange(&peer_checks, 0); + InterlockedExchange(&cancel_on_peer_check, 0); + InterlockedExchange(&delay_on_peer_check, 0); + initialize_admission(&admission); + failures += check(admission.cancel_event != NULL, "create cancel event"); + memset(buffer, 0xa5, sizeof(buffer)); + failures += check(ul_admission_read_exact(&admission, buffer, + sizeof(buffer), 500) == + UL_ADMISSION_IO_ERROR, + message); + failures += check(all_zero(buffer, sizeof(buffer)), + "failed exact read wipes output"); + failures += check(admission.cancelled != 0, + "failed exact read terminally cancels"); + finish_admission(&admission); + return failures; +} + +int main(void) +{ + ul_admission admission; + uint8_t buffer[4]; + DWORD available = 99; + int failures = 0; + + failures += failure_case(1, "event creation failure cannot report success"); + failures += failure_case(2, "zero-byte completion cannot report success"); + failures += failure_case(3, "over-count completion cannot report success"); + + fault_mode = 4; + InterlockedExchange(&peer_checks, 0); + initialize_admission(&admission); + failures += check(ul_admission_read_exact(&admission, buffer, + sizeof(buffer), 500) == + UL_ADMISSION_AUTH_OK, + "fragmented injected read completes exactly"); + failures += check(InterlockedCompareExchange(&peer_checks, 0, 0) == 10, + "retained peer checked before and after every transfer"); + finish_admission(&admission); + + fault_mode = 4; + InterlockedExchange(&peer_checks, 0); + InterlockedExchange(&cancel_on_peer_check, 10); + initialize_admission(&admission); + memset(buffer, 0xa5, sizeof(buffer)); + failures += check(ul_admission_read_exact(&admission, buffer, + sizeof(buffer), 500) == + UL_ADMISSION_CANCELLED, + "cancellation during final peer check wins over success"); + failures += check(all_zero(buffer, sizeof(buffer)), + "final peer cancellation wipes copied bytes"); + finish_admission(&admission); + + InterlockedExchange(&peer_checks, 0); + InterlockedExchange(&cancel_on_peer_check, 0); + InterlockedExchange(&delay_on_peer_check, 10); + initialize_admission(&admission); + memset(buffer, 0xa5, sizeof(buffer)); + failures += check(ul_admission_read_exact(&admission, buffer, + sizeof(buffer), 10) == + UL_ADMISSION_TIMEOUT, + "deadline after final peer check wins over success"); + failures += check(all_zero(buffer, sizeof(buffer)), + "final peer timeout wipes copied bytes"); + finish_admission(&admission); + + fault_mode = 5; + InterlockedExchange(&peer_checks, 0); + InterlockedExchange(&delay_on_peer_check, 0); + initialize_admission(&admission); + failures += check(ul_admission_probe(&admission, &available) == + UL_ADMISSION_IO_ERROR, + "PeekNamedPipe failure cannot report success"); + failures += check(available == 0u && admission.cancelled != 0, + "failed probe clears output and terminally cancels"); + finish_admission(&admission); + + puts(failures == 0 ? "admission I/O fault tests passed" + : "admission I/O fault tests failed"); + return failures == 0 ? 0 : 1; +} diff --git a/native/obs-plugin/tests/admission_io_test.c b/native/obs-plugin/tests/admission_io_test.c new file mode 100644 index 00000000..9159eded --- /dev/null +++ b/native/obs-plugin/tests/admission_io_test.c @@ -0,0 +1,672 @@ +// SPDX-License-Identifier: GPL-2.0-or-later +#ifndef _WIN32_WINNT +#define _WIN32_WINNT 0x0600 +#endif + +#include "../src/admission.h" + +#include + +#include +#include +#include +#include + +#define TEST_TIMEOUT_MS 5000u +#define PAYLOAD_BYTES 8193u + +typedef struct test_peer { + PROCESS_INFORMATION process; + HANDLE ready; + HANDLE release; + HANDLE done; + ul_admission *admission; +} test_peer; + +typedef struct read_call { + ul_admission *admission; + uint8_t buffer[64]; + int result; +} read_call; + +typedef struct write_call { + ul_admission *admission; + uint8_t buffer[UL_ADMISSION_MAX_IO_BYTES]; + int result; +} write_call; + +static int check(BOOL condition, const char *message) +{ + if (!condition) { + fprintf(stderr, "FAIL: %s (win32=%lu)\n", message, + (unsigned long)GetLastError()); + return 1; + } + return 0; +} + +static void make_session(uint8_t session[16], unsigned counter) +{ + ULONGLONG tick = GetTickCount64(); + DWORD pid = GetCurrentProcessId(); + + memcpy(session, &tick, sizeof(tick)); + memcpy(session + 8, &pid, sizeof(pid)); + memcpy(session + 12, &counter, sizeof(counter)); +} + +static void session_hex(const uint8_t session[16], WCHAR output[33]) +{ + static const WCHAR digits[] = L"0123456789abcdef"; + unsigned index; + + for (index = 0; index < 16u; ++index) { + output[index * 2u] = digits[session[index] >> 4]; + output[index * 2u + 1u] = digits[session[index] & 15u]; + } + output[32] = L'\0'; +} + +static BOOL parse_session(const WCHAR *input, uint8_t session[16]) +{ + unsigned index; + + if (input == NULL || wcslen(input) != 32u) + return FALSE; + for (index = 0; index < 16u; ++index) { + WCHAR high = input[index * 2u]; + WCHAR low = input[index * 2u + 1u]; + unsigned a; + unsigned b; + + if (high >= L'0' && high <= L'9') + a = (unsigned)(high - L'0'); + else if (high >= L'a' && high <= L'f') + a = (unsigned)(high - L'a') + 10u; + else + return FALSE; + if (low >= L'0' && low <= L'9') + b = (unsigned)(low - L'0'); + else if (low >= L'a' && low <= L'f') + b = (unsigned)(low - L'a') + 10u; + else + return FALSE; + session[index] = (uint8_t)((a << 4) | b); + } + return TRUE; +} + +static BOOL sync_write_all(HANDLE pipe, const uint8_t *buffer, DWORD size) +{ + DWORD offset = 0; + + while (offset < size) { + DWORD transferred = 0; + if (!WriteFile(pipe, buffer + offset, size - offset, &transferred, + NULL) || transferred == 0u) + return FALSE; + offset += transferred; + } + return TRUE; +} + +static BOOL sync_read_exact(HANDLE pipe, uint8_t *buffer, DWORD size) +{ + DWORD offset = 0; + + while (offset < size) { + DWORD transferred = 0; + if (!ReadFile(pipe, buffer + offset, size - offset, &transferred, + NULL) || transferred == 0u) + return FALSE; + offset += transferred; + } + return TRUE; +} + +static int child_main(int argc, WCHAR **argv) +{ + uint8_t session[16]; + uint8_t hello[56]; + uint8_t ack[56]; + uint8_t payload[PAYLOAD_BYTES]; + uint8_t response[PAYLOAD_BYTES]; + WCHAR pipe_name[64]; + HANDLE pipe = INVALID_HANDLE_VALUE; + HANDLE ready = NULL; + HANDLE release = NULL; + HANDLE done = NULL; + ULONGLONG deadline = GetTickCount64() + TEST_TIMEOUT_MS; + DWORD index; + int result = 2; + + if (argc != 7 || !parse_session(argv[3], session)) + return 2; + if (_snwprintf(pipe_name, 64, L"\\\\.\\pipe\\Utterleaf.OBS.%ls", + argv[3]) < 0) + return 2; + ready = OpenEventW(EVENT_MODIFY_STATE, FALSE, argv[4]); + release = OpenEventW(SYNCHRONIZE, FALSE, argv[5]); + done = OpenEventW(EVENT_MODIFY_STATE, FALSE, argv[6]); + if (ready == NULL || release == NULL || done == NULL) + goto cleanup; + while (GetTickCount64() < deadline) { + pipe = CreateFileW(pipe_name, GENERIC_READ | GENERIC_WRITE, 0, NULL, + OPEN_EXISTING, 0, NULL); + if (pipe != INVALID_HANDLE_VALUE) + break; + Sleep(10); + } + if (pipe == INVALID_HANDLE_VALUE) + goto cleanup; + + SecureZeroMemory(hello, sizeof(hello)); + memcpy(hello, "ULAH", 4); + hello[4] = 1; + hello[5] = 1; + memcpy(hello + 8, session, 16); + memset(hello + 24, 's', 32); + if (!sync_write_all(pipe, hello, 7) || + !sync_write_all(pipe, hello + 7, 19) || + !sync_write_all(pipe, hello + 26, 30) || + !sync_read_exact(pipe, ack, sizeof(ack))) + goto cleanup; + + if (wcscmp(argv[2], L"eof") == 0 || + wcscmp(argv[2], L"partial_eof") == 0) { + if (wcscmp(argv[2], L"partial_eof") == 0 && + !sync_write_all(pipe, (const uint8_t *)"partial", 7)) + goto cleanup; + CloseHandle(pipe); + pipe = INVALID_HANDLE_VALUE; + SetEvent(ready); + Sleep(2000); + result = 0; + goto cleanup; + } + if (wcscmp(argv[2], L"exit") == 0) { + SetEvent(ready); + result = 0; + goto cleanup; + } + if (wcscmp(argv[2], L"stall") == 0) { + SetEvent(ready); + WaitForSingleObject(release, TEST_TIMEOUT_MS * 2u); + result = 0; + goto cleanup; + } + if (wcscmp(argv[2], L"partial_stall") == 0) { + if (!sync_write_all(pipe, (const uint8_t *)"partial", 7)) + goto cleanup; + SetEvent(ready); + WaitForSingleObject(release, TEST_TIMEOUT_MS * 2u); + result = 0; + goto cleanup; + } + if (wcscmp(argv[2], L"roundtrip") != 0) + goto cleanup; + for (index = 0; index < PAYLOAD_BYTES; ++index) + payload[index] = (uint8_t)(index * 29u + 7u); + if (!sync_write_all(pipe, payload, 3)) + goto cleanup; + SetEvent(ready); + if (WaitForSingleObject(release, TEST_TIMEOUT_MS) != WAIT_OBJECT_0 || + !sync_write_all(pipe, payload + 3, 97) || + !sync_write_all(pipe, payload + 100, PAYLOAD_BYTES - 100) || + !sync_read_exact(pipe, response, sizeof(response)) || + memcmp(response, payload, sizeof(payload)) != 0) + goto cleanup; + SetEvent(done); + Sleep(TEST_TIMEOUT_MS); + result = 0; + +cleanup: + SecureZeroMemory(hello, sizeof(hello)); + SecureZeroMemory(ack, sizeof(ack)); + SecureZeroMemory(payload, sizeof(payload)); + SecureZeroMemory(response, sizeof(response)); + if (pipe != INVALID_HANDLE_VALUE) + CloseHandle(pipe); + if (ready != NULL) + CloseHandle(ready); + if (release != NULL) + CloseHandle(release); + if (done != NULL) + CloseHandle(done); + return result; +} + +static BOOL start_peer(test_peer *peer, const WCHAR *mode, + const uint8_t session[16], unsigned counter) +{ + WCHAR executable[MAX_PATH]; + WCHAR hex[33]; + WCHAR ready_name[96]; + WCHAR release_name[96]; + WCHAR done_name[96]; + WCHAR command[1024]; + STARTUPINFOW startup; + + SecureZeroMemory(peer, sizeof(*peer)); + SecureZeroMemory(&startup, sizeof(startup)); + startup.cb = sizeof(startup); + if (GetModuleFileNameW(NULL, executable, MAX_PATH) == 0u) + return FALSE; + session_hex(session, hex); + _snwprintf(ready_name, 96, L"Local\\ULAdmissionIo.%lu.%u.ready", + (unsigned long)GetCurrentProcessId(), counter); + _snwprintf(release_name, 96, L"Local\\ULAdmissionIo.%lu.%u.release", + (unsigned long)GetCurrentProcessId(), counter); + _snwprintf(done_name, 96, L"Local\\ULAdmissionIo.%lu.%u.done", + (unsigned long)GetCurrentProcessId(), counter); + peer->ready = CreateEventW(NULL, TRUE, FALSE, ready_name); + peer->release = CreateEventW(NULL, TRUE, FALSE, release_name); + peer->done = CreateEventW(NULL, TRUE, FALSE, done_name); + if (peer->ready == NULL || peer->release == NULL || peer->done == NULL) + return FALSE; + if (_snwprintf(command, 1024, L"\"%ls\" --child %ls %ls %ls %ls %ls", + executable, mode, hex, ready_name, release_name, + done_name) < 0) + return FALSE; + return CreateProcessW(executable, command, NULL, NULL, FALSE, + CREATE_NO_WINDOW, NULL, NULL, &startup, + &peer->process); +} + +static void stop_peer(test_peer *peer) +{ + if (peer->admission != NULL) { + ul_admission_cancel(peer->admission); + ul_admission_destroy(peer->admission); + } + if (peer->process.hProcess != NULL) { + SetEvent(peer->release); + if (WaitForSingleObject(peer->process.hProcess, 500) == WAIT_TIMEOUT) { + TerminateProcess(peer->process.hProcess, 99); + WaitForSingleObject(peer->process.hProcess, TEST_TIMEOUT_MS); + } + CloseHandle(peer->process.hProcess); + } + if (peer->process.hThread != NULL) + CloseHandle(peer->process.hThread); + if (peer->ready != NULL) + CloseHandle(peer->ready); + if (peer->release != NULL) + CloseHandle(peer->release); + if (peer->done != NULL) + CloseHandle(peer->done); + SecureZeroMemory(peer, sizeof(*peer)); +} + +static BOOL authenticated_peer(test_peer *peer, const WCHAR *mode, + uint8_t session[16], unsigned counter) +{ + make_session(session, counter); + if (!start_peer(peer, mode, session, counter)) + return FALSE; + peer->admission = ul_admission_create(peer->process.dwProcessId, session); + return peer->admission != NULL && + ul_admission_authenticate(peer->admission, TEST_TIMEOUT_MS) == + UL_ADMISSION_AUTH_OK; +} + +static DWORD WINAPI blocked_read(LPVOID parameter) +{ + read_call *call = (read_call *)parameter; + + memset(call->buffer, 0xa5, sizeof(call->buffer)); + call->result = ul_admission_read_exact( + call->admission, call->buffer, sizeof(call->buffer), TEST_TIMEOUT_MS); + return 0; +} + +static DWORD WINAPI blocked_write(LPVOID parameter) +{ + write_call *call = (write_call *)parameter; + + memset(call->buffer, 0x5a, sizeof(call->buffer)); + call->result = ul_admission_write_all( + call->admission, call->buffer, sizeof(call->buffer), TEST_TIMEOUT_MS); + return 0; +} + +static int test_roundtrip_and_probe(void) +{ + test_peer peer; + uint8_t session[16]; + uint8_t payload[PAYLOAD_BYTES]; + DWORD available = 0; + DWORD index; + int failures = 0; + + failures += check(authenticated_peer(&peer, L"roundtrip", session, 1), + "authenticate roundtrip child"); + if (failures != 0) { + stop_peer(&peer); + return failures; + } + failures += check(WaitForSingleObject(peer.ready, TEST_TIMEOUT_MS) == + WAIT_OBJECT_0, + "roundtrip child ready"); + failures += check(ul_admission_probe(peer.admission, &available) == 0 && + available == 3u, + "probe reports queued fragment"); + available = 0; + failures += check(ul_admission_probe(peer.admission, &available) == 0 && + available == 3u, + "probe does not consume bytes"); + SetEvent(peer.release); + failures += check(ul_admission_read_exact(peer.admission, payload, + sizeof(payload), + TEST_TIMEOUT_MS) == 0, + "fragmented exact read"); + for (index = 0; index < PAYLOAD_BYTES; ++index) { + if (payload[index] != (uint8_t)(index * 29u + 7u)) { + failures += check(FALSE, "exact read payload"); + break; + } + } + failures += check(ul_admission_write_all(peer.admission, payload, + sizeof(payload), + TEST_TIMEOUT_MS) == 0, + "bounded full write"); + failures += check(WaitForSingleObject(peer.done, TEST_TIMEOUT_MS) == + WAIT_OBJECT_0, + "child validates full write"); + SecureZeroMemory(payload, sizeof(payload)); + stop_peer(&peer); + return failures; +} + +static int test_pre_auth_is_terminal_and_wipes(void) +{ + test_peer peer; + uint8_t session[16]; + uint8_t buffer[16]; + unsigned index; + int failures = 0; + + make_session(session, 2); + failures += check(start_peer(&peer, L"stall", session, 2), + "start preauth child"); + if (failures != 0) { + stop_peer(&peer); + return failures; + } + peer.admission = ul_admission_create(peer.process.dwProcessId, session); + failures += check(peer.admission != NULL, "create preauth admission"); + memset(buffer, 0xa5, sizeof(buffer)); + failures += check(ul_admission_read_exact(peer.admission, buffer, + sizeof(buffer), 500) == + UL_ADMISSION_REJECTED, + "reject preauth read"); + for (index = 0; index < sizeof(buffer); ++index) + failures += check(buffer[index] == 0, "wipe preauth read buffer"); + failures += check(ul_admission_authenticate(peer.admission, 100) == + UL_ADMISSION_CANCELLED, + "preauth misuse terminally cancels"); + stop_peer(&peer); + return failures; +} + +static int test_timeout_and_cancel_wipe(void) +{ + test_peer peer; + uint8_t session[16]; + uint8_t buffer[64]; + read_call call; + HANDLE thread; + DWORD available = 99; + unsigned index; + int failures = 0; + + failures += check(authenticated_peer(&peer, L"stall", session, 3), + "authenticate timeout child"); + failures += check(WaitForSingleObject(peer.ready, TEST_TIMEOUT_MS) == + WAIT_OBJECT_0, + "timeout child ready"); + failures += check(ul_admission_probe(peer.admission, &available) == 0 && + available == 0u, + "zero-byte probe is live"); + memset(buffer, 0xa5, sizeof(buffer)); + failures += check(ul_admission_read_exact(peer.admission, buffer, + sizeof(buffer), 100) == + UL_ADMISSION_TIMEOUT, + "exact read timeout"); + for (index = 0; index < sizeof(buffer); ++index) + failures += check(buffer[index] == 0, "wipe timed-out read buffer"); + stop_peer(&peer); + + failures += check(authenticated_peer(&peer, L"stall", session, 4), + "authenticate cancellation child"); + failures += check(WaitForSingleObject(peer.ready, TEST_TIMEOUT_MS) == + WAIT_OBJECT_0, + "cancellation child ready"); + SecureZeroMemory(&call, sizeof(call)); + call.admission = peer.admission; + thread = CreateThread(NULL, 0, blocked_read, &call, 0, NULL); + failures += check(thread != NULL, "start blocked exact read"); + if (thread != NULL) { + Sleep(100); + ul_admission_cancel(peer.admission); + failures += check(WaitForSingleObject(thread, TEST_TIMEOUT_MS) == + WAIT_OBJECT_0, + "cancel drains exact read"); + failures += check(call.result == UL_ADMISSION_CANCELLED, + "cancel result"); + for (index = 0; index < sizeof(call.buffer); ++index) + failures += check(call.buffer[index] == 0, + "wipe cancelled read buffer"); + CloseHandle(thread); + } + stop_peer(&peer); + return failures; +} + +static BOOL all_zero(const uint8_t *buffer, DWORD size) +{ + DWORD index; + + for (index = 0; index < size; ++index) { + if (buffer[index] != 0u) + return FALSE; + } + return TRUE; +} + +static int test_invalid_arguments(void) +{ + unsigned kind; + int failures = 0; + + for (kind = 0; kind < 5u; ++kind) { + test_peer peer; + uint8_t session[16]; + uint8_t small[16]; + uint8_t *large = NULL; + int result; + + failures += check(authenticated_peer(&peer, L"stall", session, + 10u + kind), + "authenticate invalid-argument child"); + failures += check(WaitForSingleObject(peer.ready, TEST_TIMEOUT_MS) == + WAIT_OBJECT_0, + "invalid-argument child ready"); + memset(small, 0xa5, sizeof(small)); + if (kind == 0u) + result = ul_admission_read_exact(peer.admission, small, 0, 500); + else if (kind == 1u) { + large = (uint8_t *)HeapAlloc(GetProcessHeap(), 0, + UL_ADMISSION_MAX_IO_BYTES + 1u); + failures += check(large != NULL, "allocate oversized fixture buffer"); + result = large == NULL + ? UL_ADMISSION_IO_ERROR + : ul_admission_read_exact( + peer.admission, large, + UL_ADMISSION_MAX_IO_BYTES + 1u, 500); + } else if (kind == 2u) + result = ul_admission_read_exact(peer.admission, small, + sizeof(small), 0); + else if (kind == 3u) + result = ul_admission_read_exact(peer.admission, small, + sizeof(small), 30001); + else + result = ul_admission_read_exact(peer.admission, NULL, + sizeof(small), 500); + failures += check(result == UL_ADMISSION_IO_ERROR, + "invalid argument returns I/O error"); + failures += check(ul_admission_pipe(peer.admission) == NULL, + "invalid argument is terminal"); + if (kind == 2u || kind == 3u) + failures += check(all_zero(small, sizeof(small)), + "invalid timeout wipes read buffer"); + if (large != NULL) { + SecureZeroMemory(large, UL_ADMISSION_MAX_IO_BYTES + 1u); + HeapFree(GetProcessHeap(), 0, large); + } + stop_peer(&peer); + } + return failures; +} + +static int test_partial_failure_wipes(void) +{ + const WCHAR *modes[2] = {L"partial_eof", L"partial_stall"}; + unsigned mode_index; + int failures = 0; + + for (mode_index = 0; mode_index < 2u; ++mode_index) { + test_peer peer; + uint8_t session[16]; + uint8_t buffer[64]; + int result; + + failures += check(authenticated_peer(&peer, modes[mode_index], session, + 20u + mode_index), + "authenticate partial-read child"); + failures += check(WaitForSingleObject(peer.ready, TEST_TIMEOUT_MS) == + WAIT_OBJECT_0, + "partial-read child ready"); + memset(buffer, 0xa5, sizeof(buffer)); + result = ul_admission_read_exact(peer.admission, buffer, + sizeof(buffer), 150); + failures += check( + result == (mode_index == 0u ? UL_ADMISSION_IO_ERROR + : UL_ADMISSION_TIMEOUT), + "partial read returns terminal cause"); + failures += check(all_zero(buffer, sizeof(buffer)), + "partial copied bytes are wiped"); + stop_peer(&peer); + } + return failures; +} + +static int test_cancel_before_call_and_blocked_write(void) +{ + test_peer peer; + uint8_t session[16]; + uint8_t buffer[16]; + write_call *call; + HANDLE thread; + int failures = 0; + + failures += check(authenticated_peer(&peer, L"stall", session, 30), + "authenticate pre-cancel child"); + failures += check(WaitForSingleObject(peer.ready, TEST_TIMEOUT_MS) == + WAIT_OBJECT_0, + "pre-cancel child ready"); + ul_admission_cancel(peer.admission); + memset(buffer, 0xa5, sizeof(buffer)); + failures += check(ul_admission_read_exact(peer.admission, buffer, + sizeof(buffer), 500) == + UL_ADMISSION_CANCELLED, + "cancel before call preserves cancelled result"); + failures += check(all_zero(buffer, sizeof(buffer)), + "pre-cancelled read buffer wiped"); + stop_peer(&peer); + + failures += check(authenticated_peer(&peer, L"stall", session, 31), + "authenticate blocked-write child"); + failures += check(WaitForSingleObject(peer.ready, TEST_TIMEOUT_MS) == + WAIT_OBJECT_0, + "blocked-write child ready"); + call = (write_call *)HeapAlloc(GetProcessHeap(), HEAP_ZERO_MEMORY, + sizeof(*call)); + failures += check(call != NULL, "allocate blocked-write call"); + if (call != NULL) { + call->admission = peer.admission; + thread = CreateThread(NULL, 0, blocked_write, call, 0, NULL); + failures += check(thread != NULL, "start maximum blocked write"); + if (thread != NULL) { + Sleep(100); + ul_admission_cancel(peer.admission); + failures += check(WaitForSingleObject(thread, TEST_TIMEOUT_MS) == + WAIT_OBJECT_0, + "cancel drains maximum blocked write"); + failures += check(call->result == UL_ADMISSION_CANCELLED, + "blocked write reports cancellation"); + CloseHandle(thread); + } + SecureZeroMemory(call, sizeof(*call)); + HeapFree(GetProcessHeap(), 0, call); + } + stop_peer(&peer); + return failures; +} + +static int test_eof_and_dead_peer(void) +{ + test_peer peer; + uint8_t session[16]; + uint8_t buffer[16]; + DWORD available = 99; + unsigned index; + int failures = 0; + + failures += check(authenticated_peer(&peer, L"eof", session, 5), + "authenticate EOF child"); + failures += check(WaitForSingleObject(peer.ready, TEST_TIMEOUT_MS) == + WAIT_OBJECT_0, + "EOF child closed pipe"); + memset(buffer, 0xa5, sizeof(buffer)); + failures += check(ul_admission_read_exact(peer.admission, buffer, + sizeof(buffer), 500) != 0, + "EOF is terminal"); + for (index = 0; index < sizeof(buffer); ++index) + failures += check(buffer[index] == 0, "wipe EOF read buffer"); + stop_peer(&peer); + + failures += check(authenticated_peer(&peer, L"exit", session, 6), + "authenticate exiting child"); + failures += check(WaitForSingleObject(peer.ready, TEST_TIMEOUT_MS) == + WAIT_OBJECT_0, + "exiting child ready"); + failures += check(WaitForSingleObject(peer.process.hProcess, + TEST_TIMEOUT_MS) == WAIT_OBJECT_0, + "expected peer exits"); + failures += check(ul_admission_probe(peer.admission, &available) == + UL_ADMISSION_REJECTED && available == 0u, + "probe rejects dead retained peer"); + failures += check(ul_admission_pipe(peer.admission) == NULL, + "probe failure closes admission gate"); + stop_peer(&peer); + return failures; +} + +int wmain(int argc, WCHAR **argv) +{ + int failures; + + if (argc > 1 && wcscmp(argv[1], L"--child") == 0) + return child_main(argc, argv); + failures = test_roundtrip_and_probe(); + failures += test_pre_auth_is_terminal_and_wipes(); + failures += test_timeout_and_cancel_wipe(); + failures += test_invalid_arguments(); + failures += test_partial_failure_wipes(); + failures += test_cancel_before_call_and_blocked_write(); + failures += test_eof_and_dead_peer(); + puts(failures == 0 ? "admission I/O tests passed" + : "admission I/O tests failed"); + return failures == 0 ? 0 : 1; +} diff --git a/native/obs-plugin/tests/bridge_test.c b/native/obs-plugin/tests/bridge_test.c index 8ab4787c..743d6a37 100644 --- a/native/obs-plugin/tests/bridge_test.c +++ b/native/obs-plugin/tests/bridge_test.c @@ -17,6 +17,15 @@ static unsigned api_version; static char order[32]; static size_t order_size; static ul_plugin_snapshot snapshot; +static ul_arm_scheduler arm_scheduler; +static obs_task_t queued_task; +static void *queued_parameter; +static uintptr_t checked_generation; +static bool allow_scheduler, current_active, current_output_active, has_output, checked_idle; +static unsigned queue_calls, idle_checks, output_releases, stream_events; +static ul_stream_event last_stream_event; +static obs_output_t *current_output = (obs_output_t *)(uintptr_t)3; +static bool after_exit; static void record(char value) { assert(order_size + 1 < sizeof(order)); order[order_size++] = value; order[order_size] = 0; } static uint32_t test_version(void) { return LIBOBS_API_VER; } @@ -48,6 +57,24 @@ static bool test_unregister(obs_websocket_vendor handle, const char *name) } static void test_log(int level, const char *format, ...) { (void)level; (void)format; } bool ul_plugin_start(void) { ++start_calls; return allow_start; } +bool ul_plugin_set_arm_scheduler(ul_arm_scheduler scheduler) +{ arm_scheduler = scheduler; return allow_scheduler; } +void ul_plugin_arm_checked(uintptr_t generation, bool idle) +{ assert(!exiting); checked_generation = generation; checked_idle = idle; ++idle_checks; } +void ul_plugin_stream_event(ul_stream_event event) +{ assert(!exiting); ++stream_events; last_stream_event = event; } +static bool test_active(void) { assert(!exiting); return current_active; } +static obs_output_t *test_output(void) +{ assert(!after_exit); return has_output ? current_output : NULL; } +static bool test_output_active(const obs_output_t *output) +{ assert(!after_exit && output == current_output); return current_output_active; } +static void test_release(obs_output_t *output) +{ assert(!after_exit && output != NULL); ++output_releases; } +static void test_queue(enum obs_task_type type, obs_task_t task, void *parameter, bool wait) +{ + assert(!exiting && type == OBS_TASK_UI && !wait && task != NULL); + queued_task = task; queued_parameter = parameter; ++queue_calls; +} ul_plugin_snapshot ul_plugin_get_status(void) { return snapshot; } void ul_plugin_stop_accepting(void) { assert(!enabled); snapshot.status = UL_PLUGIN_CLOSED; record('S'); } void ul_plugin_close(void) { assert(!enabled); snapshot.status = UL_PLUGIN_CLOSED; record('C'); } @@ -65,6 +92,11 @@ void ul_vendor_prepare(obs_data_t *request, obs_data_t *response, void *data) #define obs_frontend_get_main_window_handle test_window #define obs_frontend_add_event_callback test_event #define obs_frontend_add_tools_menu_item test_menu +#define obs_frontend_streaming_active test_active +#define obs_frontend_get_streaming_output test_output +#define obs_output_active test_output_active +#define obs_output_release test_release +#define obs_queue_task test_queue #define obs_websocket_get_api_version test_api #define obs_websocket_register_vendor test_vendor #define obs_websocket_vendor_register_request test_register @@ -77,12 +109,20 @@ static void reset(void) start_calls = register_calls = unregister_calls = tools_calls = event_calls = 0; has_window = allow_start = allow_issue = allow_prepare = allow_unregister = true; enabled = exiting = false; + after_exit = false; api_version = OBS_WEBSOCKET_API_VERSION; snapshot = (ul_plugin_snapshot){UL_PLUGIN_UNPAIRED, UL_PAIRING_MISSING, true, false}; vendor = NULL; issue_registered = prepare_registered = false; order_size = 0; order[0] = 0; + frontend_open = false; queued_arm = 0; stream_busy = false; + arm_scheduler = NULL; queued_task = NULL; queued_parameter = NULL; + checked_generation = 0; + allow_scheduler = has_output = true; + current_active = current_output_active = checked_idle = false; + queue_calls = idle_checks = output_releases = stream_events = 0; + current_output = (obs_output_t *)(uintptr_t)3; } int main(void) @@ -122,6 +162,96 @@ int main(void) obs_module_unload(); assert(strcmp(order, "DC") == 0 && !enabled && unregister_calls == 0); assert(snapshot.status == UL_PLUGIN_CLOSED && event_calls == 0 && tools_calls == 0); - puts("bridge wrapper lifecycle tests passed"); + reset(); allow_scheduler = false; + assert(!obs_module_load() && tools_calls == 0 && event_calls == 0); + reset(); assert(obs_module_load() && arm_scheduler != NULL); + assert(!arm_scheduler(0)); + assert(arm_scheduler(7) && queue_calls == 1); + assert(!arm_scheduler(8) && queue_calls == 1); /* One outstanding slot. */ + queued_task(queued_parameter); + assert(checked_generation == 7 && checked_idle && output_releases == 1); + queued_task(queued_parameter); assert(idle_checks == 1); /* Duplicate callback. */ + assert(arm_scheduler(8)); + check_arm_on_frontend((void *)(uintptr_t)7); + assert(idle_checks == 1 && queued_arm == 8); /* Stale callback cannot consume new slot. */ + current_output_active = true; + queued_task(queued_parameter); + assert(!checked_idle && checked_generation == 8 && output_releases == 2); + current_output_active = false; current_active = true; + assert(arm_scheduler(9)); queued_task(queued_parameter); + assert(!checked_idle && output_releases == 3); + has_output = false; current_active = false; + assert(arm_scheduler(10)); queued_task(queued_parameter); + assert(checked_idle && output_releases == 3); + has_output = true; + frontend_event(OBS_FRONTEND_EVENT_STREAMING_STARTING, NULL); + assert(stream_events == 1 && last_stream_event == UL_STREAM_STARTING); + frontend_event(OBS_FRONTEND_EVENT_STREAMING_STARTED, NULL); + assert(stream_events == 2 && last_stream_event == UL_STREAM_STARTED); + frontend_event(OBS_FRONTEND_EVENT_STREAMING_STOPPING, NULL); + assert(stream_events == 3 && last_stream_event == UL_STREAM_STOPPING); + frontend_event(OBS_FRONTEND_EVENT_STREAMING_STOPPED, NULL); + assert(stream_events == 4 && last_stream_event == UL_STREAM_STOPPED); + assert(arm_scheduler(11)); + exiting = true; + frontend_event(OBS_FRONTEND_EVENT_EXIT, NULL); + queued_task(queued_parameter); /* No OBS or retired runtime access after EXIT. */ + assert(!arm_scheduler(12) && idle_checks == 4); + frontend_event(OBS_FRONTEND_EVENT_STREAMING_STARTED, NULL); + assert(stream_events == 4); + reset(); assert(obs_module_load()); assert(arm_scheduler(13)); + exiting = true; obs_module_unload(); queued_task(queued_parameter); + assert(idle_checks == 0 && !arm_scheduler(14)); + + /* STARTING, STARTED and STOPPING remain busy even while both public active + * queries are false. Only the completed STOPPED lifecycle proves idle. */ + reset(); assert(obs_module_load()); + frontend_event(OBS_FRONTEND_EVENT_STREAMING_STARTING, NULL); + assert(stream_busy); + assert(arm_scheduler(21)); queued_task(queued_parameter); + assert(checked_generation == 21 && !checked_idle && idle_checks == 1); + frontend_event(OBS_FRONTEND_EVENT_STREAMING_STARTED, NULL); + assert(arm_scheduler(22)); queued_task(queued_parameter); + assert(checked_generation == 22 && !checked_idle && idle_checks == 2); + frontend_event(OBS_FRONTEND_EVENT_STREAMING_STOPPING, NULL); + assert(arm_scheduler(23)); queued_task(queued_parameter); + assert(checked_generation == 23 && !checked_idle && idle_checks == 3); + frontend_event(OBS_FRONTEND_EVENT_STREAMING_STOPPED, NULL); + assert(arm_scheduler(24)); queued_task(queued_parameter); + assert(checked_generation == 24 && checked_idle && idle_checks == 4); + + /* A synchronous start failure has no public frontend STOPPED event, so it + * stays fail closed. An unavailable output and repeated STARTING are also + * unknown/busy until OBS reports STOPPED. */ + reset(); assert(obs_module_load()); has_output = false; + frontend_event(OBS_FRONTEND_EVENT_STREAMING_STARTING, NULL); + frontend_event(OBS_FRONTEND_EVENT_STREAMING_STARTING, NULL); + assert(arm_scheduler(25)); queued_task(queued_parameter); + assert(!checked_idle && stream_busy && output_releases == 0); + frontend_event(OBS_FRONTEND_EVENT_STREAMING_STOPPED, NULL); + assert(arm_scheduler(26)); queued_task(queued_parameter); + assert(checked_idle && !stream_busy && output_releases == 0); + + /* EXIT closes before a queued Arm callback; it performs no frontend/output + * query after the final frontend callback returns. */ + reset(); assert(obs_module_load()); + frontend_event(OBS_FRONTEND_EVENT_STREAMING_STARTING, NULL); + assert(arm_scheduler(27)); + exiting = true; + frontend_event(OBS_FRONTEND_EVENT_EXIT, NULL); + after_exit = true; + queued_task(queued_parameter); + assert(idle_checks == 0 && output_releases == 0 && !arm_scheduler(28)); + + /* Unload is the native-only fallback if frontend EXIT was missed. */ + reset(); assert(obs_module_load()); + frontend_event(OBS_FRONTEND_EVENT_STREAMING_STARTING, NULL); + assert(arm_scheduler(29)); + exiting = true; obs_module_unload(); + after_exit = true; + queued_task(queued_parameter); + assert(idle_checks == 0 && output_releases == 0 && !arm_scheduler(30)); + + puts("bridge lifecycle, queued Arm and stream-busy gate tests passed"); return 0; } diff --git a/native/obs-plugin/tests/plugin_state_test.c b/native/obs-plugin/tests/plugin_state_test.c index 025ec2ad..9b65332e 100644 --- a/native/obs-plugin/tests/plugin_state_test.c +++ b/native/obs-plugin/tests/plugin_state_test.c @@ -7,6 +7,7 @@ #include #include "../src/plugin_state.h" +#include "../src/session_protocol.h" static BOOL WINAPI shim_GetModuleHandleExW(DWORD flags, LPCWSTR address, HMODULE *module); @@ -41,6 +42,12 @@ static void shim_ul_authorizer_release(ul_authorizer *authorizer); static void shim_ul_authorizer_destroy(ul_authorizer *authorizer); static int shim_ul_admission_authenticate(ul_admission *admission, DWORD timeout_ms); +static int shim_ul_admission_read_exact(ul_admission *admission, void *buffer, + DWORD size, DWORD timeout_ms); +static int shim_ul_admission_write_all(ul_admission *admission, + const void *buffer, DWORD size, + DWORD timeout_ms); +static int shim_ul_admission_probe(ul_admission *admission, DWORD *available); static void shim_ul_admission_cancel(ul_admission *admission); static HANDLE shim_ul_admission_pipe(const ul_admission *admission); @@ -62,10 +69,14 @@ static HANDLE shim_ul_admission_pipe(const ul_admission *admission); #define ul_authorizer_release shim_ul_authorizer_release #define ul_authorizer_destroy shim_ul_authorizer_destroy #define ul_admission_authenticate shim_ul_admission_authenticate +#define ul_admission_read_exact shim_ul_admission_read_exact +#define ul_admission_write_all shim_ul_admission_write_all +#define ul_admission_probe shim_ul_admission_probe #define ul_admission_cancel shim_ul_admission_cancel #define ul_admission_pipe shim_ul_admission_pipe #define UL_PLUGIN_AUTH_TIMEOUT_MS 40u #define UL_PLUGIN_READY_TIMEOUT_MS 40u +#define UL_PLUGIN_START_TIMEOUT_MS 80u #include "../src/plugin_state.c" #undef GetModuleHandleExW #undef ul_pairing_store_open @@ -85,6 +96,9 @@ static HANDLE shim_ul_admission_pipe(const ul_admission *admission); #undef ul_authorizer_release #undef ul_authorizer_destroy #undef ul_admission_authenticate +#undef ul_admission_read_exact +#undef ul_admission_write_all +#undef ul_admission_probe #undef ul_admission_cancel #undef ul_admission_pipe @@ -95,12 +109,16 @@ struct ul_pairing_store { struct ul_admission { HANDLE cancelled; int authenticate_result; + bool command_ready; + uint8_t command[UL_SESSION_COMMAND_BYTES]; }; struct ul_authorizer { bool outstanding; bool revoked; ul_admission *admission; + uint8_t session[16]; + uint8_t mask; }; static bool fake_pin = true; @@ -119,6 +137,15 @@ static HANDLE fake_issue_release; static volatile LONG fake_release_count; static volatile LONG fake_destroy_count; static volatile LONG fake_store_destroy_count; +static bool fake_arm_command; +static int fake_arm_command_corrupt; +static bool fake_schedule_result = true; +static DWORD fake_probe_available; +static int fake_probe_result = UL_ADMISSION_AUTH_OK; +static HANDLE fake_schedule_entered; +static HANDLE fake_reply_written; +static uintptr_t fake_scheduled_generation; +static uint8_t fake_reply[UL_SESSION_COMMAND_BYTES]; static int check(bool condition, const char *message) { @@ -265,6 +292,8 @@ static bool shim_ul_authorizer_issue(ul_authorizer *authorizer, return false; memset(out_challenge, 0x33, 60); authorizer->outstanding = true; + memcpy(authorizer->session, session, sizeof(authorizer->session)); + authorizer->mask = mask; return true; } @@ -291,8 +320,22 @@ static ul_admission *shim_ul_authorizer_prepare( return NULL; } admission->authenticate_result = fake_authenticate_result; + if (fake_arm_command) { + memcpy(admission->command, "ULAC", 4); + admission->command[4] = 1; + admission->command[5] = 1; + memcpy(admission->command + 8, authorizer->session, 16); + admission->command[24] = authorizer->mask; + if (fake_arm_command_corrupt == 1) + admission->command[8] ^= 0x80; + else if (fake_arm_command_corrupt == 2) + admission->command[24] ^= 1; + admission->command_ready = true; + } authorizer->admission = admission; options->client_pid = 100; + memcpy(options->session, authorizer->session, sizeof(options->session)); + options->additional_mix_mask = authorizer->mask; return admission; } @@ -335,6 +378,49 @@ static int shim_ul_admission_authenticate(ul_admission *admission, return admission->authenticate_result; } +static int shim_ul_admission_read_exact(ul_admission *admission, void *buffer, + DWORD size, DWORD timeout_ms) +{ + DWORD waited; + + if (buffer == NULL || size != sizeof(admission->command)) + return UL_ADMISSION_IO_ERROR; + if (admission->command_ready) { + memcpy(buffer, admission->command, size); + admission->command_ready = false; + return UL_ADMISSION_AUTH_OK; + } + waited = WaitForSingleObject(admission->cancelled, timeout_ms); + SecureZeroMemory(buffer, size); + return waited == WAIT_OBJECT_0 ? UL_ADMISSION_CANCELLED + : UL_ADMISSION_TIMEOUT; +} + +static int shim_ul_admission_write_all(ul_admission *admission, + const void *buffer, DWORD size, + DWORD timeout_ms) +{ + (void)timeout_ms; + if (WaitForSingleObject(admission->cancelled, 0) == WAIT_OBJECT_0) + return UL_ADMISSION_CANCELLED; + if (buffer == NULL || size != sizeof(fake_reply)) + return UL_ADMISSION_IO_ERROR; + memcpy(fake_reply, buffer, size); + if (fake_reply_written != NULL) + SetEvent(fake_reply_written); + return UL_ADMISSION_AUTH_OK; +} + +static int shim_ul_admission_probe(ul_admission *admission, DWORD *available) +{ + if (WaitForSingleObject(admission->cancelled, 0) == WAIT_OBJECT_0) + return UL_ADMISSION_CANCELLED; + if (available == NULL) + return UL_ADMISSION_IO_ERROR; + *available = fake_probe_available; + return fake_probe_result; +} + static void shim_ul_admission_cancel(ul_admission *admission) { SetEvent(admission->cancelled); @@ -346,6 +432,14 @@ static HANDLE shim_ul_admission_pipe(const ul_admission *admission) return NULL; } +static bool queued_arm_scheduler(uintptr_t generation) +{ + fake_scheduled_generation = generation; + if (fake_schedule_entered != NULL) + SetEvent(fake_schedule_entered); + return fake_schedule_result; +} + static DWORD WINAPI export_thread(void *unused) { (void)unused; @@ -382,6 +476,337 @@ static bool bytes_are_zero(const uint8_t *bytes, size_t size) return true; } +static bool wait_phase(ul_session_phase expected, DWORD timeout_ms) +{ + ULONGLONG deadline = GetTickCount64() + timeout_ms; + + do { + if (ul_plugin_session_status() == expected) + return true; + Sleep(1); + } while (GetTickCount64() < deadline); + return ul_plugin_session_status() == expected; +} + +static bool wait_worker_inactive(DWORD timeout_ms) +{ + ULONGLONG deadline = GetTickCount64() + timeout_ms; + ul_plugin_snapshot snapshot; + + do { + snapshot = ul_plugin_get_status(); + if (!snapshot.admission_pending) { + HANDLE worker = plugin_global.runtime == NULL + ? NULL + : plugin_global.runtime->worker; + return worker == NULL || + WaitForSingleObject(worker, timeout_ms) == WAIT_OBJECT_0; + } + Sleep(1); + } while (GetTickCount64() < deadline); + if (ul_plugin_get_status().admission_pending) + return false; + return plugin_global.runtime == NULL || plugin_global.runtime->worker == NULL || + WaitForSingleObject(plugin_global.runtime->worker, timeout_ms) == + WAIT_OBJECT_0; +} + +static bool wait_current_worker_signaled(DWORD timeout_ms) +{ + ul_plugin_runtime *runtime = plugin_global.runtime; + HANDLE worker = runtime == NULL ? NULL : runtime->worker; + + return worker != NULL && + WaitForSingleObject(worker, timeout_ms) == WAIT_OBJECT_0; +} + +static int prepare_arm(const uint8_t session[16], uint8_t mask) +{ + uint8_t challenge[60]; + uint8_t proof[32] = {0x44}; + int failures = 0; + + ResetEvent(fake_schedule_entered); + ResetEvent(fake_reply_written); + SecureZeroMemory(fake_reply, sizeof(fake_reply)); + fake_scheduled_generation = 0; + failures += check(ul_plugin_issue(100, session, mask, challenge), + "Arm issue succeeds"); + failures += check(ul_plugin_prepare(challenge, sizeof(challenge), proof, + sizeof(proof)), + "Arm prepare succeeds"); + failures += check(WaitForSingleObject(fake_schedule_entered, 2000) == + WAIT_OBJECT_0, + "Arm scheduler receives queued generation"); + failures += check(fake_scheduled_generation != 0, + "Arm scheduler receives nonzero generation"); + failures += check(wait_phase(UL_SESSION_ARM_PENDING, 2000), + "Arm reaches pending phase"); + return failures; +} + +static int start_arm_runtime(void) +{ + int failures = 0; + + fake_load_result = UL_PAIRING_OK; + fake_arm_command = true; + fake_schedule_entered = CreateEventW(NULL, TRUE, FALSE, NULL); + fake_reply_written = CreateEventW(NULL, TRUE, FALSE, NULL); + failures += check(fake_schedule_entered != NULL && + fake_reply_written != NULL, + "create deterministic Arm events"); + failures += check(ul_plugin_start(), "Arm runtime starts paired"); + failures += check(ul_plugin_set_arm_scheduler(queued_arm_scheduler), + "install Arm scheduler once before worker"); + return failures; +} + +static void close_arm_events(void) +{ + if (fake_schedule_entered != NULL) + CloseHandle(fake_schedule_entered); + if (fake_reply_written != NULL) + CloseHandle(fake_reply_written); + fake_schedule_entered = NULL; + fake_reply_written = NULL; +} + +static int scenario_arm_valid(void) +{ + uint8_t session[16] = {0x21, 2, 3, 4}; + uint8_t expected[UL_SESSION_COMMAND_BYTES]; + ul_plugin_snapshot snapshot; + int failures = start_arm_runtime(); + + failures += prepare_arm(session, 37); + ul_plugin_arm_checked(fake_scheduled_generation, true); + failures += check(WaitForSingleObject(fake_reply_written, 2000) == + WAIT_OBJECT_0, + "committed Arm reply written"); + failures += check(ul_session_arm_reply(session, 37, true, expected) && + memcmp(fake_reply, expected, sizeof(expected)) == 0, + "Arm reply binds exact prepared session and mask"); + failures += check(wait_phase(UL_SESSION_ARMED, 2000), + "idle callback commits Arm"); + Sleep(120); + snapshot = ul_plugin_get_status(); + failures += check(snapshot.admission_pending && + ul_plugin_session_status() == UL_SESSION_ARMED, + "armed session persists beyond READY timeout"); + ul_plugin_stream_event(UL_STREAM_STARTING); + failures += check(ul_plugin_session_status() == UL_SESSION_STARTING, + "post-Arm STARTING accepted"); + ul_plugin_stream_event(UL_STREAM_STARTED); + failures += check(ul_plugin_session_status() == UL_SESSION_STARTED, + "ordered STARTED accepted once"); + ul_plugin_stream_event(UL_STREAM_STOPPING); + failures += check(wait_phase(UL_SESSION_TERMINAL, 2000), + "stop terminates active consent"); + ul_plugin_close(); + close_arm_events(); + return failures; +} + +static int scenario_arm_ordering(void) +{ + uint8_t session[16] = {0x31}; + uint8_t expected[UL_SESSION_COMMAND_BYTES]; + int failures = start_arm_runtime(); + + failures += prepare_arm(session, 3); + ul_plugin_arm_checked(fake_scheduled_generation, false); + failures += check(WaitForSingleObject(fake_reply_written, 2000) == + WAIT_OBJECT_0, + "busy Arm writes refusal"); + failures += check(ul_session_arm_reply(session, 3, false, expected) && + memcmp(fake_reply, expected, sizeof(expected)) == 0, + "already-busy Arm reply is exact refusal"); + failures += check(wait_worker_inactive(2000), + "busy refusal retires worker"); + + session[0]++; + failures += prepare_arm(session, 4); + ul_plugin_stream_event(UL_STREAM_STARTING); + failures += check(wait_phase(UL_SESSION_TERMINAL, 2000), + "STARTING before callback refuses pending Arm"); + failures += check(WaitForSingleObject(fake_reply_written, 2000) == + WAIT_OBJECT_0 && + ul_session_arm_reply(session, 4, false, expected) && + fake_reply[25] == 2 && + memcmp(fake_reply, expected, sizeof(expected)) == 0, + "STARTING race returns exact status-2 refusal"); + ul_plugin_arm_checked(fake_scheduled_generation, true); + failures += check(wait_worker_inactive(2000), + "late callback cannot revive STARTING race"); + + session[0]++; + failures += prepare_arm(session, 5); + ul_plugin_arm_checked(fake_scheduled_generation, true); + failures += check(WaitForSingleObject(fake_reply_written, 2000) == + WAIT_OBJECT_0 && + wait_phase(UL_SESSION_ARMED, 2000), + "fresh Arm commits after retired race"); + ul_plugin_stream_event(UL_STREAM_STARTED); + failures += check(wait_phase(UL_SESSION_TERMINAL, 2000), + "STARTED without STARTING is terminal"); + ul_plugin_close(); + close_arm_events(); + return failures; +} + +static int scenario_arm_stale_and_commands(void) +{ + uint8_t session[16] = {0x41}; + uintptr_t stale_generation; + int failures = start_arm_runtime(); + + failures += prepare_arm(session, 6); + stale_generation = fake_scheduled_generation; + failures += check(wait_worker_inactive(2000), + "unanswered Arm expires at READY deadline"); + session[0]++; + failures += prepare_arm(session, 7); + ul_plugin_arm_checked(stale_generation, true); + failures += check(ul_plugin_session_status() == UL_SESSION_ARM_PENDING, + "late callback cannot arm new generation"); + ul_plugin_arm_checked(fake_scheduled_generation, true); + failures += check(WaitForSingleObject(fake_reply_written, 2000) == + WAIT_OBJECT_0 && + wait_phase(UL_SESSION_ARMED, 2000), + "current generation still arms"); + fake_probe_available = UL_SESSION_COMMAND_BYTES; + failures += check(wait_worker_inactive(2000), + "duplicate command bytes terminate armed worker"); + + fake_probe_available = 0; + session[0]++; + failures += prepare_arm(session, 8); + ul_plugin_arm_checked(fake_scheduled_generation, true); + failures += check(WaitForSingleObject(fake_reply_written, 2000) == + WAIT_OBJECT_0, + "EOF case Arm reply written"); + fake_probe_result = UL_ADMISSION_IO_ERROR; + failures += check(wait_worker_inactive(2000), + "authenticated EOF terminates armed worker"); + ul_plugin_close(); + close_arm_events(); + return failures; +} + +static int scenario_arm_start_timeout(void) +{ + uint8_t session[16] = {0x39}; + int failures = start_arm_runtime(); + + failures += prepare_arm(session, 15); + ul_plugin_arm_checked(fake_scheduled_generation, true); + failures += check(WaitForSingleObject(fake_reply_written, 2000) == + WAIT_OBJECT_0 && + wait_phase(UL_SESSION_ARMED, 2000), + "start-timeout session arms"); + ul_plugin_stream_event(UL_STREAM_STARTING); + failures += check(ul_plugin_session_status() == UL_SESSION_STARTING, + "start-timeout session records STARTING"); + failures += check(wait_current_worker_signaled(2000) && + ul_plugin_session_status() == UL_SESSION_TERMINAL, + "missing matching STARTED terminally expires worker"); + ul_plugin_stream_event(UL_STREAM_STARTED); + failures += check(ul_plugin_session_status() == UL_SESSION_TERMINAL, + "late STARTED cannot resurrect expired generation"); + + session[0]++; + failures += prepare_arm(session, 16); + ul_plugin_arm_checked(fake_scheduled_generation, true); + failures += check(WaitForSingleObject(fake_reply_written, 2000) == + WAIT_OBJECT_0 && + wait_phase(UL_SESSION_ARMED, 2000), + "new manual session arms after STARTING timeout"); + ul_plugin_close(); + close_arm_events(); + return failures; +} + +static int scenario_arm_bound_command(void) +{ + uint8_t session[16] = {0x51}; + uint8_t challenge[60]; + uint8_t proof[32] = {0x44}; + int failures = start_arm_runtime(); + + fake_arm_command_corrupt = 1; + failures += check(ul_plugin_issue(100, session, 9, challenge) && + ul_plugin_prepare(challenge, sizeof(challenge), proof, + sizeof(proof)), + "submit wrong-session Arm command"); + failures += check(wait_worker_inactive(2000), + "wrong prepared session is refused"); + fake_arm_command_corrupt = 2; + session[0]++; + failures += check(ul_plugin_issue(100, session, 10, challenge) && + ul_plugin_prepare(challenge, sizeof(challenge), proof, + sizeof(proof)), + "submit wrong-mask Arm command"); + failures += check(wait_worker_inactive(2000), + "wrong prepared mix mask is refused"); + ul_plugin_close(); + close_arm_events(); + return failures; +} + +static int scenario_arm_teardown(void) +{ + uint8_t session[16] = {0x61}; + bool saved = false; + LONG releases; + int failures = start_arm_runtime(); + + failures += prepare_arm(session, 11); + releases = InterlockedCompareExchange(&fake_release_count, 0, 0); + failures += check(ul_plugin_forget() == UL_PAIRING_OK, + "forget completes while Arm callback pending"); + failures += check(InterlockedCompareExchange(&fake_release_count, 0, 0) > + releases, + "forget cancels and joins pending Arm worker"); + failures += check(ul_plugin_pair(L"C:\\new.ulpair", false, &saved) == + UL_PAIRING_OK && saved, + "re-pair after forget"); + + session[0]++; + failures += prepare_arm(session, 12); + releases = InterlockedCompareExchange(&fake_release_count, 0, 0); + saved = false; + failures += check(ul_plugin_pair(L"C:\\replace.ulpair", true, &saved) == + UL_PAIRING_OK && saved, + "replace completes while Arm callback pending"); + failures += check(InterlockedCompareExchange(&fake_release_count, 0, 0) > + releases, + "replace cancels and joins pending Arm worker"); + + session[0]++; + failures += prepare_arm(session, 13); + ul_plugin_stop_accepting(); + failures += check(wait_worker_inactive(2000), + "stop accepting cancels pending Arm worker"); + ul_plugin_close(); + close_arm_events(); + return failures; +} + +static int scenario_arm_close_pending(void) +{ + uint8_t session[16] = {0x71}; + int failures = start_arm_runtime(); + + failures += prepare_arm(session, 14); + ul_plugin_close(); + failures += check(ul_plugin_get_status().status == UL_PLUGIN_CLOSED && + ul_plugin_session_status() == UL_SESSION_NONE, + "close joins pending Arm and closes public gate"); + close_arm_events(); + return failures; +} + static int scenario_pin_failure(void) { ul_plugin_snapshot snapshot; @@ -770,6 +1195,20 @@ static int run_child(const char *scenario) return scenario_close_cancels_mutation(); if (strcmp(scenario, "dispatch") == 0) return scenario_close_drains_issue(); + if (strcmp(scenario, "arm-valid") == 0) + return scenario_arm_valid(); + if (strcmp(scenario, "arm-order") == 0) + return scenario_arm_ordering(); + if (strcmp(scenario, "arm-start-timeout") == 0) + return scenario_arm_start_timeout(); + if (strcmp(scenario, "arm-stale") == 0) + return scenario_arm_stale_and_commands(); + if (strcmp(scenario, "arm-bound") == 0) + return scenario_arm_bound_command(); + if (strcmp(scenario, "arm-teardown") == 0) + return scenario_arm_teardown(); + if (strcmp(scenario, "arm-close") == 0) + return scenario_arm_close_pending(); return 1; } @@ -802,7 +1241,9 @@ int main(int argc, char **argv) { static const wchar_t *scenarios[] = { L"pin", L"reload", L"fallback", L"states", L"pair", L"prepare", - L"expiry", L"cancel", L"dispatch"}; + L"expiry", L"cancel", L"dispatch", L"arm-valid", L"arm-order", + L"arm-start-timeout", L"arm-stale", L"arm-bound", L"arm-teardown", + L"arm-close"}; wchar_t executable[32768]; size_t index; int failures = 0; diff --git a/native/obs-plugin/tests/session_protocol_test.c b/native/obs-plugin/tests/session_protocol_test.c new file mode 100644 index 00000000..dc9ae50c --- /dev/null +++ b/native/obs-plugin/tests/session_protocol_test.c @@ -0,0 +1,49 @@ +// SPDX-License-Identifier: GPL-2.0-or-later +#include "../src/session_protocol.h" + +#include +#include +#include + +int main(void) +{ + const uint8_t session[16] = {1, 2, 3, 4, 5, 6, 7, 8, + 9, 10, 11, 12, 13, 14, 15, 16}; + const uint8_t fixture[28] = {'U', 'L', 'A', 'C', 1, 1, 0, 0, + 1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 13, 14, 15, 16, 37, 0, 0, 0}; + uint8_t modified[29], reply[28], expected[28], zero[28] = {0}; + size_t index; + + assert(ul_session_arm_request(fixture, sizeof(fixture), session, 37)); + assert(!ul_session_arm_request(NULL, sizeof(fixture), session, 37)); + assert(!ul_session_arm_request(fixture, sizeof(fixture), NULL, 37)); + assert(!ul_session_arm_request(fixture, sizeof(fixture), zero, 37)); + assert(!ul_session_arm_request(fixture, sizeof(fixture), session, 64)); + for (index = 0; index < sizeof(fixture); ++index) { + assert(!ul_session_arm_request(fixture, index, session, 37)); + memcpy(modified, fixture, sizeof(fixture)); + modified[index] ^= 1u; + assert(!ul_session_arm_request(modified, sizeof(fixture), session, 37)); + } + memcpy(modified, fixture, sizeof(fixture)); + modified[28] = 0; + assert(!ul_session_arm_request(modified, sizeof(modified), session, 37)); + + memcpy(expected, fixture, sizeof(expected)); + expected[5] = 2; + expected[25] = 1; + assert(ul_session_arm_reply(session, 37, true, reply)); + assert(memcmp(reply, expected, sizeof(reply)) == 0); + expected[25] = 2; + assert(ul_session_arm_reply(session, 37, false, reply)); + assert(memcmp(reply, expected, sizeof(reply)) == 0); + assert(!ul_session_arm_request(reply, sizeof(reply), session, 37)); + assert(!ul_session_arm_reply(session, 64, true, reply)); + assert(memcmp(reply, zero, sizeof(reply)) == 0); + memset(reply, 255, sizeof(reply)); + assert(!ul_session_arm_reply(zero, 0, true, reply)); + assert(memcmp(reply, zero, sizeof(reply)) == 0); + assert(!ul_session_arm_reply(session, 0, true, NULL)); + puts("fixed Arm request/reply contract passed"); + return 0; +} diff --git a/native/obs-plugin/tools/build.py b/native/obs-plugin/tools/build.py index f6de4833..02c05c98 100644 --- a/native/obs-plugin/tools/build.py +++ b/native/obs-plugin/tools/build.py @@ -15,7 +15,7 @@ ROOT = Path(__file__).resolve().parents[1] SOURCES = ( "bridge", "plugin_state", "pairing_ui", "vendor_dispatch", "pairing_store", - "authorization", "admission", "handshake", "crypto", + "authorization", "admission", "handshake", "crypto", "session_protocol", ) ISC_HEADERS = ( "callback/calldata.h", "callback/proc.h", "callback/signal.h", diff --git a/native/obs-plugin/tools/test_native.py b/native/obs-plugin/tools/test_native.py index 1a8fdb8e..b6957cec 100644 --- a/native/obs-plugin/tools/test_native.py +++ b/native/obs-plugin/tools/test_native.py @@ -44,6 +44,8 @@ def main() -> None: logs = [] sources = [ROOT / name for name in ( "src/handshake.c", "src/handshake.h", "src/admission.c", "src/admission.h", + "src/session_protocol.c", "src/session_protocol.h", "tests/session_protocol_test.c", + "tests/admission_io_test.c", "tests/admission_io_fault_test.c", "src/crypto.c", "src/crypto.h", "src/authorization.c", "src/authorization.h", "src/pairing_store.c", "src/pairing_store.h", @@ -118,7 +120,23 @@ def run(name: str, arguments: list[str | Path], timeout: int = 60) -> None: pairing_dll = output / "utterleaf-pairing-store-test.dll" pairing_state = output / "pairing_store_test.exe" plugin_state = output / "plugin_state_test.exe" + session_protocol = output / "session_protocol_test.exe" + admission_io = output / "admission_io_test.exe" + admission_io_fault = output / "admission_io_fault_test.exe" run("compiler", [compiler, "--version"]) + run("session-protocol-build", [*flags, ROOT / "src/session_protocol.c", + ROOT / "tests/session_protocol_test.c", "-o", session_protocol]) + run("session-protocol-test", [session_protocol]) + run("admission-io-build", [*flags, "-D_M_X64=100", "-municode", ROOT / "src/admission.c", + ROOT / "src/crypto.c", ROOT / "src/handshake.c", + ROOT / "tests/admission_io_test.c", "-ladvapi32", "-lbcrypt", + "-o", admission_io]) + run("admission-io-test", [admission_io]) + run("admission-io-fault-build", [*flags, "-D_M_X64=100", ROOT / "src/crypto.c", + ROOT / "src/handshake.c", + ROOT / "tests/admission_io_fault_test.c", + "-ladvapi32", "-lbcrypt", "-o", admission_io_fault]) + run("admission-io-fault-test", [admission_io_fault]) run("handshake-build", [*flags, ROOT / "src/crypto.c", ROOT / "src/handshake.c", ROOT / "tests/handshake_test.c", "-lbcrypt", "-o", fixed]) run("handshake-test", [fixed]) @@ -129,6 +147,7 @@ def run(name: str, arguments: list[str | Path], timeout: int = 60) -> None: "-o", authorization_state]) run("authorization-state-test", [authorization_state]) run("plugin-state-build", [*flags, "-D_M_X64=100", ROOT / "tests/plugin_state_test.c", + ROOT / "src/session_protocol.c", "-o", plugin_state]) run("plugin-state-test", [plugin_state]) functions = ( @@ -173,7 +192,8 @@ def run(name: str, arguments: list[str | Path], timeout: int = 60) -> None: run("pairing-interop-test", [sys._base_executable, ROOT / "tests/test_pairing_interop.py", pairing_dll, authorization_dll]) artifacts = [fixed, fault, identity, dll, crypto, authorization_state, authorization_dll, - pairing_state, pairing_dll, plugin_state] + pairing_state, pairing_dll, plugin_state, session_protocol, admission_io, + admission_io_fault] if args.build is not None: bridge_test = output / "bridge_test.exe" run("bridge-wrapper-build", [*flags, "-D_M_X64=100", f"-I{headers / 'libobs'}", @@ -201,7 +221,7 @@ def run(name: str, arguments: list[str | Path], timeout: int = 60) -> None: if any(digest(path) != expected for path, expected in dispatch_inputs.items()): raise RuntimeError("Reviewed dispatch inputs changed during verification") receipt = { - "schema": 2, "scope": "pairing/admission/runtime; optional parsed libobs dispatch; no OBS application, arming or audio", + "schema": 2, "scope": "pairing/admission/Arm runtime fixtures; optional parsed libobs dispatch; no OBS application or audio", "vendor_dispatch": "passed" if args.build is not None else "not run: supply --build and --headers", "native_dialog": "passed" if args.ui else "not run: supply --ui on a Windows desktop", "dispatch_inputs": {str(path): expected for path, expected in dispatch_inputs.items()}, diff --git a/tests/test_obs_audio_arm.py b/tests/test_obs_audio_arm.py new file mode 100644 index 00000000..928046c8 --- /dev/null +++ b/tests/test_obs_audio_arm.py @@ -0,0 +1,163 @@ +"""Focused Arm codec checks, independent of OBS or a microphone.""" +import struct +import threading +import time + +import pytest + +from utterleaf import obs_audio_pipe as audio_pipe +from utterleaf import obs_protocol as protocol +from test_obs_audio_pipe import Peer, Pipe, SESSION + + +def opened(): + return audio_pipe.ObsAudioPipe(ArmPipe(), Peer(), SESSION, lambda: False) + + +class ArmPipe(Pipe): + def write_all(self, data, *, deadline): + self.writes.append(data) + + +ARM_REPLY = struct.Struct("<4sBBH16sBBH") + + +def accepted(pipe, *, mask=0, trailing=b""): + pipe.pending.extend(ARM_REPLY.pack(b"ULAC", 1, 2, 0, SESSION, mask, 1, 0)) + pipe.pending.extend(trailing) + + +def test_arm_writes_exact_record_and_accepts_matching_reply(): + pipe = ArmPipe() + connection = audio_pipe.ObsAudioPipe(pipe, Peer(), SESSION, lambda: False) + pipe.pending.extend(struct.pack("<4sBBH16sBBH", b"ULAC", 1, 2, 0, + SESSION, 7, 1, 0)) + connection.arm(additional_mix_mask=7, deadline=time.monotonic() + 1) + assert pipe.writes == [struct.pack("<4sBBH16sBBH", b"ULAC", 1, 1, 0, + SESSION, 7, 0, 0)] + connection.close() + + +@pytest.mark.parametrize("reply", [ + struct.pack("<4sBBH16sBBH", b"NOPE", 1, 2, 0, SESSION, 0, 1, 0), + struct.pack("<4sBBH16sBBH", b"ULAC", 9, 2, 0, SESSION, 0, 1, 0), + struct.pack("<4sBBH16sBBH", b"ULAC", 1, 9, 0, SESSION, 0, 1, 0), + struct.pack("<4sBBH16sBBH", b"ULAC", 1, 2, 1, SESSION, 0, 1, 0), + struct.pack("<4sBBH16sBBH", b"ULAC", 1, 2, 0, b"x" * 16, 0, 1, 0), + struct.pack("<4sBBH16sBBH", b"ULAC", 1, 2, 0, SESSION, 9, 1, 0), + struct.pack("<4sBBH16sBBH", b"ULAC", 1, 2, 0, SESSION, 0, 2, 0), + struct.pack("<4sBBH16sBBH", b"ULAC", 1, 2, 0, SESSION, 0, 3, 0), + struct.pack("<4sBBH16sBBH", b"ULAC", 1, 2, 0, SESSION, 0, 1, 1), +]) +def test_arm_rejects_bad_reply(reply): + pipe = ArmPipe() + pipe.pending.extend(reply) + connection = audio_pipe.ObsAudioPipe(pipe, Peer(), SESSION, lambda: False) + with pytest.raises(audio_pipe.ObsAudioPipeError): + connection.arm(deadline=time.monotonic() + 1) + assert pipe.closed + + +def test_read_frames_requires_arm_and_closes(): + connection = opened() + with pytest.raises(audio_pipe.ObsAudioPipeError): + connection.read_frames(deadline=time.monotonic() + 1) + assert connection._closed.is_set() + + +def test_arm_is_single_use_and_coalesced_audio_remains_for_reader(): + pipe = ArmPipe() + connection = audio_pipe.ObsAudioPipe(pipe, Peer(), SESSION, lambda: False) + frame = protocol.encode_frame(protocol.StartFrame(SESSION, 16000, 2, 4, 9000000000)) + accepted(pipe, mask=3, trailing=frame) + connection.arm(additional_mix_mask=3, deadline=time.monotonic() + 1) + assert pipe.pending == frame + assert connection.read_frames(deadline=time.monotonic() + 1) == [ + protocol.StartFrame(SESSION, 16000, 2, 4, 9000000000) + ] + with pytest.raises(audio_pipe.ObsAudioPipeError): + connection.arm(additional_mix_mask=3, deadline=time.monotonic() + 1) + assert pipe.closed and len(pipe.writes) == 1 + + +@pytest.mark.parametrize("mask", [-1, 64, True, "0"]) +def test_arm_bad_mask_is_terminal_before_dispatch(mask): + pipe = ArmPipe() + connection = audio_pipe.ObsAudioPipe(pipe, Peer(), SESSION, lambda: False) + with pytest.raises(audio_pipe.ObsAudioPipeError): + connection.arm(additional_mix_mask=mask, deadline=time.monotonic() + 1) + assert pipe.closed and pipe.writes == [] + + +def test_arm_deadline_and_peer_failure_are_terminal(): + pipe = ArmPipe() + peer = Peer() + connection = audio_pipe.ObsAudioPipe(pipe, peer, SESSION, lambda: False) + with pytest.raises(audio_pipe.ObsAudioPipeError): + connection.arm(deadline=time.monotonic() - 1) + assert pipe.closed and pipe.writes == [] + + pipe = Pipe() + peer = Peer() + peer.fail_at = len(peer.checks) + 2 + connection = audio_pipe.ObsAudioPipe(pipe, peer, SESSION, lambda: False) + with pytest.raises(audio_pipe.ObsAudioPipeError): + connection.arm(deadline=time.monotonic() + 1) + assert pipe.closed and peer.closed + + pipe = ArmPipe() + peer = Peer() + peer.fail_at = 1 + connection = audio_pipe.ObsAudioPipe(pipe, peer, SESSION, lambda: False) + with pytest.raises(audio_pipe.ObsAudioPipeError): + connection.arm(deadline=time.monotonic() + 1) + assert pipe.closed and pipe.writes == [] + + +def test_arm_cancellation_after_dispatch_is_terminal(): + pipe = Pipe() + stopped = [] + pipe.after_write = lambda: stopped.append(True) + connection = audio_pipe.ObsAudioPipe(pipe, Peer(), SESSION, + lambda: bool(stopped)) + with pytest.raises(audio_pipe.ObsAudioPipeCancelled): + connection.arm(deadline=time.monotonic() + 1) + assert pipe.closed and len(pipe.writes) == 1 + + +class BlockingArmPipe(ArmPipe): + def __init__(self): + super().__init__() + self.released = threading.Event() + + def read(self, maximum, *, deadline): + self.released.wait(2) + raise audio_pipe.ObsAudioPipeCancelled("closed") + + def close(self): + self.closed = True + self.released.set() + + +def test_close_interrupts_blocked_arm_without_second_dispatch(): + pipe = BlockingArmPipe() + connection = audio_pipe.ObsAudioPipe(pipe, Peer(), SESSION, lambda: False) + result = [] + + def worker(): + try: + connection.arm(deadline=time.monotonic() + 5) + except BaseException as exc: + result.append(exc) + + thread = threading.Thread(target=worker) + thread.start() + for _ in range(100): + if pipe.writes: + break + time.sleep(0.01) + connection.close() + thread.join(2) + assert not thread.is_alive() + assert len(pipe.writes) == 1 + assert result and isinstance(result[0], audio_pipe.ObsAudioPipeCancelled) diff --git a/tests/test_obs_audio_pipe.py b/tests/test_obs_audio_pipe.py index b0d2df18..5449bc75 100644 --- a/tests/test_obs_audio_pipe.py +++ b/tests/test_obs_audio_pipe.py @@ -60,11 +60,24 @@ def server_pid(self): def write_all(self, data, *, deadline): self.writes.append(data) # Independent server record construction, not the client's private helpers. - assert len(data) == 56 - assert data[:8] == b"ULAH\x01\x01\x00\x00" - mac = hmac.new(data[24:], DOMAIN + data[:24], hashlib.sha256).digest() - ack = b"ULAH\x01\x02\x00\x00" + data[8:24] + mac - self.pending.extend(self.mutate_ack(ack)) + if len(data) == 56: + assert data[:8] == b"ULAH\x01\x01\x00\x00" + mac = hmac.new(data[24:], DOMAIN + data[:24], hashlib.sha256).digest() + ack = b"ULAH\x01\x02\x00\x00" + data[8:24] + mac + self.pending.extend(self.mutate_ack(ack)) + elif len(data) == 28: + magic, version, kind, reserved, session, mask, status, trailing = struct.unpack( + "<4sBBH16sBBH", data + ) + assert (magic, version, kind, reserved, session, status, trailing) == ( + b"ULAC", 1, 1, 0, SESSION, 0, 0 + ) + assert 0 <= mask <= 63 + self.pending.extend(struct.pack( + "<4sBBH16sBBH", b"ULAC", 1, 2, 0, session, mask, 1, 0 + )) + else: + raise AssertionError(f"unexpected protocol record length {len(data)}") self.after_write() def read(self, maximum, *, deadline): @@ -91,6 +104,10 @@ def open_pipe(name, **controls): return result, pipe, peer +def arm(connection): + connection.arm(deadline=time.monotonic() + 2) + + def test_fragmented_ack_has_exact_role_bound_mac_and_retains_no_secret(monkeypatch): pipe = Pipe() pipe.max_chunk = 3 @@ -202,6 +219,7 @@ def audio_frames(session=SESSION): def test_fragmented_audio_retains_receiver_consent_and_actual_primary_bus(monkeypatch): connection, pipe, peer = connect(monkeypatch) + arm(connection) frames = audio_frames() pipe.pending.extend(b"".join(protocol.encode_frame(frame) for frame in frames)) pipe.max_chunk = 31 @@ -231,6 +249,7 @@ def test_fragmented_audio_retains_receiver_consent_and_actual_primary_bus(monkey def test_authenticated_pipe_does_not_grant_permission_to_capture(monkeypatch): connection, pipe, _ = connect(monkeypatch) + arm(connection) pipe.pending.extend(protocol.encode_frame(audio_frames()[0])) receiver = ObsCaptureSession(SESSION, stream_active=False, store_factory=lambda rate: pytest.fail("Unarmed storage")) @@ -246,6 +265,7 @@ def test_authenticated_pipe_does_not_grant_permission_to_capture(monkeypatch): @pytest.mark.parametrize("suffix", [b"U", protocol.encode_frame(audio_frames()[0])]) def test_end_with_trailing_partial_or_complete_frame_fails(monkeypatch, suffix): connection, pipe, peer = connect(monkeypatch) + arm(connection) pipe.pending.extend(protocol.encode_frame(audio_frames()[-1]) + suffix) with pytest.raises(audio_pipe.ObsAudioPipeError): connection.read_frames(deadline=time.monotonic() + 1) @@ -254,6 +274,7 @@ def test_end_with_trailing_partial_or_complete_frame_fails(monkeypatch, suffix): def test_identity_loss_after_read_prevents_returning_audio(monkeypatch): connection, pipe, peer = connect(monkeypatch) + arm(connection) pipe.pending.extend(protocol.encode_frame(audio_frames()[0])) peer.fail_at = len(peer.checks) + 2 with pytest.raises(audio_pipe.ObsAudioPipeError): @@ -263,6 +284,7 @@ def test_identity_loss_after_read_prevents_returning_audio(monkeypatch): def test_foreign_audio_session_is_terminal(monkeypatch): connection, pipe, peer = connect(monkeypatch) + arm(connection) pipe.pending.extend(protocol.encode_frame(protocol.StartFrame(b"x" * 16, 16000, 0, 1, 0))) with pytest.raises(audio_pipe.ObsAudioPipeError): connection.read_frames(deadline=time.monotonic() + 1) @@ -274,7 +296,8 @@ def test_native_child_pipe_auth_and_receiver_continue_after_tcp_loss(): session = secrets.token_bytes(16) expected_frames = audio_frames(session) server = start_native_pipe_server( - "audio", b"".join(protocol.encode_frame(frame) for frame in expected_frames), + "audio", + b"".join(protocol.encode_frame(frame) for frame in expected_frames), session_id=session, ) connection = retained = None @@ -291,10 +314,11 @@ def test_native_child_pipe_auth_and_receiver_continue_after_tcp_loss(): deadline=time.monotonic() + 2) connection = audio_pipe.connect(session, retained, cancelled=lambda: False, deadline=time.monotonic() + 3) + arm(connection) client.shutdown(socket.SHUT_RDWR) client.close() - # The original TCP-based identity fails; the independent - # pipe identity must still authenticate the same child. + # The original TCP-based identity fails; the independently + # retained process lease must still authenticate the pipe. with pytest.raises(identity.PeerIdentityError): peer.revalidate(cancelled=lambda: False, deadline=time.monotonic() + 1) receiver.notify_stream_started(session) diff --git a/tests/windows_pipe_server.py b/tests/windows_pipe_server.py index bb041696..f97506ee 100644 --- a/tests/windows_pipe_server.py +++ b/tests/windows_pipe_server.py @@ -114,6 +114,15 @@ def read_exact(length): raise SystemExit(14) digest = hmac.digest(secret, b"Utterleaf OBS audio server ack v1\0" + header, "sha256") write_all(struct.pack("<4sBBH16s32s", b"ULAH", 1, 2, 0, session_id, digest)) + arm = read_exact(28) + arm_magic, arm_version, arm_kind, arm_reserved, arm_session, arm_mask, arm_status, arm_trailing = struct.unpack( + "<4sBBH16sBBH", arm + ) + if (arm_magic, arm_version, arm_kind, arm_reserved, arm_session, arm_status, arm_trailing) != ( + b"ULAC", 1, 1, 0, session_id, 0, 0 + ) or arm_mask > 63: + raise SystemExit(17) + write_all(struct.pack("<4sBBH16sBBH", b"ULAC", 1, 2, 0, session_id, arm_mask, 1, 0)) write_all(payload) buffer = ctypes.create_string_buffer(1) count = DWORD() diff --git a/utterleaf/obs_audio_pipe.py b/utterleaf/obs_audio_pipe.py index 5f72b35f..879f5377 100644 --- a/utterleaf/obs_audio_pipe.py +++ b/utterleaf/obs_audio_pipe.py @@ -25,6 +25,8 @@ _HEADER = struct.Struct("<4sBBH16s") _RECORD = struct.Struct("<4sBBH16s32s") _ACK_DOMAIN = b"Utterleaf OBS audio server ack v1\0" +_ARM_MAGIC = b"ULAC" +_ARM_RECORD = struct.Struct("<4sBBH16sBBH") HANDSHAKE_BYTES = _RECORD.size @@ -92,6 +94,7 @@ def __init__(self, pipe, peer, session_id: bytes, cancelled: Callable[[], bool]) self._decoder = obs_protocol.FrameDecoder() self._io_lock = threading.Lock() self._closed = threading.Event() + self._armed = False def __repr__(self) -> str: return "ObsAudioPipe()" @@ -113,6 +116,8 @@ def read_frames(self, *, deadline: float) -> list[obs_protocol.Frame]: with self._io_lock: if self._closed.is_set(): raise ObsAudioPipeCancelled("OBS audio connection closed") + if not self._armed: + raise ObsAudioPipeError("OBS audio session is not armed") _verify(self._pipe, self._peer, self._cancelled, deadline) data = self._pipe.read(obs_protocol.MAX_FEED_BYTES, deadline=deadline) _verify(self._pipe, self._peer, self._cancelled, deadline) @@ -137,6 +142,38 @@ def read_frames(self, *, deadline: float) -> list[obs_protocol.Frame]: self.close() _failure(exc) + def arm(self, *, additional_mix_mask: int = 0, deadline: float) -> None: + """Commit the one-use native Arm command before accepting ULAP frames.""" + try: + with self._io_lock: + if self._closed.is_set() or self._armed: + raise ObsAudioPipeError("OBS audio session cannot be armed") + if (type(additional_mix_mask) is not int + or not 0 <= additional_mix_mask <= 63): + raise ObsAudioPipeError("Invalid OBS audio arm request") + record = _ARM_RECORD.pack(_ARM_MAGIC, 1, 1, 0, self._session_id, + additional_mix_mask, 0, 0) + _verify(self._pipe, self._peer, self._cancelled, deadline) + self._pipe.write_all(record, deadline=deadline) + _verify(self._pipe, self._peer, self._cancelled, deadline) + reply = _read_exact(self._pipe, _ARM_RECORD.size, + self._cancelled, deadline) + _verify(self._pipe, self._peer, self._cancelled, deadline) + magic, version, kind, reserved, session, mask, status, trailing = _ARM_RECORD.unpack(reply) + if (magic != _ARM_MAGIC or version != 1 or kind != 2 or reserved + or session != self._session_id or mask != additional_mix_mask + or status not in (1, 2) or trailing): + raise ObsAudioPipeError("Invalid OBS audio arm response") + if status != 1: + raise ObsAudioPipeError("OBS audio arm was refused") + _check(self._cancelled, deadline) + if self._closed.is_set(): + raise ObsAudioPipeCancelled("OBS audio connection closed") + self._armed = True + except BaseException as exc: + self.close() + _failure(exc) + def close(self) -> None: self._closed.set() try: