diff --git a/CHANGELOG.md b/CHANGELOG.md index 4e1bffe..8cb9f21 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,63 @@ # Changelog +## 2.1.0 + +The "quick-wins train" from the 2026-07-28 review (`docs/IMPROVEMENT-PLANS.md` #1, +#3, #7). Three independently useful changes; no schema or layout changes. + +### Added — `load --brief`: a token-budgeted cold-start digest + +The SessionStart hook injected the full text of six files every session (measured +15.9KB and growing without bound — Observations and Scope History never shrink). +`load --brief` renders a digest instead: Fact / Pattern rules in full, the last ~10 +Observations (with a `+N older` note), scope + `scope_updated` + sessions-since, +the heartbeat-pointed session log's key sections (Goal / Decisions Made / Open +Threads / Suggested Next Step), today's Agent Log lines, and the inbox as a +**count + oldest age** — everything the load-time reconcile needs. Over +`ECHO_LOAD_BUDGET` (default ~8000 chars), sections trim lowest-priority-first +(observations → agent log → session summary) with explicit `(truncated)` markers. +The hook now injects the brief form; `/echo-load` and the manual procedure keep the +full dump. All `load` reads (plus the sessions listing, which also feeds the +heartbeat-absent fallback) are fetched **in parallel** over the pooled connections. + +### Fixed — capture is now offline-durable (the write path's inversion bug) + +Low-level verbs queued on an outage, but the *default* write (`capture`) died: the +update path's POST/PATCH and `ensure_daily_log` bypassed the queue entirely, and an +unreachable index aborted the whole op. Now: (1) a fully-offline `capture` queues +the **whole operation as one semantic record** (title/kind/body/tags/…, plus the +capture-time date); `flush` replays it **through `capture`** so create-vs-update, +the duplicate gate, and aliasing re-run against the index as it is at replay time — +byte-level replay would freeze a decision made blind. (2) A duplicate-gate stop on +replay keeps the record, flags it `needs_attention` (surfaced by `flush` and +`load`), and does NOT block later records. (3) Mid-capture outages: the update +path and `ensure_daily_log` writes ride `safe_request` (accepted edge: a queued +agent-log PATCH 400-drops loudly if the daily note never existed — the trace line, +not the memory). Same-args offline captures dedupe by idem-key. + +### Added — `session-end`: the one-call session close + +`echo.py session-end [--apply]` performs the whole end-of-session +sequence under ONE lock acquisition, in order: session log PUT → Agent Log line → +reflect proposals (via `capture`, gate-aware) → optional `scope set` → **heartbeat +LAST as the commit marker** (a part-way failure leaves the previous pointer intact, +never dangling). Dry-run by default previews the full plan including the reflect +classification. Filename HHMM from `$ECHO_NOW` (else local clock), validated +against the canonical session-log regex **before any write**. The Stop hook's nudge +now names the command, and recognizes a `session-end` invocation as +"already reflected". + +### Fixed — helper-module errors exited 2 instead of their intended codes + +When `echo.py` runs as `__main__`, helper modules' `import echo` creates a twin +module whose `EchoError` is a different class object, so `main()`'s handler missed +it and the raw traceback leaked. The handler now catches `RuntimeError` (both +classes subclass it) and honors `.code`. + +Tests: +19 end-to-end checks (brief-load digest/budget, offline-capture queue / +replay-through-capture / gate-on-replay / idem-key dedup, session-end dry-run / +apply / ordering / validation-abort). All suites green. + ## 2.0.0 **Packaging & structure major — nothing about memory behavior changes.** The major diff --git a/README.md b/README.md index 9fa48c8..e17007f 100644 --- a/README.md +++ b/README.md @@ -1,4 +1,4 @@ -# echo-memory — v2.0.0 +# echo-memory — v2.1.0 Persistent memory for Claude / CoWork sessions via the **ECHO** Obsidian vault, driven over the [Obsidian Local REST API](https://github.com/coddingtonbear/obsidian-local-rest-api). The skill makes direct REST calls through a bundled validated client (`scripts/echo.py`). The whole toolchain is **pure Python** (stdlib only), so it runs identically on Windows, macOS, and Linux — no bash, no platform-specific `date`. @@ -421,6 +421,7 @@ From the credential-free harness (`eval/run_eval.py` against the deterministic m | Version | Highlights | |---------|-----------| +| **2.1.0** | **Quick-wins train: cheaper sessions, durable capture, one-call session end.** (1) **`load --brief`** — a token-budgeted cold-start digest (Fact/Pattern in full, last ~10 Observations, scope+freshness, last session's key sections, Agent-Log lines, inbox *count*; `ECHO_LOAD_BUDGET` default ~8000 chars) now injected by the SessionStart hook, replacing the full six-file dump that grew unboundedly; all `load` reads are fetched in parallel. (2) **Offline capture durability** — `capture` on an unreachable vault queues the *whole operation* as one semantic record; `flush` replays it **through capture** so routing/gate/aliasing re-run against the current index; a gate stop on replay is kept + flagged, never landed blind; `ensure_daily_log`/update-path writes ride the queue too. (3) **`session-end`** — one call (one lock) writes session log → Agent-Log line → reflect proposals → optional scope switch → **heartbeat last as the commit marker**; dry-run by default; `ECHO_NOW` pins the HHMM. Also fixes the `__main__` twin-module trap so helper-module `EchoError`s exit with their intended codes. +19 end-to-end checks. | | **2.0.0** | **Packaging & structure major — memory behavior unchanged.** Skills-only: the legacy `commands/` directory is deleted (breaking for pre-skills clients; the 1.6.0-verified skills are the sole entry points). Artifact policy: the repo tracks only the `echo-memory.plugin` pointer; the 15 historical versioned zips leave the tree and versioned builds ship as **Gitea releases** (one per `v` tag) from `v2.0.0` on; `.gitignore` blocks `*.plugin` except the pointer. README gains the "packaging for other agent runtimes" port note (build-target rule, never a second source tree). | | **1.6.0** | **Skills-format migration + lint blind-spot fix.** The eight slash commands gain `skills//SKILL.md` twins (same names, byte-identical bodies; a same-name skill takes precedence, so `commands/` coexists until its 2.0 deletion): per-skill `allowed-tools` end the permission prompts, `disable-model-invocation: true` makes `/echo-sweep` + `/echo-triage` operator-only. `vault_lint` checks retired patterns **before** routes, so a seed README inside a retired tree (`archive/…`) is flagged instead of being absolved by the leaf-README route (+ regression test). `plugin.json` gains `homepage`/`repository`. | | **1.5.1** | **Duplicate-gate precision.** Live use caught the gate blocking a decision note because it shared the single vault-common token "echo" with an archived project (min-normalized scoring gives any single-token entity a 1.0 match). New `gate_candidates()` blocks only on **same-kind** candidates (cross-kind name collisions — a decision titled after its project — warn instead), and a **lone shared token blocks only when unique to that entity** (vault df == 1); multi-token overlaps still block. Advisory warnings (`fuzzy_candidates`) unchanged; exit-76/`--merge-into`/`--force` unchanged. +5 tests; verified against the live vault both ways (false positive creates cleanly, true lookalike still gates). | diff --git a/docs/IMPROVEMENT-PLANS.md b/docs/IMPROVEMENT-PLANS.md index 6749d42..1976898 100644 --- a/docs/IMPROVEMENT-PLANS.md +++ b/docs/IMPROVEMENT-PLANS.md @@ -345,8 +345,9 @@ manifest-verification pass). | Release | Items | Rationale | |---|---|---| | **1.6.0** | TODO-1.6: skills-format migration (+ per-skill `allowed-tools`, `disable-model-invocation` on write-heavy entries, `$ECHO` block dedupe) · lint blind-spot fix (retired-dir/leaf-README) · plugin.json `homepage` | Restructures the command/SKILL files later work edits; allowed-tools removes prompt friction for every verb added below. Vault link repointing (TODO-1.6 §3) is vault data, not plugin code — any session. | -| **1.6.x quick wins** | #1 brief load · #3 offline capture · #7 session-end (+ return-not-print refactor, MCP spec Phase 0, if it fits) | Small, correctness/efficiency, no schema changes. #7 unblocks the MCP tool. Users get these without a reinstall — don't gate them behind 2.0. | -| **2.0** | ROADMAP-2.0 unchanged: delete legacy `commands/` · artifact policy (Gitea releases, gitignore, tags) · Codex note · README layout | Packaging only, no feature riders — a bad feature must never force reverting a packaging break. | -| **2.1** | MCP server (containerized, `echomcp.alwisp.com`) — `docs/MCP-SERVER-SPEC.md` | Needs #7 + Phase 0. Lands after 2.0 so its docs target one command format and the Gitea release channel exists. Deployed via CI + PORT, not the plugin manifest. | -| **2.2 index train** | #2 local-first recall (CLI-path half) · #4 incremental sweep · #8 stemming+expansion | One schema-bump train; all touch the same index/sweep loop. The container already keeps its index warm in memory (spec §7), so #2 here covers the CLI fallback path; #4's fast sweep also becomes the container's background job. | -| **2.3 memory train** | #5 lifecycle/decay · #6 supersession · #9 resolve-miss learning | Capability tier; #6 after #5 (status vocabulary). | +| **1.6.0 — SHIPPED 2026-07-28** | Skills migration + lint fix + homepage | Verified on desktop + CoWork. | +| **2.0.0 — SHIPPED 2026-07-28** | ROADMAP-2.0: commands/ deleted · release-based artifacts · gitignore · Codex note | Packaging-only major; skills-only install re-verified. | +| **2.1.0 quick wins — SHIPPED 2026-07-28** | #1 brief load · #3 offline capture · #7 session-end | All three landed with +19 e2e checks; see CHANGELOG 2.1.0. The MCP spec's Phase 0 (return-not-print) did NOT ride along — it's the first phase of the MCP build session. | +| **2.2** | MCP server (containerized, `echomcp.alwisp.com`) — `docs/MCP-SERVER-SPEC.md` | Needs Phase 0 (return-not-print refactor, first step of the build session). #7 session-end shipped, so `echo_log_session` wraps a real verb. Deployed via CI + PORT, not the plugin manifest. | +| **2.3 index train** | #2 local-first recall (CLI-path half) · #4 incremental sweep · #8 stemming+expansion | One schema-bump train; all touch the same index/sweep loop. The container already keeps its index warm in memory (spec §7), so #2 here covers the CLI fallback path; #4's fast sweep also becomes the container's background job. | +| **2.4 memory train** | #5 lifecycle/decay · #6 supersession · #9 resolve-miss learning | Capability tier; #6 after #5 (status vocabulary). | diff --git a/docs/MCP-SERVER-SPEC.md b/docs/MCP-SERVER-SPEC.md index 79a6a41..5bc1f8d 100644 --- a/docs/MCP-SERVER-SPEC.md +++ b/docs/MCP-SERVER-SPEC.md @@ -75,7 +75,7 @@ monitor, PORT deploy/rollback. Server down ⇒ exactly today's behavior. | Tool prefix | `echo_` | Namespace safety next to other servers. | | Result shape | `structuredContent` (the `echo_output.envelope` dict) + a 1–3 line human text block | Envelope shape already exists. | | Statefulness | Stateless protocol; warm in-process caches | See §7. | -| Version target | **2.1** (after the 1.6.x train and the 2.0 packaging major) | Needs #7 + Phase 0; docs then target the single post-2.0 command format. | +| Version target | **2.2** (1.6.0, 2.0.0, and the 2.1.0 quick-wins train all shipped 2026-07-28; `session-end` exists) | Remaining prerequisite: Phase 0 (return-not-print), the build session's first step. | | Local stdio variant | **Dropped for v1** (documented, not built) | The CLI/skill fallback covers the no-server case; two transports = two test matrices for little gain. Phase 0 keeps the door open — the `*_op` core is transport-agnostic. | | Result budgets | Every read tool takes an explicit size budget (`budget_chars` / `max_chars`) and truncates with a marker + "call X for more" | The single biggest token lever: one right-sized answer instead of follow-up fetches. | | Tool profiles | `ECHO_MCP_TOOLS=core\|full` env: `core` exposes only load/recall/capture/triage_inbox/log_session/health | Tool schemas cost context on every surface that lists them; the claude.ai connector doesn't need the escape hatches. Desktop registers `full`. | diff --git a/echo-memory.plugin b/echo-memory.plugin index 7c5d32f..8659cd4 100644 Binary files a/echo-memory.plugin and b/echo-memory.plugin differ diff --git a/echo-memory.plugin.src/.claude-plugin/plugin.json b/echo-memory.plugin.src/.claude-plugin/plugin.json index eb3357b..5f57621 100644 --- a/echo-memory.plugin.src/.claude-plugin/plugin.json +++ b/echo-memory.plugin.src/.claude-plugin/plugin.json @@ -1,6 +1,6 @@ { "name": "echo-memory", - "version": "2.0.0", + "version": "2.1.0", "homepage": "https://git.alwisp.com/jason/echo", "repository": "https://git.alwisp.com/jason/echo", "description": "Persistent memory via the ECHO Obsidian vault over the Obsidian Local REST API. Cross-platform Python client: one-call capture/resolve/recall/link/triage over an entity index, hybrid BM25 + graph recall spanning entities + sessions/journal (recency/status-aware), a pre-write duplicate gate, complete-frontmatter capture, session hooks that self-fire load/reflect, offline write-ahead queue, lock-guarded concurrency, linter-enforced routing, and /echo-* commands.", diff --git a/echo-memory.plugin.src/skills/echo-memory/SKILL.md b/echo-memory.plugin.src/skills/echo-memory/SKILL.md index 0b141b0..fe60e63 100644 --- a/echo-memory.plugin.src/skills/echo-memory/SKILL.md +++ b/echo-memory.plugin.src/skills/echo-memory/SKILL.md @@ -140,7 +140,7 @@ Load at the start of any substantive conversation — anything beyond a single q ### Loading procedure -**Run `python3 "$ECHO" load`** — one call that performs the six orientation reads below and prints each labeled section (404s on today's note / inbox show as absent, not errors; an absent marker is flagged). This replaces issuing six separate GETs by hand. The table documents what `load` reads and why: +**Run `python3 "$ECHO" load`** — one call that performs the six orientation reads below (fetched in parallel) and prints each labeled section (404s on today's note / inbox show as absent, not errors; an absent marker is flagged). This replaces issuing six separate GETs by hand. **`load --brief`** renders a token-budgeted digest instead — Fact/Pattern rules in full, the last ~10 Observations, scope + freshness, the last session's key sections, today's Agent Log lines, and the inbox as a *count* (`ECHO_LOAD_BUDGET`, default ~8000 chars, trims lowest-priority sections first). The SessionStart hook injects the brief form; `/echo-load` and this procedure use the full form. The table documents what `load` reads and why: | # | GET | Notes | |---|-----|-------| @@ -405,25 +405,40 @@ At the end of substantive conversations (ones that produced decisions, artifacts **Filename format is canonical: `_agent/sessions/YYYY-MM-DD-HHMM-.md`.** Always include the four-digit local-time HHMM component — it makes filenames lexically sort in true chronological order, which is what Step 3 of loading relies on. Older session logs without the HHMM part exist; leave them alone, but every new one must use the full form. -Write the body (see `references/session-log-template.md`) to a file with the Write tool, then: +**Preferred — one call (`session-end`).** Write a bundle JSON with the Write tool and run it; it performs the whole session-end sequence under one lock, in order, with the **heartbeat written last as the commit marker** (a failure part-way leaves the previous pointer intact): + +```json +{ "slug": "", + "log_body": "", + "agent_log_line": "- : ", + "scope": "", + "reflect": [ { "title": "…", "kind": "…", "body": "…", "confidence": 0.9 } ] } +``` + +```bash +python3 "$ECHO" session-end bundle.json # dry-run: previews the full plan +ECHO_NOW= python3 "$ECHO" session-end bundle.json --apply +``` + +`--apply` writes: session log → Agent Log line → reflect proposals (via `capture`, gate-aware) → optional scope switch → heartbeat. Set `ECHO_NOW` to the conversation's local HHMM so the filename sorts truthfully. Reflect proposals still follow the reflect contract — preview with the dry-run and apply only with the operator's go-ahead. Offline, every step queues durably. + +**Manual fallback** (only if `session-end` is unavailable) — the same sequence by hand: ```bash python3 "$ECHO" put _agent/sessions/-HHMM-.md ``` -Then add a one-line entry to today's daily note via the **Daily Note — Agent Log** procedure above. - -Finally, update the heartbeat pointer so the next session can orient in one read (this closes the loop with loading-procedure Step 4 — a pointer nobody writes is a pointer nobody can read). It is a single line, ` @ `; write it to a file and PUT it: +Then add a one-line entry to today's daily note via the **Daily Note — Agent Log** procedure above, and finally update the heartbeat pointer so the next session can orient in one read (a pointer nobody writes is a pointer nobody can read) — a single line, ` @ `: ```bash python3 "$ECHO" put _agent/heartbeat/last-session.md ``` -`last-session.md` is overwritten (PUT) each session end — never appended, so it can't grow or duplicate. +`last-session.md` is overwritten (PUT) each session end — never appended, so it can't grow or duplicate. Keep the heartbeat write **last** in the manual sequence too. ## Vault Unreachable -If the API returns a connection error, timeout, or `502`, tell the operator once that the memory vault is unreachable (a `502` usually means Obsidian/the REST plugin is not running on the backend), then proceed without memory. Do not retry repeatedly. +If the API returns a connection error, timeout, or `502`, tell the operator once that the memory vault is unreachable (a `502` usually means Obsidian/the REST plugin is not running on the backend), then proceed — reads degrade to the last-known-good cache at load, and **writes queue durably**: the low-level verbs queue the request, and `capture` queues the whole operation as one semantic record that replays *through capture* (routing + duplicate gate re-run against the index at replay time) on the next reachable session. A queued capture that stops at the duplicate gate on replay is kept and flagged (`flush` lists it) rather than landed blind. Do not retry repeatedly. ## Style Rules diff --git a/echo-memory.plugin.src/skills/echo-memory/scripts/echo.py b/echo-memory.plugin.src/skills/echo-memory/scripts/echo.py index 0f3b1ec..30eed46 100644 --- a/echo-memory.plugin.src/skills/echo-memory/scripts/echo.py +++ b/echo-memory.plugin.src/skills/echo-memory/scripts/echo.py @@ -605,9 +605,136 @@ def cmd_scope(subcommand: str, text: str | None = None, as_json: bool = False) - return 0 -def cmd_load() -> int: - """Cold-start orientation: the canonical 6 reads in one call. 404s on today's - daily note and the inbox are normal (printed as absent, not errors).""" +def _fetch_statuses(paths: list[str]) -> dict[str, tuple[int, bytes]]: + """Concurrently GET vault paths keeping (status, body) per path — unlike read_many, + the 404-vs-offline distinction survives, which cmd_load's cache logic needs.""" + if not paths: + return {} + workers = min(MAX_WORKERS, len(paths)) + with ThreadPoolExecutor(max_workers=workers) as pool: + return dict(zip(paths, pool.map(lambda p: request("GET", vault_url(p)), paths))) + + +def _bullet_lines(text: str) -> list[str]: + return [ln for ln in text.splitlines() if ln.lstrip().startswith("- ")] + + +def _fm_line(text: str, field: str) -> str: + for ln in text.splitlines(): + if ln.startswith(f"{field}:"): + return ln.split(":", 1)[1].strip().strip('"').strip("'") + return "" + + +def _render_brief(texts: dict[str, str | None], listing_files: list[str], + offline: bool) -> str: + """The token-budgeted cold-start digest (2.1.0). Selection over compression: Fact / + Pattern rules in full, recent Observations only, scope + freshness (not the whole + history), the last session's key sections, today's Agent Log lines, inbox COUNT. + Sections are trimmed lowest-priority-first when over ECHO_LOAD_BUDGET.""" + budget = int(os.environ.get("ECHO_LOAD_BUDGET", "8000")) + S: list[tuple[str, str]] = [] # (section-name, text) in display order + + marker = texts.get("marker") + if marker is None: + S.append(("marker", "marker: _agent/echo-vault.md ABSENT — vault not bootstrapped " + "(run bootstrap.py before relying on memory)")) + else: + ver = _fm_line(marker, "schema_version") or "?" + S.append(("marker", f"marker: _agent/echo-vault.md — bootstrapped, schema_version {ver}")) + + ctx = texts.get("context") + if ctx: + scope = extract_heading(ctx, "Scope") or "(empty)" + upd = _fm_line(ctx, "scope_updated") or "unknown" + since = sum(1 for f in listing_files if f.endswith(".md") and f[:10] > upd) \ + if upd != "unknown" else None + head = f"scope (updated {upd}" + if since is not None: + head += f", {since} session(s) logged since" + S.append(("scope", head + "):\n" + scope)) + else: + S.append(("scope", "scope: current-context.md absent")) + + prefs = texts.get("preferences") + if prefs: + fact = extract_heading(prefs, "Fact / Pattern") + if fact: + S.append(("facts", "preferences — Fact / Pattern:\n" + fact)) + obs = _bullet_lines(extract_heading(prefs, "Observations")) + if obs: + shown = obs[-10:] + head = f"preferences — Observations (last {len(shown)} of {len(obs)}):" + S.append(("observations", head + "\n" + "\n".join(shown))) + + hb = texts.get("heartbeat") + if hb and hb.strip(): + pointer = hb.strip().splitlines()[0] + sess_path = pointer.split(" @ ", 1)[0].strip() + part = [f"last session: {pointer}"] + log_text = texts.get("_pointed_session") + if log_text: + for h in ("Goal", "Decisions Made", "Open Threads", "Suggested Next Step"): + sec = extract_heading(log_text, h) + if sec and sec.strip("() \n"): + part.append(f" {h}: " + " / ".join(ln.strip() for ln in sec.splitlines() + if ln.strip())[:400]) + elif sess_path: + part.append(" (session log not readable — see the path above)") + S.append(("session", "\n".join(part))) + elif listing_files: + recent = sorted((f for f in listing_files if f.endswith(".md")), reverse=True)[:3] + S.append(("session", "last session: heartbeat absent — recent logs:\n" + + "\n".join(f" _agent/sessions/{f}" for f in recent))) + + today_note = texts.get("today") + if today_note: + log_lines = _bullet_lines(extract_heading(today_note, "Agent Log")) + S.append(("agent-log", "today's Agent Log:\n" + ("\n".join(log_lines) or " (empty)"))) + else: + S.append(("agent-log", "today's Agent Log: (no daily note yet)")) + + inbox = texts.get("inbox") + if inbox: + items = re.findall(r"(?m)^\s*-\s*(\d{4}-\d{2}-\d{2})", inbox) + if items: + try: + oldest = (dt.date.fromisoformat(today()) - dt.date.fromisoformat(min(items))).days + S.append(("inbox", f"inbox: {len(items)} capture(s), oldest {oldest}d — " + "offer triage if any are older than ~7 days")) + except ValueError: + S.append(("inbox", f"inbox: {len(items)} capture(s)")) + else: + S.append(("inbox", "inbox: empty")) + else: + S.append(("inbox", "inbox: empty")) + + def total() -> int: + return sum(len(t) + 2 for _, t in S) + + # Over budget -> trim lowest-priority sections first, with explicit markers. + for name, keep in (("observations", 4), ("agent-log", 4), ("session", 9)): + if total() <= budget: + break + for i, (n, t) in enumerate(S): + if n == name: + lines = t.splitlines() + if len(lines) > keep: + S[i] = (n, "\n".join(lines[:keep]) + "\n (truncated — /echo-load for full)") + out = f"ECHO load (brief) — {today()}\n\n" + "\n\n".join(t for _, t in S) + if len(out) > budget: + out = out[:budget] + "\n(truncated — /echo-load for full)" + if offline: + out += ("\n\nNOTE: vault unreachable — context above is last-known-good cache; " + "writes this session will be queued and synced when the vault returns.") + return out + + +def cmd_load(brief: bool = False) -> int: + """Cold-start orientation: the canonical 6 reads in one call (fetched in parallel). + 404s on today's daily note and the inbox are normal (printed as absent, not errors). + `brief` renders the token-budgeted digest the SessionStart hook injects; the default + full mode (and /echo-load) prints the raw sections unchanged.""" if not echo_config.is_configured(_CFG): print("echo: NOT CONFIGURED — this machine has no usable ECHO key file yet.") print(f"Expected at: {echo_config.config_path()}") @@ -631,35 +758,77 @@ def cmd_load() -> int: synced = echo_queue.flush() if synced: print(f"(synced {synced} write(s) queued during a prior offline session)\n") + flagged = echo_queue.needs_attention() + if flagged: + print(f"NOTE: {len(flagged)} queued write(s) need attention (e.g. a capture " + "stopped at the duplicate gate on replay) — resolve with capture " + "--merge-into/--force; `echo.py flush` re-lists them.\n") except Exception as exc: # noqa: BLE001 — never let queue upkeep block a load print(f"(queue flush skipped: {exc})\n", file=sys.stderr) + # All orientation reads (+ the sessions listing) fetched in parallel over the + # warm connection pool — the listing feeds both the brief digest's scope-freshness + # count and the heartbeat-absent fallback, so it's no longer a second round-trip. + listing_path = "_agent/sessions/" + results = _fetch_statuses([p for _, p in targets] + [listing_path]) + marker_missing = False heartbeat_absent = False offline = False + texts: dict[str, str | None] = {} for label, path in targets: - status, body = request("GET", vault_url(path)) + status, body = results[path] if status == 200: echo_queue.cache_put(path, body) # refresh last-known-good for offline fallback + texts[label] = body.decode(errors="replace") + elif status == 0: # vault unreachable -> degrade to last-known-good cache + offline = True + cached = echo_queue.cache_get(path) + texts[label] = cached.decode(errors="replace") if cached is not None else None + else: + texts[label] = None + if status == 404: + if label == "marker": + marker_missing = True + if label == "heartbeat": + heartbeat_absent = True + + lst_status, lst_body = results[listing_path] + listing_files: list[str] = [] + if lst_status == 200: + try: + listing_files = list(json.loads(lst_body).get("files", [])) + except json.JSONDecodeError: + pass + + if brief: + hb = texts.get("heartbeat") + if hb and hb.strip(): # follow-up: the pointed session log's key sections + sess_path = hb.strip().splitlines()[0].split(" @ ", 1)[0].strip() + if sess_path: + st2, b2 = request("GET", vault_url(sess_path)) + if st2 == 200: + texts["_pointed_session"] = b2.decode(errors="replace") + print(_render_brief(texts, listing_files, offline)) + return 0 + + # ---- full mode: raw sections, output format unchanged ---- + for label, path in targets: + status, body = results[path] + if status == 200: print(f"===== {label}: {path} (HTTP 200) =====") - text = body.decode(errors="replace") + text = texts[label] or "" sys.stdout.write(text) if not text.endswith("\n"): sys.stdout.write("\n") elif status == 404: print(f"===== {label}: {path} (HTTP 404) =====") print("(absent — fine)") - if label == "marker": - marker_missing = True - if label == "heartbeat": - heartbeat_absent = True - elif status == 0: # vault unreachable -> degrade to last-known-good cache - offline = True - cached = echo_queue.cache_get(path) - if cached is not None: + elif status == 0: + if texts[label] is not None: print(f"===== {label}: {path} (OFFLINE — serving stale cache) =====") - sys.stdout.write(cached.decode(errors="replace")) - if not cached.endswith(b"\n"): + sys.stdout.write(texts[label]) + if not texts[label].endswith("\n"): sys.stdout.write("\n") else: print(f"===== {label}: {path} (OFFLINE — no cache) =====") @@ -669,20 +838,14 @@ def cmd_load() -> int: print(f"(error HTTP {status}: {body.decode(errors='replace')[:200]})") print() # M3: heartbeat pointer missing/stale -> fall back to the recent sessions listing - # (matches the documented loading procedure) so orientation works without the pointer. - if heartbeat_absent and not offline: - st, body = request("GET", vault_url("_agent/sessions/")) - if st == 200: - try: - files = sorted((f for f in json.loads(body).get("files", []) if f.endswith(".md")), - reverse=True) - except json.JSONDecodeError: - files = [] - if files: - print("===== recent sessions (heartbeat absent — fallback) =====") - for f in files[:5]: - print(f" _agent/sessions/{f}") - print() + # (already fetched in the batch) so orientation works without the pointer. + if heartbeat_absent and not offline and listing_files: + files = sorted((f for f in listing_files if f.endswith(".md")), reverse=True) + if files: + print("===== recent sessions (heartbeat absent — fallback) =====") + for f in files[:5]: + print(f" _agent/sessions/{f}") + print() if offline: print("NOTE: vault unreachable — context above is last-known-good cache; writes this " "session will be queued and synced when the vault returns.") @@ -732,7 +895,8 @@ def cmd_config(args) -> int: def build_parser() -> argparse.ArgumentParser: parser = argparse.ArgumentParser(description="Validated cross-platform ECHO vault client") sub = parser.add_subparsers(dest="cmd", required=True) - sub.add_parser("load") + p = sub.add_parser("load") + p.add_argument("--brief", action="store_true") # token-budgeted digest (hook default) sub.add_parser("flush") # H2: replay writes queued during a prior outage sub.add_parser("doctor") # M3: one-call readiness check sub.add_parser("get").add_argument("path") @@ -761,6 +925,9 @@ def build_parser() -> argparse.ArgumentParser: p = sub.add_parser("reflect") # H5: apply a JSON proposal set (dry-run unless --apply) p.add_argument("file", nargs="?") p.add_argument("--apply", action="store_true") + p = sub.add_parser("session-end") # 2.1.0: one-call session end (log+agent-log+reflect+scope+heartbeat) + p.add_argument("file", nargs="?") # bundle JSON; - or stdin + p.add_argument("--apply", action="store_true") p = sub.add_parser("triage") # one-tap inbox triage (reflect pipeline + audit log) p.add_argument("file", nargs="?") # proposals JSON; omit to list p.add_argument("--list", action="store_true", dest="list_inbox") @@ -796,21 +963,27 @@ def main(argv: list[str] | None = None) -> int: args = build_parser().parse_args(argv) try: if args.cmd == "load": - return cmd_load() + return cmd_load(brief=args.brief) if args.cmd == "doctor": import echo_doctor return echo_doctor.run() if args.cmd == "flush": import echo_queue n = echo_queue.flush() - remaining = len(echo_queue.pending()) + remaining = echo_queue.pending() + flagged = [r for r in remaining if r.get("needs_attention")] if n: print(f"ok: flushed {n} queued write(s)" - + (f"; {remaining} still queued (vault unreachable)" if remaining else "")) + + (f"; {len(remaining)} still queued" if remaining else "")) elif remaining: - print(f"ok: {remaining} write(s) still queued (vault unreachable)") + print(f"ok: {len(remaining)} write(s) still queued") else: print("ok: queue empty") + for r in flagged: + what = (f"capture '{(r.get('args') or {}).get('title')}'" if r.get("op") == "capture" + else f"{r.get('method')} {r.get('url')}") + print(f" needs attention ({r['needs_attention']}): {what} — resolve " + "(e.g. capture --merge-into / --force), it will not auto-land") return 0 if args.cmd == "config": return cmd_config(args) @@ -852,6 +1025,14 @@ def main(argv: list[str] | None = None) -> int: if not isinstance(proposals, list): raise EchoError("reflect: proposals must be a JSON array of objects", 2) return echo_reflect.apply(proposals, confirm=args.apply) + if args.cmd == "session-end": + import echo_session + raw = read_body(args.file).decode("utf-8", errors="replace") + try: + bundle = json.loads(raw) + except json.JSONDecodeError as exc: + raise EchoError(f"session-end: bundle must be a JSON object ({exc})", 2) + return echo_session.session_end(bundle, apply=args.apply) if args.cmd == "triage": import echo_triage if args.list_inbox or not args.file: @@ -880,9 +1061,12 @@ def main(argv: list[str] | None = None) -> int: date=args.date, domain=args.domain, inbox=args.inbox, no_log=args.no_log, as_json=args.json, dry_run=args.dry_run, force=args.force, merge_into=args.merge_into) - except EchoError as exc: + except RuntimeError as exc: + # Catches EchoError from THIS module and from helper modules' own `import echo` + # twin (when echo.py runs as __main__ the two class objects differ, but both + # subclass RuntimeError and carry .code). print(f"echo.py: {exc}", file=sys.stderr) - return exc.code + return getattr(exc, "code", 1) return 2 diff --git a/echo-memory.plugin.src/skills/echo-memory/scripts/echo_hook_session_start.py b/echo-memory.plugin.src/skills/echo-memory/scripts/echo_hook_session_start.py index 688cf68..642c117 100644 --- a/echo-memory.plugin.src/skills/echo-memory/scripts/echo_hook_session_start.py +++ b/echo-memory.plugin.src/skills/echo-memory/scripts/echo_hook_session_start.py @@ -42,7 +42,7 @@ def main() -> int: try: import echo with contextlib.redirect_stdout(buf): - rc = echo.cmd_load() + rc = echo.cmd_load(brief=True) # token-budgeted digest; /echo-load stays full text = buf.getvalue() if rc == 78: context = ("[echo-memory] ECHO is NOT CONFIGURED on this machine. Before " diff --git a/echo-memory.plugin.src/skills/echo-memory/scripts/echo_hook_stop.py b/echo-memory.plugin.src/skills/echo-memory/scripts/echo_hook_stop.py index dc20a35..1cbc9f8 100644 --- a/echo-memory.plugin.src/skills/echo-memory/scripts/echo_hook_stop.py +++ b/echo-memory.plugin.src/skills/echo-memory/scripts/echo_hook_stop.py @@ -30,9 +30,11 @@ MIN_TURNS = int(os.environ.get("ECHO_REFLECT_MIN_TURNS", "5")) REASON = ( "[echo-memory] Session-end reflection has not run yet. If this session produced " "durable facts, decisions, commitments, or artifacts worth remembering: run the " - "/echo-reflect flow (extract -> preview -> apply only with the operator's confirm) " - "and write the session log + heartbeat per the echo-memory skill. If nothing " - "durable emerged, simply finish — do not invent memories. " + "/echo-reflect flow (extract -> preview -> apply only with the operator's confirm), " + "then finish with ONE call — `echo.py session-end --apply` — which " + "writes the session log, Agent Log line, reflect proposals, optional scope switch, " + "and the heartbeat (last, as the commit marker) together. If nothing durable " + "emerged, simply finish — do not invent memories. " "(This reminder fires once per session.)" ) @@ -64,6 +66,7 @@ def already_reflected(raw: str) -> bool: """Transcript evidence that reflection or session logging already happened.""" return ("/echo-reflect" in raw or re.search(r'echo\.py"?\s+reflect\b', raw) is not None + or re.search(r'echo\.py"?\s+session-end\b', raw) is not None or "put _agent/sessions/" in raw or "put _agent/heartbeat/last-session.md" in raw) diff --git a/echo-memory.plugin.src/skills/echo-memory/scripts/echo_ops.py b/echo-memory.plugin.src/skills/echo-memory/scripts/echo_ops.py index 72b2e3c..b6abf6d 100644 --- a/echo-memory.plugin.src/skills/echo-memory/scripts/echo_ops.py +++ b/echo-memory.plugin.src/skills/echo-memory/scripts/echo_ops.py @@ -83,8 +83,13 @@ def recall(query, limit: int = 8, as_json: bool = False) -> int: # ----------------------------------------------------------- agent log ------- def ensure_daily_log(line: str) -> None: """Resilient: ensure today's daily note + its `## Agent Log` heading exist, then - idempotently append the line. Best-effort — never raises into the caller.""" + idempotently append the line. Best-effort — never raises into the caller. Writes go + through the offline queue (2.1.0), so a mid-session outage queues the line instead + of silently dropping it. Offline edge accepted: if the daily note still doesn't + exist at replay time, the queued PATCH 400-drops with a loud warning — the agent-log + line is a convenience trace, not the memory itself.""" try: + import echo_queue date = echo.today() path = f"journal/daily/{date}.md" status, body = echo.request("GET", echo.vault_url(path)) @@ -92,21 +97,22 @@ def ensure_daily_log(line: str) -> None: ts, tb = echo.request("GET", echo.vault_url("journal/templates/daily-note-template.md")) tmpl = tb.decode(errors="replace") if ts == 200 else f"# {date}\n\n## Agent Log\n\n## Related\n" tmpl = tmpl.replace("{{date:YYYY-MM-DD}}", date).replace("{{DATE}}", date) - echo.request("PUT", echo.vault_url(path), data=tmpl.encode(), - headers={"Content-Type": "text/markdown"}) + echo_queue.safe_request("PUT", echo.vault_url(path), data=tmpl.encode(), + headers={"Content-Type": "text/markdown"}) text = tmpl else: - text = body.decode(errors="replace") - if not re.search(r"(?m)^## Agent Log\s*$", text): - echo.request("POST", echo.vault_url(path), data=b"\n\n## Agent Log\n", - headers={"Content-Type": "text/markdown"}) - status, body = echo.request("GET", echo.vault_url(path)) - if status == 200 and line in body.decode(errors="replace"): + text = body.decode(errors="replace") if status == 200 else "" + if status == 200 and not re.search(r"(?m)^## Agent Log\s*$", text): + echo_queue.safe_request("POST", echo.vault_url(path), data=b"\n\n## Agent Log\n", + headers={"Content-Type": "text/markdown"}) + st2, body2 = echo.request("GET", echo.vault_url(path)) + if st2 == 200 and line in body2.decode(errors="replace"): return - echo.request("PATCH", echo.vault_url(path), - data=echo.normalize_patch_body((line + "\n").encode(), "append", "heading"), - headers={"Operation": "append", "Target-Type": "heading", - "Target": f"{date}::Agent Log", "Content-Type": "text/markdown"}) + echo_queue.safe_request( + "PATCH", echo.vault_url(path), + data=echo.normalize_patch_body((line + "\n").encode(), "append", "heading"), + headers={"Operation": "append", "Target-Type": "heading", + "Target": f"{date}::Agent Log", "Content-Type": "text/markdown"}) except Exception as exc: print(f"echo_ops: agent-log skipped ({exc})", file=sys.stderr) @@ -148,20 +154,25 @@ def _dated_block(today_s: str, body_text: str) -> tuple[str, str]: def _append_to_existing(path: str, kind: str, today_s: str, body_text: str) -> None: + # Writes route through the offline queue (2.1.0) so a mid-capture outage queues the + # update instead of silently losing it. (A fully-offline capture never reaches here — + # the short-circuit in capture() queues the whole op as one semantic record.) + import echo_queue text = links.get_text(path) or "" h1 = links.first_h1(text) or path.rsplit("/", 1)[-1][:-3] heading = LOG_HEADING.get(kind, "Notes") if not re.search(rf"(?m)^##\s+{re.escape(heading)}\s*$", text): - echo.request("POST", echo.vault_url(path), data=f"\n\n## {heading}\n".encode(), - headers={"Content-Type": "text/markdown"}) + echo_queue.safe_request("POST", echo.vault_url(path), data=f"\n\n## {heading}\n".encode(), + headers={"Content-Type": "text/markdown"}) bullet, block = _dated_block(today_s, body_text) status, body = echo.request("GET", echo.vault_url(path)) if status == 200 and bullet in body.decode(errors="replace"): return - echo.request("PATCH", echo.vault_url(path), - data=echo.normalize_patch_body((block + "\n").encode(), "append", "heading"), - headers={"Operation": "append", "Target-Type": "heading", - "Target": f"{h1}::{heading}", "Content-Type": "text/markdown"}) + echo_queue.safe_request( + "PATCH", echo.vault_url(path), + data=echo.normalize_patch_body((block + "\n").encode(), "append", "heading"), + headers={"Operation": "append", "Target-Type": "heading", + "Target": f"{h1}::{heading}", "Content-Type": "text/markdown"}) echo.cmd_fm(path, "updated", json.dumps(today_s)) @@ -213,7 +224,35 @@ def capture(kind: str | None, title: str, file_arg: str | None, status_v: str = return done("inbox", "inbox/captures/inbox.md", ok=rc == 0) slug = idx_mod.slugify(title) - index = idx_mod.load() + # Offline short-circuit (2.1.0). The routed path needs the entity index and several + # round-trips; when the vault is unreachable, queue the WHOLE capture as one semantic + # record — flush replays it back through this function against the index as it is at + # replay time, so routing, the duplicate gate, and aliasing re-run with fresh state. + # (Byte-level request replay would freeze a create-vs-update decision made blind.) + try: + index = idx_mod.load() + except echo.EchoError as exc: + if "unreachable" not in str(exc): + raise + if dry_run: + print("offline: vault unreachable — a real run would queue this capture " + "for replay on the next reachable session.", file=real_stdout) + return 0 + import echo_queue + echo_queue.enqueue_capture({ + "kind": kind, "title": title, "body_text": body_text, "status_v": status_v, + "aliases": aliases, "sources": sources, "tags": tags, "date": date, + "domain": domain, "inbox": False, "no_log": no_log, + "force": force, "merge_into": merge_into, "today": today_s, + }) + if as_json: + env = echo_output.envelope("queued:capture", + {"kind": kind, "title": title, "queued": True}) + print(json.dumps(env, ensure_ascii=False), file=real_stdout) + else: + print(f"queued (offline): capture {kind} '{title}' — will replay through " + "capture on the next reachable session", file=real_stdout) + return 0 match_slug, existing = idx_mod.resolve(index, title) # --merge-into: the operator has already identified the canonical entity (e.g. after # a duplicate-gate stop) — route this capture as an UPDATE to it, whatever the title. diff --git a/echo-memory.plugin.src/skills/echo-memory/scripts/echo_queue.py b/echo-memory.plugin.src/skills/echo-memory/scripts/echo_queue.py index e57c6ff..9b01f3e 100644 --- a/echo-memory.plugin.src/skills/echo-memory/scripts/echo_queue.py +++ b/echo-memory.plugin.src/skills/echo-memory/scripts/echo_queue.py @@ -81,6 +81,23 @@ def enqueue(method: str, url: str, data: bytes | None, headers: dict, idem_key: fh.write(json.dumps(rec, ensure_ascii=False) + "\n") +def enqueue_capture(args: dict) -> None: + """Queue a whole capture as ONE semantic record (2.1.0). Byte-level replay is wrong + for the routed capture path: the create-vs-update decision, duplicate gate, and + aliasing must re-run against the index AS IT IS AT REPLAY TIME, so flush re-invokes + echo_ops.capture with the recorded arguments instead of replaying stale requests. + `args` carries body_text inline (never a temp-file path) plus `today` (the capture- + time ECHO_TODAY) so replayed frontmatter/log dates reflect when it was said.""" + key = "capture:" + hashlib.sha1( + json.dumps(args, sort_keys=True, ensure_ascii=False).encode("utf-8")).hexdigest() + if any(rec.get("idem_key") == key for rec in pending()): + return + _ensure_dirs() + rec = {"op": "capture", "args": args, "idem_key": key} + with outbox_path().open("a", encoding="utf-8") as fh: + fh.write(json.dumps(rec, ensure_ascii=False) + "\n") + + def pending() -> list[dict]: p = outbox_path() if not p.exists(): @@ -88,6 +105,12 @@ def pending() -> list[dict]: return [json.loads(ln) for ln in p.read_text(encoding="utf-8").splitlines() if ln.strip()] +def needs_attention() -> list[dict]: + """Queued records that replay could not land automatically (e.g. a duplicate-gate + stop) — surfaced by `flush` and `load` so the operator resolves them deliberately.""" + return [rec for rec in pending() if rec.get("needs_attention")] + + def _rewrite(records: list[dict]) -> None: p = outbox_path() if not records: @@ -108,9 +131,57 @@ def _rebase(url: str) -> str: return echo.BASE + parts.path + (("?" + parts.query) if parts.query else "") +def _replay_capture(rec: dict) -> str: + """Replay a queued semantic capture through echo_ops.capture against the CURRENT + index. Returns 'ok', 'retry' (still offline), or 'gated' (duplicate gate / needs + the operator — the record is kept and flagged, later records still replay).""" + # Cheap reachability probe FIRST: if the vault is still down, echo_ops.capture's + # own offline short-circuit would try to re-enqueue this very record — probe and + # bail as 'retry' instead of recursing into the queue. + st, _ = echo.request("GET", echo.vault_url("_agent/echo-vault.md")) + if st == 0: + return "retry" + import echo_ops + a = dict(rec.get("args") or {}) + body_text = a.pop("body_text", "") or "" + day = a.pop("today", None) + bodyfile = echo.temp_file(body_text.encode("utf-8")) if body_text else None + prev = os.environ.get("ECHO_TODAY") + try: + if day: + os.environ["ECHO_TODAY"] = day # replayed dates reflect when it was said + rc = echo_ops.capture( + a.get("kind"), a.get("title") or "(untitled)", bodyfile, + status_v=a.get("status_v", ""), aliases=a.get("aliases") or [], + sources=a.get("sources") or [], tags=a.get("tags") or [], + date=a.get("date"), domain=a.get("domain", "business"), + inbox=bool(a.get("inbox")), no_log=bool(a.get("no_log")), + force=bool(a.get("force")), merge_into=a.get("merge_into")) + except Exception as exc: # noqa: BLE001 — a broken record must not wedge the queue + print(f"echo_queue: queued capture '{a.get('title')}' failed on replay ({exc}) " + "— kept and flagged", file=sys.stderr) + return "gated" + finally: + if day: + if prev is None: + os.environ.pop("ECHO_TODAY", None) + else: + os.environ["ECHO_TODAY"] = prev + if rc == 0: + return "ok" + if rc == 76: + print(f"echo_queue: queued capture '{a.get('title')}' stopped at the duplicate " + "gate on replay — resolve with capture --merge-into or --force, " + "then remove it from the outbox", file=sys.stderr) + return "gated" + + def _replay(rec: dict) -> str: """Re-issue one queued write idempotently. Returns 'ok' (landed/already-present), - 'drop' (permanent 4xx — can never succeed), or 'retry' (still offline / 5xx).""" + 'drop' (permanent 4xx — can never succeed), 'gated' (kept + flagged for the + operator), or 'retry' (still offline / 5xx).""" + if rec.get("op") == "capture": + return _replay_capture(rec) method = rec["method"] url = _rebase(rec["url"]) data = base64.b64decode(rec["body_b64"]) if rec.get("body_b64") else None @@ -137,21 +208,30 @@ def _replay(rec: dict) -> str: def flush() -> int: """Replay queued writes in order; stop at the first still-unreachable one (keeping it - and the rest). Returns the number actually replayed. Safe to call when empty.""" + and everything after). A 'gated' record (duplicate gate / replay error) is kept and + flagged but does NOT block later records — they are independent writes. Returns the + number actually replayed. Safe to call when empty.""" items = pending() if not items: return 0 replayed = 0 remaining: list[dict] = [] - for i, rec in enumerate(items): + stopped = False + for rec in items: + if stopped: + remaining.append(rec) + continue result = _replay(rec) if result == "ok": replayed += 1 elif result == "drop": continue # warned in _replay; remove from queue + elif result == "gated": + rec["needs_attention"] = rec.get("needs_attention") or "replay-gated" + remaining.append(rec) else: # retry -> still offline; keep this and everything after, in order - remaining = items[i:] - break + remaining.append(rec) + stopped = True _rewrite(remaining) return replayed diff --git a/echo-memory.plugin.src/skills/echo-memory/scripts/echo_session.py b/echo-memory.plugin.src/skills/echo-memory/scripts/echo_session.py new file mode 100644 index 0000000..3bc614b --- /dev/null +++ b/echo-memory.plugin.src/skills/echo-memory/scripts/echo_session.py @@ -0,0 +1,144 @@ +#!/usr/bin/env python3 +"""echo_session.py — the session-end bundle: one call ends a session correctly. [2.1.0] + +Ending a substantive session used to take four-plus separate invocations (session-log +PUT, Agent Log append, heartbeat PUT, optional scope set, optional reflect apply) — +a cost per session, and a partial failure left broken orientation state (log written +but heartbeat stale). This verb does all of it in one call, under ONE advisory-lock +acquisition, in a fixed order with the heartbeat written LAST as the commit marker: +a failure part-way leaves the previous pointer intact, never a dangling one. + +Bundle shape (JSON object): + { + "slug": "echo-mcp-spec", # required, kebab-case + "log_body": "", # required + "agent_log_line": "- 2026-07-28: ...", # optional; derived when omitted + "scope": "new scope text", # optional -> scope set + "reflect": [ { ...PROPOSAL_SCHEMA... } ] # optional -> routed via capture + } + +Dry-run by default (previews the whole plan, reflect included); --apply writes. +Times: the filename HHMM comes from $ECHO_NOW (HHMM) else the local clock; the date +from ECHO_TODAY via echo.today(). Offline: every step rides the offline queue, so a +session end during an outage queues durably instead of failing. + +CLI: echo.py session-end [--apply] +""" +from __future__ import annotations + +import datetime as dt +import json +import os +import re +import sys +from pathlib import Path + +sys.path.insert(0, str(Path(__file__).resolve().parent)) +import echo # noqa: E402 + +SESSION_RE = re.compile(r"^\d{4}-\d{2}-\d{2}-\d{4}-[a-z0-9-]+\.md$") +SLUG_RE = re.compile(r"^[a-z0-9][a-z0-9-]*$") + + +def _hhmm() -> str: + v = os.environ.get("ECHO_NOW", "").strip() + if v: + if not re.match(r"^\d{4}$", v): + raise echo.EchoError(f"session-end: $ECHO_NOW must be HHMM, got '{v}'", 2) + return v + return dt.datetime.now().strftime("%H%M") + + +def validate(bundle: dict) -> tuple[str, str]: + """Return (session_path, agent_log_line) or raise EchoError(2). Validation runs + BEFORE any write — a bad bundle writes nothing at all.""" + if not isinstance(bundle, dict): + raise echo.EchoError("session-end: bundle must be a JSON object", 2) + slug = str(bundle.get("slug", "")).strip() + if not SLUG_RE.match(slug): + raise echo.EchoError(f"session-end: slug must be kebab-case, got '{slug}'", 2) + if not str(bundle.get("log_body", "")).strip(): + raise echo.EchoError("session-end: log_body is required (the full session-log markdown)", 2) + reflect = bundle.get("reflect") + if reflect is not None and not isinstance(reflect, list): + raise echo.EchoError("session-end: reflect must be a JSON array of proposals", 2) + fname = f"{echo.today()}-{_hhmm()}-{slug}.md" + if not SESSION_RE.match(fname): + raise echo.EchoError(f"session-end: derived filename '{fname}' violates the " + "canonical YYYY-MM-DD-HHMM-.md form", 2) + path = f"_agent/sessions/{fname}" + line = str(bundle.get("agent_log_line") or "").strip() \ + or f"- {echo.today()}: session logged [[{path[:-3]}]]" + return path, line + + +def session_end(bundle: dict, apply: bool = False) -> int: + import echo_reflect + path, line = validate(bundle) + scope = str(bundle.get("scope") or "").strip() + proposals = bundle.get("reflect") or [] + + valid, errors = echo_reflect.validate(proposals) + for e in errors: + print(f"skip: {e}", file=sys.stderr) + + print(f"session-end plan ({'APPLY' if apply else 'dry-run'}):") + print(f" 1. session log -> {path}") + print(f" 2. agent-log line: {line}") + if valid: + echo_reflect.classify(valid) + print(f" 3. reflect: {len(valid)} proposal(s)") + print(echo_reflect.preview(valid)) + else: + print(" 3. reflect: (none)") + print(f" 4. scope set: {scope!r}" if scope else " 4. scope: (unchanged)") + print(f" 5. heartbeat -> {path} @ (written LAST — the commit marker)") + + if not apply: + print("\nsession-end: dry-run — re-run with --apply to write.") + return 0 + + import echo_concurrency + import echo_ops + steps: dict[str, str] = {} + with echo_concurrency.vault_lock(): + # 1. The session log itself. A hard failure here aborts the whole bundle — + # nothing after it (including the heartbeat) runs, so orientation state + # can never point at a log that was never written. + echo.cmd_put(path, echo.temp_file(str(bundle["log_body"]).encode("utf-8"))) + steps["session_log"] = "ok" + # 2. Agent-log line (best-effort by contract; queues itself when offline). + echo_ops.ensure_daily_log(line) + steps["agent_log"] = "ok" + # 3. Reflect proposals through the normal capture path (gate-aware). + if valid: + gated = 0 + for p in valid: + if p.get("_action") == "error": + continue + bodyfile = echo.temp_file((p.get("body") or "").encode()) if p.get("body") else None + rc = echo_ops.capture( + p.get("kind"), p["title"], bodyfile, + aliases=p.get("aliases") or [], sources=p.get("sources") or [], + tags=p.get("tags") or [], date=p.get("date"), + domain=p.get("domain", "business"), + inbox=bool(p.get("inbox")) or not p.get("kind")) + gated += 1 if rc == 76 else 0 + steps["reflect"] = f"{len(valid) - gated}/{len(valid)} applied" \ + + (f", {gated} gated (re-propose with --merge-into)" if gated else "") + else: + steps["reflect"] = "skipped (none)" + # 4. Scope switch (optional). + if scope: + echo.cmd_scope("set", scope) + steps["scope"] = "ok" + else: + steps["scope"] = "unchanged" + # 5. Heartbeat LAST — the commit marker for the whole bundle. + echo.cmd_put("_agent/heartbeat/last-session.md", + echo.temp_file(f"{path} @ {echo.now_iso()}\n".encode("utf-8"))) + steps["heartbeat"] = "ok" + + print("\nsession-end: done — " + + "; ".join(f"{k}: {v}" for k, v in steps.items())) + return 0 diff --git a/echo-memory.plugin.src/skills/echo-reflect/SKILL.md b/echo-memory.plugin.src/skills/echo-reflect/SKILL.md index 628d649..9500857 100644 --- a/echo-memory.plugin.src/skills/echo-reflect/SKILL.md +++ b/echo-memory.plugin.src/skills/echo-reflect/SKILL.md @@ -30,4 +30,9 @@ python3 "$ECHO" reflect proposals.json --apply # write — routes each via ca `reflect` dedups every proposal against the entity index (so it shows create-vs-update), skips below-confidence items, and routes `inbox` proposals to the capture inbox. Each applied item goes through the normal `capture` path — canonical frontmatter, auto-linking, -and the recall index — and the writes are lock-guarded. $ARGUMENTS +and the recall index — and the writes are lock-guarded. + +**Finishing the session?** Prefer the one-call bundle — `python3 "$ECHO" session-end +bundle.json --apply` — which applies the reflect proposals AND writes the session log, +Agent Log line, optional scope switch, and the heartbeat (last, as the commit marker) +together. See the echo-memory skill's **Session Logging** section for the bundle shape. $ARGUMENTS diff --git a/eval/test_features.py b/eval/test_features.py index 3046a5c..8c40610 100644 --- a/eval/test_features.py +++ b/eval/test_features.py @@ -345,6 +345,63 @@ def main(): check("v1.5 session-start hook stays quiet on resume", r.returncode == 0 and not r.stdout.strip(), r.stdout) + # 18. (2.1.0) brief load: digest form, budget respected, key facts present. + h.seed("_agent/memory/semantic/operator-preferences.md", + "---\ntype: semantic-memory\ncreated: 2026-06-01\n---\n# Operator Preferences\n\n" + "## Fact / Pattern\n- prefers concise output\n\n## Observations\n" + + "\n".join(f"- 2026-06-{i:02d}: observation {i}" for i in range(1, 15)) + "\n") + h.seed("inbox/captures/inbox.md", "- 2026-06-01: old capture\n- 2026-06-20: new capture\n") + r = h.echo(ECHO, "load", "--brief") + check("v2.1 brief load renders the digest header", + "ECHO load (brief)" in r.stdout, r.stdout[:200]) + check("v2.1 brief load keeps Fact / Pattern", + "prefers concise output" in r.stdout, r.stdout[:800]) + check("v2.1 brief load trims observations to the last 10", + "last 10 of 14" in r.stdout and "observation 14" in r.stdout + and "observation 2" not in r.stdout.replace("observation 2026", ""), r.stdout[-1200:]) + check("v2.1 brief load reports inbox as a count", + "inbox: 2 capture(s), oldest 20d" in r.stdout, r.stdout[-400:]) + check("v2.1 brief load includes the marker line", + "echo-vault.md" in r.stdout and "marker" in r.stdout, r.stdout[:300]) + + # 19. (2.1.0) session-end: dry-run writes nothing; apply lands all steps with + # the heartbeat LAST; a bad bundle aborts before any write. + se_env = dict(os.environ, ECHO_BASE=base, ECHO_KEY=KEY, ECHO_VERIFY="1", + ECHO_TODAY="2026-06-21", ECHO_NOW="2359") + bundle = {"slug": "quickwins-ship", + "log_body": "---\ntype: session-log\nstatus: complete\ncreated: 2026-06-21\n---\n" + "# Session Log\n\n## Goal\nShip the quick wins.\n", + "agent_log_line": "- 2026-06-21: shipped the quick wins", + "reflect": [{"title": "Quark Preference", "kind": "semantic", + "body": "The operator prefers quark.", "confidence": 0.9}]} + bfile = HERE / "_session_bundle.json" + bfile.write_text(json.dumps(bundle), encoding="utf-8") + sess_path = "_agent/sessions/2026-06-21-2359-quickwins-ship.md" + r = subprocess.run([sys.executable, str(ECHO), "session-end", str(bfile)], + capture_output=True, text=True, env=se_env) + check("v2.1 session-end dry-run previews the plan", + "dry-run" in r.stdout and sess_path in r.stdout, r.stdout) + check("v2.1 session-end dry-run writes nothing", h.ground(sess_path) is None) + r = subprocess.run([sys.executable, str(ECHO), "session-end", str(bfile), "--apply"], + capture_output=True, text=True, env=se_env) + check("v2.1 session-end apply lands the session log", + "Ship the quick wins" in (h.ground(sess_path) or ""), r.stdout + r.stderr) + check("v2.1 session-end apply routes the reflect proposal", + h.ground("_agent/memory/semantic/quark-preference.md") is not None) + check("v2.1 session-end apply writes the heartbeat last (commit marker)", + sess_path in (h.ground("_agent/heartbeat/last-session.md") or "")) + daily = h.ground("journal/daily/2026-06-21.md") or "" + check("v2.1 session-end apply appends the agent-log line", + "shipped the quick wins" in daily, daily[-300:]) + bad = dict(bundle, slug="Bad_Slug!") + bfile.write_text(json.dumps(bad), encoding="utf-8") + r = subprocess.run([sys.executable, str(ECHO), "session-end", str(bfile), "--apply"], + capture_output=True, text=True, env=se_env) + bfile.unlink() + check("v2.1 session-end rejects a bad slug before any write", + r.returncode == 2 and sess_path in (h.ground("_agent/heartbeat/last-session.md") or ""), + r.stderr) + print(f"\n{len(failures)} failure(s)" if failures else "\nall feature tests passed") return 1 if failures else 0 finally: diff --git a/eval/test_offline_queue.py b/eval/test_offline_queue.py index f9a864f..005fffc 100644 --- a/eval/test_offline_queue.py +++ b/eval/test_offline_queue.py @@ -115,6 +115,40 @@ def main(): check("offline load flags OFFLINE", "OFFLINE" in r.stdout, r.stdout) check("offline load serves cached marker", "EVALMARK-marker" in r.stdout, r.stdout) + # --- Phase 5 (2.1.0): vault DOWN -> capture queues a SEMANTIC record ------ + r = echo(DEAD, "capture", "Offline Person", "--kind", "person") + check("offline capture exits 0 (queued)", r.returncode == 0, r.stdout + r.stderr) + check("offline capture reports queued", "queued (offline): capture" in r.stdout, r.stdout) + outbox = Path(state) / "outbox.ndjson" + check("capture queued as an op record", + outbox.exists() and '"op": "capture"' in outbox.read_text(encoding="utf-8")) + # same capture again while offline -> deduped by idem_key, not double-queued + echo(DEAD, "capture", "Offline Person", "--kind", "person") + recs = [ln for ln in outbox.read_text(encoding="utf-8").splitlines() if '"op": "capture"' in ln] + check("offline capture is idempotent in the queue", len(recs) == 1, str(len(recs))) + + # --- Phase 6 (2.1.0): vault UP -> flush replays capture THROUGH capture --- + srv = start_mock() + r = echo(mock_base, "flush") + check("flush replays the queued capture", "flushed" in r.stdout, r.stdout + r.stderr) + note = ground("resources/people/offline-person.md") + check("replayed capture created the routed note", + note is not None and "type: person" in note, str(note)[:200]) + check("replayed capture used the capture-time date", + note is not None and "created: 2026-06-22" in note, str(note)[:200]) + + # --- Phase 7 (2.1.0): duplicate gate ON REPLAY keeps + flags the record --- + echo(mock_base, "capture", "Zebulon Quargle", "--kind", "person") # the existing entity + r = echo(DEAD, "capture", "Zebulon Quargle Junior", "--kind", "person") + check("lookalike capture queues offline", "queued (offline)" in r.stdout, r.stdout) + r = echo(mock_base, "flush") + check("gated replay is flagged, not landed", + "needs attention" in r.stdout and "replay-gated" in r.stdout, r.stdout + r.stderr) + check("gated replay did NOT create the duplicate", + ground("resources/people/zebulon-quargle-junior.md") is None) + check("gated record survives in the outbox", + '"needs_attention"' in (outbox.read_text(encoding="utf-8") if outbox.exists() else "")) + print(f"\n{len(failures)} failure(s)" if failures else "\nall offline-queue tests passed") return 1 if failures else 0 finally: