diff --git a/docs/CLI.md b/docs/CLI.md index 298be865..fd37de9e 100644 --- a/docs/CLI.md +++ b/docs/CLI.md @@ -153,7 +153,8 @@ granted. - `--allow-download` — permit downloading a missing runtime / system image (multi-GB; never implicit). Without it, a missing runtime is exit 12. iOS runtimes remain Xcode-managed in v1: `--allow-download` cannot install - them; install the runtime through Xcode first. + them; install the runtime through Xcode first. Through a gateway the flag + has no effect. - `--ttl ` — the lease's initial TTL, replacing `lease.defaultTtlMs` (15m) for this lease. Asking for more than `lease.maxTtlMs` (4h) is a `BAD_REQUEST` (exit 2), not a silent clamp. See @@ -567,12 +568,19 @@ differs. **Leasing is identical.** `simlock lease` takes the same flags and prints the same grant line, including `--ttl`, `--no-wait`, `--timeout`, and -`--allow-download` (forwarded to the chosen worker, which clamps it through -its own `downloads.policy`). The request waits in the gateway's own -fleet-wide FIFO queue, reporting `queued` with a `queuePosition` exactly as a -worker's queue does, and is dispatched to the worker best placed to serve it -— a machine with a matching warm device first, otherwise the one with the -most free capacity. You do not name a machine and there is no flag to; where +`--allow-download`, which is accepted and has no effect through a gateway: +only runtimes already installed on a worker count, and no download is +started. The request waits in the gateway's own fleet-wide FIFO queue, +reporting `queued` with a `queuePosition` exactly as a worker's queue does. + +The gateway sends a request only to a worker that can serve it. `--device` +matches a worker's model in any letter case and by any other name that +worker's catalog lists for it, such as an Android AVD id (`pixel_7`). When +`--os` is given, the worker must pair that runtime with the model; without +it, the model must pair with at least one installed runtime. A worker that +has the model and the runtime but cannot pair them is passed over. Among the +workers that can serve it, the request goes to a machine with a matching warm +device first, otherwise the one with the most free capacity. You do not name a machine and there is no flag to; where a device lives is the gateway's decision. The grant carries one additional block so you can see where it landed: @@ -628,9 +636,8 @@ runtime annotated with the workers that have it — so a `--device` the catalog lists is leasable *somewhere*, not necessarily everywhere. A model is paired with a runtime when at least one connected worker pairs them itself; one worker having the model and another having the runtime does not make a -pair. `simlock worker list --json` shows each worker's own pairings. The -gateway does not yet pick a worker by its pairings, so a listed pair can -still be sent to a worker that cannot pair them. +pair. `simlock worker list --json` shows each worker's own pairings. A +request goes only to a worker that pairs the model with the runtime. **`simlock events`** shows the fleet: every worker's business events are republished on the gateway's bus with `workerId` added to the payload, @@ -803,9 +810,9 @@ why its Android catalog looks thin, trimmed to one worker: ``` `downloads.policy` and `lease.maxTtlMs` are that worker's own effective -config, read when its uplink connects and again on every periodic refresh — -routing needs the policy to know whether a machine may install a missing -runtime before sending it a request that needs one. `catalog` is what that +config, read when its uplink connects and again on every periodic refresh. +The policy is shown for reference; routing does not read it, since no +download is started through a gateway. `catalog` is what that worker can lease, each model with the runtimes it pairs with. `host` is the machine: operating system, its version, CPU architecture, and the version of each platform tool its drivers use (`xcode` with its build; the Android diff --git a/docs/CLIENT.md b/docs/CLIENT.md index 780868af..c7275810 100644 --- a/docs/CLIENT.md +++ b/docs/CLIENT.md @@ -217,11 +217,12 @@ const { platforms } = await client.getCatalog({ platform: "ios" }); Against a gateway the catalog is the union of the connected workers'. A model is paired with a runtime when at least one worker pairs them, and `modelWorkers` and `runtimeWorkers` say which workers have each model and -runtime. The gateway does not yet use the pairings to pick a worker, so a -pair it lists can still go to a worker that has the model and the runtime -but cannot pair them, and that request fails there. `modelAliases` and -`images` are the unions of each worker's own. The gateway does not yet route -by another name, so ask it for a model by its name in `models`. +runtime. The gateway sends a request only to a worker that pairs the model +with the runtime. `modelAliases` and `images` are the unions of each +worker's own. A model may be asked for by any name a worker lists for it, in +any letter case, and the gateway sends that worker its own name for it. +`allowDownload` has no effect through a gateway: only installed runtimes +count. ## What machine answered: `getStatus().host` diff --git a/docs/CONFIGURATION.md b/docs/CONFIGURATION.md index fb4f8d1b..3b1e73af 100644 --- a/docs/CONFIGURATION.md +++ b/docs/CONFIGURATION.md @@ -25,7 +25,7 @@ a warning. Inspect the effective, merged configuration at any time with | `gateway.token` | **Worker side.** The join token (`simlock token create --role worker`, minted on the gateway) this worker presents when it opens its uplink. Required whenever `gateway.url` is set. | unset | | `gateway.label` | **Worker side.** Display name for this worker in `simlock worker list`, `status`, the console, and on the lease's `worker` block. Display-only: nothing routes on it and it need not be unique. | the worker's own id | | `exec.timeoutMs` | **Worker side.** How long one `device.exec` command (`simlock simctl` / `simlock adb` against a gateway or over HTTP) may run before the worker kills it and the operation fails with `EXEC_TIMEOUT`. Authoritative: it bounds the process that actually runs. | `10 minutes` | -| `gateway.routing` | **Gateway side.** Which routing policy places a queued request on a worker. `warm-then-free` is the only policy in v1: warm hit first, then the most free running capacity for the platform. | `warm-then-free` | +| `gateway.routing` | **Gateway side.** Which routing policy places a queued request on a worker. `warm-then-free` is the only policy in v1: among workers that can serve the request, warm hit first, then the most free running capacity for the platform. | `warm-then-free` | | `gateway.disconnectedRetentionMs` | **Gateway side.** How long a disconnected worker is kept (greyed, never dispatched to) before the gateway forgets it. The clock is held while the gateway still knows of gateway-issued leases on that worker, and that hold ends when the last of those leases passes its deadline. | `24 hours` | | `gateway.execTimeoutMs` | **Gateway side.** How long the gateway waits on a proxied `device.exec` before giving up. A backstop for a worker that never answers at all — deliberately longer than the worker's own `exec.timeoutMs`, which is authoritative because that side owns the process and can kill it, so an ordinary timeout surfaces as the worker's `EXEC_TIMEOUT` rather than racing this one. | `11 minutes` | | `gateway.leaseRequestTimeoutMs` | **Gateway side.** How long the gateway waits on a forwarded `lease.request` before giving up on that worker for this request, answering `WORKER_UNREACHABLE`. Bounds the one uplink call that otherwise had no timeout of its own, so a wedged worker cannot park a request where neither a deadline nor `lease.cancel` could ever reach it again. Generous against a cold device provision-plus-boot; well below `gateway.execTimeoutMs`, since granting a lease should never take as long as a command run against the device afterward. | `5 minutes` | @@ -207,6 +207,14 @@ does not apply to it. It reads: | `lease.*` | `defaultTtlMs`/`maxTtlMs` bound what its own clients may ask for, before a request is dispatched — see below | | `log.*`, `eventBuffer.*`, `eventLog.*` | logging and the event history, as anywhere | +**A worker's `downloads.policy` does not apply to requests through a +gateway.** The gateway sends a request only to a worker whose catalog already +has what it asks for: the model under any name the worker lists for it, in +any letter case, paired with the requested runtime (or, with none requested, +with at least one installed runtime). It never asks a worker to download, so +`--allow-download` has no effect through a gateway, whatever each worker's +policy says. + **Both ends have a `lease.*` block, and on a fleet lease the gateway's is the one that decides the width.** A request arriving at a gateway with no `ttlMs` is filled in with the *gateway's* `lease.defaultTtlMs` before it is dispatched diff --git a/docs/HTTP-API.md b/docs/HTTP-API.md index 5940790a..27a694e7 100644 --- a/docs/HTTP-API.md +++ b/docs/HTTP-API.md @@ -187,9 +187,8 @@ Each platform entry carries `modelRuntimes`: for every name in `models`, the installed runtimes that model pairs with. A pair listed there can be leased; a model and a runtime that are each listed but not paired cannot. An empty list means nothing installed pairs with that model. On a gateway a model is -paired with a runtime when at least one connected worker pairs them; the -gateway does not yet pick a worker by its pairings, so such a request can -still go to a worker that cannot pair them and fail there. +paired with a runtime when at least one connected worker pairs them, and the +gateway sends a request only to a worker that pairs them. Each entry also carries `modelAliases`: for a name in `models`, the other names a lease request may use for it, in any letter case. Only models that @@ -199,8 +198,9 @@ system image with its API level (`runtime`, a value from `runtimes`), `tag`, and `abi`, including an image whose ABI the host cannot run natively; an iOS entry has no `images`. On a gateway `modelAliases` is the union per model and `images` the union of each worker's images, absent when no worker's entry -for that platform has an `images` field. The gateway does not yet route by another name: send it the name from -`models`. +for that platform has an `images` field. A gateway accepts any name a worker +lists for a model, in any letter case, and sends that worker its own name for +it. ```json { "platforms": [ { @@ -277,6 +277,14 @@ try again, use a new key. Repeating works across a daemon restart, for belong to your token: another token sending the same key starts a request of its own. +Through a **gateway**, `allowDownload` does not change which worker is picked +and never starts a download: only runtimes already installed on a worker +count. It still makes the `POST` answer early, as described below. The +gateway sends the request only to a worker whose catalog can serve it: one +that lists `device` as a model or another name for one, in any letter case, +and pairs that model with `os` (or, without `os`, with at least one installed +runtime). The worker is sent its own name for the model. + With `allowDownload: true` the `201` is returned as soon as the request is stored — resolving a downloadable runtime can take minutes, so progress and any later failure surface on the request resource instead of on the `POST` @@ -704,10 +712,9 @@ simulated (hence the thin Android catalog), trimmed to one worker: ``` `downloads.policy` is that worker's own effective policy, read when its -uplink connects and again on every periodic refresh. Routing needs it to know -whether a worker may install a missing runtime at all before sending it a -request that depends on one; it is never an override, since the worker clamps -`allowDownload` through the same policy regardless. +uplink connects and again on every periodic refresh. It is shown for +reference; routing does not read it, since no download is started through a +gateway. `catalog` is what that worker can lease, each model with the runtimes it pairs with. `host` is the worker's machine, the same block its own diff --git a/docs/internal/ARCHITECTURE.md b/docs/internal/ARCHITECTURE.md index 24517beb..00a6e49b 100644 --- a/docs/internal/ARCHITECTURE.md +++ b/docs/internal/ARCHITECTURE.md @@ -479,11 +479,8 @@ carries none. A gateway's own `status.get` reports the gateway's machine with no tools. From `config.get` the gateway keeps two fields. The effective -`downloads.policy`, because routing has to know whether a worker is even -*allowed* to install a missing runtime before it sends that worker a request -which depends on one. It is a routing input and never an override: the -worker still clamps `allowDownload` through its own policy, whatever the view -said. And `lease.maxTtlMs`, compared against the gateway's own to warn when a +`downloads.policy`, for display only: routing counts installed runtimes and +never reads it (ADR 0009 §3). And `lease.maxTtlMs`, compared against the gateway's own to warn when a worker's cap is lower. Config is daemon input, read at start, so these change only across a worker restart; re-reading them on the tick costs one call. `config.get` is an admin operation, which the uplink session is. @@ -578,17 +575,30 @@ that removed a worker, or the last that decided when none removed any — is reported on `request.dispatched`. `gateway.routing` names a whole list; the lists are code, and no config key lists or orders stages. -The v1 policy (`warm-then-free`) is three stages: - -1. `eligible` (filter): drop workers that are disconnected, drained, - incompatible, or lacking the requested platform, model, or runtime — a - download counts as available only on a worker whose own `downloads.policy` - would allow it; -2. `warm-hit` (rank, settles): prefer a worker with an unleased `ready` device - matching the request — a **warm hit**, and a sub-second grant; -3. `free-capacity` (rank): otherwise the worker with the **most free running +The v1 policy (`warm-then-free`) is four stages: + +1. `takes-requests` (filter): drop workers that are disconnected, + incompatible, drained, or whose capacity has not been read; +2. `can-serve` (filter): drop workers whose catalog cannot serve the request + (ADR 0009 §3, `routing/request-match.ts`). The model is the first entry of + the worker's `models` whose name or `modelAliases` entry equals the + requested name, ignoring letter case. A named runtime must be in that + model's `modelRuntimes`; with none named the list must be non-empty. Only + installed runtimes count, so a download never makes a worker able to + serve; +3. `warm-hit` (rank, settles): prefer a worker with an unleased `ready` device + matching the request, compared against the worker's own name for the + model — a **warm hit**, and a sub-second grant; +4. `free-capacity` (rank): otherwise the worker with the **most free running capacity** for that platform. +The same matcher gives the name the gateway forwards: the worker is sent its +own name for the model, so it resolves exactly what routing matched, and +`allowDownload` is always forwarded as `false`. `lease.requested` keeps the +name the client sent. The `eligible` stage that predates this split stays in +the code only so the conformance tests can run the three stages that +reproduce the policy before ADR 0009. + There is no other placement rule in v1: no requester affinity, no label selectors, no per-worker platform exclusions. Each of those is a future routing policy behind `gateway.routing`, not a change to the request shape — diff --git a/e2e/gateway-fleet.test.ts b/e2e/gateway-fleet.test.ts index 6f559fd2..c510913e 100644 --- a/e2e/gateway-fleet.test.ts +++ b/e2e/gateway-fleet.test.ts @@ -451,6 +451,54 @@ describe("gateway fleet", () => { }); }); + it("grants --device 'iphone 16 pro' and --device pixel_7 through a gateway by the worker's own names", async () => { + const port = await freeLoopbackPort(); + const gateway = await withDaemon({ + configOverrides: { http: { host: "127.0.0.1", port }, mode: "gateway" }, + driver: "none", + }); + const minted = await gateway.cli(["token", "create", "--role", "worker"]); + const { secret } = minted.json as { secret: string }; + // The fake driver resolves model names exactly, so a grant proves the gateway sent the + // worker its own name for the model. + await withDaemon({ + configOverrides: { + gateway: { label: "worker", token: secret, url: `ws://127.0.0.1:${port}` }, + }, + driverScript: { + android: { + availableOsVersions: ["35"], + knownModels: ["Pixel 7"], + modelAliases: { "Pixel 7": ["pixel_7"] }, + }, + ios: { availableOsVersions: ["26.0"], knownModels: ["iPhone 16 Pro"] }, + }, + }); + await waitForWorkers( + gateway, + (views) => + views.length === 1 && views[0]?.connection === "connected" && views[0].catalog.length === 2, + "the worker connected with both catalogs", + ); + + const lease = async (platform: string, device: string) => { + const leased = await gateway.cli( + ["lease", "--platform", platform, "--device", device, "--detach", "--no-wait"], + { timeout: 30_000 }, + ); + expect(leased.code, leased.stderr).toBe(0); + const grant = leased.json as { + readonly device: { readonly spec: { readonly model: string } }; + readonly lease: { readonly id: string }; + }; + expect((await gateway.cli(["release", grant.lease.id], { timeout: 30_000 })).code).toBe(0); + return grant.device.spec.model; + }; + + await expect(lease("ios", "iphone 16 pro")).resolves.toBe("iPhone 16 Pro"); + await expect(lease("android", "pixel_7")).resolves.toBe("Pixel 7"); + }); + it("refuses an uplink whose token is not a worker join token", async () => { const port = await freeLoopbackPort(); const gateway = await withDaemon({ diff --git a/src/contract/schemas.ts b/src/contract/schemas.ts index ee9ff07e..358a4c4a 100644 --- a/src/contract/schemas.ts +++ b/src/contract/schemas.ts @@ -725,11 +725,10 @@ export const workerViewSchema = z.object({ capacity: statusCapacitySchema.optional(), /** * The worker's effective `downloads.policy`, read once with `config.get` when the uplink - * connects. It is on the view because it is a *routing input*, not decoration: ADR 0005 §13 - * says a request that would need a download is only eligible on a worker whose policy allows - * one, and #118's policy reads it from here rather than asking at dispatch time. Absent for - * a worker whose `config.get` the gateway could not read (an incompatible one, or a call - * that failed). + * connects. Display only: routing counts installed runtimes and never reads it, and the + * gateway forwards every request with `allowDownload: false` (ADR 0009 §3). Absent for a + * worker whose `config.get` the gateway could not read (an incompatible one, or a call that + * failed). */ downloads: z.object({ policy: z.enum(["never", "on-request", "always"]) }).optional(), /** diff --git a/src/gateway/dispatcher.ts b/src/gateway/dispatcher.ts index e5836366..0fc57339 100644 --- a/src/gateway/dispatcher.ts +++ b/src/gateway/dispatcher.ts @@ -370,7 +370,8 @@ export class GatewayDispatcher { ...(input.mode === undefined ? {} : { mode: input.mode }), }, { - allowDownload: input.allowDownload ?? false, + // `input.allowDownload` stays accepted and is not passed on: only installed runtimes + // count through a gateway (ADR 0009 §3). noWait: input.noWait ?? false, ownerId: input.owner ?? session.principal, requesterId: input.requesterId ?? session.principal, diff --git a/src/gateway/fleet-coordinator.test.ts b/src/gateway/fleet-coordinator.test.ts index 10342778..6ff94d5b 100644 --- a/src/gateway/fleet-coordinator.test.ts +++ b/src/gateway/fleet-coordinator.test.ts @@ -2033,3 +2033,159 @@ describe("drain and unreachable lifecycle (ADR §9/§28/§29, #119)", () => { expect(client.calls.filter((call) => call.startsWith("lease.request"))).toHaveLength(1); }); }); + +describe("FleetLeaseCoordinator sends a worker only requests its catalog can serve (ADR 0009 §3)", () => { + /** Connects a worker with one catalog entry stated in full, so pairings and aliases are explicit. */ + function connectWithCatalog( + workers: WorkerRegistry, + workerId: string, + entry: Parameters[0][number], + downloads: "never" | "on-request" | "always" = "on-request", + ): void { + workers.connected(workerId, undefined, "0.3.0"); + workers.refresh(workerId, { + capacity: statusFixture().capacity, + catalog: catalogFixture([entry]).platforms, + devices: [], + downloads: { policy: downloads }, + health: "running", + leases: [], + queueDepth: 0, + }); + } + + function leaseRequests(client: ScriptedWorkerClient): string[] { + return client.calls.filter((call) => call.startsWith("lease.request")); + } + + it("dispatches a request in another letter case with the worker's own name for the model", async () => { + const { coordinator, directory, eventBus, workers } = harness(); + const client = new ScriptedWorkerClient(); + directory.add("wrk_a", client); + client.requestLeaseQueue.push({ grant: grantFixture(), kind: "grant" }); + connectWithCatalog(workers, "wrk_a", { + models: ["iPhone 16 Pro"], + platform: "ios", + runtimes: ["26.0"], + }); + + await coordinator.request( + { model: "iphone 16 pro", platform: "ios" }, + requestOptions({ noWait: true }), + ); + + expect(client.lastRequestLeaseInput?.model).toBe("iPhone 16 Pro"); + const requested = eventBus.replay().find((event) => event.event === "lease.requested"); + expect(requested?.payload).toMatchObject({ requestSpec: { model: "iphone 16 pro" } }); + }); + + it("dispatches a request by an alias to the worker that lists it", async () => { + const { coordinator, directory, workers } = harness(); + const other = new ScriptedWorkerClient(); + const lister = new ScriptedWorkerClient(); + directory.add("wrk_a", other); + directory.add("wrk_b", lister); + lister.requestLeaseQueue.push({ grant: grantFixture(), kind: "grant" }); + // wrk_a sorts first and has the same free capacity, so only the alias can send it to wrk_b. + connectWithCatalog(workers, "wrk_a", { + models: ["Pixel 7"], + platform: "android", + runtimes: ["35"], + }); + connectWithCatalog(workers, "wrk_b", { + modelAliases: { "Pixel 7": ["pixel_7"] }, + models: ["Pixel 7"], + platform: "android", + runtimes: ["35"], + }); + + await coordinator.request( + { model: "pixel_7", platform: "android" }, + requestOptions({ noWait: true }), + ); + + expect(leaseRequests(other)).toEqual([]); + expect(lister.lastRequestLeaseInput?.model).toBe("Pixel 7"); + }); + + it("passes over a worker with the model and the runtime unpaired for one that pairs them", async () => { + const { coordinator, directory, workers } = harness(); + const unpaired = new ScriptedWorkerClient(); + const paired = new ScriptedWorkerClient(); + directory.add("wrk_a", unpaired); + directory.add("wrk_b", paired); + paired.requestLeaseQueue.push({ grant: grantFixture(), kind: "grant" }); + // wrk_a lists iPhone 17 and 26.0, but pairs iPhone 17 only with 18.0. + connectWithCatalog(workers, "wrk_a", { + modelRuntimes: { "iPhone 15": ["26.0"], "iPhone 17": ["18.0"] }, + models: ["iPhone 17", "iPhone 15"], + platform: "ios", + runtimes: ["18.0", "26.0"], + }); + connectWithCatalog(workers, "wrk_b", { + models: ["iPhone 17"], + platform: "ios", + runtimes: ["26.0"], + }); + + await coordinator.request(REQUEST, requestOptions({ noWait: true })); + + expect(leaseRequests(unpaired)).toEqual([]); + expect(leaseRequests(paired)).toHaveLength(1); + }); + + it("with no runtime named, needs a model with at least one paired runtime", async () => { + const { coordinator, directory, workers } = harness(); + const unpaired = new ScriptedWorkerClient(); + const paired = new ScriptedWorkerClient(); + directory.add("wrk_a", unpaired); + directory.add("wrk_b", paired); + paired.requestLeaseQueue.push({ grant: grantFixture(), kind: "grant" }); + connectWithCatalog(workers, "wrk_a", { + modelRuntimes: { "iPhone 17": [] }, + models: ["iPhone 17"], + platform: "ios", + runtimes: ["26.0"], + }); + connectWithCatalog(workers, "wrk_b", { + models: ["iPhone 17"], + platform: "ios", + runtimes: ["26.0"], + }); + + await coordinator.request( + { model: "iPhone 17", platform: "ios" }, + requestOptions({ noWait: true }), + ); + + expect(leaseRequests(unpaired)).toEqual([]); + expect(leaseRequests(paired)).toHaveLength(1); + }); + + it("does not make a worker eligible for allowDownload: true, and forwards allowDownload: false", async () => { + const { coordinator, directory, workers } = harness(); + const downloader = new ScriptedWorkerClient(); + const installed = new ScriptedWorkerClient(); + directory.add("wrk_a", downloader); + directory.add("wrk_b", installed); + installed.requestLeaseQueue.push({ grant: grantFixture(), kind: "grant" }); + // wrk_a would download 26.0 under its own policy; on main that made it eligible, and it + // sorts first with the same free capacity. + connectWithCatalog( + workers, + "wrk_a", + { models: ["iPhone 17"], platform: "ios", runtimes: ["18.0"] }, + "always", + ); + connectWithCatalog(workers, "wrk_b", { + models: ["iPhone 17"], + platform: "ios", + runtimes: ["26.0"], + }); + + await coordinator.request(REQUEST, requestOptions({ allowDownload: true, noWait: true })); + + expect(leaseRequests(downloader)).toEqual([]); + expect(installed.lastRequestLeaseInput?.allowDownload).toBe(false); + }); +}); diff --git a/src/gateway/fleet-coordinator.ts b/src/gateway/fleet-coordinator.ts index 75fe1393..f98a5ed0 100644 --- a/src/gateway/fleet-coordinator.ts +++ b/src/gateway/fleet-coordinator.ts @@ -69,7 +69,12 @@ import { type FleetWaiter, type LeaseRequestOptions, } from "./queue.js"; -import type { RoutableRequest, RoutingDecision, RoutingPolicy } from "./routing.js"; +import { + matchRequest, + type RoutableRequest, + type RoutingDecision, + type RoutingPolicy, +} from "./routing.js"; import { viewLoadKey } from "./routing/view-state.js"; export interface FleetExecInput { @@ -701,11 +706,12 @@ export class FleetLeaseCoordinator { const eligibleWorkers = this.options.views .views() .filter((worker) => !claimedThisPass.has(worker.id) && !this.#refused(waiter, worker)); - const decision = this.options.routing.select(routable(waiter), eligibleWorkers); + const request = routable(waiter); + const decision = this.options.routing.select(request, eligibleWorkers); if (decision === undefined) continue; claimedThisPass.add(decision.workerId); if (waiter === candidate) candidateAttempted = true; - this.#beginAttempt(waiter, decision); + this.#beginAttempt(waiter, decision, forwardedModel(request, eligibleWorkers, decision)); } return candidateAttempted; } @@ -718,9 +724,9 @@ export class FleetLeaseCoordinator { * `#settleGrant` while `queue.resolve` quietly answers `false`, leaving an orphan lease no * client holds a reference to release. */ - #beginAttempt(waiter: FleetWaiter, decision: RoutingDecision): void { + #beginAttempt(waiter: FleetWaiter, decision: RoutingDecision, model: string): void { if (!this.#queue.markProcessing(waiter)) return; - void this.#attempt(waiter, decision); + void this.#attempt(waiter, decision, model); } /** @@ -752,7 +758,7 @@ export class FleetLeaseCoordinator { * queued waiter needs, and nothing else would ever schedule that second look. */ // fallow-ignore-next-line complexity -- one attempt, every exit of which is named in the doc comment above. - async #attempt(waiter: FleetWaiter, decision: RoutingDecision): Promise { + async #attempt(waiter: FleetWaiter, decision: RoutingDecision, model: string): Promise { const workerId = decision.workerId; const target = this.options.directory.target(workerId); const client = target?.reachable === true ? target.client() : undefined; @@ -800,7 +806,7 @@ export class FleetLeaseCoordinator { client.requestLease( { platform: waiter.request.platform, - model: waiter.request.model, + model, ...(waiter.request.osVersion === undefined ? {} : { osVersion: waiter.request.osVersion }), @@ -813,7 +819,8 @@ export class FleetLeaseCoordinator { // `ownerId` instead of deriving it from the connection -- see // `daemon/dispatcher.ts`'s `#leaseRequest`. owner: waiter.options.ownerId, - allowDownload: waiter.options.allowDownload ?? false, + // ADR 0009 §3: a download is never triggered through a gateway. + allowDownload: false, // ADR §12: worker queues never hold gateway traffic. Every dispatch is `noWait` // regardless of what the original caller asked the gateway for -- the *gateway's // own* queue is where a "wait" request actually waits. @@ -1116,11 +1123,25 @@ export class FleetLeaseCoordinator { } } +/** + * ADR 0009 §3: the worker is sent its own name for the model, so it resolves exactly what routing + * matched. The client's name is kept only for a worker the catalog does not match, which no + * registered policy picks. + */ +function forwardedModel( + request: RoutableRequest, + views: readonly WorkerView[], + decision: RoutingDecision, +): string { + const view = views.find((worker) => worker.id === decision.workerId); + return (view === undefined ? undefined : matchRequest(view, request)) ?? request.model; +} + function routable(waiter: FleetWaiter): RoutableRequest { return { platform: waiter.request.platform as Platform, model: waiter.request.model, ...(waiter.request.osVersion === undefined ? {} : { osVersion: waiter.request.osVersion }), - allowDownload: waiter.options.allowDownload ?? false, + allowDownload: false, }; } diff --git a/src/gateway/routing.test.ts b/src/gateway/routing.test.ts index a2eaf6a9..eb09b462 100644 --- a/src/gateway/routing.test.ts +++ b/src/gateway/routing.test.ts @@ -1,10 +1,13 @@ import { describe, expect, it } from "vitest"; import { - createRoutingPolicy, + composeRoutingPolicy, DEFAULT_ROUTING_POLICY, + eligible, + freeCapacity, isRoutingPolicyName, routingPolicyNames, + warmHit, } from "./routing.js"; import { catalogFixture, deviceFixture, statusFixture } from "./test-support.js"; import type { WorkerView } from "./worker-registry.js"; @@ -45,7 +48,7 @@ describe("routing registry", () => { }); describe("warm-then-free policy", () => { - const policy = createRoutingPolicy("warm-then-free"); + const policy = composeRoutingPolicy([eligible, warmHit, freeCapacity]); it("prefers a worker with a matching device already in the ready state (a warm hit)", () => { const cold = view("wrk_cold"); diff --git a/src/gateway/routing.ts b/src/gateway/routing.ts index f61bc38f..1fa36007 100644 --- a/src/gateway/routing.ts +++ b/src/gateway/routing.ts @@ -6,12 +6,16 @@ * them. Nothing outside this module decides how a worker is chosen. */ import { type RoutableRequest, type RoutingStage, runStages } from "./routing/pipeline.js"; +import { canServe } from "./routing/stages/can-serve.js"; import { eligible } from "./routing/stages/eligible.js"; import { freeCapacity } from "./routing/stages/free-capacity.js"; +import { takesRequests } from "./routing/stages/takes-requests.js"; import { warmHit } from "./routing/stages/warm-hit.js"; import type { WorkerView } from "./worker-registry.js"; export type { FilterStage, RoutableRequest, RoutingStage } from "./routing/pipeline.js"; +export { matchRequest } from "./routing/request-match.js"; +// `eligible` is exported for the conformance tests only; no registered policy lists it. export { eligible, freeCapacity, warmHit }; export type RoutingReason = "warm-hit" | "free-capacity"; @@ -51,11 +55,12 @@ export function composeRoutingPolicy(stages: readonly RoutingStage[]): RoutingPo /** * The registry: adding a policy means adding one entry here, nothing else. * - * `warm-then-free` is ADR 0005 §13's v1 policy: drop the ineligible, prefer a warm hit, otherwise - * the worker with the most free running capacity, ties broken by ascending worker id. + * `warm-then-free` is ADR 0005 §13's v1 policy: keep the workers that take requests and whose + * catalog can serve this one (ADR 0009 §2, §3), prefer a warm hit, otherwise the worker with the + * most free running capacity, ties broken by ascending worker id. */ const routingPolicies = { - "warm-then-free": [eligible, warmHit, freeCapacity], + "warm-then-free": [takesRequests, canServe, warmHit, freeCapacity], } as const satisfies Record; export type RoutingPolicyName = keyof typeof routingPolicies; diff --git a/src/gateway/routing/pipeline.ts b/src/gateway/routing/pipeline.ts index aa4f6024..352a8eea 100644 --- a/src/gateway/routing/pipeline.ts +++ b/src/gateway/routing/pipeline.ts @@ -12,11 +12,9 @@ export interface RoutableRequest { readonly model: string; readonly osVersion?: string; /** - * The request's own flag, forwarded as-is: a gateway has no `downloads.policy` of its own to - * fold it against (that config key is worker-only, ADR 0005 §2), so the clamp the `eligible` - * stage applies is only ever the *worker's* declared policy, never a gateway-side one -- the - * worker still applies its own `effectiveAllowDownload` when the forwarded `lease.request` - * lands. + * Read only by the legacy `eligible` stage, which the conformance tests run. The gateway always + * routes with `false`: only installed runtimes count, and a client's `allowDownload` has no + * effect through a gateway (ADR 0009 §3). */ readonly allowDownload: boolean; } diff --git a/src/gateway/routing/request-match.test.ts b/src/gateway/routing/request-match.test.ts new file mode 100644 index 00000000..9d658ef2 --- /dev/null +++ b/src/gateway/routing/request-match.test.ts @@ -0,0 +1,131 @@ +import { describe, expect, it } from "vitest"; + +import { catalogFixture, statusFixture } from "../test-support.js"; +import type { WorkerView } from "../worker-registry.js"; +import { matchRequest } from "./request-match.js"; + +function worker(entry: Parameters[0][number]): WorkerView { + return { + capacity: statusFixture().capacity, + catalog: catalogFixture([entry]).platforms, + connection: "connected", + devices: [], + drained: false, + id: "wrk_a", + lastSeenAt: 1, + leases: [], + }; +} + +const PIXELS = { + modelAliases: { "Pixel 7": ["pixel_7"], "Pixel 7 Legacy": ["pixel 7"] }, + platform: "android" as const, + runtimes: ["35"], +}; + +describe("matchRequest", () => { + it("matches a model name in another letter case and returns the worker's own name", () => { + const view = worker({ models: ["iPhone 16 Pro"], platform: "ios", runtimes: ["26.0"] }); + + expect( + matchRequest(view, { allowDownload: false, model: "iphone 16 PRO", platform: "ios" }), + ).toBe("iPhone 16 Pro"); + }); + + it("matches a name the catalog lists as an alias, in any letter case", () => { + const view = worker({ ...PIXELS, models: ["Pixel 7"] }); + + expect( + matchRequest(view, { allowDownload: false, model: "PIXEL_7", platform: "android" }), + ).toBe("Pixel 7"); + expect( + matchRequest(view, { allowDownload: false, model: "pixel_8", platform: "android" }), + ).toBeUndefined(); + }); + + it("when two models answer to a name, the first in models wins", () => { + // "pixel 7" is the name of one model and an alias of the other. + const request = { allowDownload: false, model: "pixel 7", platform: "android" as const }; + + expect( + matchRequest(worker({ ...PIXELS, models: ["Pixel 7", "Pixel 7 Legacy"] }), request), + ).toBe("Pixel 7"); + expect( + matchRequest(worker({ ...PIXELS, models: ["Pixel 7 Legacy", "Pixel 7"] }), request), + ).toBe("Pixel 7 Legacy"); + }); + + it("rejects a named runtime the catalog lists but does not pair with the model", () => { + const view = worker({ + modelRuntimes: { "iPhone 17": ["18.0"] }, + models: ["iPhone 17"], + platform: "ios", + runtimes: ["18.0", "26.0"], + }); + const request = { allowDownload: false, model: "iPhone 17", platform: "ios" as const }; + + expect(matchRequest(view, { ...request, osVersion: "26.0" })).toBeUndefined(); + expect(matchRequest(view, { ...request, osVersion: "18.0" })).toBe("iPhone 17"); + }); + + it("with no runtime named, needs the model to pair with at least one", () => { + const request = { allowDownload: false, model: "iPhone 17", platform: "ios" as const }; + const unpaired = worker({ + modelRuntimes: { "iPhone 17": [] }, + models: ["iPhone 17"], + platform: "ios", + runtimes: ["26.0"], + }); + const paired = worker({ models: ["iPhone 17"], platform: "ios", runtimes: ["26.0"] }); + + expect(matchRequest(unpaired, request)).toBeUndefined(); + expect(matchRequest(paired, request)).toBe("iPhone 17"); + }); + + it("does not match on another platform's catalog", () => { + const view = worker({ models: ["iPhone 17"], platform: "ios", runtimes: ["26.0"] }); + + expect( + matchRequest(view, { allowDownload: false, model: "iPhone 17", platform: "android" }), + ).toBeUndefined(); + }); + + it("reads no inherited key as a pairing or an alias", () => { + const view = worker({ + modelAliases: {}, + modelRuntimes: {}, + models: ["constructor"], + platform: "ios", + runtimes: ["26.0"], + }); + + // No pairing: `modelRuntimes.constructor` is inherited, not a list of runtimes. + expect( + matchRequest(view, { allowDownload: false, model: "constructor", platform: "ios" }), + ).toBeUndefined(); + // No alias: `modelAliases.constructor` is inherited, so asking for another name reads no + // alias list for the model and moves on to the next one. + const withNext = worker({ + modelAliases: {}, + models: ["constructor", "iPhone 17"], + platform: "ios", + runtimes: ["26.0"], + }); + expect( + matchRequest(withNext, { allowDownload: false, model: "iPhone 17", platform: "ios" }), + ).toBe("iPhone 17"); + }); + + it("matches no download: allowDownload does not make an unlisted runtime pair", () => { + const view = worker({ models: ["iPhone 17"], platform: "ios", runtimes: ["26.0"] }); + + expect( + matchRequest(view, { + allowDownload: true, + model: "iPhone 17", + osVersion: "18.0", + platform: "ios", + }), + ).toBeUndefined(); + }); +}); diff --git a/src/gateway/routing/request-match.ts b/src/gateway/routing/request-match.ts new file mode 100644 index 00000000..8d98f4ff --- /dev/null +++ b/src/gateway/routing/request-match.ts @@ -0,0 +1,44 @@ +/** + * The one place the gateway decides whether a worker can serve a request (ADR 0009 §3). It reads + * only the worker's catalog (ADR 0008 §1): the `can-serve` stage, the `warm-hit` stage, and the + * forwarded `lease.request` all ask it, so the worker resolves exactly what the gateway matched. + */ +import type { WorkerView } from "../worker-registry.js"; +import type { RoutableRequest } from "./pipeline.js"; + +/** + * The worker's own name for the requested model, or `undefined` when it cannot serve the request. + * + * The model is the first entry of `models` whose name or `modelAliases` entry equals the requested + * name, ignoring letter case. A named runtime must be among that model's `modelRuntimes`; an + * unnamed one needs the list to be non-empty. Only installed runtimes are listed, so a request + * that would need a download never matches. + */ +export function matchRequest(worker: WorkerView, request: RoutableRequest): string | undefined { + const catalog = worker.catalog.find((entry) => entry.platform === request.platform); + if (catalog === undefined) return undefined; + const wanted = fold(request.model); + const model = catalog.models.find( + (name) => + fold(name) === wanted || + ownList(catalog.modelAliases, name).some((alias) => fold(alias) === wanted), + ); + if (model === undefined) return undefined; + const runtimes = ownList(catalog.modelRuntimes, model); + const pairs = + request.osVersion === undefined ? runtimes.length > 0 : runtimes.includes(request.osVersion); + return pairs ? model : undefined; +} + +/** Locale-independent, so a gateway and a worker on different locales fold a name alike. */ +function fold(name: string): string { + return name.toLowerCase(); +} + +/** A record lookup that ignores inherited keys, so a model named `constructor` reads nothing. */ +function ownList( + record: Readonly>, + key: string, +): readonly string[] { + return Object.hasOwn(record, key) ? (record[key] ?? []) : []; +} diff --git a/src/gateway/routing/stages/can-serve.ts b/src/gateway/routing/stages/can-serve.ts new file mode 100644 index 00000000..d0a4c227 --- /dev/null +++ b/src/gateway/routing/stages/can-serve.ts @@ -0,0 +1,14 @@ +import { matchRequest } from "../request-match.js"; +import type { FilterStage } from "../pipeline.js"; + +/** + * Keeps a worker whose catalog pairs the requested model with the requested runtime (ADR 0009 + * §2, stage 2; §3 for the match). A download never makes a worker able to serve. + */ +export const canServe: FilterStage = { + keeps(worker, request) { + return matchRequest(worker, request) !== undefined; + }, + kind: "filter", + name: "can-serve", +}; diff --git a/src/gateway/routing/stages/eligible.ts b/src/gateway/routing/stages/eligible.ts index 2bf62390..31efc7af 100644 --- a/src/gateway/routing/stages/eligible.ts +++ b/src/gateway/routing/stages/eligible.ts @@ -1,6 +1,10 @@ import type { FilterStage } from "../pipeline.js"; /** + * The policy before ADR 0009 §2 and §3, kept only so the conformance tests can run the three + * legacy stages (`routing.test.ts`, `routing/warm-then-free.test.ts`). No registered policy uses + * it: `takes-requests` and `can-serve` replace it. + * * Drops a worker that is disconnected, incompatible, or drained; that lacks the requested * platform in its catalog at all; that lacks the requested `model`; that lacks the requested * `osVersion` (when the request named one) among its reported runtimes; or whose own diff --git a/src/gateway/routing/stages/takes-requests.test.ts b/src/gateway/routing/stages/takes-requests.test.ts new file mode 100644 index 00000000..3c94c0c8 --- /dev/null +++ b/src/gateway/routing/stages/takes-requests.test.ts @@ -0,0 +1,43 @@ +import { describe, expect, it } from "vitest"; + +import { createRoutingPolicy } from "../../routing.js"; +import { catalogFixture, deviceFixture, statusFixture } from "../../test-support.js"; +import type { WorkerView } from "../../worker-registry.js"; + +function view(id: string, overrides: Partial = {}): WorkerView { + return { + capacity: statusFixture().capacity, + catalog: catalogFixture([{ models: ["iPhone 17"], platform: "ios", runtimes: ["26.0"] }]) + .platforms, + connection: "connected", + devices: [], + drained: false, + id, + lastSeenAt: 1, + leases: [], + ...overrides, + }; +} + +const REQUEST = { allowDownload: false, model: "iPhone 17", platform: "ios" as const }; + +describe("takes-requests in the registered policy", () => { + const policy = createRoutingPolicy("warm-then-free"); + + it.each([ + ["disconnected", { connection: "disconnected" as const }], + ["incompatible", { connection: "incompatible" as const }], + ["drained", { drained: true }], + ["whose capacity has not been read", { capacity: undefined }], + ])("passes over a worker that is %s for one that takes requests", (_label, overrides) => { + // wrk_a sorts first and holds a warm device, so only takes-requests can pass it over. + const { capacity, ...rest } = { + ...view("wrk_a", { devices: [deviceFixture("dev_1", "ready")] }), + ...overrides, + }; + const excluded = capacity === undefined ? rest : { ...rest, capacity }; + + expect(policy.select(REQUEST, [excluded, view("wrk_b")])?.workerId).toBe("wrk_b"); + expect(policy.select(REQUEST, [excluded])).toBeUndefined(); + }); +}); diff --git a/src/gateway/routing/stages/takes-requests.ts b/src/gateway/routing/stages/takes-requests.ts new file mode 100644 index 00000000..f852ffc1 --- /dev/null +++ b/src/gateway/routing/stages/takes-requests.ts @@ -0,0 +1,13 @@ +import type { FilterStage } from "../pipeline.js"; + +/** + * Keeps a worker that is connected, not drained, and whose capacity has been read (ADR 0009 §2, + * stage 1). A policy that cannot see free capacity must not guess it has some. + */ +export const takesRequests: FilterStage = { + keeps(worker) { + return worker.connection === "connected" && !worker.drained && worker.capacity !== undefined; + }, + kind: "filter", + name: "takes-requests", +}; diff --git a/src/gateway/routing/stages/warm-hit.test.ts b/src/gateway/routing/stages/warm-hit.test.ts new file mode 100644 index 00000000..b0129782 --- /dev/null +++ b/src/gateway/routing/stages/warm-hit.test.ts @@ -0,0 +1,32 @@ +import { describe, expect, it } from "vitest"; + +import { catalogFixture, deviceFixture, statusFixture } from "../../test-support.js"; +import type { WorkerView } from "../../worker-registry.js"; +import { warmHit } from "./warm-hit.js"; + +describe("warm-hit", () => { + it("counts a ready device as warm when the request names its model by an alias in another letter case", () => { + const view: WorkerView = { + capacity: statusFixture().capacity, + catalog: catalogFixture([ + { + modelAliases: { "iPhone 17": ["iphone-17"] }, + models: ["iPhone 17"], + platform: "ios", + runtimes: ["26.0"], + }, + ]).platforms, + connection: "connected", + devices: [deviceFixture("dev_1", "ready")], + drained: false, + id: "wrk_a", + lastSeenAt: 1, + leases: [], + }; + const request = { allowDownload: false, osVersion: "26.0", platform: "ios" as const }; + + expect(warmHit.score(view, { ...request, model: "IPHONE-17" })).toBe(1); + expect(warmHit.score(view, { ...request, model: "iphone 17" })).toBe(1); + expect(warmHit.score(view, { ...request, model: "iPhone 16" })).toBe(0); + }); +}); diff --git a/src/gateway/routing/stages/warm-hit.ts b/src/gateway/routing/stages/warm-hit.ts index 0187e633..46f12b15 100644 --- a/src/gateway/routing/stages/warm-hit.ts +++ b/src/gateway/routing/stages/warm-hit.ts @@ -1,18 +1,22 @@ +import { matchRequest } from "../request-match.js"; import type { RankStage } from "../pipeline.js"; /** - * Prefers a worker with an unleased `ready` device matching the request. Settles: once a warm - * worker exists, no later stage looks at capacity. + * Prefers a worker with an unleased `ready` device matching the request, compared against the + * worker's own name for the model (ADR 0009 §3). A worker the catalog does not match keeps the + * exact compare; only the legacy `eligible` stage lets such a worker reach here. Settles: once a + * warm worker exists, no later stage looks at capacity. */ export const warmHit: RankStage = { kind: "rank", name: "warm-hit", score(worker, request) { + const model = matchRequest(worker, request) ?? request.model; const warm = worker.devices.some( (device) => device.state === "ready" && device.spec.platform === request.platform && - device.spec.model === request.model && + device.spec.model === model && (request.osVersion === undefined || device.spec.osVersion === request.osVersion), ); return warm ? 1 : 0; diff --git a/src/gateway/routing/warm-then-free.test.ts b/src/gateway/routing/warm-then-free.test.ts index 4b4b83d3..59f23b47 100644 --- a/src/gateway/routing/warm-then-free.test.ts +++ b/src/gateway/routing/warm-then-free.test.ts @@ -172,7 +172,7 @@ function withFreeIos(free: number): WorkerView["capacity"] { describe("warm-then-free as a list of stages", () => { it("the composed three-stage list picks the same worker and reason as the function on main for every generated fleet", () => { - const policy = createRoutingPolicy("warm-then-free"); + const policy = composeRoutingPolicy([eligible, warmHit, freeCapacity]); const outcomes = new Set(); for (let seed = 1; seed <= 5_000; seed++) { diff --git a/src/gateway/test-support.ts b/src/gateway/test-support.ts index 65175e12..8f3fbc3e 100644 --- a/src/gateway/test-support.ts +++ b/src/gateway/test-support.ts @@ -208,7 +208,7 @@ export class ScriptedWorkerClient { status: StatusGetOutput = statusFixture(); devices: unknown[] = []; catalog: CatalogOutput = catalogFixture([]); - /** What `config.get` reports; the view carries it as a routing input (ADR 0005 §13). */ + /** What `config.get` reports; the view carries it for display (ADR 0009 §3). */ downloadPolicy: DownloadPolicy = "on-request"; /** What `config.get` reports for `lease.maxTtlMs` (ADR 0005 §15) -- the routing-adjacent * counterpart to `downloadPolicy` above. Defaults comfortably above every gateway cap this diff --git a/src/gateway/worker-link.ts b/src/gateway/worker-link.ts index 6278556f..02dca6bb 100644 --- a/src/gateway/worker-link.ts +++ b/src/gateway/worker-link.ts @@ -10,8 +10,8 @@ * those events onto the gateway's own bus with `workerId` added. * * One call §7 does not name is here too: `config.get`, read once per session for the worker's - * `downloads.policy`. It is a routing input rather than decoration (§13: a request needing a - * download is only eligible on a worker whose policy allows one), it never changes without a + * `downloads.policy` and `lease.maxTtlMs`. The policy is shown on the view for display only: + * routing counts installed runtimes and never reads it (ADR 0009 §3). Neither changes without a * worker restart -- which is a new session anyway -- and the uplink session is admin, so the * gateway may read it. */