Skip to content

workload-manager: opt-in loopback ingress for third-party workload producers - #90

Draft
dmmdea wants to merge 1 commit into
NVIDIA:developfrom
dmmdea:feat/workload-local-ingress
Draft

dmmdea wants to merge 1 commit into
NVIDIA:developfrom
dmmdea:feat/workload-local-ingress

Conversation

@dmmdea

@dmmdea dmmdea commented Sep 17, 2026 •

Copy link
Copy Markdown

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.json containing {"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/events takes the same JSON-RPC 2.0 frames as the peer port (workload:submitted|started|completed|errored with params.workloadInfo, workloads:remove with params.workloadId) through the same parseLifecycle / parseRemove validation and a 1 MiB body limit (413 above it). originatedFrom is 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 is 400. Producer mistakes are 400, other methods 405, and a broker that cannot be written is 500. Frames may carry seq, which is part of the key peers deduplicate on, and cancelled is a valid state.
  • Trust checks. Loopback does not keep a browser out: a web page can POST to 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 an Origin header (403), when its Host is not localhost, a 127.0.0.0/8 address or ::1 with the bound port (421), or when its Content-Type is not application/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 lifecycle state must be one of the five known states (400 otherwise), 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/json matches 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-s state), or two names of one object that fold together, is 400; otherwise a case variant would be decoded as the field while slipping past the exact-key origin check. A resync name in any spelling is 400 too (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.
  • Removal memory. The broker replays its store to a restarted worker, and that replay can land after a producer's 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.
  • An accepted frame is treated as local origin. It is first emitted up to the broker as 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. 200 means 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) that main calls after NewManager, so NewManager and its callers and tests are unchanged.

No desktop change is needed: the Go side treats engine as an opaque string, WorkloadItemCard renders an unknown engine without an icon, and formatModelDisplayName has a generic branch. A follow-up could add a generic engine icon and surface requesterId on the card.

Producer example:

curl -s -X POST http://127.0.0.1:14324/v1/workloads/events -H 'Content-Type: application/json' -d '{
  "jsonrpc": "2.0", "method": "workload:started",
  "params": {"workloadInfo": {"id": "job-42", "runId": "job-42", "seq": 1, "model": "qwen3.5-9b", "engine": "llamacpp",
             "state": "running", "createdAt": 1789661457000, "startedAt": 1789661457000,
             "completedAt": null, "error": null, "requesterId": "my-scheduler"}}}'

Release intent

Changelog title

Third-party local workload producers can report their jobs

Changelog body

  • The workload manager can now accept workload events from other programs on the same machine, so work that does not go through PAIR's own proxies (an external scheduler, a local inference harness) appears in every cluster member's Jobs list and counts toward the scheduler's pending work for the node that runs it.
  • The feature is off by default. It is enabled with the --local-ingress flag or an <appdir>/workload-ingress.json file, listens on loopback only, and refuses any other bind address.
  • Requests from web pages are refused, and a producer can only report workloads that run on this machine.

Bumps

  • services: minor
  • nvpair-cluster-manager: none
  • nvpair-engine-manager: none
  • nvpair-errors: none
  • nvpair-job-scheduler: none
  • nvpair-manual-nodes: none
  • nvpair-node-info: none
  • nvpair-node-scanner: none
  • nvpair-node-settings: none
  • nvpair-proxy: none
  • nvpair-tui: none
  • nvpair-ui-broker: none
  • nvpair-workload-manager: minor

An additive HTTP feature is a MINOR bump per services/VERSIONING.md; services follows the component. services/versions.json is not touched.

Scope

In:

  • The loopback ingress (localingress.go), its wiring in main.go and Manager.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 by activeMu and only touched under broadcastMu).
  • A small refactor of the local-frame path in manager.go: handleLocalLifecycle / handleLocalRemove and the ingress now share parseLocalFrame (validation, no side effects) and commitLocal (track and enqueue under broadcastMu). The stdio handlers behave as before.
  • README.md (flag row and ingress section) and the normative spec.md (interface §7.4, requirements, failure mode).

Out:

  • Any change to the inter-node port, its mTLS gate or its dedup.
  • Desktop changes (see the note above).
  • Authentication on the ingress. It is plaintext by design and stays on loopback.

Concurrency, since the ingress adds a concurrent HTTP entry point next to the single-goroutine stdio read loop:

  • Lock order is ingestMu -> broadcastMu -> activeMu. Only the ingress takes ingestMu.
  • Ingest order is validate, write up to the broker, then commitLocal (track and enqueue under broadcastMu). 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.
  • commitLocal holds broadcastMu across track and enqueue, exactly like the stdio handlers, so a re-sync snapshot cannot enqueue a workload after its removal was enqueued.
  • ingestMu keeps concurrent posts for one workload in the same order at the broker and in the peer queue. It is deliberately not broadcastMu: the broker write can block, and the stdio read loop and the re-sync heartbeat take broadcastMu.
  • No lock is held across a peer network write; the ordered broadcastLoop owns those.

Validation

