Background ConsolidatorWorker dispatches consolidation work to a Claude Code agent through harness.Harness.Execute (NOT via system DMs — per feedback_system_dm_no_trigger.md) with a one-time 15m dispatch token. The dispatched agent uses six new MCP tools, all token-gated and recording every action to memory_consolidation_jobs.actions JSON. Stores: - memory_links.go (+ test): typed edges with actor-prefix reserved-type guard; AddConsolidationLink bypass for memory_mark_duplicate / memory_supersede (their contractual writers). - memory_pins.go (+ test): owner pin overlay, bypasses score floor. - memory_status.go: queries the memory_status view. - consolidation_jobs.go: Create / Dispatch / Lease / AppendAction / Complete with ErrJobAlreadyInFlight via partial unique index. - auto_links.go: MessageListener generating mention / reply_to / channel_cooccurrence links automatically on send. Worker: - consolidator.go (+ test): ticker pattern modeled on StalemateWorker. Watermark trigger for link_gen / dedup_contradiction; daily 03:00 UTC for sleep-time core rewrite. Wallclock budget via harness Budget. Global semaphore gates concurrent owners. Mocked-harness test asserts no system DM is ever sent. - consolidator_prompts.go: four job-type prompts passed via env to the dispatched agent. MCP tools (internal/mcp/memory_tools.go + test): - memory_list_unprocessed, memory_write_reflection, memory_rewrite_core, memory_mark_duplicate, memory_supersede, memory_add_link. - Full error-code matrix tested per contracts/mcp-memory-tools.md. - Registered only when SYNAPBUS_DREAM_ENABLED=1. Injection extensions: - search/injection.go: pin overlay applied after retrieval; status filter drops soft_deleted / superseded unless pinned. New PinProvider, StatusProvider, MessageLookup hooks on InjectionOpts. Wiring: - cmd/synapbus/main.go: stores constructed, AutoLinkListener attached to MessagingService, mcpSrv.SetDream wired, ConsolidatorWorker start/stop, admin DreamRun closure. - cmd/synapbus/admin.go: synapbus memory dream-run --owner --job socket-RPC command (forces a single job bypassing trigger). Cycle workarounds (documented in code): - messaging.DreamAgent / HarnessDispatcher are local interfaces (the agents and harness packages import messaging, not the reverse). main.go wraps the real types via adapter structs. Stubbed: - Cron expression parsing (DreamDeepCron). Hardcoded daily 03:00 UTC. Adding robfig/cron deferred to keep no-new-deps. Pre-existing reactor test failures (5) are unchanged; confirmed pre-020 via stash check. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
103 lines
3.0 KiB
Go
103 lines
3.0 KiB
Go
// memory_status query helpers — read the SQL view defined in migration
|
|
// 028 and turn it into a map[message_id]MemoryStatus suitable for the
|
|
// injection retrieval filter. The view itself derives state from
|
|
// `memory_consolidation_jobs.actions` rows (data-model.md §`memory_status`)
|
|
// so callers never have to mutate a status column directly.
|
|
package messaging
|
|
|
|
import (
|
|
"context"
|
|
"database/sql"
|
|
"fmt"
|
|
"strings"
|
|
"time"
|
|
)
|
|
|
|
// Memory status constants — value of `memory_status.status`.
|
|
const (
|
|
MemoryStatusActive = "active"
|
|
MemoryStatusSoftDeleted = "soft_deleted"
|
|
MemoryStatusSuperseded = "superseded"
|
|
)
|
|
|
|
// MemoryStatus is one row derived from the `memory_status` view.
|
|
// SupersededBy / SoftDeletedAt are nil when the message is active.
|
|
type MemoryStatus struct {
|
|
Status string `json:"status"`
|
|
SupersededBy *int64 `json:"superseded_by,omitempty"`
|
|
SoftDeletedAt *time.Time `json:"soft_deleted_at,omitempty"`
|
|
Reason string `json:"reason,omitempty"`
|
|
}
|
|
|
|
// MemoryStatuses returns the status of each message id in msgIDs. Ids
|
|
// absent from the result map are implicitly `active` (the view only
|
|
// contains rows that have at least one consolidation action against
|
|
// them).
|
|
func MemoryStatuses(ctx context.Context, db *sql.DB, msgIDs []int64) (map[int64]MemoryStatus, error) {
|
|
out := map[int64]MemoryStatus{}
|
|
if len(msgIDs) == 0 {
|
|
return out, nil
|
|
}
|
|
if db == nil {
|
|
return out, fmt.Errorf("memory status: nil db")
|
|
}
|
|
|
|
placeholders := strings.Repeat("?,", len(msgIDs))
|
|
placeholders = placeholders[:len(placeholders)-1]
|
|
args := make([]any, 0, len(msgIDs))
|
|
for _, id := range msgIDs {
|
|
args = append(args, id)
|
|
}
|
|
|
|
q := `SELECT message_id, status, superseded_by, soft_deleted_at, COALESCE(reason, '')
|
|
FROM memory_status
|
|
WHERE message_id IN (` + placeholders + `)`
|
|
rows, err := db.QueryContext(ctx, q, args...)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("memory status: query: %w", err)
|
|
}
|
|
defer rows.Close()
|
|
|
|
for rows.Next() {
|
|
var (
|
|
id int64
|
|
status string
|
|
supersede sql.NullInt64
|
|
deletedAt sql.NullTime
|
|
reason string
|
|
)
|
|
if err := rows.Scan(&id, &status, &supersede, &deletedAt, &reason); err != nil {
|
|
return nil, fmt.Errorf("memory status: scan: %w", err)
|
|
}
|
|
ms := MemoryStatus{Status: status, Reason: reason}
|
|
if supersede.Valid {
|
|
v := supersede.Int64
|
|
ms.SupersededBy = &v
|
|
}
|
|
if deletedAt.Valid {
|
|
t := deletedAt.Time
|
|
ms.SoftDeletedAt = &t
|
|
}
|
|
out[id] = ms
|
|
}
|
|
if err := rows.Err(); err != nil {
|
|
return nil, fmt.Errorf("memory status: iterate: %w", err)
|
|
}
|
|
return out, nil
|
|
}
|
|
|
|
// StatusByID returns just the string status for each id; callers that
|
|
// only need active/non-active (e.g. the injection retrieval filter) can
|
|
// use this to avoid pulling in MemoryStatus's optional fields.
|
|
func StatusByID(ctx context.Context, db *sql.DB, msgIDs []int64) (map[int64]string, error) {
|
|
full, err := MemoryStatuses(ctx, db, msgIDs)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
out := make(map[int64]string, len(full))
|
|
for id, st := range full {
|
|
out[id] = st.Status
|
|
}
|
|
return out, nil
|
|
}
|