Skip to content

feat(orchestration): substrate fixes for SPA-side adaptive reaping (v1.0.27) - #39

Merged
KodBena merged 6 commits into
mainfrom
feat/orchestration-output-channel
May 19, 2026
Merged

KodBena merged 6 commits into
mainfrom
feat/orchestration-output-channel

Conversation

@KodBena

@KodBena KodBena commented May 19, 2026

Copy link
Copy Markdown
Owner

Summary

Two orchestration-substrate fixes that closed the SPA-side adaptive
reaping symptom surfaced 2026-05-18: the SPA's range-adaptive query
delivered partial data but no protocol-mandated
is_during_search=False finals, leaving the query state stranded
in-flight forever.

  • fix(orchestration): snapshot parent.opaque before Hub-side strip
    (d918bf8, originally added during the v1.0.26 testing cycle).
    Latent since v1.0.16: OrchestrationMiddleware.on_query schedules
    the driver task to run asynchronously; meanwhile the Hub's
    subscribe() pops capabilities from the same opaque dict the
    coroutine still holds a reference to. By the time the coro reads
    parent.opaque["capabilities"], the Hub has stripped it and the
    substrate silently falls back to closure defaults (no Phase 3
    fields, no per-query worst_quantile / extra_visits
    overrides). Fixed by deep-copying the opaque dict at on_query
    time and passing the snapshot to the coro factory.

  • feat(orchestration): push-based output via caps.send_response
    (this arc, 5 commits). Diagnosed and fixed a drain/driver race
    in OrchestrationMiddleware.handle_response: the prior
    await asyncio.sleep(0) + drain ctx._output_queue pattern only
    yielded one event-loop iteration. When an intermediate task layer
    sits between the input queue and the driver (ctx.parallel's
    collect tasks, or adaptive_reevaluate's _stream_parallel_spawns
    pump tasks), the driver's wake-up defers to iteration N+1 and the
    drain in N runs against an empty queue. Yields produced after the
    last input response — Stage 3 finalizations, fork-join summary
    emissions, framework-synthesised error responses — were stranded
    indefinitely.

    The fix: _drive_coroutine calls
    SessionCapabilities.send_response(parent_id, resp) for each
    yield, which wire-encodes under the parent's orig_id, logs
    lifecycle.forward, and ws.sends directly. handle_response
    becomes input-only for orchestration-managed orig_ids;
    ctx._output_queue and the drain loop are removed.

    Design analysis in proxy/docs/roadmap-orchestration-output-channel.md
    (status: design-note: implemented). Regression coverage at
    tests/test_orchestration_middleware.py::TestTrailingYieldAfterSpawnPrimitive
    (control via ctx.spawn direct iteration, regression via
    ctx.parallel).

Commits

aaeabf5 docs(orchestration): describe push-based output channel; mark design-note implemented
2db1718 feat(orchestration): push-based output via caps.send_response
57e7891 feat(orchestration): SessionCapabilities.send_response API
2120130 test(orchestration): regression for trailing-yield stranding
0416063 docs(orchestration): design note for the output-channel race
d918bf8 fix(orchestration): snapshot parent.opaque before Hub-side strip

Validation

  • Unit suite: 485 pass, 0 fail. mypy --strict clean on all
    touched production modules (middleware/orchestration.py,
    middleware/session_middleware.py, proxy_server.py).
  • Live wire test: SPA's range-adaptive query (62 deepening
    turns across 8 jerry-rig rounds during diagnosis; user-set
    parameters in production) now reaps all
    is_during_search=False packets per the KataGo protocol
    contract. User verified end-to-end.
  • Pinned regression (test_trailing_yield_after_parallel_reaches):
    was red on the pre-fix substrate, green under the new
    push-based output channel. Same coroutine skeleton, single-axis
    diff between ctx.spawn (control) and ctx.parallel
    (regression) isolates the bug to the scheduling asymmetry the
    design note analyses.

Wire-compatibility

Byte-identical on the wire to v1.0.26 under default and adaptive
configurations. The fix is purely internal to the orchestration
framework's output delivery; no client-visible protocol change.

The SessionCapabilities surface gains one optional field
(send_response). Existing constructor sites that don't wire it
get a NotImplementedError default per ADR-0002 — loud failure on
any path that actually needs it. Test fixtures across four files
updated to thread send_response through their fake caps; no
production caller outside the proxy itself uses the type alias.

Tag

v1.0.27 on merge.

🤖 Generated with Claude Code

bork and others added 6 commits May 18, 2026 23:25
Latent bug introduced by the v1.0.16 orchestration refactor:

OrchestrationMiddleware.on_query creates an async task that drives
the coroutine. The task runs LATER (next event-loop tick), but the
client query meanwhile flows through Layer 2 — and pubsub_hub.py:494
pops `capabilities` from the same opaque dict the coroutine still
holds a reference to. By the time the coro reads
`parent.opaque["capabilities"]`, the Hub has stripped it and the
substrate silently falls back to closure defaults
(worst_quantile=0.25, extra_visits=800, no Phase 3 fields).

