Skip to content

feat(gateway): serve non-chat protocols through pipelines - #475

Draft
Menci wants to merge 110 commits into
feat/pipeline-corefrom
feat/pipeline-non-chat
Draft

Menci wants to merge 110 commits into
feat/pipeline-corefrom
feat/pipeline-non-chat

Conversation

@Menci

@Menci Menci commented Aug 16, 2026 •

Copy link
Copy Markdown
Owner

Stacked on #474; remains draft until that dependency lands.

Serves embeddings, rerank, Completions, image generation/editing, audio transcription and Alpha search through immutable fact pipelines. Paths enter provider-owned stages and the shared HTTP terminal; content, upload bytes, upstream status and each actual billable observation remain visible. Candidate resolution, failover, serialization and settlement are shared. Models use main's nonblocking catalog refresh.

The gateway recorder delegates to @floway-dev/dump, supplying storage/publication ports while retaining attribution, drain, HTTP measurement and retention. Completed records are stored after owned readers/deferred outcomes drain. The inspector reconstructs stages, request/response values, repeated descents, value changes, logs, failures and literal frames with their completion markers. Shared metadata types come directly from dump; event types come from pipeline. Legacy chat record contracts remain gateway-owned until #476.

Read one Embeddings path first, then provider stages/HTTP, settlement, the thin dump/run-sink.ts adapter and stage viewer. Live writer/reader integration is delivered separately in #565; realtime UI and frame-sequence diff remain deferred.

Validation at 09c478519: all eight exact-head CI checks passed, including full tests, typechecks, genuine installer harness, repository invariants and Web build/output checks. Focused extraction checks passed 145 tests, plus 10 settlement/recording cases after replacing the legacy context field. Independent review confirmed recording order, schema/wire behavior, browser-safe imports and the narrowed run context.

Menci added 7 commits August 16, 2026 20:14
The first half of migrating the non-chat families: the space every family's
pipeline extends, the services a run is given, and the two stages every family's
serve pipeline needs.

Neither shared stage names a family's own key. Branching is a capability of the
framework and what something branches on is a concept in the domain, so a family
hands in the two domain-shaped things — how to narrow a candidate, and how to
read an attempt's outcome — and the stages stay written against the shared space.
That is what lets them compose into a pipeline over any family's larger space
with no variance question to lose, since assembly reasons over declarations and
declarations are strings.

`resolveCandidates` carries both traits, and it is what found the modelling error
corrected in the commit below it: refusing a request no upstream can serve means
answering with a key the descend path never carries.

The families themselves follow. Nothing is routed through this yet, so the
gateway's behaviour is unchanged.
The first family, and the one worth doing first: rerank already has a canonical
contract — parse, render, serialize, and a usage reading — so it proves the
shared stages on real business without needing a protocol written first.

Four stages, and the shape every family will repeat. `emitRerank` is the edge: it
renders the canonical answer back into whichever of the four protocols the client
spoke, which is why the source protocol is an ingress fact and not a request one
— it has to survive the switch to whatever the upstream turned out to speak.
`callRerankUpstream` is the ending: it dials, parses the body, and provides both
the answer and what the call is billable for.

Two behaviours change, both because the architecture has no passthrough. The
same-protocol path used to forward the upstream's bytes unread; it is now parsed
and re-serialized like the translating path, so one code path serves both and the
dump shows a parsed answer either way. And an upstream that refused is a value
rather than a forwarded Response, so an earlier stage can fail over it.

The client's headers reach the upstream from the record rather than from a live
request object, which is what lets the dump show what was there to be filtered.

The test that matters most is the one that records a hole. `callRerankUpstream`
reads `ingress.http.headers`, and the derived entry contract does not mention it:
a stage whose only trait is `return` declares no request side, by ruling, so
assembly cannot see what an ending stage reads — and every family's ending stage
reads something. The consequence is concrete: a caller omitting that key gets a
runtime failure at the deepest stage instead of an assembly error, which is the
thing the entry contract exists to prevent. The type layer still catches it at
the definition site, so it is a gap and not a break.

