← Workstreams

Workstream: Observables

Status: Implemented (JVM), graduation pending · Component: Maximize developer productivity

Implementation status & remaining work: implementation history is recorded in the epic Observable#87. The JVM url:// implementation is feature-complete and entirely on main: the read contract, runtime proxies, collections-aware sync protocol, negotiated push transport, command-lifecycle handles, idempotent command dispatch (UrlResolver PR #707), both JVM Compose boundary forms, the reusable resilience fixture and lower-revision server-lineage recovery (UrlResolver PR #680, merged 2026-07-24), and the CqrsObservableExample reference adoption are all merged. UrlResolver main currently declares foundation.url:resolver:0.0.1066. Web/wss, JS/Wasm materialization, and the wallet ambient remain deliberately deferred in Observable#86; ContainerNursery push forwarding is now in progress rather than deferred (see Current state). The capability is therefore feature-complete on JVM but has not graduated.

Goal

Make reactive projections a solved primitive across process and network boundaries. A developer defines a projection of their data as a plain Kotlin interface, serves it at a URL, and any consumer — a Jetpack Compose UI, a CLI, another service's projection — resolves that URL into a local object whose properties read like a POJO and update automatically whenever the underlying data changes. No hand-written client wrappers, no marshaling code, no cache-invalidation plumbing, no read-after-write polling.

Vision: URL-addressed live projections

