Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
17 commits
Select commit Hold shift + click to select a range
5390069
refactor(native): move client phase and rejection policy into the core
RoboticHuman Oct 4, 2026
23c382e
refactor(native): move backlog hold and contiguity into the receiver …
RoboticHuman Oct 4, 2026
d827444
refactor(native): add the sans-IO receiver endpoint
RoboticHuman Oct 4, 2026
91da5c7
refactor(native): run ReceiverThread on the native receiver driver
RoboticHuman Oct 4, 2026
cc939fa
refactor(native): add the sans-IO producer endpoint
RoboticHuman Oct 4, 2026
6aa03bb
refactor(native): tidy the engine layer and make the driver loop prod…
RoboticHuman Oct 4, 2026
ad89ec9
refactor(native): drop the Python test seam and test ReceiverThread a…
RoboticHuman Oct 4, 2026
0e68097
refactor(native): run EventSender on a native producer driver
RoboticHuman Oct 5, 2026
42b4092
refactor(native): drop the primitive bindings from the Python module
RoboticHuman Oct 9, 2026
3af2ce7
refactor(native): collapse the driver layer and settle the thread's o…
RoboticHuman Oct 9, 2026
957bd94
refactor(client): rename ReceiverThread to EventReceiver and add close()
RoboticHuman Oct 9, 2026
b171a67
docs: tighten the client docs after the native engine work
RoboticHuman Oct 9, 2026
814b9cd
refactor(client): drain one notification queue inside update()
RoboticHuman Oct 9, 2026
2ed0740
test: keep the tests that guard regressions
RoboticHuman Oct 10, 2026
3c7b9a4
refactor(unreal): run the plugin on the client engine endpoints
RoboticHuman Oct 10, 2026
21cd6a8
fix(unreal): capture edit values when the edit is observed
RoboticHuman Oct 10, 2026
0c707f8
fix(client): write queued transactions before a sender closes
RoboticHuman Oct 10, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 0 additions & 3 deletions .github/workflows/pr-tests.yml
Original file line number Diff line number Diff line change
Expand Up @@ -40,9 +40,6 @@ jobs:
- name: Run unit tests
run: uv run --frozen pytest tests/unit -q

- name: Install native FlatBuffers headers
run: uv run --frozen python integrations/unreal/OpenUSDConnect/setup_flatbuffers.py

- name: Build native client tests
run: cmake -S tests/native -B build/native-tests

Expand Down
44 changes: 36 additions & 8 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -64,7 +64,8 @@ Return values and defaults:
Behavior:

- Token, metadata, and playback notifications run during `update()` or
`close()` on the calling thread instead of on network threads.
`close()` on the calling thread instead of on network threads; stage
metadata is delivered when it changes.
- Stage edits made in `on_resync` are no longer published, matching
`on_applied`.
- `UsdPublisher.update()` raises before `start()`. While disconnected it
Expand All @@ -78,8 +79,17 @@ Behavior:
like `UsdPublisher.flush()`, instead of returning `False`.
- `UsdReceiver.status` reports `CONNECTING` instead of `READY` while it
reconnects, like the other clients.

The low-level `EventSender`, `ReceiverThread`, and `EventDispatcher` keep their
- `ReceiverThread` is `EventReceiver` and no longer a `threading.Thread`:
call `close(timeout)` instead of `stop()` and `join()`, and read `running`
instead of `is_alive()`. On it and on `EventSender`, settings and state are
read-only properties (`token`, and the receiver's `reconnect`, stay
assignable) and `sock` is gone; read `connected`. Both close on leaving a
`with` block, and a collected one closes too, so keep the handle while it
should run. `EventSender.close()` writes the transactions already queued
and the Quit message, and does not wait for acknowledgements; call
`flush()` first for those.

The low-level `EventSender`, `EventReceiver`, and `EventDispatcher` keep their
callable arguments and properties.

### Added
Expand All @@ -98,25 +108,43 @@ callable arguments and properties.
`no_pending_recovery_stage`).
- `claim_playback()` and `send_playback_control()` on `SharedStageClient` and
`UsdPublisher`.
- Opt-in `background_send=True` moves transaction writes to a worker thread.
The worker needs the GIL, adding about 5 ms per write while the host's main
thread runs Python.
- `token_provider=` on `EventSender` and `ReceiverThread` supplies the token
- `token_provider=` on `EventSender` and `EventReceiver` supplies the token
for each connection attempt.
- `EventDispatcher.drained_message_count`, `ReceiverThread.stopped`, and
- `notifications=` on `EventSender` and `EventReceiver` pushes their
notifications into a queue the owner drains, and `snapshot()` returns the
native status in one call.
- `EventDispatcher.drained_message_count`, `EventReceiver.stopped`, and
`NoticeEmitter.has_local_changes`.
- Receiver replay identity and optional post-commit transaction checkpoints.

