From ccd3ce672285a462bf55a468a57589fed9ffd71d Mon Sep 17 00:00:00 2001 From: Gregory Buchenberger Date: Thu, 20 Aug 2026 09:00:17 -0600 Subject: [PATCH 1/2] feat: harden SQLite persistence Apply per-connection SQLite pragmas (busy_timeout, foreign_keys), replace fuzzy FTS memory dedup with exact SHA-256 digests, and make message persistence idempotent via an upsert and deterministic IDs. --- .agents/plans/consolidate-persistence/PLAN.md | 561 ++++++++++++++++++ .agents/plans/persistence-hardening/PLAN.md | 254 ++++++++ internal/agent/persist.go | 17 +- internal/memory/debounce_test.go | 4 + internal/memory/memory.go | 90 ++- internal/memory/memory_test.go | 159 +++++ internal/memory/message_repo.go | 15 +- 7 files changed, 1077 insertions(+), 23 deletions(-) create mode 100644 .agents/plans/consolidate-persistence/PLAN.md create mode 100644 .agents/plans/persistence-hardening/PLAN.md diff --git a/.agents/plans/consolidate-persistence/PLAN.md b/.agents/plans/consolidate-persistence/PLAN.md new file mode 100644 index 00000000..421fcd24 --- /dev/null +++ b/.agents/plans/consolidate-persistence/PLAN.md @@ -0,0 +1,561 @@ +--- +name: consolidate-persistence +description: Consolidate overlapping persistence systems - session storage, memories, OTel traces, and Shepherd traces +status: draft +--- + +# Plan: Consolidate yaah Persistence and Observability Systems + +## Goal + +Reduce redundancy across yaah's four persistence/observability systems by consolidating overlapping data and unifying storage backends where practical. The primary targets are: + +1. **Session storage** (`internal/memory/` SQLite `state.db`) - conversation history +2. **Shepherd traces** (`shepherd-kernel-go` SQLite `trace.sqlite`) - tool execution traces +3. **OTel traces** (`internal/observability/` OTLP backend) - distributed tracing +4. **Memories** (`internal/memory/` SQLite `state.db`) - user knowledge base + +## Problem Statement + +Currently yaah maintains **two separate SQLite databases** (`~/.yaah/state.db` and `~/.yaah/traces/trace.sqlite`) that store overlapping data about tool executions, token usage, and session metadata. This creates: + +- **Storage fragmentation**: Two files to backup/migrate/manage +- **Query fragmentation**: Different CLIs for similar data (`yaah session` vs `yaah shepherd-trace`) +- **Write duplication**: Token counts and tool call details written to multiple places +- **Conceptual overlap**: Session messages and Shepherd facts both record tool execution history +- **No cross-links**: OTel spans carry no `session_id`, and `messages` carry no `trace_id`, so the four systems cannot be joined to correlate a single turn +- **Hidden duplication**: verbose OTel tracing re-records full message content, reasoning, and tool-call args already persisted in `messages` + +## Non-Goals + +- Do NOT consolidate **Memories** - user knowledge is conceptually separate from execution data. (But the FTS + embedding + vector-search *plumbing* shared by `memory` and `messages` can be unified without merging the data — see Finding F7.) +- Do NOT eliminate **OTel traces** - distributed tracing has different query semantics (span trees vs fact DAGs) +- Do NOT change or fork the **shepherd-kernel-go** library - upstream dependency. This constrains Phase 1 — see Finding F1. + +--- + +## Current Architecture Overview + +### System 1: Session Storage + Memories (SQLite `state.db`) + +**Location**: `internal/memory/` + +**Tables**: +- `sessions` - session metadata (id, started_at, ended_at, cwd, model, tokens_in, tokens_out, system_prompt, compacted_summary, compaction_cooldown_until, ineffective_compactions) +- `messages` - conversation messages (session_id, idx, role, content, reasoning_content, tool_name, tool_call_id, tool_calls, ts, id, embedding) +- `memory` - long-term memory notes (id, text, tags, source, created_at, accessed_at, access_count, embedding) +- `todos` - todo persistence +- FTS5 virtual tables: `messages_fts`, `memory_fts` + +**CLI**: `yaah session`, `yaah memory` + +### System 2: Shepherd Traces (SQLite `trace.sqlite`) + +**Location**: `internal/agent/pipeline/trace.go` + `shepherd-kernel-go` + +**Storage**: Separate SQLite file at `~/.yaah/traces/trace.sqlite` + +**Data Model**: +- Facts with Declaration/Capture modes +- Causal parent-child links (Fact DAG) +- Frontiers for stable read points +- Schema: `yaah.tool.{name}.v1` for declarations, `yaah.tool.{name}.v1.applied` for captures +- Turn lifecycle: `yaah.execution.created.v1`, `yaah.execution.started.v1`, `yaah.execution.completed.v1`, `yaah.execution.failed.v1` + +**CLI**: `yaah shepherd-trace list/show/profile` + +### System 3: OTel Traces (External or In-Memory) + +**Location**: `internal/observability/` + +**Storage**: +- External OTLP backend (Jaeger, SigNoz, etc.) when `endpoint` is configured +- In-memory `BufferingSpanProcessor` in serve mode when no endpoint + +**Data Model**: OpenTelemetry span tree with attributes for tokens, tool names, etc. + +**Control**: Gated by `OtelEnabled` flag in loop config + +### System 4: Memories (Already in `state.db`) + +**No overlap** - memories are user knowledge, not execution data. Keep separate. + +--- + +## Overlap Analysis + +### Session Storage vs Shepherd Traces + +| Data | Session Storage | Shepherd Traces | Overlap | +|------|---------------|----------------|---------| +| Session ID | `sessions.id` PK | `TraceOwnerID` | ✅ Same concept | +| Tool calls | `messages.tool_name`, `messages.tool_calls` | Declaration facts with `yaah.tool.{name}.v1` | ✅ Both record tool invocations | +| Tool args | `messages.tool_calls` JSON | Declaration fact payload `args` | ✅ Same data | +| Tool results | Subsequent messages | Capture fact payload `success`, `error`, `duration` | ✅ Same data | +| Token counts | `sessions.tokens_in/out` | Turn capture facts `prompt_tokens`, `completion_tokens` | ✅ Same data | +| Timestamps | `sessions.started_at`, `messages.ts` | Fact envelope timestamps | ✅ Same data | +| Start time | `sessions.started_at` | Turn declaration fact | ✅ Same data | + +**Conclusion**: High overlap. Both store the same execution history from different perspectives. + +### Session Storage vs OTel Traces + +| Data | Session Storage | OTel Traces | Overlap | +|------|---------------|-------------|---------| +| Token counts | ✅ | ✅ (span attributes) | ✅ Partial | +| Tool execution | ✅ | ✅ (span per tool) | ✅ Partial | +| LLM calls | ❌ | ✅ (detailed) | ❌ Different | +| Session concept | ✅ | ❌ | ❌ Different | +| Persistence | ✅ SQLite | ❌ External/In-memory | ❌ Different | + +**Conclusion**: Partial overlap but different purposes (durable history vs ephemeral debugging). OTel is for **observability**, session storage is for **history/replay**. + +### Shepherd Traces vs OTel Traces + +| Aspect | Shepherd | OTel | Relationship | +|--------|----------|------|--------------| +| Purpose | Supervised execution (rollback) | Distributed debugging | Complementary | +| Data model | Fact DAG with causal links | Span tree with parent-child | Different | +| Granularity | Intent + Capture per operation | Span per operation | Similar | +| Supervision | ✅ ScopeManager for rollback | ❌ | Unique to Shepherd | +| Persistence | ✅ SQLite | ❌ (unless external) | Different | + +**Conclusion**: **Overlapping on the turn/tool/token lifecycle; complementary on the model.** Shepherd's fact DAG (causal parents, frontiers, witnesses) and OTel's span tree both record turn boundaries, tool invocations, token counts, and success/error. What is genuinely distinct is Shepherd's supervision capability (checkpoints, tree states, rollback/fork) and OTel's cross-service distributed semantics. + +--- + +## Findings (blocking / reshaping the plan) + +### F1. Shepherd exposes no `TraceStore` interface — Phase 1 as written is not implementable + +The plan's Phase 1 assumes yaah can implement a `shepherd.TraceStore` interface over the shared `*sql.DB`. It cannot. In `shepherd-kernel-go@v0.3.2` (the version currently in `go.mod`): + +- `NewSQLiteTraceStore(path)` returns a **concrete** `*shepherd.SQLiteTraceStore` that **owns its own `*sql.DB` connection** and opens its own file. There is no exported interface for the store. +- The store's schema is far richer than "facts + frontiers". Real tables: `records`, `path_entries`, `record_edges`, `contexts`, `append_intents` (idempotency), `owner_ordinals`, `meta`, `frontiers`, plus witness records. Record IDs are content-addressed SHA-256 digests; append is idempotent via `append_intents`; there is a witness/authorization trust chain (`ensureAppendAuthorized` / `ensureReadAuthorized`). +- `ScopeManager`, `Scope`, checkpoints, and `CaptureTree`/`ApplyTree` all build on that concrete store plus git operations (`scope_manager.go`, `scope.go`). + +So merging Shepherd into `state.db` means one of: + +1. **`ATTACH DATABASE`** — keep the shepherd store on its own file but attach it to the shared connection. Still two files; cosmetic only. +2. **Fork/vendor** the store schema and logic — explicitly out of scope (Non-Goals) and a large, fragile re-implementation of digest/idempotency/witness/scope. +3. **Leave the two files separate** and link them with shared IDs (F3, Phase 0). **This is the recommendation.** + +The plan's proposed `shepherd_facts` / `shepherd_frontiers` schema does not match the real schema and would break idempotency, digests, and the supervisor's checkpoint/rollback. + +### F2. Version and path assumptions are stale + +- `go.mod` requires `github.com/buchenberg/shepherd-kernel-go v0.3.2` (not v0.1.1). +- `shepherd_trace_dir` is empty by default (Shepherd tracing is opt-in); there is no hardcoded `~/.yaah/traces/trace.sqlite`. The trace store path is `filepath.Join(traceDir, "trace.sqlite")` (`internal/agent/pipeline/scope_init.go:26`), with `traceDir` read from `cfg.Agent.Default.ShepherdTraceDir`. + +### F3. The four systems cannot be joined today (highest-leverage gap) + +- OTel spans carry **no `session_id`** attribute (`internal/observability/` has no session reference; the only IDs are the W3C `trace_id`/`span_id`). +- `messages` and `sessions` carry **no `trace_id`**. +- Shepherd facts are keyed by `TraceOwnerID`, which for sub-agents is **not** the session ID — see F4. + +Result: a Jaeger trace, a `messages` row, and a Shepherd fact describing the same turn cannot be correlated. One shared identifier fixes this cheaply without merging any storage. This is Phase 0. + +### F4. Sub-agent trace owners differ from session IDs + +`internal/agent/runner/runner.go:276` assigns sub-agents a synthetic owner `sub-{role}-{parentSession}-{unixnano}` (`subTraceID`), while messages are persisted under the parent `sessionID`. So `shepherd_facts.trace_owner_id = sessions.id` is false for all sub-agent work — any "sum tokens from Shepherd by session ID" (Phase 2 Option C) silently misses sub-agent tokens. + +### F5. Token counts are written in three places, and sessions is the only unconditional one + +- `sessions.tokens_in/out` — written unconditionally by `EndSession` (`cmd/yaah/session.go:132`, sums `totalUsage`). +- Shepherd `turn:completed` capture (`prompt_tokens`/`completion_tokens`) — **opt-in**, only when Shepherd tracing is enabled. +- OTel `llm.prompt_tokens`/`llm.completion_tokens` span events (`FinishLLM`/`FinishStream`) — **observational**, only when OTel is enabled. + +Phase 2 as written makes Shepherd the source of truth, but Shepherd is disabled by default — that would regress the common case. Sessions should remain authoritative; Shepherd/OTel are derived. + +### F6. Verbose OTel re-records full message content already in `messages` + +When `OtelVerbose` is on, `RecordAssistantResponse`, `RecordConversation`, `RecordSystemPrompt`, and `RecordTUIView` (`internal/observability/trace.go`) write full content, reasoning, and tool-call args as span attributes/events — the same data already durably stored in `messages`. This is the clearest "same data, two sinks" and should be either removed (point users at `yaah session show`) or explicitly labeled a debug-only mirror. + +### F7. Memory and session-message search machinery is duplicated + +`memory_fts` vs `messages_fts`, and `SearchMemory`/`SearchMessages` plus `SearchMemoryVector`/`SearchMessagesVector`, are near-identical (FTS + `embedding` BLOB + cosine). The *data* should stay separate (different retention/scoping), but the *plumbing* can be one `SearchableText` abstraction. This is orthogonal to, and cheaper than, Phase 1. + +--- + +## Consolidation Strategy + +### Phase 0: Cross-link the systems with shared IDs (NEW — HIGHEST LEVERAGE, LOW RISK) + +**Goal**: Make every sink queryable against the others without merging any storage (addresses F3/F4). + +**Implementation**: +1. Add a `session_id` attribute to the root `prompt`/`agent.turn` span in `internal/observability/trace.go` (thread the session ID into `StartPrompt`/`StartTurn`, or set it once per run on the span context). +2. Add a nullable `trace_id` column to `messages` (set when OTel is active) so a message row can be joined to a Jaeger trace. +3. Introduce a stable `turn_id` shared by the turn span, the Shepherd `turn:*` facts (`payload["session_id"]` / `payload["turn_id"]`), and the `messages` rows of that turn. +4. Persist the `subTraceID → parentSession` mapping (already computed at `internal/agent/runner/runner.go:276`) so Shepherd owners and parent sessions are joinable (F4). + +**Success**: `trace_id` → `session_id` → `messages` and `trace_owner_id` all resolve to the same conversation. + +**Files**: `internal/observability/trace.go`, `internal/agent/persist.go`, `internal/memory/memory.go` (column), `internal/agent/pipeline/trace.go`. + +--- + +### Phase 1: Co-locate Shepherd Store with Session DB — REVISED (BLOCKED by F1) + +> ⚠ **Superseded.** The implementation below assumes a `shepherd.TraceStore` interface and a `facts`/`frontiers` schema that do not exist (Finding F1). Keep this section as reference only. The revised approach is: +> +> - **Option A — `ATTACH DATABASE`**: one connection, still two files (cosmetic only). +> - **Option B — fork/vendor the store**: single file, but re-implements digest/idempotency/witness/scope; violates Non-Goals. +> - **Option C — keep separate + Phase 0 linkage (recommended)**: two files, joined via `session_id`/`trace_id`; no schema risk. +> +> **Recommendation**: drop the "single file" goal; pursue Phase 0. If a single file is later required, Option A is the only change that respects the upstream schema. + +**Goal (original)**: Move Shepherd trace store from separate `trace.sqlite` into the existing `state.db`. + +**Rationale**: +- Eliminates the two-database problem +- Atomic transactions across sessions + traces +- Simpler backup, migration, and management +- No functional changes to either system + +**Implementation**: + +1. **Add Shepherd tables to `memory/memory.go` migration** + ```go + // In migrate() function, add: + CREATE TABLE IF NOT EXISTS shepherd_facts ( + id TEXT PRIMARY KEY, + trace_owner_id TEXT NOT NULL, + mode TEXT NOT NULL, -- 'declaration' or 'capture' + schema_ref TEXT NOT NULL, + kind_label TEXT NOT NULL, + payload BLOB NOT NULL, + caused_by_ids BLOB, -- JSON array of parent fact IDs + created_at INTEGER NOT NULL + ); + + CREATE TABLE IF NOT EXISTS shepherd_frontiers ( + id TEXT PRIMARY KEY, + target_trace_owner_id TEXT NOT NULL, + through_fact_id TEXT NOT NULL, + created_at INTEGER NOT NULL + ); + + CREATE INDEX IF NOT EXISTS idx_shepherd_trace_owner ON shepherd_facts(trace_owner_id); + CREATE INDEX IF NOT EXISTS idx_shepherd_mode ON shepherd_facts(mode); + ``` + +2. **Create SQLiteTraceStore wrapper** + + Create `internal/memory/shepherd_store.go`: + ```go + package memory + + import shepherd "github.com/buchenberg/shepherd-kernel-go" + + // NewShepherdTraceStore creates a shepherd.SQLiteTraceStore backed by the + // shared yaah DB connection instead of opening a separate SQLite file. + func NewShepherdTraceStore(db *DB) *shepherd.SQLiteTraceStore { + // Use db.sql for all operations + // Implement shepherd.TraceStore interface + } + ``` + +3. **Update `InitShepherdInfrastructure`** + + In `internal/agent/pipeline/scope_init.go`: + ```go + // Accept *memory.DB instead of traceDir path + func InitShepherdInfrastructure(db *memory.DB, busBuffer int) (*shepherd.SQLiteTraceStore, *shepherd.EffectBus, *shepherd.ScopeManager, error) + ``` + +4. **Update wiring** + + In `cmd/yaah/wiring.go`: + ```go + // Pass memory DB to shepherd init + store, bus, mgr, err := pipeline.InitShepherdInfrastructure(db, cfg.Agent.Default.ShepherdBusBuffer) + tools.SharedTraceStore = store + tools.SharedScopeManager = mgr + ``` + +5. **Migration for existing users** + + On first run with new code: + - Detect if `~/.yaah/traces/trace.sqlite` exists + - Migrate data to `state.db` shepherd tables + - Rename old file to backup + +**Files to modify**: +- `internal/memory/memory.go` - add Shepherd tables to migration +- `internal/memory/shepherd_store.go` - NEW, SQLiteTraceStore wrapper +- `internal/agent/pipeline/scope_init.go` - accept *memory.DB +- `cmd/yaah/wiring.go` - pass memory DB to shepherd init +- `cmd/yaah/trace.go` - update `openShepherdTraceStore()` to use `state.db` + +**Test files to update**: +- `internal/agent/pipeline/config_test.go` +- `internal/agent/runner/checkpoint_integration_test.go` + +**CLI impact**: None - transparent to users + +--- + +### Phase 2: Deduplicate Token Storage — REVISED (direction corrected by F4/F5) + +> ⚠ **Superseded.** The original plan makes Shepherd the source of truth, but Shepherd is opt-in (disabled by default) and sub-agent facts use `sub-*` trace owners (F4), so that would regress the default path and miss sub-agent tokens. OTel also writes per-call token events (`FinishLLM`/`FinishStream`), but only when enabled (F5). + +**Revised goal**: keep `sessions.tokens_in/out` as the **authoritative, unconditional** total (already written by `EndSession` at `cmd/yaah/session.go:132`). Treat Shepherd `turn:*` token payloads and OTel `llm.*_tokens` events as **derived/observational**. + +**Revised implementation**: +1. Make no change to `sessions.tokens_in/out` writes. +2. Optionally add `SessionTokenStats()` as a *read model* summing Shepherd facts when present — but never the primary store. +3. If a single source of truth is desired later, the direction is the opposite of the original: derive Shepherd/OTel from sessions, not sessions from Shepherd. + +**Original rationale (for reference)**: +- Token counts stored in both `sessions.tokens_in/out` and Shepherd turn capture facts +- Single source of truth reduces write duplication and potential inconsistency + +**Implementation Options**: + +**Option A: Computed columns (SQLite 3.35.0+, 2021-03-12)** +```go +// In migrate(), add generated columns: +ALTER TABLE sessions ADD COLUMN tokens_in_computed INTEGER GENERATED ALWAYS AS ( + SELECT COALESCE(SUM(CASE WHEN mode = 'capture' AND schema_ref LIKE '%completed.v1' THEN json_extract(payload, '$.prompt_tokens') ELSE 0 END), 0) + FROM shepherd_facts + WHERE trace_owner_id = sessions.id +) STORED; +``` +*Problem*: SQLite GENERATED columns don't support subqueries. Not feasible. + +**Option B: Trigger-based sync** +```sql +CREATE TRIGGER IF NOT EXISTS update_session_tokens_after_shepherd_capture +AFTER INSERT ON shepherd_facts +FOR EACH ROW +WHEN NEW.mode = 'capture' AND NEW.schema_ref LIKE '%completed.v1%' +BEGIN + UPDATE sessions + SET tokens_in = tokens_in + COALESCE(json_extract(NEW.payload, '$.prompt_tokens'), 0), + tokens_out = tokens_out + COALESCE(json_extract(NEW.payload, '$.completion_tokens'), 0) + WHERE id = NEW.trace_owner_id; +END; +``` +*Problem*: JSON extraction in triggers is complex and fragile. + +**Option C: Remove from sessions, compute on read (RECOMMENDED)** +1. Add `SessionTokenStats()` method to `*DB`: + ```go + func (d *DB) SessionTokenStats(sessionID string) (tokensIn, tokensOut int, err error) { + // Sum from Shepherd facts + row := d.sql.QueryRow(` + SELECT + COALESCE(SUM(CASE WHEN json_extract(payload, '$.prompt_tokens') THEN json_extract(payload, '$.prompt_tokens') ELSE 0 END), 0), + COALESCE(SUM(CASE WHEN json_extract(payload, '$.completion_tokens') THEN json_extract(payload, '$.completion_tokens') ELSE 0 END), 0) + FROM shepherd_facts + WHERE trace_owner_id = ? AND mode = 'capture' AND schema_ref LIKE '%completed.v1%' + `, sessionID) + err = row.Scan(&tokensIn, &tokensOut) + return + } + ``` + +2. Deprecate `sessions.tokens_in/tokens_out` columns (keep for backward compat, don't write) +3. Update all readers to use `SessionTokenStats()` when Shepherd is enabled +4. Fallback to legacy columns when Shepherd is disabled + +**Files to modify**: +- `internal/memory/session_repo.go` - add `SessionTokenStats()` method +- `internal/memory/memory.go` - add index on `shepherd_facts(trace_owner_id, mode, schema_ref)` +- `internal/agent/loop.go` - use new method for token reporting + +**Backward compatibility**: Keep columns, just don't write to them. Read falls back to legacy if Shepherd disabled. + +--- + +### Phase 3: Unified CLI Commands (LOW PRIORITY) + +**Goal**: Merge `yaah session` and `yaah shepherd-trace` CLIs into a cohesive interface. + +**Rationale**: Users shouldn't need to know which system stores which data. + +**Proposed CLI**: + +```bash +# List all sessions (combines sessions + trace owners) +yaah session list + +# Show session details (messages + traces) +yaah session show + +# Show execution profile (aggregate stats from both) +yaah session profile + +# Legacy aliases (kept for compatibility) +yaah shepherd-trace list # alias for yaah session list +yaah shepherd-trace show # alias for yaah session show +``` + +**Implementation**: + +1. Rename `cmd/yaah/trace.go` → `cmd/yaah/session_trace.go` +2. Add `session show` subcommand that combines: + - Session metadata from `memory.DB` + - Messages from `memory.DB` + - Shepherd facts from same DB +3. Keep `shepherd-trace` as deprecated aliases + +**Files to modify**: +- `cmd/yaah/trace.go` - refactor into unified commands +- `cmd/yaah/session.go` - add show/profile subcommands + +**Test files**: Minimal - CLI refactor + +--- + +## Migration Path + +> ⚠ **Superseded by F1.** This section assumes Phase 1's single-file merge, which is not feasible without forking `shepherd-kernel-go`. The real store schema is `records`/`path_entries`/`record_edges`/`contexts`/`append_intents`/`owner_ordinals`/`meta`/`frontiers` — not `facts`/`frontiers`. Keep this section as reference only if option B (fork) is later authorized. + +### For Existing Users + +**Before consolidation** (current state): +``` +~/.yaah/ +├── state.db # sessions, messages, memories, todos +└── traces/ + └── trace.sqlite # shepherd facts, frontiers +``` + +**After Phase 1**: +``` +~/.yaah/ +├── state.db # sessions, messages, memories, todos, shepherd_facts, shepherd_frontiers +└── traces/ + └── trace.sqlite # DEPRECATED - read-only for backward compat +``` + +**Migration on first run**: +```go +// In memory.Open() or separate migration function +func (d *DB) MigrateShepherdIfNeeded() error { + oldPath := filepath.Join(filepath.Dir(d.path), "traces", "trace.sqlite") + if _, err := os.Stat(oldPath); err == nil { + // Migrate data from old trace.sqlite to shepherd_* tables + oldDB, err := sql.Open("sqlite", oldPath) + if err != nil { + return err + } + defer oldDB.Close() + + // Copy facts + rows, err := oldDB.Query("SELECT * FROM facts") + // ... insert into d.sql shepherd_facts + + // Copy frontiers + // ... + + // Rename old file + os.Rename(oldPath, oldPath+".migrated") + } + return nil +} +``` + +--- + +## File Change Summary + +### New Files +| File | Purpose | +|------|---------| +| `internal/memory/shepherd_store.go` | SQLiteTraceStore wrapper using shared DB | + +### Modified Files +| File | Changes | +|------|---------| +| `internal/memory/memory.go` | Add Shepherd tables to migration, MigrateShepherdIfNeeded | +| `internal/memory/session_repo.go` | Add SessionTokenStats() method | +| `internal/agent/pipeline/scope_init.go` | Accept *memory.DB instead of traceDir | +| `cmd/yaah/wiring.go` | Pass memory DB to shepherd init | +| `cmd/yaah/trace.go` | Use unified DB, update openShepherdTraceStore() | +| `internal/agent/pipeline/config_test.go` | Update tests for new init signature | +| `internal/agent/runner/checkpoint_integration_test.go` | Update tests for new init signature | + +> **Note**: the above reflects the original Phase 1/2. After revision, the concrete changes are Phase 0 files: `internal/observability/trace.go`, `internal/agent/persist.go`, `internal/memory/memory.go`, `internal/agent/pipeline/trace.go`. + +### Deprecated (future cleanup) +| File/Feature | Status | +|--------------|--------| +| `~/.yaah/traces/trace.sqlite` | Deprecated after migration | +| `sessions.tokens_in/out` | Deprecated, read-only after Phase 2 | + +--- + +## Testing Strategy + +### Unit Tests +1. **Shepherd store wrapper**: Verify all `shepherd.TraceStore` interface methods work against `memory.DB` +2. **Migration**: Test data migration from old `trace.sqlite` to new tables +3. **Token stats**: Verify `SessionTokenStats()` correctly sums from Shepherd facts + +### Integration Tests +1. **End-to-end**: Run yaah session, verify data appears in both session tables and Shepherd tables +2. **CLI**: Verify `yaah session show` displays combined data correctly +3. **Backward compat**: Verify old `trace.sqlite` is migrated and new code works with migrated data + +### Performance Tests +1. **Concurrent writes**: Verify no SQLite contention between session writes and Shepherd fact writes +2. **Query performance**: Verify joined queries across sessions + Shepherd tables perform acceptably + +--- + +## Rollback Plan + +If issues arise during migration: + +1. **Phase 1 rollback**: Revert to opening separate `trace.sqlite` if migration fails + ```go + // In InitShepherdInfrastructure, fallback: + if migrationFailed { + return oldInitShepherdInfrastructure(traceDir, busBuffer) + } + ``` + +2. **Data rollback**: Restore `trace.sqlite` from `.migrated` backup if needed + +3. **Feature flags**: Keep old token columns readable for fallback + +--- + +## Success Criteria + +| Phase | Success Metric | +|-------|----------------| +| Phase 0 | `trace_id` → `session_id` → `messages` / `trace_owner_id` resolve to the same conversation; a Jaeger trace and a `messages` row are mutually navigable | +| Phase 1 (revised) | Two files retained but joined via shared IDs; `yaah shepherd-trace` CLI still works unchanged | +| Phase 2 (revised) | `sessions.tokens_in/out` remains authoritative; Shepherd/OTel treated as derived | +| Phase 3 | Unified `yaah session` CLI provides access to all session data in one place | + +--- + +## Open Questions + +1. **Should `sessions.tokens_in/out` be removed entirely or kept as computed cache?** + - *Recommendation (revised)*: Keep as the authoritative total; do not compute from Shepherd (F4/F5). + +2. **Should Shepherd frontiers be exposed via the unified session API?** + - *Recommendation*: Yes, for supervised task rollback functionality + +3. **How to handle users who disable Shepherd tracing?** + - *Recommendation (revised)*: Irrelevant once sessions stays authoritative (Phase 2 revision); Shepherd/OTel are derived. + +4. **Should the migration be automatic or require explicit user action?** + - *Recommendation*: Moot while Phase 1 is rescoped to "keep separate + link" (F1). Revisit only if option B (fork) is authorized. + +5. **What is the join key?** (`session_id` on spans, `trace_id` on messages, or a new `turn_id`?) + - *Recommendation*: A new stable `turn_id` shared across the span attribute, Shepherd fact payload, and `messages` row; plus `session_id` on the root span as the coarse key. + +--- + +## Related Work + +- `docs/supervised-task-plan.md` - Shepherd infrastructure wiring +- `internal/agent/pipeline/scope_init.go` - Current Shepherd initialization +- `internal/tools/supervisor_shared.go` - Global SharedTraceStore/SharedScopeManager diff --git a/.agents/plans/persistence-hardening/PLAN.md b/.agents/plans/persistence-hardening/PLAN.md new file mode 100644 index 00000000..8efdcb18 --- /dev/null +++ b/.agents/plans/persistence-hardening/PLAN.md @@ -0,0 +1,254 @@ +--- +name: persistence-hardening +description: Adopt the useful parts of shepherd-kernel-go's append-only store (SQLite hardening, exact content dedup, idempotent writes) in yaah's session/memory storage +status: draft +--- + +# Plan: Harden yaah's SQLite persistence (adopt shepherd-store techniques) + +## Goal + +Bring three properties from `shepherd-kernel-go`'s `SQLiteTraceStore` into yaah's +`internal/memory/` (`~/.yaah/state.db`) without importing the parts that are +specific to supervised execution: + +1. **SQLite hardening** — `busy_timeout`, `foreign_keys=ON`, applied per-connection. +2. **Exact memory dedup** — content-addressed dedup (SHA-256 digest) instead of fuzzy FTS matching. +3. **Idempotent message persistence** — a re-persist of the same message is a no-op, not a duplicate or a UNIQUE error. + +These map to the durable wins identified in the shepherd-store analysis; the +provenance/witness/authorization machinery and the fact-DAG/scope/checkpoint +layers are explicitly out of scope. + +## Background + +`../shepherd-kernel-go/store.go` demonstrates three techniques that yaah's +`state.db` currently lacks: + +| Technique | Shepherd | yaah today | +|-----------|----------|------------| +| Busy-timeout / FK hardening | `store.go:62-76` (`SetMaxOpenConns(1)`, WAL, `busy_timeout=5000`, `foreign_keys=ON`) | WAL only (`internal/memory/memory.go:98`); FK on `messages.session_id` is declared (`memory.go:143`) but never enforced | +| Exact dedup | `record_id` = canonical SHA-256 (`canonical.go:170`); `insertRecordIfMissing` short-circuits on identical content (`store.go:1395-1408`) | `AddMemoryDedup` uses FTS `MATCH` (`memory.go:323-335`), which is fuzzy, not exact | +| Idempotent writes | `append_intents(append_intent_id → batch_digest → receipt)` returns the prior receipt on retry (`store.go:499-516`) | `AddMessage` is an unconditional `INSERT` keyed by `(session_id, idx)` (`internal/memory/message_repo.go:8-17`); a retry of the same position errors | + +## Non-Goals + +- Do NOT import the witness/authorization chain (`ensureAppendAuthorized`/`ensureReadAuthorized`, root witness, `PresentedAuthorityRefs`). +- Do NOT import `RetainedContext` / schema-environment refs, or the declaration/capture taxonomy. +- Do NOT add a `record_edges` causal-DAG to `messages` — `tool_call_id`/`tool_name`/`tool_calls` (`memory.go:58-61`) already carry the tool→result linkage. +- Do NOT make `memory` immutable/append-only — `UpdateMemory`/`DeleteMemory` are legitimate for a user-editable knowledge base. +- Do NOT merge the two SQLite files (`state.db` vs `trace.sqlite`) — that is blocked by the shepherd schema (see `consolidate-persistence` plan, Finding F1). + +--- + +## Item 1 — SQLite hardening (per-connection pragmas) + +**Problem**: `memory.Open` sets only `PRAGMA journal_mode=WAL`. yaah has +concurrent writers (the debounced writer, FTS triggers, and the background +embedding goroutines `EmbedMemoryAsync`/`embedMessageAsync`), so a concurrent +write can return `SQLITE_BUSY`; and the `messages.session_id → sessions(id)` +FK is dead because `foreign_keys` is never enabled. + +**Key subtlety**: `PRAGMA` set via `db.Exec(...)` only affects *one* pooled +connection. shepherd works around this with `SetMaxOpenConns(1)`, so its single +connection always has the pragmas. yaah benefits from WAL read concurrency +(search during a turn + background embedding UPDATEs), so a single connection +is too restrictive. Instead use modernc's `_pragma` DSN query parameters, +which are applied to **every** new connection (verified in +`modernc.org/sqlite@v1.53.0/sqlite.go:applyQueryParams`; `busy_timeout` is +applied first by the driver). + +**Change** — `internal/memory/memory.go` `Open()` (currently `sql.Open("sqlite", path)`): + +```go +dsn := path +if !strings.Contains(dsn, "?") { + dsn += "?_pragma=busy_timeout(5000)&_pragma=foreign_keys(1)" +} +db, err := sql.Open("sqlite", dsn) +``` + +Notes: +- `busy_timeout(5000)` mirrors shepherd (`store.go:72`). +- `foreign_keys(1)` enables enforcement of the already-declared FK. +- Leave the connection pool at its default; do **not** copy `SetMaxOpenConns(1)` (see Open Questions). + +**Files**: `internal/memory/memory.go`. + +**Tests to update**: any test that inserts a `messages` row without first +`CreateSession`-ing the parent will now fail on the FK — add the parent session +(or use a helper). Audit `memory_test.go`, `session_repo`/`message_repo` tests. + +--- + +## Item 2 — Exact memory dedup via content digest + +**Problem**: `AddMemoryDedup` (`memory.go:323-335`) detects duplicates with FTS +`MATCH`, which tokenizes and can match near-duplicates — broader and fuzzier +than the documented "identical text" intent. A content digest gives exact, +idempotent dedup like shepherd's content-addressed records. + +**Changes**: + +1. **Schema**: add a nullable `digest` column (no unique index — soft dedup via + `AddMemoryDedup` keeps `AddMemory`/`UpdateMemory` semantics unchanged): + + ```go + // in migrate(), alongside the existing embedding migration + ALTER TABLE memory ADD COLUMN digest TEXT; + ``` + +2. **Digest helper** in `internal/memory/memory.go`: + + ```go + func memoryDigest(text string) string { + sum := sha256.Sum256([]byte(text)) + return hex.EncodeToString(sum[:]) + } + ``` + + (yaah's memory text is a plain string, so a plain SHA-256 suffices — no need + for shepherd's canonical-JSON key-sorting, which exists only for `map[string]any` payloads.) + +3. **`AddMemory`** stores the digest; **`AddMemoryDedup`** becomes an exact lookup: + + ```go + func (d *DB) AddMemoryDedup(e Entry) (string, error) { + digest := memoryDigest(e.Text) + var dupID string + err := d.sql.QueryRow(`SELECT id FROM memory WHERE digest = ? LIMIT 1`, digest).Scan(&dupID) + if err == nil { + return dupID, nil + } + if err != sql.ErrNoRows { + return "", err + } + e.Digest = digest + return "", d.AddMemory(e) + } + ``` + +4. **`UpdateMemory`** recomputes the digest when text changes. + +5. **Backfill** existing rows (mirrors `ReconcileEmbeddings`): a + `ReconcileMemoryDigests` pass that sets `digest` for rows where it is `NULL`. + +**Normalization**: decide whether `digest` is over the raw text or +`strings.TrimSpace(text)` — see Open Questions. Recommend `TrimSpace` to match +"identical text" intent while tolerating trailing whitespace. + +**Files**: `internal/memory/memory.go`. + +--- + +## Item 3 — Idempotent message persistence + +**Problem**: `SessionPersister.Persist` (`internal/agent/persist.go:36-75`) mints +a random `newMessageID()` and a counter `idx` per call, then writes through +`DebouncedWriter` (`internal/memory/debounce.go`) which coalesces by `m.ID`. +Two consequences: + +- A retry of the same logical message mints a **new random ID and a new idx**, + so the debouncer can't coalesce it and the DB gets a duplicate/error. +- `AddMessage` does a bare `INSERT` against `PRIMARY KEY (session_id, idx)`, so + re-persisting the same position returns a UNIQUE error rather than a no-op. + +**Changes**: + +1. **DB layer — upsert on `(session_id, idx)`** in `message_repo.go` `AddMessage`: + + ```sql + INSERT INTO messages (...) VALUES (...) + ON CONFLICT(session_id, idx) DO NOTHING + ``` + + Semantics: first write to a position wins; a retry is a no-op (matches + shepherd's "return prior receipt, no new records"). Skip the background + embed when `RowsAffected() == 0`. + +2. **ID layer — stable message ID** in `persist.go`: replace `newMessageID()` + with a deterministic ID derived from position + content, e.g. + + ```go + func messageID(sessionID string, idx int, role, content string) string { + h := sha256.Sum256([]byte(fmt.Sprintf("%s\x00%d\x00%s\x00%s", sessionID, idx, role, content))) + return hex.EncodeToString(h[:]) + } + ``` + + Benefits: the debouncer's `pending[m.ID]` map now coalesces a re-submitted + message before flush, and the embedding goroutine's + `UPDATE messages SET embedding = ? WHERE id = ?` targets a stable row. + + Historical rows keep their existing random IDs; no backfill required + (ID is not part of any uniqueness constraint — `(session_id, idx)` is). + +**Explicitly NOT doing**: content-dedup of messages. Identical assistant +content at different turns/positions is legitimate, so dedup is position-keyed +(idempotency), not content-keyed — unlike `memory`. + +**Files**: `internal/memory/message_repo.go`, `internal/agent/persist.go`. + +--- + +## Migration & backfill + +- `memory.digest` is additive (`ALTER TABLE ADD COLUMN`). No breaking change; + `ReconcileMemoryDigests` backfills legacy rows on open (wired into `migrate()`). +- `messages` upsert needs no schema change — `PRIMARY KEY (session_id, idx)` + already exists. +- `foreign_keys=ON` is the only behavior change with test fallout: inserts that + reference a missing session now fail. + +--- + +## File Change Summary + +| File | Change | +|------|--------| +| `internal/memory/memory.go` | `_pragma` DSN (busy_timeout, foreign_keys); `digest` column + helper; `AddMemory`/`AddMemoryDedup`/`UpdateMemory` digest support; `ReconcileMemoryDigests` (wired into `migrate`) | +| `internal/memory/message_repo.go` | `AddMessage` → `ON CONFLICT(session_id, idx) DO NOTHING`; skip embed on no-op | +| `internal/agent/persist.go` | deterministic `messageID(...)` replacing `newMessageID()` | + +No new dependencies (`crypto/sha256`, `encoding/hex` are stdlib). + +--- + +## Testing Strategy + +### Unit +1. **Pragmas**: open a DB, assert `PRAGMA foreign_keys` is 1 and `PRAGMA busy_timeout` is 5000 from a *freshly pooled* connection (not just the `Open` caller's connection). +2. **FK enforcement**: inserting a `messages` row for a missing session now errors. +3. **Digest dedup**: `AddMemoryDedup` returns the existing ID for identical text; adds when text differs; `UpdateMemory` re-dedups correctly; backfill sets digests. +4. **Message idempotency**: calling `AddMessage` twice with the same `(session_id, idx)` yields one row and no error; the debouncer coalesces a re-submitted message with the same deterministic ID. + +### Integration +- Run a short agent session; confirm messages persist, `yaah session show` unchanged, and `yaah memory` dedup still works. + +--- + +## Rollback + +- Item 1: drop the `_pragma` query param to revert to WAL-only (no data change). +- Item 2: the `digest` column is additive; revert the `AddMemoryDedup` path to + FTS if the exact-match semantics regress any workflow. +- Item 3: `ON CONFLICT DO NOTHING` is strictly safer than the current + error-on-conflict; revert to bare `INSERT` if position-overwrite semantics are + ever required. + +--- + +## Open Questions + +1. **`SetMaxOpenConns(1)`?** shepherd uses it; yaah's background embedding + writes + concurrent search reads suggest keeping a small pool under WAL with + `busy_timeout`. Confirm via a quick concurrency smoke test. *Recommendation*: + keep default pool, rely on `_pragma busy_timeout`. +2. **Digest normalization**: raw text vs `TrimSpace`? *Recommendation*: + `TrimSpace` (matches "identical text" intent). +3. **Deterministic message ID length**: full SHA-256 hex (64 chars) vs truncated + (e.g. 32 chars)? Purely cosmetic; *Recommendation*: keep full hex, consistent + with `memory.digest`. +4. **Should `AddMemoryDedup` keep the FTS path as a fuzzy fallback** for legacy + rows that predate the `digest` column? *Recommendation*: no — backfill covers + them; keep one exact code path. diff --git a/internal/agent/persist.go b/internal/agent/persist.go index 504fc0c1..7fea25ba 100644 --- a/internal/agent/persist.go +++ b/internal/agent/persist.go @@ -2,7 +2,8 @@ package agent import ( "context" - "crypto/rand" + "crypto/sha256" + "encoding/hex" "encoding/json" "fmt" "os" @@ -56,7 +57,7 @@ func (p *SessionPersister) Persist(msg types.Message) { toolName = msg.Name } m := memory.Message{ - ID: newMessageID(), + ID: messageID(p.sessionID, p.msgIdx, msg.Role, content), SessionID: p.sessionID, Idx: p.msgIdx, Role: msg.Role, @@ -116,8 +117,12 @@ func (p *SessionPersister) writeMsg(m memory.Message) error { return nil } -func newMessageID() string { - b := make([]byte, 16) - rand.Read(b) - return fmt.Sprintf("%x-%x-%x-%x-%x", b[0:4], b[4:6], b[6:8], b[8:10], b[10:]) +// messageID derives a stable, deterministic ID from a message's session, +// position, role, and content. This makes persistence idempotent: the +// debounced writer coalesces a re-submitted message by ID, and the embedding +// goroutine targets a stable row. Position-keyed (not content-keyed) so +// identical content at different positions stays distinct. +func messageID(sessionID string, idx int, role, content string) string { + sum := sha256.Sum256([]byte(fmt.Sprintf("%s\x00%d\x00%s\x00%s", sessionID, idx, role, content))) + return hex.EncodeToString(sum[:]) } diff --git a/internal/memory/debounce_test.go b/internal/memory/debounce_test.go index dde87659..115bec73 100644 --- a/internal/memory/debounce_test.go +++ b/internal/memory/debounce_test.go @@ -14,6 +14,10 @@ func openTestDB(t *testing.T) *DB { t.Fatalf("Open() error: %v", err) } t.Cleanup(func() { db.Close() }) + // Foreign keys are now enforced, so messages need a parent session. + if err := db.CreateSession(Session{ID: "ses-test", StartedAt: time.Now().Unix()}); err != nil { + t.Fatalf("CreateSession() error: %v", err) + } return db } diff --git a/internal/memory/memory.go b/internal/memory/memory.go index 4454beab..9567ab87 100644 --- a/internal/memory/memory.go +++ b/internal/memory/memory.go @@ -12,7 +12,9 @@ package memory import ( "context" + "crypto/sha256" "database/sql" + "encoding/hex" "fmt" "os" "path/filepath" @@ -90,7 +92,7 @@ func Open(path string) (*DB, error) { return nil, fmt.Errorf("create db dir: %w", err) } - db, err := sql.Open("sqlite", path) + db, err := sql.Open("sqlite", sqliteDSN(path)) if err != nil { return nil, fmt.Errorf("open db: %w", err) } @@ -114,6 +116,17 @@ func (d *DB) Close() error { return d.sql.Close() } +// sqliteDSN appends per-connection pragmas to the SQLite DSN. The +// modernc.org/sqlite driver applies `_pragma` query parameters to every new +// pooled connection (busy_timeout first), unlike a one-off `PRAGMA` Exec +// which only affects a single connection. +func sqliteDSN(path string) string { + if strings.Contains(path, "?") { + return path + } + return path + "?_pragma=busy_timeout(5000)&_pragma=foreign_keys(1)" +} + // sanitizeFTSQuery escapes special FTS5 characters and wraps each word // in quotes for exact matching. func sanitizeFTSQuery(query string) string { @@ -160,7 +173,8 @@ func (d *DB) migrate() error { source TEXT, created_at INTEGER NOT NULL, accessed_at INTEGER, - access_count INTEGER DEFAULT 0 + access_count INTEGER DEFAULT 0, + digest TEXT ); CREATE TABLE IF NOT EXISTS todos ( @@ -275,6 +289,13 @@ func (d *DB) migrate() error { d.sql.Exec("ALTER TABLE memory ADD COLUMN embedding BLOB") } + // Migration: add digest column to memory for exact content dedup. + row = d.sql.QueryRow("SELECT COUNT(*) FROM pragma_table_info('memory') WHERE name = 'digest'") + row.Scan(&hasColumn) + if !hasColumn { + d.sql.Exec("ALTER TABLE memory ADD COLUMN digest TEXT") + } + // Migration: add embedding column to messages for vector search. row = d.sql.QueryRow("SELECT COUNT(*) FROM pragma_table_info('messages') WHERE name = 'embedding'") row.Scan(&hasColumn) @@ -282,15 +303,26 @@ func (d *DB) migrate() error { d.sql.Exec("ALTER TABLE messages ADD COLUMN embedding BLOB") } + // Backfill digests for memory rows created before the digest migration. + d.ReconcileMemoryDigests() + return nil } +// memoryDigest returns the exact content digest for a memory entry's text. +// Used for exact dedup; a plain SHA-256 suffices because memory text is a +// plain string (no need for shepherd's canonical-JSON key sorting). +func memoryDigest(text string) string { + sum := sha256.Sum256([]byte(text)) + return hex.EncodeToString(sum[:]) +} + // AddMemory inserts a new memory entry. The caller should separately call // EmbedMemoryAsync to produce the embedding vector. func (d *DB) AddMemory(e Entry) error { _, err := d.sql.Exec( - `INSERT INTO memory (id, text, tags, source, created_at) VALUES (?, ?, ?, ?, ?)`, - e.ID, e.Text, e.Tags, e.Source, e.CreatedAt, + `INSERT INTO memory (id, text, tags, source, created_at, digest) VALUES (?, ?, ?, ?, ?, ?)`, + e.ID, e.Text, e.Tags, e.Source, e.CreatedAt, memoryDigest(e.Text), ) return err } @@ -321,16 +353,15 @@ func (d *DB) EmbedMemoryAsync(id, text string) <-chan struct{} { // existing entry. Returns the ID of the duplicate if found, or empty string // if the entry was added. func (d *DB) AddMemoryDedup(e Entry) (string, error) { - safeQuery := sanitizeFTSQuery(e.Text) - row := d.sql.QueryRow(` - SELECT m.id FROM memory m - JOIN memory_fts ON memory_fts.rowid = m.rowid - WHERE memory_fts MATCH ? LIMIT 1 - `, safeQuery) + digest := memoryDigest(e.Text) var dupID string - if err := row.Scan(&dupID); err == nil { + err := d.sql.QueryRow(`SELECT id FROM memory WHERE digest = ? LIMIT 1`, digest).Scan(&dupID) + if err == nil { return dupID, nil } + if err != sql.ErrNoRows { + return "", err + } return "", d.AddMemory(e) } @@ -554,7 +585,7 @@ func (d *DB) DeleteMemory(id string) error { // UpdateMemory updates the text of an existing memory entry. func (d *DB) UpdateMemory(id string, text string) error { - result, err := d.sql.Exec(`UPDATE memory SET text = ? WHERE id = ?`, text, id) + result, err := d.sql.Exec(`UPDATE memory SET text = ?, digest = ? WHERE id = ?`, text, memoryDigest(text), id) if err != nil { return err } @@ -565,6 +596,41 @@ func (d *DB) UpdateMemory(id string, text string) error { return nil } +// ReconcileMemoryDigests backfills the digest column for rows where it is +// NULL (e.g. entries added before the digest migration). Returns the number +// of rows updated. +func (d *DB) ReconcileMemoryDigests() (int, error) { + rows, err := d.sql.Query(`SELECT id, text FROM memory WHERE digest IS NULL`) + if err != nil { + return 0, fmt.Errorf("reconcile digests: query: %w", err) + } + defer rows.Close() + + type row struct { + id, text string + } + var pending []row + for rows.Next() { + var r row + if err := rows.Scan(&r.id, &r.text); err != nil { + return 0, err + } + pending = append(pending, r) + } + if err := rows.Err(); err != nil { + return 0, err + } + + updated := 0 + for _, r := range pending { + if _, err := d.sql.Exec(`UPDATE memory SET digest = ? WHERE id = ?`, memoryDigest(r.text), r.id); err != nil { + return updated, fmt.Errorf("reconcile digests: update %s: %w", r.id, err) + } + updated++ + } + return updated, nil +} + // ListMemory returns the most recent memory entries. func (d *DB) ListMemory(limit int) ([]Entry, error) { rows, err := d.sql.Query(` diff --git a/internal/memory/memory_test.go b/internal/memory/memory_test.go index 61534ff2..d2adcbcb 100644 --- a/internal/memory/memory_test.go +++ b/internal/memory/memory_test.go @@ -767,3 +767,162 @@ func TestDB_SearchMessagesVectorNoEmbedder(t *testing.T) { t.Fatalf("expected empty when embedder is nil, got %d results", len(results)) } } + +func TestDB_SQLiteDSN(t *testing.T) { + want := "/tmp/a.db?_pragma=busy_timeout(5000)&_pragma=foreign_keys(1)" + if got := sqliteDSN("/tmp/a.db"); got != want { + t.Errorf("sqliteDSN(plain) = %q, want %q", got, want) + } + if got := sqliteDSN("/tmp/a.db?cache=shared"); got != "/tmp/a.db?cache=shared" { + t.Errorf("sqliteDSN(existing query) = %q", got) + } +} + +func TestDB_PragmasApplied(t *testing.T) { + tmp := t.TempDir() + db, err := Open(filepath.Join(tmp, "test.db")) + if err != nil { + t.Fatal(err) + } + defer db.Close() + + var fk int + if err := db.sql.QueryRow(`PRAGMA foreign_keys`).Scan(&fk); err != nil { + t.Fatal(err) + } + if fk != 1 { + t.Errorf("foreign_keys = %d, want 1", fk) + } + + var bt int + if err := db.sql.QueryRow(`PRAGMA busy_timeout`).Scan(&bt); err != nil { + t.Fatal(err) + } + if bt != 5000 { + t.Errorf("busy_timeout = %d, want 5000", bt) + } +} + +func TestDB_ForeignKeysEnforced(t *testing.T) { + tmp := t.TempDir() + db, err := Open(filepath.Join(tmp, "test.db")) + if err != nil { + t.Fatal(err) + } + defer db.Close() + + m := Message{SessionID: "missing-session", Idx: 0, Role: "user", Content: "hi", Timestamp: 1, ID: "m1"} + if err := db.AddMessage(m); err == nil { + t.Fatal("expected foreign key error for message with no parent session") + } +} + +func TestDB_AddMessage_IdempotentOnSamePosition(t *testing.T) { + tmp := t.TempDir() + db, err := Open(filepath.Join(tmp, "test.db")) + if err != nil { + t.Fatal(err) + } + defer db.Close() + + if err := db.CreateSession(Session{ID: "sess-1", StartedAt: 1}); err != nil { + t.Fatal(err) + } + + m := Message{SessionID: "sess-1", Idx: 0, Role: "user", Content: "hello", Timestamp: 1, ID: "m1"} + if err := db.AddMessage(m); err != nil { + t.Fatalf("first AddMessage: %v", err) + } + if err := db.AddMessage(m); err != nil { + t.Fatalf("second AddMessage (idempotent) should not error: %v", err) + } + + msgs, err := db.GetMessages("sess-1") + if err != nil { + t.Fatal(err) + } + if len(msgs) != 1 { + t.Fatalf("expected 1 message after duplicate insert, got %d", len(msgs)) + } +} + +func TestDB_AddMemoryDedup_ExactOnly(t *testing.T) { + tmp := t.TempDir() + db, err := Open(filepath.Join(tmp, "test.db")) + if err != nil { + t.Fatal(err) + } + defer db.Close() + + if err := db.AddMemory(Entry{ID: "mem-1", Text: "User prefers dark mode", Source: "agent", CreatedAt: 1}); err != nil { + t.Fatal(err) + } + // Near-duplicate (trailing punctuation) must NOT be deduped: dedup is exact. + dupID, err := db.AddMemoryDedup(Entry{ID: "mem-2", Text: "User prefers dark mode!", Source: "agent", CreatedAt: 2}) + if err != nil { + t.Fatal(err) + } + if dupID != "" { + t.Errorf("near-duplicate text should not be deduped, got %q", dupID) + } + all, _ := db.ListMemory(10) + if len(all) != 2 { + t.Errorf("expected 2 memories, got %d", len(all)) + } +} + +func TestDB_UpdateMemory_RecomputesDigest(t *testing.T) { + tmp := t.TempDir() + db, err := Open(filepath.Join(tmp, "test.db")) + if err != nil { + t.Fatal(err) + } + defer db.Close() + + if err := db.AddMemory(Entry{ID: "mem-1", Text: "alpha", Source: "agent", CreatedAt: 1}); err != nil { + t.Fatal(err) + } + if err := db.AddMemory(Entry{ID: "mem-2", Text: "beta", Source: "agent", CreatedAt: 2}); err != nil { + t.Fatal(err) + } + if err := db.UpdateMemory("mem-2", "gamma"); err != nil { + t.Fatal(err) + } + + dupID, err := db.AddMemoryDedup(Entry{ID: "mem-3", Text: "gamma", Source: "agent", CreatedAt: 3}) + if err != nil { + t.Fatal(err) + } + if dupID != "mem-2" { + t.Errorf("expected dedup against updated mem-2, got %q", dupID) + } +} + +func TestDB_ReconcileMemoryDigests(t *testing.T) { + tmp := t.TempDir() + db, err := Open(filepath.Join(tmp, "test.db")) + if err != nil { + t.Fatal(err) + } + defer db.Close() + + if _, err := db.sql.Exec(`INSERT INTO memory (id, text, tags, source, created_at) VALUES ('legacy-1', 'legacy text', NULL, 'agent', 1)`); err != nil { + t.Fatal(err) + } + + n, err := db.ReconcileMemoryDigests() + if err != nil { + t.Fatal(err) + } + if n != 1 { + t.Errorf("expected 1 backfilled row, got %d", n) + } + + var digest string + if err := db.sql.QueryRow(`SELECT digest FROM memory WHERE id = 'legacy-1'`).Scan(&digest); err != nil { + t.Fatal(err) + } + if digest != memoryDigest("legacy text") { + t.Errorf("digest = %q, want %q", digest, memoryDigest("legacy text")) + } +} diff --git a/internal/memory/message_repo.go b/internal/memory/message_repo.go index 3e60b7d5..35d5537e 100644 --- a/internal/memory/message_repo.go +++ b/internal/memory/message_repo.go @@ -6,14 +6,19 @@ import "context" // configured and the role is "user" or "assistant", the content is // embedded in a background goroutine so the caller is not blocked. func (d *DB) AddMessage(m Message) error { - _, err := d.sql.Exec( - `INSERT INTO messages (session_id, idx, role, content, reasoning_content, tool_name, tool_call_id, tool_calls, ts, id) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`, + res, err := d.sql.Exec( + `INSERT INTO messages (session_id, idx, role, content, reasoning_content, tool_name, tool_call_id, tool_calls, ts, id) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + ON CONFLICT(session_id, idx) DO NOTHING`, m.SessionID, m.Idx, m.Role, m.Content, m.ReasoningContent, m.ToolName, m.ToolCallID, m.ToolCalls, m.Timestamp, m.ID, ) - if err == nil { - d.embedMessageAsync(m.ID, m.Role, m.Content) + if err != nil { + return err + } + if n, err := res.RowsAffected(); err == nil && n == 0 { + return nil } - return err + d.embedMessageAsync(m.ID, m.Role, m.Content) + return nil } // embedMessageAsync embeds the content in a background goroutine and From e77e7da3e3108dd187ac6330f826a8a836d72554 Mon Sep 17 00:00:00 2001 From: Gregory Buchenberger Date: Thu, 20 Aug 2026 09:38:43 -0600 Subject: [PATCH 2/2] fix: address review feedback on persistence hardening Make memory dedup atomic via a write transaction (_txlock=immediate), include all immutable fields in the deterministic message ID, validate message-position conflicts instead of silently dropping, preserve existing DSN query params, return migration errors, and clean obsolete Shepherd material from the plans. --- .agents/plans/consolidate-persistence/PLAN.md | 93 +++++-------------- .agents/plans/persistence-hardening/PLAN.md | 87 ++++++++++------- internal/agent/persist.go | 20 ++-- internal/memory/memory.go | 47 ++++++++-- internal/memory/memory_test.go | 40 +++++++- internal/memory/message_repo.go | 25 ++++- 6 files changed, 184 insertions(+), 128 deletions(-) diff --git a/.agents/plans/consolidate-persistence/PLAN.md b/.agents/plans/consolidate-persistence/PLAN.md index 421fcd24..5043b251 100644 --- a/.agents/plans/consolidate-persistence/PLAN.md +++ b/.agents/plans/consolidate-persistence/PLAN.md @@ -17,7 +17,7 @@ Reduce redundancy across yaah's four persistence/observability systems by consol ## Problem Statement -Currently yaah maintains **two separate SQLite databases** (`~/.yaah/state.db` and `~/.yaah/traces/trace.sqlite`) that store overlapping data about tool executions, token usage, and session metadata. This creates: +Currently yaah maintains **two separate SQLite databases** (`~/.yaah/state.db` and a Shepherd trace store at `/trace.sqlite` — configurable, opt-in, empty by default) that store overlapping data about tool executions, token usage, and session metadata. This creates: - **Storage fragmentation**: Two files to backup/migrate/manage - **Query fragmentation**: Different CLIs for similar data (`yaah session` vs `yaah shepherd-trace`) @@ -53,7 +53,7 @@ Currently yaah maintains **two separate SQLite databases** (`~/.yaah/state.db` a **Location**: `internal/agent/pipeline/trace.go` + `shepherd-kernel-go` -**Storage**: Separate SQLite file at `~/.yaah/traces/trace.sqlite` +**Storage**: Separate SQLite file at `filepath.Join(shepherd_trace_dir, "trace.sqlite")` — the directory is configurable (`cfg.Agent.Default.ShepherdTraceDir`) and tracing is opt-in (empty by default). **Data Model**: - Facts with Declaration/Capture modes @@ -199,11 +199,11 @@ When `OtelVerbose` is on, `RecordAssistantResponse`, `RecordConversation`, `Reco > ⚠ **Superseded.** The implementation below assumes a `shepherd.TraceStore` interface and a `facts`/`frontiers` schema that do not exist (Finding F1). Keep this section as reference only. The revised approach is: > -> - **Option A — `ATTACH DATABASE`**: one connection, still two files (cosmetic only). +> - **Option A — `ATTACH DATABASE`**: `ATTACH` is scoped to a single SQLite connection, but `NewSQLiteTraceStore` opens and owns its own `*sql.DB`, so it cannot reach an attachment on yaah's connection. This requires an upstream connection-injection API or a fork, so it is not viable as-is. > - **Option B — fork/vendor the store**: single file, but re-implements digest/idempotency/witness/scope; violates Non-Goals. > - **Option C — keep separate + Phase 0 linkage (recommended)**: two files, joined via `session_id`/`trace_id`; no schema risk. > -> **Recommendation**: drop the "single file" goal; pursue Phase 0. If a single file is later required, Option A is the only change that respects the upstream schema. +> **Recommendation**: drop the "single file" goal; pursue Phase 0. A single file would require an upstream connection-injection API or a fork (Option B); keep Option C unless that dependency change is made. **Goal (original)**: Move Shepherd trace store from separate `trace.sqlite` into the existing `state.db`. @@ -431,54 +431,21 @@ yaah shepherd-trace show # alias for yaah session show └── trace.sqlite # DEPRECATED - read-only for backward compat ``` -**Migration on first run**: -```go -// In memory.Open() or separate migration function -func (d *DB) MigrateShepherdIfNeeded() error { - oldPath := filepath.Join(filepath.Dir(d.path), "traces", "trace.sqlite") - if _, err := os.Stat(oldPath); err == nil { - // Migrate data from old trace.sqlite to shepherd_* tables - oldDB, err := sql.Open("sqlite", oldPath) - if err != nil { - return err - } - defer oldDB.Close() - - // Copy facts - rows, err := oldDB.Query("SELECT * FROM facts") - // ... insert into d.sql shepherd_facts - - // Copy frontiers - // ... - - // Rename old file - os.Rename(oldPath, oldPath+".migrated") - } - return nil -} -``` +**Migration on first run**: not applicable — there is no `facts`/`frontiers` +schema to migrate (the real store uses `records`/`path_entries`/`record_edges`/ +`contexts`/`append_intents`/`owner_ordinals`/`meta`/`frontiers`). No +`MigrateShepherdIfNeeded` step is planned. --- ## File Change Summary -### New Files -| File | Purpose | -|------|---------| -| `internal/memory/shepherd_store.go` | SQLiteTraceStore wrapper using shared DB | - -### Modified Files -| File | Changes | -|------|---------| -| `internal/memory/memory.go` | Add Shepherd tables to migration, MigrateShepherdIfNeeded | -| `internal/memory/session_repo.go` | Add SessionTokenStats() method | -| `internal/agent/pipeline/scope_init.go` | Accept *memory.DB instead of traceDir | -| `cmd/yaah/wiring.go` | Pass memory DB to shepherd init | -| `cmd/yaah/trace.go` | Use unified DB, update openShepherdTraceStore() | -| `internal/agent/pipeline/config_test.go` | Update tests for new init signature | -| `internal/agent/runner/checkpoint_integration_test.go` | Update tests for new init signature | - -> **Note**: the above reflects the original Phase 1/2. After revision, the concrete changes are Phase 0 files: `internal/observability/trace.go`, `internal/agent/persist.go`, `internal/memory/memory.go`, `internal/agent/pipeline/trace.go`. +> The original Phase 1/2 file list (a `shepherd_store.go` wrapper, +> `MigrateShepherdIfNeeded`, `SessionTokenStats`, and their tests) is +> **non-actionable** — it depends on a `shepherd.TraceStore` interface and a +> `facts`/`frontiers` schema that do not exist (F1). The revised Phase 0 scope +> touches: `internal/observability/trace.go`, `internal/agent/persist.go`, +> `internal/memory/memory.go`, `internal/agent/pipeline/trace.go`. ### Deprecated (future cleanup) | File/Feature | Status | @@ -490,19 +457,13 @@ func (d *DB) MigrateShepherdIfNeeded() error { ## Testing Strategy -### Unit Tests -1. **Shepherd store wrapper**: Verify all `shepherd.TraceStore` interface methods work against `memory.DB` -2. **Migration**: Test data migration from old `trace.sqlite` to new tables -3. **Token stats**: Verify `SessionTokenStats()` correctly sums from Shepherd facts - -### Integration Tests -1. **End-to-end**: Run yaah session, verify data appears in both session tables and Shepherd tables -2. **CLI**: Verify `yaah session show` displays combined data correctly -3. **Backward compat**: Verify old `trace.sqlite` is migrated and new code works with migrated data +The original Phase 1/2 test plan (a `shepherd.TraceStore` wrapper test, a +`trace.sqlite` migration test, and `SessionTokenStats()` tests) is +**non-actionable** (F1). The revised Phase 0 scope is: -### Performance Tests -1. **Concurrent writes**: Verify no SQLite contention between session writes and Shepherd fact writes -2. **Query performance**: Verify joined queries across sessions + Shepherd tables perform acceptably +1. Root OTel span carries `session_id`; `messages` carry `trace_id`/`turn_id`. +2. `subTraceID → parentSession` mapping is joinable. +3. `yaah session` / `yaah shepherd-trace` still resolve a conversation end-to-end. --- @@ -510,17 +471,9 @@ func (d *DB) MigrateShepherdIfNeeded() error { If issues arise during migration: -1. **Phase 1 rollback**: Revert to opening separate `trace.sqlite` if migration fails - ```go - // In InitShepherdInfrastructure, fallback: - if migrationFailed { - return oldInitShepherdInfrastructure(traceDir, busBuffer) - } - ``` - -2. **Data rollback**: Restore `trace.sqlite` from `.migrated` backup if needed - -3. **Feature flags**: Keep old token columns readable for fallback +There is no data migration to roll back — the two files stay separate +(Option C). The only reversible change is the Phase 0 linkage columns, which +are additive. --- diff --git a/.agents/plans/persistence-hardening/PLAN.md b/.agents/plans/persistence-hardening/PLAN.md index 8efdcb18..27d0bb6e 100644 --- a/.agents/plans/persistence-hardening/PLAN.md +++ b/.agents/plans/persistence-hardening/PLAN.md @@ -61,16 +61,19 @@ applied first by the driver). **Change** — `internal/memory/memory.go` `Open()` (currently `sql.Open("sqlite", path)`): ```go -dsn := path -if !strings.Contains(dsn, "?") { - dsn += "?_pragma=busy_timeout(5000)&_pragma=foreign_keys(1)" +sep := "?" +if strings.Contains(path, "?") { + sep = "&" } +dsn := path + sep + "_pragma=busy_timeout(5000)&_pragma=foreign_keys(1)&_txlock=immediate" db, err := sql.Open("sqlite", dsn) ``` Notes: - `busy_timeout(5000)` mirrors shepherd (`store.go:72`). - `foreign_keys(1)` enables enforcement of the already-declared FK. +- Existing query parameters are preserved (separator is `&` when a `?` is already present). +- `_txlock=immediate` makes `Begin()` take the write lock up front, so the `AddMemoryDedup` transaction (Item 2) is atomic. - Leave the connection pool at its default; do **not** copy `SetMaxOpenConns(1)` (see Open Questions). **Files**: `internal/memory/memory.go`. @@ -91,11 +94,17 @@ idempotent dedup like shepherd's content-addressed records. **Changes**: 1. **Schema**: add a nullable `digest` column (no unique index — soft dedup via - `AddMemoryDedup` keeps `AddMemory`/`UpdateMemory` semantics unchanged): + `AddMemoryDedup` keeps `AddMemory`/`UpdateMemory` semantics unchanged). The + migration is guarded so repeated `Open()` is idempotent, and returns errors: ```go // in migrate(), alongside the existing embedding migration - ALTER TABLE memory ADD COLUMN digest TEXT; + row := d.sql.QueryRow("SELECT COUNT(*) FROM pragma_table_info('memory') WHERE name = 'digest'") + var hasDigest bool + if err := row.Scan(&hasDigest); err != nil { return err } + if !hasDigest { + if _, err := d.sql.Exec("ALTER TABLE memory ADD COLUMN digest TEXT"); err != nil { return err } + } ``` 2. **Digest helper** in `internal/memory/memory.go`: @@ -110,21 +119,23 @@ idempotent dedup like shepherd's content-addressed records. (yaah's memory text is a plain string, so a plain SHA-256 suffices — no need for shepherd's canonical-JSON key-sorting, which exists only for `map[string]any` payloads.) -3. **`AddMemory`** stores the digest; **`AddMemoryDedup`** becomes an exact lookup: +3. **`AddMemory`** stores the digest; **`AddMemoryDedup`** does an atomic + lookup+insert inside one write transaction (the DSN's `_txlock=immediate` + makes `Begin()` take the write lock, so concurrent dedup calls serialize): ```go func (d *DB) AddMemoryDedup(e Entry) (string, error) { digest := memoryDigest(e.Text) + tx, err := d.sql.Begin() + if err != nil { return "", err } + defer tx.Rollback() var dupID string - err := d.sql.QueryRow(`SELECT id FROM memory WHERE digest = ? LIMIT 1`, digest).Scan(&dupID) - if err == nil { - return dupID, nil - } - if err != sql.ErrNoRows { - return "", err - } - e.Digest = digest - return "", d.AddMemory(e) + err = tx.QueryRow(`SELECT id FROM memory WHERE digest = ? LIMIT 1`, digest).Scan(&dupID) + if err == nil { return dupID, nil } + if err != sql.ErrNoRows { return "", err } + if _, err := tx.Exec(`INSERT INTO memory (...) VALUES (...)`, ...); err != nil { return "", err } + if err := tx.Commit(); err != nil { return "", err } + return "", nil } ``` @@ -133,9 +144,10 @@ idempotent dedup like shepherd's content-addressed records. 5. **Backfill** existing rows (mirrors `ReconcileEmbeddings`): a `ReconcileMemoryDigests` pass that sets `digest` for rows where it is `NULL`. -**Normalization**: decide whether `digest` is over the raw text or -`strings.TrimSpace(text)` — see Open Questions. Recommend `TrimSpace` to match -"identical text" intent while tolerating trailing whitespace. +**Normalization**: `memoryDigest` hashes the raw text exactly as stored (no +`TrimSpace`). Every digest producer (`AddMemory`, `AddMemoryDedup`, +`UpdateMemory`, `ReconcileMemoryDigests`) calls the same `memoryDigest` helper, +so the rule is applied consistently. **Files**: `internal/memory/memory.go`. @@ -162,23 +174,28 @@ Two consequences: ON CONFLICT(session_id, idx) DO NOTHING ``` - Semantics: first write to a position wins; a retry is a no-op (matches - shepherd's "return prior receipt, no new records"). Skip the background - embed when `RowsAffected() == 0`. + Semantics: a retry with the same deterministic ID is a no-op; a different + message at the same position returns a conflict error (never silently + dropped). Only a newly inserted row starts background embedding. 2. **ID layer — stable message ID** in `persist.go`: replace `newMessageID()` - with a deterministic ID derived from position + content, e.g. + with a deterministic ID over **all** immutable persisted fields (session, + position, role, content, reasoning, tool name, tool-call ID, tool-calls): ```go - func messageID(sessionID string, idx int, role, content string) string { - h := sha256.Sum256([]byte(fmt.Sprintf("%s\x00%d\x00%s\x00%s", sessionID, idx, role, content))) + func messageID(sessionID string, idx int, role, content, reasoning, toolName, toolCallID, toolCalls string) string { + h := sha256.Sum256([]byte(fmt.Sprintf( + "%s\x00%d\x00%s\x00%s\x00%s\x00%s\x00%s\x00%s", + sessionID, idx, role, content, reasoning, toolName, toolCallID, toolCalls))) return hex.EncodeToString(h[:]) } ``` Benefits: the debouncer's `pending[m.ID]` map now coalesces a re-submitted message before flush, and the embedding goroutine's - `UPDATE messages SET embedding = ? WHERE id = ?` targets a stable row. + `UPDATE messages SET embedding = ? WHERE id = ?` targets a stable row. The ID + doubles as a content fingerprint, so `AddMessage` can detect a *different* + message occupying the same position (compare stored `id` on conflict). Historical rows keep their existing random IDs; no backfill required (ID is not part of any uniqueness constraint — `(session_id, idx)` is). @@ -206,9 +223,9 @@ content at different turns/positions is legitimate, so dedup is position-keyed | File | Change | |------|--------| -| `internal/memory/memory.go` | `_pragma` DSN (busy_timeout, foreign_keys); `digest` column + helper; `AddMemory`/`AddMemoryDedup`/`UpdateMemory` digest support; `ReconcileMemoryDigests` (wired into `migrate`) | -| `internal/memory/message_repo.go` | `AddMessage` → `ON CONFLICT(session_id, idx) DO NOTHING`; skip embed on no-op | -| `internal/agent/persist.go` | deterministic `messageID(...)` replacing `newMessageID()` | +| `internal/memory/memory.go` | `_pragma` + `_txlock` DSN (preserving existing params); guarded `digest` migration with error handling; `AddMemory`/`AddMemoryDedup`/`UpdateMemory` digest support; atomic transactional dedup; `ReconcileMemoryDigests` (wired into `migrate`) | +| `internal/memory/message_repo.go` | `AddMessage` → idempotent `ON CONFLICT(session_id, idx) DO NOTHING` with conflict validation; embed only on insert | +| `internal/agent/persist.go` | deterministic `messageID(...)` over all immutable fields, replacing `newMessageID()` | No new dependencies (`crypto/sha256`, `encoding/hex` are stdlib). @@ -217,10 +234,11 @@ No new dependencies (`crypto/sha256`, `encoding/hex` are stdlib). ## Testing Strategy ### Unit -1. **Pragmas**: open a DB, assert `PRAGMA foreign_keys` is 1 and `PRAGMA busy_timeout` is 5000 from a *freshly pooled* connection (not just the `Open` caller's connection). +1. **Pragmas**: open a DB, assert `PRAGMA foreign_keys` is 1 and `PRAGMA busy_timeout` is 5000 from a *freshly pooled* connection (not just the `Open` caller's connection); `sqliteDSN` preserves existing query params and appends the pragmas. 2. **FK enforcement**: inserting a `messages` row for a missing session now errors. 3. **Digest dedup**: `AddMemoryDedup` returns the existing ID for identical text; adds when text differs; `UpdateMemory` re-dedups correctly; backfill sets digests. -4. **Message idempotency**: calling `AddMessage` twice with the same `(session_id, idx)` yields one row and no error; the debouncer coalesces a re-submitted message with the same deterministic ID. +4. **Message idempotency**: calling `AddMessage` twice with the same `(session_id, idx)` and content yields one row and no error; a different message at the same position returns a conflict error; the debouncer coalesces a re-submitted message with the same deterministic ID. +5. **Migration idempotency**: `Open()` twice on the same file succeeds (guarded digest column add). ### Integration - Run a short agent session; confirm messages persist, `yaah session show` unchanged, and `yaah memory` dedup still works. @@ -232,9 +250,9 @@ No new dependencies (`crypto/sha256`, `encoding/hex` are stdlib). - Item 1: drop the `_pragma` query param to revert to WAL-only (no data change). - Item 2: the `digest` column is additive; revert the `AddMemoryDedup` path to FTS if the exact-match semantics regress any workflow. -- Item 3: `ON CONFLICT DO NOTHING` is strictly safer than the current - error-on-conflict; revert to bare `INSERT` if position-overwrite semantics are - ever required. +- Item 3: reverting `AddMessage` to a bare `INSERT` restores uniqueness errors + on conflict (not position-overwrite). If overwrite-on-conflict is ever + required, use an explicit `ON CONFLICT(session_id, idx) DO UPDATE` policy. --- @@ -244,8 +262,7 @@ No new dependencies (`crypto/sha256`, `encoding/hex` are stdlib). writes + concurrent search reads suggest keeping a small pool under WAL with `busy_timeout`. Confirm via a quick concurrency smoke test. *Recommendation*: keep default pool, rely on `_pragma busy_timeout`. -2. **Digest normalization**: raw text vs `TrimSpace`? *Recommendation*: - `TrimSpace` (matches "identical text" intent). +2. **Digest normalization**: raw text (as stored) — resolved; no `TrimSpace`. 3. **Deterministic message ID length**: full SHA-256 hex (64 chars) vs truncated (e.g. 32 chars)? Purely cosmetic; *Recommendation*: keep full hex, consistent with `memory.digest`. diff --git a/internal/agent/persist.go b/internal/agent/persist.go index 7fea25ba..02754d6f 100644 --- a/internal/agent/persist.go +++ b/internal/agent/persist.go @@ -57,7 +57,7 @@ func (p *SessionPersister) Persist(msg types.Message) { toolName = msg.Name } m := memory.Message{ - ID: messageID(p.sessionID, p.msgIdx, msg.Role, content), + ID: messageID(p.sessionID, p.msgIdx, msg.Role, content, msg.ReasoningContent, toolName, msg.ToolCallID, toolCallsJSON), SessionID: p.sessionID, Idx: p.msgIdx, Role: msg.Role, @@ -117,12 +117,16 @@ func (p *SessionPersister) writeMsg(m memory.Message) error { return nil } -// messageID derives a stable, deterministic ID from a message's session, -// position, role, and content. This makes persistence idempotent: the -// debounced writer coalesces a re-submitted message by ID, and the embedding -// goroutine targets a stable row. Position-keyed (not content-keyed) so -// identical content at different positions stays distinct. -func messageID(sessionID string, idx int, role, content string) string { - sum := sha256.Sum256([]byte(fmt.Sprintf("%s\x00%d\x00%s\x00%s", sessionID, idx, role, content))) +// messageID derives a stable, deterministic ID from all immutable persisted +// message fields. This makes persistence idempotent: the debounced writer +// coalesces a re-submitted message by ID, and the embedding goroutine targets +// a stable row. Position-keyed (session + idx) so identical content at +// different positions stays distinct; content-keyed (all remaining fields) so +// a changed message at the same position gets a different fingerprint. +func messageID(sessionID string, idx int, role, content, reasoning, toolName, toolCallID, toolCalls string) string { + sum := sha256.Sum256([]byte(fmt.Sprintf( + "%s\x00%d\x00%s\x00%s\x00%s\x00%s\x00%s\x00%s", + sessionID, idx, role, content, reasoning, toolName, toolCallID, toolCalls, + ))) return hex.EncodeToString(sum[:]) } diff --git a/internal/memory/memory.go b/internal/memory/memory.go index 9567ab87..d7ef3ecf 100644 --- a/internal/memory/memory.go +++ b/internal/memory/memory.go @@ -119,12 +119,15 @@ func (d *DB) Close() error { // sqliteDSN appends per-connection pragmas to the SQLite DSN. The // modernc.org/sqlite driver applies `_pragma` query parameters to every new // pooled connection (busy_timeout first), unlike a one-off `PRAGMA` Exec -// which only affects a single connection. +// which only affects a single connection. Existing query parameters are +// preserved. `_txlock=immediate` makes Begin() take the write lock up front, +// so the AddMemoryDedup lookup+insert transaction is atomic. func sqliteDSN(path string) string { + sep := "?" if strings.Contains(path, "?") { - return path + sep = "&" } - return path + "?_pragma=busy_timeout(5000)&_pragma=foreign_keys(1)" + return path + sep + "_pragma=busy_timeout(5000)&_pragma=foreign_keys(1)&_txlock=immediate" } // sanitizeFTSQuery escapes special FTS5 characters and wraps each word @@ -291,9 +294,13 @@ func (d *DB) migrate() error { // Migration: add digest column to memory for exact content dedup. row = d.sql.QueryRow("SELECT COUNT(*) FROM pragma_table_info('memory') WHERE name = 'digest'") - row.Scan(&hasColumn) + if err := row.Scan(&hasColumn); err != nil { + return fmt.Errorf("check memory.digest column: %w", err) + } if !hasColumn { - d.sql.Exec("ALTER TABLE memory ADD COLUMN digest TEXT") + if _, err := d.sql.Exec("ALTER TABLE memory ADD COLUMN digest TEXT"); err != nil { + return fmt.Errorf("add memory.digest column: %w", err) + } } // Migration: add embedding column to messages for vector search. @@ -304,7 +311,9 @@ func (d *DB) migrate() error { } // Backfill digests for memory rows created before the digest migration. - d.ReconcileMemoryDigests() + if _, err := d.ReconcileMemoryDigests(); err != nil { + return fmt.Errorf("reconcile memory digests: %w", err) + } return nil } @@ -351,18 +360,38 @@ func (d *DB) EmbedMemoryAsync(id, text string) <-chan struct{} { // AddMemoryDedup adds a memory entry, skipping if text is identical to an // existing entry. Returns the ID of the duplicate if found, or empty string -// if the entry was added. +// if the entry was added. The lookup and insert run in a single write +// transaction (the DSN sets _txlock=immediate) so concurrent dedup calls +// cannot both observe "no match" and insert duplicate rows. func (d *DB) AddMemoryDedup(e Entry) (string, error) { digest := memoryDigest(e.Text) + + tx, err := d.sql.Begin() + if err != nil { + return "", err + } + defer tx.Rollback() + var dupID string - err := d.sql.QueryRow(`SELECT id FROM memory WHERE digest = ? LIMIT 1`, digest).Scan(&dupID) + err = tx.QueryRow(`SELECT id FROM memory WHERE digest = ? LIMIT 1`, digest).Scan(&dupID) if err == nil { return dupID, nil } if err != sql.ErrNoRows { return "", err } - return "", d.AddMemory(e) + + if _, err := tx.Exec( + `INSERT INTO memory (id, text, tags, source, created_at, digest) VALUES (?, ?, ?, ?, ?, ?)`, + e.ID, e.Text, e.Tags, e.Source, e.CreatedAt, digest, + ); err != nil { + return "", err + } + + if err := tx.Commit(); err != nil { + return "", err + } + return "", nil } // SearchMemory searches memory entries using FTS5 and bumps access counters diff --git a/internal/memory/memory_test.go b/internal/memory/memory_test.go index d2adcbcb..2ab8adfa 100644 --- a/internal/memory/memory_test.go +++ b/internal/memory/memory_test.go @@ -769,12 +769,13 @@ func TestDB_SearchMessagesVectorNoEmbedder(t *testing.T) { } func TestDB_SQLiteDSN(t *testing.T) { - want := "/tmp/a.db?_pragma=busy_timeout(5000)&_pragma=foreign_keys(1)" + want := "/tmp/a.db?_pragma=busy_timeout(5000)&_pragma=foreign_keys(1)&_txlock=immediate" if got := sqliteDSN("/tmp/a.db"); got != want { t.Errorf("sqliteDSN(plain) = %q, want %q", got, want) } - if got := sqliteDSN("/tmp/a.db?cache=shared"); got != "/tmp/a.db?cache=shared" { - t.Errorf("sqliteDSN(existing query) = %q", got) + wantQ := "/tmp/a.db?cache=shared&_pragma=busy_timeout(5000)&_pragma=foreign_keys(1)&_txlock=immediate" + if got := sqliteDSN("/tmp/a.db?cache=shared"); got != wantQ { + t.Errorf("sqliteDSN(existing query) = %q, want %q", got, wantQ) } } @@ -846,6 +847,39 @@ func TestDB_AddMessage_IdempotentOnSamePosition(t *testing.T) { } } +func TestDB_AddMessage_ConflictWithDifferentContent(t *testing.T) { + tmp := t.TempDir() + db, err := Open(filepath.Join(tmp, "test.db")) + if err != nil { + t.Fatal(err) + } + defer db.Close() + + if err := db.CreateSession(Session{ID: "sess-1", StartedAt: 1}); err != nil { + t.Fatal(err) + } + + if err := db.AddMessage(Message{SessionID: "sess-1", Idx: 0, Role: "user", Content: "hello", Timestamp: 1, ID: "id-a"}); err != nil { + t.Fatalf("first AddMessage: %v", err) + } + // A different message at the same position must not be silently dropped. + err = db.AddMessage(Message{SessionID: "sess-1", Idx: 0, Role: "user", Content: "world", Timestamp: 2, ID: "id-b"}) + if err == nil { + t.Fatal("expected conflict error for different content at the same position") + } + + msgs, err := db.GetMessages("sess-1") + if err != nil { + t.Fatal(err) + } + if len(msgs) != 1 { + t.Fatalf("expected 1 message, got %d", len(msgs)) + } + if msgs[0].Content != "hello" { + t.Errorf("expected first write to win, got %q", msgs[0].Content) + } +} + func TestDB_AddMemoryDedup_ExactOnly(t *testing.T) { tmp := t.TempDir() db, err := Open(filepath.Join(tmp, "test.db")) diff --git a/internal/memory/message_repo.go b/internal/memory/message_repo.go index 35d5537e..000b2be1 100644 --- a/internal/memory/message_repo.go +++ b/internal/memory/message_repo.go @@ -1,6 +1,9 @@ package memory -import "context" +import ( + "context" + "fmt" +) // AddMessage inserts a message into a session. When an embedder is // configured and the role is "user" or "assistant", the content is @@ -14,10 +17,26 @@ func (d *DB) AddMessage(m Message) error { if err != nil { return err } - if n, err := res.RowsAffected(); err == nil && n == 0 { + n, err := res.RowsAffected() + if err != nil { + return err + } + if n == 1 { + d.embedMessageAsync(m.ID, m.Role, m.Content) return nil } - d.embedMessageAsync(m.ID, m.Role, m.Content) + + // The position is already occupied. Treat it as an idempotent retry only + // when the stored row carries the same content fingerprint (the + // deterministic ID covers all immutable fields); otherwise it is a + // different message and must not be silently dropped. + var existingID string + if err := d.sql.QueryRow(`SELECT id FROM messages WHERE session_id = ? AND idx = ?`, m.SessionID, m.Idx).Scan(&existingID); err != nil { + return err + } + if existingID != m.ID { + return fmt.Errorf("message conflict at (%s, %d): stored content differs", m.SessionID, m.Idx) + } return nil }