The full design is the Observable repository's REMOTE_OBSERVABLES.md; this workstream tracks the intent and the build order. The committed shape:

  • Projection interfaces extend RemoteObservable (getRemoteObservableAddress(), awaitAvailable()) in the shared Api module — the schema, not bytecode. The server side stays the CQRS projection class already mandated to be observable.
  • resolver.resolve("url://todos/") returns immediately, before the first server response, as a runtime-generated proxy (ASM on the JVM, JavaScript Proxy objects on Kotlin/JS, ahead-of-time generation for Kotlin/Wasm) — extending the proxy-generator family that UrlResolver already ships. Property reads never block on the network.
  • One data-access contract, no placeholders, no status booleans: a read returns the value or throws a structured RemoteObservableDataUnavailable (reason NEVER_LOADED, CONSTRAINT_NOT_MET, or SCHEMA_UNAVAILABLE), after registering the read as a dependency — so presentation boundaries self-heal via normal invalidation. Per-object availability checks were rejected as unsound under composition (check-then-act).
  • Read-after-write for free: the HierarchicalClock ClockPad becomes the ambient read constraint. Command proxies record returned timestamps automatically; a read that does not yet reflect the caller's writes throws CONSTRAINT_NOT_MET, and Compose boundaries keep showing last-known-good content, dimmed, until the deltas catch up.
  • Commands are observable and cancelable, not fire-and-forget: a command (a parameterized proxy method) returns a RemoteCommandDerivedObservable — an Observable-API interface declared in observable-core alongside the RemoteObservable marker, which UrlResolver implements via its schema-driven live-projection proxy machinery. It is a reactive, subscriber-driven handle to the in-flight invocation, exposing its observable surface by composition (val state: DerivedObservable<CommandState> — DerivedObservable is deliberately final, so the handle wraps one rather than extending it). It exposes cancel() and a lifecycle status (PENDING / RUNNING / CANCEL_REQUESTED / CANCELED / SUCCEEDED / FAILED / NOT_FOUND), and yields the command's HierarchicalTimestamp on success — so the read-your-writes token above is preserved while callers also get progress, cancellation, and failure visibility. See REMOTE_OBSERVABLES.md §5.1. This subsumes the imperative submit → poll-to-terminal → cancel → renew-lease orchestration that long-running backend operations hand-roll today: a CI build, a deploy, a provisioning job becomes a single observed command — its status watched to a terminal state, cancel()ed if abandoned, kept alive by the GC lease below — with no separate status-poll RPC that can throw on a transient reconnect and abort the whole operation. The command's status survives a connection reconnect while the serving process retains its record; if that bookkeeping is gone the handle settles NOT_FOUND. The default in-memory command store does not make in-flight work durable across a server-process restart; durable terminal replay requires an application-supplied shared CommandOutcomeStore or event log.
  • Proxy calls consume the standard ambients without polluting the projection interface: freshness is built today through the resolver-owned session ClockPad plus RequiredClockPad / withRequiredClockPad. The settled wallet design establishes W3WalletSpenderApi at an entry point and propagates its capability across url://, but that carrier remains deferred in Observable#86.
  • High update granularity end to end: per-key server observables → per-property revisioned deltas over one subscription per object URL → per-slot client invalidation → per-read Compose recomposition. Collections split membership from content; one-level inline child embedding kills N+1 loading waterfalls.
  • Transport-agnostic, one authorization model: url:// push over the existing PersistentRpcConnection, and a web transport of HTTPS GET-snapshot + a wss sync connection carrying the same SUBSCRIBE → SNAPSHOT → DELTA protocol, multiplexed over one WebSocket — a second binding of one protocol, not a differently-shaped sibling, which is what the lazy dynamic subscription set, the connection-scoped lease below, and eventual command parity all want. Authorization is per-projection and enforced identically on both transports: over url:// the wallet ambient marshals the required capability; on the web the client presents the same W3Wallet token capability as a bearer credential — an Authorization header on the GET, a first auth frame after the wss handshake (browsers cannot set WebSocket headers, and keeping the token out of the handshake URL keeps it out of access logs); token capabilities verify offline either way. A projection served without a capability requirement is thereby a free public JSON read API for non-Kotlin clients — a plain curl-able GET — and serving a projection never silently publishes it. The web transport is read-only in v1 (commands ride url://); web command parity over the same wss connection is a planned follow-up milestone on the roadmap below.
  • Server resources track client GC via a connection-scoped lease: a consumer that drops a proxy (reclaimed by GC), crashes, or loses its connection is auto-unsubscribed, so the server frees the subscription's state without needing a clean goodbye — the server-side complement to the client-side subscriber-driven lifecycle. The lease is the persistent connection's liveness: every subscription on a live connection stays live — an idle-but-observed subscription never lapses, so a quiet projection can never silently miss its next delta. When the connection dies, its subscriptions enter a reconnect grace window and are freed if the client does not return; on reconnect the sync engine re-subscribes every live proxy automatically and receives a fresh atomic snapshot. A proxy the client GC'd on a healthy connection is reaped lazily at the next delta push (each update is implicitly "still reading?"), so an idle subscription pays for no separate heartbeat.
  • Transient outages degrade, they never abort — and a lifecycle is never a deadline: the read contract (§3) and the command lifecycle (§5.1) are a reliability contract as much as an ergonomics one. A loaded projection rides connection loss, peer restarts, and failed re-syncs as stale-last-known-good that re-syncs on recovery — per REMOTE_OBSERVABLES.md §3 a loaded read "never throws again," and §11 routes genuine transport failures to last-known-good + backoff, not to an escaping exception. The resolver's bounded persistent-connection reconnect is an internal signal the sync engine absorbs into stale-but-recovering — never a thrown timeout that surfaces to the consumer. A read's and a command's lifetime is bounded by liveness and explicit cancel() (the GC lease above), never by a request deadline whose expiry throws. Resilience is therefore the framework's default and only behavior, not per-consumer retry plumbing that every caller must re-implement — and that a long-running consumer (a CI runner polling a singleton build service over a persistent connection, a deploy watcher, a downstream projection chaining an upstream one) invariably gets wrong, turning one transient blip mid-operation into a hard abort.

Current state

The as-built Observables project provides the CQRS read-side foundation (MutableObservable / DerivedObservable, automatic dependency tracking, Compose integration, and the failure/retry/lifecycle model). Its live-projection JVM path is now implemented across Observable and UrlResolver's foundation.url.resolver.projection package: resolve() immediately materializes a URL-addressed projection, property reads use the structured availability contract, the ambient session ClockPad supplies read-your-writes, and the sync engine maintains per-property and per-key collection updates over negotiated push or its compatible poll fallback.

CqrsObservableExample PR #4 deleted its hand-written TodoClient in favor of resolve<TodoList>() and is the reference adoption. It proves the integration, but it is not a production service and therefore does not satisfy the graduation criterion. The Phase C resilience fixture landed with UrlResolver PR #680 on 2026-07-24 and now lives in the resolver's projection/testing package as the standing ResilienceAcceptanceHarness.