### Changed

- `EventDispatcher` starts its cursor at `receiver.sync_from - 1`, so
integrations no longer seed `last_seq` for continuation.
- `EventSender` and `EventReceiver` run their connections on native threads
in the client core: transaction writes leave the calling thread and need no
GIL, callbacks and the token provider run on the connection thread, and
`close()` interrupts a pending connect or read at once.
- `EventSender.connect()` waits for an attempt already in flight before
making its own; a token provider that raises is logged and fails that
attempt instead of raising. While recovery is required, `rejection_reason`
names the failure.
- Building the native extension fetches the pinned FlatBuffers headers on the
first configure, which needs network access unless
`FETCHCONTENT_SOURCE_DIR_FLATBUFFERS` names a local copy.
- The Unreal plugin runs on the native client engine: its receiver and emitter
threads drive the client core's `ReceiverEndpoint` and `ProducerEndpoint`,
built as the plugin's `OpenUSDConnectClientCore` module. Auth tokens stay in
the user's Unreal config under the same keys.

### Fixed

- A receiver continuing from a live-open snapshot replayed the full history
over it when the integration did not seed the dispatcher cursor.
- MCP writes were confirmed before the mirror applied them.
- `EventSender.send_events()` returns `False` for a transaction above the
16 MiB frame limit instead of queueing one that made the server close the
connection on every replay.
- Replay completion markers were lost when a resync reset applied progress.
- The emitter dropped property edits absorbed by a prim resync.
- Bidirectional clients read the token file on every `update()` while their
Expand Down
5 changes: 4 additions & 1 deletion CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -26,8 +26,11 @@ add_subdirectory(native/client_core)

nanobind_add_module(_native_client STABLE_ABI
native/python/client_module.cpp
native/python/driver_bindings.cpp
native/python/producer_bindings.cpp
native/python/receiver_bindings.cpp
)
target_link_libraries(_native_client PRIVATE OpenUSDConnectClientCore)
target_link_libraries(_native_client PRIVATE OpenUSDConnect::ClientDriver)
target_compile_features(_native_client PRIVATE cxx_std_17)

if(MSVC)
Expand Down
2 changes: 2 additions & 0 deletions docs/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,8 @@ linked from the shorter workflow guides.
stage ownership, native adapters, publishers, receivers, and resolver behavior.
- [Shared stage architecture](shared-stage-architecture.md): exact file-layer
synchronization and its protocol.
- [Native client core](../native/client_core/README.md): C++ targets, the
sans-IO endpoints, and the reference driver.
- [MCP integration layout](../integrations/mcp/README.md#layout): extension
points and module ownership.
- [Unreal plugin developer notes](../integrations/unreal/OpenUSDConnect/PLUGIN_DEV.md):
Expand Down
4 changes: 2 additions & 2 deletions docs/mcp-server-usage.md
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@ operations as local stdio tools. A client can author USD transactions and
inspect the composed result through an in-memory mirror.

The MCP process is a network client built on the core library (`EventSender` +
`ReceiverThread` + `EventDispatcher` + `UsdStageAdapter`), the same shape as the
`EventReceiver` + `EventDispatcher` + `UsdStageAdapter`), the same shape as the
`usdview` integration. Every scene event it sends uses the core protocol. Its
USD mirror also negotiates the optional layered-replay capability so authored
logical-layer opinions retain their server strength ordering during live sync
Expand Down Expand Up @@ -229,7 +229,7 @@ Verify any network with
## Implementation notes

- **Emit + mirror.** The MCP emits via `EventSender` and keeps a read-only
`Usd.Stage` mirror through `ReceiverThread`, replaying from sequence 1. The
`Usd.Stage` mirror through `EventReceiver`, replaying from sequence 1. The
server broadcasts committed records to every receiver, including the
producer's receiver, so the mirror contains the authoritative result of both
local and remote edits. The emitter and receiver use distinct diagnostic
Expand Down
52 changes: 26 additions & 26 deletions docs/usd-native-integration.md
Original file line number Diff line number Diff line change
@@ -1,9 +1,9 @@
# Python client and host-integration API