In services/nvpair-workload-manager:

  • go vet ./... && go test ./... -count=1 on Windows: pass.
  • go test ./... -race -count=5 on Linux (Go 1.27.1, amd64): pass.
  • New tests (localingress_test.go): origin stamping for lifecycle and removal frames, the node's own origin accepted and a foreign one rejected (at ingestLocal and as 400 over HTTP), missing workloadInfo rejected, unknown states rejected and every known state accepted, a resync name in any spelling rejected and the tracked event never re-sync tagged, the stdio path still accepting any state, 415 / 403 / 421 with 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), 413 for a whitespace-padded body over the cap, server header/read/idle timeouts set, the server stopping on context cancellation, Manager.Run returning the listen error for a taken port, non-loopback addresses refused (newLocalIngress and EnableLocalIngress), an end-to-end HTTP test over a real loopback listener (200 / 400 / 405, re-sync tracking, upward emission), 500 when 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:
    • a failed broker write leaves no tracked state and enqueues nothing;
    • a failed removal keeps the workload tracked and enqueues nothing;
    • 64 concurrent posts for one workload reach the broker and the peer queue in the same order;
    • a removal ingested while a re-sync snapshot is in flight is never followed by that snapshot's upsert.
  • Hardening tests: every case-variant bypass found in review (originatedfrom beside or instead of the canonical key, originatedfRom, OriginatedFrom, a second top-level workloadinfo, a remove frame's originatedfrom, ID, a long-s state, two unknown names that fold together, Resync / RESYNC / long-s resync on lifecycle and remove frames) is an error from ingestLocal that leaves nothing at the broker, in the re-sync set, in the peer queue or in the removal memory, and 400 over HTTP; a key-folding test checks foldKey against what encoding/json decodes; the canonical harness-shaped frame, with a field of the producer's own, is 200 and 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 are 400 for 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.
  • Each ordering test was checked to fail with its guard removed (no ingestMu, no broadcastMu in commitLocal, commit before the broker write). The same was done for each new check: with the Origin, 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, the resync rejection, 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.
  • TestLocalIngressStopsWhenContextCancelled now 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 ./... and node scripts/spdx-headers.mjs are clean in services/nvpair-workload-manager and at the repository root.
  • Before the trust checks were added, in services/tests: GOOS=linux go build ./... is clean. GOOS=linux go vet ./... reports two existing "unreachable code" findings in main_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), so TestWorkloadManagerInboundRelay, TestWorkloadManagerOutboundBroadcast and TestBrokerRestoresClusterIdentityAfterRestart fail on that collision before reaching this code. It runs in CI. npm run service-contracts:check was not run locally (desktop dependencies not installed); the change adds no notification and no method.
  • Manual, built binary with --local-ingress 127.0.0.1:14325, before the trust checks were added: a started event, a cancelled event with seq 2 and a workloads:remove each returned 200 and reached the broker side as workloads:upsert / workloads:remove with the node UUID stamped; an empty workloadInfo returned 400; GET returned 405; 0.0.0.0:14326 was refused at startup; a second instance on a taken ingress port exited with a bind error.
  • An earlier revision of this change was also run live on a four-member cluster (two Windows nodes running the v0.1.1 desktop, one Linux node with the .deb and nvpair-tui), with the patched worker dropped into cli-bin/: a workload:started + workload:completed pair posted to 127.0.0.1:14324 appeared as workloads:upsert in the local broker log, persisted to workloads-history.json, and arrived on a peer node's broker log and history through the existing broadcast; workloads:remove retired it on both; a connection to the node's LAN address on that port was refused. That run predates the rebase onto develop and the emit-first ordering.

Risk

  • Security: the ingress is plaintext and unauthenticated. It is the same boundary the proxies' plaintext loopback personality documents in SECURITY.md ("a process that can reach a local endpoint may be able to submit inference or observe behavior"): a local process that can already submit work through a proxy may now report work it ran elsewhere, a strictly smaller privilege. It never binds beyond loopback, is off by default, and the peer port's cluster-mTLS gate is untouched. Because a browser can also reach loopback, the ingress refuses requests with an Origin header, a non-loopback Host or 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 a resync name 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 inside workloadInfo or params is 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.
  • Compatibility: additive. With the ingress off, the listener does not exist and the local-frame path behaves as before. No JSON-RPC method or payload changes.
  • Delivery: 200 does 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).
  • Load: a stalled broker write holds 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.
  • Ordering: the broker and the peers see ingress frames for one workload in the same order. That holds among ingress frames; a Jobs-list removal that the broker itself sends for the same workload while an ingress update for it is being written can reach peers before that update, and the producer's next frame or removal converges them. Documented in the README and the spec.
  • Removal memory: a workload removed through the ingress also ignores broker lifecycle frames for the same (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.
  • Data and packaging: no data migration, no new dependency, no new binary.

Checklist

  • I have read the Contributing Guidelines.
  • Every commit is signed off (git commit -s), certifying the Developer Certificate of Origin.
  • New or existing tests cover the change.
  • Relevant documentation is updated.
  • I checked the diff, changed filenames, and commit messages for credentials, private data, internal URLs, internal issue identifiers, and generated artifacts.
  • I recorded the validation commands and results above.
  • I declared version bumps in the release-intent block above. services/versions.json is written by automation — do not edit it by hand.

🤖 Generated with Claude Code

…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
dmmdea force-pushed the feat/workload-local-ingress branch from cd31a1c to 78bc20b Compare September 29, 2026 22:29
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