Why it stayed latent through v1.0.24/v1.0.25:

  - The original benchmark intentionally opted out of capability
    middleware via `capabilities: {}`, never exercising the
    per-query-metadata path.
  - SPA users running with `PROXY_ADVERTISE_CAPABILITIES=false`
    take the legacy auto-engage path with closure defaults — the
    defaults happen to match what the SPA's UI would have sent
    (worst_quantile=0.05 default, extra_visits=800 default).
  - Phase 3 engagement requires `allocation_algorithm` to be
    non-empty in cap_meta, which only happens with the v1.0.26
    SPA dropdown. Phase 3.5's end-to-end test exposed the gap.

Fix:

OrchestrationMiddleware.on_query now deep-copies the query's opaque
dict before storing it as the parent_query and passing it to the
coro factory. The coro reads a frozen snapshot taken at on_query
time (BEFORE the Hub's post-subscribe strip). Original `query`
object continues to flow through Layer 2 unmodified.

Discovered by debug-instrumenting adaptive_reevaluate's cap_meta
read with the dropdown's `learned_v1` + `extra_visits=4800`
selection. Pre-fix dump:
  cap_meta={} opaque_caps=None  → phase3=off, extra_visits=800
Post-fix dump:
  cap_meta={'worst_quantile': 0.25, 'extra_visits': 4800,
            'value_binding': 'learned_v1',
            'allocation_algorithm': 'learned_piecewise'}
  → phase3=on, extra_visits=4800 ✓

Tests: all 496 existing tests pass; mypy --strict clean.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
Adds proxy/docs/roadmap-orchestration-output-channel.md
(status: design-note: planned). Companion to v1.0.16's
roadmap-orchestration-middleware.md.

Diagnoses the omitted-finals symptom surfaced 2026-05-19 under
adaptive_reevaluate's Phase 3 dispatch path: SPA sees no
is_during_search=False per analyzed turn, despite the
orchestration coroutine returning normal and KataGo having
emitted the originals plus all deepening sub-query responses.

Root cause: handle_response's drain heuristic uses a single
await asyncio.sleep(0) before non-blocking drain of
ctx._output_queue. CPython's _run_once snapshots ntodo at
iteration start; callbacks scheduled by call_soon during
the iteration land in _ready but defer to N+1. When the
driver task wakes directly from the input queue (the
async-for-over-ctx.spawn path), it runs in iteration N
alongside handle_response_task and its yields land before
the drain runs. When an intermediate task layer sits between
input queue and driver (ctx.parallel's collect tasks; the
analogous pump tasks in _stream_parallel_spawns), the
driver's wake-up defers to N+1; the drain in N runs against
an empty queue. For yields produced after the LAST input
response, no future handle_response invocation ever drains
the trailing window — Stage 3 finalizations and error-path
synthesised responses are stranded until GC.

The note covers:
- the per-stage scheduling trace (Stage 1 reaches; Stage 2
  lags one response; Stage 3 strands)
- why the v1.0.16 design didn't surface this (output-side
  timing was an implementation detail, not a contract)
- substitution-test severity calibration (error-path
  responses are the worst-case stranding)
- three fix options with tradeoffs; Option 1 (push-based
  output via caps.send_response) recommended
- six-commit migration sketch

The v1.0.21+ identity-type branding and the v1.0.24
multi-round substrate are unaffected. Wire shape unchanged.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
Adds TestTrailingYieldAfterSpawnPrimitive with two single-axis-
discriminating tests that pin the output-channel race the design
note at docs/roadmap-orchestration-output-channel.md analyses.

  test_trailing_yield_after_direct_spawn_reaches — control.
  Coroutine: `async for r in ctx.spawn(sub): yield r` then a
  trailing `yield canary`. Driver wakes from record.queue
  directly in the same event-loop iteration as
  handle_response's sleep(0); drain catches both the
  passed-through response and the trailing canary. PASSES.

  test_trailing_yield_after_parallel_reaches — regression.
  Same coroutine skeleton, single-axis change: `await
  ctx.parallel(a, b)` then trailing `yield canary`.
  ctx.parallel wraps each spawn in a collect task via
  asyncio.gather; sub-query response push wakes the collect
  task, not the driver. Driver's wake-up defers to iteration
  N+1; handle_response's drain in N runs empty.
  EXPECTED RED until the output-channel fix arc lands.
  Currently fails with yields_a=[], yields_b=[] — gather
  aggregation strands not just the canary but every
  sub-query response.

This commit deliberately leaves the regression test red as
the arc's pinning artefact. The design note's Option 1
(push-based output via caps.send_response) is the proposed
fix; that arc turns the test green.

Suite status with this commit: 28 existing pass, control
passes, regression fails. mypy --strict clean.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
Commit 1 of the orchestration output-channel arc (see
docs/roadmap-orchestration-output-channel.md).

