Three things that decide whether memory improves with use or rots with it. A RELEVANCE FLOOR. Recall returned its top five whatever they scored, so a run about shift cover was handed five memories about certifications simply because nothing better existed — and the block tells the model these are things the workspace remembered, so it reads them as pertinent. Embeddings are unit-normalised, so knowledge_dot is cosine, and 0.30 is where text is usually about something else. A judgement rather than a measurement, and the honest way to tune it is to watch what gets carried on real questions. NO DUPLICATES. 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. Five recall slots spent on one fact restated five ways is the normal failure, not a rare one. A write with the same normalised text, in the same org and about the same subject, pushes the existing memory's expiry out instead of adding a row — matched on the same sentence rather than a similar one, because collapsing two genuinely different facts is the worse error. EXPIRY THAT DELETES. expires_at was set and filtered on read, and nothing ever removed anything: the row was invisible and still retained. "We keep it ninety days" has to be true of the table, not only of the query. Prune is batched, and a redaction is kept for a thirty-day grace period so an erasure stays provable shortly afterwards. It runs in the maintenance sweeper that already exists rather than a second scheduler — same ticker, same cancellation, same failure isolation. That forced one honest change: Maintenance() used to be nil without OAuth, on the reasoning that there was nothing to sweep. There is now, and a retention promise enforced only when an unrelated feature happens to be enabled is not a promise. The test that asserted the old behaviour now asserts the new one and says why. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
223 lines
8.0 KiB
Go
223 lines
8.0 KiB
Go
package httpserver
|
|
|
|
import (
|
|
"context"
|
|
"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"
|
|
)
|
|
|
|
// Scheduled maintenance for the OAuth and rate-limit tables.
|
|
//
|
|
// WHY THIS SHAPE AND NOT A NEW ONE
|
|
//
|
|
// The process already has a scheduled maintenance mechanism: sweepSessions in
|
|
// cmd/api/main.go, a ticker goroutine whose context is the server's, which runs
|
|
// once at startup and then on an interval, logs a failure and retries at the
|
|
// next tick. It is bounded, cancellable, non-blocking and failure-isolated, and
|
|
// it has been in production.
|
|
//
|
|
// So this is the same thing for two more tables rather than a second kind of
|
|
// thing. No new process, no cron dependency, no leader election, no library.
|
|
// The one addition is that both sweeps live behind a single type, so
|
|
// cmd/api/main.go gains one line rather than two more goroutines.
|
|
//
|
|
// MULTI-INSTANCE SAFETY COMES FROM THE STATEMENTS, NOT FROM COORDINATION
|
|
//
|
|
// Every instance runs this, on its own schedule, with no lock between them —
|
|
// deliberately. A lease or an advisory lock would be state to hold, to expire
|
|
// and to recover when the holder dies mid-sweep, in exchange for avoiding work
|
|
// that is already harmless: each sweep is a bounded DELETE whose predicate no
|
|
// longer matches once a row is gone. Two instances sweeping at the same moment
|
|
// delete disjoint sets and neither errors. A row deleted twice is not an error;
|
|
// it is a row that was already deleted.
|
|
//
|
|
// That is the same property Phase 5's concurrent-cleanup test asserts directly:
|
|
// four workers, six dead tokens, exactly six removed between them.
|
|
|
|
// maintenanceInterval is how often the sweep runs.
|
|
//
|
|
// Hourly. The grace period before anything is deleted is also an hour, so a
|
|
// row becomes eligible and is collected within roughly two — soon enough that
|
|
// nothing accumulates, and far enough apart that a DELETE never lands on a hot
|
|
// path. Shorter would buy nothing: nothing here is a correctness deadline.
|
|
//
|
|
// Deliberately NOT sweepInterval's fifteen minutes. Sessions churn with every
|
|
// sign-in; authorization codes live sixty seconds and tokens fifteen minutes,
|
|
// so an hour still collects them promptly while running a quarter as often.
|
|
const maintenanceInterval = time.Hour
|
|
|
|
// maintenanceTimeout bounds one pass.
|
|
//
|
|
// Generous for three bounded deletes and short enough that a wedged statement
|
|
// cannot hold this goroutine past shutdown. Matches sweepSessions' own bound in
|
|
// spirit; longer only because there are more statements.
|
|
const maintenanceTimeout = 60 * time.Second
|
|
|
|
// Maintenance sweeps the OAuth and rate-limit tables.
|
|
//
|
|
// Nil when the deployment does not serve MCP, which is why Server.Maintenance
|
|
// returns a pointer and the caller checks it — the same way routeOAuth simply
|
|
// registers nothing.
|
|
type Maintenance struct {
|
|
store *oauth.Store
|
|
limiter *ratelimit.Limiter
|
|
|
|
// 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.
|
|
//
|
|
// 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 {
|
|
/* 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
|
|
}
|
|
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.
|
|
type MaintenanceResult struct {
|
|
Grants int64
|
|
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 + r.Memories
|
|
}
|
|
|
|
// Sweep runs one maintenance pass.
|
|
//
|
|
// The two halves are independent on purpose: a failure sweeping OAuth rows must
|
|
// not prevent the rate-limit sweep, because the second is the one that would
|
|
// otherwise grow without bound. The first error is returned, after both have
|
|
// been attempted.
|
|
func (m *Maintenance) Sweep(ctx context.Context) (MaintenanceResult, error) {
|
|
var out MaintenanceResult
|
|
var firstErr 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.
|
|
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 {
|
|
swept, err := m.limiter.Sweep(ctx, 0) // 0 = the package's own batch size
|
|
if err != nil && firstErr == nil {
|
|
firstErr = err
|
|
}
|
|
out.RateLimits = swept
|
|
}
|
|
|
|
return out, firstErr
|
|
}
|
|
|
|
// SweepMaintenance runs the sweep until the context is cancelled.
|
|
//
|
|
// Deliberately identical in shape to sweepSessions: one pass immediately so a
|
|
// process that has been down does not carry a backlog for a further hour, then
|
|
// on the ticker. A failed pass is logged and retried at the next tick — the
|
|
// tables being briefly larger than they should be is not worth stopping the API
|
|
// for, and it is certainly not worth a panic in a goroutine nobody is watching.
|
|
//
|
|
// Exported because cmd/api owns the process's goroutines and this package owns
|
|
// what they do.
|
|
func SweepMaintenance(ctx context.Context, m *Maintenance, log *slog.Logger) {
|
|
if m == nil {
|
|
// No OAuth surface, nothing to sweep. Returning rather than ticking
|
|
// uselessly for the life of the process.
|
|
return
|
|
}
|
|
|
|
ticker := time.NewTicker(maintenanceInterval)
|
|
defer ticker.Stop()
|
|
|
|
pass := func() {
|
|
// A deadline of its own, so a slow DELETE cannot leave this goroutine
|
|
// blocked past shutdown.
|
|
sweepCtx, cancel := context.WithTimeout(ctx, maintenanceTimeout)
|
|
defer cancel()
|
|
|
|
// A panic in a background goroutine takes the process with it, and
|
|
// this one runs unattended for the life of the deployment. Recovering
|
|
// turns a bug here into a logged failure and a retry at the next tick.
|
|
defer func() {
|
|
if p := recover(); p != nil {
|
|
log.Error("maintenance sweep panicked", "panic", p)
|
|
}
|
|
}()
|
|
|
|
result, err := m.Sweep(sweepCtx)
|
|
switch {
|
|
case err != nil && ctx.Err() != nil:
|
|
// Shutting down; the cancellation is expected, not a failure.
|
|
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,
|
|
"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,
|
|
"memories", result.Memories)
|
|
default:
|
|
log.Debug("maintenance sweep found nothing to delete")
|
|
}
|
|
}
|
|
|
|
pass()
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
log.Debug("maintenance sweeper stopped")
|
|
return
|
|
case <-ticker.C:
|
|
pass()
|
|
}
|
|
}
|
|
}
|