a9f014f9fa
Close the loop so dispatched work actually completes on bl: - on_no_result fallback in followup-sweeper (op + job forms via exec-constrained registry / job-dispatch), seeded on the three autonomy-pulse jobs; fallback_due() dedupes the gravity path - gravity.py: add __main__ entry (loop-remediator.timer was a no-op), 300s re-arm budget, fallback firing + stamp/skip logic - harvester: proof-of-result followups (result_has_evidence), acted-variant NACK, emit-model tool-hint wording - envelope: RESPONSE RULE states the emit model (agents EMIT directives verbatim; runtime executes; works from bare containers) - completion-audit.py + systemd 15-min timer: per-family funnel, swarm drain, followup backlog; digest DM when degraded, 6h heartbeat - tests/test_completion.py (29 tests), JOB-SPEC.md docs Tests: 67/67 focused green (completion + tool_calls).
320 lines
12 KiB
Markdown
320 lines
12 KiB
Markdown
# JOB-SPEC.md: Hosted Job Scheduler and Distributor
|
|
|
|
> **Box is the main surface.** All operator work goes through Box (box.muse-dev.online). The web UI, `box` CLI, and agents share the same API endpoints. No UI-only powers.
|
|
|
|
## Overview
|
|
|
|
A hosted system on bl that processes and distributes jobs to agents via DM.
|
|
Cron jobs and system scripts inject prompts/jobs; agents execute them.
|
|
Loops run as loops on the server, not in agent heads.
|
|
|
|
**Principle:** Limit agency to get smarter. The server decides *what* and *when*;
|
|
agents decide *how*. Deterministic, auditable, debuggable.
|
|
|
|
## Architecture
|
|
|
|
```
|
|
┌──────────────────────────────────────────────────┐
|
|
│ Hosted on bl (systemd timers + scripts) │
|
|
│ │
|
|
│ ┌────────────┐ ┌─────────────┐ │
|
|
│ │ Scheduler │→ │ Dispatcher │→ dm.py send │
|
|
│ │ (cron) │ │ (render+log)│ │
|
|
│ └────────────┘ └─────────────┘ │
|
|
│ ↑ ↓ │
|
|
│ │ ┌─────────────┐ │
|
|
│ └─────────│ Collector │← DM [RESULT] │
|
|
│ │ (log+chain) │ │
|
|
│ └─────────────┘ │
|
|
└──────────────────────────────────────────────────┘
|
|
```
|
|
|
|
## Components
|
|
|
|
### 1. Job Definition (YAML)
|
|
|
|
Location: `/home/super/Projects/NetVM/jobs/<name>.yaml`
|
|
|
|
```yaml
|
|
name: board-watch
|
|
description: "Check board for new posts every 5 minutes"
|
|
schedule: "*/5 * * * *" # cron format
|
|
agent: muse # which agent executes
|
|
sidechat:
|
|
create: false # use main chat
|
|
# OR:
|
|
# create: true
|
|
# name_template: "job-{name}-{date}"
|
|
# reuse_pattern: "job-{name}-*" # for chaining
|
|
prompt_template: |
|
|
Check the board for posts since {last_run}.
|
|
Summarize new items in 3 bullet points.
|
|
Reply with [RESULT {job_id}] followed by your summary.
|
|
timeout: 300 # seconds before marking failed
|
|
on_failure: retry # retry | alert | ignore
|
|
chain_next: null # job to trigger after success
|
|
```
|
|
|
|
**Fields:**
|
|
- `name`: Unique job identifier (used in logs, sidechat names)
|
|
- `schedule`: Cron expression (systemd timer or cron)
|
|
- `agent`: Target agent (`muse`, `pip`, `646`, `opm`)
|
|
- `sidechat`: Side chat configuration (see below)
|
|
- `prompt_template`: Jinja-style template with `{variables}`
|
|
- `timeout`: Max seconds to wait for result
|
|
- `on_failure`: What to do on timeout/failure
|
|
- `chain_next`: Next job to trigger (for chains)
|
|
|
|
### 2. Scheduler
|
|
|
|
Uses systemd timers (preferred) or cron. Each job gets a timer unit.
|
|
|
|
**Systemd timer example:**
|
|
```ini
|
|
# /etc/systemd/user/job-board-watch.timer
|
|
[Unit]
|
|
Description=Run board-watch job every 5 minutes
|
|
|
|
[Timer]
|
|
OnCalendar=*:0/5
|
|
Persistent=true
|
|
|
|
[Install]
|
|
WantedBy=timers.target
|
|
```
|
|
|
|
**Service:**
|
|
```ini
|
|
# /etc/systemd/user/job-board-watch.service
|
|
[Unit]
|
|
Description=Dispatch board-watch job
|
|
|
|
[Service]
|
|
Type=oneshot
|
|
ExecStart=/home/super/Projects/NetVM/bin/job-dispatch.sh board-watch
|
|
```
|
|
|
|
### 3. Dispatcher (`bin/job-dispatch.sh`)
|
|
|
|
Responsibilities:
|
|
1. Load job YAML
|
|
2. Render prompt template with variables (`{job_id}`, `{last_run}`, `{date}`, etc.)
|
|
3. Create side chat if specified
|
|
4. Send DM via `dm.py`:
|
|
- Format: `[JOB {job_id}] {rendered_prompt}`
|
|
- Use `--raw` if prompt is pre-signed
|
|
5. Log to `job-log.jsonl`: `{job_id, job_name, agent, sidechat, sent_at}`
|
|
6. If `chain_next`, schedule the next job (or trigger immediately on result)
|
|
|
|
**Job ID format:** `{name}-{YYYYMMDD-HHMMSS}-{short_uuid}`
|
|
Example: `board-watch-20261004-023000-a1b2c3d4`
|
|
|
|
### 4. Agent Job Handler (Convention)
|
|
|
|
Agents MUST recognize job DMs and respond in format.
|
|
|
|
**Job DM format:**
|
|
```
|
|
[JOB board-watch-20261004-023000-a1b2c3d4] Check the board for posts
|
|
since 2026-10-04T02:25:00Z. Summarize new items in 3 bullet points.
|
|
Reply with [RESULT board-watch-20261004-023000-a1b2c3d4] followed by
|
|
your summary.
|
|
```
|
|
|
|
**Agent responsibilities:**
|
|
1. Parse `[JOB {job_id}]` from DM
|
|
2. Execute the prompt
|
|
3. If `sidechat` specified, work in that side chat
|
|
4. Reply via DM with: `[RESULT {job_id}] {result_text}`
|
|
5. If unable, reply: `[RESULT {job_id}] FAILED: {reason}`
|
|
|
|
**Teaching:** New agents get `JOB-HANDLER.md` with examples. The format is
|
|
simple enough to learn from 2-3 examples.
|
|
|
|
### 5. Collector
|
|
|
|
Watches for `[RESULT {job_id}]` in DM logs or via `dm.py log`.
|
|
|
|
**Responsibilities:**
|
|
1. Parse result DMs
|
|
2. Log to `job-log.jsonl`: `{job_id, completed_at, success, result_preview}`
|
|
3. If `chain_next` specified and result was success, trigger next job
|
|
4. On timeout (no result within `timeout`), mark failed, apply `on_failure`
|
|
|
|
**Timeout handling:** A background sweeper checks for jobs with `sent_at`
|
|
older than `timeout` and no result. Marks them failed.
|
|
|
|
## Side Chat Integration
|
|
|
|
### Job → Side Chat Mapping
|
|
|
|
Jobs can specify side chat behavior:
|
|
|
|
**Option A: No side chat** (use main chat)
|
|
```yaml
|
|
sidechat:
|
|
create: false
|
|
```
|
|
|
|
**Option B: Create new side chat per run**
|
|
```yaml
|
|
sidechat:
|
|
create: true
|
|
name_template: "job-{name}-{date}" # e.g., "job-board-watch-20261004"
|
|
```
|
|
|
|
**Option C: Reuse side chat (for chains)**
|
|
```yaml
|
|
sidechat:
|
|
create: false
|
|
reuse_pattern: "job-{name}-*" # find most recent
|
|
# OR:
|
|
# reuse_name: "{prev_job_sidechat}" # from chain
|
|
```
|
|
|
|
### Side Chat as Job Workspace
|
|
|
|
When `create: true`:
|
|
1. Dispatcher calls `dm.py` sidechat create (via `muse-chat-api.py`)
|
|
2. Gets the side chat name/ID
|
|
3. Includes it in the job DM: "Work in side chat 'job-board-watch-20261004'"
|
|
4. Logs the mapping: `{job_id → sidechat_name}`
|
|
5. Agent does all work in that side chat (full context, isolated)
|
|
|
|
**Benefits:**
|
|
- Each job run has isolated context
|
|
- The side chat IS the audit log
|
|
- Chains: Job B can continue in Job A's side chat
|
|
- No cross-talk between concurrent jobs
|
|
|
|
## Logging
|
|
|
|
### `job-log.jsonl` (on bl)
|
|
|
|
Append-only, one JSON per line:
|
|
```json
|
|
{"ts": "2026-10-04T02:30:00Z", "type": "job_sent", "job_id": "...", "job_name": "board-watch", "agent": "muse", "sidechat": null}
|
|
{"ts": "2026-10-04T02:32:15Z", "type": "job_result", "job_id": "...", "success": true, "duration_s": 135}
|
|
{"ts": "2026-10-04T02:35:00Z", "type": "job_timeout", "job_id": "...", "job_name": "board-watch"}
|
|
```
|
|
|
|
### `sidechat-log.jsonl` (on bl)
|
|
|
|
From `sidechat_manager.py`:
|
|
```json
|
|
{"ts": "...", "agent": "muse", "op": "create", "details": "job-board-watch-20261004"}
|
|
{"ts": "...", "agent": "muse", "op": "navigate", "details": "main"}
|
|
```
|
|
|
|
## Rate Limiting
|
|
|
|
Use the shared `bin/rate_limiter.py` module:
|
|
```python
|
|
from rate_limiter import rate_limit_wait
|
|
rate_limit_wait(agent) # blocks if too frequent
|
|
```
|
|
|
|
Default: 1 DM per 3s per agent, burst 5, max 20/min. Prevents accidental
|
|
spam if a job misfires in a loop.
|
|
|
|
## Security
|
|
|
|
- Job definitions are in git (auditable)
|
|
- Only operators can create/edit jobs (file permissions)
|
|
- Agents cannot create jobs (they execute, not schedule)
|
|
- DMs are logged (dm-log.jsonl) for audit
|
|
- Side chats are per-job, not shared across trust boundaries
|
|
|
|
## Completion Enforcement
|
|
|
|
Sending a job DM does not complete work. Four mechanisms close the loop on bl:
|
|
|
|
### 1. `on_no_result` fallback (server-side guarantee)
|
|
|
|
A job may declare a fallback effect that runs when its followup reaches
|
|
terminal expiry with no agent result (`bin/followup-sweeper.py`):
|
|
|
|
```json
|
|
"on_no_result": {"op": "swarm.spawn", "args": {"count": 2, "task": "...", "label": "..."}}
|
|
"on_no_result": {"job": "<another-job-name>"}
|
|
```
|
|
|
|
- `op` form: validated + built through the `exec-constrained` op registry
|
|
(same validators the daemon uses), then executed as a subprocess.
|
|
- `job` form: dispatches the named job via `job-dispatch.py` with
|
|
`CHAIN_PREV_*` timeout context.
|
|
- Outcomes log as `fallback_executed` / `fallback_failed` in `job-log.jsonl`
|
|
and stamp `rec["fallback"]` on the followup record. Failures never break
|
|
the sweep. Jobs without the key behave exactly as before.
|
|
- Firing paths (either; never both): the sweeper terminal branch
|
|
(`nudges_sent >= nudges_allowed`, expired) and `gravity.py:remediate_breaks`
|
|
(same terminal condition, runs on `loop-remediator.timer` ~every 15m).
|
|
`fallback_due()` dedupes: fires once per record, retries a failed attempt
|
|
after 1h. Gravity also revives the local sweep loop (it previously had no
|
|
`__main__`, so the timer was a no-op) with a 300s re-arm budget (was 10s,
|
|
which strangled every sweep mid-first-send — median nudge send is ~10s).
|
|
- Split-brain note: the VM board sweeper (`box-request-sweeper.timer`) sends
|
|
the live `[NUDGE <id>]` DMs from remote `dm_followup` requests; the local
|
|
sweeper sends `[nudge N/M]` from `followups.json`. Both fire per deadline
|
|
until a reply resolves both sides. Accept the duplicate nudge as the cost
|
|
of a guaranteed local path; cross-system dedup is future work.
|
|
- Seeded on: `autonomy-pulse-646`, `autonomy-pulse-pip`, `autonomy-pulse-opm`
|
|
(each spawns 2 standing-work subagents if the agent naps through the pulse).
|
|
- Race note: set the job's followup timeout longer than the expected work
|
|
time, or a slow-but-working agent can double-fire alongside the fallback.
|
|
|
|
### 2. NACK for directive-less replies (existing, sharpened)
|
|
|
|
`response-harvester.py:maybe_nudge_untagged_sidechat` already rejects
|
|
conversational replies in job-backed threads (1-turn strict nudge, then a
|
|
5-minute escalation timer to opm; tracked in
|
|
`conversation-nudge-tracker.json`). Two refinements:
|
|
|
|
- Acted-variant: when the agent emitted directives but never closed with
|
|
`[RESULT]`, the nudge acknowledges the action and demands the close
|
|
instead of crying "commentary rejected".
|
|
- Emit-model wording: envelope (`prompt_envelope._response_rule`), tool
|
|
hint, and strict nudge now state plainly that agents EMIT `[TOOL]` /
|
|
`[DM]` lines verbatim and the Box runtime on bl executes them — this
|
|
works from containers with no box CLI or tmux socket. (Root cause of the
|
|
2026-10-06 zero-`tool_exec` stretch: agents believed they had to execute
|
|
tools locally and declined for lack of a "container equivalent".)
|
|
|
|
### 3. Proof-of-result followups
|
|
|
|
A success `[RESULT]` with no checkable artifact (swarm/timer IDs, UUIDs,
|
|
paths, `n/m` completion counts — see `result_has_evidence`) triggers a
|
|
one-shot `followup.create` (+30m, same thread) asking for the evidence, and
|
|
logs `proof_requested`. One per (thread, job). Failures and declines skip
|
|
proof (they already chain via `on_failure`).
|
|
|
|
### 4. Completion auditor (15-minute timer)
|
|
|
|
`bin/completion-audit.py` via `systemd/completion-audit.timer`
|
|
(`OnCalendar=*:4/15`, installed to `/etc/systemd/system`):
|
|
|
|
- Computes the funnel per job family over `--hours` (default 24) from
|
|
`job-log.jsonl`: sent / dispatched / `tool_exec` / results, plus
|
|
`fallback_*`, `proof_requested`, and `job_failed` counts.
|
|
- Adds point-in-time swarm slot drain (flags running slots idle >90m) and
|
|
followup backlog (pending / overdue / escalated).
|
|
- Writes `logs/completion-audit-<ts>.json` + `logs/completion-audit-latest.json`.
|
|
- Posts the digest to `opm` / `heartbeat` only when DEGRADED
|
|
(silent families, failures, tool errors, stale slots, overdue/escalated
|
|
followups); when green, posts a heartbeat at most every 6h
|
|
(`logs/completion-audit-state.json`). Silent-when-healthy otherwise.
|
|
- Read-only except the digest DM and its own log/state files.
|
|
|
|
## Future Expansions
|
|
|
|
1. **Conditional jobs**: Run Job B only if Job A succeeds with specific output
|
|
2. **Parallel jobs**: Fan-out to multiple agents, collect all results
|
|
3. **Human approval**: Certain jobs require human sign-off before dispatch
|
|
4. **Web UI**: View job status, logs, side chats from box.muse-dev.online
|
|
5. **Agent-proposed jobs**: Agents can suggest jobs via `[PROPOSE]` DM, human approves
|
|
|
|
## Status
|
|
|
|
Spec v0.1. Owner: operator-main. Sanctioned by the human 2026-10-04.
|
|
Implementation: dispatcher script + example jobs pending.
|