Production services have begun serving live projections, but none consumes one exclusively yet. GithubProxy serves its representative tranche — repository metadata, pull-request list, pull-request details plus checks — as live projections with per-resource freshness watermarks and strictest-subscriber required-ClockPad refresh, merged 2026-08-15/16 across GithubProxyApi PR #21, GithubProxyServerService PR #37, GithubProxyCli PR #19, and GithubProxyWui PR #2. The CLI and WUI establish the ambient pad and handle the freshness effects, but the pages that render that tranche through resolve() are still feature work — the projections are served and not yet consumed. Separately, BuildTest publishes a BuildTestProjectionSignal : RemoteObservable journal-cursor signal that BuildTestWui resolves with resolveRemoteObservable, then drains changes over its existing RPC read path — a live projection used as a change notifier alongside the imperative reads, not in place of them.

ContainerNursery push forwarding — deferred item 4 of Observable#86, and the reason a CN-fronted projection falls back to polling — is now being built: foundation.url:protocol:0.0.486 is published and ContainerNursery PR #595 carries the duplex push bridge through UrlFacadeProvider.

The production-facing kotlin-build-ci BuildRunner still consumes build status via an imperative poll; kotlin-build-ci PR #160 made that poll tolerate the transient-failure class that motivated this work, but migrating the read path to a live projection is separate production-graduation work tracked in Observable#86.

Why it accelerates developers

  • The read side of every service becomes one line. resolver.resolve(url) replaces the per-service client wrapper; the projection interface is the only contract.
  • UIs get loading, staleness, and read-your-writes for free. The boundary pattern renders the unavailable state once; everything else is plain property access that recomposes exactly when the data it read changes.
  • Long-running reads stop dying on transient blips. Loaded projections degrade to stale-last-known-good and recover across connection loss and a dependency restart; command handles survive a connection reconnect while the server retains their record. A consumer that holds a projection across minutes — a CI runner, a deploy watcher, a backend chaining another service's projection — therefore does not need to reimplement read-side reconnect/retry logic. This does not promise that in-flight command work survives a restart of the process executing that work; durable terminal replay needs the shared outcome store described above.
  • It composes the rest of component 1. Projections chain across services (a downstream service's calculator reads upstream remote projections with zero plumbing), Event Log projections become directly servable, the state plane pairs with the EventSourcingStreams event plane (a projection is the fold of a stream: consumers who need every event subscribe to the stream, consumers who need current state resolve the projection), and HierarchicalClock consistency stops being per-service plumbing.

Plan / roadmap

These follow the build order implied by the design; items marked (core) land in the Observable repository, (resolver) in UrlResolver.

  • [x] The read contract (core): RemoteObservableDataUnavailable with structured reason/address/watermark fields, dependency registration before the throw, and the rule that unavailability never feeds the retry/backoff scheduler.
  • [x] RemoteObservable + ProjectionState (core; common marker, JVM store): the marker interface and the per-URL store holding one observable slot per property.
  • [x] JVM ASM proxy generator (resolver): thin generated shims over ProjectionState, default methods left to run (locally computed reactive properties), command shims allowed before load, address-based identity, per-interface class cache.
  • [x] Sync protocol: SUBSCRIBE → atomic SNAPSHOT → per-property revisioned DELTAs with HierarchicalTimestamp watermarks, key-based list ops, one-level inline child embedding, weak identity map.
  • [x] url:// push transport: negotiated server push over the persistent direct connection, with connection-scoped leases, reconnect grace, idle-hold, and compatible level-triggered poll fallback when no push path is available.
  • [ ] Web transport: HTTPS GET-snapshot plus a wss sync connection carrying the same SUBSCRIBE → SNAPSHOT → DELTA protocol, read-only at first and with the same per-projection capability model as url://. Tracked in the deferred-work issue. This is the application-level binding for non-resolver clients; carrying url:// itself over HTTPS for resolver clients is the separate UrlResolver HTTPS Transport workstream, and a service serves both from one HTTPS origin. On a host that runs only while a request is open, that workstream makes an observed command the thing that keeps the instance awake, and leans on the shared CommandOutcomeStore above for one specific job: a command whose request is delivered twice, or to two instances, still runs once. Pausing and resuming in-flight work is a separate requirement on each command's own durable progress, not something that store provides.
  • [ ] Web command parity (follow-up milestone, after the read-only web transport ships): command invocation and cancellation as frames on the same established wss sync connection, with the returned command's CommandState readable as a projection over that connection — so non-Kotlin clients become first-class writers, not just readers. Tracked in the deferred-work issue.
  • [x] Ambient RequiredClockPad: resolver-owned auto-advancing session pad, withRequiredClockPad / ProvideClockPad, pad requirement carried on SUBSCRIBE — the non-blocking mode the HierarchicalClock workstream calls for.
  • [ ] Wallet ambient on proxy calls (resolver): establish W3WalletSpenderApi at the entry point and have the resolver propagate its capability on every read and command without adding it to projection signatures, marshalled across url:// as a capability reference per the ambients plan — shared work item with the W3Wallet workstream, tracked in the deferred-work issue.
  • [x] Command lifecycle observables (resolver): command shims return a RemoteCommandDerivedObservable whose observed value is a CommandState(status, timestamp?, failure?) and which exposes cancel(), with server-side command bookkeeping so status survives reconnects (NOT_FOUND when it does not) and cancellation is acknowledged. The type is declared as an interface in observable-core alongside the RemoteObservable marker, with the CommandState/CommandStatus/CommandFailure value types beside it, exposing its reactive surface by composition (val state: DerivedObservable<CommandState> plus cancel()) since DerivedObservable is deliberately final; UrlResolver implements it via its proxy machinery, so projection Api modules depend on lightweight observable-core alone and never pull the resolver implementation. The terminal-status set (SUCCEEDED/FAILED/CANCELED/NOT_FOUND, no TIMED_OUT/REJECTED) and observed-value shape are now settled — see the resolved REMOTE_OBSERVABLES.md §14.8. Idempotent dispatch is merged in UrlResolver PR #707.
  • [x] Compose RemoteObservableBoundary: read-at-boundary v1 (fallback on NEVER_LOADED, keep-last-good-dimmed on CONSTRAINT_NOT_MET) is built; the deep-throw form is also built for desktop/JVM on the Compose ErrorBoundary primitive. Non-JVM materialization remains deferred with the platform work below.
  • [ ] Web materialization: Kotlin/JS Proxy objects (with the @JsExport authoring convention for stable names); ahead-of-time generation tool for Kotlin/Wasm. Depends on the portable value-serialization facade (a common JsonObject with an org.json JVM actual) tracked in the Kotlin Multiplatform workstream, so toJson()/fromJson() value objects can materialize on JS/Wasm. Tracked in the deferred-work issue.
  • [x] Resilience acceptance fixture implementation: a shared end-to-end fixture — kill the backing service mid-subscription, then assert the consumer goes stale and recovers to live, and never throws — that every served projection and command must pass, exercising the real reconnecting url:// transport (not an in-process cooperative fake) and the relay path, not only direct localhost. This is the standing guard for the §11 failure semantics: the class of outage where a transient dependency restart aborts a long-running consumer is then caught by construction rather than left to per-caller retry that no test exercises. Generalizes the "kill the backing service and observe recovery" shape into a reusable harness so resilience is verified once per primitive, not re-litigated per consumer. Per the testing standards. Merged in UrlResolver PR #680 on 2026-07-24 and available on main as ResilienceAcceptanceHarness.
  • [x] Reference adoption: rewrite the CqrsObservableExample client on resolve() (deleting TodoClient), proving per-key granularity over the wire, with end-to-end tests per the testing standards, completed by CqrsObservableExample PR #4.
  • [ ] ContainerNursery push forwarding: relay server-initiated frames through UrlFacadeProvider.forwardViaRpc so a CN-fronted projection pushes instead of falling back to poll. In progress — foundation.url:protocol:0.0.486 published, ContainerNursery PR #595 open with the duplex bridge.
  • [ ] Production graduation: migrate a real service's read path exclusively to resolve(); the reference example does not count, and neither does serving projections alongside an unchanged RPC read path. The nearest candidate is GithubProxy, whose tranche is already served — the remaining step is rendering the CLI and WUI read paths through resolve() and retiring the operations they replace. Tracked in the deferred-work issue.

Graduation

The JVM implementation is feature-complete, but the workstream has not graduated: CqrsObservableExample is the reference adoption, not a production service, and the production adoption that has happened since — GithubProxy serving its tranche, BuildTest publishing a cursor signal — adds projections beside existing RPC read paths rather than replacing one. Graduation occurs when a real service's read path is consumed exclusively through resolve() in production — UI updating reactively, read-your-writes via the ambient pad, no hand-written client — as tracked in Observable#86. At that point this graduates to a first-class capability: mark it so on the Observables project page, whose Implementation status section already describes the as-built JVM capability and tracks what remains.