From 0e7fe8766af53b5c8ea5e707553d7e026431d643 Mon Sep 17 00:00:00 2001 From: spinloop-agent Date: Mon, 5 Oct 2026 23:41:17 +0100 Subject: [PATCH 1/3] docs(openspec): propose a per-environment lock on remote starts --- openspec/changes/wake-lock/.openspec.yaml | 2 + openspec/changes/wake-lock/design.md | 45 +++++++++++++++ openspec/changes/wake-lock/proposal.md | 31 +++++++++++ .../specs/endpoint-lifecycle/spec.md | 55 +++++++++++++++++++ openspec/changes/wake-lock/tasks.md | 22 ++++++++ 5 files changed, 155 insertions(+) create mode 100644 openspec/changes/wake-lock/.openspec.yaml create mode 100644 openspec/changes/wake-lock/design.md create mode 100644 openspec/changes/wake-lock/proposal.md create mode 100644 openspec/changes/wake-lock/specs/endpoint-lifecycle/spec.md create mode 100644 openspec/changes/wake-lock/tasks.md diff --git a/openspec/changes/wake-lock/.openspec.yaml b/openspec/changes/wake-lock/.openspec.yaml new file mode 100644 index 00000000..e3966d7a --- /dev/null +++ b/openspec/changes/wake-lock/.openspec.yaml @@ -0,0 +1,2 @@ +schema: spec-driven +created: 2026-10-05 diff --git a/openspec/changes/wake-lock/design.md b/openspec/changes/wake-lock/design.md new file mode 100644 index 00000000..af17417d --- /dev/null +++ b/openspec/changes/wake-lock/design.md @@ -0,0 +1,45 @@ +## Context + +See proposal.md for why. `wake()` in `remote/lambda/start/index.ts` reads the deploy config, finds the environment's Elastic IP and security group, runs the weights check (which can start a seed), looks the instance up by tag, then launches or re-wakes it and polls until the model answers. That poll can run for the Lambda's whole 15-minute limit. The tag lookup is `DescribeInstances`, which is eventually consistent, so it is not a safe way to tell whether another start has already launched. + +The Lambda's role already has `ssm:GetParameter` and `ssm:PutParameter` on `/cloud-vm-llm/*`. It has no `ssm:DeleteParameter`. Issue 224 says no new grant is needed; releasing a lock by deleting the parameter does need one. + +## Goals / Non-Goals + +**Goals:** +- Two starts for one environment never launch two instances, from any caller. +- A start that cannot take the lock does nothing but reply that it should retry. +- A crashed start cannot block an environment for longer than the time that start could have run. + +**Non-Goals:** +- Serialising stops, pauses or the idle sweep against starts. Those race a start in a different way (a stop between a start's lookup and its poll is already handled in `wake()`) and are left as they are. +- Making the lock safe against two starts that both find the same expired lock in the same few milliseconds. See Risks. +- A queue, so that a refused start waits its turn on the server. The caller retries. +- Any change to the CLI or gateway. + +## Decisions + +**An SSM parameter as the lock.** `/cloud-vm-llm//wake-lock`, created with `Overwrite: false`, so creation fails with `ParameterAlreadyExists` when another start holds it. This needs no new infrastructure. DynamoDB with a conditional put would give true compare-and-set but adds a table, a construct and grants to a stack that has none today; the SSM lock is enough for the case that matters (starts racing within seconds of each other, where creation is atomic). + +**What the lock holds.** A JSON value `{"owner": , "expiresAt": }`. The owner lets a release remove only a lock this start created. `expiresAt` is now plus the time remaining on this invocation plus a 30 second margin, so a lock lives exactly as long as the start that took it could run. A fixed TTL longer than the 900 second limit would leave a killed start's environment blocked for longer than needed. + +**Where it is taken and released.** `wake()` keeps its signature and becomes a wrapper: it runs the existing read-only checks that reply "unconfigured" and "undeployed" (they change nothing, so they need no lock), takes the lock, runs the rest of the existing body as `wakeLocked()`, and releases the lock in a `finally`. The lock is taken before the weights check, because that check can start a seed and two starts racing there could start two. The lock is held through the whole poll, not only the launch: the lookup cannot be trusted to see a just-launched instance for a while, so releasing at launch would reopen the race. + +**What a refused start returns.** HTTP 503 with state `starting`, a message saying another start for the environment is in progress, and `retry_after_seconds: 15` with the matching `retry-after` header. The Go client already retries any 503 using `retry_after_seconds` and prints "instance starting; retrying in 15s", and the gateway wake path treats a 503 the same way, so neither changes. When the holder finishes, the retry finds the running instance and returns ready. + +**Taking over an expired lock.** When creation fails because the parameter exists, the start reads it. If `expiresAt` is in the future, it replies as above. If it has passed, the start reads the lock again, checks the owner and expiry are unchanged, deletes it, and tries the create-if-absent again; whichever start's create succeeds holds the lock, and the other replies as above. An unreadable or malformed lock value counts as expired, since nothing else could ever clear it. + +**Failures other than "already exists".** Any other error creating the lock is logged and answered with a 503 `starting`, retryable, with nothing launched. A failure releasing the lock is logged and ignored: the reply to the caller is already decided, and the lock expires on its own. + +**IAM.** One statement on the start Lambda's role: `ssm:DeleteParameter` on `parameter/cloud-vm-llm/*/wake-lock`. Create and read use the existing grants. + +## Risks / Trade-offs + +- **Two starts taking over the same expired lock together.** SSM has no compare-and-set, so between one start's re-read and its delete, another can delete and re-create the lock, and the first then deletes the new one. This needs an abandoned lock (a killed start) and two starts arriving within milliseconds after its expiry. The result is the old race for that one start, no worse than today. A DynamoDB lock would remove it; this change accepts it and records it. +- **A killed start blocks the environment until the lock expires.** At most the remaining time of that invocation plus 30 seconds, which is at most about 15 minutes. Callers see "another start is in progress" for that time. +- **A refused start does not wait on the server.** Callers poll every 15 seconds, so a start can learn the instance is ready up to 15 seconds late. +- **Deployments that have not taken the new grant.** The start cannot delete its lock, so a start would hold the environment until expiry. The release failure is logged. Deploying the stack fixes it; the proposal says so. + +## Open Questions + +None that change what gets built. diff --git a/openspec/changes/wake-lock/proposal.md b/openspec/changes/wake-lock/proposal.md new file mode 100644 index 00000000..c7848bd6 --- /dev/null +++ b/openspec/changes/wake-lock/proposal.md @@ -0,0 +1,31 @@ +## Why + +The start Lambda decides whether to launch an instance by looking one up by tag. That lookup is eventually consistent, so two starts for the same environment that arrive close together can each miss the other's new instance and each launch one: two GPU instances for one environment, both billed (issue 224). + +The gateway already coalesces concurrent wakes inside one gateway process, but the control plane is open to every other caller: a second gateway, `spinloop remote start` racing a gateway, and a scheduled start (issue 178) firing while someone starts the environment by hand. + +## What Changes + +- The start Lambda takes a per-environment lock before it looks anything up and releases it on every way a start can end. +- A start that finds the lock held does not look up, launch or re-wake anything. It replies with a retryable 503 saying another start is in progress, which the CLI and gateway already retry. +- A lock left behind by a start that never finished (the Lambda was killed) expires when that start would have timed out, and the next start takes it over. +- The lock is an SSM parameter created only if absent, `/cloud-vm-llm//wake-lock`. The start Lambda's role gains `ssm:DeleteParameter` on that parameter name only. +- Locks are per environment: starts for different environments do not wait for each other. + +## Capabilities + +### New Capabilities + +None. + +### Modified Capabilities + +- `endpoint-lifecycle`: adds a requirement that concurrent starts of one environment are serialised and that a held or abandoned lock is handled as above. The existing "Starting on demand" behaviour is otherwise unchanged. + +## Impact + +- `remote/lambda/start/index.ts`: `wake()` becomes a lock wrapper around the existing body; a new shared module holds the lock. +- `remote/lib/llm-stack.ts`: one IAM statement on the start Lambda's role. +- `remote/test/`: new tests for the lock and for `wake()` under contention. +- `docs/maintainer/internals.md` and `remote/README.md`: a short note on the lock. +- No CLI, config or API change. Existing deployments need a `spinloop remote bootstrap` re-run (or `pnpm run deploy`) to get the new grant; until then a start fails at the lock instead of launching unlocked. diff --git a/openspec/changes/wake-lock/specs/endpoint-lifecycle/spec.md b/openspec/changes/wake-lock/specs/endpoint-lifecycle/spec.md new file mode 100644 index 00000000..1f3f0fb3 --- /dev/null +++ b/openspec/changes/wake-lock/specs/endpoint-lifecycle/spec.md @@ -0,0 +1,55 @@ +## ADDED Requirements + +### Requirement: Concurrent starts of one environment are serialised + +At most one start SHALL be working on an environment at a time, whoever calls it: the CLI, a gateway, a scheduled start, or any other client of the start endpoint. A start SHALL take the environment's lock before it checks the weights, looks up the instance, launches or re-wakes it, and SHALL release the lock on every way the start can end: ready, a retryable reply, a failure reply, or an error. A start that finds the lock held SHALL NOT check the weights, look up, launch or re-wake anything. It SHALL reply with a retryable state that says another start is in progress and when to retry. A start made after the holder has finished SHALL find the instance the holder started and report it as it would for any existing instance, so two starts never produce two launches. + +The lock SHALL be per environment: a start for one environment SHALL NOT wait for, or be refused because of, a start for another. A read of an environment's state SHALL NOT take the lock. + +#### Scenario: Two starts arrive together + +- **WHEN** two starts for the same environment reach the control plane at the same moment and the environment has no instance +- **THEN** exactly one instance is launched, and the start that did not take the lock replies with a retryable "another start is in progress" state + +#### Scenario: The refused start is retried + +- **WHEN** a start refused because the lock was held is retried after the first start has finished +- **THEN** it finds the first start's instance and reports it, and launches nothing + +#### Scenario: Different environments do not wait for each other + +- **WHEN** a start for environment `a` holds its lock and a start for environment `b` arrives +- **THEN** the start for `b` proceeds without waiting + +#### Scenario: Every ending releases the lock + +- **WHEN** a start ends because the weights are absent, no capacity was found, the instance went into a terminal state, the boot failed, the deadline passed, or an unexpected error was raised +- **THEN** the lock is released, and the next start for that environment is not refused because of it + +#### Scenario: Status does not take the lock + +- **WHEN** a start holds the environment's lock and the state of the environment is read +- **THEN** the read is answered normally + +### Requirement: An abandoned lock expires + +A lock SHALL record when it expires, no earlier than the moment the start holding it would itself run out of time. A start that finds an expired lock SHALL take it over and proceed as if it had found none. A lock that has not expired SHALL NOT be taken over, so a start that is still legitimately running keeps its lock for as long as it could run. + +#### Scenario: A killed start does not block the environment for ever + +- **WHEN** a start took the lock and its Lambda was killed before releasing it, and its time limit has since passed +- **THEN** the next start for that environment takes over the lock and proceeds + +#### Scenario: A running start keeps its lock + +- **WHEN** a start has held the lock for longer than a typical wake but not past its time limit +- **THEN** another start does not take the lock over + +### Requirement: A start does not proceed without its lock + +When the control plane cannot tell whether the lock is held, because the lock could not be read or written, the start SHALL NOT proceed to launch or re-wake. It SHALL reply with a retryable state and SHALL log the cause. + +#### Scenario: The lock cannot be written + +- **WHEN** a start cannot create the lock for a reason other than the lock already existing +- **THEN** nothing is launched, and the reply is retryable diff --git a/openspec/changes/wake-lock/tasks.md b/openspec/changes/wake-lock/tasks.md new file mode 100644 index 00000000..8c4a128d --- /dev/null +++ b/openspec/changes/wake-lock/tasks.md @@ -0,0 +1,22 @@ +## 1. The lock + +- [ ] 1.1 Add `remote/lambda/shared/wake-lock.ts`: the parameter name, the value format, `acquireWakeLock(env, owner, expiresAt)` (create if absent, take over an expired or unreadable one, report held), and `releaseWakeLock(env, owner)` (delete only when the owner matches) +- [ ] 1.2 Tests with a mocked SSM client: acquire when free, refuse when held, take over when expired, take over when malformed, release only own lock, error other than "already exists" is surfaced + +## 2. The start Lambda + +- [ ] 2.1 Split `wake()` into a wrapper and `wakeLocked()`: the wrapper runs the unconfigured and undeployed checks, takes the lock, calls `wakeLocked()` and releases in a `finally` +- [ ] 2.2 Reply 503 `starting` with `retry_after_seconds: 15` and the `retry-after` header when the lock is held, and when the lock cannot be taken for another reason (logged, nothing launched) +- [ ] 2.3 Set the lock's expiry from the invocation's remaining time plus 30 seconds +- [ ] 2.4 Tests through the handler: two concurrent starts launch once; a refused start launches and looks up nothing; the retry finds the instance; the lock is released after ready, no-capacity, terminal state, boot failure, deadline and a thrown error; different environments do not block each other; status does not take the lock; a scheduled start takes the lock + +## 3. Stack + +- [ ] 3.1 Grant the start Lambda's role `ssm:DeleteParameter` on `parameter/cloud-vm-llm/*/wake-lock` only +- [ ] 3.2 Stack test: the delete grant exists, is scoped to the lock parameter, and no other Lambda has it; run `scripts/check-no-cloud-identifiers.sh` + +## 4. Documentation + +- [ ] 4.1 Add a note to `docs/maintainer/internals.md` on the lock, its expiry and the takeover window +- [ ] 4.2 Note the lock and the "another start is in progress" reply in `remote/README.md` and the troubleshooting page if it lists start replies +- [ ] 4.3 Run `go test ./...`, `go vet ./...`, `gofmt`, `pnpm test` in `remote/`, `scripts/check-spec-purposes.sh` and `openspec validate wake-lock` From 01a6b3e8662d1f478d560a7c08c88b8f94937133 Mon Sep 17 00:00:00 2001 From: spinloop-agent Date: Mon, 5 Oct 2026 23:49:18 +0100 Subject: [PATCH 2/3] docs(openspec): end a start whose instance is stopped under it --- openspec/changes/wake-lock/design.md | 6 +++++- openspec/changes/wake-lock/proposal.md | 3 ++- .../specs/endpoint-lifecycle/spec.md | 19 +++++++++++++++++++ openspec/changes/wake-lock/tasks.md | 4 +++- 4 files changed, 29 insertions(+), 3 deletions(-) diff --git a/openspec/changes/wake-lock/design.md b/openspec/changes/wake-lock/design.md index af17417d..5d409d7b 100644 --- a/openspec/changes/wake-lock/design.md +++ b/openspec/changes/wake-lock/design.md @@ -12,7 +12,7 @@ The Lambda's role already has `ssm:GetParameter` and `ssm:PutParameter` on `/clo - A crashed start cannot block an environment for longer than the time that start could have run. **Non-Goals:** -- Serialising stops, pauses or the idle sweep against starts. Those race a start in a different way (a stop between a start's lookup and its poll is already handled in `wake()`) and are left as they are. +- Serialising stops, pauses or the idle sweep against starts. Stop is the way out of a start that has hung, so a lock held for up to 15 minutes would refuse it exactly when it is needed. The stop Lambda times out at 120 seconds, so it could not wait for the lock either. It would also need get, put and delete on the lock parameter, and the 5-minute sweep would make extra SSM calls for every environment. Stops race a start in a different way, and the start now handles that by ending promptly when it sees the stop (see Decisions). - Making the lock safe against two starts that both find the same expired lock in the same few milliseconds. See Risks. - A queue, so that a refused start waits its turn on the server. The caller retries. - Any change to the CLI or gateway. @@ -29,6 +29,10 @@ The Lambda's role already has `ssm:GetParameter` and `ssm:PutParameter` on `/clo **Taking over an expired lock.** When creation fails because the parameter exists, the start reads it. If `expiresAt` is in the future, it replies as above. If it has passed, the start reads the lock again, checks the owner and expiry are unchanged, deletes it, and tries the create-if-absent again; whichever start's create succeeds holds the lock, and the other replies as above. An unreadable or malformed lock value counts as expired, since nothing else could ever clear it. +**Ending a start whose instance was stopped.** Today, after a start has issued its own start command, an instance seen `stopped` has no branch in the first poll loop, and the later loops (agent online, daemon answering, health) never look at the instance state, so a stop mid-wake leaves the start polling until its deadline. With the lock that would also block the environment for those minutes. Each iteration of those loops now checks the instance state with `getInstance`, tolerating `InvalidInstanceID.NotFound` as the first loop already does, and when it sees `stopped`, `stopping`, `shutting-down` or `terminated` after the start command was issued it returns 503 with the observed state, a message that the instance was stopped while starting, and `retry_after_seconds: 15`. The `finally` in `wake()` releases the lock. The existing first-loop behaviours are unchanged: a `stopped` instance seen before the start command was issued is re-woken, and a dying instance in that loop keeps its existing reply. + +**Retrying is the client's choice, and it can undo a deliberate stop.** The reply is retryable, so a `spinloop remote start` that is still waiting will ask again and re-wake the instance a user just paused. That matches what happens today once the deadline passes, only sooner. A non-retryable reply would honour the stop but would make a scheduled 18:00 stop that lands during a slow start fail the start with an error rather than being a hiccup. This change keeps the reply retryable and does not try to decide whose intent wins. + **Failures other than "already exists".** Any other error creating the lock is logged and answered with a 503 `starting`, retryable, with nothing launched. A failure releasing the lock is logged and ignored: the reply to the caller is already decided, and the lock expires on its own. **IAM.** One statement on the start Lambda's role: `ssm:DeleteParameter` on `parameter/cloud-vm-llm/*/wake-lock`. Create and read use the existing grants. diff --git a/openspec/changes/wake-lock/proposal.md b/openspec/changes/wake-lock/proposal.md index c7848bd6..4f23d45b 100644 --- a/openspec/changes/wake-lock/proposal.md +++ b/openspec/changes/wake-lock/proposal.md @@ -11,6 +11,7 @@ The gateway already coalesces concurrent wakes inside one gateway process, but t - A lock left behind by a start that never finished (the Lambda was killed) expires when that start would have timed out, and the next start takes it over. - The lock is an SSM parameter created only if absent, `/cloud-vm-llm//wake-lock`. The start Lambda's role gains `ssm:DeleteParameter` on that parameter name only. - Locks are per environment: starts for different environments do not wait for each other. +- A start that sees its instance stopped, stopping or terminated after it has issued its start command now ends at once with a retryable reply, instead of polling until its deadline. Without this, a stop mid-start would leave the environment locked for up to 15 minutes. Stops themselves are not locked and never wait for a start. ## Capabilities @@ -20,7 +21,7 @@ None. ### Modified Capabilities -- `endpoint-lifecycle`: adds a requirement that concurrent starts of one environment are serialised and that a held or abandoned lock is handled as above. The existing "Starting on demand" behaviour is otherwise unchanged. +- `endpoint-lifecycle`: adds requirements that concurrent starts of one environment are serialised, that a held or abandoned lock is handled as above, and that a start ends promptly when its instance is stopped under it. The existing "Starting on demand" behaviour is otherwise unchanged. ## Impact diff --git a/openspec/changes/wake-lock/specs/endpoint-lifecycle/spec.md b/openspec/changes/wake-lock/specs/endpoint-lifecycle/spec.md index 1f3f0fb3..64208619 100644 --- a/openspec/changes/wake-lock/specs/endpoint-lifecycle/spec.md +++ b/openspec/changes/wake-lock/specs/endpoint-lifecycle/spec.md @@ -45,6 +45,25 @@ A lock SHALL record when it expires, no earlier than the moment the start holdin - **WHEN** a start has held the lock for longer than a typical wake but not past its time limit - **THEN** another start does not take the lock over +### Requirement: A stop during a start ends the start + +Once a start has launched the instance or issued the command to start it, and then sees the instance stopped, stopping, shutting down or terminated, the start SHALL end at once with a retryable reply that names the state it saw, and SHALL release the lock. It SHALL NOT keep polling an instance that is not coming up until its time limit. The reply SHALL be retryable so that a client still waiting on its start can ask again, which re-wakes the instance or launches a fresh one; a stop does not wait for, or take, the start's lock. + +#### Scenario: The instance is stopped after the start command + +- **WHEN** a start has issued the command to start a stopped instance and the instance is then stopped by a pause, a scheduled stop or the idle sweep +- **THEN** the start ends promptly with a retryable reply naming the state, and the lock is released + +#### Scenario: The instance is terminated while the engine loads + +- **WHEN** an instance is terminated after it has reached running, while the start waits for the model to answer +- **THEN** the start ends promptly with a retryable reply naming the state, and the lock is released + +#### Scenario: A stop is never refused because a start is running + +- **WHEN** a start holds the environment's lock and a pause, stop or scheduled stop for the environment arrives +- **THEN** the stop is carried out and is not refused or delayed by the lock + ### Requirement: A start does not proceed without its lock When the control plane cannot tell whether the lock is held, because the lock could not be read or written, the start SHALL NOT proceed to launch or re-wake. It SHALL reply with a retryable state and SHALL log the cause. diff --git a/openspec/changes/wake-lock/tasks.md b/openspec/changes/wake-lock/tasks.md index 8c4a128d..a25ec128 100644 --- a/openspec/changes/wake-lock/tasks.md +++ b/openspec/changes/wake-lock/tasks.md @@ -8,7 +8,9 @@ - [ ] 2.1 Split `wake()` into a wrapper and `wakeLocked()`: the wrapper runs the unconfigured and undeployed checks, takes the lock, calls `wakeLocked()` and releases in a `finally` - [ ] 2.2 Reply 503 `starting` with `retry_after_seconds: 15` and the `retry-after` header when the lock is held, and when the lock cannot be taken for another reason (logged, nothing launched) - [ ] 2.3 Set the lock's expiry from the invocation's remaining time plus 30 seconds -- [ ] 2.4 Tests through the handler: two concurrent starts launch once; a refused start launches and looks up nothing; the retry finds the instance; the lock is released after ready, no-capacity, terminal state, boot failure, deadline and a thrown error; different environments do not block each other; status does not take the lock; a scheduled start takes the lock +- [ ] 2.4 Check the instance state on each iteration of the agent-online, daemon-answering and health loops and in the first loop once the start command was issued; on `stopped`, `stopping`, `shutting-down` or `terminated` return 503 with the observed state and `retry_after_seconds: 15` (tolerating `InvalidInstanceID.NotFound`), leaving the first loop's existing pre-start behaviour as it is +- [ ] 2.5 Tests through the handler: two concurrent starts launch once; a refused start launches and looks up nothing; the retry finds the instance; the lock is released after ready, no-capacity, terminal state, boot failure, deadline and a thrown error; different environments do not block each other; status does not take the lock; a scheduled start takes the lock +- [ ] 2.6 Tests for a stop mid-start: stopped after the start command, stopped while waiting for the agent, terminated while the engine loads, each ends promptly with a retryable reply naming the state and releasing the lock; a `stopped` instance seen before the start command is still re-woken; a stop request is not affected by a held lock ## 3. Stack From d8cb56bfa7d6f7c8e04b8568ca84d40388622d22 Mon Sep 17 00:00:00 2001 From: spinloop-agent Date: Mon, 5 Oct 2026 23:57:24 +0100 Subject: [PATCH 3/3] fix(remote): serialise starts of one environment with a lock A start took no lock, so two starts close together could each miss the other's not-yet-visible instance and launch one apiece. Starts now take a per-environment SSM lock that expires with the invocation, a refused start replies with a retryable 503, and a start whose instance is stopped under it ends at once and releases the lock. Status reports a start in progress. Closes #224 --- docs/commands/remote.md | 11 + docs/maintainer/internals.md | 2 + docs/troubleshooting.md | 6 + openspec/changes/wake-lock/design.md | 4 + openspec/changes/wake-lock/proposal.md | 3 +- .../specs/endpoint-lifecycle/spec.md | 38 +- openspec/changes/wake-lock/tasks.md | 28 +- remote/README.md | 11 + remote/lambda/shared/wake-lock.ts | 150 +++++ remote/lambda/start/index.ts | 152 ++++- remote/lib/llm-stack.ts | 11 + remote/test/stack.test.ts | 26 + remote/test/start-boot-failure.test.ts | 7 + remote/test/start-launch.test.ts | 7 + remote/test/start-rewake.test.ts | 7 + remote/test/start-seeding.test.ts | 7 + remote/test/start-status.test.ts | 8 + remote/test/wake-lock-start.test.ts | 546 ++++++++++++++++++ remote/test/wake-lock.test.ts | 184 ++++++ 19 files changed, 1191 insertions(+), 17 deletions(-) create mode 100644 remote/lambda/shared/wake-lock.ts create mode 100644 remote/test/wake-lock-start.test.ts create mode 100644 remote/test/wake-lock.test.ts diff --git a/docs/commands/remote.md b/docs/commands/remote.md index 9bcb9440..e54aa7e8 100644 --- a/docs/commands/remote.md +++ b/docs/commands/remote.md @@ -254,6 +254,17 @@ output degrades rather than goes blank. `--format=table` is a key-value table, and `--format=json` is the raw reply. A stopped engine keeps its readings: the sparkline runs to the stop, ending at it. +While a start is under way, `status` shows **`starting`** for the endpoint, +whichever machine began the start, until the instance is running. Once it is +running but still loading the model it reads `running` and not ready, as +before. Two starts for one endpoint never launch two instances: the second is +told "another start is in progress" and retries on its own, so running +`start` twice, or from two machines, is safe. A start that is waiting for GPU +capacity is between attempts, not in progress, and `status` does not show it. +Stopping an endpoint during a start is always allowed; the start ends and +reports that the instance was stopped, and a client that keeps retrying will +wake it again. + Both report **`active`** — how long since the endpoint's engine last did any work. It comes from the activity the on-instance daemon tracks, so it is one answer decided on the box rather than something each command re-derives diff --git a/docs/maintainer/internals.md b/docs/maintainer/internals.md index 982722bf..65e4b5a8 100644 --- a/docs/maintainer/internals.md +++ b/docs/maintainer/internals.md @@ -35,6 +35,8 @@ These are mistakes already made here; each was silent rather than loud, which is **An IAM user's inline policies are capped at 2,048 characters in aggregate.** The control plane's seven functions each take a `grantInvokeUrl` pair — two actions, the auth-type conditions, the function's ARN — and with the log-reading, stack-discovery, pricing and self-service statements the document far exceeds that; the first deploy of the `RemoteCliUser` inline policy failed with `ServiceLimitExceeded`, which CDK does not warn about ahead of time. It is now a stack-owned `AWS::IAM::ManagedPolicy` (`RemoteCliPolicy`, 6,144 cap; the deployed document measures ~2 KB, so the grant list has room to grow). Keep it managed rather than re-inlining it, and keep the iam self-service ARN built from the `AWS::Partition`/`AWS::AccountId` pseudo parameters instead of the user's `Arn`: the policy attaches to that user, so referencing the user from inside it is a dependency cycle. +**A start holds a per-environment lock for its whole run, and it is an SSM parameter, not a conditional write.** `wake()` in the start Lambda takes `/cloud-vm-llm//wake-lock` (created with `Overwrite: false`) after the read-only checks and before the weights check, and releases it in a `finally`. It has to cover the whole poll and not only the launch, because the tag lookup that decides whether to launch lags a launch: a second start that got in after the lock was released could still miss the first one's instance. The value is `{owner, expiresAt}` with the expiry set to the invocation's remaining time plus 30 seconds, so a start killed before its `finally` blocks the environment for no longer than it could have run. A refused start replies 503 `starting` with 15 seconds to retry, which the CLI and gateway already retry. Taking over an expired lock is read, re-read, delete, create-if-absent; SSM has no compare-and-set, so two starts taking over the same abandoned lock within milliseconds can still both proceed. That needs a killed start plus a near-simultaneous pair, and a DynamoDB conditional write would remove it. Releasing needs `ssm:DeleteParameter`, granted on `parameter/cloud-vm-llm/*/wake-lock` only. Stops do not take the lock: a lock held for up to 15 minutes would refuse the stop that gets you out of a hung start, and the stop Lambda's 120 second limit could not wait for it. Instead each polling loop in a start checks the instance state and ends the start, releasing the lock, if the instance is stopped, stopping or terminated. The status read reports the lock without taking it: `starting` when the lock is held and the instance is absent, stopped or pending, and `start_in_progress` in every reply. A start waiting for capacity has released its lock between attempts, so it is not shown. + **The file credential store's index is non-secret by design.** OS keystores offer no way to list entries, so the file store — used where no keystore is reachable or `SPINLOOP_REMOTE_KEYSTORE=file` — keeps a plain-text index of the stored regions beside the `0600` per-region files under `/keystore/`. The report (`spinloop remote auth`) reads the index, so a file added, removed or renamed by hand is reported wrong until the index matches; and a corrupt index is reported, not silently reset, because a report that misleads about what is stored misleads about which access keys exist on the AWS side. **The two local model caches are separate.** A model `llama-server` downloaded sits in llama.cpp's cache (`$LLAMA_CACHE`, else the platform's user cache directory) as flat filenames; one fetched with the hub's tools sits in the Hugging Face hub cache (`$HF_HUB_CACHE`, else `$HF_HOME/hub`, else `~/.cache/huggingface/hub`) as `models----/snapshots//` of symlinks into a content-addressed blob store. Neither tool looks in the other's, so a model already on the machine is "already on the machine" in only one of the two senses. `spinloop hf` therefore resolves both roots up front (`hf.ResolveRoots`) and checks both before touching the network; a cache-aware lookup that consults only one side re-downloads what is already there. The hub-cache shape has two more traps: a snapshot entry whose symlink dangles is an interrupted or half-finished download and must count as absent, and `refs/` holds a commit sha, so a revision name is only a snapshot once it has been resolved through `refs/`. diff --git a/docs/troubleshooting.md b/docs/troubleshooting.md index 43734c08..5f2b494c 100644 --- a/docs/troubleshooting.md +++ b/docs/troubleshooting.md @@ -86,6 +86,12 @@ scaled; see [`spinloop serve`](commands/serve.md#parallelism). stderr; `--timeout` (default 15m) bounds the wait. `status` and `logs` answer while it boots and after it is gone — logs are readable even from a terminated instance. +- **`start` says another start is in progress.** Another command, machine or + gateway is already starting that endpoint, and only one start works on an + endpoint at a time. `start` retries on its own and finishes when the other + one does. If nothing is really starting, the lock a crashed start left + behind expires within about fifteen minutes. `status` shows `starting` + meanwhile. - **Quota.** Bootstrap needs enough GPU vCPU quota for a later launch; a launch that can't get an instance reports the AWS error. diff --git a/openspec/changes/wake-lock/design.md b/openspec/changes/wake-lock/design.md index 5d409d7b..17be7090 100644 --- a/openspec/changes/wake-lock/design.md +++ b/openspec/changes/wake-lock/design.md @@ -33,6 +33,10 @@ The Lambda's role already has `ssm:GetParameter` and `ssm:PutParameter` on `/clo **Retrying is the client's choice, and it can undo a deliberate stop.** The reply is retryable, so a `spinloop remote start` that is still waiting will ask again and re-wake the instance a user just paused. That matches what happens today once the deadline passes, only sooner. A non-retryable reply would honour the stop but would make a scheduled 18:00 stop that lands during a slow start fail the start with an error rather than being a hiccup. This change keeps the reply retryable and does not try to decide whose intent wins. +**Showing a start in progress.** The GET status branch of the start Lambda reads the lock (`wakeLockHeld`: the parameter exists, parses, and has not expired) and adds `start_in_progress` to its reply. When the lock is held and the instance is absent, `stopped` or `pending`, `state` is `starting`, so `spinloop status`, which prints the state string as it is, shows it with no change to the Go client. When the instance is already `running` the state stays `running` with `healthy: false` and only the flag marks the start, so nothing that keys off `running` changes. A failure to read the lock is logged and treated as no start in progress, because a status read that fails over the lock is worse than one that omits it. Reads take no lock and write nothing. + +**What status does not show.** A start waiting for capacity holds no lock between attempts (it released the lock when it replied no-capacity), so an environment in that wait looks idle. Showing it would need a record of the last start result, which is a separate feature. + **Failures other than "already exists".** Any other error creating the lock is logged and answered with a 503 `starting`, retryable, with nothing launched. A failure releasing the lock is logged and ignored: the reply to the caller is already decided, and the lock expires on its own. **IAM.** One statement on the start Lambda's role: `ssm:DeleteParameter` on `parameter/cloud-vm-llm/*/wake-lock`. Create and read use the existing grants. diff --git a/openspec/changes/wake-lock/proposal.md b/openspec/changes/wake-lock/proposal.md index 4f23d45b..d3f60d24 100644 --- a/openspec/changes/wake-lock/proposal.md +++ b/openspec/changes/wake-lock/proposal.md @@ -11,6 +11,7 @@ The gateway already coalesces concurrent wakes inside one gateway process, but t - A lock left behind by a start that never finished (the Lambda was killed) expires when that start would have timed out, and the next start takes it over. - The lock is an SSM parameter created only if absent, `/cloud-vm-llm//wake-lock`. The start Lambda's role gains `ssm:DeleteParameter` on that parameter name only. - Locks are per environment: starts for different environments do not wait for each other. +- The status read shows a start in progress: while a start holds the lock, the state is `starting` (until the instance is running) and the reply carries `start_in_progress: true`, so a second client can see that another client began a start. A start waiting for GPU capacity holds no lock between attempts and is not shown. - A start that sees its instance stopped, stopping or terminated after it has issued its start command now ends at once with a retryable reply, instead of polling until its deadline. Without this, a stop mid-start would leave the environment locked for up to 15 minutes. Stops themselves are not locked and never wait for a start. ## Capabilities @@ -21,7 +22,7 @@ None. ### Modified Capabilities -- `endpoint-lifecycle`: adds requirements that concurrent starts of one environment are serialised, that a held or abandoned lock is handled as above, and that a start ends promptly when its instance is stopped under it. The existing "Starting on demand" behaviour is otherwise unchanged. +- `endpoint-lifecycle`: adds requirements that concurrent starts of one environment are serialised, that a held or abandoned lock is handled as above, that a start ends promptly when its instance is stopped under it, and that status shows a start in progress. The existing "Starting on demand" behaviour is otherwise unchanged. ## Impact diff --git a/openspec/changes/wake-lock/specs/endpoint-lifecycle/spec.md b/openspec/changes/wake-lock/specs/endpoint-lifecycle/spec.md index 64208619..ffa5c893 100644 --- a/openspec/changes/wake-lock/specs/endpoint-lifecycle/spec.md +++ b/openspec/changes/wake-lock/specs/endpoint-lifecycle/spec.md @@ -31,6 +31,42 @@ The lock SHALL be per environment: a start for one environment SHALL NOT wait fo - **WHEN** a start holds the environment's lock and the state of the environment is read - **THEN** the read is answered normally +### Requirement: Status shows a start in progress + +The state of an environment SHALL show that a start is in progress while an unexpired start lock exists, whichever client began the start. While the lock is held and the instance is absent, stopped or still pending, the reported state SHALL be `starting`. In every reply while the lock is held, including when the instance is already running and its model is still loading, the reply SHALL carry `start_in_progress: true`; it SHALL carry `start_in_progress: false` otherwise. A lock that has expired SHALL NOT count. A failure to read the lock SHALL NOT fail the read of the environment's state: the reply is given as if no start were in progress. + +A start that is waiting for GPU capacity, or for any other retry, holds no lock between attempts and is not shown as in progress. + +#### Scenario: A second client sees a start another client began + +- **WHEN** one client has started an environment that has no instance yet and a second client reads the environment's state +- **THEN** the state is `starting` and `start_in_progress` is true + +#### Scenario: A loading model is still reported as a start + +- **WHEN** the instance is running and its model is still loading under a start that holds the lock +- **THEN** the state is the instance's own (`running`), healthy is false, and `start_in_progress` is true + +#### Scenario: No start, no flag + +- **WHEN** no start holds the lock +- **THEN** the state is the instance's own and `start_in_progress` is false + +#### Scenario: An abandoned lock is not a start in progress + +- **WHEN** the lock's expiry has passed +- **THEN** the state read does not report a start in progress + +#### Scenario: A start waiting for capacity is not in progress + +- **WHEN** a start found no capacity and is waiting to retry +- **THEN** the environment is not reported as starting between attempts + +#### Scenario: An unreadable lock does not fail the read + +- **WHEN** the lock cannot be read +- **THEN** the state is still returned, without a start in progress + ### Requirement: An abandoned lock expires A lock SHALL record when it expires, no earlier than the moment the start holding it would itself run out of time. A start that finds an expired lock SHALL take it over and proceed as if it had found none. A lock that has not expired SHALL NOT be taken over, so a start that is still legitimately running keeps its lock for as long as it could run. @@ -47,7 +83,7 @@ A lock SHALL record when it expires, no earlier than the moment the start holdin ### Requirement: A stop during a start ends the start -Once a start has launched the instance or issued the command to start it, and then sees the instance stopped, stopping, shutting down or terminated, the start SHALL end at once with a retryable reply that names the state it saw, and SHALL release the lock. It SHALL NOT keep polling an instance that is not coming up until its time limit. The reply SHALL be retryable so that a client still waiting on its start can ask again, which re-wakes the instance or launches a fresh one; a stop does not wait for, or take, the start's lock. +Once a start has launched the instance, issued the command to start it, or seen it running, and then sees the instance stopped, stopping, shutting down or terminated, the start SHALL end at once with a retryable reply that names the state it saw, and SHALL release the lock. It SHALL NOT keep polling an instance that is not coming up until its time limit. The reply SHALL be retryable so that a client still waiting on its start can ask again, which re-wakes the instance or launches a fresh one; a stop does not wait for, or take, the start's lock. #### Scenario: The instance is stopped after the start command diff --git a/openspec/changes/wake-lock/tasks.md b/openspec/changes/wake-lock/tasks.md index a25ec128..8df13856 100644 --- a/openspec/changes/wake-lock/tasks.md +++ b/openspec/changes/wake-lock/tasks.md @@ -1,24 +1,26 @@ ## 1. The lock -- [ ] 1.1 Add `remote/lambda/shared/wake-lock.ts`: the parameter name, the value format, `acquireWakeLock(env, owner, expiresAt)` (create if absent, take over an expired or unreadable one, report held), and `releaseWakeLock(env, owner)` (delete only when the owner matches) -- [ ] 1.2 Tests with a mocked SSM client: acquire when free, refuse when held, take over when expired, take over when malformed, release only own lock, error other than "already exists" is surfaced +- [x] 1.1 Add `remote/lambda/shared/wake-lock.ts`: the parameter name, the value format, `acquireWakeLock(env, owner, expiresAt)` (create if absent, take over an expired or unreadable one, report held), and `releaseWakeLock(env, owner)` (delete only when the owner matches) +- [x] 1.2 Tests with a mocked SSM client: acquire when free, refuse when held, take over when expired, take over when malformed, release only own lock, error other than "already exists" is surfaced ## 2. The start Lambda -- [ ] 2.1 Split `wake()` into a wrapper and `wakeLocked()`: the wrapper runs the unconfigured and undeployed checks, takes the lock, calls `wakeLocked()` and releases in a `finally` -- [ ] 2.2 Reply 503 `starting` with `retry_after_seconds: 15` and the `retry-after` header when the lock is held, and when the lock cannot be taken for another reason (logged, nothing launched) -- [ ] 2.3 Set the lock's expiry from the invocation's remaining time plus 30 seconds -- [ ] 2.4 Check the instance state on each iteration of the agent-online, daemon-answering and health loops and in the first loop once the start command was issued; on `stopped`, `stopping`, `shutting-down` or `terminated` return 503 with the observed state and `retry_after_seconds: 15` (tolerating `InvalidInstanceID.NotFound`), leaving the first loop's existing pre-start behaviour as it is -- [ ] 2.5 Tests through the handler: two concurrent starts launch once; a refused start launches and looks up nothing; the retry finds the instance; the lock is released after ready, no-capacity, terminal state, boot failure, deadline and a thrown error; different environments do not block each other; status does not take the lock; a scheduled start takes the lock -- [ ] 2.6 Tests for a stop mid-start: stopped after the start command, stopped while waiting for the agent, terminated while the engine loads, each ends promptly with a retryable reply naming the state and releasing the lock; a `stopped` instance seen before the start command is still re-woken; a stop request is not affected by a held lock +- [x] 2.1 Split `wake()` into a wrapper and `wakeLocked()`: the wrapper runs the unconfigured and undeployed checks, takes the lock, calls `wakeLocked()` and releases in a `finally` +- [x] 2.2 Reply 503 `starting` with `retry_after_seconds: 15` and the `retry-after` header when the lock is held, and when the lock cannot be taken for another reason (logged, nothing launched) +- [x] 2.3 Set the lock's expiry from the invocation's remaining time plus 30 seconds +- [x] 2.4 Check the instance state on each iteration of the agent-online, daemon-answering and health loops and in the first loop once the start command was issued; on `stopped`, `stopping`, `shutting-down` or `terminated` return 503 with the observed state and `retry_after_seconds: 15` (tolerating a lookup that fails), leaving the first loop's existing pre-start behaviour as it is +- [x] 2.5 Tests through the handler: two concurrent starts launch once; a refused start launches and looks up nothing; the retry finds the instance; the lock is released after ready, a retryable refusal, a terminal state and a thrown error, and only when still held by this start; different environments do not block each other; an expired lock is taken over and a valid one is not; a lock that cannot be written stops the start; a status read and a pause are not affected by a held lock +- [x] 2.6 Tests for a stop mid-start: stopped after the start command, right after a fresh launch, while waiting for the agent, terminated while the engine loads, each ends promptly with a retryable reply naming the state and releasing the lock; a lookup that fails keeps waiting; a `stopped` instance seen before the start command is still re-woken +- [x] 2.7 Add `wakeLockHeld(env)` to the lock module and make the GET status reply read it: `state` becomes `starting` while the lock is held and the instance is absent, stopped or pending, every reply carries `start_in_progress`, and a failure to read the lock is logged and treated as no start in progress +- [x] 2.8 Tests for status: a second client sees `starting` and the flag during another client's start; a running instance whose model is loading keeps `running` with the flag; no lock means no flag; an expired lock is ignored; an unreadable lock does not fail the read; a refused start's lock holder is unchanged by the read ## 3. Stack -- [ ] 3.1 Grant the start Lambda's role `ssm:DeleteParameter` on `parameter/cloud-vm-llm/*/wake-lock` only -- [ ] 3.2 Stack test: the delete grant exists, is scoped to the lock parameter, and no other Lambda has it; run `scripts/check-no-cloud-identifiers.sh` +- [x] 3.1 Grant the start Lambda's role `ssm:DeleteParameter` on `parameter/cloud-vm-llm/*/wake-lock` only +- [x] 3.2 Stack test: the delete grant exists, is scoped to the lock parameter, and no other Lambda has it; run `scripts/check-no-cloud-identifiers.sh` ## 4. Documentation -- [ ] 4.1 Add a note to `docs/maintainer/internals.md` on the lock, its expiry and the takeover window -- [ ] 4.2 Note the lock and the "another start is in progress" reply in `remote/README.md` and the troubleshooting page if it lists start replies -- [ ] 4.3 Run `go test ./...`, `go vet ./...`, `gofmt`, `pnpm test` in `remote/`, `scripts/check-spec-purposes.sh` and `openspec validate wake-lock` +- [x] 4.1 Add a note to `docs/maintainer/internals.md` on the lock, its expiry, the takeover window and the status reading +- [x] 4.2 Note the lock, the "another start is in progress" reply and the `starting` state in `remote/README.md` and the troubleshooting page if it lists start replies +- [x] 4.3 Run `go test ./...`, `go vet ./...`, `gofmt`, `pnpm test` in `remote/`, `scripts/check-spec-purposes.sh` and `openspec validate wake-lock` diff --git a/remote/README.md b/remote/README.md index 8953ca95..482be8ff 100644 --- a/remote/README.md +++ b/remote/README.md @@ -512,6 +512,17 @@ deregister the AMIs, and delete their snapshots by hand to reclaim that storage. `spinloop remote deploy`. - **`start` returns `no-ami`**: no AMI is tagged for the engine you asked for. Run `pnpm bake ` and wait for it to reach `AVAILABLE`. +- **`start` returns `starting` with "another start … is in progress"**: a + start for that environment already holds its lock (the SSM parameter + `/cloud-vm-llm//wake-lock`), so this one launched nothing. It is + retryable and the CLI retries it. A start killed before it could release the + lock blocks the environment until the lock's expiry, at most the Lambda's + remaining time plus 30 seconds; deleting the parameter by hand clears it + sooner. The `start` Lambda's log shows `"phase":"lock-held"`, + `"lock-error"` (the lock could not be taken; nothing is launched) or + `"stopped-under-start"` (the instance was stopped while the start waited). + The start Lambda's role needs `ssm:DeleteParameter` on the lock, so a + control plane deployed before it needs `pnpm run deploy`. - **`start` returns `no-capacity`**: every configured AZ was out of g6e capacity at that moment. The start Lambda already tried them all; wait a few minutes and retry, or widen/adjust `availabilityZones`. diff --git a/remote/lambda/shared/wake-lock.ts b/remote/lambda/shared/wake-lock.ts new file mode 100644 index 00000000..63ad4199 --- /dev/null +++ b/remote/lambda/shared/wake-lock.ts @@ -0,0 +1,150 @@ +/** + * The per-environment lock a start holds while it works on an environment. + * + * The start Lambda decides whether to launch by looking the instance up by tag, + * which is eventually consistent, so two starts close together can each miss the + * other's new instance and each launch one. The lock is an SSM parameter created + * only if absent, so exactly one of two simultaneous starts creates it. + * + * The value records who holds the lock and when it expires. The expiry is the + * time the holder could run until, so a start that was killed before it could + * release the lock blocks its environment for no longer than that. + */ + +import { + DeleteParameterCommand, + GetParameterCommand, + PutParameterCommand, + SSMClient, +} from '@aws-sdk/client-ssm'; +import { errorName } from './aws'; + +const ssm = new SSMClient({}); + +/** SSM parameter holding an environment's start lock. */ +export function wakeLockParam(env: string): string { + return `/cloud-vm-llm/${env}/wake-lock`; +} + +interface LockValue { + owner: string; + expiresAt: string; +} + +/** The lock a stored value describes, or null when it cannot be read as one. */ +function parseLock(raw: string | undefined): LockValue | null { + if (!raw) { + return null; + } + try { + const parsed = JSON.parse(raw) as Partial; + if (typeof parsed.owner === 'string' && typeof parsed.expiresAt === 'string') { + return { owner: parsed.owner, expiresAt: parsed.expiresAt }; + } + } catch { + // Fall through: a value that is not JSON is not a lock anyone could release. + } + return null; +} + +/** Create the lock if no parameter exists. True when this call created it. */ +async function createIfAbsent(env: string, value: string): Promise { + try { + await ssm.send( + new PutParameterCommand({ + Name: wakeLockParam(env), + Value: value, + Type: 'String', + Overwrite: false, + }), + ); + return true; + } catch (err) { + if (errorName(err) === 'ParameterAlreadyExists') { + return false; + } + throw err; + } +} + +async function readRaw(env: string): Promise { + try { + const out = await ssm.send(new GetParameterCommand({ Name: wakeLockParam(env) })); + return out.Parameter?.Value ?? ''; + } catch (err) { + if (errorName(err) === 'ParameterNotFound') { + return null; + } + throw err; + } +} + +async function remove(env: string): Promise { + try { + await ssm.send(new DeleteParameterCommand({ Name: wakeLockParam(env) })); + } catch (err) { + if (errorName(err) !== 'ParameterNotFound') { + throw err; + } + } +} + +/** + * Take the environment's lock for `owner`, expiring at `expiresAt`. Returns + * true when this call holds the lock, false when another start does. An expired + * lock, or one whose value cannot be read, is taken over. Any other failure to + * read or write the lock is thrown, so the caller never proceeds unlocked. + * + * SSM has no compare-and-set, so two starts that both find the same expired + * lock within milliseconds can both pass the re-read below; the create that + * follows lets only one of them win unless a third start deletes in between. + */ +export async function acquireWakeLock( + env: string, + owner: string, + expiresAt: Date, + now: Date = new Date(), +): Promise { + const value = JSON.stringify({ owner, expiresAt: expiresAt.toISOString() } satisfies LockValue); + if (await createIfAbsent(env, value)) { + return true; + } + const current = await readRaw(env); + if (current === null) { + // Released between the failed create and the read. + return createIfAbsent(env, value); + } + const lock = parseLock(current); + if (lock && new Date(lock.expiresAt).getTime() > now.getTime()) { + return false; + } + // Expired, or not a lock at all. Confirm it is still the same value, so a + // lock another start has just taken over is not deleted from under it. + if ((await readRaw(env)) !== current) { + return false; + } + await remove(env); + return createIfAbsent(env, value); +} + +/** + * Whether a start currently holds the environment's lock: the parameter exists, + * reads as a lock, and has not expired. Writes nothing. A failure to read it is + * thrown, and the caller decides what that means for what it is reporting. + */ +export async function wakeLockHeld(env: string, now: Date = new Date()): Promise { + const lock = parseLock((await readRaw(env)) ?? undefined); + return lock !== null && new Date(lock.expiresAt).getTime() > now.getTime(); +} + +/** Release the lock if `owner` still holds it. A lock someone else holds is left alone. */ +export async function releaseWakeLock(env: string, owner: string): Promise { + const current = await readRaw(env); + if (current === null) { + return; + } + if (parseLock(current)?.owner !== owner) { + return; + } + await remove(env); +} diff --git a/remote/lambda/start/index.ts b/remote/lambda/start/index.ts index 00c82b96..baf71592 100644 --- a/remote/lambda/start/index.ts +++ b/remote/lambda/start/index.ts @@ -1,3 +1,4 @@ +import { randomUUID } from 'node:crypto'; import type { Context, LambdaFunctionURLEvent, LambdaFunctionURLResult } from 'aws-lambda'; import { RETAIN_UNTIL_TAG, @@ -29,6 +30,7 @@ import { baseUrlFor, deployConfigParam, ENV_TAG_KEY, + type EnvEip, environmentFrom, findEnvEip, findEnvSecurityGroup, @@ -36,6 +38,7 @@ import { } from '../shared/environments'; import { DAEMON_STATUS_CMD, parseDaemonStatus } from '../shared/daemon'; import { jsonResponse } from '../shared/http'; +import { acquireWakeLock, releaseWakeLock, wakeLockHeld } from '../shared/wake-lock'; import { weightsPresent } from '../shared/seed'; import { findSeedInstances, seedAlive } from '../shared/seed/discovery'; import { seedIdFor } from '../shared/seed/identity'; @@ -75,6 +78,15 @@ const HEALTH_POLL_MS = 10_000; // now, so a wake must find the one it is trying to revive rather than fail on // it. Only states headed for the scrapyard end a wake. const TERMINAL_STATES = new Set(['shutting-down', 'terminated']); +// Once a start has launched an instance or issued its start command, any of +// these means a stop (or the sweep) got there first and the instance is not +// coming up, so the start ends rather than polling it until the deadline. +const GONE_STATES = new Set(['stopped', 'stopping', 'shutting-down', 'terminated']); + +// How long a refused start is told to wait, and how far past its own time +// limit a lock stays valid. +const LOCK_RETRY_SECONDS = 15; +const LOCK_MARGIN_MS = 30_000; const HEALTH_COMMAND = `curl -s -o /dev/null -w "%{http_code}" --max-time 5 http://localhost:${ENGINE_PORT}/health || true`; @@ -139,19 +151,39 @@ export async function handler( return wake(env, context, retainUntil); } +/** + * Whether a start holds the environment's lock. A lock that cannot be read is + * reported as no start in progress: a status read must not fail over it. + */ +async function startHoldsLock(env: string): Promise { + try { + return await wakeLockHeld(env); + } catch (err) { + console.log(JSON.stringify({ phase: 'lock-read', environment: env, error: errorName(err) })); + return false; + } +} + /** GET — report one environment's state without side effects. */ async function status(env: string): Promise { const eip = await findEnvEip(env); const baseUrl = eip ? baseUrlFor(eip.publicIp, ENGINE_PORT) : ''; const instance = await findManagedInstance(TAG_KEY, TAG_VALUE, envFilter(env)); + const starting = await startHoldsLock(env); if (!instance || instance.state !== 'running') { // No SSM call on this branch: reaching the daemon needs a running box, so // a stopped environment reports no activity rather than a made-up one. + // While a start holds the lock, an instance that is absent, stopped or + // still coming up is `starting`: the instance lookup alone cannot show a + // start another client began, because it lags a launch. + const instanceState = instance?.state ?? (eip ? 'stopped' : 'undeployed'); + const comingUp = !instance || instance.state === 'stopped' || instance.state === 'pending'; return jsonResponse(200, { - state: instance?.state ?? (eip ? 'stopped' : 'undeployed'), + state: starting && comingUp ? 'starting' : instanceState, environment: env, healthy: false, base_url: baseUrl, + start_in_progress: starting, }); } if (!(await isSsmAgentOnline(instance.instanceId))) { @@ -160,6 +192,7 @@ async function status(env: string): Promise { environment: env, healthy: false, base_url: baseUrl, + start_in_progress: starting, }); } // Concurrently, not in sequence: status is what you type repeatedly while @@ -176,6 +209,7 @@ async function status(env: string): Promise { environment: env, healthy, base_url: baseUrl, + start_in_progress: starting, ...deploy, ...activity, }; @@ -335,7 +369,64 @@ async function readDeployFacts(env: string): Promise<{ } } -/** POST — launch the environment's instance if needed and block until serving. */ +/** + * What a held lock, and a lock that could not be taken, are both answered + * with: a retryable "starting", so the CLI and gateway ask again and find the + * holder's instance. Nothing has been looked up or launched. + */ +function startInProgress(env: string, message: string): LambdaFunctionURLResult { + return jsonResponse( + 503, + { state: 'starting', environment: env, message, retry_after_seconds: LOCK_RETRY_SECONDS }, + { 'retry-after': String(LOCK_RETRY_SECONDS) }, + ); +} + +/** + * The reply for a start whose instance was stopped or terminated under it. + * Retryable, so a client still waiting asks again and re-wakes the instance + * (or launches a fresh one); ending here is what releases the lock. + */ +function stoppedUnderStart(env: string, instanceId: string, state: string): LambdaFunctionURLResult { + console.log(JSON.stringify({ phase: 'stopped-under-start', environment: env, instanceId, state })); + return jsonResponse( + 503, + { + state, + environment: env, + instance_id: instanceId, + message: `the instance was ${state === 'stopping' ? 'stopped' : state} while starting`, + retry_after_seconds: LOCK_RETRY_SECONDS, + }, + { 'retry-after': String(LOCK_RETRY_SECONDS) }, + ); +} + +/** + * Whether the instance this start is waiting on has been stopped or terminated, + * as the reply to end the start with, or null to keep waiting. A lookup that + * fails (the instance not yet visible, a transient error) keeps waiting: the + * poll loops have their own deadline. + */ +async function goneUnderStart(env: string, instanceId: string): Promise { + let state: string; + try { + state = (await getInstance(instanceId)).state; + } catch { + return null; + } + return GONE_STATES.has(state) ? stoppedUnderStart(env, instanceId, state) : null; +} + +/** + * POST — launch the environment's instance if needed and block until serving. + * + * Only one start works on an environment at a time: the instance lookup that + * decides whether to launch is eventually consistent, so two starts close + * together could each launch one. The checks that only read (the deploy config, + * the environment's address and security group) come first and take no lock, + * so an environment that cannot start says so without contending for one. + */ async function wake( env: string, context: Context, @@ -376,6 +467,43 @@ async function wake( } const baseUrl = baseUrlFor(eip.publicIp, ENGINE_PORT); + // The lock lasts as long as this invocation can run, so a start killed + // before it could release the lock blocks the environment for no longer. + // A lock that cannot be taken for a reason other than being held is + // answered the same way as a held one: nothing is launched unlocked. + const owner = context.awsRequestId ?? randomUUID(); + const expiresAt = new Date(Date.now() + context.getRemainingTimeInMillis() + LOCK_MARGIN_MS); + try { + if (!(await acquireWakeLock(env, owner, expiresAt))) { + console.log(JSON.stringify({ phase: 'lock-held', environment: env })); + return startInProgress(env, `another start for environment ${JSON.stringify(env)} is in progress`); + } + } catch (err) { + console.log(JSON.stringify({ phase: 'lock-error', environment: env, error: errorName(err) })); + return startInProgress(env, 'could not take the start lock; retrying'); + } + try { + return await wakeLocked(env, deadline, retainUntil, deployConfig, eip, securityGroupId, baseUrl); + } finally { + try { + await releaseWakeLock(env, owner); + } catch (err) { + // The reply is already decided; the lock expires on its own. + console.log(JSON.stringify({ phase: 'lock-release', environment: env, error: errorName(err) })); + } + } +} + +/** The body of a start, run while this invocation holds the environment's lock. */ +async function wakeLocked( + env: string, + deadline: number, + retainUntil: string | null, + deployConfig: DeployConfig, + eip: EnvEip, + securityGroupId: string, + baseUrl: string, +): Promise { // Weights first: a launch against an incomplete prefix would boot the engine // on it, and a re-wake would keep an old one alive on nothing. const gate = await seedingGate(env, deployConfig); @@ -386,6 +514,7 @@ async function wake( const existing = await findManagedInstance(TAG_KEY, TAG_VALUE, envFilter(env)); let instanceId: string; let startIssued = false; + let launchedFresh = false; if (existing) { // Idempotent: this environment's instance already exists (up, coming up, // or stopped — the sweep stops idle ones now, so this is the normal @@ -413,6 +542,7 @@ async function wake( return launched.error; } instanceId = launched.instanceId; + launchedFresh = true; } // Phase 1: EC2 state -> running (then pin the env's EIP so its URL resolves). @@ -441,6 +571,12 @@ async function wake( { 'retry-after': '300' }, ); } + if ((state === 'stopped' || state === 'stopping') && (startIssued || launchedFresh)) { + // This start launched the instance or issued its start command, and it + // has since been stopped: it is not coming up, so end the wake and let + // the lock go rather than poll it until the deadline. + return stoppedUnderStart(env, instanceId, state); + } if (state === 'stopped' && !startIssued) { // A stop raced us between discovery and now; issue the re-wake here so // one wake owns at most one start call. @@ -468,6 +604,10 @@ async function wake( // Phase 2: SSM agent online (registers 30-60 s after boot). while (Date.now() < deadline) { + const gone = await goneUnderStart(env, instanceId); + if (gone) { + return gone; + } if (await isSsmAgentOnline(instanceId)) { break; } @@ -481,6 +621,10 @@ async function wake( // is no one to take a start, so this wait converts a lost start (and a // full-deadline health timeout) into a short pause. while (Date.now() < deadline) { + const gone = await goneUnderStart(env, instanceId); + if (gone) { + return gone; + } if (await daemonAnswers(instanceId)) { break; } @@ -532,6 +676,10 @@ async function wake( } return ready(env, baseUrl, retainUntil); } + const gone = await goneUnderStart(env, instanceId); + if (gone) { + return gone; + } await sleep(HEALTH_POLL_MS); } diff --git a/remote/lib/llm-stack.ts b/remote/lib/llm-stack.ts index 8dff4d76..b7b2cfec 100644 --- a/remote/lib/llm-stack.ts +++ b/remote/lib/llm-stack.ts @@ -435,6 +435,17 @@ export class LlmStack extends cdk.Stack { }), ); startFn.addToRolePolicy(envParamsStatement); + // The per-environment start lock is created with the grants above and + // removed here; deleting is limited to the lock parameter itself so the + // start cannot delete an environment's deploy-config. + startFn.addToRolePolicy( + new iam.PolicyStatement({ + actions: ['ssm:DeleteParameter'], + resources: [ + `arn:${cdk.Aws.PARTITION}:ssm:${cdk.Aws.REGION}:${cdk.Aws.ACCOUNT_ID}:parameter/cloud-vm-llm/*/wake-lock`, + ], + }), + ); startFn.addToRolePolicy( new iam.PolicyStatement({ actions: ['secretsmanager:GetSecretValue'], diff --git a/remote/test/stack.test.ts b/remote/test/stack.test.ts index 12b4226e..9b7d27b1 100644 --- a/remote/test/stack.test.ts +++ b/remote/test/stack.test.ts @@ -438,6 +438,32 @@ describe('LlmStack (control plane)', () => { expect(JSON.stringify(terminate!.Condition)).toContain('cloud-vm-llm'); }); + it('lets only the start Lambda delete SSM parameters, and only the start lock', () => { + const policies = template.findResources('AWS::IAM::Policy') as Record; + const withDelete = Object.values(policies).filter((p) => + (p.Properties.PolicyDocument.Statement as Statement[]).some((s) => + [s.Action].flat().includes('ssm:DeleteParameter'), + ), + ); + expect(withDelete).toHaveLength(1); + + const statements = (withDelete[0].Properties.PolicyDocument.Statement as Statement[]).filter((s) => + [s.Action].flat().includes('ssm:DeleteParameter'), + ); + expect(statements).toHaveLength(1); + expect([statements[0].Action].flat()).toEqual(['ssm:DeleteParameter']); + const resources = JSON.stringify(statements[0].Resource); + expect(resources).toContain('parameter/cloud-vm-llm/*/wake-lock'); + expect(resources).not.toContain('"*"'); + + // The role it is attached to belongs to the start Lambda, found by the + // AMI settings only that Lambda is given. + const roleId = (withDelete[0].Properties.Roles as { Ref: string }[])[0].Ref; + const fns = template.findResources('AWS::Lambda::Function') as Record; + const owner = Object.values(fns).find((f) => f.Properties.Role['Fn::GetAtt'][0] === roleId); + expect(owner?.Properties.Environment.Variables.AMI_ROLE_TAG_KEY).toBeDefined(); + }); + it('passes the AMI role tag, weights bucket and subnet list to the start Lambda', () => { const fns = template.findResources('AWS::Lambda::Function'); const start = Object.values(fns).find((f) => diff --git a/remote/test/start-boot-failure.test.ts b/remote/test/start-boot-failure.test.ts index 890eb0ba..955869dc 100644 --- a/remote/test/start-boot-failure.test.ts +++ b/remote/test/start-boot-failure.test.ts @@ -61,6 +61,13 @@ vi.mock('../lambda/shared/seed', () => ({ weightsPresent: async () => true, })); +// These tests are about the wake itself, so the environment's lock is always free. +vi.mock('../lambda/shared/wake-lock', () => ({ + acquireWakeLock: async () => true, + releaseWakeLock: async () => undefined, + wakeLockHeld: async () => false, +})); + let handler: (event: LambdaFunctionURLEvent, context: Context) => Promise; let buildInferenceUserData: (env: string, cfg: DeployConfig) => string; let BOOT_FAILED_MARKER: string; diff --git a/remote/test/start-launch.test.ts b/remote/test/start-launch.test.ts index db7d1859..17e87eb1 100644 --- a/remote/test/start-launch.test.ts +++ b/remote/test/start-launch.test.ts @@ -66,6 +66,13 @@ vi.mock('../lambda/shared/seed', () => ({ weightsPresent: async () => true, })); +// These tests are about the launch itself, so the environment's lock is always free. +vi.mock('../lambda/shared/wake-lock', () => ({ + acquireWakeLock: async () => true, + releaseWakeLock: async () => undefined, + wakeLockHeld: async () => false, +})); + let handler: (event: LambdaFunctionURLEvent, context: Context) => Promise; beforeAll(async () => { diff --git a/remote/test/start-rewake.test.ts b/remote/test/start-rewake.test.ts index 5dcbd005..beb3284e 100644 --- a/remote/test/start-rewake.test.ts +++ b/remote/test/start-rewake.test.ts @@ -67,6 +67,13 @@ vi.mock('../lambda/shared/seed', () => ({ weightsPresent: async () => true, })); +// These tests are about the wake itself, so the environment's lock is always free. +vi.mock('../lambda/shared/wake-lock', () => ({ + acquireWakeLock: async () => true, + releaseWakeLock: async () => undefined, + wakeLockHeld: async () => false, +})); + let handler: (event: LambdaFunctionURLEvent, context: Context) => Promise; beforeAll(async () => { diff --git a/remote/test/start-seeding.test.ts b/remote/test/start-seeding.test.ts index df205355..d1d3a391 100644 --- a/remote/test/start-seeding.test.ts +++ b/remote/test/start-seeding.test.ts @@ -77,6 +77,13 @@ vi.mock('../lambda/shared/seed/launch', () => ({ launchSeedInstance: (...args: unknown[]) => launchSeedInstance(...args), })); +// These tests are about the weights gate, so the environment's lock is always free. +vi.mock('../lambda/shared/wake-lock', () => ({ + acquireWakeLock: async () => true, + releaseWakeLock: async () => undefined, + wakeLockHeld: async () => false, +})); + let handler: (event: LambdaFunctionURLEvent, context: Context) => Promise; beforeAll(async () => { diff --git a/remote/test/start-status.test.ts b/remote/test/start-status.test.ts index a7ce6f77..8544edaa 100644 --- a/remote/test/start-status.test.ts +++ b/remote/test/start-status.test.ts @@ -51,6 +51,14 @@ vi.mock('../lambda/shared/seed', () => ({ weightsPresent: async () => true, })); +// These tests are about what the status reports from the instance and daemon, +// so no start holds the environment's lock. +vi.mock('../lambda/shared/wake-lock', () => ({ + acquireWakeLock: async () => true, + releaseWakeLock: async () => undefined, + wakeLockHeld: async () => false, +})); + let handler: (event: LambdaFunctionURLEvent, context: Context) => Promise; beforeAll(async () => { diff --git a/remote/test/wake-lock-start.test.ts b/remote/test/wake-lock-start.test.ts new file mode 100644 index 00000000..ca22f8df --- /dev/null +++ b/remote/test/wake-lock-start.test.ts @@ -0,0 +1,546 @@ +import { beforeAll, beforeEach, describe, expect, it, vi } from 'vitest'; +import type { Context, LambdaFunctionURLEvent, LambdaFunctionURLResult } from 'aws-lambda'; +import { DAEMON_STATUS_CMD } from '../lambda/shared/daemon'; +import { wakeLockParam } from '../lambda/shared/wake-lock'; + +// Starts for one environment take a lock, so two starts close together launch +// one instance, and a start whose instance is stopped under it ends promptly +// and lets the lock go. The lock is the real one, over an in-memory stand-in +// for SSM; every other AWS call is stubbed. + +const LAMBDA_ENV = { + TAG_KEY: 'cloud-vm-llm:managed', + TAG_VALUE: 'true', + ENGINE_PORT: '8000', + AMI_ROLE_TAG_KEY: 'cloud-vm-llm:role', + AMI_ROLE_TAG_VALUE: 'runtime-ami', + AMI_RUNNER_TAG_KEY: 'cloud-vm-llm:runner', + INSTANCE_TYPE: 'g6e.xlarge', + SUBNET_IDS: 'subnet-test', + INSTANCE_PROFILE_ARN: 'arn:aws:iam::0:instance-profile/test', + WEIGHTS_BUCKET: 'test-bucket', + MAX_CONCURRENT_SEEDS: '2', + AWS_REGION: 'us-east-1', + BOOT_LOG_GROUP: '/test/boot', + LLAMACPP_LOG_GROUP: '/test/llamacpp', + VLLM_LOG_GROUP: '/test/vllm', + IDLE_THRESHOLD_MINUTES: '15', + GRACE_PERIOD_MINUTES: '10', + MAX_RUNTIME_MINUTES: '240', + STOP_RETENTION_MINUTES: '60', + MAX_SEED_MINUTES: '60', + SEED_STALL_MINUTES: '10', +}; + +const store = new Map(); +const ssmFailures = new Map(); + +const findManagedInstance = vi.fn(); +const getInstance = vi.fn(); +const startEngineDaemon = vi.fn(); +const startInstance = vi.fn(); +const runInstance = vi.fn(); +const findLatestAmi = vi.fn(); +const tagInstance = vi.fn(); +const associateEip = vi.fn(); +const isSsmAgentOnline = vi.fn(); +const runShellCommand = vi.fn(); +const readDeployConfig = vi.fn(); +const stopEngineDaemon = vi.fn(); +const stopInstance = vi.fn(); +const findEnvEip = vi.fn(); +const findEnvSecurityGroup = vi.fn(); +const readEnvApiKey = vi.fn(); + +function awsError(name: string): Error { + return Object.assign(new Error(name), { name }); +} + +vi.mock('@aws-sdk/client-ssm', async (importOriginal) => ({ + ...(await importOriginal()), + SSMClient: class { + async send(cmd: { constructor: { name: string }; input: { Name: string; Value?: string; Overwrite?: boolean } }) { + const kind = cmd.constructor.name; + const failure = ssmFailures.get(kind); + if (failure) { + throw failure; + } + const { Name, Value, Overwrite } = cmd.input; + if (kind === 'PutParameterCommand') { + if (!Overwrite && store.has(Name)) { + throw awsError('ParameterAlreadyExists'); + } + store.set(Name, Value as string); + return {}; + } + if (kind === 'GetParameterCommand') { + if (!store.has(Name)) { + throw awsError('ParameterNotFound'); + } + return { Parameter: { Value: store.get(Name) } }; + } + if (kind === 'DeleteParameterCommand') { + if (!store.has(Name)) { + throw awsError('ParameterNotFound'); + } + store.delete(Name); + return {}; + } + throw new Error(`unexpected ${kind}`); + } + }, +})); + +vi.mock('../lambda/shared/aws', async (importOriginal) => ({ + ...(await importOriginal()), + findManagedInstance: (...args: unknown[]) => findManagedInstance(...args), + getInstance: (...args: unknown[]) => getInstance(...args), + startEngineDaemon: (...args: unknown[]) => startEngineDaemon(...args), + startInstance: (...args: unknown[]) => startInstance(...args), + runInstance: (...args: unknown[]) => runInstance(...args), + findLatestAmi: (...args: unknown[]) => findLatestAmi(...args), + tagInstance: (...args: unknown[]) => tagInstance(...args), + associateEip: (...args: unknown[]) => associateEip(...args), + isSsmAgentOnline: (...args: unknown[]) => isSsmAgentOnline(...args), + runShellCommand: (...args: unknown[]) => runShellCommand(...args), + readDeployConfig: (...args: unknown[]) => readDeployConfig(...args), + stopEngineDaemon: (...args: unknown[]) => stopEngineDaemon(...args), + stopInstance: (...args: unknown[]) => stopInstance(...args), + // The wake sleeps between polls; these tests step through the polls instead. + sleep: async () => undefined, +})); + +vi.mock('../lambda/shared/environments', async (importOriginal) => ({ + ...(await importOriginal()), + findEnvEip: (...args: unknown[]) => findEnvEip(...args), + findEnvSecurityGroup: (...args: unknown[]) => findEnvSecurityGroup(...args), + readEnvApiKey: (...args: unknown[]) => readEnvApiKey(...args), +})); + +vi.mock('../lambda/shared/seed', () => ({ weightsPresent: async () => true })); + +type Handler = (event: LambdaFunctionURLEvent, context: Context) => Promise; +let start: Handler; +let stop: Handler; + +beforeAll(async () => { + Object.assign(process.env, LAMBDA_ENV); + ({ handler: start } = await import('../lambda/start/index')); + stop = (await import('../lambda/stop/index')).handler as unknown as Handler; +}); + +function event(env: string, method = 'POST', query: Record = {}): LambdaFunctionURLEvent { + return { + queryStringParameters: { env, ...query }, + requestContext: { http: { method } }, + } as unknown as LambdaFunctionURLEvent; +} + +function contextOf(requestId: string): Context { + return { awsRequestId: requestId, getRemainingTimeInMillis: () => 600_000 } as unknown as Context; +} + +function structured(result: LambdaFunctionURLResult): { statusCode: number; body: string } { + return result as { statusCode: number; body: string }; +} + +function bodyOf(result: LambdaFunctionURLResult): Record { + return JSON.parse(structured(result).body); +} + +const FUTURE = new Date(Date.now() + 3_600_000); +const PAST = new Date(Date.now() - 3_600_000); + +function lockOf(owner: string, expiresAt: Date): string { + return JSON.stringify({ owner, expiresAt: expiresAt.toISOString() }); +} + +const wait = (ms: number) => new Promise((resolve) => setTimeout(resolve, ms)); + +beforeEach(() => { + vi.clearAllMocks(); + store.clear(); + ssmFailures.clear(); + readDeployConfig.mockResolvedValue({ + runner: 'llamacpp', + modelId: 'org/model', + quant: 'Q4_K_M', + weightsPrefix: 'llamacpp/org/model/Q4_K_M', + contextSize: 32768, + servedModelName: 'friendly', + serveArgs: [], + companions: {}, + }); + findEnvEip.mockResolvedValue({ publicIp: '198.51.100.7', allocationId: 'eipalloc-test' }); + findEnvSecurityGroup.mockResolvedValue('sg-test'); + readEnvApiKey.mockResolvedValue('sk-test'); + isSsmAgentOnline.mockResolvedValue(true); + runShellCommand.mockImplementation((_id: string, command: string) => + command === DAEMON_STATUS_CMD + ? Promise.resolve({ status: 'Success', stdout: JSON.stringify({ state: 'stopped' }) }) + : Promise.resolve({ status: 'Success', stdout: '200' }), + ); + startEngineDaemon.mockResolvedValue(true); + findLatestAmi.mockResolvedValue({ imageId: 'ami-test1', rootVolumeSizeGb: 80 }); + // The lookup the lock exists to protect: it never sees an instance another + // start is in the middle of launching. + findManagedInstance.mockResolvedValue(null); + getInstance.mockResolvedValue({ instanceId: 'i-new', state: 'running', launchTime: new Date() }); + runInstance.mockImplementation(async () => { + await wait(30); + return 'i-new'; + }); + stopEngineDaemon.mockResolvedValue(undefined); +}); + +describe('two starts for one environment', () => { + it('launch one instance; the other is told a start is in progress', async () => { + const [a, b] = await Promise.all([ + start(event('dev'), contextOf('req-1')), + start(event('dev'), contextOf('req-2')), + ]); + + expect(runInstance).toHaveBeenCalledTimes(1); + const results = [a, b].map((r) => structured(r).statusCode).sort(); + expect(results).toEqual([200, 503]); + const refused = bodyOf([a, b].find((r) => structured(r).statusCode === 503)!); + expect(refused.state).toBe('starting'); + expect(refused.message).toContain('another start'); + expect(refused.retry_after_seconds).toBe(15); + }); + + it('leave the refused start having looked up and launched nothing', async () => { + store.set(wakeLockParam('dev'), lockOf('req-holder', FUTURE)); + + const result = await start(event('dev'), contextOf('req-2')); + + expect(structured(result).statusCode).toBe(503); + expect(bodyOf(result).message).toContain('another start'); + expect(findManagedInstance).not.toHaveBeenCalled(); + expect(runInstance).not.toHaveBeenCalled(); + expect(startInstance).not.toHaveBeenCalled(); + // The holder's lock is untouched. + expect(JSON.parse(store.get(wakeLockParam('dev'))!).owner).toBe('req-holder'); + }); + + it('are followed by a retry that finds the first start’s instance and launches nothing more', async () => { + await start(event('dev'), contextOf('req-1')); + expect(runInstance).toHaveBeenCalledTimes(1); + + findManagedInstance.mockResolvedValue({ instanceId: 'i-new', state: 'running' }); + const retry = await start(event('dev'), contextOf('req-2')); + + expect(structured(retry).statusCode).toBe(200); + expect(bodyOf(retry).state).toBe('ready'); + expect(runInstance).toHaveBeenCalledTimes(1); + }); +}); + +describe('different environments', () => { + it('do not wait for each other', async () => { + const [a, b] = await Promise.all([ + start(event('a'), contextOf('req-1')), + start(event('b'), contextOf('req-2')), + ]); + expect(structured(a).statusCode).toBe(200); + expect(structured(b).statusCode).toBe(200); + expect(runInstance).toHaveBeenCalledTimes(2); + }); +}); + +describe('the lock is released on every ending', () => { + it('after ready', async () => { + await start(event('dev'), contextOf('req-1')); + expect(store.has(wakeLockParam('dev'))).toBe(false); + }); + + it('after a retryable refusal (no AMI baked)', async () => { + findLatestAmi.mockResolvedValue(null); + const result = await start(event('dev'), contextOf('req-1')); + expect(bodyOf(result).state).toBe('no-ami'); + expect(store.has(wakeLockParam('dev'))).toBe(false); + }); + + it('after an instance goes terminal', async () => { + getInstance.mockResolvedValue({ instanceId: 'i-new', state: 'terminated' }); + const result = await start(event('dev'), contextOf('req-1')); + expect(structured(result).statusCode).toBe(503); + expect(store.has(wakeLockParam('dev'))).toBe(false); + }); + + it('after an unexpected error', async () => { + runInstance.mockRejectedValue(new Error('boom')); + await expect(start(event('dev'), contextOf('req-1'))).rejects.toThrow('boom'); + expect(store.has(wakeLockParam('dev'))).toBe(false); + }); + + it('so the next start is not refused', async () => { + runInstance.mockRejectedValueOnce(new Error('boom')); + await expect(start(event('dev'), contextOf('req-1'))).rejects.toThrow('boom'); + const next = await start(event('dev'), contextOf('req-2')); + expect(structured(next).statusCode).toBe(200); + }); + + it('only when this start still holds it', async () => { + findLatestAmi.mockImplementation(async () => { + // Another start took the lock over while this one was working. + store.set(wakeLockParam('dev'), lockOf('req-other', FUTURE)); + return null; + }); + await start(event('dev'), contextOf('req-1')); + expect(JSON.parse(store.get(wakeLockParam('dev'))!).owner).toBe('req-other'); + }); +}); + +describe('an abandoned lock', () => { + it('is taken over once it has expired', async () => { + store.set(wakeLockParam('dev'), lockOf('req-killed', PAST)); + const result = await start(event('dev'), contextOf('req-2')); + expect(structured(result).statusCode).toBe(200); + expect(runInstance).toHaveBeenCalledTimes(1); + expect(store.has(wakeLockParam('dev'))).toBe(false); + }); + + it('is not taken over while it is still valid', async () => { + store.set(wakeLockParam('dev'), lockOf('req-running', FUTURE)); + const result = await start(event('dev'), contextOf('req-2')); + expect(structured(result).statusCode).toBe(503); + expect(runInstance).not.toHaveBeenCalled(); + }); + + it('is given an expiry just past this invocation’s own time limit', async () => { + let written: { owner: string; expiresAt: string } | undefined; + findLatestAmi.mockImplementation(async () => { + written = JSON.parse(store.get(wakeLockParam('dev'))!); + return null; + }); + const before = Date.now(); + await start(event('dev'), contextOf('req-1')); + const margin = Date.parse(written!.expiresAt) - before - 600_000; + expect(written!.owner).toBe('req-1'); + expect(margin).toBeGreaterThanOrEqual(29_000); + expect(margin).toBeLessThan(35_000); + }); +}); + +describe('a lock that cannot be taken', () => { + it('stops the start before it launches anything, and says to retry', async () => { + ssmFailures.set('PutParameterCommand', awsError('AccessDeniedException')); + const result = await start(event('dev'), contextOf('req-1')); + expect(structured(result).statusCode).toBe(503); + expect(bodyOf(result).state).toBe('starting'); + expect(findManagedInstance).not.toHaveBeenCalled(); + expect(runInstance).not.toHaveBeenCalled(); + }); +}); + +describe('status while a start is in progress', () => { + const read = async () => bodyOf(await start(event('dev', 'GET'), contextOf('req-read'))); + + it('shows a second client the start another client began', async () => { + store.set(wakeLockParam('dev'), lockOf('req-holder', FUTURE)); + // The launch is not visible to the lookup yet, which is the window the flag covers. + findManagedInstance.mockResolvedValue(null); + + const body = await read(); + + expect(body.state).toBe('starting'); + expect(body.start_in_progress).toBe(true); + }); + + it.each(['stopped', 'pending'])('reports %s as starting while a start holds the lock', async (state) => { + store.set(wakeLockParam('dev'), lockOf('req-holder', FUTURE)); + findManagedInstance.mockResolvedValue({ instanceId: 'i-x', state }); + + const body = await read(); + + expect(body.state).toBe('starting'); + expect(body.start_in_progress).toBe(true); + }); + + it('keeps a running instance running, and flags the start while its model loads', async () => { + store.set(wakeLockParam('dev'), lockOf('req-holder', FUTURE)); + findManagedInstance.mockResolvedValue({ instanceId: 'i-run', state: 'running' }); + isSsmAgentOnline.mockResolvedValue(false); + + const body = await read(); + + expect(body.state).toBe('running'); + expect(body.healthy).toBe(false); + expect(body.start_in_progress).toBe(true); + }); + + it('flags a start on a fully reporting running instance too', async () => { + store.set(wakeLockParam('dev'), lockOf('req-holder', FUTURE)); + findManagedInstance.mockResolvedValue({ instanceId: 'i-run', state: 'running' }); + + const body = await read(); + + expect(body.state).toBe('running'); + expect(body.start_in_progress).toBe(true); + }); + + it('reports no start when none holds the lock', async () => { + findManagedInstance.mockResolvedValue(null); + const body = await read(); + expect(body.state).toBe('stopped'); + expect(body.start_in_progress).toBe(false); + }); + + it('ignores a lock that has expired', async () => { + store.set(wakeLockParam('dev'), lockOf('req-killed', PAST)); + findManagedInstance.mockResolvedValue(null); + + const body = await read(); + + expect(body.state).toBe('stopped'); + expect(body.start_in_progress).toBe(false); + }); + + it('still answers when the lock cannot be read', async () => { + ssmFailures.set('GetParameterCommand', awsError('ThrottlingException')); + findManagedInstance.mockResolvedValue(null); + + const body = await read(); + + expect(body.state).toBe('stopped'); + expect(body.start_in_progress).toBe(false); + }); + + it('shows the start from the moment it takes the lock, and not after it ends', async () => { + let during: Record | undefined; + findLatestAmi.mockImplementation(async () => { + during = bodyOf(await start(event('dev', 'GET'), contextOf('req-read'))); + return null; + }); + + await start(event('dev'), contextOf('req-1')); + const after = await read(); + + expect(during?.state).toBe('starting'); + expect(during?.start_in_progress).toBe(true); + expect(after.start_in_progress).toBe(false); + }); + + it('does not show a start that is only waiting to retry for capacity', async () => { + findLatestAmi.mockResolvedValue({ imageId: 'ami-test1', rootVolumeSizeGb: 80 }); + runInstance.mockRejectedValue(awsError('InsufficientInstanceCapacity')); + + const refused = await start(event('dev'), contextOf('req-1')); + expect(bodyOf(refused).state).toBe('no-capacity'); + + const body = await read(); + expect(body.state).toBe('stopped'); + expect(body.start_in_progress).toBe(false); + }); + + it('leaves the lock as the holder had it', async () => { + store.set(wakeLockParam('dev'), lockOf('req-holder', FUTURE)); + findManagedInstance.mockResolvedValue(null); + + await read(); + + expect(JSON.parse(store.get(wakeLockParam('dev'))!).owner).toBe('req-holder'); + }); +}); + +describe('reads and stops', () => { + it('a status read does not take, or wait for, the lock', async () => { + store.set(wakeLockParam('dev'), lockOf('req-holder', FUTURE)); + const result = await start(event('dev', 'GET'), contextOf('req-2')); + expect(structured(result).statusCode).toBe(200); + expect(JSON.parse(store.get(wakeLockParam('dev'))!).owner).toBe('req-holder'); + }); + + it('a pause is carried out while a start holds the lock', async () => { + store.set(wakeLockParam('dev'), lockOf('req-holder', FUTURE)); + findManagedInstance.mockResolvedValue({ instanceId: 'i-run', state: 'running' }); + const result = await stop(event('dev', 'POST', { action: 'pause' }), contextOf('req-stop')); + expect(structured(result).statusCode).toBe(200); + expect(stopInstance).toHaveBeenCalledWith('i-run'); + }); +}); + +describe('a stop during a start', () => { + it('after the start command ends the start, and releases the lock', async () => { + findManagedInstance.mockResolvedValue({ instanceId: 'i-off', state: 'stopped' }); + getInstance.mockResolvedValue({ instanceId: 'i-off', state: 'stopped' }); + + const result = await start(event('dev'), contextOf('req-1')); + + expect(startInstance).toHaveBeenCalledWith('i-off'); + expect(structured(result).statusCode).toBe(503); + const body = bodyOf(result); + expect(body.state).toBe('stopped'); + expect(body.message).toContain('while starting'); + expect(body.retry_after_seconds).toBe(15); + expect(startEngineDaemon).not.toHaveBeenCalled(); + expect(store.has(wakeLockParam('dev'))).toBe(false); + }); + + it('right after a fresh launch ends the start, rather than re-waking the instance it just launched', async () => { + getInstance.mockResolvedValue({ instanceId: 'i-new', state: 'stopping' }); + + const result = await start(event('dev'), contextOf('req-1')); + + expect(structured(result).statusCode).toBe(503); + expect(bodyOf(result).message).toContain('stopped while starting'); + expect(startInstance).not.toHaveBeenCalled(); + expect(store.has(wakeLockParam('dev'))).toBe(false); + }); + + it('while waiting for the agent ends the start', async () => { + isSsmAgentOnline.mockResolvedValue(false); + getInstance + .mockResolvedValueOnce({ instanceId: 'i-new', state: 'running' }) // first loop + .mockResolvedValueOnce({ instanceId: 'i-new', state: 'stopped' }); // agent loop + + const result = await start(event('dev'), contextOf('req-1')); + + expect(structured(result).statusCode).toBe(503); + expect(bodyOf(result).state).toBe('stopped'); + expect(store.has(wakeLockParam('dev'))).toBe(false); + }); + + it('terminated while the engine loads ends the start', async () => { + runShellCommand.mockImplementation((_id: string, command: string) => + command === DAEMON_STATUS_CMD + ? Promise.resolve({ status: 'Success', stdout: JSON.stringify({ state: 'stopped' }) }) + : Promise.resolve({ status: 'Success', stdout: '503' }), + ); + getInstance + .mockResolvedValueOnce({ instanceId: 'i-new', state: 'running' }) // first loop + .mockResolvedValueOnce({ instanceId: 'i-new', state: 'running' }) // agent loop + .mockResolvedValueOnce({ instanceId: 'i-new', state: 'running' }) // daemon loop + .mockResolvedValueOnce({ instanceId: 'i-new', state: 'terminated' }); // health loop + + const result = await start(event('dev'), contextOf('req-1')); + + expect(structured(result).statusCode).toBe(503); + expect(bodyOf(result).state).toBe('terminated'); + expect(store.has(wakeLockParam('dev'))).toBe(false); + }); + + it('is not mistaken for a stop when the instance cannot be looked up yet', async () => { + getInstance + .mockResolvedValueOnce({ instanceId: 'i-new', state: 'running' }) // first loop + .mockRejectedValueOnce(new Error('not visible yet')) // agent loop + .mockResolvedValue({ instanceId: 'i-new', state: 'running' }); + + const result = await start(event('dev'), contextOf('req-1')); + + expect(structured(result).statusCode).toBe(200); + }); + + it('still re-wakes an instance that was stopped before this start issued its command', async () => { + findManagedInstance.mockResolvedValue({ instanceId: 'i-run', state: 'running' }); + getInstance + .mockResolvedValueOnce({ instanceId: 'i-run', state: 'stopped' }) // first loop: the stop raced discovery + .mockResolvedValue({ instanceId: 'i-run', state: 'running' }); + + const result = await start(event('dev'), contextOf('req-1')); + + expect(startInstance).toHaveBeenCalledWith('i-run'); + expect(structured(result).statusCode).toBe(200); + }); +}); diff --git a/remote/test/wake-lock.test.ts b/remote/test/wake-lock.test.ts new file mode 100644 index 00000000..29a84302 --- /dev/null +++ b/remote/test/wake-lock.test.ts @@ -0,0 +1,184 @@ +import { beforeEach, describe, expect, it, vi } from 'vitest'; +import { acquireWakeLock, releaseWakeLock, wakeLockHeld, wakeLockParam } from '../lambda/shared/wake-lock'; + +// The lock is an SSM parameter created only if absent. The SSM client is +// replaced by an in-memory store with the same create/read/delete behaviour, +// so the tests cover the lock's rules, not SSM. + +const store = new Map(); +const failures = new Map(); + +function awsError(name: string): Error { + return Object.assign(new Error(name), { name }); +} + +const sendHook = vi.fn(); + +vi.mock('@aws-sdk/client-ssm', async (importOriginal) => ({ + ...(await importOriginal()), + SSMClient: class { + async send(cmd: { constructor: { name: string }; input: { Name: string; Value?: string; Overwrite?: boolean } }) { + sendHook(cmd); + const kind = cmd.constructor.name; + const failure = failures.get(kind); + if (failure) { + throw failure; + } + const { Name, Value, Overwrite } = cmd.input; + if (kind === 'PutParameterCommand') { + if (!Overwrite && store.has(Name)) { + throw awsError('ParameterAlreadyExists'); + } + store.set(Name, Value as string); + return {}; + } + if (kind === 'GetParameterCommand') { + if (!store.has(Name)) { + throw awsError('ParameterNotFound'); + } + return { Parameter: { Value: store.get(Name) } }; + } + if (kind === 'DeleteParameterCommand') { + if (!store.has(Name)) { + throw awsError('ParameterNotFound'); + } + store.delete(Name); + return {}; + } + throw new Error(`unexpected ${kind}`); + } + }, +})); + +const NOW = new Date('2026-10-05T12:00:00Z'); +const FUTURE = new Date('2026-10-05T12:15:00Z'); +const PAST = new Date('2026-10-05T11:00:00Z'); + +function lockOf(owner: string, expiresAt: Date): string { + return JSON.stringify({ owner, expiresAt: expiresAt.toISOString() }); +} + +beforeEach(() => { + store.clear(); + failures.clear(); + sendHook.mockReset(); +}); + +describe('acquireWakeLock', () => { + it('takes a free lock and records the owner and expiry', async () => { + expect(await acquireWakeLock('dev', 'req-1', FUTURE, NOW)).toBe(true); + expect(JSON.parse(store.get(wakeLockParam('dev'))!)).toEqual({ + owner: 'req-1', + expiresAt: FUTURE.toISOString(), + }); + }); + + it('names the lock parameter under the environment', () => { + expect(wakeLockParam('dev')).toBe('/cloud-vm-llm/dev/wake-lock'); + }); + + it('refuses a lock another start holds, and leaves it as it was', async () => { + store.set(wakeLockParam('dev'), lockOf('req-1', FUTURE)); + expect(await acquireWakeLock('dev', 'req-2', FUTURE, NOW)).toBe(false); + expect(JSON.parse(store.get(wakeLockParam('dev'))!).owner).toBe('req-1'); + }); + + it('lets only one of two simultaneous starts take a free lock', async () => { + const results = await Promise.all([ + acquireWakeLock('dev', 'req-1', FUTURE, NOW), + acquireWakeLock('dev', 'req-2', FUTURE, NOW), + ]); + expect(results.filter(Boolean)).toHaveLength(1); + }); + + it('takes over a lock that has expired', async () => { + store.set(wakeLockParam('dev'), lockOf('req-old', PAST)); + expect(await acquireWakeLock('dev', 'req-2', FUTURE, NOW)).toBe(true); + expect(JSON.parse(store.get(wakeLockParam('dev'))!).owner).toBe('req-2'); + }); + + it.each([['not json'], ['{"owner":1}'], ['']])('takes over a value that is not a lock: %j', async (raw) => { + store.set(wakeLockParam('dev'), raw); + expect(await acquireWakeLock('dev', 'req-2', FUTURE, NOW)).toBe(true); + expect(JSON.parse(store.get(wakeLockParam('dev'))!).owner).toBe('req-2'); + }); + + it('does not delete a lock that changed between its two reads', async () => { + store.set(wakeLockParam('dev'), lockOf('req-old', PAST)); + // After the first read finds the expired lock, another start replaces it. + let reads = 0; + sendHook.mockImplementation((cmd: { constructor: { name: string } }) => { + if (cmd.constructor.name === 'GetParameterCommand' && ++reads === 1) { + queueMicrotask(() => store.set(wakeLockParam('dev'), lockOf('req-new', FUTURE))); + } + }); + expect(await acquireWakeLock('dev', 'req-2', FUTURE, NOW)).toBe(false); + expect(JSON.parse(store.get(wakeLockParam('dev'))!).owner).toBe('req-new'); + }); + + it('keeps locks of different environments apart', async () => { + expect(await acquireWakeLock('a', 'req-1', FUTURE, NOW)).toBe(true); + expect(await acquireWakeLock('b', 'req-2', FUTURE, NOW)).toBe(true); + }); + + it('throws on a failure other than "already exists", rather than proceeding unlocked', async () => { + failures.set('PutParameterCommand', awsError('AccessDeniedException')); + await expect(acquireWakeLock('dev', 'req-1', FUTURE, NOW)).rejects.toThrow('AccessDeniedException'); + expect(store.size).toBe(0); + }); + + it('throws when the held lock cannot be read', async () => { + store.set(wakeLockParam('dev'), lockOf('req-1', FUTURE)); + failures.set('GetParameterCommand', awsError('ThrottlingException')); + await expect(acquireWakeLock('dev', 'req-2', FUTURE, NOW)).rejects.toThrow('ThrottlingException'); + }); +}); + +describe('wakeLockHeld', () => { + it('is true for a valid lock and false once it has expired', async () => { + store.set(wakeLockParam('dev'), lockOf('req-1', FUTURE)); + expect(await wakeLockHeld('dev', NOW)).toBe(true); + expect(await wakeLockHeld('dev', new Date(FUTURE.getTime() + 1))).toBe(false); + }); + + it('is false with no lock, and for a value that is not a lock', async () => { + expect(await wakeLockHeld('dev', NOW)).toBe(false); + store.set(wakeLockParam('dev'), 'not json'); + expect(await wakeLockHeld('dev', NOW)).toBe(false); + }); + + it('writes nothing', async () => { + store.set(wakeLockParam('dev'), lockOf('req-1', PAST)); + await wakeLockHeld('dev', NOW); + expect(JSON.parse(store.get(wakeLockParam('dev'))!).owner).toBe('req-1'); + }); + + it('throws when the lock cannot be read, leaving the caller to decide', async () => { + failures.set('GetParameterCommand', awsError('ThrottlingException')); + await expect(wakeLockHeld('dev', NOW)).rejects.toThrow('ThrottlingException'); + }); +}); + +describe('releaseWakeLock', () => { + it('removes the lock its owner holds', async () => { + await acquireWakeLock('dev', 'req-1', FUTURE, NOW); + await releaseWakeLock('dev', 'req-1'); + expect(store.has(wakeLockParam('dev'))).toBe(false); + }); + + it('leaves a lock another start holds', async () => { + store.set(wakeLockParam('dev'), lockOf('req-2', FUTURE)); + await releaseWakeLock('dev', 'req-1'); + expect(JSON.parse(store.get(wakeLockParam('dev'))!).owner).toBe('req-2'); + }); + + it('does nothing when there is no lock', async () => { + await expect(releaseWakeLock('dev', 'req-1')).resolves.toBeUndefined(); + }); + + it('lets the next start take the lock after a release', async () => { + await acquireWakeLock('dev', 'req-1', FUTURE, NOW); + await releaseWakeLock('dev', 'req-1'); + expect(await acquireWakeLock('dev', 'req-2', FUTURE, NOW)).toBe(true); + }); +});