Repository navigation
refactor(native): move client phase and rejection policy into the core - #49
Draft
RoboticHuman wants to merge 17 commits into
Draft
RoboticHuman wants to merge 17 commits into
RoboticHuman wants to merge 17 commits into
Conversation
RoboticHuman
commented
Oct 4, 2026
Owner
- Fetch the pinned FlatBuffers headers with CMake FetchContent on the ClientProtocol target and link it into the Python module and the Blender build; the native tests and CI no longer depend on the Unreal plugin's staged copy.
- Add engine/status.h with ComputePhase, RejectionCodeName and RejectionDisposition, bind them, and route compute_phase, TransactionFailure and the sender's rejection handling through the bindings so one policy table replaces the three copies.
- Give ProducerRecoveryDisposition its own header, add the engine ctest target with static_asserts against the wire enum, and pin the FetchContent tag in check_versions.py.
Fetch the pinned FlatBuffers headers with CMake FetchContent on the ClientProtocol target and link it into the Python module and the Blender build; the native tests and CI no longer depend on the Unreal plugin's staged copy. Add engine/status.h with ComputePhase, RejectionCodeName and RejectionDisposition, bind them, and route compute_phase, TransactionFailure and the sender's rejection handling through the bindings so one policy table replaces the three copies. Give ProducerRecoveryDisposition its own header, add the engine ctest target with static_asserts against the wire enum, and pin the FetchContent tag in check_versions.py.
…inbox Add FreezeMarker and DrainedThrough to OrderedReceiverSession and bind them, replacing the Python BacklogHold in ManagedClient and SharedStageClient. Enable require_contiguous for the Python inbox so duplicates are dropped and gaps replay on the receive thread, and report the consumer's applied cursor after each successful batch so replay starts from what the stage applied.
Add the OpenUSDConnect::ClientEngine target with ReceiverEndpoint, which turns socket bytes, read timeouts, disconnects, and host time into connect, send, close, wake, and log actions, and queues handshake, token, stage metadata, and playback notifications. It reproduces the Python receiver's handshake, negotiation, control routing, sequence gap and overflow handling, and reconnect backoff with the overflow drain wait. Reconcile ReceiverReplayIdentity with the Python receiver: claim the received identity, keep it across a replay request unless a reset is queued or drained but not yet reported applied (the inbox tracks this as ResetPending), and queue the receiver's own reset before a replay from one. Reject a ReplayComplete whose head was not received as a sequence gap. Keep the Unreal receiver's existing claims by requesting the replay on its identity before each connection. Add ctest coverage for every handshake outcome, control message, accept result, timeout, backoff, and replay identity scenario, and a TLA+ model of the Hello and replay identity flow across sequence domains. The model ties both identities to the prefix the receiver holds, checks that a receiver which knows its prefix's domain is resumed rather than reset, and records the two remaining reset causes: a snapshot prefix and a live reset whose ReplayComplete never arrived.
Add OpenUSDConnect::ClientDriver, a reference host loop that drives a ReceiverEndpoint on one thread with Winsock or BSD sockets whose blocking calls an event interrupts, and a scripted socket factory as a test seam. Bind the endpoint, driver, and factories, and make ReceiverThread a wrapper over them that keeps its public API. Receivers may omit client_id and origin, and reconnect can change while running. Tests drive receivers through scripted sockets and relay real server connections frame by frame.
ProducerEndpoint owns the producer side of the connection protocol that EventSender implements today: one-shot attempts with request backoff and the server's retry-after, the emitter Hello and its committed highwater check, ordered replay of the outbox on publication, cumulative acknowledgements with mirror checkpoints, rejection dispositions, and the repair and abandon recovery flows. An attempt the host has not taken is withdrawn when cancelled, and a reported socket end voids the frames and close still queued for it, so a stale close never ends the next attempt. SendAction shares its bytes with the outbox, so a replay copies nothing. TransactionFailure joins status.h with the description hosts show. The emitter Hello requires a client and producer session id, as the server does, and no longer requires an origin. Both endpoints share their message decoding. ProducerConnection.tla checks attempts, cancellation, and the host loop's report ordering against an honest and a diverging server.
…ucer-capable Fold single-declaration headers into their first users. producer_recovery.h carries the rejection policy and TransactionFailure; status.h holds only the client phase. Move the scripted sockets into OpenUSDConnect::ClientDriverTesting under driver/testing, so the reference driver no longer ships the test seam. Use one close vocabulary (Stop, Disconnect, Close, Quit, EndConnection, OnDisconnected) in both endpoints and the driver loop, and share the Hello, framing, handshake classification, and HelloOk notifications through src/engine/endpoint_common.h. Drive DriverLoop through the host I/O both endpoints share plus a LoopRole. Give Socket::SendAll a deadline and Socket::Receive an optional one, so a stalled write is a TransportError and reads end at NextWake. Trim the endpoint tests to behavior, with one shared test host.
…gainst a live server Remove the scripted socket bindings from the extension and the _SOCKET_FACTORY hook from ReceiverThread, which builds its TCP factory when it starts. ClientDriverTesting stays for the C++ driver tests and builds only when a target links it. Rewrite test_receiver.py around the wrapper contract against an in-process server: settings, property pass-throughs, lifecycle, token provider and legacy callbacks, drain and replay wiring, connection errors, and rejections. The protocol cases it replaced are covered by the endpoint, driver, and socket ctest cases. Connect the managed and receiver client tests to in-process servers, give the shared-stage recovery tests a receiver stub, and drive the budget tests with peer traffic through the server. Integration tests force reconnects by restarting the server listener on the same port and read hellos on the server instead of relaying frames.
Add ThreadedProducerDriver, which drives a ProducerEndpoint on the shared DriverLoop and offers the blocking Connect and Flush, and report Closing in ProducerStatus so Connect can wait out a close in flight. Bind the producer endpoint and driver, and share the Python driver plumbing and common types between both role bindings. Make EventSender a wrapper that keeps its public API and remove background_send from it and the high-level clients, since the native writer needs no GIL. The sock attribute is gone; Blender capture checks connected. Settle handshake parity: both endpoints close with ProtocolError on an unexpected handshake answer and keep an empty HelloRejected reason, which the wrapper words. Rewrite the sender tests around the wrapper contract against live servers; protocol cases are covered by the endpoint and new driver ctest cases.
RoboticHuman
marked this pull request as draft
October 6, 2026 23:57
Remove FrameDecoder, encode_frame, ReceiverInbox, ProducerSession, and the enums and glue that only they used from _native_client; the Python client reaches the client core through the endpoints and drivers alone. Port the tests that drove those primitives to ctest cases on the frame codec, the receiver inbox, and the producer session, and keep only the packaging tests in test_native_client.py. Write ThreadedProducerDriver::Connect as a switch over ConnectResult and shorten the ProtocolError note in actions.h.
…wner Make the driver loop the public ThreadedDriver<Endpoint> and derive the two role drivers from it, so the lifecycle exists once and both Python drivers bind it the same way. Drop WakeAction: NextWake() is the only wake contract. Fold ReconnectPolicy into actions.h, Quit into Close, and remove the scripted socket's timeout delivery in favor of a real read deadline. The Python wrapper owns its native thread: a collected receiver or sender stops its connection, which removes the driver's self-reference and the join-from-its-own-thread path. A callback that turns out to hold the last reference stops the driver and releases it on the main thread. Share one GIL trampoline and one duration property helper between the bindings; merge the two driver test harnesses, the two Kinds helpers, and the duplicated backoff tests.
EventReceiver pairs with EventSender, and both own their native thread the same way: close(timeout) stops the thread and the connection, waits for the thread to exit, and reports whether it did; both are context managers. The receiver's stop(), join(), is_alive(), and ident are gone, since it is not a thread; running and stopped remain. The high-level clients close both endpoints the same way and drop the thread check that stopped being meaningful when the receiver stopped being a thread. The Blender add-on, the editor bridge, and the example use close() in place of stop() and join().
State the native client core for a C++ integrator in one pass: the four targets, how to link them, the protocol layer's borrow and builder rules, the four-step endpoint host loop, and the reference driver's thread rules. Drop the per-method semantics the headers state and the test-seam notes. In the integration doc, say the wrapper lifecycle once, remove the second statements of the live-open and native-scene-rebuild contracts, and describe the managed client's republish rule instead of the dispatcher internals. Merge the 0.5.0 changelog bullets that repeated each other and point the docs index at the client core README.
Both roles of a high-level client push into one native NotificationQueue that ClientBase drains in update() and close(), converting each notification straight to its observer value; the Python callback queue, the dict converters, and the observer hooks module are gone. Stage metadata reaches the observer when it changes, since both handshakes carry it. A driver hook reports an issued token on the driver thread before the next actions, so the shared credential is updated and persisted before the other role's handshake without draining notifications into Python. EventReceiver and EventSender accept notifications= to push into a queue their owner drains, and snapshot() returns the native status once; status reads two snapshots instead of fifteen native calls. One function words handshake rejections for the wrappers and for status.
Remove tests that pin a transition another test already pins with the same assertions, and fold runs of near-identical tests into tables: the five highwater contradictions, the handshake rejections of both endpoints, the receiver's start and first hello. Inside kept tests drop the incidental counters, log-text checks, validation rows beyond the rule they pin, and parametrizations over paths that do not branch. The wrapper tests share one client_registered helper, and the phase precedence table is checked once, in C++.
Make native/client_core a UBT module the test harness stages (headers, the engine sources, and the frame codec), exported for Unreal's one DLL per module through OPENUSDCONNECT_CLIENT_API, which is empty elsewhere. One FEndpointRunner<Endpoint> drives either endpoint on an FRunnable with an FSocket: it applies the endpoint's actions, reads with no added latency on the receiver and wakes at once for a submit on the producer, and reports bytes, timeouts, and time. The subsystem owns a ReceiverEndpoint, a ProducerEndpoint, a notification queue per endpoint drained in Tick, and builds its status from the two snapshots; tokens stay in the editor config under the existing key. FEmitClient, FSyncClient, FProducerEndpointState, the wire framing helpers, and the plugin's own protocol tests are gone, as is the replay identity ctest that mirrored them.
The plugin recorded changed prim paths and read their values when it emitted, so a replay applied in between sent the server's values back as local edits; an edit made while the server was down was lost on reconnect. Read the values at the start of the tick that observes the change, before any received frame applies, keep the latest per prim and per input, and append them once connected. A fresh connection re-sends previously emitted transforms only for prims with no captured edit.
Closing a sender stopped its thread at once, before the loop had written the transactions already accepted into the outbox, so a send followed by close, disconnect, or collection lost the transaction. ThreadedProducerDriver::Close disconnects first, which queues the Quit and the close after the pending sends, waits for the connection to end within the timeout or a one-second grace, then stops and joins; the receiver's Close stops and joins. Explicit close, collection, and interpreter exit all take this path. Closing still does not wait for acknowledgements; flush does. The material zoo flushes before closing because its viewers wait for the published sequence; the send command and the examples close instead of disconnecting.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.