From 1c8b49db2dccfc493bec955691ac8525e168071b Mon Sep 17 00:00:00 2001 From: adrianwedd Date: Tue, 25 Aug 2026 10:30:50 +1000 Subject: [PATCH] fix(memory): run nightly consolidation off the px-mind tick (#291) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Consolidation could not succeed under load, and the three reasons were mutually reinforcing — fixing any one alone would not have helped. It ran inline on px-mind's ~60s awareness tick while its declared deadline is 600s, twice px-mind's own 300s staleness window. Honouring the budget and keeping the mind loop alive were mutually exclusive, so the call carried an ad-hoc timeout=180 that silently overrode the declared 600s: every live failure timed out at exactly 180.1s while every success took 30-65s. And attempt 2 was unspendable — memory.MAX_ATTEMPTS_PER_DAY promised two tries while _TYPE_QUOTAS["consolidate"] was 1 and its cooldown 20h. - mind: `_consolidation_tick` now only supervises. It reaps a finished worker, heartbeats the in-flight marker, clears what a previous process left behind, and starts at most one `px-mind-consolidation` daemon thread. It always returns promptly; awareness, reflection and the battery check run behind it. - memory: `state/consolidation_job.json`, keyed on pid + heartbeat. The worker dies with its process, so a marker outliving its owner is always a lie — a restart cannot leave a false in-progress claim, and the cleanup is recorded as a failure rather than silently reset. An unfinished run past JOB_OVERRUN_AFTER_S is reported once, not every 60s. - memory: `consolidate()` passes no `timeout`, and `run_claude_session`'s default is now None, so brain._DEADLINE_S is the single deadline source and the declared 600s is what reaches `ask_brain`. - claude_session: consolidate quota 1 -> 2 (== MAX_ATTEMPTS_PER_DAY) and cooldown 72000 -> 2400 (== RETRY_SPACING_S). 40min spaces the retry *past* the 30-min global cooldown rather than exempting it from one — nobody is waiting on a 3am retry. - px-motd renders an in-flight run additively; no new health status value. Resident-only Claude is untouched: same ask_brain, same single-flight lock, same mailbox, no second process and no cold fallback. Only which thread blocks on it changed. Tests are inert and every duration is synthetic: the ">300s run" is a threading.Event plus a backdated marker, and the 600s deadline is asserted as a plumbed number read out of the request file. No live brain call. Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01ThqC6Gq4mZyXWnnvGC2a57 --- CLAUDE.md | 53 ++- bin/px-motd | 20 +- src/pxh/claude_session.py | 28 +- src/pxh/memory.py | 271 ++++++++++++- src/pxh/mind.py | 130 +++++- tests/test_claude_session.py | 35 +- tests/test_consolidation_background_job.py | 434 +++++++++++++++++++++ tests/test_health_memory_formation.py | 42 ++ tests/test_memory.py | 13 +- tests/test_mind_consolidation_health.py | 62 ++- tests/test_mind_utils.py | 60 ++- 11 files changed, 1080 insertions(+), 68 deletions(-) create mode 100644 tests/test_consolidation_background_job.py diff --git a/CLAUDE.md b/CLAUDE.md index abea3a11..1321f68f 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -153,6 +153,39 @@ threshold while the store has not grown for days. That is why `overdue` (>48h) is a separate axis from `status`, and why `summarize()` returns a non-empty string for overdue memory on an otherwise clean bill. +**Consolidation runs on a background daemon thread, never on the tick (#291).** +`_consolidation_tick` only *supervises*: it heartbeats the in-flight marker, +reaps what a previous process left behind, and decides whether to start a +worker. It must always return promptly, because awareness, reflection and the +battery check run behind it. The reason is arithmetic, not taste — the kind's +declared deadline is **600s** and px-mind's own staleness window is **300s**, so +an inline call that honoured its budget guaranteed px-mind read `stale`. That +made the three defects mutually reinforcing: the pass could not be given enough +time, so it was given an ad-hoc `timeout=180` that silently overrode the +declared deadline, so every failure timed out at exactly 180.1s. + +**`state/consolidation_job.json` is the in-flight marker**, and it is keyed on +**pid plus a heartbeat** — `{"status","pid","attempt","started_ts","heartbeat_ts"}`. +A health record says what the last *finished* attempt did; nothing else answers +"is one running right now". The worker thread dies with its process, so a marker +outliving its owner is always a lie: `consolidation_job_is_stale()` calls it +stale when `/proc/` is gone (or is not px-mind) **or** the heartbeat has +been quiet past `JOB_HEARTBEAT_STALE_S`. A restart therefore cannot leave a +false in-progress claim behind, and the tick records that cleanup as a +*failure* — no memory formed that night — rather than silently resetting. The +heartbeat is written by the **tick**, not the worker: the worker spends its +whole life blocked in `ask_brain` and could not beat if it wanted to. Past +`JOB_OVERRUN_AFTER_S` an unfinished run is reported **once**, not every 60s. + +**Two attempts a night, spaced 40 min apart** (`memory.RETRY_SPACING_S`). All +three numbers used to disagree: `MAX_ATTEMPTS_PER_DAY` promised 2 while the +`consolidate` quota was 1 and its type cooldown 20h, so attempt 2 was +structurally unreachable. The spacing deliberately **clears** the 30-min global +cooldown rather than adding `consolidate` to `_GLOBAL_COOLDOWN_EXEMPT` — nobody +is waiting on a 3am retry, and each exemption is one more way for a background +job to crowd a session someone *is* waiting on. Success **or a correct skip** +marks the date done; only a failure leaves the later retry open. + **Claude spend visibility:** `token_log.log_usage()` takes a `backend` argument and splits totals under `by_backend` in `state/token_usage.json`. The top-level totals mix free Ollama with paid Claude and cannot answer "what am I spending". `call_llm()` also sets `result["backend"]` to the tier that actually served — the `backend=` reflection log line shows the *configured* primary, not the one that answered. ### Idle-Alive Daemon @@ -193,7 +226,7 @@ Three-layer architecture: Ollama. Writes to `state/thoughts.jsonl` only after a valid M5 response. - **Layer 3 — Expression** (30min cooldown; `greet_arrival` bypasses it on a real arrival, 120s anti-flap): dispatches to tool-voice/tool-look/tool-remember and cognitive tools. Valid actions include (wait, greet, greet_arrival, comment, remember, look_at, weather_comment, scan, play_sound, photograph, emote, look_around, time_check, calendar_check, introspect, evolve, morning_fact, research, compose, self_debug, blog_essay, message_obi, set_goal, update_goal, complete_goal). Suppressed during school, quiet time, bedtime (all calendar-driven). **Hardcoded night silence: 19:00–07:00 Hobart time — no speech/audio/motion. Silent cognitive actions (`NIGHT_ALLOWED_ACTIONS`: wait, remember, research, compose, introspect, self_debug, set_goal, update_goal, complete_goal) are exempt and run overnight.** - **`message_obi` action**: SPARK initiates a direct message to Obi via the dashboard. Exponential backoff: starts at 10min, doubles on unanswered nudge, caps at 4h, resets when Obi replies. Respects all suppressors. Thoughts with `action=message_obi` are **redacted** in `thoughts-spark.jsonl` (written as `[private message to Obi]`) so the private DM content never reaches the public `/api/v1/public/thoughts` endpoint. -- **Memory consolidation**: nightly Haiku pass (02:00–06:00 Hobart, ≤2 attempts/day, state/consolidation_meta.json) distills the last 24h of thoughts into state/memories-spark.jsonl; reflection retrieves the top-3 relevant memories by keyword/tag overlap. Goal persistence in state/intention-spark.json (7-day expiry, one active at a time). +- **Memory consolidation**: nightly Haiku pass (02:00–06:00 Hobart, ≤2 attempts/day ≥40min apart, state/consolidation_meta.json) distills the last 24h of thoughts into state/memories-spark.jsonl; reflection retrieves the top-3 relevant memories by keyword/tag overlap. **Runs on a background daemon thread with a pid-keyed job marker (`state/consolidation_job.json`), never inline on the tick** — see the health section. Goal persistence in state/intention-spark.json (7-day expiry, one active at a time). **Critical gotchas:** - All time-of-day logic uses `ZoneInfo("Australia/Hobart")` — never hardcoded UTC offsets @@ -244,7 +277,23 @@ Watches `state/thoughts-spark.jsonl` (salience ≥0.7 or spoken action), runs Cl | `compose` | Haiku | 4h | 2/day | | `conversation` | Sonnet | 15min | 4/day | | `blog` | Haiku | 30min | 5/day | -| `consolidate` | Haiku | 20h | 1/day | +| `consolidate` | Haiku | 40min | 2/day | + +`consolidate` is 40min/2 rather than 20h/1 so `memory.MAX_ATTEMPTS_PER_DAY`'s +second nightly attempt can actually be spent (#291). 40min also clears the +30-min global cooldown, so the retry is *spaced past* it rather than exempted +from it. + +**Deadlines are declared once, in `brain._DEADLINE_S`, and `timeout=` on +`run_claude_session` defaults to `None` so that table is what reaches +`ask_brain`.** Pass a number only when you mean to override the kind's declared +budget — an override that is *tighter* silently wins and makes the declared +value unreachable, which is exactly what an ad-hoc `timeout=180` did to +`consolidate`'s 600s: every live failure timed out at 180.1s while every +success took 30–65s. Callers still passing ad-hoc values for classified kinds +(`mind.py` self_debug 600 vs. declared 900; `bin/px-blog`, `bin/tool-research`, +`bin/tool-compose`, `bin/tool-blog` passing 300, which matches) are redundant at +best and drift at worst. Global: 30min cooldown between sessions (except `self_debug`/`blog`), 8/day cap. When ≤2 remaining: only `self_debug`/`evolve` allowed. Bypass: `PX_CLAUDE_BUDGET_DISABLED=1`. Session log: `state/claude_sessions.jsonl`. diff --git a/bin/px-motd b/bin/px-motd index cd76fef4..e70f682c 100755 --- a/bin/px-motd +++ b/bin/px-motd @@ -141,8 +141,12 @@ def _trunc(s: str, maxlen: int = 72) -> str: # because px-motd runs as root from PAM and deliberately imports no pxh module. MEMORY_OVERDUE_S = 48 * 3600 +# Mirrors pxh.memory.JOB_HEARTBEAT_STALE_S, duplicated for the same reason and +# pinned by test_motd_job_heartbeat_threshold_matches_memory_module. +JOB_HEARTBEAT_STALE_S = 300 -def _memory_formation_line(rec: dict) -> str: + +def _memory_formation_line(rec: dict, job: dict | None = None) -> str: """Render "long-term memory" from the px-mind-consolidation health record. Keyed on ``last_success_ts``, never ``updated_ts``: a failed nightly pass @@ -150,10 +154,21 @@ def _memory_formation_line(rec: dict) -> str: *success* is the only number that answers "when did SPARK last form a long-term memory". `rec` is the raw record so this stays pure and testable; px-motd's own path clamping decides which file it comes from. + + `job` is `state/consolidation_job.json`, the background worker's in-flight + marker (#291). A run in progress is appended as a hint rather than + replacing the line: "last formed 26h ago" is still the honest answer while + tonight's pass is mid-flight, and a marker whose heartbeat has gone quiet + is not evidence of anything, so it is ignored. """ formed = (rec or {}).get("last_success_ts") age = _age_secs(formed) label = f" \ud83e\udde0 {DIM}long-term memory{RESET} " + beat = _age_secs((job or {}).get("heartbeat_ts")) + if beat is not None and beat <= JOB_HEARTBEAT_STALE_S: + started = _age_secs((job or {}).get("started_ts")) + for_s = f" for {started // 60}m" if started and started >= 60 else "" + label += f"{DIM}(consolidating now{for_s}){RESET} " if age is None: err = _trunc(str((rec or {}).get("last_error") or ""), 44) tail = f" {DIM}{err}{RESET}" if err else "" @@ -636,7 +651,8 @@ def section_spark() -> list[str]: # all until now, so it could fail every night while this banner showed a # busy, healthy-looking robot. out.append(_memory_formation_line( - _json(STATE / "health" / "px-mind-consolidation.json"))) + _json(STATE / "health" / "px-mind-consolidation.json"), + _json(STATE / "consolidation_job.json"))) # Weather weather = aw.get("weather", {}) if aw else {} diff --git a/src/pxh/claude_session.py b/src/pxh/claude_session.py index 35e854ef..33ad39dd 100644 --- a/src/pxh/claude_session.py +++ b/src/pxh/claude_session.py @@ -75,7 +75,13 @@ def _model_for_type(session_type: str) -> str: "compose": 14400, # 4 hours "conversation": 900, # 15 min "blog": 1800, # 30 min - "consolidate": 72000, # 20 hours + # 40 min, matching memory.RETRY_SPACING_S (#291). It was 20 hours, which + # meant the *first* attempt of a night consumed the only slot the second + # one could ever have used: `memory.MAX_ATTEMPTS_PER_DAY` promised two + # tries between 02:00 and 06:00 and this made the second structurally + # unreachable. 40 min also clears the 30-min global cooldown, so attempt 2 + # is spaced past it rather than exempted from it. + "consolidate": 2400, } _TYPE_QUOTAS: dict[str, int] = { @@ -85,7 +91,9 @@ def _model_for_type(session_type: str) -> str: "compose": 2, "conversation": 4, "blog": 5, - "consolidate": 1, + # Two, to match memory.MAX_ATTEMPTS_PER_DAY (#291). A quota of 1 made the + # retry that module offers impossible to spend. + "consolidate": 2, } # Higher number = higher priority. Used for budget-tight gating. @@ -325,7 +333,7 @@ def brain_kinds() -> frozenset[str]: def _run_via_brain( session_type: str, prompt: str, - timeout: int, + timeout: int | None, model: str, ) -> RunResult: """Serve a session from the resident Claude session instead of a subprocess. @@ -339,6 +347,14 @@ def _run_via_brain( resident session's envelope is fixed when it launches and cannot be widened for one request. That is a security property, not a limitation to work around. + + `timeout=None` means "use the deadline this kind declares" — + `brain._DEADLINE_S`, which is the single source of truth for how long a + classified kind may take. A caller that passes a number overrides it, and + #291 is what that costs when the number is wrong: `consolidate` declares + 600s, memory.py passed an ad-hoc 180, and because the tighter value always + wins the declared budget was unreachable — every live failure timed out at + exactly 180.1s. """ from . import brain # local import keeps the tmux dependency off the hot path @@ -376,7 +392,7 @@ def _run_via_brain( def run_claude_session( session_type: str, prompt: str, - timeout: int = 300, + timeout: int | None = None, allowed_tools: str = "", skip_permissions: bool = False, cwd: str | Path | None = None, @@ -385,6 +401,10 @@ def run_claude_session( ) -> RunResult: """Run a Claude session with budget checking, model routing, and logging. + timeout: seconds, or None (the default) to use the deadline the kind + declares in `brain._DEADLINE_S`. Prefer None — the declared per-kind + deadline is the one source of truth, and an ad-hoc override that is + tighter silently replaces it (see #291). model_override: use this model instead of the session-type default. skip_budget_check: skip rate-limit check (use for sub-phases of an already-checked session). Raises SessionBudgetExhausted if rate-limited (unless skip_budget_check=True). diff --git a/src/pxh/memory.py b/src/pxh/memory.py index 9fe575fc..a0e324c9 100644 --- a/src/pxh/memory.py +++ b/src/pxh/memory.py @@ -46,6 +46,36 @@ MIN_THOUGHTS = 5 MAX_MEMORIES_PER_DAY = 8 +# Minimum gap between the two nightly attempts (#291). +# +# Both attempts used to land ~60s apart, which made the retry decorative twice +# over: `claude_session.COOLDOWN_S` (the 30-minute global session cooldown) +# rejected it outright, and even without that a second Haiku turn started one +# minute after the first failed is a retry of the same load conditions. 40 +# minutes clears the global cooldown with ten minutes of margin and still fits +# twice inside the four-hour 02:00-06:00 window. +# +# Spacing past the cooldown is deliberate, rather than adding `consolidate` to +# `_GLOBAL_COOLDOWN_EXEMPT`. The cooldown exists because two Claude sessions +# close together contend on a 4-core Pi, and that reason applies to +# consolidation exactly as written — the attempt is 600s of resident-session +# time. An exemption would buy the retry by denying the premise. +RETRY_SPACING_S = 2400 + +# A consolidation worker that has held the job marker this long has overrun its +# own budget (brain._DEADLINE_S["consolidate"] is 600s, plus prompt assembly and +# the single-flight lock wait). Reported once, then left alone: the worker is a +# daemon thread blocked inside ask_brain, which has its own deadline, so the +# honest thing is to make the overrun visible rather than to pretend it can be +# cancelled. +JOB_OVERRUN_AFTER_S = 900 + +# How long the marker's heartbeat may lag before its owner is presumed gone. +# The owning px-mind refreshes it on every ~60s awareness tick, so five minutes +# of silence means the process that claimed it is no longer ticking — which is +# what a restart mid-consolidation looks like from the next process's side. +JOB_HEARTBEAT_STALE_S = 300 + _STOPWORDS = frozenset( """a about after again all am an and any are as at be because been before but by can did do does for from had has have he her his how i if in into is it its just me more @@ -325,8 +355,15 @@ def consolidate(dry: bool = False, persona: str = "spark", import pxh.claude_session as claude_session try: + # No `timeout=` on purpose (#291). The deadline for a classified + # brain kind is declared once, in brain._DEADLINE_S — 600s for + # `consolidate` — and `ask_brain` only reaches for it when the + # caller passes nothing. The ad-hoc 180 that used to sit here was + # tighter, so it won every time and the declared budget was never + # once reachable: every live failure measured exactly 180.1s while + # the successes measured 30-65s. result = claude_session.run_claude_session( - "consolidate", prompt, timeout=180, allowed_tools="") + "consolidate", prompt, allowed_tools="") except claude_session.SessionBudgetExhausted as exc: return {"status": "failed", "error": str(exc)} except Exception as exc: @@ -368,26 +405,89 @@ def consolidation_meta_file() -> Path: return _state_dir() / "consolidation_meta.json" +def _read_consolidation_meta() -> dict: + try: + meta = json.loads(consolidation_meta_file().read_text(encoding="utf-8")) + except (OSError, json.JSONDecodeError): + return {} + return meta if isinstance(meta, dict) else {} + + +def _parse_ts(value: object) -> dt.datetime | None: + try: + parsed = dt.datetime.fromisoformat(str(value).replace("Z", "+00:00")) + except (ValueError, TypeError): + return None + if parsed.tzinfo is None: + parsed = parsed.replace(tzinfo=dt.timezone.utc) + return parsed + + +def consolidation_due(now: dt.datetime | None = None) -> str | None: + """The read-only half of `maybe_consolidate`'s gate. + + Returns ``None`` when an attempt is due right now, otherwise a short reason + it is not. Writes nothing and costs one small file read, because px-mind + calls it on every ~60s awareness tick to decide whether to spend a worker + thread — the authoritative, state-mutating gate is still inside + `maybe_consolidate`, which calls this first. + + Splitting it out is what lets the tick stay cheap without duplicating the + window/date/attempt logic in two places that could then disagree. + """ + local = (now or dt.datetime.now(HOBART_TZ)).astimezone(HOBART_TZ) + if not (CONSOLIDATION_WINDOW[0] <= local.hour < CONSOLIDATION_WINDOW[1]): + return "outside the 02:00-06:00 window" + meta = _read_consolidation_meta() + if meta.get("last_date") != local.strftime("%Y-%m-%d"): + return None # a fresh Hobart date: first attempt + if meta.get("done"): + return "already done for this date" + attempts = meta.get("attempts", 0) + if attempts >= MAX_ATTEMPTS_PER_DAY: + return f"attempt cap reached ({attempts}/{MAX_ATTEMPTS_PER_DAY})" + last_attempt = _parse_ts(meta.get("last_attempt_ts")) + if last_attempt is not None: + elapsed = (local - last_attempt).total_seconds() + if 0 <= elapsed < RETRY_SPACING_S: + return (f"retry spacing ({int(elapsed)}s / {RETRY_SPACING_S}s " + f"since attempt {attempts})") + return None + + +def next_consolidation_attempt() -> int: + """1 or 2 — which attempt `maybe_consolidate` would spend next. + + Recorded on the job marker so an operator looking at a stuck run can tell + "tonight's first try" from "the last chance before morning". + """ + return _read_consolidation_meta().get("attempts", 0) + 1 + + def maybe_consolidate(dry: bool = False, persona: str = "spark", now: dt.datetime | None = None) -> dict | None: - """Once-per-Hobart-date gate for consolidate(). None = not now, or dry mode.""" + """Once-per-Hobart-date gate for consolidate(). None = not now, or dry mode. + + Runs on px-mind's consolidation worker thread, never on the awareness tick + itself (#291): a single attempt is allowed up to `brain._DEADLINE_S`'s 600s, + which is twice px-mind's own 300s health-staleness window. + """ if dry: return None local = (now or dt.datetime.now(HOBART_TZ)).astimezone(HOBART_TZ) - if not (CONSOLIDATION_WINDOW[0] <= local.hour < CONSOLIDATION_WINDOW[1]): + if consolidation_due(now=local) is not None: return None today = local.strftime("%Y-%m-%d") meta_f = consolidation_meta_file() - meta = {} - try: - meta = json.loads(meta_f.read_text(encoding="utf-8")) - except (OSError, json.JSONDecodeError): - pass + meta = _read_consolidation_meta() if meta.get("last_date") != today: meta = {"last_date": today, "attempts": 0, "done": False} - if meta.get("done") or meta.get("attempts", 0) >= MAX_ATTEMPTS_PER_DAY: - return None meta["attempts"] = meta.get("attempts", 0) + 1 + # Stamped from `local`, not the wall clock, so the spacing gate reads the + # same clock the caller passed in — otherwise a test (or a replayed window) + # measures the gap against `now()` and every retry looks either overdue or + # impossible. + meta["last_attempt_ts"] = local.isoformat() meta_f.parent.mkdir(parents=True, exist_ok=True) try: atomic_write(meta_f, json.dumps(meta) + "\n") @@ -404,3 +504,154 @@ def maybe_consolidate(dry: bool = False, persona: str = "spark", except OSError: pass return result + + +# --------------------------------------------------------------------------- +# Job marker — the in-flight record for the background consolidation worker +# --------------------------------------------------------------------------- +# +# Consolidation runs on a daemon thread inside px-mind (#291), so "is it +# running right now?" is not answerable from any file health.py writes: a +# component record says what the *last finished* attempt did. The marker below +# is that missing state, and it is deliberately a file rather than a process +# variable — px-motd, the operator and a later px-mind all need to read it. +# +# It is keyed on pid, and liveness is `/proc/{pid}` plus a cmdline check, the +# same idiom as px-mind's own single-instance guard and its px-alive pid check. +# The worker thread dies with its process, so a marker outliving its owner is +# always a lie; making the lie detectable is the whole job of `pid` and +# `heartbeat_ts`. Nothing here ever treats a marker as authoritative evidence +# that work is in progress without first checking that its owner still exists. + +def consolidation_job_file() -> Path: + return _state_dir() / "consolidation_job.json" + + +def read_consolidation_job() -> dict: + """The in-flight marker, or {} when no worker holds one.""" + try: + job = json.loads(consolidation_job_file().read_text(encoding="utf-8")) + except (OSError, json.JSONDecodeError, ValueError): + return {} + return job if isinstance(job, dict) else {} + + +def _write_consolidation_job(job: dict) -> bool: + f = consolidation_job_file() + try: + f.parent.mkdir(parents=True, exist_ok=True) + atomic_write(f, json.dumps(job) + "\n") + return True + except OSError: + return False + + +def clear_consolidation_job() -> None: + """Drop the marker. Called from the worker's `finally`, so it must not + raise — a lost marker is recoverable (staleness), a raise in a daemon + thread's teardown is not.""" + try: + consolidation_job_file().unlink(missing_ok=True) + except OSError: + pass + + +def _pid_is_px_mind(pid: object) -> bool: + if not isinstance(pid, int) or pid <= 0: + return False + if pid == os.getpid(): + return True + if not Path(f"/proc/{pid}").is_dir(): + return False + try: + cmdline = Path(f"/proc/{pid}/cmdline").read_bytes().decode(errors="replace") + except OSError: + return False + return "px-mind" in cmdline + + +def consolidation_job_age_s(job: dict, + now: dt.datetime | None = None) -> float | None: + started = _parse_ts((job or {}).get("started_ts")) + if started is None: + return None + return ((now or dt.datetime.now(dt.timezone.utc)).astimezone(dt.timezone.utc) + - started).total_seconds() + + +def consolidation_job_is_stale(job: dict, + now: dt.datetime | None = None) -> bool: + """True when nothing is plausibly still working on this marker. + + Two independent ways to be stale, and both are needed. The owning process + may be gone outright (a restart — the thread died with it), or it may still + exist while no longer ticking, in which case only the heartbeat says so. + A missing heartbeat reads as stale rather than fresh: an unreadable marker + must never be able to block tonight's attempt forever. + """ + if not job: + return False + if not _pid_is_px_mind(job.get("pid")): + return True + beat = _parse_ts(job.get("heartbeat_ts")) + if beat is None: + return True + now = (now or dt.datetime.now(dt.timezone.utc)).astimezone(dt.timezone.utc) + return (now - beat).total_seconds() > JOB_HEARTBEAT_STALE_S + + +def claim_consolidation_job(attempt: int = 0, + now: dt.datetime | None = None) -> bool: + """Take the marker for this process. False means somebody else holds it. + + Under `FileLock` because the check and the write must not interleave, even + though px-mind's single-instance PID guard should already make a second + claimant impossible — a guard that is only correct because another guard + exists is one refactor away from being wrong. + """ + f = consolidation_job_file() + stamp = ((now or dt.datetime.now(dt.timezone.utc)) + .astimezone(dt.timezone.utc).isoformat()) + try: + f.parent.mkdir(parents=True, exist_ok=True) + with FileLock(str(f) + ".lock", timeout=LOCK_TIMEOUT_S): + existing = read_consolidation_job() + if existing and not consolidation_job_is_stale(existing, now=now): + return False + return _write_consolidation_job({ + "status": "running", + "pid": os.getpid(), + "attempt": attempt, + "started_ts": stamp, + "heartbeat_ts": stamp, + }) + except Exception: # noqa: BLE001 — a failed claim is "don't start", never a raise + return False + + +def touch_consolidation_job(now: dt.datetime | None = None) -> dict: + """Refresh the heartbeat from the owning process and report an overrun once. + + The heartbeat is written by px-mind's tick, not by the worker: the worker + spends its whole life blocked inside `ask_brain` and could not beat if it + wanted to. "px-mind still observes this thread alive" is exactly the fact a + later reader needs, so that is what gets recorded. + + The returned dict carries a transient ``newly_overran`` flag — not + persisted — the first tick on which the run crosses JOB_OVERRUN_AFTER_S, so + the caller can report a stuck worker exactly once instead of every 60s. + """ + job = read_consolidation_job() + if not job: + return {} + now = (now or dt.datetime.now(dt.timezone.utc)).astimezone(dt.timezone.utc) + job["heartbeat_ts"] = now.isoformat() + age = consolidation_job_age_s(job, now=now) + newly_overran = (age is not None and age > JOB_OVERRUN_AFTER_S + and not job.get("overrun_reported")) + if newly_overran: + job["overrun_reported"] = True + _write_consolidation_job(job) + if newly_overran: + job = dict(job, newly_overran=True, age_s=age) + return job diff --git a/src/pxh/mind.py b/src/pxh/mind.py index 3ed2b5ee..28110430 100644 --- a/src/pxh/mind.py +++ b/src/pxh/mind.py @@ -26,6 +26,7 @@ import subprocess import sys import tempfile +import threading import time import urllib.error import urllib.request @@ -1143,7 +1144,6 @@ def _fetch_ha_context(dry: bool = False) -> dict | None: log("ha_context: skipped — no PX_HA_TOKEN") return None - import threading result_box: dict = {} def _do_fetch(): @@ -3575,12 +3575,19 @@ def expression(thought: dict, dry: bool, awareness: dict | None = None) -> bool: return True -def _consolidation_tick(session: dict, dry: bool) -> None: - """Nightly memory consolidation (02:00-06:00 Hobart, once per date, SPARK only). - Never raises — the mind loop must survive any consolidation failure. +CONSOLIDATION_THREAD_NAME = "px-mind-consolidation" + +# The one in-flight consolidation worker, or None. Process-local by design: a +# daemon thread dies with its process, so this variable and the on-disk marker +# answer two different questions — "is this process running one right now" and +# "did some process claim one and not finish it". +_consolidation_worker: threading.Thread | None = None + + +def _record_consolidation_outcome(result: dict | None) -> None: + """Report one finished consolidation attempt to health (#289). - Health reporting distinguishes three outcomes, and the distinction is the - whole point of reporting at all: + Three outcomes, and the distinction is the whole point of reporting at all: * ``None`` — not attempted, because it is not the window, it is already done for this Hobart date, or this is a dry run. **Records nothing.** A @@ -3594,26 +3601,107 @@ def _consolidation_tick(session: dict, dry: bool) -> None: * anything else — attempted and failed. This is the case that was silent before: the store simply stopped growing while every other dial read ok. """ + if result is None: + return + component = health_mod.CONSOLIDATION_COMPONENT + status = result.get("status") + detail = (result.get("error") or result.get("reason") + or f"wrote {result.get('written', 0)}") + log(f"consolidation: {status} — {detail}") + if status == "dry": + return # a dry run forms no memory; it must not refresh the record + if status in ("ok", "skipped"): + health_mod.record_success( + component, + detail={"status": status, "written": result.get("written", 0), + "note": detail}) + else: + health_mod.record_failure(component, detail, detail={"status": status}) + + +def _consolidation_worker_main(dry: bool) -> None: + """The background consolidation run. Never raises, always clears the marker. + + This is the body that used to run inline on the awareness tick. A single + attempt is allowed `brain._DEADLINE_S["consolidate"]` = 600s, which is twice + px-mind's own 300s staleness window — inline, honouring that budget would + mean ten minutes with no awareness snapshot, no reflection and no battery + check, which is why #291's one-line deadline fix was not safe on its own. + """ + try: + _record_consolidation_outcome(spark_memory.maybe_consolidate(dry=dry)) + except Exception as exc: # noqa: BLE001 — a daemon thread's death is silent + log(f"consolidation error: {exc}") + health_mod.record_failure(health_mod.CONSOLIDATION_COMPONENT, + f"{type(exc).__name__}: {exc}") + finally: + spark_memory.clear_consolidation_job() + + +def _consolidation_tick(session: dict, dry: bool) -> None: + """Supervise the nightly consolidation worker. Returns promptly, always. + + Everything expensive happens on `_consolidation_worker_main`'s thread; this + function only decides whether to start one, refreshes the in-flight + marker's heartbeat, and reaps what the last process left behind. It must + never block, because the mind loop's awareness/reflection/battery cadence + runs behind it. + """ + global _consolidation_worker if (session.get("persona") or "").lower().strip() != "spark": return component = health_mod.CONSOLIDATION_COMPONENT try: - _cons = spark_memory.maybe_consolidate(dry=dry) - if _cons is None: + worker = _consolidation_worker + if worker is not None and not worker.is_alive(): + worker.join(timeout=0) + _consolidation_worker = worker = None + + if worker is not None: + # In flight in this process. Heartbeat the marker so a live-slow + # run stays distinguishable from an abandoned one, then get out of + # the way — this early return is what keeps awareness ticking + # through a ten-minute consolidation. + job = spark_memory.touch_consolidation_job() + if job.get("newly_overran"): + age = int(job.get("age_s") or 0) + log(f"consolidation: worker still running after {age}s " + f"(budget {spark_memory.JOB_OVERRUN_AFTER_S}s) — overran") + health_mod.record_failure( + component, f"worker overran ({age}s)", + detail={"status": "overran", "pid": job.get("pid")}) return - _status = _cons.get("status") - _detail = _cons.get("error") or _cons.get("reason") or f"wrote {_cons.get('written', 0)}" - log(f"consolidation: {_status} — {_detail}") - if _status == "dry": - return # a dry run forms no memory; it must not refresh the record - if _status in ("ok", "skipped"): - health_mod.record_success( - component, - detail={"status": _status, "written": _cons.get("written", 0), - "note": _detail}) - else: - health_mod.record_failure( - component, _detail, detail={"status": _status}) + + job = spark_memory.read_consolidation_job() + if job: + # A marker with no worker of ours behind it. Either px-mind + # restarted mid-run (the thread went with the process) or another + # px-mind owns it. Only the first is ours to clean up, and it is a + # night on which no memory formed — so it is a failure, recorded + # once, not a silent reset. + if spark_memory.consolidation_job_is_stale(job): + spark_memory.clear_consolidation_job() + log(f"consolidation: cleared abandoned job " + f"(pid={job.get('pid')}, started {job.get('started_ts')})") + health_mod.record_failure( + component, "abandoned mid-run (px-mind restarted?)", + detail={"status": "abandoned", "pid": job.get("pid")}) + return + + if dry: + return + reason = spark_memory.consolidation_due() + if reason is not None: + return + attempt = spark_memory.next_consolidation_attempt() + if not spark_memory.claim_consolidation_job(attempt=attempt): + return + _consolidation_worker = threading.Thread( + target=_consolidation_worker_main, args=(dry,), + name=CONSOLIDATION_THREAD_NAME, daemon=True) + _consolidation_worker.start() + log(f"consolidation: started in background (attempt {attempt}" + f"/{spark_memory.MAX_ATTEMPTS_PER_DAY})") except Exception as exc: log(f"consolidation error: {exc}") health_mod.record_failure(component, f"{type(exc).__name__}: {exc}") diff --git a/tests/test_claude_session.py b/tests/test_claude_session.py index 29accd1d..19e08d7a 100644 --- a/tests/test_claude_session.py +++ b/tests/test_claude_session.py @@ -502,26 +502,51 @@ def test_empty_log(self, tmp_path): def test_consolidate_session_type_registered(): from pxh import claude_session as cs + from pxh import memory assert cs._model_for_type("consolidate").startswith("claude-haiku") - assert cs._TYPE_QUOTAS["consolidate"] == 1 - assert cs._TYPE_COOLDOWNS["consolidate"] == 72000 + # Two per night, 40 min apart (#291). Was 1/72000, which made the second of + # memory.MAX_ATTEMPTS_PER_DAY's two attempts structurally unspendable — + # attempt 1 consumed the only slot attempt 2 could ever have used. + assert cs._TYPE_QUOTAS["consolidate"] == memory.MAX_ATTEMPTS_PER_DAY == 2 + assert cs._TYPE_COOLDOWNS["consolidate"] == memory.RETRY_SPACING_S == 2400 assert cs._PRIORITY["consolidate"] == 2 assert cs._ENV_OVERRIDES["consolidate"] == "PX_CLAUDE_MODEL_CONSOLIDATE" -def test_consolidate_quota_one_per_day(tmp_path, monkeypatch): +def test_consolidate_quota_is_two_per_day(tmp_path, monkeypatch): import datetime as dt import json from pxh import claude_session as cs log = tmp_path / "claude_sessions.jsonl" - now = dt.datetime.now(dt.timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ") - log.write_text(json.dumps({"ts": now, "type": "consolidate"}) + "\n", encoding="utf-8") + # Two attempts already spent tonight, spaced far enough apart that neither + # cooldown is what refuses the third — the quota must be. + now = dt.datetime.now(dt.timezone.utc) + log.write_text("".join( + json.dumps({"ts": (now - dt.timedelta(hours=h)).strftime("%Y-%m-%dT%H:%M:%SZ"), + "type": "consolidate"}) + "\n" for h in (3, 2)), encoding="utf-8") monkeypatch.setattr(cs, "SESSION_LOG", log) monkeypatch.setattr(cs, "BUDGET_DISABLED", False) reason = cs.check_budget("consolidate") assert reason is not None and "quota" in reason +def test_consolidate_second_attempt_is_admitted(tmp_path, monkeypatch): + """The retry memory.py schedules must actually get past the budget gate.""" + import datetime as dt + import json + from pxh import claude_session as cs + from pxh import memory + log = tmp_path / "claude_sessions.jsonl" + then = dt.datetime.now(dt.timezone.utc) - dt.timedelta( + seconds=memory.RETRY_SPACING_S) + log.write_text(json.dumps( + {"ts": then.strftime("%Y-%m-%dT%H:%M:%SZ"), "type": "consolidate"}) + "\n", + encoding="utf-8") + monkeypatch.setattr(cs, "SESSION_LOG", log) + monkeypatch.setattr(cs, "BUDGET_DISABLED", False) + assert cs.check_budget("consolidate") is None + + # --------------------------------------------------------------------------- # Routing to the resident Claude session (the `claude -p` retirement) # --------------------------------------------------------------------------- diff --git a/tests/test_consolidation_background_job.py b/tests/test_consolidation_background_job.py new file mode 100644 index 00000000..1b997ab3 --- /dev/null +++ b/tests/test_consolidation_background_job.py @@ -0,0 +1,434 @@ +"""Nightly consolidation must be able to succeed (#291). + +Three defects made success structurally impossible, and each gets its own +section here: + +1. **The pass ran inline on px-mind's ~60s awareness tick.** Its declared + budget is 600s — twice px-mind's own 300s health-staleness window — so + honouring the budget and keeping the mind loop alive were mutually + exclusive. Now it runs on a daemon thread with a pid-keyed job marker. +2. **An ad-hoc `timeout=180` overrode the declared 600s deadline.** The tighter + number always won, so the declared budget was never once reachable. +3. **Attempt 2 could not be spent.** `MAX_ATTEMPTS_PER_DAY` promised two tries + a night while the `consolidate` quota was 1 and its type cooldown 20 hours. + +Inert by construction: no service is touched, no tmux session is reached, no +Claude call is made, and every duration in here is synthetic — a "600s" +deadline is asserted as a *plumbed number*, never waited for. The one real +thread that runs blocks on an Event the test controls. +""" +from __future__ import annotations + +import datetime as dt +import json +import threading +import time + +import pytest + +import pxh.mind as mind +from pxh import brain, claude_session, memory + +SPARK = {"persona": "spark"} +UTC = dt.timezone.utc + +# This Pi routinely sits above a load average of 10; margins are sized to that +# rather than to an idle machine (a join that finished in 0.06s idle has been +# observed taking 12s here under load). +JOIN_TIMEOUT_S = 60 +# Used only where a *fresh* write has to be distinguishable from a backdated +# one — never as a "the tick was fast enough" budget. See the note in +# test_tick_stays_responsive_while_consolidation_exceeds_300s. +FRESH_WITHIN_S = 120.0 + + +@pytest.fixture(autouse=True) +def _isolate(monkeypatch, tmp_path): + """Keep the job marker, the meta file and the worker global out of live state.""" + monkeypatch.setenv("PX_STATE_DIR", str(tmp_path)) + mind._consolidation_worker = None + yield + worker = mind._consolidation_worker + if worker is not None: + worker.join(timeout=JOIN_TIMEOUT_S) + mind._consolidation_worker = None + + +@pytest.fixture +def due(monkeypatch): + """Report consolidation as due, without depending on the wall clock.""" + monkeypatch.setattr(mind.spark_memory, "consolidation_due", + lambda *a, **kw: None) + + +def _backdate_marker(started_ago_s: float, beat_ago_s: float = 0.0) -> None: + """Age the in-flight marker synthetically. No test ever waits real time.""" + now = dt.datetime.now(UTC) + job = memory.read_consolidation_job() + job["started_ts"] = (now - dt.timedelta(seconds=started_ago_s)).isoformat() + job["heartbeat_ts"] = (now - dt.timedelta(seconds=beat_ago_s)).isoformat() + memory._write_consolidation_job(job) + + +# --------------------------------------------------------------------------- +# 1. The tick cannot be stalled by a long consolidation +# --------------------------------------------------------------------------- + +def test_tick_stays_responsive_while_consolidation_exceeds_300s(due, monkeypatch): + """px-mind's loop keeps running while a >300s consolidation is in flight. + + 300s is not an arbitrary number: it is `health.STALE_AFTER_S["px-mind"]`, + the point at which px-mind itself reads stale. Under the old inline call a + consolidation allowed to use its real 600s budget guaranteed that reading — + no awareness snapshot, no reflection, no battery check, for ten minutes. + + The "600 seconds" here is synthetic: the worker blocks on an Event the test + owns, and the marker is backdated. Nothing sleeps. + + The proof is structural rather than a stopwatch, deliberately. The gate is + held shut for the whole loop below, so a tick that waited on the worker + could not return **at all** — reaching the end of the loop is itself the + assertion. A wall-clock budget would prove nothing extra and would fail on + this host for the wrong reason: it is the live robot, routinely above a load + average of 10, where a single fsync has been measured taking tens of + seconds. Slow is not the property under test; blocked is. + """ + gate = threading.Event() + calls = [] + + def _slow(**kw): + calls.append(kw) + gate.wait(JOIN_TIMEOUT_S) + return {"status": "ok", "written": 1} + + monkeypatch.setattr(mind.spark_memory, "maybe_consolidate", _slow) + try: + mind._consolidation_tick(SPARK, dry=False) + worker = mind._consolidation_worker + assert worker is not None and worker.is_alive() + + # Five further ticks, each standing in for one 60s awareness cycle, + # with the run backdated well past px-mind's own 300s staleness window + # and past the 600s deadline. + for i in range(5): + _backdate_marker(started_ago_s=60 * (i + 1) + 600) + mind._consolidation_tick(SPARK, dry=False) + assert not gate.is_set() + assert worker.is_alive(), ( + f"tick {i} returned only because the worker had finished — " + "it must return while the worker is still running") + + assert len(calls) == 1, "the tick started more than one consolidation" + finally: + gate.set() + if mind._consolidation_worker is not None: + mind._consolidation_worker.join(timeout=JOIN_TIMEOUT_S) + + +def test_no_duplicate_concurrent_consolidation(due, monkeypatch): + """A second tick must not start a second run — nor a second px-mind.""" + gate = threading.Event() + calls = [] + + def _slow(**kw): + calls.append(kw) + gate.wait(JOIN_TIMEOUT_S) + return {"status": "ok", "written": 1} + + monkeypatch.setattr(mind.spark_memory, "maybe_consolidate", _slow) + try: + mind._consolidation_tick(SPARK, dry=False) + first = mind._consolidation_worker + for _ in range(3): + mind._consolidation_tick(SPARK, dry=False) + assert mind._consolidation_worker is first + assert len(calls) == 1 + # And the marker itself refuses a second claimant, so the guard does + # not rest solely on a process-local variable. + assert memory.claim_consolidation_job(attempt=2) is False + finally: + gate.set() + if mind._consolidation_worker is not None: + mind._consolidation_worker.join(timeout=JOIN_TIMEOUT_S) + + +def test_the_tick_heartbeats_the_marker_while_the_worker_runs(due, monkeypatch): + gate = threading.Event() + monkeypatch.setattr(mind.spark_memory, "maybe_consolidate", + lambda **kw: (gate.wait(JOIN_TIMEOUT_S), {"status": "ok"})[1]) + try: + mind._consolidation_tick(SPARK, dry=False) + _backdate_marker(started_ago_s=400, beat_ago_s=400) + mind._consolidation_tick(SPARK, dry=False) + job = memory.read_consolidation_job() + beat = memory._parse_ts(job["heartbeat_ts"]) + age = (dt.datetime.now(UTC) - beat).total_seconds() + assert age < FRESH_WITHIN_S, "the tick did not refresh the heartbeat" + assert not memory.consolidation_job_is_stale(job) + finally: + gate.set() + if mind._consolidation_worker is not None: + mind._consolidation_worker.join(timeout=JOIN_TIMEOUT_S) + + +def test_an_overrunning_worker_is_reported_exactly_once(due, monkeypatch): + """A stuck worker must be visible, and must not spam a failure every 60s.""" + gate = threading.Event() + monkeypatch.setattr(mind.spark_memory, "maybe_consolidate", + lambda **kw: (gate.wait(JOIN_TIMEOUT_S), {"status": "ok"})[1]) + from pxh import health + try: + mind._consolidation_tick(SPARK, dry=False) + _backdate_marker(started_ago_s=memory.JOB_OVERRUN_AFTER_S + 60) + for _ in range(4): + mind._consolidation_tick(SPARK, dry=False) + rec = json.loads( + health._component_path(health.CONSOLIDATION_COMPONENT).read_text()) + assert rec["consecutive_failures"] == 1 + assert "overran" in rec["last_error"] + finally: + gate.set() + if mind._consolidation_worker is not None: + mind._consolidation_worker.join(timeout=JOIN_TIMEOUT_S) + + +# --------------------------------------------------------------------------- +# 1b. A marker cannot survive a restart as a false "in progress" claim +# --------------------------------------------------------------------------- + +def test_a_marker_from_a_dead_owner_is_not_a_running_claim(due, monkeypatch): + """px-mind restarted mid-run: the thread went with it, the file did not. + + The marker is keyed on pid, so a marker whose owner is gone is detectably a + lie. It is cleared and recorded as a *failure* — no memory formed that + night — rather than silently reset. + """ + from pxh import health + memory._write_consolidation_job({ + "status": "running", "pid": 999999, "attempt": 1, + "started_ts": dt.datetime.now(UTC).isoformat(), + "heartbeat_ts": dt.datetime.now(UTC).isoformat(), + }) + ran = [] + monkeypatch.setattr(mind.spark_memory, "maybe_consolidate", + lambda **kw: ran.append(1) or {"status": "ok"}) + + mind._consolidation_tick(SPARK, dry=False) + assert not memory.read_consolidation_job(), "the stale marker was not cleared" + rec = json.loads( + health._component_path(health.CONSOLIDATION_COMPONENT).read_text()) + assert "abandoned" in rec["last_error"] + # Cleaning up is one tick's work; the next tick is free to try again. + mind._consolidation_tick(SPARK, dry=False) + worker = mind._consolidation_worker + if worker is not None: + worker.join(timeout=JOIN_TIMEOUT_S) + assert ran == [1] + + +def test_a_marker_with_a_silent_heartbeat_is_stale(): + """A live pid is not enough: the process may exist and no longer be ticking.""" + import os + fresh = dt.datetime.now(UTC) + old = fresh - dt.timedelta(seconds=memory.JOB_HEARTBEAT_STALE_S + 60) + live_and_beating = {"pid": os.getpid(), "started_ts": old.isoformat(), + "heartbeat_ts": fresh.isoformat()} + live_but_quiet = {"pid": os.getpid(), "started_ts": old.isoformat(), + "heartbeat_ts": old.isoformat()} + no_beat_at_all = {"pid": os.getpid(), "started_ts": fresh.isoformat()} + assert not memory.consolidation_job_is_stale(live_and_beating) + assert memory.consolidation_job_is_stale(live_but_quiet) + assert memory.consolidation_job_is_stale(no_beat_at_all) + assert not memory.consolidation_job_is_stale({}) # absent is not stale + + +def test_a_stale_marker_does_not_block_a_fresh_claim(): + memory._write_consolidation_job({ + "status": "running", "pid": 999999, "attempt": 1, + "started_ts": dt.datetime.now(UTC).isoformat(), + "heartbeat_ts": dt.datetime.now(UTC).isoformat(), + }) + assert memory.claim_consolidation_job(attempt=1) is True + assert memory.read_consolidation_job()["pid"] == __import__("os").getpid() + + +# --------------------------------------------------------------------------- +# 2. One deadline source: the kind's declared 600s +# --------------------------------------------------------------------------- + +def test_consolidate_passes_no_ad_hoc_timeout(monkeypatch, tmp_path): + """memory.consolidate() must let the declared deadline stand. + + The ad-hoc 180 that used to sit here was tighter than the declared 600, and + the tighter number always wins — so the budget the kind declares was never + once reachable. + """ + (tmp_path / "thoughts-spark.jsonl").write_text("".join( + json.dumps({"ts": dt.datetime.now(UTC).isoformat(), + "thought": f"thought {i}"}) + "\n" for i in range(8)), + encoding="utf-8") + seen = {} + + def _fake(session_type, prompt, **kw): + seen["type"] = session_type + seen["kw"] = kw + return claude_session.RunResult( + stdout="[]", stderr="", returncode=0, duration_s=0.0, + model_used="haiku") + + monkeypatch.setattr(claude_session, "run_claude_session", _fake) + memory.consolidate() + assert seen["type"] == "consolidate" + assert "timeout" not in seen["kw"], ( + f"consolidate() still overrides the declared deadline: {seen['kw']}") + + +def test_run_claude_session_defaults_to_the_declared_deadline(monkeypatch): + """`timeout=None` is the default and reaches ask_brain as None. + + None is what makes brain.py's per-kind table authoritative — ask_brain + substitutes `deadline_for_kind(kind)` only when the caller passed nothing. + """ + seen = {} + + def _fake_ask(kind, payload, timeout_s=None, model=None): + seen["kind"] = kind + seen["timeout_s"] = timeout_s + return {"reply": "ok"} + + monkeypatch.setattr(brain, "ask_brain", _fake_ask) + monkeypatch.setattr(claude_session, "BUDGET_DISABLED", True) + monkeypatch.setattr(claude_session, "SESSION_LOG", + claude_session.PROJECT_ROOT / "state" / "nonexistent.jsonl") + monkeypatch.setattr(claude_session, "_log_session", lambda *a, **kw: None) + claude_session.run_claude_session("consolidate", "prompt", allowed_tools="") + assert seen["kind"] == "consolidate" + assert seen["timeout_s"] is None + assert brain.deadline_for_kind("consolidate") == 600 + + +def test_the_declared_600s_is_what_reaches_the_request(monkeypatch): + """End of the plumbing: the request the brain would answer carries 600s. + + Fully inert — the mailbox is conftest's tmp dir, the session state is + stubbed validated, and `inject` is stubbed to fail so ask_brain returns + immediately after writing the request. Nothing waits; the assertion is on + the number in the file. + """ + session = brain.session_for_kind("consolidate") + captured = [] + + def _capture_then_fail(*a, **kw): + # ask_brain cleans the inbox up in its `finally`, so read the request + # here — at the one moment it exists — and then refuse the injection so + # the call returns without waiting for any reply. + for f in brain.inbox_dir(session).glob("*.json"): + captured.append(json.loads(f.read_text())) + return False + + monkeypatch.setattr(brain, "session_state", lambda *a, **kw: brain.VALIDATED) + monkeypatch.setattr(brain.tmux_claude, "inject", _capture_then_fail) + before = time.time() + assert brain.ask_brain("consolidate", {"prompt": "x"}) is None + + assert len(captured) == 1 + request = captured[0] + assert request["kind"] == "consolidate" + budget = request["deadline"] - before + # A window, not an equality: ask_brain deducts the time already spent + # waiting for validation and the lock from the caller's budget. + assert 580 <= budget <= 601, f"deadline carried {budget:.0f}s, expected ~600" + + +# --------------------------------------------------------------------------- +# 3. Attempt 2 is actually reachable +# --------------------------------------------------------------------------- + +def test_quota_matches_the_attempts_memory_promises(): + """A quota of 1 against MAX_ATTEMPTS_PER_DAY=2 made the retry unspendable.""" + assert (claude_session._TYPE_QUOTAS["consolidate"] + == memory.MAX_ATTEMPTS_PER_DAY == 2) + + +def test_retry_spacing_clears_the_global_cooldown(): + """Attempt 2 is spaced past the 30-min global cooldown, not exempted from it. + + Exempting `consolidate` (as `self_debug` and `blog` are) would buy nothing + here — nobody is waiting on a 3am retry — and every exemption is one more + way for the nightly job to crowd a session someone *is* waiting on. Spacing + is the cheaper answer, but it only works if the gap really is larger. + """ + assert memory.RETRY_SPACING_S > claude_session.COOLDOWN_S + assert memory.RETRY_SPACING_S >= claude_session._TYPE_COOLDOWNS["consolidate"] + assert "consolidate" not in claude_session._GLOBAL_COOLDOWN_EXEMPT + # Two spaced attempts still fit inside the 02:00-06:00 window. + window_s = (memory.CONSOLIDATION_WINDOW[1] + - memory.CONSOLIDATION_WINDOW[0]) * 3600 + assert memory.RETRY_SPACING_S * (memory.MAX_ATTEMPTS_PER_DAY - 1) < window_s + + +def test_a_failed_attempt_one_leaves_attempt_two_reachable(monkeypatch): + """The end-to-end gate: fail at 02:00, and 02:40 is a real second attempt.""" + at2 = dt.datetime(2026, 7, 11, 2, 0, tzinfo=memory.HOBART_TZ) + calls = [] + + def _fail(**kw): + calls.append(kw) + return {"status": "failed", "error": "brain unavailable"} + + monkeypatch.setattr(memory, "consolidate", _fail) + assert memory.maybe_consolidate(now=at2)["status"] == "failed" + # Immediately after, the spacing gate holds it back... + assert memory.consolidation_due(now=at2 + dt.timedelta(minutes=5)) is not None + # ...and once the gap has passed, attempt 2 really runs. + later = at2 + dt.timedelta(seconds=memory.RETRY_SPACING_S) + assert memory.consolidation_due(now=later) is None + assert memory.maybe_consolidate(now=later)["status"] == "failed" + assert len(calls) == 2 + meta = json.loads(memory.consolidation_meta_file().read_text()) + assert meta["attempts"] == 2 and meta["done"] is False + + +def test_the_budget_gate_admits_attempt_two(monkeypatch, tmp_path): + """claude_session's own quota/cooldown must not be what blocks the retry.""" + log = tmp_path / "claude_sessions.jsonl" + then = dt.datetime.now(UTC) - dt.timedelta(seconds=memory.RETRY_SPACING_S) + log.write_text(json.dumps({ + "ts": then.isoformat().replace("+00:00", "Z"), + "session_type": "consolidate", "model": "haiku", + "duration_s": 1.0, "returncode": 1, "status": "brain_unavailable", + }) + "\n", encoding="utf-8") + monkeypatch.setattr(claude_session, "SESSION_LOG", log) + monkeypatch.setattr(claude_session, "BUDGET_DISABLED", False) + assert claude_session.check_budget("consolidate") is None + + +def test_success_marks_the_day_done(monkeypatch): + at2 = dt.datetime(2026, 7, 11, 2, 0, tzinfo=memory.HOBART_TZ) + monkeypatch.setattr(memory, "consolidate", + lambda **kw: {"status": "ok", "written": 3}) + assert memory.maybe_consolidate(now=at2)["status"] == "ok" + later = at2 + dt.timedelta(seconds=memory.RETRY_SPACING_S) + assert memory.consolidation_due(now=later) == "already done for this date" + + +def test_a_correct_skip_also_marks_the_day_done(monkeypatch): + # Too few thoughts in 24h is a correct outcome; retrying it at 03:40 would + # spend a session to reach the same answer. + at2 = dt.datetime(2026, 7, 11, 2, 0, tzinfo=memory.HOBART_TZ) + monkeypatch.setattr(memory, "consolidate", + lambda **kw: {"status": "skipped", "reason": "3 thoughts"}) + assert memory.maybe_consolidate(now=at2)["status"] == "skipped" + assert memory.consolidation_due( + now=at2 + dt.timedelta(hours=1)) == "already done for this date" + + +def test_the_window_and_meta_shape_are_unchanged(): + """The existing contract other code and the operator read.""" + assert memory.CONSOLIDATION_WINDOW == (2, 6) + assert memory.consolidation_due( + now=dt.datetime(2026, 7, 11, 12, 0, tzinfo=memory.HOBART_TZ) + ) == "outside the 02:00-06:00 window" + assert memory.consolidation_meta_file().name == "consolidation_meta.json" + at2 = dt.datetime(2026, 7, 11, 2, 30, tzinfo=memory.HOBART_TZ) + assert memory.consolidation_due(now=at2) is None diff --git a/tests/test_health_memory_formation.py b/tests/test_health_memory_formation.py index aefd3364..f7ad558f 100644 --- a/tests/test_health_memory_formation.py +++ b/tests/test_health_memory_formation.py @@ -240,3 +240,45 @@ def test_motd_overdue_threshold_matches_health_module(motd): # Duplicated constant (px-motd imports no pxh module by design) — pin the # two together so they cannot silently drift. assert motd.MEMORY_OVERDUE_S == health.MEMORY_FORMATION_OVERDUE_S + + +# --- the in-flight job marker (#291) ---------------------------------------- + + +def _job(started_ago_s: int, beat_ago_s: int) -> dict: + now = dt.datetime.now(dt.timezone.utc) + return { + "status": "running", "pid": 1234, "attempt": 1, + "started_ts": (now - dt.timedelta(seconds=started_ago_s)).isoformat(), + "heartbeat_ts": (now - dt.timedelta(seconds=beat_ago_s)).isoformat(), + } + + +def test_motd_shows_a_consolidation_in_flight(motd): + ts = dt.datetime.now(dt.timezone.utc) - dt.timedelta(hours=20) + line = _plain(motd._memory_formation_line( + {"last_success_ts": ts.strftime("%Y-%m-%dT%H:%M:%SZ")}, + _job(started_ago_s=180, beat_ago_s=5))) + assert "consolidating now for 3m" in line + # The hint is additive: the honest age of the last real memory stays. + assert "last formed 20h ago" in line + + +def test_motd_ignores_a_marker_whose_heartbeat_went_quiet(motd): + # A marker outliving its owner is a lie; px-motd must not repeat it. + ts = dt.datetime.now(dt.timezone.utc) - dt.timedelta(hours=20) + line = _plain(motd._memory_formation_line( + {"last_success_ts": ts.strftime("%Y-%m-%dT%H:%M:%SZ")}, + _job(started_ago_s=9000, beat_ago_s=8000))) + assert "consolidating" not in line + + +def test_motd_line_renders_without_a_job_argument(motd): + # The marker is absent almost all the time; the default must stay valid. + assert "never formed" in _plain(motd._memory_formation_line({})) + assert "never formed" in _plain(motd._memory_formation_line({}, {})) + + +def test_motd_job_heartbeat_threshold_matches_memory_module(motd): + from pxh import memory + assert motd.JOB_HEARTBEAT_STALE_S == memory.JOB_HEARTBEAT_STALE_S diff --git a/tests/test_memory.py b/tests/test_memory.py index 639bbf5e..ee0fcbc4 100644 --- a/tests/test_memory.py +++ b/tests/test_memory.py @@ -387,10 +387,16 @@ def test_maybe_consolidate_runs_once_then_stamps(): def test_maybe_consolidate_two_failures_stop_for_the_day(): at3 = dt.datetime(2026, 7, 11, 3, 0, tzinfo=HOBART) + # Attempt 2 is deliberately spaced past attempt 1 (#291), so the retry has + # to be asked for at a later clock — an immediate re-ask is refused. + at3b = at3 + dt.timedelta(seconds=memory.RETRY_SPACING_S) with patch.object(memory, "consolidate", return_value={"status": "failed", "error": "x"}) as mc: assert memory.maybe_consolidate(now=at3)["status"] == "failed" - assert memory.maybe_consolidate(now=at3)["status"] == "failed" - assert memory.maybe_consolidate(now=at3) is None # attempt cap + assert memory.maybe_consolidate(now=at3) is None # too soon + assert memory.maybe_consolidate(now=at3b)["status"] == "failed" + assert memory.maybe_consolidate( + now=at3b + dt.timedelta(seconds=memory.RETRY_SPACING_S) + ) is None # attempt cap assert mc.call_count == 2 @@ -399,7 +405,8 @@ def test_maybe_consolidate_fresh_date_resets_attempts(): day2 = dt.datetime(2026, 7, 12, 3, 0, tzinfo=HOBART) with patch.object(memory, "consolidate", return_value={"status": "failed", "error": "x"}): memory.maybe_consolidate(now=day1) - memory.maybe_consolidate(now=day1) + memory.maybe_consolidate( + now=day1 + dt.timedelta(seconds=memory.RETRY_SPACING_S)) with patch.object(memory, "consolidate", return_value={"status": "ok", "written": 1}) as mc: assert memory.maybe_consolidate(now=day2)["status"] == "ok" assert mc.call_count == 1 diff --git a/tests/test_mind_consolidation_health.py b/tests/test_mind_consolidation_health.py index a57e4df2..7c4c23f7 100644 --- a/tests/test_mind_consolidation_health.py +++ b/tests/test_mind_consolidation_health.py @@ -4,8 +4,13 @@ "attempted and failed" must not read the same as "not due tonight", and "not due" must not refresh the record every 60s forever. +Since #291 the pass runs on a background daemon thread, so `_consolidation_tick` +only *starts* it; `run_tick` below waits for the worker so these tests keep +asserting on the finished outcome rather than on a race. + Inert: conftest's autouse `_isolate_health_writes` redirects health_dir() into -tmp, and maybe_consolidate is stubbed — no LLM call, no live state. +tmp, `PX_STATE_DIR` is redirected per-test, and maybe_consolidate is stubbed — +no LLM call, no live state. """ from __future__ import annotations @@ -25,6 +30,37 @@ def _record() -> dict: return json.loads(path.read_text()) if path.exists() else {} +@pytest.fixture(autouse=True) +def _isolate(monkeypatch, tmp_path): + """Keep the job marker and the due-gate out of the live state dir.""" + monkeypatch.setenv("PX_STATE_DIR", str(tmp_path)) + # The gate is exercised on its own in test_memory.py; here it is always + # "due" so each test is about the outcome reporting. + monkeypatch.setattr(mind.spark_memory, "consolidation_due", + lambda *a, **kw: None) + mind._consolidation_worker = None + yield + worker = mind._consolidation_worker + if worker is not None: + worker.join(timeout=JOIN_TIMEOUT_S) + mind._consolidation_worker = None + + +# Generous on purpose: this Pi routinely sits at a load average above 10, and a +# join margin sized to an idle machine turns a passing test into a flake. +JOIN_TIMEOUT_S = 60 + + +def run_tick(session=SPARK, dry=False, timeout=JOIN_TIMEOUT_S): + """Start the tick and wait for whatever worker it launched.""" + mind._consolidation_tick(session, dry) + worker = mind._consolidation_worker + if worker is not None: + worker.join(timeout=timeout) + assert not worker.is_alive(), "consolidation worker did not finish" + mind._consolidation_worker = None + + @pytest.fixture def stub(monkeypatch): """Replace maybe_consolidate with a canned result.""" @@ -36,7 +72,7 @@ def _set(result): def test_success_is_recorded(stub): stub({"status": "ok", "written": 6}) - mind._consolidation_tick(SPARK, dry=False) + run_tick() rec = _record() assert rec["success_count"] == 1 assert rec["last_success_detail"]["written"] == 6 @@ -47,7 +83,7 @@ def test_skipped_counts_as_a_healthy_night(stub): # maybe_consolidate marks the date done for it. Recording a success keeps a # quiet-but-working night from ageing into "stale". stub({"status": "skipped", "reason": "only 3 thoughts in 24h"}) - mind._consolidation_tick(SPARK, dry=False) + run_tick() rec = _record() assert rec["success_count"] == 1 assert rec["consecutive_failures"] == 0 @@ -55,7 +91,7 @@ def test_skipped_counts_as_a_healthy_night(stub): def test_failure_is_recorded_with_its_error(stub): stub({"status": "failed", "error": "claude exit 1"}) - mind._consolidation_tick(SPARK, dry=False) + run_tick() rec = _record() assert rec["consecutive_failures"] == 1 assert rec["last_error"] == "claude exit 1" @@ -70,16 +106,18 @@ def test_not_due_records_nothing(stub): change exists to remove. Staleness already covers true silence. """ stub(None) - mind._consolidation_tick(SPARK, dry=False) + run_tick() assert not health._component_path(CONS).exists() def test_dry_run_records_nothing(stub): - # A dry run forms no memory, so it must not refresh the record. (Live - # maybe_consolidate returns None under dry; this pins the belt as well as - # the braces, since the status is what the tick actually branches on.) + # A dry run forms no memory, so it must not refresh the record. Since #291 + # the tick refuses to even launch a worker under dry, so nothing runs at + # all — but the worker still branches on the status, so pin both. stub({"status": "dry"}) - mind._consolidation_tick(SPARK, dry=True) + run_tick(dry=True) + assert not health._component_path(CONS).exists() + mind._record_consolidation_outcome({"status": "dry"}) assert not health._component_path(CONS).exists() @@ -87,7 +125,7 @@ def test_a_raising_consolidation_is_recorded_not_swallowed(stub, monkeypatch): def _boom(**kw): raise RuntimeError("mailbox gone") monkeypatch.setattr(mind.spark_memory, "maybe_consolidate", _boom) - mind._consolidation_tick(SPARK, dry=False) # must not raise + run_tick() # must not raise, and must not lose the error to the thread rec = _record() assert rec["consecutive_failures"] == 1 assert "mailbox gone" in rec["last_error"] @@ -96,12 +134,12 @@ def _boom(**kw): def test_non_spark_persona_reports_nothing(stub): # Consolidation only runs for SPARK; GREMLIN's night is not a missed one. stub({"status": "ok", "written": 3}) - mind._consolidation_tick({"persona": "gremlin"}, dry=False) + run_tick({"persona": "gremlin"}) assert not health._component_path(CONS).exists() def test_three_failed_nights_read_failing(stub): stub({"status": "failed", "error": "claude exit 1"}) for _ in range(3): - mind._consolidation_tick(SPARK, dry=False) + run_tick() assert health.read_health()["components"][CONS]["status"] == "failing" diff --git a/tests/test_mind_utils.py b/tests/test_mind_utils.py index 4a553f88..a5401128 100644 --- a/tests/test_mind_utils.py +++ b/tests/test_mind_utils.py @@ -1607,39 +1607,81 @@ def test_reflection_injects_active_intention(tmp_path, monkeypatch): assert "current intention" in ctx.lower() -def test_consolidation_tick_spark_logs_result(monkeypatch): +_SPARK_SESSION = {"persona": "spark"} + + +def _run_consolidation_tick(session=_SPARK_SESSION, dry=False): + """Tick, then wait for the worker it launched (#291 made the pass async). + + Default is a sentinel, not `None`-or-truthy: `{}` is a real session under + test here (no persona at all) and must not be swapped for SPARK's. + """ + pxh.mind._consolidation_tick(session, dry=dry) + worker = pxh.mind._consolidation_worker + if worker is not None: + worker.join(timeout=60) + pxh.mind._consolidation_worker = None + + +@pytest.fixture +def consolidation_due(monkeypatch, tmp_path): + """Report the pass as due, with the job marker isolated into tmp.""" + monkeypatch.setenv("PX_STATE_DIR", str(tmp_path)) + monkeypatch.setattr(pxh.mind.spark_memory, "consolidation_due", + lambda *a, **kw: None) + pxh.mind._consolidation_worker = None + yield + pxh.mind._consolidation_worker = None + + +def test_consolidation_tick_spark_logs_result(monkeypatch, consolidation_due): calls = [] monkeypatch.setattr(pxh.mind.spark_memory, "maybe_consolidate", lambda dry: {"status": "ok", "written": 2}) monkeypatch.setattr(pxh.mind, "log", lambda msg: calls.append(msg)) - pxh.mind._consolidation_tick({"persona": "spark"}, dry=False) + _run_consolidation_tick() assert any("consolidation: ok" in c and "wrote 2" in c for c in calls) -def test_consolidation_tick_none_is_silent(monkeypatch): +def test_consolidation_tick_none_is_silent(monkeypatch, consolidation_due): calls = [] monkeypatch.setattr(pxh.mind.spark_memory, "maybe_consolidate", lambda dry: None) monkeypatch.setattr(pxh.mind, "log", lambda msg: calls.append(msg)) + _run_consolidation_tick() + # The launch line is the tick's own; the *outcome* is what must stay silent. + assert not any("consolidation:" in c and "started" not in c for c in calls) + + +def test_consolidation_tick_not_due_launches_nothing(monkeypatch, tmp_path): + """The cheap read-only gate runs before any thread is spent.""" + monkeypatch.setenv("PX_STATE_DIR", str(tmp_path)) + pxh.mind._consolidation_worker = None + called = [] + monkeypatch.setattr(pxh.mind.spark_memory, "consolidation_due", + lambda *a, **kw: "outside the 02:00-06:00 window") + monkeypatch.setattr(pxh.mind.spark_memory, "maybe_consolidate", + lambda dry: called.append(1)) pxh.mind._consolidation_tick({"persona": "spark"}, dry=False) - assert calls == [] + assert pxh.mind._consolidation_worker is None + assert called == [] -def test_consolidation_tick_never_raises(monkeypatch): +def test_consolidation_tick_never_raises(monkeypatch, consolidation_due): def boom(dry): raise RuntimeError("disk exploded") calls = [] monkeypatch.setattr(pxh.mind.spark_memory, "maybe_consolidate", boom) monkeypatch.setattr(pxh.mind, "log", lambda msg: calls.append(msg)) - pxh.mind._consolidation_tick({"persona": "spark"}, dry=False) # must not raise + _run_consolidation_tick() # must not raise, on this thread or the worker's assert any("consolidation error" in c for c in calls) -def test_consolidation_tick_skips_other_personas(monkeypatch): +def test_consolidation_tick_skips_other_personas(monkeypatch, consolidation_due): called = [] monkeypatch.setattr(pxh.mind.spark_memory, "maybe_consolidate", lambda dry: called.append(1)) - pxh.mind._consolidation_tick({"persona": "gremlin"}, dry=False) - pxh.mind._consolidation_tick({}, dry=False) + _run_consolidation_tick({"persona": "gremlin"}) + _run_consolidation_tick({}) assert called == []