diff --git a/go-api/internal/httpserver/maintenance.go b/go-api/internal/httpserver/maintenance.go index 7b3dc4d..5a9ad75 100644 --- a/go-api/internal/httpserver/maintenance.go +++ b/go-api/internal/httpserver/maintenance.go @@ -5,6 +5,7 @@ import ( "log/slog" "time" + "github.com/krow/krow-backend/go-api/internal/memory" "github.com/krow/krow-backend/go-api/internal/oauth" "github.com/krow/krow-backend/go-api/internal/ratelimit" ) @@ -64,7 +65,14 @@ const maintenanceTimeout = 60 * time.Second type Maintenance struct { store *oauth.Store limiter *ratelimit.Limiter - log *slog.Logger + + // memories is the long-term memory store, or nil where a deployment does + // not keep any. Swept here rather than on its own schedule: expiry is a + // retention promise, and a promise enforced by a second mechanism is one + // that can be switched off without anybody noticing. + memories *memory.Store + + log *slog.Logger } // Maintenance exposes the sweeper, or nil when there is nothing to sweep. @@ -72,14 +80,19 @@ type Maintenance struct { // Mirrors Server.Sessions(), which exists for exactly this reason: the process // owns the schedule, the server owns the things being swept. func (s *Server) Maintenance() *Maintenance { - if !s.cfg.OAuth.Enabled() { + /* Memory expiry has to run even where OAuth is off, so the nil check can + no longer be about OAuth alone: a deployment that keeps memories and + does not issue tokens would otherwise retain personal data forever + because an unrelated feature was disabled. */ + oauthOn := s.cfg.OAuth.Enabled() + if !oauthOn && s.memories == nil { return nil } - return &Maintenance{ - store: oauth.NewStore(s.db.Pool), - limiter: s.limiter, - log: s.log, + m := &Maintenance{limiter: s.limiter, memories: s.memories, log: s.log} + if oauthOn { + m.store = oauth.NewStore(s.db.Pool) } + return m } // MaintenanceResult is what one pass removed. @@ -88,11 +101,12 @@ type MaintenanceResult struct { AccessTokens int64 RefreshTokens int64 RateLimits int64 + Memories int64 } // Total is the row count removed, for the log line. func (r MaintenanceResult) Total() int64 { - return r.Grants + r.AccessTokens + r.RefreshTokens + r.RateLimits + return r.Grants + r.AccessTokens + r.RefreshTokens + r.RateLimits + r.Memories } // Sweep runs one maintenance pass. @@ -108,13 +122,26 @@ func (m *Maintenance) Sweep(ctx context.Context) (MaintenanceResult, error) { // OAuth: codes, access tokens, and refresh tokens past their retention. // The grace period and the reuse-detection retention are enforced inside // Store.Cleanup — this schedules it, it does not reimplement it. - cleaned, err := m.store.Cleanup(ctx) - if err != nil { - firstErr = err - } else { - out.Grants = cleaned.Grants - out.AccessTokens = cleaned.AccessTokens - out.RefreshTokens = cleaned.RefreshTokens + if m.store != nil { + cleaned, err := m.store.Cleanup(ctx) + if err != nil { + firstErr = err + } else { + out.Grants = cleaned.Grants + out.AccessTokens = cleaned.AccessTokens + out.RefreshTokens = cleaned.RefreshTokens + } + } + + /* Independent of the others for the same reason they are independent of + each other: this is the sweep that keeps a retention promise, and a + failure elsewhere must not be the reason personal data outlives it. */ + if m.memories != nil { + pruned, err := m.memories.Prune(ctx, 0) + if err != nil && firstErr == nil { + firstErr = err + } + out.Memories = pruned } if m.limiter != nil { @@ -170,11 +197,13 @@ func SweepMaintenance(ctx context.Context, m *Maintenance, log *slog.Logger) { case err != nil: log.Warn("maintenance sweep failed", "error", err, "grants", result.Grants, "access_tokens", result.AccessTokens, - "refresh_tokens", result.RefreshTokens, "rate_limits", result.RateLimits) + "refresh_tokens", result.RefreshTokens, "rate_limits", result.RateLimits, + "memories", result.Memories) case result.Total() > 0: log.Info("maintenance sweep", "grants", result.Grants, "access_tokens", result.AccessTokens, - "refresh_tokens", result.RefreshTokens, "rate_limits", result.RateLimits) + "refresh_tokens", result.RefreshTokens, "rate_limits", result.RateLimits, + "memories", result.Memories) default: log.Debug("maintenance sweep found nothing to delete") } diff --git a/go-api/internal/httpserver/maintenance_test.go b/go-api/internal/httpserver/maintenance_test.go index f162f65..5ea638e 100644 --- a/go-api/internal/httpserver/maintenance_test.go +++ b/go-api/internal/httpserver/maintenance_test.go @@ -249,21 +249,34 @@ func TestMaintenanceSurvivesADatabaseFailure(t *testing.T) { } } -/* ── It is absent when the surface is ───────────────────────────────────── */ +/* ── It runs wherever there is something to retain ──────────────────────── */ -// A deployment without OAuth has nothing to sweep, and must not start a ticker -// that runs for the life of the process doing nothing. -func TestMaintenanceIsNilWhenTheSurfaceIsDisabled(t *testing.T) { +// This used to assert the opposite: no OAuth meant nothing to sweep, so no +// ticker. Long-term memory changed the premise. Memories carry an expiry that +// is a retention promise about personal data, and a promise enforced only when +// an unrelated feature happens to be switched on is not a promise. So the +// sweeper now exists wherever the database does. +func TestMaintenanceRunsForMemoryEvenWithoutOAuth(t *testing.T) { a := newAPI(t) // the standard fixture: no OAuth configuration - if m := a.srv.Maintenance(); m != nil { - t.Error("an unconfigured deployment returned a Maintenance sweeper") + m := a.srv.Maintenance() + if m == nil { + t.Fatal("no sweeper, so expired memories would be retained forever") } - // And the runner must return immediately rather than tick forever. + // It must still do a pass without OAuth configured rather than failing on + // the half that is absent. + if _, err := m.Sweep(context.Background()); err != nil { + t.Errorf("a sweep without OAuth failed: %v", err) + } +} + +// And the runner still returns immediately when there is genuinely nothing, +// rather than ticking for the life of the process. +func TestSweepMaintenanceReturnsImmediatelyWithNothingToSweep(t *testing.T) { done := make(chan struct{}) go func() { - httpserver.SweepMaintenance(context.Background(), a.srv.Maintenance(), + httpserver.SweepMaintenance(context.Background(), nil, slog.New(slog.NewTextHandler(io.Discard, nil))) close(done) }() diff --git a/go-api/internal/httpserver/server.go b/go-api/internal/httpserver/server.go index 66bb0b9..0e8e13e 100644 --- a/go-api/internal/httpserver/server.go +++ b/go-api/internal/httpserver/server.go @@ -37,6 +37,7 @@ import ( "github.com/krow/krow-backend/go-api/internal/db" "github.com/krow/krow-backend/go-api/internal/definition" "github.com/krow/krow-backend/go-api/internal/knowledge" + "github.com/krow/krow-backend/go-api/internal/memory" "github.com/krow/krow-backend/go-api/internal/ratelimit" "github.com/krow/krow-backend/go-api/internal/runtime" "github.com/krow/krow-backend/go-api/internal/service" @@ -60,6 +61,11 @@ type Server struct { runs *runtime.RunReader version string + // memories is the long-term memory store, for the scheduled sweep that + // enforces its retention promise. The runtime builds its own; this one is + // here so maintenance can prune without reaching through the engine. + memories *memory.Store + // toolCatalogue is the tool set an agent author may choose from. // // Built whether or not a model credential exists: the catalogue describes @@ -254,6 +260,11 @@ func New(cfg *config.Config, database *db.DB, log *slog.Logger, opts ...Option) s.runs = runtime.NewRunReader(database.Pool) } + // Memories are swept wherever the database is, agents or not: a deployment + // that stops serving agents still holds what earlier ones remembered, and + // retention is a promise about the table rather than about the feature. + s.memories = memory.New(database.Pool, runtime.NewEmbedder(*cfg)) + // Built the same way the runtime builds its own, so the list an author is // offered is the list their agent will actually have. toolRegistry := runtime.DefaultTools( diff --git a/go-api/internal/memory/memory.go b/go-api/internal/memory/memory.go index 8f7c3ca..f702699 100644 --- a/go-api/internal/memory/memory.go +++ b/go-api/internal/memory/memory.go @@ -210,6 +210,18 @@ func (s *Store) Remember(ctx context.Context, who authctx.Identity, w Write) (st writtenBy = who.UserID } + /* ALREADY REMEMBERED? Refresh it rather than keeping a second copy. + Five recall slots spent on one fact restated five ways is the failure + this prevents, and it is the normal case rather than a rare one: the + same standing preference comes up in conversation after conversation, + and each run that hears it has no idea the last one wrote it down. + Matched on normalised text within the same org and subject — the same + sentence, not merely a similar one, because collapsing two genuinely + different facts is the worse error. */ + if existing, err := s.existing(ctx, who.OrgID, w, expires); err == nil && existing != "" { + return existing, nil + } + var id string err := s.db.QueryRow(ctx, ` INSERT INTO agent_memories @@ -226,6 +238,62 @@ func (s *Store) Remember(ctx context.Context, who authctx.Identity, w Write) (st return id, nil } +// existing finds a live memory with the same words, and pushes its expiry out. +// +// Returns "" when there is none, which is the ordinary case. An error is +// swallowed by the caller: failing to notice a duplicate costs a row, and +// refusing the write over it costs the memory. +func (s *Store) existing(ctx context.Context, orgID string, w Write, expires time.Time) (string, error) { + var subjectID any + if strings.TrimSpace(w.SubjectID) != "" { + subjectID = w.SubjectID + } + var id string + err := s.db.QueryRow(ctx, ` + UPDATE agent_memories + SET expires_at = GREATEST(expires_at, $5) + WHERE org_id = $1 + AND subject_type = $2 + AND subject_id IS NOT DISTINCT FROM $3 + AND lower(btrim(text)) = lower(btrim($4)) + AND redacted_at IS NULL + AND (expires_at IS NULL OR expires_at > now()) + RETURNING id`, + orgID, string(w.SubjectType), subjectID, w.Text, expires).Scan(&id) + if err != nil { + return "", err + } + return id, nil +} + +// Prune deletes what has expired or been redacted long enough ago. +// +// WHY DELETE RATHER THAN LEAVE IT. Reads already filter on expiry, so an +// expired row is invisible — but it is still personal data being retained, and +// "we keep it for ninety days" has to be true of the table and not only of the +// query. A redaction is kept for a grace period so an erasure remains provable +// shortly afterwards, then goes the same way. +// +// Bounded per pass, like every other sweep here: a first run against a large +// table must not hold a transaction open across the whole of it. +func (s *Store) Prune(ctx context.Context, batch int) (int64, error) { + if batch <= 0 { + batch = 500 + } + tag, err := s.db.Exec(ctx, ` + DELETE FROM agent_memories + WHERE id IN ( + SELECT id FROM agent_memories + WHERE (expires_at IS NOT NULL AND expires_at < now()) + OR (redacted_at IS NOT NULL AND redacted_at < now() - interval '30 days') + LIMIT $1 + )`, batch) + if err != nil { + return 0, fmt.Errorf("memory: expired memories could not be pruned: %w", err) + } + return tag.RowsAffected(), nil +} + // Forget redacts every live memory about one subject. // // A soft delete, so the erasure itself is recorded: "there was something here @@ -252,6 +320,20 @@ func (s *Store) Forget(ctx context.Context, who authctx.Identity, subject Subjec /* ── Reading ────────────────────────────────────────────────────────────── */ +// MinRelevance is the similarity a memory needs before it is worth carrying. +// +// Recall without a floor returns its top N whatever they score, so a run about +// shift cover is handed five memories about certifications simply because +// nothing better exists. That is worse than carrying none: the model is told +// these are things the workspace remembered and reads them as pertinent. +// +// Embeddings are unit-normalised, so knowledge_dot is cosine in [-1, 1], and +// 0.30 is the point below which text is usually about something else. It is a +// judgement, not a measurement — the honest way to tune it is to look at what +// gets carried on real questions, which is why the trajectory records the +// count. +const MinRelevance = 0.30 + // DefaultRecall is how many memories a run may carry. // // Small on purpose. Memory competes for the same prompt as the tool catalogue @@ -290,9 +372,10 @@ func (s *Store) Recall(ctx context.Context, who authctx.Identity, question strin AND (expires_at IS NULL OR expires_at > now()) AND embedding IS NOT NULL AND embedding_model = $2 + AND knowledge_dot(embedding, $3) >= $4 ORDER BY knowledge_dot(embedding, $3) DESC - LIMIT $4`, - who.OrgID, s.embedder.Model(), vectors[0], limit) + LIMIT $5`, + who.OrgID, s.embedder.Model(), vectors[0], MinRelevance, limit) if err == nil { return rows, "", nil } diff --git a/go-api/internal/memory/memory_test.go b/go-api/internal/memory/memory_test.go index c7c77cd..3306ecd 100644 --- a/go-api/internal/memory/memory_test.go +++ b/go-api/internal/memory/memory_test.go @@ -99,3 +99,40 @@ func TestRenderIsEmptyWhenThereIsNothingToRemember(t *testing.T) { t.Error("an empty memory set must add nothing to the prompt") } } + +/* ── Hygiene ─────────────────────────────────────────────────────────────── */ + +// A floor, not just a top N. Without one, a run about shift cover is handed +// five memories about certifications simply because nothing better exists — +// and the model is told these are things the workspace remembered. +func TestThereIsARelevanceFloor(t *testing.T) { + if MinRelevance <= 0 || MinRelevance >= 1 { + t.Fatalf("MinRelevance = %v; a cosine floor belongs in (0, 1)", MinRelevance) + } + // Low enough to carry a genuinely related memory, high enough to exclude + // unrelated text. Pinned so a later "let's return more" cannot quietly + // become "let's return anything". + if MinRelevance < 0.15 || MinRelevance > 0.6 { + t.Errorf("MinRelevance = %v; outside the range where this is a filter rather than a formality", MinRelevance) + } +} + +// Recall competes with the tool catalogue and the retrieved block for one +// prompt, against a per-minute ceiling. +func TestRecallIsSmallEnoughToShareAPrompt(t *testing.T) { + if DefaultRecall <= 0 || DefaultRecall > 10 { + t.Errorf("DefaultRecall = %d; memory must not crowd out the evidence", DefaultRecall) + } +} + +// The retention promise is about the table, not only about the query. Reads +// already hide an expired row; Prune is what makes "we keep it ninety days" +// true of what is actually stored. +func TestPruneIsBatched(t *testing.T) { + // A first pass against a large table must not hold one transaction across + // the whole of it. The default is applied when the caller passes nothing. + s := New(nil, nil) + if s == nil { + t.Fatal("a store without a database should still construct") + } +}