Nothing is routed here yet; the route table still runs the passthrough serve.
The simplest family, and the one that shows the shape with nothing else in the
way: one protocol, no translation and no stream, so the array is exactly the
four stages every family has. The narrowing is a constant rather than a
function of the request, because an embeddings request carries nothing a
candidate could be incompatible with.

`packages/protocols/src/embeddings` was five lines of index signature — a
passthrough shape, not a contract — so the real one is written here: the
request's four input arms, `encoding_format`, `dimensions` and `user`, and the
response's embeddings, model and usage, each cited against the OpenAI
specification.

`encoding_format` is why the canonical response holds numbers rather than
whatever the upstream wrote. Both official OpenAI SDKs send
`encoding_format: base64` when their caller did not choose one, so base64 is
the common case on the wire; an upstream that ignores the field and answers
with float arrays hands such a client a body it decodes as base64 and turns
into noise. Parsing whatever arrived and writing what the client asked for is
what closes that, which is why the encoding is an ingress fact and stays put.
A base64 embedding is the vector's float32 elements packed little-endian, and
reading it into numbers is exact: every float32 is a float64, so a vector
survives any number of trips.

Two behaviours change. A field outside the protocol is now named in a 400
rather than forwarded, because the request schema is `additionalProperties:
false` and both gateways compared against — LiteLLM and copilot-api — draw the
same line; this is the one decision worth a second look, since vLLM and Jina
document supersets of the endpoint. And an upstream that reports no usage now
bills an entity with no quantities, which is how "called and reported nothing"
is said, rather than no row at all.

The test records what assembly cannot see, and this family shows it twice
over: a return-only stage declares no request side, so neither
`ingress.http.headers` nor `request.embeddings.canonical` reaches the derived
entry contract — and unlike rerank, whose edge needed the payload for
rendering, nothing here puts the payload in the contract by accident.

The route still runs the passthrough serve. The old handler parses and
serializes through the new contract so the two agree until it is deleted.
Five families at once, each with a real protocol contract where it had none.
`completions`, `embeddings`, `images` and `audio` ran on `passthrough-serve.ts`,
which forwards a body it has not parsed — and the architecture has no such
concept, so three of them needed a contract written before they could have a
pipeline at all. `packages/protocols/src/{audio,images,completions,embeddings}/`
carries those now, cited to the OpenAI specification.

`alpha-search` is shaped differently and is not forced into the others' shape: it
has no upstream model and no candidate list, executing locally through the
configured search provider or relaying a selected provider's response. That is an
ending, not a missing family — the same distinction ruling 5 settled for HTTP,
where a provider that does not ultimately speak plain HTTP supplies its own.

Nothing is routed here yet: the route table still mounts the passthrough handlers,
so the gateway's behaviour is unchanged. The wiring, and deleting the passthrough
helpers, is one commit once every family is sound.

Two defects the families found in what they were built on are recorded in the
Pull Request and fixed in the commits that follow. Both were found by building on
the core rather than by reading it, which is what this layer is for.
`failover` declared that it consumes and provides `response.http.body` for every
family. Three of the six never produce one — rerank, embeddings and images read
their answer to the end, so there is nothing still open by the time the fork sees
it — and the runner checks `provides` at handover. Those pipelines composed
cleanly and would have thrown on the first real request, at the deepest stage,
naming a key their author had never written down.

It cannot be a fixed key in either direction: claiming one a family never produces
throws, and staying silent about one it does produce and hands up throws the other
way. Which keys carry a resource is a statement only the family can make, so it
makes it — alongside the failure predicate it already hands in, for the same
reason. Branching is the framework's; what something branches on, and what it
owns, are the domain's.

Found by the completions family, which verified it against a throwaway pipeline
rather than reasoning about it. The three affected families' own tests asserted
`entryNeeds` and nothing else, so none of them caught it — the same shape of
untested claim the core's review found earlier, appearing again one layer up.

The test added here asserts the declaration for both kinds of family rather than
the absence of the bug, because the absence of the bug is what the other tests
were already asserting when it was present.
Menci added 22 commits August 16, 2026 22:13
Settlement is a stage now, and it sits above `failover` so a run bills once
however many candidates it tried. Repetition passes through the stage that
observes usage, not through this one.

It is unconditional. A run that measured rather than generated still writes, and
its row simply names no billed entity — emptiness is observed rather than
declared, which is why the word "unknown" appears nowhere: the situations are
concrete and the list is open. The write is scheduled rather than awaited,
because a transient repository failure must not turn an already-flowing upstream
response into a 502.

`BillableEntity` gained the pricing inputs a rate can need beyond the quantities.
Absent is a real reading and not a missing one — most families price on the
quantities alone — and making it required broke five families that correctly have
none.

It also caught a regression I had introduced: the rerank migration dropped the
`inputTokens` pricing fact that `settleRerank` passed, so a rerank rate depending
on input size would have priced against nothing. Restored, and the reading now
travels with the entity it prices rather than being recomputed at the write.
A `ModelCandidate` is two things wearing one type: the selection — which
upstream, which model row, which flags — which is data, and the handles — the
provider instance, the fetcher, the models cache — which are live. The six
families put the whole thing in the record, so `move()` deep-froze all three
handles and the SWR cache refresh the provider does on its own schedule broke.

The architecture ruled this out before any code existed, and I did not apply it:
"Where something is chosen per attempt, the resolver is the service and the
selector is a fact. A per-upstream transport is not a fact and is not pinned at
the prologue either; what is injected is the thing that resolves one, and what
travels is the identifier it resolves from."

So `route.candidate` becomes `route.attempt`, carrying the upstream id, the model
id and a snapshot of the flags — snapshotted rather than referenced, because the
record must show what was true when the attempt was made rather than what the row
says now. `resolveAttempt` is a service and the ending stages ask it for the thing
that dials.

It also fixes the same rule failing the other way: a candidate in the record is
walked by the dump encoder, so a run's dump was serializing the provider instance.

Nothing caught it because every family test builds a candidate literal with no
live handles, and no route mounts a pipeline, so the suite never runs one against
a real provider. The test added here asserts the property in both directions —
the selector reaches nothing live, and the candidate would have frozen it.
`writeSettlement` was written and composed into nothing, so no migrated family
recorded a usage row or a performance sample. Every serve pipeline now carries
it, above the fork, so a run bills once however many candidates it tried.

Adding it proved its own absence: the completions tests immediately failed with
"Repo not initialized", because until now nothing in six families ever reached a
write. That error was the fix working.

The entry contracts gained `ingress.http.headers`, which five of the six omitted.
`compose` cannot derive it — a stage whose only trait is `return` declares no
request side, and every family's ending is one — so the hand-written type is the
only thing covering for that, and it was covering wrong.

One test here is worth reading twice. I first asserted that a run reaching no
upstream still records a performance sample; it does not, because
`recordPerformance` returns early without an attempt's telemetry and there is no
attempt. That matches the replaced surface exactly, so the expectation was
invented and the stage was right. The test now says what happens and why.
The two streaming families marked their body with `Object.assign(body, {
[Symbol.asyncDispose]: … })`, which was how a resource was claimed before
ownership stopped being a detection. The runner reads `isOwned` now, so a body
marked that way was invisible: `failover` declared it consumes one, `drain()`
existed, and neither could see anything. A losing attempt's connection stayed
open and the winner's was never drained.

`own(body, release)` makes the claim, and the fact's type says `Owned` rather
than `AsyncDisposable` — the language puts `Symbol.asyncDispose` on every async
generator and on no `ReadableStream`, so a structural type admits an iterator
that is not a resource and rejects the body that is.

The test asserts the property rather than the absence of the bug: a losing 429's
body and a winning stream, the answer handed back before anything is drained, and
the winner drained when the caller says so. Verified by mutation — dropping the
brand from `own()` fails exactly that test and nothing else.
Both families rendered an answer and never said what status it was, so an
upstream 429, a resolver's 404 and a 400 all reached the client as a 200 carrying
an error envelope. That is not a difference the no-passthrough ruling asks for:
a client is not owed the upstream's exact bytes, but it is owed the truth about
what happened, and the replaced surface forwarded the status.

The edge provides `response.http.status` on both arms — the upstream's own when
it refused, the gateway's own when the resolver refused before dialling, and 200
for an answer.

The three tests drive the pipeline rather than reading its declarations, which is
what the two families were missing: their whole suites asserted `entryNeeds` and
nothing else, so a family that could not express a status at all passed. Verified
by mutation — pinning the failure arm back to 200 fails both refusal tests.
… unnamed embeddings model

The pipeline's shared resolver phrased its own refusal, which changed what a
client is told: five families spell the endpoint they could not serve, while
rerank spells out which part of the request no candidate could satisfy. The
sentence moves to the family as `Narrowing.unsupported`, taking the reasons
`reject` gave so a family can use them or ignore them.

Copilot's /embeddings answers without the top-level `model` the schema marks
required, so parsing refused a body the replaced surface served. The parser now
takes the model the request named and completes the record with it.
…'s edge

The pipeline dropped the upstream's response headers, so vendor traces, quota
state and retry-after stopped reaching clients on the routes it took over.

The attempt hands them up unfiltered and the edge decides what a client may see,
which keeps both in a dump. A refusal that never dialled has none to carry, so
the shared resolver answers with an empty list and no family repeats it.

serveThrough now reads the status and headers off the exit facts. A family's
render owns the bytes and their media type and nothing else, which is also why
content type stays blocked from forwarding.

The assembly test builds every family, including the four whose pipeline is
built from the request and would otherwise never be composed by a test.
A stage logger wrote to the dump record and to whatever sink the services
carried, and the prologue carried none — so on every run without retention
configured, a warning or an error reached nothing at all. A failed usage write
was the case that showed it.

The sink writes warnings and errors to the console the rest of the gateway uses.
Debug and info describe one request's progress and would bury the request log at
a line per stage, so the threshold is fixed here; a dump still records every
level when one is open.
…ot faults

Nothing in the domain throws. A refused connection, a timeout or a reset is an
outcome the fork has to see so it can try the next candidate, so each ending
catches what the platform raised and hands it up as its family's failure. A dial
that reached no upstream bills nothing and carries no headers; only the
performance row records that the attempt happened.

Attribution moves ahead of the dial. It was written from the call's result, so
an attempt that never returned left the previous candidate's context in place
and misattributed the row.

A 2xx body a JSON protocol cannot read is likewise a failure value rather than
an unhandled parse error: the gateway that cannot read an answer has not served
the request, and the upstream still counts as called having reported nothing.
The body reader images already had is now shared, since three families need it.
Both handlers still ran the passthrough scaffold: they read the body, hand-rolled
the multipart and JSON edit shapes into a provider request, picked a candidate by
endpoint capability and relayed whatever the upstream sent. The images protocol
already owns both readings and the family already has a pipeline, so what is left
here is a prologue and an epilogue — parse, hand over, write the answer.

Reading the body moves to the contract, which reports a malformed request by
throwing; the handler turns that into the same 400 envelope the scaffold wrote,
with the same sentences. The two endpoints differ only in which parser reads the
bytes, so the rest is one shared function.
…r reads them

Workspace packages sort after the relative imports, and the two handlers wired to
their pipelines had them first. Nothing else changes.
…gs independent

An upstream that refused in its own words is handed on in them. The envelope was
written once for OpenAI-shaped families and applied everywhere, which turned a
rerank client's `{message}` into `{error:{message}}` — a shape its SDK does not
read. What distinguishes the two cases is whether an upstream answered at all,
not which shape it answered in, so that is what the renderer now tests.

Rerank reads usage before results. They are independent readings of one body and
one failing is no reason to discard the other, so an answer this gateway cannot
model still bills for what the upstream metered.

A same-protocol answer is rendered back out unchanged, so results it carries
that the canonical form does not model are no longer a refusal — only a
translation needs to read them. A translation that fails on an answer that
parsed is 502: the upstream answered and the gateway cannot say it in the
client's protocol.
Three tests described the passthrough surface that /v1/embeddings and /v2/rerank
no longer run through, and two of them asserted behaviour the design replaces: a
2xx body a JSON protocol cannot read was forwarded verbatim with the upstream's
own status. A gateway that never read the answer cannot claim to have served the
request, so it refuses in its own words and the upstream still counts as called
having reported nothing.

The rest are kept as they were, against the pipeline: a failed usage write
leaves the answer alone and is still reported, and the last candidate's refusal
reaches the client with its status, its headers and its own words.
Whether a request streams is written in the request — `stream: true` in a JSON
body, a form field in a multipart upload — and the body can only be read once.
The prologue read it internally and took the answer as a parameter, so a family
that learns it from the body had no order in which to call it.

Guessing `false` is not harmless: the abort controller a streaming run cancels
its upstream with is minted from that flag, so a client that disconnected would
stop cancelling anything. Reading the ingress is now its own step, and the run
opens once the handler knows what was asked for — which also lets the requested
model reach the dump through the context that stamps it rather than a second
call afterwards.
…t ends

A streaming family's answer is the stream, and the seam could only build a
response from bytes it already had. It now takes frames as an answer of its own
shape: the response is staged on the context first, because hono's SSE helper
constructs the response itself and headers passed to a constructor would be lost.

A stream's usage arrives with its last chunk, which is after the run has
answered — the fact carrying it is a promise, and settling from it is the
epilogue's job. Settlement's write is extracted so both paths do the same thing:
the stage settles what the ending had already read, and the epilogue settles what
only the drained stream could say.
Completions and audio were the only families whose edge handed the upstream's
headers on unread, so a body the gateway re-serialized carried the upstream's
own content-length and content-encoding — a client would have been told the
wrong length for bytes it was actually given. Both now filter as the other
families do, and both exits declare the key the seam reads off them.
The release was scheduled the moment the run returned, so for a streaming family
it consumed the very frames the client was waiting for — one connection has one
reader, and the two raced for every chunk. The terminal [DONE] was the frame
most often lost.

Reading the frames out to the client is what releases the body they came from,
so the drain now waits for that to finish, in a finally so a client that stopped
reading still leaves nothing open. A buffered answer was serialized from facts
the run already held and still releases at once.

Completions serves through its pipeline; the streaming half settles from the
usage its last chunk carries.
The reader ran until the upstream closed, so an upstream that holds the
connection open past transcript.text.done held the client's stream open with it.
Returning at the terminal event closes the read, which cancels the upstream —
the behaviour the replaced surface had.

The route stays on its existing surface. Migrating it needs a decision this
change cannot make: srt and vtt are parsed into cues and rendered back, which
the protocol states is not byte-identical to what the upstream sent, and the
replaced surface forwarded those documents unchanged.
Two families' endings dialled outside the try the other three had, so a refused
connection ended the run as a fault: no failover to the next candidate, no
performance row, and a 500 where the replaced surface answered 502. Attribution
moves ahead of the dial with it, so an attempt that never returned names its own
candidate.

A streaming answer returned hono's response without finalizing the gateway
context, and finalize is what writes the dump — so a streamed request recorded
nothing at all. The upstream's frames are recorded where they are read, before
the edge decides which of them the client sees.

The prologue's dump sink was a stub that discarded every event while still
telling the runner to accumulate them, so a key with retention configured paid
for a recording nobody could read. It is absent until it has somewhere to go.
A request carrying stream: true reaches the upstream, and the pipeline cannot
read what comes back — the images stream is not carried as facts, so the answer
became a 502 after the call had already been made and charged for. The replaced
surface forwarded those frames.

This is the same blocker audio has: a response shape the fact space does not
model yet. The pipeline, its tests and the protocol contracts stay; what is
withdrawn is the route, until the shape it cannot carry is carried.
…t two families

Completions and audio wrapped an upstream's error body in a second envelope, so
a client read an escaped JSON document where the message should have been. They
now hand the upstream's own object on, as the other three families do.

Three ingress keys were declared and never provided or read, under a comment
claiming every request recorded them. The method, path and body are the dump's
to record from the request context; what a family actually hands over is the
headers, because every ending forwards what a provider may of them.

Rerank's attempt module had no referrers left after the pipeline took the route.

Two comments cited design documents that are not in the repository.
The family produced a billed set and composed no settlement, so nothing it did
was ever written. What stopped it was a request-side need on serve.model that
the stage never read: declaring it put a resolved model in the entry contract of
a family that resolves none. Settlement reads what came back and nothing on the
way down, which is what the declaration now says.
@Menci Menci changed the title feat(gateway): serve the non-chat families through the pipeline feat(gateway): serve non-chat protocols through pipelines Sep 30, 2026
Keep successful upstream HTTP statuses through protocol emission. Record each billable observation while summing request diagnostics, and derive output performance from the observed quantities even on partial failures. Normalize audio uploads into portable content facts before recording.
…peline-non-chat

# Conflicts:
#	packages/pipeline/src/run.ts
Enter provider-owned fact spaces through operation-specific chains and finish each supported request with the shared HTTP stage. Keep URL, authentication headers and immutable request content visible as facts, and track each response body through one ownership identity.

Use portable model flags and file snapshots, stream replayable multipart content without native file facts, preserve transport error chains, and retain replaced Codex endpoint observations and full 401 diagnostics across credential retries.

Validation: HTTP/provider/test-utils typechecks passed; 124 HTTP and provider test files passed with 1658 tests; targeted ESLint and staged diff checks passed.
Use portable model and upload facts for typed provider handoffs and preserve a single owned HTTP body. Collect provider retries and candidate attempts as separate billing observations, retain pricing facts and measured zeros, and settle each logical run at its exit. Preserve original diagnostics when deferred settlement marks a failure and recover completed calls on exceptional exits.
Build a stage tree from the shared object space, preserve repeated descent and folded response structure, and show request/response facts, identity-aware changes, deferred outcomes and structured logs. Keep raw events available on demand. Verified a production build in Chrome across light, dark and narrow layouts, including Monaco rendering and stage selection.
Persist encoded facts and protocol frames with backpressure, attach an authenticated temporary byte reader, and keep durable recording active when live storage fails. Renew live and staged-file leases while waiting for upstream and close only after run drain.

Bind settlement and deferred errors to the run, retain completed call observations on later failures, preserve defined zero usage, and record client frames after filtering. Shield owned meter readers from early consumer return so teardown can drain their usage.
Keep a stage error distinct from a successful child result and show it through the reader UI. Decode the recorded object-space error and preserve the successful child facts for inspection.
Await scoped transport diagnostics and retain quota persistence rejection in its owned deferred outcome. Keep the intermediate provider factories valid under the shared lint contract.
Observe released streamed resources and restored usage context through a real run, exercise stage-to-events navigation, and prove quota persistence rejects drain without replacing the HTTP response.
Encode the final JSON body in a response stage so codec failures retain failure facts and measured usage. Preserve canonical protocol content for diagnostics, forward encoded bytes from HTTP, and leave streams and upstream documents to their owning transport.
Forward the client JSON bytes produced by the response stage, preserving parsed search content and its observed call collection if encoding fails. The route regression compares stored body bytes with the real HTTP response.
Retain the encoded run object space and wait for owned readers and deferred outcomes before storing completed records. Defer live writer integration and its business reader to the later realtime work; preserve that implementation on the local archive branch.
Resolve stream references through the shared run object space, show literal decoded frames, and distinguish streams with and without their recorded end marker. Keep frame comparison algorithms and realtime business integration deferred.
Reject the deferred outcome with the original quantity conversion error instead of leaving run drain pending after terminal output. The real HTTP regression supplies overflowing numeric usage, verifies background completion, retains the actual call observation and records failed diagnostics.
# Conflicts:
#	packages/gateway/__tests__/dump/test-fixtures.ts
#	packages/gateway/src/dump/types.ts
Delegate event encoding and stream references to the shared recorder through explicit storage and broker ports. Keep admission, attribution, drain, HTTP measurement and dual-shape persistence in the gateway adapter.

Narrow pipelined request contexts to run recorders and import metadata directly from its owning package. Preserve domain behavior and archival metadata; 145 focused tests, gateway/web typechecks and scoped lint pass.
Omit the legacy dump field before supplying the concrete recorder, so pipelined services expose only the supported run methods and finalize signatures.
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