Repository navigation
feat(orchestration): substrate fixes for SPA-side adaptive reaping (v1.0.27) - #39
Merged
Merged
Conversation
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>
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.
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=Falsefinals, leaving the query state strandedin-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_queryschedulesthe driver task to run asynchronously; meanwhile the Hub's
subscribe()popscapabilitiesfrom the same opaque dict thecoroutine still holds a reference to. By the time the coro reads
parent.opaque["capabilities"], the Hub has stripped it and thesubstrate silently falls back to closure defaults (no Phase 3
fields, no per-query
worst_quantile/extra_visitsoverrides). 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 priorawait asyncio.sleep(0) + drain ctx._output_queuepattern onlyyielded one event-loop iteration. When an intermediate task layer
sits between the input queue and the driver (
ctx.parallel'scollect tasks, or
adaptive_reevaluate's_stream_parallel_spawnspump 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_coroutinecallsSessionCapabilities.send_response(parent_id, resp)for eachyield, which wire-encodes under the parent's orig_id, logs
lifecycle.forward, andws.sends directly.handle_responsebecomes input-only for orchestration-managed orig_ids;
ctx._output_queueand the drain loop are removed.Design analysis in
proxy/docs/roadmap-orchestration-output-channel.md(status:
design-note: implemented). Regression coverage attests/test_orchestration_middleware.py::TestTrailingYieldAfterSpawnPrimitive(control via
ctx.spawndirect iteration, regression viactx.parallel).Commits
Validation
mypy --strictclean on alltouched production modules (
middleware/orchestration.py,middleware/session_middleware.py,proxy_server.py).turns across 8 jerry-rig rounds during diagnosis; user-set
parameters in production) now reaps all
is_during_search=Falsepackets per the KataGo protocolcontract. User verified end-to-end.
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) andctx.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 itget 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