mirror of
https://github.com/razzant/ouroboros.git
synced 2026-10-03 04:07:04 +00:00
Keep live turns visible and honor worker startup progress
Reserve chat composer space with a flex item for native WebKit. Combine queue and live activity identities through the existing task controls. Give a worker with its own early progress one readiness extension bounded to 300 seconds from spawn. Wait for the actual chat socket before the large-attachment browser test sends its files. Co-authored-by: Ouroboros <311266734+ouroboros-agent@users.noreply.github.com>
This commit is contained in:
parent
fe251b89e8
commit
7e4a7796b7
20 changed files with 588 additions and 31 deletions
|
|
@ -174,7 +174,7 @@ Dashboard groups Logs, Evolution, Costs, Updates and Activity under the common t
|
|||
|
||||
Logs merges live WebSocket log frames with bounded REST backfill from the events, tools, progress and supervisor logs: chronological order, dedupe against live and reconnect overlap, bounded grouped task cards, raw record on demand, and shared category/severity/review presentation from `log_events.js`. A failed backfill names the unavailable history sources in a separate status row while live events continue. Clearing the visible panel does not delete the underlying logs.
|
||||
|
||||
Activity shows running and pending queue entries, the consciousness alarm clock (enabled, status line, autonomy level, next and last wake-up with its outcome, the rolling-24h allowance and how many of its tasks run) and scheduled work, with mechanical controls only where that surface is authoritative: typed cascade cancellation, start/stop for background consciousness, enable/disable/delete for owner-managed schedules. Each section reports its own failed read; a failed queue read neither erases the background/schedule sections nor claims an empty queue. A schedule reconciled from a skill manifest is read-only here — a direct edit would be overwritten by the skill lifecycle and falsely appear durable.
|
||||
Activity combines running and pending queue entries with the live activity census from the same refresh's `/api/state`: an identity already in the queue keeps its title, runtime and budget controls; a census-only direct turn shows its kind, phase, elapsed time and the shared task controls. A partial census supplies positive rows but cannot prove nothing is running. The view also shows the consciousness alarm clock (enabled, status line, autonomy level, next and last wake-up with its outcome, the rolling-24h allowance and how many of its tasks run) and scheduled work, with mechanical controls only where that surface is authoritative: typed cascade cancellation, start/stop for background consciousness, enable/disable/delete for owner-managed schedules. Each section reports its own failed read; a failed queue read neither erases the background/schedule sections nor claims an empty queue. A schedule reconciled from a skill manifest is read-only here — a direct edit would be overwritten by the skill lifecycle and falsely appear durable.
|
||||
|
||||
Costs is a projection of the physical-attempt ledger distinguishing confirmed, reserved, unresolved upper-bound, unknown/unmetered usage, open rows and finality; unavailable data renders as unavailable, not `$0`. Breakdowns by model, key, model category and task category are views over the same ledger. Total budget can hot-apply; the per-task value is a hard cost cap over the whole root tree of the next task, not an own-task soft warning (§7 Default settings), and raising a cap does not resume work that already finalized or paused.
|
||||
|
||||
|
|
|
|||
|
|
@ -43,7 +43,7 @@ The event bus is process-lifetime rather than worker-generation-lifetime: respaw
|
|||
|
||||
Heartbeat and progress are different evidence: a heartbeat proves a process or loop is alive; owner-visible progress and model-usage events prove the task advanced. Fresh descendant progress or queued descendants can keep an orchestrator alive, while an explicit deadline, absolute ceiling, cancellation and budget stop remain hard. After the typed finalization episode (§6), timeout handling freezes its decision under the queue lock, marks the worker `reaping`, and hands kill, join, salvage, retry and respawn to the single off-loop reaper; an orchestrator with live descendants is not blindly retried, because a retry would replay its plan and spawn a competing tree. No retry or new assignment may occupy a timed-out slot until the original process is provably dead: if kill and join cannot establish death, the reaper keeps a low-rank RUNNING result and the `reaping` slot, emits a visible `task_reaper_wedged` receipt and restart hint, and writes no terminal, `task_done`, retry or respawn — one slot is sacrificed rather than letting a still-running process race a replacement and overwrite its result; the next supervisor generation reconciles the record after old-generation process custody.
|
||||
|
||||
A spawned or respawned slot is not assignable until its child's PID-bound `worker_ready` row arrives (`supervisor/worker_pool_lifecycle.py`). After `WORKER_READY_MAX_ATTEMPTS` windows of `WORKER_READY_WINDOW_SEC` (`runtime_limits.py`), `Worker.readiness_exhausted` is final for that exact slot — late events cannot reopen it. Total exhaustion, distinguished from busy/booting/reaping capacity and from a live owner-wait stack, closes pooled ingress (owner `/review` included) without blocking direct chat/control or boot/update recovery; once RUNNING completion custody has settled, `disable_exhausted_worker_pool` fails unstarted PENDING work honestly with a Restart hint, and a new task cannot clear the latch. Readiness stays separate from liveness and task idle time; a watcher error releases only still-booting, non-exhausted slots to the crash detector (`worker_ready_released`). Linux workers use forkserver; macOS and Windows use spawn.
|
||||
A spawned or respawned slot is not assignable until its child's PID-bound `worker_ready` row arrives (`supervisor/worker_pool_lifecycle.py`). A live child's own `worker_starting` row, emitted before extension loading and agent construction, permits one extension of `WORKER_READY_WINDOW_SEC` to `WORKER_READY_CEILING_SEC` (300 seconds from birth, both in `runtime_limits.py`); foreign or pre-spawn rows cannot extend another slot. `worker_ready_window_extended` records that decision. A silent child keeps the original window, and logging failure cannot block startup. After `WORKER_READY_MAX_ATTEMPTS` failed attempts, `Worker.readiness_exhausted` is final for that exact slot — late events cannot reopen it. Total exhaustion, distinguished from busy/booting/reaping capacity and from a live owner-wait stack, closes pooled ingress (owner `/review` included) without blocking direct chat/control or boot/update recovery; once RUNNING completion custody has settled, `disable_exhausted_worker_pool` fails unstarted PENDING work honestly with a Restart hint, and a new task cannot clear the latch. Readiness stays separate from liveness and task idle time; a watcher error releases only still-booting, non-exhausted slots to the crash detector (`worker_ready_released`). Linux workers use forkserver; macOS and Windows use spawn.
|
||||
|
||||
Unexpected worker death reserves exact custody under the queue lock and enqueues `confirmed_dead_worker` on the reaper (`worker_health.recover_confirmed_dead_worker`). A saved terminal source wins even after signal death; unknown or incomplete file publication keeps the same job (`TerminalFileRecoveryPending`); only confirmed absence of one reaches the crash policy: a signal is an infrastructure failure, an otherwise eligible non-signal crash retries within `QUEUE_MAX_RETRIES`, preserving owner-wait replay restrictions and cost. A crash storm suppresses respawn while terminal sources settle, then its fence stops pooled admission; direct chat stays available. Startup runs the same terminal-file recovery in `_run_supervisor` after process custody and before `_startup_prune_sweeps` (the no-provider lifespan runs it too, spawning nothing); unknown or still-live ownership defers it rather than racing a writer, and any unresolved or protected source, or an ownership/read error, sets `preserve_task_sources`, skipping task-drive deletion for that pass. For older canonical scheduled rows, `_recover_terminal_task_files` restores that start binding only from a known non-direct child's positive running/started-at record when the existing fresh-queue and later-worker-boot checks prove it orphaned, with no pending queue owner or active cancel; the normal orphan reconciler and terminal guards retain authority, without resuming work. The recovery report includes `rebound`.
|
||||
|
||||
|
|
@ -52,4 +52,3 @@ Startup and throttled maintenance reconcile three residue classes. Process custo
|
|||
Cooperative project checkpointing has two equivalent quiescence triggers: a host-minted genesis or cooperative tree is checked when its root settles with no live descendants, and again when the last child settles beneath an already-terminal root — a root-scope budget stop terminalizes the root before its children, so a root-only trigger would see a live tree once and never return. The bounded git chain runs on a daemon thread, revalidates quiescence under the queue lock immediately before mutation, and replays a trigger that arrives during an in-flight check. Only host-minted project roots are eligible: owner-attached folders are never auto-committed, credential-shaped files stay excluded and disclosed, and every material success, skip or error receives a durable receipt.
|
||||
|
||||
The bridge recognizes `/panic`, `/restart`, `/review`, `/evolve [on|off]`, `/bg [start|stop|status]` and `/status`; all other text enters ordinary agent routing. External transports may invoke these commands only with positive owner identity and a transport-specific owner-chat binding, and the commands reuse runtime-mode, queue, cancellation and typed-result authority rather than implementing parallel control paths. Runtime logs rotate on the same supervisor tick and archive readers preserve their retained timelines; only explicitly isolated devtool roots may use the narrow rotation sentinel from §1.
|
||||
|
||||
|
|
|
|||
|
|
@ -803,12 +803,14 @@ and what enforces each.
|
|||
`SETTINGS_DEFAULTS`, the clamped getter in `runtime_limits.py`, both re-exported
|
||||
through `ouroboros.config`, the one import surface; register the env key; no magic
|
||||
wait numbers at call sites (`tests/test_timeout_policy.py`).
|
||||
- Worker readiness keeps its own `WORKER_READY_WINDOW_SEC` and
|
||||
- Worker readiness keeps `WORKER_READY_WINDOW_SEC`, `WORKER_READY_CEILING_SEC` and
|
||||
`WORKER_READY_MAX_ATTEMPTS` in `runtime_limits.py`, re-exported by config
|
||||
(ARCHITECTURE §5 "Supervisor Loop"). Reuse the lifecycle-owned execution-state
|
||||
reader (workers facade) at reserve, final enqueue and snapshot; keep the separate
|
||||
repository-writer policy at public admission and the boot/update exceptions;
|
||||
readiness, process liveness and idle deadlines stay independent; a failed write keeps
|
||||
the child's own `worker_starting` row before extension loading permits one
|
||||
readiness extension to 300 seconds from birth, never a sliding deadline or a
|
||||
fresh window at observation. Readiness, process liveness and idle deadlines stay independent; a failed write keeps
|
||||
terminalization retry, never a false Done or a fresh startup budget.
|
||||
- Nested process wrappers are ordered, never tied: provider bound before its killable
|
||||
child, child before the generic ToolEntry envelope (the settlement margin from
|
||||
|
|
@ -1062,4 +1064,3 @@ Enforcement: review-only — CHECKLISTS item 2(f) scores the no-`[:N]` rule in
|
|||
commit review.
|
||||
|
||||
---
|
||||
|
||||
|
|
|
|||
|
|
@ -28,7 +28,7 @@ This chapter owns the engineering rules that preserve the visual and interaction
|
|||
- **Narration leads; routine execution evidence stays compact; exceptions keep their explanation and controls** (semantics: DESIGN "Conversation activity block"; mechanism: ARCHITECTURE §3 "Direct turns and the activity block"). A note is promoted on the typed `narration` fact its producer stamps, never on a client reading of its text; execution evidence is one bounded row per block, never a row per event and never a client tool-name list; host snapshots of the same turn merge field-wise, because an absent field is unknown rather than zero. A change that draws one row or line per event is shown at a realistic burst size (a multi-call turn, collapsed and expanded, at desktop and phone width) before it is accepted; a two-event fixture proves layout only for itself. Enforced by `web/tests/chat_activity_block.test.js`, `web/tests/chat_addressing_card.test.js`, `web/tests/progress_narration_voice.test.js` and `tests/test_progress_narration_voice.py`; the visual check follows "Responsive and accessible behavior" below.
|
||||
- **Executor presentation** consumes the existing task/run attempt facts (shape and producer: ARCHITECTURE §3 "Child cards and executor presentation"). Keep `executor_observation` event-local through Agent, supervisor delivery, progress history and both Chat metadata paths; ordinary coordinator notes inherit none. Do not borrow a final-attempt model or parse progress prose to fill an absent live observation. Label last activity separately from current computation, configured/coordinator model and settled observed-model history. Preserve terminal-only `execution_evidence` and `actual_substrate`, with no new poller or execution-state store. Render the chip/model facts through `harness_presentation.js::executorIdentityMarkup`, keep execution-evidence selection in `log_events.js`, and add no second label builder in Chat. Tests: `tests/test_executor_observation.py`, `web/tests/wire_contract.test.js`.
|
||||
- **Model-call provenance** owned by the task survives result storage, copy-back and terminal/history rendering. Show the last usable solve response separately from initial routing, executor observations and final-answer authorship; a post-task or cost-only update cannot erase it. Fan-out counters describe emissions and wall-clock intervals, never inferred execution waves. Enforced by `tests/test_task_model_execution.py` (storage, the terminal projection and the fan-out key vocabulary) and `web/tests/task_result_reconciliation.test.js` (terminal and history rendering).
|
||||
- **Chat viewport invariant.** Sample live-edge intent before an ordinary transcript mutation — native scroll anchoring is not proof that the owner's visible message stays stable, so focused regressions disable it. The stable-viewport seam (ARCHITECTURE §3 "Timeline ownership and ordering") decides when the transcript follows the bottom; otherwise preserve the visible keyed message, nested card or Reviews anchor. Route late application-controlled DOM writes through that seam, and keep awaited Load-older, reconnect reconciliation and cross-instance restoration as explicit lifecycle transactions. Browser coverage is chosen by risk; this WebKit-sensitive contract requires the engines exercised by its marker-gated UI smoke.
|
||||
- **Chat viewport invariant.** Sample live-edge intent before an ordinary transcript mutation — native scroll anchoring is not proof that the owner's visible message stays stable, so focused regressions disable it. The stable-viewport seam (ARCHITECTURE §3 "Timeline ownership and ordering") decides when the transcript follows the bottom; otherwise preserve the visible keyed message, nested card or Reviews anchor. Route late application-controlled DOM writes through that seam, and keep awaited Load-older, reconnect reconciliation and cross-instance restoration as explicit lifecycle transactions. Main and Project transcripts reserve their overlaid composer with one non-shrinking flex spacer, because bottom padding can fall outside native WKWebView's scrollable overflow; the keyboard-in-flow layout removes that item and its gap. Preserve each surface's measured safe-area reserve. Browser coverage is chosen by risk; this WebKit-sensitive contract requires the engines exercised by its marker-gated UI smoke, plus the system WKWebView when the defect is specific to that engine.
|
||||
|
||||
History pages and reconnect merge into the existing keyed card/row owners. Preserve the actual selected/focused/expanded nodes — rebuilding an equivalent node is not preservation. Timeline patches compare generated markup so unchanged enhanced markdown keeps its controls; the Reviews reconciler separately owns lazy attempt-detail state and cannot replace that behaviour. Physical source identity orders equal-time archive rows without rewriting JSONL. A historical frame never grants current activity or replaces newer terminal evidence. Eviction releases only its own page's media, markdown and decision views and protects visible reading, focus and selection. Exact page handles retain return navigation; read gaps and sparse empty pages never become false EOF; readable recent rows survive an unavailable archive with an explicit gap/retry and no fabricated physical cursor. Flush historical timeline changes once per card and skip idle scroll cleanup when no work is pending.
|
||||
|
||||
|
|
@ -79,4 +79,3 @@ Enforcement: `tests/test_widgets_ui_static.py` at commit tier; in the release-ti
|
|||
Use `ouroboros.server_web.read_author_kit_assets(request.app.state.repo_dir)` to read the fixed installed `web/ui.css` and `web/modules/ui_primitives.js` sources for an author-owned page (what the kit is and is not: ARCHITECTURE §3 "Author UI kit"). Resolve at the page/kit GET that serves a new mount, not at extension registration; the request root is propagated by both in-process and out-of-process dispatch. No bundle cache, new endpoint or auth exception belongs in the helper.
|
||||
|
||||
`docs/examples/author_ui_kit/` holds the two ordinary extension recipes. Revoke temporary Blob URLs; keep the kit optional and author-overridable, with no theme poller or forced remount; add no opaque `/static` request, bridge message or widget schema flag. `test_author_ui_kit.py` and `test_author_ui_kit_browser.py` cover source-root delivery and actual framed consumers; they do not certify an arbitrary author's CSP or application.
|
||||
|
||||
|
|
|
|||
|
|
@ -2,14 +2,14 @@
|
|||
|
||||
AST-derived inventory of compatibility facades, regenerated by `python scripts/regenerate_inventories.py`. Do not edit. A facade row is any runtime module whose top-level `from <population module> import ...` statements carry the `noqa: F401` re-export marker — the codebase's declared "this binding exists for its binding, not for this module's own use" convention (reference FACADE_CONSUMERS method). Leaf domains come from `ouroboros/domains.toml`; a leaf outside the facade's domain is marked ✗ (that edge also appears in the manifest's pinned direction matrix). `tests/test_generated_inventories.py` pins byte-identity, so any re-export surface change must regenerate this file.
|
||||
|
||||
- facade modules: **58**; marked re-export bindings: **2357**; cross-domain facade→leaf pairs: **131**
|
||||
- facade modules: **58**; marked re-export bindings: **2358**; cross-domain facade→leaf pairs: **131**
|
||||
|
||||
| facade | domain | bindings | leaves |
|
||||
|---|---|---:|---|
|
||||
| `launcher.py` | D18 | 2 | `ouroboros/launcher_windows_runtime.py` (2) |
|
||||
| `ouroboros/agent.py` | D01 | 32 | `ouroboros/agent_dispatch.py` (15)<br>`ouroboros/agent_startup_checks.py` (4)<br>`ouroboros/config.py` (2 ✗D12)<br>`ouroboros/subagent_dispatch_notes.py` (4 ✗D07)<br>`ouroboros/subagents.py` (7 ✗D07) |
|
||||
| `ouroboros/agent_task_pipeline.py` | D01 | 28 | `ouroboros/dialogue_provenance.py` (2 ✗D15)<br>`ouroboros/post_task_synthesis.py` (11)<br>`ouroboros/synthesis_cost_text.py` (5)<br>`ouroboros/task_finalization.py` (10) |
|
||||
| `ouroboros/config.py` | D12 | 130 | `ouroboros/model_slots.py` (17)<br>`ouroboros/provider_models.py` (6 ✗D02)<br>`ouroboros/review_model_routes.py` (10)<br>`ouroboros/runtime_limits.py` (59)<br>`ouroboros/settings_defaults.py` (19)<br>`ouroboros/settings_integrity.py` (4)<br>`ouroboros/settings_scales.py` (13)<br>`ouroboros/update_channels.py` (2) |
|
||||
| `ouroboros/config.py` | D12 | 131 | `ouroboros/model_slots.py` (17)<br>`ouroboros/provider_models.py` (6 ✗D02)<br>`ouroboros/review_model_routes.py` (10)<br>`ouroboros/runtime_limits.py` (60)<br>`ouroboros/settings_defaults.py` (19)<br>`ouroboros/settings_integrity.py` (4)<br>`ouroboros/settings_scales.py` (13)<br>`ouroboros/update_channels.py` (2) |
|
||||
| `ouroboros/context.py` | D03 | 4 | `ouroboros/context_runtime_facts.py` (4) |
|
||||
| `ouroboros/delegate_custody.py` | D07 | 10 | `ouroboros/delegate_custody_reconcile.py` (9)<br>`ouroboros/delegate_evidence.py` (1) |
|
||||
| `ouroboros/extension_loader.py` | D14 | 99 | `ouroboros/contracts/plugin_api.py` (7 ✗D19)<br>`ouroboros/extension_child_catalog.py` (8)<br>`ouroboros/extension_companion.py` (3)<br>`ouroboros/extension_import_staging.py` (6)<br>`ouroboros/extension_isolated_deps.py` (4)<br>`ouroboros/extension_liveness.py` (8)<br>`ouroboros/extension_plugin_api.py` (6)<br>`ouroboros/extension_registry_state.py` (20)<br>`ouroboros/extension_surface_names.py` (12)<br>`ouroboros/extension_ui_validation.py` (5)<br>`ouroboros/gateway/host_service.py` (1 ✗D11)<br>`ouroboros/provider_models.py` (1 ✗D02)<br>`ouroboros/skill_loader.py` (13)<br>`ouroboros/skill_token.py` (1)<br>`ouroboros/tools/skill_exec.py` (1)<br>`ouroboros/utils.py` (3 ✗D18) |
|
||||
|
|
|
|||
|
|
@ -96,7 +96,7 @@ from ouroboros.review_model_routes import (
|
|||
from ouroboros.runtime_limits import (
|
||||
WORKER_SPAWN_GRACE_SEC, # noqa: F401
|
||||
WORKER_READY_WINDOW_SEC, # noqa: F401
|
||||
WORKER_READY_MAX_ATTEMPTS, # noqa: F401
|
||||
WORKER_READY_MAX_ATTEMPTS, WORKER_READY_CEILING_SEC, # noqa: F401
|
||||
EXTENSION_STREAM_CHUNK_BYTES, # noqa: F401
|
||||
EXTENSION_CHILD_CLEANUP_GRACE_SEC, # noqa: F401
|
||||
NESTED_SETTLEMENT_MARGIN_SEC, # noqa: F401
|
||||
|
|
|
|||
|
|
@ -51,13 +51,15 @@ WS_RELAY_REFILL_PER_SEC = 1.0
|
|||
# detector counts dead workers (up to ~60s to init: spawn + pip); workers.py binds it as `_SPAWN_GRACE_SEC`, the extension import-staging sweep reads it too.
|
||||
WORKER_SPAWN_GRACE_SEC = 90.0
|
||||
# Readiness window for ONE spawned/respawned slot: unassignable until the child's own `worker_ready` row lands; alive
|
||||
# but silent past this = torn down and replaced. Sized to the spawn grace (the pool's existing init budget): a warm
|
||||
# forkserver child boots in ~3-4s (G13 mock lane: 3.5-4.9s startup, 2.5-3.2s respawn), a cold 4-vCPU CI runner well under 60s (its 21-scenario mock lane runs in ~80s), and the E2E
|
||||
# scenarios wait 240s per task, so a wedged child is a fast, named failure. A contract distinct from process liveness
|
||||
# but silent past this = torn down and replaced. A child's own entry progress permits one longer window for
|
||||
# expensive extension loading; an empty mock install does not establish production startup latency.
|
||||
# Readiness is a contract distinct from process liveness
|
||||
# (`proc.is_alive`, worker_health.py) and from the task idle rail (queue_timeouts.py): a deadlocked child is alive.
|
||||
WORKER_READY_WINDOW_SEC = 90.0
|
||||
# Consecutive readiness failures of one slot before it is parked and reported (three strikes, like the crash-storm fence).
|
||||
WORKER_READY_MAX_ATTEMPTS = 3
|
||||
# One extension for a child that wrote its own entry progress, measured from birth, never from the last poll.
|
||||
WORKER_READY_CEILING_SEC = 300.0
|
||||
|
||||
|
||||
def _clamped_number_setting(key: str, *, low, high=float("inf"), cast=float):
|
||||
|
|
|
|||
|
|
@ -4,6 +4,7 @@ A spawned or respawned slot is installed unassignable (``reaping=True``) and ope
|
|||
only when the child's own ``worker_ready`` row is observed, which is also where the
|
||||
SHA it booted is verified; a child alive but silent past the readiness window is
|
||||
torn down and replaced through the same respawn path, a bounded number of times.
|
||||
Its own entry-progress row permits one extension, still bounded from its birth.
|
||||
The pids workers ran under are recorded durably so an orphan surviving a restart
|
||||
can be reaped; a replaced worker's queue is closed under the lock before the new
|
||||
one takes its slot.
|
||||
|
|
@ -24,7 +25,7 @@ import sys
|
|||
import threading
|
||||
import time
|
||||
from typing import Any, Dict, List, Optional, Tuple
|
||||
from ouroboros.config import WORKER_READY_MAX_ATTEMPTS, WORKER_READY_WINDOW_SEC
|
||||
from ouroboros.config import WORKER_READY_CEILING_SEC, WORKER_READY_MAX_ATTEMPTS, WORKER_READY_WINDOW_SEC
|
||||
from ouroboros.review_owner_custody import reconcile_confirmed_dead_review_owners
|
||||
from supervisor.state import append_jsonl
|
||||
from ouroboros.outcomes import EXECUTION_INFRA_FAILED, terminal_outcome_axes
|
||||
|
|
@ -237,7 +238,9 @@ def _verify_worker_sha_after_spawn(
|
|||
it here. A slot opens only when the child's own ``worker_ready`` row
|
||||
(supervisor/worker_process.py) names its pid, and that row's ``git_sha`` is
|
||||
verified against ``current_sha`` in the same step. A child that is alive
|
||||
but silent past ``WORKER_READY_WINDOW_SEC`` is torn down and replaced
|
||||
but silent past ``WORKER_READY_WINDOW_SEC`` is torn down and replaced;
|
||||
its own ``worker_starting`` row permits one extension to
|
||||
``WORKER_READY_CEILING_SEC`` from birth. Replacement still runs
|
||||
through ``respawn_worker`` — at most ``WORKER_READY_MAX_ATTEMPTS``
|
||||
consecutive times for one slot, then the slot is parked and reported. A
|
||||
child that DIED during boot is released to the crash detector, which
|
||||
|
|
@ -284,6 +287,8 @@ def _watch_booting_slots(
|
|||
if not expected_sha:
|
||||
_supervisor_row({"type": "worker_sha_verify_skipped", "reason": "missing_current_sha"})
|
||||
deadline = started + max(float(WORKER_READY_WINDOW_SEC), 1.0)
|
||||
ceiling = started + max(float(WORKER_READY_CEILING_SEC), float(WORKER_READY_WINDOW_SEC), 1.0)
|
||||
extended = False
|
||||
while pending:
|
||||
ready_rows: Dict[int, Dict[str, Any]] = {}
|
||||
for row in _worker_events_since(events_cursor, "worker_ready"):
|
||||
|
|
@ -302,11 +307,32 @@ def _watch_booting_slots(
|
|||
exitcode=getattr(slot.proc, "exitcode", None),
|
||||
)
|
||||
pending.pop(wid)
|
||||
if not pending or time.time() >= deadline:
|
||||
if not pending:
|
||||
break
|
||||
if time.time() >= deadline:
|
||||
if extended or time.time() >= ceiling:
|
||||
break
|
||||
# The cursor belongs to this spawn attempt; foreign or older progress
|
||||
# cannot buy capacity for a silent slot in the same wave.
|
||||
starting_pids = {str(row.get("pid") or "")
|
||||
for row in _worker_events_since(events_cursor, "worker_starting")}
|
||||
for wid, slot in list(pending.items()):
|
||||
if str(_slot_pid(slot)) not in starting_pids:
|
||||
_replace_unready_slot(wid, slot, owner_chat_id, started, attempt)
|
||||
pending.pop(wid)
|
||||
if not pending:
|
||||
break
|
||||
extended = True
|
||||
deadline = ceiling
|
||||
_supervisor_row({
|
||||
"type": "worker_ready_window_extended", "attempt": attempt,
|
||||
"worker_ids": sorted(pending), "window_sec": float(WORKER_READY_WINDOW_SEC),
|
||||
"ceiling_sec": float(WORKER_READY_CEILING_SEC),
|
||||
})
|
||||
time.sleep(0.25)
|
||||
for wid, slot in list(pending.items()):
|
||||
_replace_unready_slot(wid, slot, owner_chat_id, started, attempt)
|
||||
_replace_unready_slot(wid, slot, owner_chat_id, started, attempt,
|
||||
window_sec=float(WORKER_READY_CEILING_SEC if extended else WORKER_READY_WINDOW_SEC))
|
||||
pending.pop(wid)
|
||||
|
||||
|
||||
|
|
@ -386,7 +412,8 @@ def kill_worker_tree(pid: int, *, keep_services: bool = False) -> None:
|
|||
|
||||
|
||||
@_serialized_worker_lifecycle
|
||||
def _replace_unready_slot(wid: int, slot: Any, owner_chat_id: int, started: float, attempt: int) -> None:
|
||||
def _replace_unready_slot(wid: int, slot: Any, owner_chat_id: int, started: float, attempt: int,
|
||||
window_sec: float = 0.0) -> None:
|
||||
"""No worker_ready inside the window: tear the child down and replace the slot, bounded.
|
||||
|
||||
One lifecycle transaction (lifecycle -> queue lock order, like every pool
|
||||
|
|
@ -407,7 +434,7 @@ def _replace_unready_slot(wid: int, slot: Any, owner_chat_id: int, started: floa
|
|||
"worker_id": wid,
|
||||
"pid": pid,
|
||||
"waited_sec": round(time.time() - started, 2),
|
||||
"window_sec": float(WORKER_READY_WINDOW_SEC),
|
||||
"window_sec": float(window_sec or WORKER_READY_WINDOW_SEC),
|
||||
"attempt": attempt,
|
||||
"max_attempts": int(WORKER_READY_MAX_ATTEMPTS),
|
||||
"action": action,
|
||||
|
|
@ -437,7 +464,7 @@ def _replace_unready_slot(wid: int, slot: Any, owner_chat_id: int, started: floa
|
|||
_pool().send_with_budget(
|
||||
owner_chat_id,
|
||||
f"⚠️ Worker slot {wid} never confirmed ready in {attempt} attempts "
|
||||
f"({WORKER_READY_WINDOW_SEC:.0f}s window each); the slot is parked. Use /restart.",
|
||||
f"(waited {time.time() - started:.0f}s on the last attempt); the slot is parked. Use /restart.",
|
||||
)
|
||||
_pool().disable_exhausted_worker_pool()
|
||||
|
||||
|
|
|
|||
|
|
@ -115,6 +115,17 @@ def worker_main(wid: int, in_q: Any, out_q: Any, repo_dir: str, drive_root: str,
|
|||
# Before ANY import that resolves the update-tx marker through git_ops (see
|
||||
# _bind_worker_repo_root): a spawned child would otherwise gate on the hardcoded default repo.
|
||||
_bind_worker_repo_root(repo_dir, drive_root)
|
||||
# Entry progress precedes extension loading and agent construction. If logging
|
||||
# fails, the parent retains the ordinary readiness window rather than losing the child.
|
||||
try:
|
||||
from ouroboros.utils import append_jsonl, utc_now_iso
|
||||
|
||||
append_jsonl(pathlib.Path(drive_root) / "logs" / "events.jsonl", {
|
||||
"ts": utc_now_iso(), "type": "worker_starting",
|
||||
"worker_id": wid, "pid": _os.getpid(), "phase": "entry",
|
||||
})
|
||||
except Exception:
|
||||
log.debug("Worker entry progress unavailable", exc_info=True)
|
||||
# Adopt the server's custody session id. Under the 'spawn' start method this
|
||||
# process re-imported process_custody and minted a fresh _SESSION_ID; without
|
||||
# adopting the server's id, every service/process this worker records looks
|
||||
|
|
@ -191,7 +202,6 @@ def worker_main(wid: int, in_q: Any, out_q: Any, repo_dir: str, drive_root: str,
|
|||
if pytest_default_real_data_dir:
|
||||
extensions_owned = False
|
||||
try:
|
||||
from ouroboros.utils import append_jsonl, utc_now_iso
|
||||
append_jsonl(_drive / "logs" / "supervisor.jsonl", {
|
||||
"ts": utc_now_iso(),
|
||||
"type": "worker_extension_reload_skipped",
|
||||
|
|
|
|||
|
|
@ -492,6 +492,7 @@ def queue_snapshot(data_root) -> dict:
|
|||
|
||||
_READINESS_ROW_TYPES = (
|
||||
"worker_sha_verify", "worker_ready_timeout", "worker_ready_released",
|
||||
"worker_ready_window_extended",
|
||||
"worker_dead_detected", "worker_crash",
|
||||
)
|
||||
|
||||
|
|
|
|||
121
tests/test_activity_census_browser.py
Normal file
121
tests/test_activity_census_browser.py
Normal file
|
|
@ -0,0 +1,121 @@
|
|||
"""Activity census rows use real direct actors and the shared task controls."""
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import os
|
||||
import time
|
||||
import threading
|
||||
from pathlib import Path
|
||||
|
||||
import pytest
|
||||
|
||||
from tests.test_ui_smoke_playwright import direct_server_with_data as direct_server_with_data
|
||||
from tests.ui_chat_viewport_smoke import _CAPTURE_TEST_SOCKET
|
||||
|
||||
|
||||
@pytest.mark.ui_browser
|
||||
@pytest.mark.parametrize("action", ["hurry", "finalize", "stop_now"])
|
||||
def test_activity_direct_turn_reaches_shared_control_endpoint(direct_server_with_data, monkeypatch, action):
|
||||
from playwright.sync_api import sync_playwright
|
||||
from tests import fixtures_mock_llm
|
||||
|
||||
monkeypatch.setattr(fixtures_mock_llm, "HOLD_SECONDS", 90)
|
||||
fixtures_mock_llm.HOLD_RELEASE.clear()
|
||||
entered = threading.Event()
|
||||
|
||||
def held_completion(handler):
|
||||
payload = json.loads(handler.rfile.read(int(handler.headers.get("Content-Length", 0))))
|
||||
entered.set()
|
||||
fixtures_mock_llm.HOLD_RELEASE.wait(90)
|
||||
streaming = payload.get("stream") is True
|
||||
answer = {"id": "held-completion", "choices": [{"index": 0, "finish_reason": "stop",
|
||||
"delta" if streaming else "message": {"role": "assistant", "content": "OK"}}],
|
||||
"usage": {"prompt_tokens": 1, "completion_tokens": 1, "total_tokens": 2}}
|
||||
content = ("data: " + json.dumps(answer) + "\n\ndata: [DONE]\n\n" if streaming else json.dumps(answer)).encode()
|
||||
handler.send_response(200)
|
||||
handler.send_header("Content-Type", "text/event-stream" if streaming else "application/json")
|
||||
handler.send_header("Content-Length", str(len(content)))
|
||||
handler.end_headers()
|
||||
handler.wfile.write(content)
|
||||
|
||||
monkeypatch.setattr(fixtures_mock_llm._Handler, "do_POST", held_completion)
|
||||
url = direct_server_with_data["url"]
|
||||
evidence = Path(os.environ.get("OUROBOROS_UI_EVIDENCE_DIR", direct_server_with_data["data_dir"].parent))
|
||||
evidence.mkdir(parents=True, exist_ok=True)
|
||||
try:
|
||||
with sync_playwright() as pw:
|
||||
browser = pw.chromium.launch()
|
||||
try:
|
||||
page = browser.new_page(viewport={"width": 1440, "height": 900})
|
||||
page.add_init_script(f"({_CAPTURE_TEST_SOCKET})()")
|
||||
page.goto(url, wait_until="domcontentloaded")
|
||||
page.wait_for_function("() => window.__testSockets?.some(s => s.readyState === 1)")
|
||||
page.fill("#chat-input", "Respond with exactly OK")
|
||||
page.click("#chat-send")
|
||||
deadline = time.monotonic() + 30
|
||||
actor = None
|
||||
while time.monotonic() < deadline and (actor is None or not entered.is_set()):
|
||||
state = page.request.get(url + "/api/state").json()
|
||||
actor = next((a for a in state.get("active_chat_activities", []) if a["kind"] == "direct_chat"), None)
|
||||
if actor is None or not entered.is_set():
|
||||
page.wait_for_timeout(100)
|
||||
assert actor and entered.is_set(), state
|
||||
task_id = actor["activity_id"]
|
||||
queue = page.request.get(url + "/api/tasks?queue_only=1").json()["queue"]
|
||||
assert all((q.get("id") or q.get("task", {}).get("id")) != task_id
|
||||
for q in queue["running"] + queue["pending"])
|
||||
page.click('[data-nav-page="dashboard"]')
|
||||
page.click('[data-dashboard-tab="activity"]')
|
||||
button = page.locator(f'[data-activity-section="queue"] [data-id="{task_id}"]')
|
||||
button.wait_for(state="visible")
|
||||
assert button.count() == 1
|
||||
row = button.locator("xpath=../..")
|
||||
assert "Direct turn" in row.inner_text()
|
||||
assert "Nothing running" not in page.locator('[data-activity-section="queue"]').inner_text()
|
||||
button.click()
|
||||
menu = page.locator('body > .task-control-menu')
|
||||
menu.wait_for(state="visible")
|
||||
assert menu.locator('[data-task-control]').all_text_contents() == ["Wrap up", "Hurry up", "Stop now"]
|
||||
page.screenshot(path=str(evidence / f"activity-{action}.png"), full_page=True)
|
||||
endpoint = "hurry" if action == "hurry" else "cancel"
|
||||
# Direct actors stop cooperatively. Let the held fake provider
|
||||
# complete after cancellation is submitted, so the real endpoint
|
||||
# can observe settlement instead of an artificial never-returning call.
|
||||
release = None
|
||||
if action == "stop_now":
|
||||
def release_after_request(request):
|
||||
nonlocal release
|
||||
if request.method == "POST" and request.url.endswith(f"/{task_id}/cancel"):
|
||||
release = threading.Timer(0.3, fixtures_mock_llm.HOLD_RELEASE.set)
|
||||
release.start()
|
||||
page.on("request", release_after_request)
|
||||
with page.expect_response(lambda r: f"/api/tasks/{task_id}/{endpoint}" in r.url and r.request.method == "POST") as response:
|
||||
menu.locator(f'[data-task-control="{action}"]').click()
|
||||
if release is not None:
|
||||
release.join(timeout=1)
|
||||
receipt = response.value
|
||||
body = receipt.json()
|
||||
assert receipt.status in (200, 202), body
|
||||
assert body.get("ok") is True, body
|
||||
request = receipt.request.post_data_json
|
||||
if action == "hurry":
|
||||
assert request["request_id"], request
|
||||
else:
|
||||
assert request.get("stop_policy", "immediate") == ("immediate" if action == "stop_now" else "finalize_then_cancel")
|
||||
(evidence / f"activity-{action}.json").write_text(json.dumps({
|
||||
"actor": actor, "queue": queue, "status": receipt.status, "request": request, "response": body,
|
||||
}, indent=2), encoding="utf-8")
|
||||
fixtures_mock_llm.HOLD_RELEASE.set()
|
||||
deadline = time.monotonic() + 30
|
||||
while time.monotonic() < deadline:
|
||||
state = page.request.get(url + "/api/state").json()
|
||||
if state.get("active_chat_activities_complete") is True and not any(
|
||||
a["activity_id"] == task_id for a in state["active_chat_activities"]):
|
||||
break
|
||||
page.wait_for_timeout(100)
|
||||
else:
|
||||
raise AssertionError(state)
|
||||
finally:
|
||||
browser.close()
|
||||
finally:
|
||||
fixtures_mock_llm.HOLD_RELEASE.set()
|
||||
116
tests/test_chat_composer_reserve_browser.py
Normal file
116
tests/test_chat_composer_reserve_browser.py
Normal file
|
|
@ -0,0 +1,116 @@
|
|||
"""Main/Project composer clearance and wizard checkbox geometry in real UI flows."""
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import os
|
||||
from pathlib import Path
|
||||
|
||||
import pytest
|
||||
|
||||
from tests.test_ui_smoke_playwright import direct_server_with_data as direct_server_with_data
|
||||
from tests.test_subscription_setup_browser import subscription_ui as subscription_ui
|
||||
from tests.ui_chat_viewport_smoke import _CAPTURE_TEST_SOCKET, _emit_ws_frame, _SETTLE_RESTORE_FRAMES
|
||||
|
||||
|
||||
_MEASURE = """messages => {
|
||||
messages.scrollTop = messages.scrollHeight;
|
||||
const composer = messages.parentElement.querySelector('#chat-input-area, .chat-input-area');
|
||||
const last = [...messages.children].filter(el => el.getBoundingClientRect().height > 0).at(-1);
|
||||
const end = getComputedStyle(messages, '::after');
|
||||
return {height: messages.clientHeight, scrollHeight: messages.scrollHeight,
|
||||
scrollTop: messages.scrollTop, padding: getComputedStyle(messages).paddingBottom,
|
||||
spacer: end.content, basis: end.flexBasis, shrink: end.flexShrink,
|
||||
lastBottom: last?.getBoundingClientRect().bottom ?? null,
|
||||
composerTop: composer.getBoundingClientRect().top,
|
||||
reserve: getComputedStyle(messages).getPropertyValue('--chat-input-reserve')};
|
||||
}"""
|
||||
|
||||
|
||||
@pytest.mark.ui_browser
|
||||
@pytest.mark.parametrize("engine", ["chromium", "webkit"])
|
||||
def test_main_and_project_composer_clearance(direct_server_with_data, engine):
|
||||
from playwright.sync_api import sync_playwright
|
||||
from ouroboros.projects_registry import create_project
|
||||
|
||||
data = direct_server_with_data["data_dir"]
|
||||
project = create_project(data, "composer-room", name="Composer room")
|
||||
(data / "logs" / "chat.jsonl").write_text(json.dumps({
|
||||
"ts": "2026-09-17T10:00:00Z", "direction": "out", "chat_id": project["chat_id"],
|
||||
"text": "Saved Project answer.\n" + "\n".join(f"History line {i}" for i in range(12)),
|
||||
"format": "markdown",
|
||||
}) + "\n", encoding="utf-8")
|
||||
evidence = Path(os.environ.get("OUROBOROS_UI_EVIDENCE_DIR", data.parent))
|
||||
evidence.mkdir(parents=True, exist_ok=True)
|
||||
observations = []
|
||||
with sync_playwright() as pw:
|
||||
browser = getattr(pw, engine).launch()
|
||||
try:
|
||||
page = browser.new_page(viewport={"width": 1187, "height": 734})
|
||||
page.add_init_script(f"({_CAPTURE_TEST_SOCKET})()")
|
||||
page.goto(direct_server_with_data["url"], wait_until="domcontentloaded")
|
||||
page.wait_for_function("() => window.__testSockets?.some(s => s.readyState === 1)")
|
||||
for width in (1187, 390):
|
||||
page.set_viewport_size({"width": width, "height": 734})
|
||||
for surface in ("main", "project"):
|
||||
if surface == "project":
|
||||
page.locator('.nav-project-row[data-project-id="composer-room"]').evaluate("el => el.click()")
|
||||
messages = page.locator('#project-panel .chat-messages')
|
||||
messages.locator('.chat-bubble').filter(has_text="Saved Project answer").wait_for(state="visible")
|
||||
else:
|
||||
page.locator('[data-nav-page="chat"]').evaluate("el => el.click()")
|
||||
messages = page.locator('#chat-messages')
|
||||
page.evaluate(_SETTLE_RESTORE_FRAMES)
|
||||
for text in ("", "one\ntwo\nthree\nfour\nfive"):
|
||||
composer_input = messages.locator('xpath=..').locator('textarea')
|
||||
composer_input.fill(text)
|
||||
page.evaluate(_SETTLE_RESTORE_FRAMES)
|
||||
metrics = messages.evaluate(_MEASURE)
|
||||
observations.append({"width": width, "surface": surface, "multiline": bool(text), **metrics})
|
||||
assert metrics["padding"] == "0px" and metrics["shrink"] == "0", metrics
|
||||
assert metrics["lastBottom"] is None or metrics["lastBottom"] <= metrics["composerTop"] + 1, metrics
|
||||
chat_id = project["chat_id"] if surface == "project" else 1
|
||||
_emit_ws_frame(page, {"type": "chat", "role": "assistant", "chat_id": chat_id,
|
||||
"content": "Growing answer.\n" + "\n".join(f"New line {i}" for i in range(50)),
|
||||
"ts": "2026-09-18T10:00:00Z"})
|
||||
page.evaluate(_SETTLE_RESTORE_FRAMES)
|
||||
metrics = messages.evaluate(_MEASURE)
|
||||
assert metrics["lastBottom"] <= metrics["composerTop"] + 1, metrics
|
||||
page.screenshot(path=str(evidence / f"composer-{engine}-{surface}-{width}.png"), full_page=True)
|
||||
if width == 390:
|
||||
page.evaluate("() => {document.documentElement.style.setProperty('--vvh','500px'); document.body.classList.add('keyboard-open'); document.documentElement.classList.add('keyboard-open');}")
|
||||
metrics = messages.evaluate(_MEASURE)
|
||||
assert metrics["spacer"] == "none", metrics
|
||||
assert metrics["lastBottom"] <= metrics["composerTop"] + 1, metrics
|
||||
page.evaluate("() => {document.body.classList.remove('keyboard-open'); document.documentElement.classList.remove('keyboard-open');}")
|
||||
(evidence / f"composer-{engine}.json").write_text(json.dumps(observations, indent=2), encoding="utf-8")
|
||||
finally:
|
||||
browser.close()
|
||||
|
||||
|
||||
@pytest.mark.ui_browser
|
||||
@pytest.mark.serial
|
||||
def test_wizard_enabled_checkbox_stays_inline(subscription_ui):
|
||||
"""The real wizard keeps its already-fixed checkbox width at phone scale."""
|
||||
from tests.test_ui_smoke_agents_panel import _wizard_step_until, _WIZARD_ON_AGENTS_JS
|
||||
from tests.test_subscription_setup_browser import capture
|
||||
|
||||
page = subscription_ui["page"]
|
||||
page.goto(subscription_ui["url"] + "/onboarding")
|
||||
_wizard_step_until(page, _WIZARD_ON_AGENTS_JS, forward=True)
|
||||
label = page.locator('#onboarding-available-subagents .available-subagents-toolbar .local-toggle')
|
||||
for width in (1440, 390):
|
||||
page.set_viewport_size({"width": width, "height": 900})
|
||||
label.scroll_into_view_if_needed()
|
||||
geometry = label.evaluate("""label => {
|
||||
const input = label.querySelector('input');
|
||||
const box = input.getBoundingClientRect();
|
||||
const text = [...label.childNodes].find(n => n.nodeType === Node.TEXT_NODE && n.textContent.trim());
|
||||
const range = new Range(); range.selectNodeContents(text);
|
||||
const word = range.getBoundingClientRect();
|
||||
return {checkboxWidth:box.width, height:label.getBoundingClientRect().height,
|
||||
textTop:word.top, textBottom:word.bottom, checkboxTop:box.top, checkboxBottom:box.bottom};
|
||||
}""")
|
||||
assert geometry["checkboxWidth"] == 14, geometry
|
||||
assert geometry["height"] <= 28, geometry
|
||||
assert geometry["textBottom"] <= geometry["checkboxBottom"] + 6, geometry
|
||||
capture(page, f"wizard-inline-checkbox-{width}")
|
||||
|
|
@ -23,6 +23,7 @@ _LEAVES = (settings_defaults, settings_scales, model_slots, review_model_routes,
|
|||
# New subscription capabilities belong to the same leaves, but did not exist on
|
||||
# the historical extraction's facade and need not add compatibility re-exports.
|
||||
_ADDED_OWNERS = {
|
||||
"WORKER_READY_CEILING_SEC": runtime_limits,
|
||||
"IMMEDIATE_SETTINGS": settings_scales,
|
||||
"RESTART_REQUIRED_SETTINGS": settings_scales,
|
||||
"get_finalization_grace_sec": runtime_limits,
|
||||
|
|
|
|||
|
|
@ -15,6 +15,7 @@ import urllib.parse
|
|||
import urllib.request
|
||||
|
||||
import pytest
|
||||
from tests.ui_chat_viewport_smoke import _CAPTURE_TEST_SOCKET
|
||||
|
||||
pytest_plugins = ("tests.test_ui_smoke_playwright",)
|
||||
|
||||
|
|
@ -132,6 +133,7 @@ def test_large_attachment_returns_through_real_document_handler_and_download(
|
|||
browser = playwright.chromium.launch()
|
||||
try:
|
||||
page = browser.new_page(viewport={"width": 1440, "height": 1000}, accept_downloads=True)
|
||||
page.add_init_script(f"({_CAPTURE_TEST_SOCKET})()")
|
||||
document_frames = []
|
||||
def websocket(socket):
|
||||
def frame(payload):
|
||||
|
|
@ -147,6 +149,8 @@ def test_large_attachment_returns_through_real_document_handler_and_download(
|
|||
page.on("request", lambda request: requests.append((request.method, request.url)))
|
||||
page.goto(url, wait_until="domcontentloaded")
|
||||
page.locator("#chat-input").wait_for(state="visible")
|
||||
page.wait_for_function(
|
||||
"() => window.__testSockets?.some(socket => socket.readyState === WebSocket.OPEN)")
|
||||
page.locator("#chat-file-input").set_input_files([str(path) for path in attachments])
|
||||
assert page.locator(".attach-badge").count() == 28
|
||||
page.locator("#chat-input").fill("Return the large attached dataset as a downloadable document.")
|
||||
|
|
|
|||
|
|
@ -512,3 +512,24 @@ def test_webkit_scrollbar_recipe_covers_both_axes() -> None:
|
|||
"the horizontal bar must be as thin as the vertical one: "
|
||||
f"width {declarations['width']} vs height {declarations['height']}"
|
||||
)
|
||||
|
||||
|
||||
def test_chat_transcript_reserves_composer_space_as_one_flex_spacer():
|
||||
"""End space is a flex item; keyboard flow removes it and its extra gap."""
|
||||
css = _decommented(_read("web/style.css"))
|
||||
rules = [(selector.strip(), body) for selector, body in RULE.findall(css)]
|
||||
assert not re.search(r"padding-bottom:\s*(?:calc\()?var\(--chat-input-reserve", css)
|
||||
spacers = [(selector, body) for selector, body in rules if "chat-messages::after" in selector]
|
||||
bases = [body for _, body in spacers if "flex:" in body]
|
||||
assert len(bases) == 1
|
||||
assert "content: '';" in bases[0]
|
||||
assert "flex: 0 0 calc(var(--chat-input-reserve, 108px) - 8px);" in bases[0]
|
||||
mobile = [body for _, body in spacers if "env(safe-area-inset-bottom" in body]
|
||||
assert len(mobile) == 1
|
||||
assert "- 8px" in mobile[0]
|
||||
panel = [body for selector, body in spacers if selector == ".chat-instance-panel .chat-messages::after"]
|
||||
assert len(panel) == 1
|
||||
assert "flex-basis: calc(var(--chat-input-reserve, 108px) - 8px);" in panel[0]
|
||||
keyboard = [body for selector, body in spacers if "body.keyboard-open" in selector]
|
||||
assert len(keyboard) == 1
|
||||
assert "content: none;" in keyboard[0]
|
||||
|
|
|
|||
|
|
@ -8,7 +8,7 @@ What is pinned, on fake process objects (no child is ever forked here):
|
|||
booted SHA in the same step; a foreign pid's row does not open it;
|
||||
* no ``worker_ready`` inside the window -> the child is torn down (process tree), the slot is
|
||||
replaced through ``respawn_worker`` and a typed ``worker_ready_timeout`` row names slot, pid,
|
||||
wait and reason;
|
||||
wait and reason; its own entry progress permits one extension to the birth-relative ceiling;
|
||||
* the replacement loop is bounded: at ``WORKER_READY_MAX_ATTEMPTS`` the slot is parked (kept
|
||||
``reaping``, no further respawn) and the owner is told;
|
||||
* a child that DIED during boot is released to the crash detector, which already owns death;
|
||||
|
|
@ -224,6 +224,107 @@ def test_a_foreign_pid_row_does_not_open_the_slot(pool, seam):
|
|||
assert seam.killed == [5001] and seam.respawned == [(0, {"ready_attempt": 2})]
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def boot_clock(pool, monkeypatch):
|
||||
"""Advance only the readiness clock; no real wait or fabricated ready signal."""
|
||||
clock = SimpleNamespace(now=1000.0, on_tick=lambda: None)
|
||||
|
||||
def sleep(seconds):
|
||||
clock.now += seconds
|
||||
clock.on_tick()
|
||||
|
||||
monkeypatch.setattr(pool.lifecycle, "time", SimpleNamespace(time=lambda: clock.now, sleep=sleep))
|
||||
monkeypatch.setattr(pool.lifecycle, "WORKER_READY_WINDOW_SEC", 1.0)
|
||||
monkeypatch.setattr(pool.lifecycle, "WORKER_READY_CEILING_SEC", 3.0)
|
||||
return clock
|
||||
|
||||
|
||||
def test_own_entry_progress_allows_ready_after_the_initial_window(pool, seam, boot_clock):
|
||||
slot = _booting_slot(pool, 0, 5001)
|
||||
append_jsonl(seam.events, {"type": "worker_starting", "worker_id": 0, "pid": 5001})
|
||||
|
||||
def ready_later():
|
||||
if boot_clock.now == 1002.0:
|
||||
append_jsonl(seam.events, _READY_ROW)
|
||||
|
||||
boot_clock.on_tick = ready_later
|
||||
seam.run({0: slot}, seam.cursor, 1, 1000.0)
|
||||
|
||||
assert not slot.reaping and boot_clock.now == 1002.0
|
||||
assert seam.killed == [] and seam.respawned == []
|
||||
extension = _rows(seam.supervisor, "worker_ready_window_extended")
|
||||
assert len(extension) == 1 and extension[0]["worker_ids"] == [0]
|
||||
assert extension[0]["ceiling_sec"] == 3.0
|
||||
assert _rows(seam.supervisor, "worker_sha_verify")[0]["ok"] is True
|
||||
|
||||
|
||||
def test_progress_extends_once_and_never_moves_the_birth_relative_ceiling(pool, seam, boot_clock):
|
||||
slot = _booting_slot(pool, 0, 5001)
|
||||
row = {"type": "worker_starting", "worker_id": 0, "pid": 5001}
|
||||
append_jsonl(seam.events, row)
|
||||
boot_clock.on_tick = lambda: append_jsonl(seam.events, row)
|
||||
|
||||
seam.run({0: slot}, seam.cursor, 1, 1000.0)
|
||||
|
||||
assert boot_clock.now == 1003.0
|
||||
assert seam.killed == [5001] and seam.respawned == [(0, {"ready_attempt": 2})]
|
||||
assert len(_rows(seam.supervisor, "worker_ready_window_extended")) == 1
|
||||
timeout = _rows(seam.supervisor, "worker_ready_timeout")[0]
|
||||
assert timeout["window_sec"] == timeout["waited_sec"] == 3.0
|
||||
|
||||
|
||||
@pytest.mark.parametrize("progress", ["absent", "foreign_pid", "before_cursor", "dead_child"])
|
||||
def test_only_current_own_progress_from_a_live_child_extends(pool, seam, boot_clock, progress):
|
||||
slot = _booting_slot(pool, 0, 5001, alive=progress != "dead_child")
|
||||
cursor = seam.cursor
|
||||
if progress != "absent":
|
||||
append_jsonl(seam.events, {"type": "worker_starting", "worker_id": 0,
|
||||
"pid": 9999 if progress == "foreign_pid" else 5001})
|
||||
if progress == "before_cursor":
|
||||
cursor = pool.lifecycle.events_log_cursor()
|
||||
|
||||
seam.run({0: slot}, cursor, 1, 1000.0)
|
||||
|
||||
assert _rows(seam.supervisor, "worker_ready_window_extended") == []
|
||||
if progress == "dead_child":
|
||||
assert boot_clock.now == 1000.0 and not slot.reaping
|
||||
assert seam.killed == [] and seam.respawned == []
|
||||
assert _rows(seam.supervisor, "worker_ready_released")[0]["reason"] == "died_during_boot"
|
||||
else:
|
||||
assert boot_clock.now == 1001.0 and seam.killed == [5001]
|
||||
assert _rows(seam.supervisor, "worker_ready_timeout")[0]["window_sec"] == 1.0
|
||||
|
||||
|
||||
def test_one_childs_progress_does_not_extend_a_silent_peer_in_the_same_wave(pool, seam, boot_clock):
|
||||
first, second = _booting_slot(pool, 0, 5001), _booting_slot(pool, 1, 5002)
|
||||
append_jsonl(seam.events, {"type": "worker_starting", "worker_id": 0, "pid": 5001})
|
||||
|
||||
def ready_later():
|
||||
if boot_clock.now == 1002.0:
|
||||
append_jsonl(seam.events, _READY_ROW)
|
||||
|
||||
boot_clock.on_tick = ready_later
|
||||
seam.run({0: first, 1: second}, seam.cursor, 1, 1000.0)
|
||||
|
||||
assert not first.reaping and second.reaping
|
||||
assert seam.killed == [5002] and seam.respawned == [(1, {"ready_attempt": 2})]
|
||||
timeout = _rows(seam.supervisor, "worker_ready_timeout")[0]
|
||||
assert timeout["waited_sec"] == 1.0 and timeout["pid"] == 5002
|
||||
assert _rows(seam.supervisor, "worker_ready_window_extended")[0]["worker_ids"] == [0]
|
||||
|
||||
|
||||
def test_a_late_watcher_cannot_grant_a_fresh_window_after_the_ceiling(pool, seam, boot_clock):
|
||||
slot = _booting_slot(pool, 0, 5001)
|
||||
append_jsonl(seam.events, {"type": "worker_starting", "worker_id": 0, "pid": 5001})
|
||||
boot_clock.now = 1004.0
|
||||
|
||||
seam.run({0: slot}, seam.cursor, 1, 1000.0)
|
||||
|
||||
assert boot_clock.now == 1004.0 and seam.killed == [5001]
|
||||
assert _rows(seam.supervisor, "worker_ready_window_extended") == []
|
||||
assert _rows(seam.supervisor, "worker_ready_timeout")[0]["waited_sec"] == 4.0
|
||||
|
||||
|
||||
def test_sha_mismatch_on_the_ready_row_opens_the_slot_and_tells_the_owner(pool, seam, monkeypatch):
|
||||
monkeypatch.setattr(pool.workers, "load_state", lambda: {"current_sha": "abc123", "owner_chat_id": 7})
|
||||
slot = _booting_slot(pool, 0, 5001)
|
||||
|
|
|
|||
47
tests/test_worker_startup_progress.py
Normal file
47
tests/test_worker_startup_progress.py
Normal file
|
|
@ -0,0 +1,47 @@
|
|||
"""The real worker reports entry before the potentially slow extension load."""
|
||||
|
||||
import json
|
||||
import os
|
||||
|
||||
import pytest
|
||||
|
||||
from tests.test_terminal_file_boundary import worker as worker
|
||||
|
||||
|
||||
@pytest.mark.parametrize("progress_write_fails", [False, True])
|
||||
def test_entry_progress_precedes_extensions_and_logging_failure_does_not_block_ready(
|
||||
worker, monkeypatch, progress_write_fails,
|
||||
):
|
||||
import ouroboros.extension_loader as extensions
|
||||
import ouroboros.utils as utils
|
||||
|
||||
worker.task["type"] = "shutdown"
|
||||
events = worker.root / "logs" / "events.jsonl"
|
||||
real_append = utils.append_jsonl
|
||||
attempts, loads = [], []
|
||||
|
||||
def append(path, row, **kwargs):
|
||||
if row.get("type") == "worker_starting":
|
||||
attempts.append(row)
|
||||
if progress_write_fails:
|
||||
raise OSError("entry log unavailable")
|
||||
return real_append(path, row, **kwargs)
|
||||
|
||||
def load(*_args, **_kwargs):
|
||||
loads.append(True)
|
||||
assert len(attempts) == 1 and attempts[0]["pid"] == os.getpid()
|
||||
assert attempts[0]["worker_id"] == 0 and attempts[0]["phase"] == "entry"
|
||||
rows = events.read_text(encoding="utf-8") if events.exists() else ""
|
||||
assert '"worker_ready"' not in rows
|
||||
assert ('"worker_starting"' in rows) is not progress_write_fails
|
||||
|
||||
monkeypatch.setattr(utils, "append_jsonl", append)
|
||||
monkeypatch.setattr(extensions, "reload_all", load)
|
||||
worker.run()
|
||||
|
||||
rows = [json.loads(line) for line in events.read_text(encoding="utf-8").splitlines()]
|
||||
assert loads == [True] and worker.crashes == [] and worker.calls == []
|
||||
assert [row["type"] for row in rows] == (
|
||||
["worker_ready"] if progress_write_fails else ["worker_starting", "worker_ready"]
|
||||
)
|
||||
assert rows[-1]["pid"] == os.getpid() and rows[-1]["git_sha"] == "test-sha"
|
||||
|
|
@ -56,7 +56,7 @@ export function initActivity({ mount, ws } = {}) {
|
|||
return { root, status, content, loaded: false };
|
||||
});
|
||||
|
||||
function renderQueue(queue) {
|
||||
function renderQueue(queue, census) {
|
||||
if (!Array.isArray(queue?.running) || !Array.isArray(queue?.pending)) throw new Error('Queue unavailable');
|
||||
const { running, pending } = queue;
|
||||
// #322: the snapshot already carries the pause truth — a member's own
|
||||
|
|
@ -86,8 +86,30 @@ export function initActivity({ mount, ws } = {}) {
|
|||
</div>
|
||||
</div>`;
|
||||
};
|
||||
const parts = [...running.map((q) => row(q, 'running')), ...pending.map((q) => row(q, 'pending'))];
|
||||
return parts.length ? parts.join('') : '<div class="activity-empty">Nothing running or queued.</div>';
|
||||
// Queue facts keep their runtime/budget controls. The census adds only
|
||||
// missing identities; queued direct turns can appear in both sources.
|
||||
const known = new Set([...running, ...pending].map((q) => String(q.id || q.task?.id || '')));
|
||||
const live = (Array.isArray(census?.active_chat_activities) ? census.active_chat_activities : [])
|
||||
.filter((a) => a && !known.has(String(a.activity_id || '')));
|
||||
const names = { direct_chat: 'Direct turn', managed_task: 'Managed task' };
|
||||
const liveRow = (a) => {
|
||||
const started = Number(a.started_at) || 0;
|
||||
const elapsed = started > 0 ? ` · ${Math.max(0, Math.round(Date.now() / 1000 - started))}s` : '';
|
||||
return `<div class="activity-row">
|
||||
<div class="activity-row-main">
|
||||
<span class="activity-name">${esc(names[a.kind] || 'Live turn')}</span>
|
||||
<span class="activity-sub">${esc(a.phase || '')}${elapsed}</span>
|
||||
</div>
|
||||
<div class="activity-row-actions">
|
||||
<button type="button" class="btn btn-xs btn-danger" data-act="task-control" data-id="${esc(a.activity_id || '')}">${esc(TASK_CONTROL_TRIGGER_LABEL)}</button>
|
||||
</div>
|
||||
</div>`;
|
||||
};
|
||||
const parts = [...running.map((q) => row(q, 'running')), ...pending.map((q) => row(q, 'pending')), ...live.map(liveRow)];
|
||||
if (parts.length) return parts.join('');
|
||||
return census?.active_chat_activities_complete === true
|
||||
? '<div class="activity-empty">Nothing running or queued.</div>'
|
||||
: '<div class="activity-empty">Queue empty; live turns unknown.</div>';
|
||||
}
|
||||
|
||||
function formatWhen(value) {
|
||||
|
|
@ -164,7 +186,8 @@ export function initActivity({ mount, ws } = {}) {
|
|||
getJson('/api/schedules'),
|
||||
]);
|
||||
if (revision !== refreshRevision) return;
|
||||
const renderers = [(data) => renderQueue(data?.queue), renderBg, renderSchedules];
|
||||
const census = results[1].status === 'fulfilled' ? results[1].value : null;
|
||||
const renderers = [(data) => renderQueue(data?.queue, census), renderBg, renderSchedules];
|
||||
sections.forEach((section, index) => {
|
||||
const { root, status, content } = section;
|
||||
root.removeAttribute('aria-busy');
|
||||
|
|
|
|||
|
|
@ -1297,7 +1297,7 @@ body.resizing-panels { user-select: none; }
|
|||
overflow-x: hidden;
|
||||
padding: 20px;
|
||||
padding-top: var(--chat-header-reserve, 56px);
|
||||
padding-bottom: var(--chat-input-reserve, 108px);
|
||||
padding-bottom: 0;
|
||||
display: flex;
|
||||
flex-direction: column;
|
||||
gap: 8px;
|
||||
|
|
@ -1316,6 +1316,14 @@ body.resizing-panels { user-select: none; }
|
|||
flex-shrink: 0;
|
||||
}
|
||||
|
||||
/* An actual flex item keeps the overlay reserve inside scrollable overflow in
|
||||
WKWebView, including short histories. The column gap supplies the final 8px. */
|
||||
#chat-messages::after,
|
||||
.chat-messages::after {
|
||||
content: '';
|
||||
flex: 0 0 calc(var(--chat-input-reserve, 108px) - 8px);
|
||||
}
|
||||
|
||||
/* Quiet transcript navigation, centered above the complete two-row composer. */
|
||||
.chat-scroll-bottom-btn {
|
||||
position: absolute;
|
||||
|
|
@ -5594,10 +5602,15 @@ textarea.chat-input {
|
|||
#chat-messages,
|
||||
.chat-messages {
|
||||
padding-top: var(--chat-header-reserve, 80px);
|
||||
padding-bottom: calc(var(--chat-input-reserve, 108px) + env(safe-area-inset-bottom, 0px));
|
||||
padding-bottom: 0;
|
||||
overflow-x: hidden;
|
||||
}
|
||||
|
||||
#chat-messages::after,
|
||||
.chat-messages::after {
|
||||
flex-basis: calc(var(--chat-input-reserve, 108px) + env(safe-area-inset-bottom, 0px) - 8px);
|
||||
}
|
||||
|
||||
/* Keyboard-open layout: only #chat-messages scrolls within visual viewport. */
|
||||
html.keyboard-open,
|
||||
body.keyboard-open {
|
||||
|
|
@ -5654,6 +5667,12 @@ textarea.chat-input {
|
|||
-webkit-overflow-scrolling: touch;
|
||||
}
|
||||
|
||||
/* The composer is in flow: remove the item, including its column gap. */
|
||||
body.keyboard-open #page-chat.active #chat-messages::after,
|
||||
body.keyboard-open .chat-instance-panel .chat-messages::after {
|
||||
content: none;
|
||||
}
|
||||
|
||||
body.keyboard-open #page-chat.active #chat-input-area,
|
||||
body.keyboard-open .chat-instance-panel .chat-input-area {
|
||||
flex-shrink: 0;
|
||||
|
|
@ -6139,7 +6158,13 @@ textarea.chat-input {
|
|||
/* No absolute overlay HEADER here: drop the main chat's phantom 56px
|
||||
top reserve; the composer reserve matches the shared dock contract. */
|
||||
padding: 8px 14px;
|
||||
padding-bottom: var(--chat-input-reserve, 108px);
|
||||
padding-bottom: 0;
|
||||
}
|
||||
|
||||
/* The Project composer already includes its safe-area inset in the measured
|
||||
reserve; preserve the panel's own reserve on narrow screens as well. */
|
||||
.chat-instance-panel .chat-messages::after {
|
||||
flex-basis: calc(var(--chat-input-reserve, 108px) - 8px);
|
||||
}
|
||||
|
||||
/* Shared dock contract (v6.71.0): the panel composer docks exactly like the
|
||||
|
|
|
|||
|
|
@ -134,7 +134,7 @@ const message = (node) => node.querySelector('.ui-status').textContent;
|
|||
|
||||
function emptyActivity(routes) {
|
||||
routes.set(queueUrl, response({ queue: { running: [], pending: [] } }));
|
||||
routes.set(backgroundUrl, response({ bg_consciousness_enabled: false }));
|
||||
routes.set(backgroundUrl, response({ bg_consciousness_enabled: false, active_chat_activities: [], active_chat_activities_complete: true }));
|
||||
routes.set(schedulesUrl, response({ tasks: [] }));
|
||||
}
|
||||
|
||||
|
|
@ -198,6 +198,65 @@ test('Activity invalid success payload is unavailable and stale requests cannot
|
|||
assert.match(section(mount, 'queue').textContent, /Current work/);
|
||||
});
|
||||
|
||||
test('Activity unites live direct turns with queue identities and preserves queue controls', async (t) => {
|
||||
const { mount, routes, ws, calls } = setup(t);
|
||||
emptyActivity(routes);
|
||||
const now = Date.now;
|
||||
Date.now = () => 100000;
|
||||
t.after(() => { Date.now = now; });
|
||||
routes.set(backgroundUrl, response({ bg_consciousness_enabled: true,
|
||||
active_chat_activities_complete: false, active_chat_activities: [
|
||||
{ activity_id: 'direct-1', kind: 'direct_chat', phase: 'thinking', started_at: 88 },
|
||||
{ activity_id: 'queued-1', kind: 'direct_chat', phase: 'thinking', started_at: 80 },
|
||||
{ activity_id: 'paused-1', kind: 'managed_task', phase: 'running', started_at: 80 },
|
||||
] }));
|
||||
routes.set(queueUrl, response({ queue: {
|
||||
running: [{ id: 'queued-1', type: 'task', runtime_sec: 9, task: { title: 'Queue title' } }],
|
||||
pending: [{ task: { id: 'paused-1', title: 'Paused work', _budget_pause: true } }],
|
||||
} }));
|
||||
await initActivity({ mount, ws }).refresh();
|
||||
const queue = section(mount, 'queue');
|
||||
const rows = queue.querySelectorAll('.activity-row');
|
||||
assert.equal(rows.length, 3, 'partial positive census adds exactly the missing direct identity');
|
||||
assert.match(rows[0].textContent, /Queue title[\s\S]*running · task · 9s/);
|
||||
assert.match(rows[1].textContent, /Paused work[\s\S]*paused \(budget\)/);
|
||||
assert.equal(rows[1].querySelector('button').dataset.budgetPaused, '1');
|
||||
assert.match(rows[2].textContent, /Direct turn[\s\S]*thinking · 12s/);
|
||||
assert.equal(rows[2].querySelector('button').dataset.id, 'direct-1');
|
||||
assert.equal(rows[2].querySelector('button').dataset.act, 'task-control');
|
||||
assert.doesNotMatch(queue.textContent, /Nothing running/);
|
||||
assert.deepEqual(calls, [queueUrl, backgroundUrl, schedulesUrl], 'reuse the existing state read');
|
||||
});
|
||||
|
||||
test('Activity complete empty census differs from unknown and failed state preserves queue facts', async (t) => {
|
||||
const { mount, routes, ws } = setup(t);
|
||||
emptyActivity(routes);
|
||||
const activity = initActivity({ mount, ws });
|
||||
const queue = section(mount, 'queue');
|
||||
await activity.refresh();
|
||||
assert.match(queue.textContent, /Nothing running or queued/);
|
||||
routes.set(backgroundUrl, response({ bg_consciousness_enabled: false,
|
||||
active_chat_activities: [], active_chat_activities_complete: false }));
|
||||
await activity.refresh();
|
||||
assert.match(queue.textContent, /Queue empty; live turns unknown/);
|
||||
routes.set(backgroundUrl, response({}, 503));
|
||||
await activity.refresh();
|
||||
assert.match(queue.textContent, /Queue empty; live turns unknown/);
|
||||
assert.equal(message(queue), '');
|
||||
routes.set(queueUrl, response({ queue: { running: [{ id: 'running-1', task: { title: 'Still working' } }], pending: [] } }));
|
||||
await activity.refresh();
|
||||
assert.match(queue.textContent, /Still working/);
|
||||
assert.equal(message(queue), '');
|
||||
routes.set(backgroundUrl, response({ bg_consciousness_enabled: false,
|
||||
active_chat_activities: [{ activity_id: 'direct-2', kind: 'direct_chat', phase: 'thinking' }],
|
||||
active_chat_activities_complete: true }));
|
||||
routes.set(queueUrl, response({}, 503));
|
||||
await activity.refresh();
|
||||
assert.match(message(queue), /Could not refresh.*unknown/);
|
||||
assert.match(queue.textContent, /Still working/);
|
||||
assert.doesNotMatch(queue.textContent, /Nothing running/);
|
||||
});
|
||||
|
||||
test('Logs reports partial history without losing live rows or deduplication, then clears the gap on reconnect', async (t) => {
|
||||
const { mount, routes, ws, calls } = setup(t);
|
||||
const live = { type: 'test_live_event', ts: '2026-09-09T12:00:00Z', message: 'live' };
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue