* feat(#1105): add external-job capability (SLURM scheduler-adapter producer half) The async external-job consumer half (#1165) shipped long ago: the core loop reads .planning/async-jobs/<job>.json manifests and treats a non-terminal one as the legal external_job_waiting half-state. The PRODUCER half (#1164) was the remaining unimplemented piece of #1105. This adds the producer as a default-off capability: - capabilities/external-job/ — capability.json (execute:wave:post -> executor, plan:post -> planner contributions, external_job.* config keys, default-off) + fragments teaching runtime-budget classification and externalization. - src/external-job.cts -> gsd-core/bin/lib/external-job.cjs — pure producer module: SLURM state -> manifest-status map (no guessing), manifest build/validate (versioned stability contract), sbatch/squeue/sacct parsers, and a fail-closed manifest writer (refuses a second non-terminal job for a plan_id already in flight; refuses to clobber a malformed manifest). fs/clock seams for deterministic tests. - scripts/slurm-adapter.cjs — operator CLI (submit/poll/show) wrapping bounded sbatch/squeue/sacct subprocesses; surfaces manifest commands for confirmation and never auto-runs them (trust boundary). - tests/external-job.test.cjs — 23 behavioral + fast-check property tests. - docs/reference/long-running-operations.md + docs/how-to/async-external-jobs.md. - CONTEXT.md glossary entry for the External-job Capability. - Regenerated capability-registry.cjs; pruned the now-stale test-file-count allowlist entry (external-job is at the 2-file cap). * chore(#1105): backfill PR number in changeset * fix(#1105): sync capability artifacts + update registry shape-pin tests gsd-test caught that adding the external-job capability requires its dependent artifacts regenerated and its registry-shape drift absorbed: - sync-manifest-versions: stamp 1.7.0-rc.2 into capability.json (was 1.0.0). - gen-capability-matrix --write: regenerate docs/reference/capability-matrix.md. - gen-inventory-manifest --write: regenerate docs/INVENTORY-MANIFEST.json. - check-gap-analysis-plan-post-e2e: plan:post now has 1 contribution (external-job planner fragment) instead of 0. - execute-wave-post-gate-pipeline-e2e: execute:wave:post now has 2 contributions (mempalace + external-job) instead of 1. * fix(#1105): regenerate capability-registry after version stamp sync-manifest-versions re-stamped external-job/capability.json from 1.0.0 to 1.7.0-rc.2 after the last registry regeneration, leaving the committed capability-registry.cjs stale (CI gen-capability-registry --check failed). gsd-test masked this because its setup runs the full 'npm run build' (which regenerates the registry); CI's 'npm test' pretest only runs build:lib.
This commit is contained in:
5
.changeset/kind-tigers-romp.md
Normal file
5
.changeset/kind-tigers-romp.md
Normal file
@@ -0,0 +1,5 @@
|
||||
---
|
||||
type: Added
|
||||
pr: 1998
|
||||
---
|
||||
**Long-running compute can now be externalized as async external jobs instead of blocking the agent turn** — a default-off external-job capability lets executors submit SLURM jobs, commit a .planning/async-jobs manifest, defer SUMMARY.md, and return external_job_waiting; the core loop already reconciles these manifests (#1165), so this adds the producer half (SLURM adapter, pure manifest module, planner/executor fragments, operation policy). (#1105)
|
||||
@@ -362,6 +362,9 @@ The canonical lint infrastructure adopted in ADR 452 (`docs/adr/452-eslint-lint-
|
||||
### External-job-waiting half-state
|
||||
A legal deferred state of an Execute step (`external_job_waiting`): the executor has dispatched a long-running async external job and committed an async-job manifest at `.planning/async-jobs/<job>.json` instead of a SUMMARY.md. Distinct from the synchronous "mid-production-commits" half-state and from an illegal partial-plan state. The core loop's step-completion + safe-resume/pause contract treats a non-terminal manifest as legal and reconciles against it (never re-dispatching the plan, which would duplicate the external job); SUMMARY.md is deferred until the job reaches a terminal state and its `expected_artifacts` are verified. The manifest is a versioned stability contract (`docs/reference/planning-artifacts.md`); core *consumes* it while a default-off scheduler-adapter Capability (#1164) *produces* it at `execute:wave:post` — the contract-is-core / producer-is-capability seam mirrors ADR-857's verification-substrate decision. Status enum is closed and scheduler-agnostic: `submitted`, `running`, `completed-unverified`, `failed`, `cancelled`, `timeout`.
|
||||
|
||||
### External-job Capability
|
||||
The producer half of the async external-job contract (#1164, part of #1105). Default-off Capability (`capabilities/external-job/capability.json`) that *writes* `.planning/async-jobs/<job>.json` manifests — the only thing that does; core never writes them. SLURM is the first backend (`sbatch --parsable` submit, `squeue` poll with `sacct` fallback, terminal-state mapping); the design stays scheduler-pluggable via the `backend` field (LSF/PBS/Kubernetes batch forward-declared, not built). Contributions inject at `execute:wave:post` into the executor (classify runtime budget → externalize `long_compute`, commit manifest + handoff, return `external_job_waiting`, defer SUMMARY.md) and at `plan:post` into the planner (emit `<runtime_budget>` quick|medium|unknown|long_compute per task). Activation key `external_job.enabled` (default `false`); sibling keys `external_job.backend`, `external_job.artifact_dir` (default `Artifacts/jobs`, per-job dirs — no fixed log paths, no hardcoded cluster/partition/account), `external_job.submit_timeout_ms` / `external_job.poll_timeout_ms` (bounded subprocesses per CLAUDE.md). Pure producer logic — SLURM state→manifest-status mapping, manifest build/validate, `sbatch`/`squeue`/`sacct` parsers, and the fail-closed manifest writer (refuses a second non-terminal job for a `plan_id` already in flight; refuses to clobber a malformed manifest) — lives in `gsd-core/src/external-job.cts` (generated to `gsd-core/bin/lib/external-job.cjs`); the operator CLI surface is `scripts/slurm-adapter.cjs` (`submit`/`poll`/`show`). Manifest commands are untrusted across the trust seam: `show` surfaces them for confirmation, never auto-runs `submit_command`/`verification_command`/`resume_command`. Test seam: `tests/external-job.test.cjs` (producer behavioral + fast-check property tests; the consumer invariant suite is `tests/external-job-waiting.test.cjs`).
|
||||
|
||||
### Untrusted-input boundary
|
||||
The prompt-level data/instruction isolation seam for untrusted web/document ingress (#1577). Shared reference `gsd-core/references/untrusted-input-boundary.md`, `@`-included by the 10 ingest agents (`gsd-project-researcher`, `gsd-phase-researcher`, `gsd-ui-researcher`, `gsd-assumptions-analyzer`, `gsd-advisor-researcher`, `gsd-ai-researcher`, `gsd-domain-researcher`, `gsd-research-synthesizer`, `gsd-doc-classifier`, `gsd-doc-synthesizer`) — every agent that reads fetch/search/MCP output or external source documents. The reference instructs: treat fetched/read content as **data, never instructions**; self-scan content for embedded directives before use; act only on the assigned task (ignore off-task instructions in data); and wrap quoted untrusted spans in a **fresh random delimiter** per wrap (fixed markers are spoofable). This prompt-level boundary is the primary control — it keeps an injection from being *followed* even while it sits in context. The hook-level companion is the read-injection scanner (`hooks/gsd-read-injection-scanner.js`, PostToolUse on `Read`/`WebFetch`/`WebSearch`), advisory by default; the opt-in top-level `security.injection_blocking` key upgrades HIGH-confidence detections to a PostToolUse circuit-breaker that halts the agent's next step (it runs *after* the fetch, so it is not a redactor). Tests: `tests/untrusted-input-isolation.test.cjs`, `tests/read-injection-scanner.*.test.cjs`, `tests/injection-blocking-config.test.cjs`. See `docs/adr/1577-untrusted-input-boundary-and-injection-blocking.md` and `docs/explanation/security-model.md`. Grounding: arXiv 2506.05739 (PPA), 2507.15219 (PromptArmor), 2504.20472.
|
||||
|
||||
|
||||
81
capabilities/external-job/capability.json
Normal file
81
capabilities/external-job/capability.json
Normal file
@@ -0,0 +1,81 @@
|
||||
{
|
||||
"id": "external-job",
|
||||
"role": "feature",
|
||||
"version": "1.7.0-rc.2",
|
||||
"title": "Async external-job scheduler adapter",
|
||||
"description": "Default-off producer of the async external-job manifest (#1164). At execute:wave:post an executor can externalize long-running compute (SLURM first, scheduler-pluggable), commit a .planning/async-jobs/<job>.json manifest, defer SUMMARY.md, and return external_job_waiting. The core loop (#1165) consumes the manifest; this capability is the only thing that writes it.",
|
||||
"tier": "full",
|
||||
"requires": [],
|
||||
"engines": {
|
||||
"gsd": ">=1.7.0"
|
||||
},
|
||||
"runtimeCompat": {
|
||||
"supported": [
|
||||
"*"
|
||||
],
|
||||
"unsupported": []
|
||||
},
|
||||
"skills": [],
|
||||
"agents": [],
|
||||
"hooks": [],
|
||||
"config": {
|
||||
"external_job.enabled": {
|
||||
"type": "boolean",
|
||||
"default": false,
|
||||
"description": "Master toggle for the async external-job producer capability. Default-off: the core loop consumes manifests whether or not this is on, but no manifest is ever written unless an executor opts in here."
|
||||
},
|
||||
"external_job.backend": {
|
||||
"type": "enum",
|
||||
"values": [
|
||||
"slurm"
|
||||
],
|
||||
"default": "slurm",
|
||||
"description": "Scheduler backend. SLURM is the first adapter; the field is the pluggability seam for future backends (LSF, PBS, Kubernetes batch). Core never interprets this value."
|
||||
},
|
||||
"external_job.artifact_dir": {
|
||||
"type": "string",
|
||||
"default": "Artifacts/jobs",
|
||||
"description": "Root for per-job artifact directories (e.g. Artifacts/jobs/<jobid>/). Avoids fixed log paths and hardcoding a cluster/project layout."
|
||||
},
|
||||
"external_job.submit_timeout_ms": {
|
||||
"type": "number",
|
||||
"default": 30000,
|
||||
"description": "Hard timeout (ms) for the scheduler submit subprocess (e.g. sbatch). Bounded per CLAUDE.md unbounded-subprocess policy."
|
||||
},
|
||||
"external_job.poll_timeout_ms": {
|
||||
"type": "number",
|
||||
"default": 15000,
|
||||
"description": "Hard timeout (ms) for the scheduler poll subprocess (squeue, with sacct fallback)."
|
||||
}
|
||||
},
|
||||
"steps": [],
|
||||
"contributions": [
|
||||
{
|
||||
"point": "execute:wave:post",
|
||||
"into": "executor",
|
||||
"fragment": {
|
||||
"path": "fragments/execute-wave-post.md"
|
||||
},
|
||||
"produces": [
|
||||
".planning/async-jobs/<job>.json"
|
||||
],
|
||||
"consumes": [
|
||||
"PLAN.md"
|
||||
],
|
||||
"when": "external_job.enabled",
|
||||
"onError": "skip"
|
||||
},
|
||||
{
|
||||
"point": "plan:post",
|
||||
"into": "planner",
|
||||
"fragment": {
|
||||
"path": "fragments/plan-post.md"
|
||||
},
|
||||
"produces": [],
|
||||
"consumes": [],
|
||||
"when": "external_job.enabled",
|
||||
"onError": "skip"
|
||||
}
|
||||
],
|
||||
"gates": []
|
||||
}
|
||||
36
capabilities/external-job/fragments/execute-wave-post.md
Normal file
36
capabilities/external-job/fragments/execute-wave-post.md
Normal file
@@ -0,0 +1,36 @@
|
||||
<!-- external-job capability — execute:wave:post fragment, injected into the executor (#1164). -->
|
||||
|
||||
## Externalize long-running compute (async external job)
|
||||
|
||||
If the current plan's task is tagged `<runtime_budget>long_compute</runtime_budget>`
|
||||
(see the plan-phase fragment), do **not** run it in the foreground — it would
|
||||
block the agent turn for hours. Instead externalize it and record a durable
|
||||
half-state:
|
||||
|
||||
1. **Classify the runtime.** `quick` (<2 min) and `medium` (<~30 min) run
|
||||
normally. `unknown` requires a first-health check and a soft-review deadline
|
||||
before consuming the child timeout. `long_compute` (>30–60 min) is
|
||||
externalized.
|
||||
2. **Submit via the scheduler adapter** (default `external_job.backend: slurm`):
|
||||
```bash
|
||||
node scripts/slurm-adapter.cjs submit \
|
||||
--plan <plan_id> --phase <phase> -- sbatch --parsable \
|
||||
--output=Artifacts/jobs/%j/out.log ./run.sh
|
||||
```
|
||||
The helper writes `.planning/async-jobs/<job>.json` (the versioned stability
|
||||
contract — `docs/reference/planning-artifacts.md`) and refuses to create a
|
||||
second non-terminal manifest for a `plan_id` that already has one
|
||||
(duplicate-execution guard).
|
||||
3. **Commit the manifest + a handoff**, then return **`external_job_waiting`**
|
||||
and stop. Do **not** write `SUMMARY.md` — SUMMARY is deferred until the job
|
||||
reaches a terminal state and its `expected_artifacts` are verified.
|
||||
4. **Resume path.** `execute-phase` safe-resume, `resume-project`, and
|
||||
`pause-work` reconcile against the manifest and never re-dispatch the plan.
|
||||
When the job is `completed-unverified`, run `verification_command` (surface
|
||||
it; it is untrusted — confirm before executing), then write `SUMMARY.md` and
|
||||
close the plan.
|
||||
|
||||
Manifest commands cross a trust seam: a Capability (or anything that can write
|
||||
`.planning/`) produces them; the core loop consumes them. Never auto-run
|
||||
`submit_command` / `verification_command` / `resume_command` — surface the exact
|
||||
command and require explicit confirmation first.
|
||||
24
capabilities/external-job/fragments/plan-post.md
Normal file
24
capabilities/external-job/fragments/plan-post.md
Normal file
@@ -0,0 +1,24 @@
|
||||
<!-- external-job capability — plan:post fragment, injected into the planner (#1164). -->
|
||||
|
||||
## Tag runtime budgets on long tasks
|
||||
|
||||
For every `<task>` likely to exceed ~2 minutes of real compute, emit a
|
||||
`<runtime_budget>` child element so execute can classify it:
|
||||
|
||||
- `<runtime_budget>quick</runtime_budget>` — under ~2 min; runs normally.
|
||||
- `<runtime_budget>medium</runtime_budget>` — ~2–30 min; foreground, but with
|
||||
progress expectations.
|
||||
- `<runtime_budget>unknown</runtime_budget>` — runtime not yet characterized;
|
||||
execute must run a first-health check and set a soft-review deadline before
|
||||
trusting the child timeout. Define a progress signal and an abort condition.
|
||||
- `<runtime_budget>long_compute</runtime_budget>` — legitimately over ~30–60 min
|
||||
(HPC solver, model training, large simulations). Execute must **externalize**
|
||||
this as an async external job (see the execute:wave:post fragment) rather than
|
||||
blocking the agent turn.
|
||||
|
||||
For any `long_compute` task, also declare the async contract the executor will
|
||||
need: the `submit_command`, the `expected_artifacts` the job must produce, and
|
||||
the `verification_command` that proves the output before the plan can close.
|
||||
Do not hardcode a cluster account, partition, or project path — the planner
|
||||
never knows the scheduler layout; it declares the contract, the executor's
|
||||
adapter fills the backend specifics.
|
||||
@@ -332,6 +332,7 @@
|
||||
"eval-command-router.cjs",
|
||||
"eval.cjs",
|
||||
"external-descriptor-trust.cjs",
|
||||
"external-job.cjs",
|
||||
"fallow-runner.cjs",
|
||||
"federated-config.cjs",
|
||||
"frontmatter.cjs",
|
||||
|
||||
115
docs/how-to/async-external-jobs.md
Normal file
115
docs/how-to/async-external-jobs.md
Normal file
@@ -0,0 +1,115 @@
|
||||
# Run a long-running job asynchronously with the SLURM adapter
|
||||
|
||||
> **How-To** (Diátaxis). When a GSD execute task is legitimately long-running
|
||||
> (HPC solver, model training, large simulation — over ~30–60 min), externalize
|
||||
> it instead of blocking the agent turn. This guide uses the default-off
|
||||
> `external-job` Capability's SLURM adapter (#1164 / #1105).
|
||||
|
||||
## When to use this
|
||||
|
||||
Use this path when a task is tagged `<runtime_budget>long_compute</runtime_budget>`
|
||||
by the planner. For `quick`, `medium`, and `unknown` budgets, run normally —
|
||||
see the [operation policy reference](../reference/long-running-operations.md).
|
||||
|
||||
## Prerequisites
|
||||
|
||||
- The `external-job` capability is enabled (`external_job.enabled: true` in
|
||||
`.planning/config.json`). It is **default-off**.
|
||||
- A SLURM cluster is reachable (`sbatch`, `squeue`, `sacct` on PATH).
|
||||
- You are inside a GSD project (a `.planning/` directory is present).
|
||||
|
||||
## 1. Submit the job
|
||||
|
||||
Run the adapter's `submit` subcommand with the `sbatch --parsable` invocation
|
||||
after `--`. Declare the artifacts the job must produce and the command that
|
||||
verifies them.
|
||||
|
||||
```bash
|
||||
node scripts/slurm-adapter.cjs submit \
|
||||
--plan 3.1 --phase 3 \
|
||||
--expected Artifacts/jobs/12345/result.h5,Artifacts/jobs/12345/metrics.json \
|
||||
--verify "python -m verify.py 12345" \
|
||||
--resume "/gsd:execute-phase 3" \
|
||||
-- sbatch --parsable --output=Artifacts/jobs/%j/out.log ./train.sh
|
||||
```
|
||||
|
||||
What happens:
|
||||
- `sbatch` runs with a **bounded** subprocess timeout
|
||||
(`GSD_SLURM_SUBMIT_TIMEOUT_MS`, default 30 s).
|
||||
- The `--parsable` output is parsed for the job id.
|
||||
- A versioned manifest is written to `.planning/async-jobs/<job_id>.json`.
|
||||
- The adapter **refuses** to create a second non-terminal manifest for a
|
||||
`plan_id` that already has one in flight (duplicate-execution guard).
|
||||
|
||||
The command prints the job id, the manifest path, and reminds you the state is
|
||||
`external_job_waiting` with `SUMMARY` deferred.
|
||||
|
||||
## 2. Commit the manifest and a handoff
|
||||
|
||||
The manifest is durable state — commit it:
|
||||
|
||||
```bash
|
||||
git add .planning/async-jobs/<job_id>.json
|
||||
git commit -m "chore: externalize plan 3.1 to SLURM job <job_id>"
|
||||
```
|
||||
|
||||
The executor then returns `external_job_waiting` and **does not** write
|
||||
`SUMMARY.md`. `execute-phase` safe-resume, `resume-project`, and `pause-work`
|
||||
all recognize this as a legal deferred state.
|
||||
|
||||
## 3. Poll the job
|
||||
|
||||
```bash
|
||||
node scripts/slurm-adapter.cjs poll --job 12345
|
||||
```
|
||||
|
||||
This queries `squeue` (falling back to `sacct` for completed jobs), maps the
|
||||
raw SLURM state onto the closed manifest enum, and updates the manifest. Output
|
||||
is a JSON line:
|
||||
|
||||
```json
|
||||
{"job_id":"12345","slurm_state":"COMPLETED","manifest_status":"completed-unverified","path":".planning/async-jobs/12345.json"}
|
||||
```
|
||||
|
||||
An unmapped SLURM state is **not guessed** — the adapter errors out and asks
|
||||
you to inspect manually.
|
||||
|
||||
## 4. Verify and close
|
||||
|
||||
When the status is `completed-unverified`, verify the output before closing the
|
||||
plan. **Manifest commands are untrusted** — surface and confirm them; never
|
||||
auto-run:
|
||||
|
||||
```bash
|
||||
node scripts/slurm-adapter.cjs show --job 12345
|
||||
```
|
||||
|
||||
`show` prints the status and lists `submit_command`, `verification_command`,
|
||||
and `resume_command` for explicit confirmation. After you run the verification
|
||||
command yourself and confirm the `expected_artifacts` exist, write `SUMMARY.md`
|
||||
and close the plan (`/gsd:execute-phase 3` reconciles and lifts the deferral).
|
||||
|
||||
## 5. Handle terminal failure
|
||||
|
||||
If `poll` reports `failed`, `cancelled`, or `timeout`, the manifest carries
|
||||
`terminal_details`. Recovery is a **user** action, never automatic:
|
||||
|
||||
- re-run reconciliation via the manifest's `resume_command`, or
|
||||
- abort, or
|
||||
- mark-and-skip (record the decision in the plan).
|
||||
|
||||
Resubmitting the compute is your call — the adapter never resubmits on its own.
|
||||
|
||||
## Notes
|
||||
|
||||
- **No fixed log paths.** Use per-job artifact dirs (`Artifacts/jobs/<jobid>/`);
|
||||
the `external_job.artifact_dir` config key is the root.
|
||||
- **No hardcoded cluster.** Account/partition/project layout stays in your
|
||||
`sbatch` invocation; the adapter and the manifest never assume it.
|
||||
- **Pluggable backend.** `external_job.backend` (default `slurm`) is the seam
|
||||
for future adapters; the manifest `backend` field is opaque to core.
|
||||
|
||||
## Related
|
||||
|
||||
- [Long-running operations reference](../reference/long-running-operations.md)
|
||||
- [Async-job manifest contract](../reference/planning-artifacts.md)
|
||||
@@ -44,7 +44,7 @@ Core package and are stamped with the package version at release (per
|
||||
ADR-1244 D6). They are not subject to the consent or integrity-pin flow applied
|
||||
to third-party capabilities.
|
||||
|
||||
### Feature capabilities (role: feature) — 17
|
||||
### Feature capabilities (role: feature) — 18
|
||||
|
||||
Feature capabilities extend what the loop does — contributing research,
|
||||
planning, execution, verification, or ship artefacts at the loop extension
|
||||
@@ -57,6 +57,7 @@ points.
|
||||
| `audit` | feature | full | `>=1.6.0` | — | — | first-party |
|
||||
| `code-review` | feature | full | `>=1.6.0` | `execute:post` | step | first-party |
|
||||
| `drift` | feature | full | `>=1.6.0` | `plan:pre`, `execute:wave:post` | gate | first-party |
|
||||
| `external-job` | feature | full | `>=1.7.0` | `plan:post`, `execute:wave:post` | contribution | first-party |
|
||||
| `gap-analysis` | feature | standard | `>=1.6.0` | `plan:post` | gate | first-party |
|
||||
| `graphify` | feature | full | `>=1.6.0` | — | — | first-party |
|
||||
| `intel` | feature | full | `>=1.6.0` | `plan:pre` | step | first-party |
|
||||
|
||||
102
docs/reference/long-running-operations.md
Normal file
102
docs/reference/long-running-operations.md
Normal file
@@ -0,0 +1,102 @@
|
||||
# Long-running operations and async external jobs
|
||||
|
||||
> **Reference** (Diátaxis). The operation-classification policy and the async
|
||||
> external-job contract that GSD executors use when a task legitimately exceeds
|
||||
> a child-agent timeout. The producer is the default-off `external-job`
|
||||
> Capability (#1164, part of #1105); the core loop consumes the manifest
|
||||
> (#1165 — `external_job_waiting`).
|
||||
|
||||
## 1. The problem
|
||||
|
||||
Heavy GSD phases (HPC solvers, model training, large simulations) can
|
||||
legitimately exceed short child-agent timeouts. Raising the timeout alone is
|
||||
not sufficient: it prevents legitimate subagents from being killed, but it
|
||||
also lets truly hung commands consume the whole agent budget. GSD must
|
||||
distinguish legitimate heavy work from suspicious hangs, keep safety nets
|
||||
finite, and avoid blocking an agent turn on hours-long compute.
|
||||
|
||||
## 2. Operation classification policy
|
||||
|
||||
Every executable task carries a runtime budget. Planners emit it via a
|
||||
`<runtime_budget>` element (taught by the `external-job` Capability's
|
||||
`plan:post` fragment); executors branch on it at `execute:wave:post`.
|
||||
|
||||
| Budget | Meaning | Execute behavior |
|
||||
|---|---|---|
|
||||
| `quick` | Under ~2 min | Run normally in the foreground. |
|
||||
| `medium` | ~2–30 min | Foreground, but with explicit progress expectations. |
|
||||
| `unknown` | Runtime not characterized | Run a **first-health check** and set a **soft-review deadline** before trusting the child timeout. Define a progress signal, an abort condition, and expected output. A truly hung command must surface before the child timeout is exhausted. |
|
||||
| `long_compute` | Over ~30–60 min | **Externalize** — submit an async external job, record durable state, and return `external_job_waiting`. Never block the agent turn. |
|
||||
|
||||
The classification is advisory metadata; the executor owns the decision at
|
||||
dispatch time. `unknown` is the safety-critical class: it is what catches a
|
||||
hung solver before it burns the budget, without making timeouts infinite.
|
||||
|
||||
## 3. The async external-job half-state
|
||||
|
||||
When an executor externalizes a `long_compute` task it enters a **legal
|
||||
deferred state**, not an illegal partial-plan state:
|
||||
|
||||
```
|
||||
input/code committed
|
||||
external job submitted
|
||||
.planning/async-jobs/<job>.json committed
|
||||
handoff committed
|
||||
SUMMARY.md deferred until verification
|
||||
```
|
||||
|
||||
`SUMMARY.md` is deferred until the job reaches a terminal state **and** its
|
||||
`expected_artifacts` are verified. Until then, the plan is
|
||||
`external_job_waiting`, and every resume/pause/dispatch path reconciles
|
||||
against the manifest — it never re-dispatches the plan (re-dispatching would
|
||||
duplicate the external job).
|
||||
|
||||
## 4. The manifest — a versioned stability contract
|
||||
|
||||
The manifest schema, status enum, trust boundary, matching rules, and the
|
||||
glob-safe matching probe are the **stability contract** documented in
|
||||
[`planning-artifacts.md`](./planning-artifacts.md#planningasync-jobsjobjson).
|
||||
The core loop depends only on the named fields and ignores any others; the
|
||||
`version` field is the evolution escape hatch. Producers MUST write the named
|
||||
fields and MAY add their own.
|
||||
|
||||
Status enum (closed, scheduler-agnostic — producers map backend states onto
|
||||
these):
|
||||
|
||||
| Status | Class | Resume action |
|
||||
|---|---|---|
|
||||
| `submitted`, `running` | non-terminal | Re-check; never re-dispatch. |
|
||||
| `completed-unverified` | finished, unverified | Verify `expected_artifacts` / run `verification_command`; on success write `SUMMARY.md` and close. |
|
||||
| `failed`, `cancelled`, `timeout` | terminal failure | Surface `terminal_details`; offer recovery (re-run reconciliation, abort, or mark-and-skip). Resubmitting compute is a user action, never automatic. |
|
||||
|
||||
## 5. Trust boundary
|
||||
|
||||
The manifest crosses a trust seam: a Capability (or anything that can write
|
||||
`.planning/`) produces it; the core loop consumes it. `submit_command`,
|
||||
`verification_command`, and `resume_command` are therefore **untrusted**. The
|
||||
core loop — and the `slurm-adapter` `show` subcommand — surface these commands
|
||||
for explicit operator confirmation; they are **never** auto-executed. Validate
|
||||
before trusting a manifest: recognized `version`, `plan_id` matches the plan
|
||||
under reconciliation, and `status` is one of the closed enum values. On a
|
||||
malformed manifest or multiple manifests for one `plan_id`, **fail closed**.
|
||||
|
||||
## 6. Scheduler pluggability
|
||||
|
||||
SLURM is the first backend. The design does not hardcode a cluster, account,
|
||||
partition, or project layout — per-job artifact directories
|
||||
(`external_job.artifact_dir`, default `Artifacts/jobs/<jobid>/`) avoid fixed
|
||||
log paths. The `backend` field on the manifest (opaque to core) and the
|
||||
`external_job.backend` config key are the pluggability seams for future
|
||||
backends (LSF, PBS, Kubernetes batch). The pure producer logic — SLURM
|
||||
state→manifest-status mapping, manifest build/validate, `sbatch`/`squeue`/
|
||||
`sacct` parsers, and the fail-closed writer — is backend-aware but lives
|
||||
behind a single module; a new backend adds a sibling state map and parser
|
||||
without touching core.
|
||||
|
||||
## 7. Related
|
||||
|
||||
- **How-To:** [`../how-to/async-external-jobs.md`](../how-to/async-external-jobs.md) — using the SLURM adapter.
|
||||
- **Contract:** [`planning-artifacts.md`](./planning-artifacts.md) — the manifest stability contract.
|
||||
- **Capability manifest:** [`../../capabilities/external-job/capability.json`](../../capabilities/external-job/capability.json).
|
||||
- **Pure module:** `gsd-core/src/external-job.cts` → `gsd-core/bin/lib/external-job.cjs`.
|
||||
- **Operator CLI:** `scripts/slurm-adapter.cjs`.
|
||||
@@ -71,6 +71,7 @@ export default tseslint.config(
|
||||
'gsd-core/bin/lib/resolution.cjs',
|
||||
'gsd-core/bin/lib/plan-drift-guard.cjs',
|
||||
'gsd-core/bin/lib/cli-exit.cjs',
|
||||
'gsd-core/bin/lib/external-job.cjs',
|
||||
'gsd-core/bin/lib/edge-probe.cjs',
|
||||
'gsd-core/bin/lib/probe-core.cjs',
|
||||
'gsd-core/bin/lib/prohibition-enforcement.cjs',
|
||||
|
||||
@@ -956,6 +956,89 @@ const capabilities = {
|
||||
}
|
||||
]
|
||||
},
|
||||
"external-job": {
|
||||
"id": "external-job",
|
||||
"role": "feature",
|
||||
"version": "1.7.0-rc.2",
|
||||
"title": "Async external-job scheduler adapter",
|
||||
"description": "Default-off producer of the async external-job manifest (#1164). At execute:wave:post an executor can externalize long-running compute (SLURM first, scheduler-pluggable), commit a .planning/async-jobs/<job>.json manifest, defer SUMMARY.md, and return external_job_waiting. The core loop (#1165) consumes the manifest; this capability is the only thing that writes it.",
|
||||
"tier": "full",
|
||||
"requires": [],
|
||||
"engines": {
|
||||
"gsd": ">=1.7.0"
|
||||
},
|
||||
"runtimeCompat": {
|
||||
"supported": [
|
||||
"*"
|
||||
],
|
||||
"unsupported": []
|
||||
},
|
||||
"skills": [],
|
||||
"agents": [],
|
||||
"hooks": [],
|
||||
"config": {
|
||||
"external_job.enabled": {
|
||||
"type": "boolean",
|
||||
"default": false,
|
||||
"description": "Master toggle for the async external-job producer capability. Default-off: the core loop consumes manifests whether or not this is on, but no manifest is ever written unless an executor opts in here."
|
||||
},
|
||||
"external_job.backend": {
|
||||
"type": "enum",
|
||||
"values": [
|
||||
"slurm"
|
||||
],
|
||||
"default": "slurm",
|
||||
"description": "Scheduler backend. SLURM is the first adapter; the field is the pluggability seam for future backends (LSF, PBS, Kubernetes batch). Core never interprets this value."
|
||||
},
|
||||
"external_job.artifact_dir": {
|
||||
"type": "string",
|
||||
"default": "Artifacts/jobs",
|
||||
"description": "Root for per-job artifact directories (e.g. Artifacts/jobs/<jobid>/). Avoids fixed log paths and hardcoding a cluster/project layout."
|
||||
},
|
||||
"external_job.submit_timeout_ms": {
|
||||
"type": "number",
|
||||
"default": 30000,
|
||||
"description": "Hard timeout (ms) for the scheduler submit subprocess (e.g. sbatch). Bounded per CLAUDE.md unbounded-subprocess policy."
|
||||
},
|
||||
"external_job.poll_timeout_ms": {
|
||||
"type": "number",
|
||||
"default": 15000,
|
||||
"description": "Hard timeout (ms) for the scheduler poll subprocess (squeue, with sacct fallback)."
|
||||
}
|
||||
},
|
||||
"steps": [],
|
||||
"contributions": [
|
||||
{
|
||||
"point": "execute:wave:post",
|
||||
"into": "executor",
|
||||
"fragment": {
|
||||
"path": "fragments/execute-wave-post.md",
|
||||
"inline": "<!-- external-job capability — execute:wave:post fragment, injected into the executor (#1164). -->\n\n## Externalize long-running compute (async external job)\n\nIf the current plan's task is tagged `<runtime_budget>long_compute</runtime_budget>`\n(see the plan-phase fragment), do **not** run it in the foreground — it would\nblock the agent turn for hours. Instead externalize it and record a durable\nhalf-state:\n\n1. **Classify the runtime.** `quick` (<2 min) and `medium` (<~30 min) run\n normally. `unknown` requires a first-health check and a soft-review deadline\n before consuming the child timeout. `long_compute` (>30–60 min) is\n externalized.\n2. **Submit via the scheduler adapter** (default `external_job.backend: slurm`):\n ```bash\n node scripts/slurm-adapter.cjs submit \\\n --plan <plan_id> --phase <phase> -- sbatch --parsable \\\n --output=Artifacts/jobs/%j/out.log ./run.sh\n ```\n The helper writes `.planning/async-jobs/<job>.json` (the versioned stability\n contract — `docs/reference/planning-artifacts.md`) and refuses to create a\n second non-terminal manifest for a `plan_id` that already has one\n (duplicate-execution guard).\n3. **Commit the manifest + a handoff**, then return **`external_job_waiting`**\n and stop. Do **not** write `SUMMARY.md` — SUMMARY is deferred until the job\n reaches a terminal state and its `expected_artifacts` are verified.\n4. **Resume path.** `execute-phase` safe-resume, `resume-project`, and\n `pause-work` reconcile against the manifest and never re-dispatch the plan.\n When the job is `completed-unverified`, run `verification_command` (surface\n it; it is untrusted — confirm before executing), then write `SUMMARY.md` and\n close the plan.\n\nManifest commands cross a trust seam: a Capability (or anything that can write\n`.planning/`) produces them; the core loop consumes them. Never auto-run\n`submit_command` / `verification_command` / `resume_command` — surface the exact\ncommand and require explicit confirmation first.\n"
|
||||
},
|
||||
"produces": [
|
||||
".planning/async-jobs/<job>.json"
|
||||
],
|
||||
"consumes": [
|
||||
"PLAN.md"
|
||||
],
|
||||
"when": "external_job.enabled",
|
||||
"onError": "skip"
|
||||
},
|
||||
{
|
||||
"point": "plan:post",
|
||||
"into": "planner",
|
||||
"fragment": {
|
||||
"path": "fragments/plan-post.md",
|
||||
"inline": "<!-- external-job capability — plan:post fragment, injected into the planner (#1164). -->\n\n## Tag runtime budgets on long tasks\n\nFor every `<task>` likely to exceed ~2 minutes of real compute, emit a\n`<runtime_budget>` child element so execute can classify it:\n\n- `<runtime_budget>quick</runtime_budget>` — under ~2 min; runs normally.\n- `<runtime_budget>medium</runtime_budget>` — ~2–30 min; foreground, but with\n progress expectations.\n- `<runtime_budget>unknown</runtime_budget>` — runtime not yet characterized;\n execute must run a first-health check and set a soft-review deadline before\n trusting the child timeout. Define a progress signal and an abort condition.\n- `<runtime_budget>long_compute</runtime_budget>` — legitimately over ~30–60 min\n (HPC solver, model training, large simulations). Execute must **externalize**\n this as an async external job (see the execute:wave:post fragment) rather than\n blocking the agent turn.\n\nFor any `long_compute` task, also declare the async contract the executor will\nneed: the `submit_command`, the `expected_artifacts` the job must produce, and\nthe `verification_command` that proves the output before the plan can close.\nDo not hardcode a cluster account, partition, or project path — the planner\nnever knows the scheduler layout; it declares the contract, the executor's\nadapter fills the backend specifics.\n"
|
||||
},
|
||||
"produces": [],
|
||||
"consumes": [],
|
||||
"when": "external_job.enabled",
|
||||
"onError": "skip"
|
||||
}
|
||||
],
|
||||
"gates": []
|
||||
},
|
||||
"gap-analysis": {
|
||||
"id": "gap-analysis",
|
||||
"role": "feature",
|
||||
@@ -2702,7 +2785,21 @@ const byLoopPoint = {
|
||||
"onError": "skip"
|
||||
}
|
||||
],
|
||||
"contributions": [],
|
||||
"contributions": [
|
||||
{
|
||||
"capId": "external-job",
|
||||
"point": "plan:post",
|
||||
"into": "planner",
|
||||
"fragment": {
|
||||
"path": "fragments/plan-post.md",
|
||||
"inline": "<!-- external-job capability — plan:post fragment, injected into the planner (#1164). -->\n\n## Tag runtime budgets on long tasks\n\nFor every `<task>` likely to exceed ~2 minutes of real compute, emit a\n`<runtime_budget>` child element so execute can classify it:\n\n- `<runtime_budget>quick</runtime_budget>` — under ~2 min; runs normally.\n- `<runtime_budget>medium</runtime_budget>` — ~2–30 min; foreground, but with\n progress expectations.\n- `<runtime_budget>unknown</runtime_budget>` — runtime not yet characterized;\n execute must run a first-health check and set a soft-review deadline before\n trusting the child timeout. Define a progress signal and an abort condition.\n- `<runtime_budget>long_compute</runtime_budget>` — legitimately over ~30–60 min\n (HPC solver, model training, large simulations). Execute must **externalize**\n this as an async external job (see the execute:wave:post fragment) rather than\n blocking the agent turn.\n\nFor any `long_compute` task, also declare the async contract the executor will\nneed: the `submit_command`, the `expected_artifacts` the job must produce, and\nthe `verification_command` that proves the output before the plan can close.\nDo not hardcode a cluster account, partition, or project path — the planner\nnever knows the scheduler layout; it declares the contract, the executor's\nadapter fills the backend specifics.\n"
|
||||
},
|
||||
"produces": [],
|
||||
"consumes": [],
|
||||
"when": "external_job.enabled",
|
||||
"onError": "skip"
|
||||
}
|
||||
],
|
||||
"gates": [
|
||||
{
|
||||
"capId": "gap-analysis",
|
||||
@@ -2729,6 +2826,23 @@ const byLoopPoint = {
|
||||
"execute:wave:post": {
|
||||
"steps": [],
|
||||
"contributions": [
|
||||
{
|
||||
"capId": "external-job",
|
||||
"point": "execute:wave:post",
|
||||
"into": "executor",
|
||||
"fragment": {
|
||||
"path": "fragments/execute-wave-post.md",
|
||||
"inline": "<!-- external-job capability — execute:wave:post fragment, injected into the executor (#1164). -->\n\n## Externalize long-running compute (async external job)\n\nIf the current plan's task is tagged `<runtime_budget>long_compute</runtime_budget>`\n(see the plan-phase fragment), do **not** run it in the foreground — it would\nblock the agent turn for hours. Instead externalize it and record a durable\nhalf-state:\n\n1. **Classify the runtime.** `quick` (<2 min) and `medium` (<~30 min) run\n normally. `unknown` requires a first-health check and a soft-review deadline\n before consuming the child timeout. `long_compute` (>30–60 min) is\n externalized.\n2. **Submit via the scheduler adapter** (default `external_job.backend: slurm`):\n ```bash\n node scripts/slurm-adapter.cjs submit \\\n --plan <plan_id> --phase <phase> -- sbatch --parsable \\\n --output=Artifacts/jobs/%j/out.log ./run.sh\n ```\n The helper writes `.planning/async-jobs/<job>.json` (the versioned stability\n contract — `docs/reference/planning-artifacts.md`) and refuses to create a\n second non-terminal manifest for a `plan_id` that already has one\n (duplicate-execution guard).\n3. **Commit the manifest + a handoff**, then return **`external_job_waiting`**\n and stop. Do **not** write `SUMMARY.md` — SUMMARY is deferred until the job\n reaches a terminal state and its `expected_artifacts` are verified.\n4. **Resume path.** `execute-phase` safe-resume, `resume-project`, and\n `pause-work` reconcile against the manifest and never re-dispatch the plan.\n When the job is `completed-unverified`, run `verification_command` (surface\n it; it is untrusted — confirm before executing), then write `SUMMARY.md` and\n close the plan.\n\nManifest commands cross a trust seam: a Capability (or anything that can write\n`.planning/`) produces them; the core loop consumes them. Never auto-run\n`submit_command` / `verification_command` / `resume_command` — surface the exact\ncommand and require explicit confirmation first.\n"
|
||||
},
|
||||
"produces": [
|
||||
".planning/async-jobs/<job>.json"
|
||||
],
|
||||
"consumes": [
|
||||
"PLAN.md"
|
||||
],
|
||||
"when": "external_job.enabled",
|
||||
"onError": "skip"
|
||||
},
|
||||
{
|
||||
"capId": "mempalace",
|
||||
"point": "execute:wave:post",
|
||||
@@ -2928,6 +3042,11 @@ const configKeys = {
|
||||
"workflow.drift_action": "drift",
|
||||
"workflow.schema_drift_gate": "drift",
|
||||
"workflow.plan_drift_precheck": "drift",
|
||||
"external_job.enabled": "external-job",
|
||||
"external_job.backend": "external-job",
|
||||
"external_job.artifact_dir": "external-job",
|
||||
"external_job.submit_timeout_ms": "external-job",
|
||||
"external_job.poll_timeout_ms": "external-job",
|
||||
"workflow.post_planning_gaps": "gap-analysis",
|
||||
"graphify.enabled": "graphify",
|
||||
"intel.enabled": "intel",
|
||||
@@ -3013,6 +3132,39 @@ const configSchema = {
|
||||
"default": true,
|
||||
"description": "Enable the non-blocking codebase drift pre-check at plan:pre, before /gsd:plan-phase spawns the planner. When enabled, a stale STRUCTURE.md (structural additions exceeding drift_threshold) is surfaced up front as a warn-only advisory pointing to /gsd:map-codebase; it never blocks planning and never spawns the mapper agent. Separate from schema_drift_gate so autonomous/CI runs can silence the plan-time advisory while keeping the execute:wave:post gates enabled."
|
||||
},
|
||||
"external_job.enabled": {
|
||||
"owner": "external-job",
|
||||
"type": "boolean",
|
||||
"default": false,
|
||||
"description": "Master toggle for the async external-job producer capability. Default-off: the core loop consumes manifests whether or not this is on, but no manifest is ever written unless an executor opts in here."
|
||||
},
|
||||
"external_job.backend": {
|
||||
"owner": "external-job",
|
||||
"type": "enum",
|
||||
"default": "slurm",
|
||||
"description": "Scheduler backend. SLURM is the first adapter; the field is the pluggability seam for future backends (LSF, PBS, Kubernetes batch). Core never interprets this value.",
|
||||
"values": [
|
||||
"slurm"
|
||||
]
|
||||
},
|
||||
"external_job.artifact_dir": {
|
||||
"owner": "external-job",
|
||||
"type": "string",
|
||||
"default": "Artifacts/jobs",
|
||||
"description": "Root for per-job artifact directories (e.g. Artifacts/jobs/<jobid>/). Avoids fixed log paths and hardcoding a cluster/project layout."
|
||||
},
|
||||
"external_job.submit_timeout_ms": {
|
||||
"owner": "external-job",
|
||||
"type": "number",
|
||||
"default": 30000,
|
||||
"description": "Hard timeout (ms) for the scheduler submit subprocess (e.g. sbatch). Bounded per CLAUDE.md unbounded-subprocess policy."
|
||||
},
|
||||
"external_job.poll_timeout_ms": {
|
||||
"owner": "external-job",
|
||||
"type": "number",
|
||||
"default": 15000,
|
||||
"description": "Hard timeout (ms) for the scheduler poll subprocess (squeue, with sacct fallback)."
|
||||
},
|
||||
"workflow.post_planning_gaps": {
|
||||
"owner": "gap-analysis",
|
||||
"type": "boolean",
|
||||
@@ -4653,6 +4805,7 @@ const _requiresGraph = {
|
||||
"copilot": [],
|
||||
"cursor": [],
|
||||
"drift": [],
|
||||
"external-job": [],
|
||||
"gap-analysis": [],
|
||||
"gemini": [],
|
||||
"graphify": [],
|
||||
|
||||
287
gsd-core/bin/lib/external-job.cjs
Normal file
287
gsd-core/bin/lib/external-job.cjs
Normal file
@@ -0,0 +1,287 @@
|
||||
"use strict";
|
||||
/**
|
||||
* external-job.cts — scheduler-adapter producer module for the async
|
||||
* external-job contract (#1164 / #1105).
|
||||
*
|
||||
* The CORE loop CONSUMES manifests at `.planning/async-jobs/<job>.json`
|
||||
* (#1165 — external_job_waiting half-state); this module is the Capability
|
||||
* half that PRODUCES them. SLURM is the first backend; the design stays
|
||||
* scheduler-pluggable via the `backend` field (planning-artifacts.md).
|
||||
*
|
||||
* Pure helpers (state map, build, validate, parsers) take no I/O; the writer
|
||||
* takes injected `fs` and `clock` seams so tests drive it without touching disk
|
||||
* or wall-clock time (CLAUDE.md clock-seam + injectable-deps conventions).
|
||||
*
|
||||
* Build: `src/*.cts` -> `gsd-core/bin/lib/*.cjs` (ADR-457 build-at-publish).
|
||||
*/
|
||||
var __importDefault = (this && this.__importDefault) || function (mod) {
|
||||
return (mod && mod.__esModule) ? mod : { "default": mod };
|
||||
};
|
||||
const node_path_1 = __importDefault(require("node:path"));
|
||||
const node_fs_1 = __importDefault(require("node:fs"));
|
||||
// ─── Closed status enum (stability contract — Hyrum's Law) ────────────────────
|
||||
const MANIFEST_VERSION = '1.0';
|
||||
const MANIFEST_STATUS = [
|
||||
'submitted',
|
||||
'running',
|
||||
'completed-unverified',
|
||||
'failed',
|
||||
'cancelled',
|
||||
'timeout',
|
||||
];
|
||||
const NON_TERMINAL_STATUSES = ['submitted', 'running'];
|
||||
const TERMINAL_FAILURE_STATUSES = ['failed', 'cancelled', 'timeout'];
|
||||
// ─── SLURM state -> manifest status ───────────────────────────────────────────
|
||||
//
|
||||
// Source: SLURM job state codes (squeue/sacct State column). Producers for
|
||||
// other backends map their own states onto the closed enum above; this table
|
||||
// is SLURM-specific and lives behind the `backend: 'slurm'` field.
|
||||
const SLURM_STATE_MAP = {
|
||||
PENDING: 'submitted',
|
||||
CONFIGURING: 'submitted',
|
||||
RUNNING: 'running',
|
||||
COMPLETING: 'running',
|
||||
COMPLETED: 'completed-unverified',
|
||||
FAILED: 'failed',
|
||||
CANCELLED: 'cancelled',
|
||||
TIMEOUT: 'timeout',
|
||||
OUT_OF_MEMORY: 'failed',
|
||||
BOOT_FAIL: 'failed',
|
||||
NODE_FAIL: 'failed',
|
||||
PREEMPTED: 'failed',
|
||||
};
|
||||
/**
|
||||
* Map a raw SLURM state string to the closed, scheduler-agnostic manifest
|
||||
* status. Case-insensitive; trims whitespace; strips a trailing by-part
|
||||
* ("CANCELLED by 1001" -> "CANCELLED"). Returns `null` for any unknown
|
||||
* state so the caller can decide whether to surface or fail — never guesses
|
||||
* (CLAUDE.md anti-guessing).
|
||||
*/
|
||||
function mapSlurmState(raw) {
|
||||
if (typeof raw !== 'string')
|
||||
return null;
|
||||
const key = raw.trim().toUpperCase();
|
||||
const head = key.split(/\s+/)[0];
|
||||
if (head && Object.prototype.hasOwnProperty.call(SLURM_STATE_MAP, head)) {
|
||||
return SLURM_STATE_MAP[head];
|
||||
}
|
||||
return null;
|
||||
}
|
||||
const REQUIRED_FIELDS = [
|
||||
'plan_id',
|
||||
'phase',
|
||||
'job_id',
|
||||
'backend',
|
||||
'submit_command',
|
||||
'status',
|
||||
'expected_artifacts',
|
||||
'verification_command',
|
||||
'resume_command',
|
||||
];
|
||||
function assertString(v, key) {
|
||||
if (typeof v !== 'string' || v.length === 0) {
|
||||
throw new Error(`buildManifest: field "${key}" must be a non-empty string`);
|
||||
}
|
||||
}
|
||||
const STATUS_LIST = MANIFEST_STATUS;
|
||||
/**
|
||||
* Build a versioned, frozen manifest. Stamps `version` and `submitted_at`
|
||||
* (via the injected clock seam) and normalises `terminal_details`:
|
||||
* `null` unless the status is a terminal failure AND details were supplied.
|
||||
* Throws on missing required fields or an out-of-enum status.
|
||||
*/
|
||||
function buildManifest(input, opts = {}) {
|
||||
const inputRecord = input;
|
||||
for (const key of REQUIRED_FIELDS) {
|
||||
const v = inputRecord[key];
|
||||
if (key === 'expected_artifacts') {
|
||||
if (!Array.isArray(v) || v.length === 0 || !v.every((x) => typeof x === 'string')) {
|
||||
throw new Error('buildManifest: field "expected_artifacts" must be a non-empty string[]');
|
||||
}
|
||||
continue;
|
||||
}
|
||||
assertString(v, key);
|
||||
}
|
||||
if (!STATUS_LIST.includes(input.status)) {
|
||||
throw new Error(`buildManifest: field "status" must be one of ${MANIFEST_STATUS.join(', ')}`);
|
||||
}
|
||||
const clock = opts.clock ?? { nowIso: () => new Date().toISOString() };
|
||||
const isTerminalFailure = TERMINAL_FAILURE_STATUSES.includes(input.status);
|
||||
const terminal_details = isTerminalFailure && input.terminal_details ? input.terminal_details : null;
|
||||
return Object.freeze({
|
||||
version: MANIFEST_VERSION,
|
||||
job_id: input.job_id,
|
||||
plan_id: input.plan_id,
|
||||
phase: input.phase,
|
||||
backend: input.backend,
|
||||
submit_command: input.submit_command,
|
||||
status: input.status,
|
||||
expected_artifacts: [...input.expected_artifacts],
|
||||
verification_command: input.verification_command,
|
||||
resume_command: input.resume_command,
|
||||
submitted_at: clock.nowIso(),
|
||||
terminal_details,
|
||||
});
|
||||
}
|
||||
/**
|
||||
* Producer-side schema validator — the mirror of the consumer trust boundary
|
||||
* (planning-artifacts.md). Producers MUST emit a manifest this accepts; the
|
||||
* core loop re-validates defensively on read.
|
||||
*/
|
||||
function validateManifest(value) {
|
||||
if (!value || typeof value !== 'object') {
|
||||
return { ok: false, errors: ['manifest must be an object'] };
|
||||
}
|
||||
const m = value;
|
||||
const errors = [];
|
||||
if (m.version !== MANIFEST_VERSION)
|
||||
errors.push(`version must be "${MANIFEST_VERSION}"`);
|
||||
for (const f of ['job_id', 'plan_id', 'phase', 'backend', 'submit_command', 'verification_command', 'resume_command', 'submitted_at']) {
|
||||
if (typeof m[f] !== 'string' || m[f].length === 0) {
|
||||
errors.push(`field "${f}" must be a non-empty string`);
|
||||
}
|
||||
}
|
||||
if (typeof m.status !== 'string' || !STATUS_LIST.includes(m.status)) {
|
||||
errors.push(`status must be one of ${MANIFEST_STATUS.join(', ')}`);
|
||||
}
|
||||
if (!Array.isArray(m.expected_artifacts) || !m.expected_artifacts.every((x) => typeof x === 'string')) {
|
||||
errors.push('expected_artifacts must be a string[]');
|
||||
}
|
||||
if (m.terminal_details !== null && typeof m.terminal_details !== 'object') {
|
||||
errors.push('terminal_details must be null or an object');
|
||||
}
|
||||
return errors.length === 0 ? { ok: true } : { ok: false, errors };
|
||||
}
|
||||
/**
|
||||
* Parse `sbatch --parsable` output. Accepts either a bare job id
|
||||
* (`"12345"`) or the `"12345;clustername"` form. Rejects prose like
|
||||
* `"Submitted batch job 12345"` (that is the non-parsable default format).
|
||||
*/
|
||||
function parseSbatchParsable(stdout) {
|
||||
if (typeof stdout !== 'string')
|
||||
return { ok: false, kind: 'non_string', raw: String(stdout) };
|
||||
const trimmed = stdout.trim();
|
||||
if (!trimmed)
|
||||
return { ok: false, kind: 'empty', raw: stdout };
|
||||
const head = trimmed.split(';')[0].split(/\s+/)[0];
|
||||
if (!/^\d+$/.test(head))
|
||||
return { ok: false, kind: 'not_numeric', raw: trimmed };
|
||||
return { ok: true, job_id: head };
|
||||
}
|
||||
/**
|
||||
* Parse a single `squeue` line of the form `"<jobid> <state>"`. Returns `null`
|
||||
* for a header or any row that does not have at least two tokens.
|
||||
*/
|
||||
function parseSqueueLine(line) {
|
||||
if (typeof line !== 'string')
|
||||
return null;
|
||||
const parts = line.trim().split(/\s+/);
|
||||
if (parts.length < 2)
|
||||
return null;
|
||||
const [job_id, state] = parts;
|
||||
if (!/^\d+$/.test(job_id))
|
||||
return null;
|
||||
return { job_id, state };
|
||||
}
|
||||
/**
|
||||
* Parse a `sacct -P` row given as pre-split columns where index 0 is the job
|
||||
* id and index 1 is the state. Returns `null` for malformed rows.
|
||||
*/
|
||||
function parseSacctRow(cols) {
|
||||
if (cols.length < 2)
|
||||
return null;
|
||||
const job_id = cols[0];
|
||||
const state = cols[1];
|
||||
if (typeof job_id !== 'string' || typeof state !== 'string')
|
||||
return null;
|
||||
if (!/^\d+$/.test(job_id))
|
||||
return null;
|
||||
return { job_id, state };
|
||||
}
|
||||
/**
|
||||
* Pure path projection: `.planning/async-jobs/<job_id>.json`.
|
||||
*/
|
||||
function manifestPath(planningDir, jobId) {
|
||||
return node_path_1.default.join(planningDir, 'async-jobs', `${jobId}.json`);
|
||||
}
|
||||
function _isNonTerminal(status) {
|
||||
return typeof status === 'string' && NON_TERMINAL_STATUSES.includes(status);
|
||||
}
|
||||
/**
|
||||
* Write a manifest to `.planning/async-jobs/<job_id>.json`.
|
||||
*
|
||||
* Fail-closed rules (mirror of the consumer contract, planning-artifacts.md):
|
||||
* - If any existing manifest in the dir shares `plan_id` but has a different
|
||||
* `job_id` AND is non-terminal -> refuse (`duplicate_plan_id`); dispatching
|
||||
* again would duplicate the external job.
|
||||
* - If the target file exists but is not valid JSON -> refuse
|
||||
* (`malformed_existing`); never silently clobber.
|
||||
* - Same `job_id` for the same `plan_id` -> allowed (status progression).
|
||||
* - A prior job for the same `plan_id` that is already terminal -> allowed
|
||||
* (the duplicate guard only protects against re-dispatching live work).
|
||||
*/
|
||||
function writeManifest(manifest, planningDir, opts = {}) {
|
||||
const fs = opts.fs ?? node_fs_1.default;
|
||||
const dir = node_path_1.default.join(planningDir, 'async-jobs');
|
||||
const target = manifestPath(planningDir, manifest.job_id);
|
||||
let names;
|
||||
try {
|
||||
fs.mkdirSync(dir, { recursive: true });
|
||||
names = fs.readdirSync(dir);
|
||||
}
|
||||
catch (e) {
|
||||
return { ok: false, kind: 'io_error', message: e.message };
|
||||
}
|
||||
for (const name of names) {
|
||||
if (!name.endsWith('.json'))
|
||||
continue;
|
||||
const p = node_path_1.default.join(dir, name);
|
||||
let raw;
|
||||
try {
|
||||
raw = String(fs.readFileSync(p));
|
||||
}
|
||||
catch {
|
||||
continue;
|
||||
}
|
||||
let existing;
|
||||
try {
|
||||
existing = JSON.parse(raw);
|
||||
}
|
||||
catch {
|
||||
if (p === target) {
|
||||
return { ok: false, kind: 'malformed_existing', message: `target manifest ${p} is not valid JSON` };
|
||||
}
|
||||
continue;
|
||||
}
|
||||
const samePlan = existing.plan_id === manifest.plan_id;
|
||||
const sameJob = existing.job_id === manifest.job_id;
|
||||
if (samePlan && !sameJob && _isNonTerminal(existing.status)) {
|
||||
return {
|
||||
ok: false,
|
||||
kind: 'duplicate_plan_id',
|
||||
message: `plan_id "${manifest.plan_id}" already has non-terminal job "${String(existing.job_id)}" at ${p}; dispatching again would duplicate the external job`,
|
||||
};
|
||||
}
|
||||
}
|
||||
try {
|
||||
fs.writeFileSync(target, JSON.stringify(manifest, null, 2) + '\n');
|
||||
}
|
||||
catch (e) {
|
||||
return { ok: false, kind: 'io_error', message: e.message };
|
||||
}
|
||||
return { ok: true, path: target };
|
||||
}
|
||||
module.exports = {
|
||||
MANIFEST_VERSION,
|
||||
MANIFEST_STATUS,
|
||||
NON_TERMINAL_STATUSES,
|
||||
TERMINAL_FAILURE_STATUSES,
|
||||
mapSlurmState,
|
||||
buildManifest,
|
||||
validateManifest,
|
||||
parseSbatchParsable,
|
||||
parseSqueueLine,
|
||||
parseSacctRow,
|
||||
writeManifest,
|
||||
manifestPath,
|
||||
};
|
||||
@@ -106,12 +106,6 @@
|
||||
],
|
||||
"issue": "TBD"
|
||||
},
|
||||
"external-job": {
|
||||
"files": [
|
||||
"external-job-waiting.test.cjs"
|
||||
],
|
||||
"issue": "#1165"
|
||||
},
|
||||
"host-integration": {
|
||||
"files": [
|
||||
"host-integration.test.cjs",
|
||||
|
||||
195
scripts/slurm-adapter.cjs
Normal file
195
scripts/slurm-adapter.cjs
Normal file
@@ -0,0 +1,195 @@
|
||||
#!/usr/bin/env node
|
||||
'use strict';
|
||||
|
||||
/**
|
||||
* slurm-adapter.cjs — SLURM scheduler-adapter helper for the external-job
|
||||
* capability (#1164 / #1105).
|
||||
*
|
||||
* Thin CLI: runs bounded sbatch / squeue / sacct subprocesses and delegates
|
||||
* all parsing, manifest build/validate, and fail-closed writing to the pure
|
||||
* module (gsd-core/bin/lib/external-job.cjs). The pure module is fully unit-
|
||||
* tested; this script is the operator surface that needs a real cluster.
|
||||
*
|
||||
* Subcommands:
|
||||
* submit --plan <plan_id> --phase <phase> --expected <path>[,<path>] \
|
||||
* --verify <cmd> --resume <cmd> -- <sbatch...>
|
||||
* poll --job <job_id> [--plan <plan_id>]
|
||||
* show --job <job_id> (surface manifest status + commands; no auto-run)
|
||||
*
|
||||
* Trust boundary: this script never auto-runs verification_command or
|
||||
* resume_command from a manifest (planning-artifacts.md). `show` prints them
|
||||
* for explicit operator confirmation.
|
||||
*/
|
||||
|
||||
const { execFileSync } = require('node:child_process');
|
||||
const fs = require('node:fs');
|
||||
const path = require('node:path');
|
||||
|
||||
const { ExitError, runMain } = require('./lib/cli-exit.cjs');
|
||||
const m = require('../gsd-core/bin/lib/external-job.cjs');
|
||||
|
||||
const SUBMIT_TIMEOUT_MS = Number(process.env.GSD_SLURM_SUBMIT_TIMEOUT_MS || 30000);
|
||||
const POLL_TIMEOUT_MS = Number(process.env.GSD_SLURM_POLL_TIMEOUT_MS || 15000);
|
||||
|
||||
function usage() {
|
||||
return [
|
||||
'usage: slurm-adapter.cjs <submit|poll|show> ...',
|
||||
' submit --plan <id> --phase <n> --expected <p1,p2> --verify <cmd> --resume <cmd> -- sbatch --parsable ...',
|
||||
' poll --job <job_id> [--plan <plan_id>]',
|
||||
' show --job <job_id>',
|
||||
].join('\n');
|
||||
}
|
||||
|
||||
function parseFlags(argv) {
|
||||
const out = {};
|
||||
const rest = [];
|
||||
for (let i = 0; i < argv.length; i++) {
|
||||
const a = argv[i];
|
||||
if (a === '--') { out['--'] = argv.slice(i + 1); break; }
|
||||
if (a.startsWith('--')) {
|
||||
const v = argv[i + 1];
|
||||
out[a.slice(2)] = v;
|
||||
i++;
|
||||
} else {
|
||||
rest.push(a);
|
||||
}
|
||||
}
|
||||
out._ = rest;
|
||||
return out;
|
||||
}
|
||||
|
||||
function findPlanningDir(start) {
|
||||
let dir = path.resolve(start || process.cwd());
|
||||
for (let i = 0; i < 10; i++) {
|
||||
if (fs.existsSync(path.join(dir, '.planning'))) return path.join(dir, '.planning');
|
||||
const parent = path.dirname(dir);
|
||||
if (parent === dir) break;
|
||||
dir = parent;
|
||||
}
|
||||
throw new ExitError(1, 'could not locate a .planning directory (walked up 10 levels)');
|
||||
}
|
||||
|
||||
function cmdSubmit(flags) {
|
||||
const plan = flags.plan;
|
||||
const phase = flags.phase;
|
||||
const sbatchCmd = flags['--'];
|
||||
if (!plan || phase === undefined || !Array.isArray(sbatchCmd) || sbatchCmd.length === 0) {
|
||||
throw new ExitError(1, 'submit requires --plan, --phase, and an sbatch command after --\n' + usage());
|
||||
}
|
||||
const expected = (flags.expected || '').split(',').map((s) => s.trim()).filter(Boolean);
|
||||
if (expected.length === 0) throw new ExitError(1, 'submit requires --expected (comma-separated artifact paths)');
|
||||
const verify = flags.verify;
|
||||
const resume = flags.resume || ('/gsd:execute-phase ' + phase);
|
||||
if (!verify) throw new ExitError(1, 'submit requires --verify (the command that verifies job output)');
|
||||
|
||||
let stdout;
|
||||
try {
|
||||
stdout = execFileSync(sbatchCmd[0], sbatchCmd.slice(1), {
|
||||
encoding: 'utf8',
|
||||
timeout: SUBMIT_TIMEOUT_MS,
|
||||
maxBuffer: 1024 * 1024,
|
||||
});
|
||||
} catch (e) {
|
||||
throw new ExitError(1, 'sbatch failed: ' + (e.message || String(e)));
|
||||
}
|
||||
const parsed = m.parseSbatchParsable(stdout);
|
||||
if (!parsed.ok) {
|
||||
throw new ExitError(1, 'could not parse sbatch --parsable output (kind=' + parsed.kind + '): ' + parsed.raw);
|
||||
}
|
||||
const manifest = m.buildManifest({
|
||||
plan_id: plan,
|
||||
phase,
|
||||
job_id: parsed.job_id,
|
||||
backend: 'slurm',
|
||||
submit_command: sbatchCmd.join(' '),
|
||||
status: 'submitted',
|
||||
expected_artifacts: expected,
|
||||
verification_command: verify,
|
||||
resume_command: resume,
|
||||
});
|
||||
const planningDir = findPlanningDir();
|
||||
const res = m.writeManifest(manifest, planningDir);
|
||||
if (!res.ok) {
|
||||
throw new ExitError(1, 'writeManifest refused (' + res.kind + '): ' + res.message);
|
||||
}
|
||||
process.stdout.write('submitted job ' + parsed.job_id + ' for plan ' + plan + '\n');
|
||||
process.stdout.write('manifest: ' + res.path + '\n');
|
||||
process.stdout.write('state: external_job_waiting (SUMMARY deferred)\n');
|
||||
}
|
||||
|
||||
function cmdPoll(flags) {
|
||||
const jobId = flags.job;
|
||||
if (!jobId) throw new ExitError(1, 'poll requires --job <job_id>\n' + usage());
|
||||
const planningDir = findPlanningDir();
|
||||
const manifestFile = m.manifestPath(planningDir, jobId);
|
||||
if (!fs.existsSync(manifestFile)) {
|
||||
throw new ExitError(1, 'no manifest for job ' + jobId + ' at ' + manifestFile);
|
||||
}
|
||||
const existing = JSON.parse(fs.readFileSync(manifestFile, 'utf8'));
|
||||
|
||||
let rawState = null;
|
||||
try {
|
||||
const out = execFileSync('squeue', ['-h', '-j', jobId, '-o', '%i %T'], {
|
||||
encoding: 'utf8', timeout: POLL_TIMEOUT_MS, maxBuffer: 1024 * 1024,
|
||||
}).trim();
|
||||
const line = out.split('\n')[0];
|
||||
const parsed = m.parseSqueueLine(line || '');
|
||||
if (parsed) rawState = parsed.state;
|
||||
} catch (_e) { /* squeue empty/failed — fall back to sacct */ }
|
||||
if (!rawState) {
|
||||
try {
|
||||
const out = execFileSync('sacct', ['-X', '-P', '-j', jobId, '-o', 'JobID,State'], {
|
||||
encoding: 'utf8', timeout: POLL_TIMEOUT_MS, maxBuffer: 1024 * 1024,
|
||||
}).trim();
|
||||
for (const line of out.split('\n').slice(1)) {
|
||||
const parsed = m.parseSacctRow(line.split('|'));
|
||||
if (parsed) { rawState = parsed.state; break; }
|
||||
}
|
||||
} catch (e) {
|
||||
throw new ExitError(1, 'both squeue and sacct failed: ' + (e.message || String(e)));
|
||||
}
|
||||
}
|
||||
const mapped = m.mapSlurmState(rawState || '');
|
||||
if (!mapped) throw new ExitError(1, 'unmapped SLURM state "' + rawState + '" — not guessing; inspect manually');
|
||||
|
||||
const updated = m.buildManifest(
|
||||
Object.assign({}, existing, { status: mapped, terminal_details: existing.terminal_details || null }),
|
||||
{ clock: { nowIso: () => existing.submitted_at } },
|
||||
);
|
||||
const res = m.writeManifest(updated, planningDir);
|
||||
if (!res.ok) throw new ExitError(1, 'writeManifest refused (' + res.kind + '): ' + res.message);
|
||||
process.stdout.write(JSON.stringify({ job_id: jobId, slurm_state: rawState, manifest_status: mapped, path: res.path }) + '\n');
|
||||
}
|
||||
|
||||
function cmdShow(flags) {
|
||||
const jobId = flags.job;
|
||||
if (!jobId) throw new ExitError(1, 'show requires --job <job_id>\n' + usage());
|
||||
const planningDir = findPlanningDir();
|
||||
const manifestFile = m.manifestPath(planningDir, jobId);
|
||||
if (!fs.existsSync(manifestFile)) {
|
||||
throw new ExitError(1, 'no manifest for job ' + jobId + ' at ' + manifestFile);
|
||||
}
|
||||
const manifest = JSON.parse(fs.readFileSync(manifestFile, 'utf8'));
|
||||
process.stdout.write('job ' + manifest.job_id + ' (plan ' + manifest.plan_id + ', backend ' + manifest.backend + ')\n');
|
||||
process.stdout.write('status: ' + manifest.status + '\n');
|
||||
if (manifest.terminal_details) {
|
||||
process.stdout.write('terminal_details: ' + JSON.stringify(manifest.terminal_details) + '\n');
|
||||
}
|
||||
// Trust boundary: surface commands for confirmation, never auto-run.
|
||||
process.stdout.write('\nManifest commands (UNTRUSTED — confirm before running):\n');
|
||||
process.stdout.write(' submit_command: ' + manifest.submit_command + '\n');
|
||||
process.stdout.write(' verification_command: ' + manifest.verification_command + '\n');
|
||||
process.stdout.write(' resume_command: ' + manifest.resume_command + '\n');
|
||||
}
|
||||
|
||||
function main() {
|
||||
const [, , sub, ...rest] = process.argv;
|
||||
const flags = parseFlags(rest);
|
||||
if (sub === 'submit') cmdSubmit(flags);
|
||||
else if (sub === 'poll') cmdPoll(flags);
|
||||
else if (sub === 'show') cmdShow(flags);
|
||||
else throw new ExitError(1, usage());
|
||||
return 0;
|
||||
}
|
||||
|
||||
runMain(main);
|
||||
360
src/external-job.cts
Normal file
360
src/external-job.cts
Normal file
@@ -0,0 +1,360 @@
|
||||
/**
|
||||
* external-job.cts — scheduler-adapter producer module for the async
|
||||
* external-job contract (#1164 / #1105).
|
||||
*
|
||||
* The CORE loop CONSUMES manifests at `.planning/async-jobs/<job>.json`
|
||||
* (#1165 — external_job_waiting half-state); this module is the Capability
|
||||
* half that PRODUCES them. SLURM is the first backend; the design stays
|
||||
* scheduler-pluggable via the `backend` field (planning-artifacts.md).
|
||||
*
|
||||
* Pure helpers (state map, build, validate, parsers) take no I/O; the writer
|
||||
* takes injected `fs` and `clock` seams so tests drive it without touching disk
|
||||
* or wall-clock time (CLAUDE.md clock-seam + injectable-deps conventions).
|
||||
*
|
||||
* Build: `src/*.cts` -> `gsd-core/bin/lib/*.cjs` (ADR-457 build-at-publish).
|
||||
*/
|
||||
|
||||
import path from 'node:path';
|
||||
import nodeFs from 'node:fs';
|
||||
import type { PathLike } from 'node:fs';
|
||||
|
||||
// ─── Closed status enum (stability contract — Hyrum's Law) ────────────────────
|
||||
|
||||
const MANIFEST_VERSION = '1.0';
|
||||
|
||||
const MANIFEST_STATUS = [
|
||||
'submitted',
|
||||
'running',
|
||||
'completed-unverified',
|
||||
'failed',
|
||||
'cancelled',
|
||||
'timeout',
|
||||
] as const;
|
||||
type ManifestStatus = (typeof MANIFEST_STATUS)[number];
|
||||
|
||||
const NON_TERMINAL_STATUSES: ReadonlyArray<ManifestStatus> = ['submitted', 'running'];
|
||||
const TERMINAL_FAILURE_STATUSES: ReadonlyArray<ManifestStatus> = ['failed', 'cancelled', 'timeout'];
|
||||
|
||||
// ─── SLURM state -> manifest status ───────────────────────────────────────────
|
||||
//
|
||||
// Source: SLURM job state codes (squeue/sacct State column). Producers for
|
||||
// other backends map their own states onto the closed enum above; this table
|
||||
// is SLURM-specific and lives behind the `backend: 'slurm'` field.
|
||||
|
||||
const SLURM_STATE_MAP: Readonly<Record<string, ManifestStatus>> = {
|
||||
PENDING: 'submitted',
|
||||
CONFIGURING: 'submitted',
|
||||
RUNNING: 'running',
|
||||
COMPLETING: 'running',
|
||||
COMPLETED: 'completed-unverified',
|
||||
FAILED: 'failed',
|
||||
CANCELLED: 'cancelled',
|
||||
TIMEOUT: 'timeout',
|
||||
OUT_OF_MEMORY: 'failed',
|
||||
BOOT_FAIL: 'failed',
|
||||
NODE_FAIL: 'failed',
|
||||
PREEMPTED: 'failed',
|
||||
};
|
||||
|
||||
/**
|
||||
* Map a raw SLURM state string to the closed, scheduler-agnostic manifest
|
||||
* status. Case-insensitive; trims whitespace; strips a trailing by-part
|
||||
* ("CANCELLED by 1001" -> "CANCELLED"). Returns `null` for any unknown
|
||||
* state so the caller can decide whether to surface or fail — never guesses
|
||||
* (CLAUDE.md anti-guessing).
|
||||
*/
|
||||
function mapSlurmState(raw: string): ManifestStatus | null {
|
||||
if (typeof raw !== 'string') return null;
|
||||
const key = raw.trim().toUpperCase();
|
||||
const head = key.split(/\s+/)[0];
|
||||
if (head && Object.prototype.hasOwnProperty.call(SLURM_STATE_MAP, head)) {
|
||||
return SLURM_STATE_MAP[head];
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
// ─── Manifest shape ───────────────────────────────────────────────────────────
|
||||
|
||||
interface ManifestTerminalDetails {
|
||||
readonly reason?: string;
|
||||
readonly exit_code?: number;
|
||||
readonly [k: string]: unknown;
|
||||
}
|
||||
|
||||
interface Manifest {
|
||||
readonly version: string;
|
||||
readonly job_id: string;
|
||||
readonly plan_id: string;
|
||||
readonly phase: string;
|
||||
readonly backend: string;
|
||||
readonly submit_command: string;
|
||||
readonly status: ManifestStatus;
|
||||
readonly expected_artifacts: ReadonlyArray<string>;
|
||||
readonly verification_command: string;
|
||||
readonly resume_command: string;
|
||||
readonly submitted_at: string;
|
||||
readonly terminal_details: ManifestTerminalDetails | null;
|
||||
}
|
||||
|
||||
interface BuildManifestInput {
|
||||
readonly plan_id: string;
|
||||
readonly phase: string;
|
||||
readonly job_id: string;
|
||||
readonly backend: string;
|
||||
readonly submit_command: string;
|
||||
readonly status: ManifestStatus;
|
||||
readonly expected_artifacts: ReadonlyArray<string>;
|
||||
readonly verification_command: string;
|
||||
readonly resume_command: string;
|
||||
readonly terminal_details?: ManifestTerminalDetails | null;
|
||||
}
|
||||
|
||||
interface Clock {
|
||||
nowIso(): string;
|
||||
}
|
||||
|
||||
const REQUIRED_FIELDS: ReadonlyArray<keyof BuildManifestInput> = [
|
||||
'plan_id',
|
||||
'phase',
|
||||
'job_id',
|
||||
'backend',
|
||||
'submit_command',
|
||||
'status',
|
||||
'expected_artifacts',
|
||||
'verification_command',
|
||||
'resume_command',
|
||||
];
|
||||
|
||||
function assertString(v: unknown, key: string): void {
|
||||
if (typeof v !== 'string' || v.length === 0) {
|
||||
throw new Error(`buildManifest: field "${key}" must be a non-empty string`);
|
||||
}
|
||||
}
|
||||
|
||||
const STATUS_LIST: ReadonlyArray<string> = MANIFEST_STATUS;
|
||||
|
||||
/**
|
||||
* Build a versioned, frozen manifest. Stamps `version` and `submitted_at`
|
||||
* (via the injected clock seam) and normalises `terminal_details`:
|
||||
* `null` unless the status is a terminal failure AND details were supplied.
|
||||
* Throws on missing required fields or an out-of-enum status.
|
||||
*/
|
||||
function buildManifest(
|
||||
input: BuildManifestInput,
|
||||
opts: { clock?: Clock } = {},
|
||||
): Manifest {
|
||||
const inputRecord = input as unknown as Record<string, unknown>;
|
||||
for (const key of REQUIRED_FIELDS) {
|
||||
const v = inputRecord[key];
|
||||
if (key === 'expected_artifacts') {
|
||||
if (!Array.isArray(v) || v.length === 0 || !v.every((x) => typeof x === 'string')) {
|
||||
throw new Error('buildManifest: field "expected_artifacts" must be a non-empty string[]');
|
||||
}
|
||||
continue;
|
||||
}
|
||||
assertString(v, key);
|
||||
}
|
||||
if (!STATUS_LIST.includes(input.status)) {
|
||||
throw new Error(`buildManifest: field "status" must be one of ${MANIFEST_STATUS.join(', ')}`);
|
||||
}
|
||||
const clock = opts.clock ?? { nowIso: () => new Date().toISOString() };
|
||||
const isTerminalFailure = (TERMINAL_FAILURE_STATUSES as ReadonlyArray<string>).includes(input.status);
|
||||
const terminal_details =
|
||||
isTerminalFailure && input.terminal_details ? input.terminal_details : null;
|
||||
return Object.freeze({
|
||||
version: MANIFEST_VERSION,
|
||||
job_id: input.job_id,
|
||||
plan_id: input.plan_id,
|
||||
phase: input.phase,
|
||||
backend: input.backend,
|
||||
submit_command: input.submit_command,
|
||||
status: input.status,
|
||||
expected_artifacts: [...input.expected_artifacts],
|
||||
verification_command: input.verification_command,
|
||||
resume_command: input.resume_command,
|
||||
submitted_at: clock.nowIso(),
|
||||
terminal_details,
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* Producer-side schema validator — the mirror of the consumer trust boundary
|
||||
* (planning-artifacts.md). Producers MUST emit a manifest this accepts; the
|
||||
* core loop re-validates defensively on read.
|
||||
*/
|
||||
function validateManifest(
|
||||
value: unknown,
|
||||
): { ok: true } | { ok: false; errors: string[] } {
|
||||
if (!value || typeof value !== 'object') {
|
||||
return { ok: false, errors: ['manifest must be an object'] };
|
||||
}
|
||||
const m = value as Record<string, unknown>;
|
||||
const errors: string[] = [];
|
||||
if (m.version !== MANIFEST_VERSION) errors.push(`version must be "${MANIFEST_VERSION}"`);
|
||||
for (const f of ['job_id', 'plan_id', 'phase', 'backend', 'submit_command', 'verification_command', 'resume_command', 'submitted_at']) {
|
||||
if (typeof m[f] !== 'string' || m[f].length === 0) {
|
||||
errors.push(`field "${f}" must be a non-empty string`);
|
||||
}
|
||||
}
|
||||
if (typeof m.status !== 'string' || !STATUS_LIST.includes(m.status)) {
|
||||
errors.push(`status must be one of ${MANIFEST_STATUS.join(', ')}`);
|
||||
}
|
||||
if (!Array.isArray(m.expected_artifacts) || !m.expected_artifacts.every((x) => typeof x === 'string')) {
|
||||
errors.push('expected_artifacts must be a string[]');
|
||||
}
|
||||
if (m.terminal_details !== null && typeof m.terminal_details !== 'object') {
|
||||
errors.push('terminal_details must be null or an object');
|
||||
}
|
||||
return errors.length === 0 ? { ok: true } : { ok: false, errors };
|
||||
}
|
||||
|
||||
// ─── SLURM CLI output parsers ─────────────────────────────────────────────────
|
||||
|
||||
interface ParseErr { ok: false; kind: string; raw: string }
|
||||
|
||||
/**
|
||||
* Parse `sbatch --parsable` output. Accepts either a bare job id
|
||||
* (`"12345"`) or the `"12345;clustername"` form. Rejects prose like
|
||||
* `"Submitted batch job 12345"` (that is the non-parsable default format).
|
||||
*/
|
||||
function parseSbatchParsable(stdout: string): { ok: true; job_id: string } | ParseErr {
|
||||
if (typeof stdout !== 'string') return { ok: false, kind: 'non_string', raw: String(stdout) };
|
||||
const trimmed = stdout.trim();
|
||||
if (!trimmed) return { ok: false, kind: 'empty', raw: stdout };
|
||||
const head = trimmed.split(';')[0].split(/\s+/)[0];
|
||||
if (!/^\d+$/.test(head)) return { ok: false, kind: 'not_numeric', raw: trimmed };
|
||||
return { ok: true, job_id: head };
|
||||
}
|
||||
|
||||
/**
|
||||
* Parse a single `squeue` line of the form `"<jobid> <state>"`. Returns `null`
|
||||
* for a header or any row that does not have at least two tokens.
|
||||
*/
|
||||
function parseSqueueLine(line: string): { job_id: string; state: string } | null {
|
||||
if (typeof line !== 'string') return null;
|
||||
const parts = line.trim().split(/\s+/);
|
||||
if (parts.length < 2) return null;
|
||||
const [job_id, state] = parts;
|
||||
if (!/^\d+$/.test(job_id)) return null;
|
||||
return { job_id, state };
|
||||
}
|
||||
|
||||
/**
|
||||
* Parse a `sacct -P` row given as pre-split columns where index 0 is the job
|
||||
* id and index 1 is the state. Returns `null` for malformed rows.
|
||||
*/
|
||||
function parseSacctRow(cols: ReadonlyArray<string>): { job_id: string; state: string } | null {
|
||||
if (cols.length < 2) return null;
|
||||
const job_id = cols[0];
|
||||
const state = cols[1];
|
||||
if (typeof job_id !== 'string' || typeof state !== 'string') return null;
|
||||
if (!/^\d+$/.test(job_id)) return null;
|
||||
return { job_id, state };
|
||||
}
|
||||
|
||||
// ─── Manifest writer (fail-closed duplicate guard) ────────────────────────────
|
||||
|
||||
interface FsLike {
|
||||
mkdirSync(p: PathLike, opts?: unknown): unknown;
|
||||
readdirSync(p: PathLike): ReadonlyArray<string>;
|
||||
readFileSync(p: PathLike): string | Buffer;
|
||||
writeFileSync(p: PathLike, data: string): unknown;
|
||||
existsSync(p: PathLike): boolean;
|
||||
}
|
||||
|
||||
type WriteResult =
|
||||
| { ok: true; path: string }
|
||||
| { ok: false; kind: 'malformed_existing' | 'duplicate_plan_id' | 'io_error'; message: string };
|
||||
|
||||
/**
|
||||
* Pure path projection: `.planning/async-jobs/<job_id>.json`.
|
||||
*/
|
||||
function manifestPath(planningDir: string, jobId: string): string {
|
||||
return path.join(planningDir, 'async-jobs', `${jobId}.json`);
|
||||
}
|
||||
|
||||
function _isNonTerminal(status: unknown): boolean {
|
||||
return typeof status === 'string' && (NON_TERMINAL_STATUSES as ReadonlyArray<string>).includes(status);
|
||||
}
|
||||
|
||||
/**
|
||||
* Write a manifest to `.planning/async-jobs/<job_id>.json`.
|
||||
*
|
||||
* Fail-closed rules (mirror of the consumer contract, planning-artifacts.md):
|
||||
* - If any existing manifest in the dir shares `plan_id` but has a different
|
||||
* `job_id` AND is non-terminal -> refuse (`duplicate_plan_id`); dispatching
|
||||
* again would duplicate the external job.
|
||||
* - If the target file exists but is not valid JSON -> refuse
|
||||
* (`malformed_existing`); never silently clobber.
|
||||
* - Same `job_id` for the same `plan_id` -> allowed (status progression).
|
||||
* - A prior job for the same `plan_id` that is already terminal -> allowed
|
||||
* (the duplicate guard only protects against re-dispatching live work).
|
||||
*/
|
||||
function writeManifest(
|
||||
manifest: Manifest,
|
||||
planningDir: string,
|
||||
opts: { fs?: FsLike; clock?: Clock } = {},
|
||||
): WriteResult {
|
||||
const fs = opts.fs ?? nodeFs;
|
||||
const dir = path.join(planningDir, 'async-jobs');
|
||||
const target = manifestPath(planningDir, manifest.job_id);
|
||||
|
||||
let names: ReadonlyArray<string>;
|
||||
try {
|
||||
fs.mkdirSync(dir, { recursive: true });
|
||||
names = fs.readdirSync(dir);
|
||||
} catch (e) {
|
||||
return { ok: false, kind: 'io_error', message: (e as Error).message };
|
||||
}
|
||||
|
||||
for (const name of names) {
|
||||
if (!name.endsWith('.json')) continue;
|
||||
const p = path.join(dir, name);
|
||||
let raw: string;
|
||||
try {
|
||||
raw = String(fs.readFileSync(p));
|
||||
} catch {
|
||||
continue;
|
||||
}
|
||||
let existing: Record<string, unknown>;
|
||||
try {
|
||||
existing = JSON.parse(raw) as Record<string, unknown>;
|
||||
} catch {
|
||||
if (p === target) {
|
||||
return { ok: false, kind: 'malformed_existing', message: `target manifest ${p} is not valid JSON` };
|
||||
}
|
||||
continue;
|
||||
}
|
||||
const samePlan = existing.plan_id === manifest.plan_id;
|
||||
const sameJob = existing.job_id === manifest.job_id;
|
||||
if (samePlan && !sameJob && _isNonTerminal(existing.status)) {
|
||||
return {
|
||||
ok: false,
|
||||
kind: 'duplicate_plan_id',
|
||||
message: `plan_id "${manifest.plan_id}" already has non-terminal job "${String(existing.job_id)}" at ${p}; dispatching again would duplicate the external job`,
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
try {
|
||||
fs.writeFileSync(target, JSON.stringify(manifest, null, 2) + '\n');
|
||||
} catch (e) {
|
||||
return { ok: false, kind: 'io_error', message: (e as Error).message };
|
||||
}
|
||||
return { ok: true, path: target };
|
||||
}
|
||||
|
||||
export = {
|
||||
MANIFEST_VERSION,
|
||||
MANIFEST_STATUS,
|
||||
NON_TERMINAL_STATUSES,
|
||||
TERMINAL_FAILURE_STATUSES,
|
||||
mapSlurmState,
|
||||
buildManifest,
|
||||
validateManifest,
|
||||
parseSbatchParsable,
|
||||
parseSqueueLine,
|
||||
parseSacctRow,
|
||||
writeManifest,
|
||||
manifestPath,
|
||||
};
|
||||
@@ -493,7 +493,7 @@ describe('resolveLoopHooks plan:post — pure function against real registry', (
|
||||
assert.strictEqual(result.activeHooks[0].capId, 'gap-analysis');
|
||||
});
|
||||
|
||||
test('[happy] real registry byLoopPoint plan:post has 1 step (mempalace), 0 contributions, and 1 gate (gap-analysis)', () => {
|
||||
test('[happy] real registry byLoopPoint plan:post has 1 step (mempalace), 1 contribution (external-job planner fragment), and 1 gate (gap-analysis)', () => {
|
||||
const entry = realRegistry.byLoopPoint['plan:post'];
|
||||
assert.ok(entry, 'plan:post must exist in byLoopPoint');
|
||||
assert.ok(Array.isArray(entry.steps), 'steps must be an array');
|
||||
@@ -501,7 +501,8 @@ describe('resolveLoopHooks plan:post — pure function against real registry', (
|
||||
assert.ok(Array.isArray(entry.gates), 'gates must be an array');
|
||||
assert.strictEqual(entry.steps.length, 1, 'plan:post must have 1 step (mempalace capture)');
|
||||
assert.strictEqual(entry.steps[0].capId, 'mempalace', 'plan:post step must be from mempalace');
|
||||
assert.strictEqual(entry.contributions.length, 0, 'plan:post must have zero contributions');
|
||||
assert.strictEqual(entry.contributions.length, 1, 'plan:post must have 1 contribution (external-job planner runtime-budget fragment)');
|
||||
assert.strictEqual(entry.contributions[0].capId, 'external-job', 'plan:post contribution must be from external-job');
|
||||
assert.strictEqual(entry.gates.length, 1, 'plan:post must have exactly one gate');
|
||||
assert.strictEqual(entry.gates[0].capId, 'gap-analysis');
|
||||
});
|
||||
|
||||
@@ -629,14 +629,15 @@ describe('F. Real registry execute:wave:post shape — guard against accidental
|
||||
`ui.safety-gate onError must be 'halt'; got ${uiGate.onError}`);
|
||||
});
|
||||
|
||||
test('[happy] real registry: execute:wave:post has no steps and 1 contribution (mempalace capture-problems) — gate point with mempalace contribution', () => {
|
||||
test('[happy] real registry: execute:wave:post has no steps and 2 contributions (mempalace capture-problems + external-job executor fragment)', () => {
|
||||
const point = realRegistry.byLoopPoint['execute:wave:post'];
|
||||
assert.strictEqual(point.steps.length, 0,
|
||||
`execute:wave:post steps must be empty; got ${point.steps.length}`);
|
||||
assert.strictEqual(point.contributions.length, 1,
|
||||
`execute:wave:post must have 1 contribution (mempalace); got ${point.contributions.length}`);
|
||||
assert.strictEqual(point.contributions[0].capId, 'mempalace',
|
||||
`execute:wave:post contribution must be from mempalace; got ${point.contributions[0].capId}`);
|
||||
assert.strictEqual(point.contributions.length, 2,
|
||||
`execute:wave:post must have 2 contributions (mempalace + external-job); got ${point.contributions.length}`);
|
||||
const capIds = point.contributions.map(c => c.capId).sort();
|
||||
assert.deepStrictEqual(capIds, ['external-job', 'mempalace'],
|
||||
`execute:wave:post contributions must be mempalace + external-job; got ${capIds.join(',')}`);
|
||||
});
|
||||
|
||||
});
|
||||
|
||||
311
tests/external-job.test.cjs
Normal file
311
tests/external-job.test.cjs
Normal file
@@ -0,0 +1,311 @@
|
||||
'use strict';
|
||||
|
||||
// Producer-half of the async external-job contract (#1164 / #1105).
|
||||
// The CORE consumer half (external_job_waiting) is covered by
|
||||
// external-job-waiting.test.cjs. These tests assert the scheduler-adapter
|
||||
// Capability's pure module: SLURM state mapping, manifest build/validate,
|
||||
// sbatch/squeue/sacct parsers, and the fail-closed manifest writer that
|
||||
// mirrors the consumer's duplicate-execution guard
|
||||
// (docs/reference/planning-artifacts.md).
|
||||
|
||||
const { test } = require('node:test');
|
||||
const assert = require('node:assert');
|
||||
const path = require('node:path');
|
||||
const fc = require('fast-check');
|
||||
|
||||
const m = require('../gsd-core/bin/lib/external-job.cjs');
|
||||
const {
|
||||
MANIFEST_VERSION,
|
||||
MANIFEST_STATUS,
|
||||
NON_TERMINAL_STATUSES,
|
||||
TERMINAL_FAILURE_STATUSES,
|
||||
mapSlurmState,
|
||||
buildManifest,
|
||||
validateManifest,
|
||||
parseSbatchParsable,
|
||||
parseSqueueLine,
|
||||
parseSacctRow,
|
||||
writeManifest,
|
||||
manifestPath,
|
||||
} = m;
|
||||
|
||||
// ─── Closed status enum (Hyrum's Law: stability contract) ─────────────────────
|
||||
|
||||
test('MANIFEST_STATUS is the closed scheduler-agnostic enum from the contract', () => {
|
||||
assert.deepStrictEqual([...MANIFEST_STATUS].sort(), [
|
||||
'cancelled',
|
||||
'completed-unverified',
|
||||
'failed',
|
||||
'running',
|
||||
'submitted',
|
||||
'timeout',
|
||||
]);
|
||||
assert.strictEqual(MANIFEST_VERSION, '1.0');
|
||||
});
|
||||
|
||||
test('NON_TERMINAL / TERMINAL_FAILURE partition the enum without overlap', () => {
|
||||
for (const s of MANIFEST_STATUS) {
|
||||
const inNon = NON_TERMINAL_STATUSES.includes(s);
|
||||
const inTerm = TERMINAL_FAILURE_STATUSES.includes(s);
|
||||
// completed-unverified is neither non-terminal nor a failure — its own bucket.
|
||||
if (s === 'completed-unverified') {
|
||||
assert.ok(!inNon && !inTerm, 'completed-unverified is its own bucket');
|
||||
} else {
|
||||
assert.ok(inNon !== inTerm, `${s} must sit in exactly one partition`);
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
// ─── SLURM state mapping ──────────────────────────────────────────────────────
|
||||
|
||||
test('mapSlurmState maps every documented SLURM state to the closed enum', () => {
|
||||
const cases = {
|
||||
PENDING: 'submitted',
|
||||
CONFIGURING: 'submitted',
|
||||
RUNNING: 'running',
|
||||
COMPLETED: 'completed-unverified',
|
||||
COMPLETING: 'running',
|
||||
FAILED: 'failed',
|
||||
CANCELLED: 'cancelled',
|
||||
TIMEOUT: 'timeout',
|
||||
OUT_OF_MEMORY: 'failed',
|
||||
BOOT_FAIL: 'failed',
|
||||
NODE_FAIL: 'failed',
|
||||
PREEMPTED: 'failed',
|
||||
};
|
||||
for (const [slurm, expected] of Object.entries(cases)) {
|
||||
assert.strictEqual(mapSlurmState(slurm), expected, `${slurm} -> ${expected}`);
|
||||
}
|
||||
});
|
||||
|
||||
test('mapSlurmState is case-insensitive and trims whitespace', () => {
|
||||
assert.strictEqual(mapSlurmState('running'), 'running');
|
||||
assert.strictEqual(mapSlurmState(' PENDING '), 'submitted');
|
||||
assert.strictEqual(mapSlurmState('Cancelled'), 'cancelled');
|
||||
});
|
||||
|
||||
test('mapSlurmState returns null for unknown states (no guessing)', () => {
|
||||
// Boundary: unknown must not collapse to a terminal failure silently.
|
||||
assert.strictEqual(mapSlurmState('NO_SUCH_STATE'), null);
|
||||
assert.strictEqual(mapSlurmState(''), null);
|
||||
assert.strictEqual(mapSlurmState('COMPLETED2'), null);
|
||||
});
|
||||
|
||||
// ─── Manifest build ───────────────────────────────────────────────────────────
|
||||
|
||||
function baseInput() {
|
||||
return {
|
||||
plan_id: '3.1',
|
||||
phase: '3',
|
||||
job_id: '12345',
|
||||
backend: 'slurm',
|
||||
submit_command: 'sbatch --parsable ./run.sh',
|
||||
status: 'submitted',
|
||||
expected_artifacts: ['Artifacts/jobs/12345/result.h5'],
|
||||
verification_command: 'python -m verify.py 12345',
|
||||
resume_command: '/gsd:execute-phase 3',
|
||||
};
|
||||
}
|
||||
|
||||
test('buildManifest stamps version and submitted_at via the clock seam', () => {
|
||||
const clock = { nowIso: () => '2020-06-15T12:00:00.000Z' };
|
||||
const out = buildManifest(baseInput(), { clock });
|
||||
assert.strictEqual(out.version, '1.0');
|
||||
assert.strictEqual(out.submitted_at, '2020-06-15T12:00:00.000Z');
|
||||
assert.strictEqual(out.terminal_details, null, 'non-terminal -> null terminal_details');
|
||||
assert.strictEqual(out.plan_id, '3.1');
|
||||
});
|
||||
|
||||
test('buildManifest rejects missing required fields', () => {
|
||||
for (const key of ['plan_id', 'phase', 'job_id', 'backend', 'submit_command', 'status', 'expected_artifacts', 'verification_command', 'resume_command']) {
|
||||
const bad = baseInput();
|
||||
delete bad[key];
|
||||
assert.throws(() => buildManifest(bad), { message: new RegExp(key) }, `missing ${key} must throw`);
|
||||
}
|
||||
});
|
||||
|
||||
test('buildManifest rejects an out-of-enum status', () => {
|
||||
const bad = baseInput();
|
||||
bad.status = 'done';
|
||||
assert.throws(() => buildManifest(bad), /status/i);
|
||||
});
|
||||
|
||||
test('buildManifest sets terminal_details when status is a terminal failure', () => {
|
||||
const clock = { nowIso: () => '2020-06-15T12:00:00.000Z' };
|
||||
for (const status of TERMINAL_FAILURE_STATUSES) {
|
||||
const input = { ...baseInput(), status, terminal_details: { reason: 'oom', exit_code: 137 } };
|
||||
const out = buildManifest(input, { clock });
|
||||
assert.deepStrictEqual(out.terminal_details, { reason: 'oom', exit_code: 137 }, `${status} carries terminal_details`);
|
||||
}
|
||||
});
|
||||
|
||||
// ─── validateManifest (producer-side mirror of the trust boundary) ────────────
|
||||
|
||||
test('validateManifest accepts a well-formed manifest', () => {
|
||||
const clock = { nowIso: () => '2020-06-15T12:00:00.000Z' };
|
||||
const res = validateManifest(buildManifest(baseInput(), { clock }));
|
||||
assert.strictEqual(res.ok, true);
|
||||
});
|
||||
|
||||
test('validateManifest rejects bad version, status, and missing plan_id', () => {
|
||||
const good = buildManifest(baseInput(), { clock: { nowIso: () => '2020-06-15T12:00:00.000Z' } });
|
||||
const badVersion = { ...good, version: '9.9' };
|
||||
assert.strictEqual(validateManifest(badVersion).ok, false);
|
||||
const badStatus = { ...good, status: 'finished' };
|
||||
assert.strictEqual(validateManifest(badStatus).ok, false);
|
||||
const noPlan = { ...good };
|
||||
delete noPlan.plan_id;
|
||||
assert.strictEqual(validateManifest(noPlan).ok, false);
|
||||
});
|
||||
|
||||
// ─── Parsers ──────────────────────────────────────────────────────────────────
|
||||
|
||||
test('parseSbatchParsable parses a bare number and a number;cluster form', () => {
|
||||
assert.deepStrictEqual(parseSbatchParsable('12345'), { ok: true, job_id: '12345' });
|
||||
assert.deepStrictEqual(parseSbatchParsable('12345;mycluster\n'), { ok: true, job_id: '12345' });
|
||||
});
|
||||
|
||||
test('parseSbatchParsable fails closed on empty or non-numeric output', () => {
|
||||
assert.strictEqual(parseSbatchParsable('').ok, false);
|
||||
assert.strictEqual(parseSbatchParsable('Submitted batch job 12345').ok, false, 'non-parsable prose rejected');
|
||||
assert.strictEqual(parseSbatchParsable('abc;cluster').ok, false);
|
||||
});
|
||||
|
||||
test('parseSqueueLine parses "<jobid> <state>" and returns null for malformed', () => {
|
||||
assert.deepStrictEqual(parseSqueueLine('12345 RUNNING'), { job_id: '12345', state: 'RUNNING' });
|
||||
assert.strictEqual(parseSqueueLine('header'), null);
|
||||
assert.strictEqual(parseSqueueLine(''), null);
|
||||
});
|
||||
|
||||
test('parseSacctRow parses [jobid, state] columns', () => {
|
||||
assert.deepStrictEqual(parseSacctRow(['12345', 'COMPLETED']), { job_id: '12345', state: 'COMPLETED' });
|
||||
assert.strictEqual(parseSacctRow(['x']), null);
|
||||
assert.strictEqual(parseSacctRow([], ), null);
|
||||
});
|
||||
|
||||
// ─── writeManifest (fail-closed duplicate guard + fs injection) ────────────────
|
||||
|
||||
function memFs(files = {}) {
|
||||
const store = new Map(Object.entries(files));
|
||||
return {
|
||||
mkdirSync: () => undefined,
|
||||
readdirSync: (d) => {
|
||||
const set = store.get(d);
|
||||
return Array.isArray(set) ? set : [];
|
||||
},
|
||||
readFileSync: (p) => {
|
||||
if (!store.has(p)) { const e = new Error('enoent'); e.code = 'ENOENT'; throw e; }
|
||||
return store.get(p);
|
||||
},
|
||||
writeFileSync: (p, c) => { store.set(p, c); },
|
||||
existsSync: (p) => store.has(p),
|
||||
};
|
||||
}
|
||||
|
||||
test('manifestPath projects to .planning/async-jobs/<job>.json', () => {
|
||||
assert.strictEqual(
|
||||
manifestPath('.planning', '12345'),
|
||||
path.join('.planning', 'async-jobs', '12345.json'),
|
||||
);
|
||||
});
|
||||
|
||||
test('writeManifest writes a new manifest and returns its path', () => {
|
||||
const fs = memFs({ [path.join('.planning', 'async-jobs')]: [] });
|
||||
const clock = { nowIso: () => '2020-06-15T12:00:00.000Z' };
|
||||
const manifest = buildManifest(baseInput(), { clock });
|
||||
const res = writeManifest(manifest, '.planning', { fs, clock });
|
||||
assert.strictEqual(res.ok, true);
|
||||
assert.ok(res.path.endsWith(path.join('async-jobs', '12345.json')));
|
||||
const written = JSON.parse(fs.readFileSync(res.path));
|
||||
assert.strictEqual(written.plan_id, '3.1');
|
||||
});
|
||||
|
||||
test('writeManifest allows updating the SAME job_id (status progression)', () => {
|
||||
const dir = path.join('.planning', 'async-jobs');
|
||||
const existingPath = path.join(dir, '12345.json');
|
||||
const clock = { nowIso: () => '2020-06-15T12:00:00.000Z' };
|
||||
const submitted = buildManifest(baseInput(), { clock });
|
||||
const existing = { ...submitted };
|
||||
const fs = memFs({ [dir]: ['12345.json'], [existingPath]: JSON.stringify(existing) });
|
||||
const running = buildManifest({ ...baseInput(), status: 'running' }, { clock });
|
||||
const res = writeManifest(running, '.planning', { fs, clock });
|
||||
assert.strictEqual(res.ok, true, 'same job_id progression must be allowed');
|
||||
});
|
||||
|
||||
test('writeManifest FAILS CLOSED when a different non-terminal job exists for the same plan_id', () => {
|
||||
// Duplicate-execution guard: a second dispatch for the same plan would
|
||||
// duplicate the external job. Mirror of planning-artifacts.md fail-closed.
|
||||
const dir = path.join('.planning', 'async-jobs');
|
||||
const otherPath = path.join(dir, '99999.json');
|
||||
const clock = { nowIso: () => '2020-06-15T12:00:00.000Z' };
|
||||
const other = buildManifest({ ...baseInput(), job_id: '99999' }, { clock });
|
||||
const fs = memFs({ [dir]: ['99999.json'], [otherPath]: JSON.stringify(other) });
|
||||
const second = buildManifest({ ...baseInput(), job_id: '12345' }, { clock });
|
||||
const res = writeManifest(second, '.planning', { fs, clock });
|
||||
assert.strictEqual(res.ok, false);
|
||||
assert.strictEqual(res.kind, 'duplicate_plan_id');
|
||||
});
|
||||
|
||||
test('writeManifest allows a NEW job once the prior plan_id job is terminal', () => {
|
||||
const dir = path.join('.planning', 'async-jobs');
|
||||
const otherPath = path.join(dir, '99999.json');
|
||||
const clock = { nowIso: () => '2020-06-15T12:00:00.000Z' };
|
||||
const dead = buildManifest({ ...baseInput(), job_id: '99999', status: 'failed', terminal_details: { code: 1 } }, { clock });
|
||||
const fs = memFs({ [dir]: ['99999.json'], [otherPath]: JSON.stringify(dead) });
|
||||
const next = buildManifest({ ...baseInput(), job_id: '12345' }, { clock });
|
||||
const res = writeManifest(next, '.planning', { fs, clock });
|
||||
assert.strictEqual(res.ok, true, 'terminal prior job must not block a new dispatch');
|
||||
});
|
||||
|
||||
test('writeManifest fails closed on a malformed existing manifest', () => {
|
||||
const dir = path.join('.planning', 'async-jobs');
|
||||
const brokenPath = path.join(dir, '12345.json');
|
||||
const fs = memFs({ [dir]: ['12345.json'], [brokenPath]: '{not json' });
|
||||
const clock = { nowIso: () => '2020-06-15T12:00:00.000Z' };
|
||||
const manifest = buildManifest(baseInput(), { clock });
|
||||
const res = writeManifest(manifest, '.planning', { fs, clock });
|
||||
assert.strictEqual(res.ok, false);
|
||||
assert.strictEqual(res.kind, 'malformed_existing');
|
||||
});
|
||||
|
||||
// ─── Property-based (CLAUDE.md: parsers/contracts need a fast-check test) ─────
|
||||
|
||||
test('property: mapSlurmState is total and idempotent over the known alphabet', () => {
|
||||
fc.assert(
|
||||
fc.property(fc.constantFrom(
|
||||
'PENDING', 'CONFIGURING', 'RUNNING', 'COMPLETING', 'COMPLETED',
|
||||
'FAILED', 'CANCELLED', 'TIMEOUT', 'OUT_OF_MEMORY', 'BOOT_FAIL', 'NODE_FAIL', 'PREEMPTED',
|
||||
), (state) => {
|
||||
const a = mapSlurmState(state);
|
||||
const b = mapSlurmState(state);
|
||||
return a !== null && a === b && MANIFEST_STATUS.includes(a);
|
||||
}),
|
||||
{ numRuns: 200 },
|
||||
);
|
||||
});
|
||||
|
||||
test('property: buildManifest -> validateManifest round-trips for valid generated input', () => {
|
||||
fc.assert(
|
||||
fc.property(
|
||||
fc.record({
|
||||
plan_id: fc.stringMatching(/^[0-9]+\.[0-9]+$/),
|
||||
phase: fc.stringMatching(/^[0-9]+$/),
|
||||
job_id: fc.stringMatching(/^[0-9]{1,8}$/),
|
||||
backend: fc.constantFrom('slurm'),
|
||||
submit_command: fc.constantFrom('sbatch --parsable ./run.sh'),
|
||||
status: fc.constantFrom(...MANIFEST_STATUS),
|
||||
expected_artifacts: fc.array(fc.constantFrom('Artifacts/jobs/x/out.h5'), { minLength: 1 }),
|
||||
verification_command: fc.constantFrom('python -m verify.py'),
|
||||
resume_command: fc.constantFrom('/gsd:execute-phase 3'),
|
||||
terminal_details: fc.oneof(fc.constant(null), fc.record({ code: fc.integer() })),
|
||||
}),
|
||||
(input) => {
|
||||
const clock = { nowIso: () => '2020-06-15T12:00:00.000Z' };
|
||||
const td = input.status === 'completed-unverified' ? null : input.terminal_details;
|
||||
const built = buildManifest({ ...input, terminal_details: td }, { clock });
|
||||
return validateManifest(built).ok === true;
|
||||
},
|
||||
),
|
||||
{ numRuns: 100 },
|
||||
);
|
||||
});
|
||||
Reference in New Issue
Block a user