Adds the push-based output channel that decouples a middleware's
output emissions from incoming-wire arrivals at handle_response.
The orchestration framework's driver task will use this in a
follow-on commit to deliver trailing yields (Stage 3
finalizations, error-path synthesised responses) that are
currently stranded by the sleep(0)+drain race.

Three additive changes:

  - middleware/session_middleware.py: new SendResponse type
    alias matching SubmitQuery/TerminateQuery's shape; new
    send_response field on SessionCapabilities with a
    NotImplementedError default per ADR-0002 (a constructor
    site that fails to wire it gets a loud error, not a silent
    drop).

  - proxy_server.py: new ClientSession._send_response method
    that wire-encodes the response under the given orig_id,
    logs lifecycle.forward at the kind-driven level, and
    ws.sends. Bypasses both the transformer chain (the data
    was already processed upstream of the middleware
    producing it) and the inner+outer middleware chain (the
    producing middleware is the one driving the call;
    routing back would recurse or require tagged-bypass
    logic). ConnectionClosed during send is non-fatal —
    trailing-window emissions may race with the client's
    disconnect; log and continue.

  - proxy_server.py:_handle_session: SessionCapabilities
    construction now threads self._send_response into the
    send_response field, completing the wiring on the
    real-session path.

No callers yet. The orchestration substrate still uses
ctx._output_queue + handle_response drain; the rewrite to
caps.send_response lands in commit 3. The regression test in
TestTrailingYieldAfterSpawnPrimitive::test_trailing_yield_after_parallel_reaches
remains red.

Suite: 484 pass, 1 fail (the pinned regression). mypy --strict
clean on the touched files.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
Commit 2 of the orchestration output-channel arc (see
docs/roadmap-orchestration-output-channel.md).

Closes the drain/driver race in OrchestrationMiddleware:

  - middleware/orchestration.py: _drive_coroutine no longer pushes
    yields to ctx._output_queue + sentinel; it calls
    caps.send_response(parent_id, resp) directly for each yield from
    the user's coroutine. The per-context output queue and the
    handle_response drain loop are removed. handle_response becomes
    input-only — it routes incoming responses to the parent's
    _original_queue or a sub-query's record.queue and yields nothing
    for orchestration-managed orig_ids; pass-through (non-
    orchestrated) orig_ids still yield through as before.

    The exception-path error response is also delivered via
    send_response — trailing error envelopes that previously could
    have been stranded in the post-drain trailing window now reach
    the wire reliably.

  - tests/test_orchestration_middleware.py: _FakeSessionCapabilities
    gains synthetic_sends + send method + send_response in
    as_session_capabilities; _drive_response harvests both
    handle_response yields and new synthetic sends via a brief
    settle window. test_sub_query_response_relabels_through_gate
    updated to assert relabel lands in caps.synthetic_sends rather
    than in handle_response yields.

  - tests/test_phase3_dispatch.py, tests/test_capability_negotiation.py,
    tests/test_multi_round_adaptation.py, tests/test_allocation_refusal.py:
    each file's _make_caps and _drive_response/_drive helpers updated
    to thread send_response through and harvest from synthetic_sends.
    _settle_and_drain in test_phase3_dispatch.py now returns ALL
    synthetic sends for the requested orig_id (no more direct
    _output_queue poking — that abstraction is gone).

The previously-pinned regression at
TestTrailingYieldAfterSpawnPrimitive::test_trailing_yield_after_parallel_reaches
turns green: the trailing canary now reaches caps.send_response and
is therefore delivered to the wire. Same fix applies to Stage 3
finalization emissions in adaptive_reevaluate — the original
omitted-finals symptom is structurally closed.

Suite: 485 pass, 0 fail. mypy --strict clean on the three touched
production modules.

This commit is internally cohesive: the substrate change required
the test-harness updates to stay green. Splitting them would have
left an intermediate state with ~30 tests red, which the design
note acknowledges but isn't worth the per-commit churn.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
…note implemented

Commit 3 of the orchestration output-channel arc (closing).

- FRAMEWORK.md: orchestration section gains a new subsection
  describing the push-based output channel — middleware authors
  write `yield resp` as before; the framework's driver task
  delivers via caps.send_response. handle_response is correspondingly
  input-only for orchestration-managed orig_ids; pass-through
  yields preserved.
- ARCHITECTURE.md: orchestration section names output delivery as
  one of the things "the framework owns"; SessionCapabilities
  description names the three callbacks (submit_query,
  terminate_query, send_response) explicitly.
- docs/roadmap-orchestration-output-channel.md: sibling-revised
  per ADR-0005 Rule 8 from `design-note: planned` to
  `design-note: implemented`. Adds implementation footprint
  (which files changed, what wasn't needed from the planning-time
  arc, which planning-time commits were merged for cohesion).

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
@KodBena
KodBena merged commit b871127 into main May 19, 2026
13 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant