Repository navigation
Conversation
This was referenced Sep 17, 2026
…oducers
A second, plaintext, loopback-only HTTP listener on nvpair-workload-manager
accepts the same workload:* / workloads:remove JSON-RPC frames as the
cluster-mTLS peer port. It stamps originatedFrom with this node's UUID when the
producer leaves it empty, emits the frame to the broker as workloads:upsert /
workloads:remove, and then tracks it for re-sync and queues it for peers as a
local-origin event. Work that never went through PAIR's own proxies (an
external scheduler, a local inference harness) therefore reaches the local
store, the Jobs list, the persisted history and the scheduler's pending counts.
Off by default. Enabled by --local-ingress <host:port> or, when the broker
launches the worker without the flag, by <appdir>/workload-ingress.json
({"listen": "127.0.0.1:14324"}). A non-loopback address is refused.
The ingress is added through Manager.EnableLocalIngress, so NewManager keeps its
signature. Ingress frames follow the broadcast ordering invariants: the broker
write comes first and a failed write leaves no tracked state and no queued peer
frame; ingestMu keeps concurrent posts for one workload in the same order at the
broker and in the peer queue; and the track+enqueue step takes broadcastMu, so a
re-sync snapshot cannot enqueue a workload after its removal. The stdio handlers
share that step and keep their behavior.
The listener is not trusted merely for being on loopback, since a web page can
reach it. Before it reads a body it refuses a request that carries an Origin
header (403), whose Host is not a loopback name or address with the bound port
(421), or whose Content-Type is not application/json (415); the body is capped
(413) and the server has header, read and idle timeouts (no write timeout, which
would also bound the wait for the broker and drop the response of an accepted
frame). It also accepts only
frames about workloads that run here: originatedFrom must be empty or this
node's UUID, a lifecycle state must be one of the known states, and a workload
id is at most 256 bytes. Because encoding/json matches field names without
regard to case, a field name spelt other than exactly, or two names of one
object that fold together, is refused (400), as is a resync name in any
spelling, and the decoded origin, state and id are validated again after
parsing. Fields the ingress does not know are passed through, so producers send
metadata only. The stdio and peer paths validate as before.
A workload removed through the ingress is remembered for a minute by its
(origin, id) pair, so the broker's replay of its store to a restarted worker
cannot track it again and re-assert it to peers; a new ingress frame for the
pair clears the memory. At most 4096 removals are remembered, the oldest
forgotten first.
Co-Authored-By: Claude Fable 5.1 <[email protected]>
Co-Authored-By: Claude Opus 5.5 <[email protected]>
Signed-off-by: Daniel Martinez <[email protected]>
dmmdea
force-pushed
the
feat/workload-local-ingress
branch
from
September 29, 2026 22:29
cd31a1c to
78bc20b
Compare
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.
Description
The Jobs list only shows work that entered the cluster through PAIR's own two proxies. Anything else running inference on a member (an external scheduler, a local harness that talks to its engines directly) is invisible: not in any node's Jobs list, and not in
nvpair-job-scheduler's pending counts, so PAIR routes its own traffic onto a GPU that is already busy with that work. Related asks: #7 (host-wide admission signals, external placement hooks) and #67 (third-party hosts integrating with the services).The workload manager's only listener is the inter-node port (
:14320, cluster mTLS, pinned peers only), so there is no door for a third-party producer on the same machine, and the broker's stdio is owned by the desktop application.This adds an opt-in, loopback-only, plaintext ingress on
nvpair-workload-manager:--local-ingress <host:port>, or, when the broker launches the worker without the flag,<appdir>/workload-ingress.jsoncontaining{"listen": "127.0.0.1:14324"}(file-registered, like an engine manifest under<appdir>/engines/, so it can be enabled without a desktop change). Absent, empty or malformed means off. A non-loopback bind address is refused at startup; a port already in use fails startup loudly.POST /v1/workloads/eventstakes the same JSON-RPC 2.0 frames as the peer port (workload:submitted|started|completed|erroredwithparams.workloadInfo,workloads:removewithparams.workloadId) through the sameparseLifecycle/parseRemovevalidation and a 1 MiB body limit (413above it).originatedFromis stamped with this node's UUID when the producer leaves it empty, the same courtesy the broker extends to the proxies; naming any other node is400. Producer mistakes are400, other methods405, and a broker that cannot be written is500. Frames may carryseq, which is part of the key peers deduplicate on, andcancelledis a validstate.127.0.0.1, and a DNS-rebinding page reaches it under its own hostname. Before the body is read, a request is refused when it carries anOriginheader (403), when itsHostis notlocalhost, a127.0.0.0/8address or::1with the bound port (421), or when itsContent-Typeis notapplication/json(415). The server sets header, read and idle timeouts and deliberately no write timeout: net/http's clock covers the whole handler, so a slow broker write would drop the response while the frame still lands and the producer would retry an accepted frame. A lifecyclestatemust be one of the five known states (400otherwise), since any unknown value would be tracked as active with no expiry and re-asserted to peers forever, and a workload id (workloadInfo.id,workloadId) is at most 256 bytes.encoding/jsonmatches field names without regard to case (with Unicode simple folding), so a name that differs from the canonical spelling only that way (originatedfrom,WorkloadInfo, a long-sstate), or two names of one object that fold together, is400; otherwise a case variant would be decoded as the field while slipping past the exact-key origin check. Aresyncname in any spelling is400too (it is the peers' own marker for a re-assertion and would bypass their dedup), and the decoded origin, state and id are validated again after parsing, so the rules hold even for a spelling the key check missed. These rules apply to the ingress only; the stdio and peer paths validate as before.workloads:remove, which would track the workload again and re-assert it to peers on every heartbeat. A workload removed through the ingress is remembered for one minute (the terminal-retention window), at most 4096 removals with the oldest forgotten first and expired ones pruned in amortized constant time, by its(originatedFrom, id)pair; a non-re-sync lifecycle frame from the broker for it is neither tracked nor queued, and a new ingress lifecycle frame for it clears the memory.workloads:upsert/workloads:remove(that upward emission is what updates the local store, the Jobs list, the persisted history and the scheduler's counts, and what the stdio path deliberately skips because the broker has already applied its own frames), and only then recorded in the re-sync set and queued for peers.200means the frame was written to the broker and queued for peers; the peer side stays best-effort, like every broadcast.The ingress is wired through a new
Manager.EnableLocalIngress(addr)thatmaincalls afterNewManager, soNewManagerand its callers and tests are unchanged.No desktop change is needed: the Go side treats
engineas an opaque string,WorkloadItemCardrenders an unknown engine without an icon, andformatModelDisplayNamehas a generic branch. A follow-up could add a generic engine icon and surfacerequesterIdon the card.Producer example:
Release intent
Changelog title
Third-party local workload producers can report their jobs
Changelog body
--local-ingressflag or an<appdir>/workload-ingress.jsonfile, listens on loopback only, and refuses any other bind address.Bumps
An additive HTTP feature is a MINOR bump per
services/VERSIONING.md;servicesfollows the component.services/versions.jsonis not touched.Scope
In:
localingress.go), its wiring inmain.goandManager.EnableLocalIngress, including the request trust checks (field-name spelling, decoded-value validation, id length, no write timeout) and the bounded removal memory described above (removedIngress/removedOrder, guarded byactiveMuand only touched underbroadcastMu).manager.go:handleLocalLifecycle/handleLocalRemoveand the ingress now shareparseLocalFrame(validation, no side effects) andcommitLocal(track and enqueue underbroadcastMu). The stdio handlers behave as before.README.md(flag row and ingress section) and the normativespec.md(interface §7.4, requirements, failure mode).Out:
Concurrency, since the ingress adds a concurrent HTTP entry point next to the single-goroutine stdio read loop:
ingestMu->broadcastMu->activeMu. Only the ingress takesingestMu.commitLocal(track and enqueue underbroadcastMu). A failed broker write leaves nothing tracked and nothing queued, so peers never show a workload the origin broker does not hold, and the producer can retry.commitLocalholdsbroadcastMuacross track and enqueue, exactly like the stdio handlers, so a re-sync snapshot cannot enqueue a workload after its removal was enqueued.ingestMukeeps concurrent posts for one workload in the same order at the broker and in the peer queue. It is deliberately notbroadcastMu: the broker write can block, and the stdio read loop and the re-sync heartbeat takebroadcastMu.broadcastLoopowns those.Validation
In
services/nvpair-workload-manager:go vet ./... && go test ./... -count=1on Windows: pass.go test ./... -race -count=5on Linux (Go 1.27.1, amd64): pass.localingress_test.go): origin stamping for lifecycle and removal frames, the node's own origin accepted and a foreign one rejected (atingestLocaland as400over HTTP), missingworkloadInforejected, unknown states rejected and every known state accepted, aresyncname in any spelling rejected and the tracked event never re-sync tagged, the stdio path still accepting any state,415/403/421with accepted counterparts (charset, loopback host names), a check that none of them reads the body, a check over a raw connection that a rejection whose declared body is never sent is answered within a second and not after the 10 s read timeout (a rejection closes the connection),413for a whitespace-padded body over the cap, server header/read/idle timeouts set, the server stopping on context cancellation,Manager.Runreturning the listen error for a taken port, non-loopback addresses refused (newLocalIngressandEnableLocalIngress), an end-to-end HTTP test over a real loopback listener (200 / 400 / 405, re-sync tracking, upward emission),500when the broker is down, flag-vs-file resolution, the removal memory (blocks a stale stdio replay both with the workload tracked and on a freshly started worker, is cleared by a new ingress frame, expires, covers only its own(origin, id), is not created by a stdio removal, and does not block a re-sync-tagged frame), and one test per ordering property:originatedfrombeside or instead of the canonical key,originatedfRom,OriginatedFrom, a second top-levelworkloadinfo, a remove frame'soriginatedfrom,ID, a long-sstate, two unknown names that fold together,Resync/RESYNC/ long-sresyncon lifecycle and remove frames) is an error fromingestLocalthat leaves nothing at the broker, in the re-sync set, in the peer queue or in the removal memory, and400over HTTP; a key-folding test checksfoldKeyagainst whatencoding/jsondecodes; the canonical harness-shaped frame, with a field of the producer's own, is200and the extra field is passed through; the decoded-value check is tested on frames that skipped the key check; ids of 256 bytes are accepted and 257 are400for upserts and removals; the removal memory is capped, evicts the oldest first, refreshes a repeated removal at the cap, keeps one entry for a repeated removal below it (and that entry keeps blocking a stale replay after the first removal's own expiry has passed), holds a tombstone for the terminal-retention window and prunes expired entries; the server has no write timeout.ingestMu, nobroadcastMuincommitLocal, commit before the broker write). The same was done for each new check: with theOrigin,Host, content-type, foreign-origin, unknown-state, body-cap, removal-memory check, removal-memory record and removal-memory clear disabled one at a time, the matching test fails; the same for the field-name spelling check, theresyncrejection, the decoded-value validation, the id length check, the removal-memory cap, the replacement of a repeated removal's entry, the tombstone lifetime and the connection close on a rejection.TestLocalIngressStopsWhenContextCancellednow checks that a request after shutdown gets no answer instead of that a bare TCP connect is refused: the connect form flaked on Linux (16 of 300 runs on the unchanged earlier revision under-race) because on some hosts a connect to a port nobody listens on can complete. The request form failed 0 of 300 and asserts the stronger property, that the ingress no longer answers.GOOS=linux go vet ./...andnode scripts/spdx-headers.mjsare clean inservices/nvpair-workload-managerand at the repository root.services/tests:GOOS=linux go build ./...is clean.GOOS=linux go vet ./...reports two existing "unreachable code" findings inmain_test.go, a file this change does not touch. The cross-process suite was run on Windows and Linux and cannot pass on the machine used here: a running installed instance already holds the fixed inter-node ports (listen on :14320/:14321: address already in use), soTestWorkloadManagerInboundRelay,TestWorkloadManagerOutboundBroadcastandTestBrokerRestoresClusterIdentityAfterRestartfail on that collision before reaching this code. It runs in CI.npm run service-contracts:checkwas not run locally (desktop dependencies not installed); the change adds no notification and no method.--local-ingress 127.0.0.1:14325, before the trust checks were added: a started event, acancelledevent withseq2 and aworkloads:removeeach returned200and reached the broker side asworkloads:upsert/workloads:removewith the node UUID stamped; an emptyworkloadInforeturned400;GETreturned405;0.0.0.0:14326was refused at startup; a second instance on a taken ingress port exited with a bind error..debandnvpair-tui), with the patched worker dropped intocli-bin/: aworkload:started+workload:completedpair posted to127.0.0.1:14324appeared asworkloads:upsertin the local broker log, persisted toworkloads-history.json, and arrived on a peer node's broker log and history through the existing broadcast;workloads:removeretired it on both; a connection to the node's LAN address on that port was refused. That run predates the rebase ontodevelopand the emit-first ordering.Risk
Originheader, a non-loopbackHostor a non-JSON content type before reading the body; a frame naming another node as its origin (in any spelling of the key), an unknown state, an overlong id or aresyncname is rejected, so a producer can report only workloads that run here. The ingress does not strip fields it does not know: any other field insideworkloadInfoorparamsis passed through to peers and kept in the re-sync set, so producers must send metadata only (documented), and prompts, messages and response bodies must never be added.200does not promise peer delivery. The outbound queue drops the newest frame with a warning when full, and a dropped removal is not re-synced (already documented for the stdio path).ingestMu, so it delays other ingress posts only; the stdio read loop and the heartbeat are not blocked by it. A broker pipe that stops reading blocks ingress requests, as it blocks every other path that writes to the broker.(originatedFrom, id)for one minute, so an ingress id that collides with a built-in proxy's id hides that proxy workload from peers for that minute. The docs ask producers to pick ids that cannot collide. The memory is capped at 4096 entries; past that the oldest removal stops suppressing a replay.Checklist
git commit -s), certifying the Developer Certificate of Origin.services/versions.jsonis written by automation — do not edit it by hand.🤖 Generated with Claude Code