These APIs attach OpenUSDConnect to an application-owned `pxr.Usd.Stage`.
Call `update()` from the stage-owning thread. Socket reads and reconnects run
on background threads; encoding, USD work, and (by default) transaction writes
run on the calling thread.
Call `update()` from the stage-owning thread. Socket reads, transaction
writes, and reconnects run on native threads that do not need the GIL;
encoding and USD work run on the calling thread.

## Choose an API

Expand Down Expand Up @@ -82,7 +82,7 @@ Pass one `ClientObserver` subclass as `observer=` and override only what the
host needs. Methods never run on a network thread: notifications arrive in
`update()` or `close()`, and delivery methods run wherever the client applies
authoritative state (`update()`, `refresh_asset_dependency()`, recovery). The
client wires only overridden methods, so unused notifications cost nothing:
client calls only overridden methods:

```python
class HostObserver(ClientObserver):
Expand Down Expand Up @@ -132,11 +132,6 @@ Receiving pauses while a `ManagedClient` edit target is foreign, because its
accepts any edit target (session-layer edits stay local), so its loop calls
`update()` unconditionally.

`background_send=True` moves transaction writes to a worker so a full socket
buffer cannot block the UI thread. The worker needs the GIL: while the host's
main thread runs Python, each write waits for Python's thread switch interval
(about 5 ms), so keep the default for latency-sensitive editing on fast links.

Before closing, stop authoring and call `submit_and_wait()`. Success means the
edits are durable, not that their echo has been applied locally.

Expand Down Expand Up @@ -168,8 +163,7 @@ instead of silently degrading to flat replay.

Open the original base scene. A generated live-open snapshot already contains
composed server state and is rejected because replaying the complete managed
history over it would duplicate opinions. Snapshot continuation is a separate
flat integration path used by the live-open host plugins.
history over it would duplicate opinions.

Use `rebind_stage(new_stage)` when a host replaces its stage. Passing `None`
parks stage application (phase `PARKED`) while the network queue continues to
Expand Down Expand Up @@ -285,9 +279,8 @@ with ManagedClient(
Construction creates `client.authoring_layer`, inserts it below the
authoritative managed block, and makes it the edit target. Keep that target
while the client is active. `update()` freezes local edits, applies the queued
authoritative prefix, then submits the frozen local batch. The dispatcher
suppresses and invalidates the emitter while applying server records, so
authoritative echoes do not become new local submissions.
authoritative prefix, then submits the frozen batch; applied server records are
never republished as local edits.

`publish_current_edit_target()` queues a snapshot of the authoring layer for
the next `update()` that can publish; a zero return can mean it is still queued.
Expand Down Expand Up @@ -328,8 +321,8 @@ happens next depends on `DCCAdapter.targets_stage()`:

Custom stage-backed adapters must override `targets_stage()` explicitly.

Shader mapping interfaces live in `openusdconnect.shader_mapping`; existing
imports from `openusdconnect.adapters` remain supported. Integrations that
Shader mapping interfaces live in `openusdconnect.shader_mapping` and are also
importable from `openusdconnect.adapters`. Integrations that
author shader inputs directly can use `set_connectable_input_value` and
`resolve_shader_port_type` from `openusdconnect.usd_authoring`. They operate
under the stage's current edit target and do not send network events.
Expand Down Expand Up @@ -424,14 +417,10 @@ managed receiver, call `refresh_asset_dependency(path)` after an asset becomes
available or its resolver mapping changes; omit the path to retry all pending
dependencies.

A context-only resolver remap is a special case for adapters targeting a
non-USD native scene. It can recompose both the live and previous-state stages
before projection observes the old topology. The dispatcher then sets
`native_scene_rebuild_required` and stops incremental delivery. The high-level
receiver reports it as `RECOVERY_REQUIRED` in `client.status`. Rebuild the
native destination and call
`client.acknowledge_native_scene_rebuilt()` before resuming. An ordinary
reconnect does not clear this guard.
For an adapter targeting a non-USD native scene, a context-only resolver remap
can recompose both the live and previous-state stages before projection
observes the old topology. That is the `RECOVERY_REQUIRED` case in
[Observing the client](#observing-the-client); a reconnect does not clear it.

## Identity and authentication

Expand All @@ -446,9 +435,9 @@ store.

## Low-level APIs

`NoticeEmitter`, `EventSender`, `ReceiverThread`, and `EventDispatcher` remain
`NoticeEmitter`, `EventSender`, `EventReceiver`, and `EventDispatcher` are
public for integrations whose scheduling or continuation requirements cannot
use the high-level clients. `ReceiverThread` requests layered replay by default;
use the high-level clients. `EventReceiver` requests layered replay by default;
passing `layered_replay=False` selects the single-layer flat contract. Ordinary
native-scene integrations should use `UsdReceiver(adapter=...)` instead of
assembling these components.
Expand All @@ -464,6 +453,17 @@ rejection/recovery status on subsequent ticks. `cancel_connect()` invalidates
pending attempts and reports whether they have finished; `disconnect()` also
closes an established connection. Neither discards the transaction outbox.

A sender's connection thread starts on the first connection request, a
receiver's on `start()`; callbacks run on that thread. `close(timeout=None)`,
or leaving a `with` block, stops the thread, which cannot be restarted, and
returns whether it exited in time (`False` at once from a callback). A sender
first writes the transactions already queued and the Quit message; closing does
not wait for acknowledgements, which `flush()` does. Keep a reference while the
object should run: a collected one closes.
Either object also takes `notifications=`, a `NotificationQueue` its owner
drains instead of every callback but `on_token_issued` (combining them raises
`ValueError`), and offers `snapshot()`, its native status read in one call.

## Embed a server

Use `ServerRuntime` when the application owns startup and shutdown:
Expand Down
3 changes: 2 additions & 1 deletion examples/fourier_waves/author.py
Original file line number Diff line number Diff line change
Expand Up @@ -143,7 +143,8 @@ def run_author(args) -> int:
print("\nstopping.")
return 0
finally:
sender.disconnect()
sender.flush(timeout=5.0)
sender.close()


def main() -> int:
Expand Down
8 changes: 4 additions & 4 deletions examples/fourier_waves/wave_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,7 @@
from openusdconnect.adapters import UsdStageAdapter # noqa: E402
from openusdconnect.cli_common import add_sync_endpoint_args # noqa: E402
from openusdconnect.dispatcher import EventDispatcher # noqa: E402
from openusdconnect.receiver import ReceiverThread # noqa: E402
from openusdconnect.receiver import EventReceiver # noqa: E402
from openusdconnect.sender import EventSender # noqa: E402

DEFAULTS = {
Expand Down Expand Up @@ -91,7 +91,7 @@ def __init__(self, host, port, proc_path):
from pxr import Usd

self.mirror = Usd.Stage.CreateInMemory()
self.receiver = ReceiverThread(
self.receiver = EventReceiver(
host=host, port=port, sync_from=1,
client_id="fourier-wave-client", origin=f"{origin}-recv",
)
Expand Down Expand Up @@ -171,8 +171,8 @@ def run(self):
time.sleep(1.0 / 30.0)

def stop(self):
self.receiver.stop()
self.sender.disconnect()
self.receiver.close()
self.sender.close()


def build_parser(add_help: bool = True) -> argparse.ArgumentParser:
Expand Down
5 changes: 3 additions & 2 deletions examples/instancing_dance/dance.py
Original file line number Diff line number Diff line change
Expand Up @@ -143,7 +143,7 @@ def run_dance(args: argparse.Namespace) -> int:
print(f"sending setup events ({args.instances} instances)...")
if not sender.send_events(setup_events(asset, args.instances)):
print("setup send failed")
sender.disconnect()
sender.close()
return 1
print("setup complete.")

Expand Down Expand Up @@ -173,7 +173,8 @@ def run_dance(args: argparse.Namespace) -> int:
except KeyboardInterrupt:
print("\nstopping.")
finally:
sender.disconnect()
sender.flush(timeout=5.0)
sender.close()
return 0


Expand Down
4 changes: 2 additions & 2 deletions integrations/blender/capture.py
Original file line number Diff line number Diff line change
Expand Up @@ -1136,7 +1136,7 @@ def _try_send_dirty_events():
"""Build and send dirty events if emitter and sender are both connected."""
if _state.author is not None and _state.author._applying_remote:
return
if _state.notice_emitter is None or _state.sender is None or _state.sender.sock is None:
if _state.notice_emitter is None or _state.sender is None or not _state.sender.connected:
return
events = _state.notice_emitter.prepare_events_for_send()
if events:
Expand Down Expand Up @@ -1580,7 +1580,7 @@ def execute(self, context):
if not _cancel_emitter_reconnect():
self.report({"WARNING"}, "Emitter reconnect cancellation is still finishing")
return {"CANCELLED"}
if _state.sender is not None and _state.sender.sock is not None:
if _state.sender is not None and _state.sender.connected:
self.report({"INFO"}, "Already connected")
return {"CANCELLED"}
sender = _state.sender
Expand Down
Loading
Loading