Files
krow_backend/go-api/internal/httpserver/server.go
Aravind e90bc33d0f
Some checks failed
CI / test (push) Failing after 4m40s
CI / fixture (push) Failing after 8s
Memory hygiene: a relevance floor, no duplicates, and expiry that actually deletes
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>
2026-10-07 20:29:28 +05:30

621 lines
23 KiB
Go

// Package httpserver holds the HTTP surface.
//
// It serves /health, the sign-in endpoints, the entity endpoints described in
// docs/api-contract.md, the current-user endpoints, the agent and skill
// definition endpoints, and the Owliver panel's suggestion endpoint.
//
// Authentication replaced the development identity: every request outside the
// small public allowlist in auth.go must carry a session cookie; the middleware
// resolves it to a user row and puts that user, and their organization, on the
// request context. Nothing downstream changed — every service and repository
// already took the organization as a parameter, which is what devOrgMiddleware
// existed to make true.
//
// Authorization is here, in Server.authorize: it reads the role off the
// authenticated identity, consults the deny-by-default policy table in
// internal/domain/policy.go, and answers 403 before any query runs. Row
// visibility — organization scope, and ownership for talent callers — is a SQL
// predicate in internal/repo instead, so an invisible row answers 404 rather
// than 403. The definition endpoints are the exception: they are not
// domain.Resource values, so their role checks are written inline in
// internal/service/definitions.go rather than in the policy table.
package httpserver
import (
"context"
"encoding/json"
"errors"
"fmt"
"log/slog"
"net"
"net/http"
"strconv"
"time"
"github.com/krow/krow-backend/go-api/internal/auth"
"github.com/krow/krow-backend/go-api/internal/config"
"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"
"github.com/krow/krow-backend/go-api/internal/tools"
)
// Server binds the router, the pool, authentication and the lifecycle together.
type Server struct {
cfg *config.Config
db *db.DB
api *service.Registry
definitions *service.DefinitionsService
workflows *service.WorkflowService
suggestions *service.SuggestionsService
// The agent runtime. Nil when no model credential is configured — the run
// routes are then not registered at all, so the deployment answers 404
// ("this deployment does not serve agents") rather than 500 ("this
// deployment is broken"). Only one of those is true.
agents *runtime.Engine
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
// what the tools ARE, and a deployment that cannot currently run agents can
// still be one where somebody is authoring them.
toolCatalogue []tools.ToolInfo
log *slog.Logger
http *http.Server
started time.Time
endpoints int
// The authentication surface. sessions owns the lifecycle, users is the
// read side of the users table, credentials verifies a password against it,
// and the two limiters bound how often that may be attempted.
//
// Two limiters, not one, because the budgets are different sizes on
// purpose: an email is one account and gets a tight budget, while an
// address may be a whole office behind NAT and gets a loose one. Sharing a
// limiter would force the office to live within one person's budget.
sessions *auth.Manager
users auth.UserStore
credentials *auth.Credentials
loginByEmail *attemptLimiter
// limiter bounds the OAuth and MCP routes, shared across instances via
// Postgres. Nil when those routes are not registered — see mcplimit.go,
// where a nil limiter means the middleware is not installed at all rather
// than installed and permissive.
limiter *ratelimit.Limiter
loginByAddr *attemptLimiter
// trust resolves a request to the address its per-address limits are keyed
// by, reading a forwarded address only from a configured proxy. Wired once
// here so no handler can be given a different notion of who called.
trust proxyTrust
// now is injectable so tests can drive expiry without sleeping.
now func() time.Time
}
// Option adjusts the server before it is wired. Production passes none.
type Option func(*serverOptions)
type serverOptions struct {
policy auth.Policy
now func() time.Time
perEmail int
perAddress int
loginWindow time.Duration
// agents replaces the engine New would otherwise build from configuration.
//
// For tests, and only for tests: production wires a real gateway from a
// real key, and an option that let a deployment substitute the runtime
// would be a way to run agents against something nobody configured.
agents *runtime.Engine
// version is the build identifier, stamped into the binary at link time.
// Not configuration: it describes the artefact, not the deployment, and an
// environment variable could disagree with the code it claims to describe.
version string
// curatedAgents replaces the set New would otherwise read from disk.
// nil means "read the configured directory"; an empty non-nil set means
// "protect nothing", which is a thing a test needs to be able to say.
curatedAgents map[string]bool
}
// WithBuildVersion records which build this is.
//
// Unlike the options above this one is for production. Without it there is no
// way to answer "did my deploy land?" — the symptom is pushing an image,
// redeploying, and having nobody, including the operator, able to tell whether
// the running process is the new one.
func WithBuildVersion(v string) Option {
return func(o *serverOptions) {
if v != "" {
o.version = v
}
}
}
// WithSessionPolicy overrides the session lifetimes. For tests that need to
// reach an expiry without waiting twelve hours for it.
func WithSessionPolicy(p auth.Policy) Option {
return func(o *serverOptions) { o.policy = p }
}
// WithClock replaces the clock used for session expiry and last_login_at.
func WithClock(now func() time.Time) Option {
return func(o *serverOptions) {
if now != nil {
o.now = now
}
}
}
// WithLoginRateLimit overrides the failed-attempt budgets and their window.
//
// perEmail bounds attempts against one account; perAddress bounds attempts from
// one client address across all accounts. Both are consulted on every attempt.
// WithAgentEngine substitutes the agent runtime.
//
// The seam that lets the HTTP layer be tested without a model credential —
// which matters more than it sounds, because the alternative is that the run
// endpoint is the one part of this service no test can reach until somebody
// pays for a key.
//
// It does not weaken anything: the engine still loads agents through the same
// loader, still runs them under the same budgets, and still authorizes through
// the same principal. Only the model behind it changes.
// WithCuratedAgents names the delete-protected agent ids directly.
//
// Production loads these from disk; this exists so a test can state its own
// protected set without a directory, exactly as WithAgentEngine lets one
// supply an engine without a model credential.
func WithCuratedAgents(ids ...string) Option {
return func(o *serverOptions) {
set := make(map[string]bool, len(ids))
for _, id := range ids {
set[id] = true
}
o.curatedAgents = set
}
}
func WithAgentEngine(e *runtime.Engine) Option {
return func(o *serverOptions) { o.agents = e }
}
func WithLoginRateLimit(perEmail, perAddress int, window time.Duration) Option {
return func(o *serverOptions) {
o.perEmail, o.perAddress, o.loginWindow = perEmail, perAddress, window
}
}
// New wires the routes and returns a server that has not yet been started.
//
// Authentication is built here rather than passed in, so there is exactly one
// construction of the session manager and no way to start a server with the
// middleware wired to a different store than the login handler.
func New(cfg *config.Config, database *db.DB, log *slog.Logger, opts ...Option) (*Server, error) {
o := serverOptions{
policy: auth.DefaultPolicy,
now: time.Now,
perEmail: loginAttemptLimit,
perAddress: loginAddressLimit,
loginWindow: loginAttemptWindow,
version: "unknown",
}
for _, opt := range opts {
opt(&o)
}
sessions, err := auth.NewManager(auth.NewPGStore(database.Pool), o.policy)
if err != nil {
return nil, fmt.Errorf("build session manager: %w", err)
}
sessions.WithClock(o.now)
users := auth.NewPGUserStore(database.Pool)
s := &Server{
cfg: cfg, db: database, log: log,
version: o.version,
api: service.NewRegistry(database.Pool),
definitions: service.NewDefinitions(database.Pool),
workflows: service.NewWorkflows(database.Pool).WithClock(o.now),
suggestions: service.NewSuggestions(database.Pool),
started: o.now(),
sessions: sessions,
users: users,
credentials: auth.NewCredentials(users),
loginByEmail: newAttemptLimiter(o.perEmail, o.loginWindow, o.now),
loginByAddr: newAttemptLimiter(o.perAddress, o.loginWindow, o.now),
trust: newProxyTrust(cfg.HTTP.TrustedProxies),
now: o.now,
}
// The agent runtime, wired only when there is a model to reach.
//
// Registering the routes without a credential would accept runs and fail
// every one of them at the gateway — an outage shaped like a feature. A
// deployment without a key is a deployment that does not serve agents, and
// saying so at boot is kinder than saying it once per request.
switch {
case o.agents != nil:
s.agents = o.agents
s.runs = runtime.NewRunReader(database.Pool)
case cfg.Model.APIKey != "":
s.agents = runtime.NewModelEngine(database.Pool, *cfg)
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(
database.Pool,
knowledge.NewRetriever(database.Pool, runtime.NewEmbedder(*cfg)),
)
s.toolCatalogue = toolRegistry.Catalogue()
// So a definition naming a tool that does not exist is refused at publish
// rather than becoming an agent that silently cannot do what it claims.
s.definitions = s.definitions.WithToolCheck(toolRegistry.Known)
// The agents this deployment ships specs for, so DELETE refuses them at the
// endpoint rather than only in the list that renders the button.
//
// Read from the directory `importagents` publishes from, so the protected
// set is the published set by construction. A deployment without that
// directory protects nothing and says so here, once, at boot: silence would
// leave an operator believing in a guard that is not running.
curated := o.curatedAgents
if curated == nil {
loaded, err := definition.CuratedIDs(cfg.Agents.CuratedPath)
if err != nil {
return nil, fmt.Errorf("load curated agents: %w", err)
}
curated = loaded
}
if len(curated) == 0 {
log.Warn("no curated agent specs found; built-in agents are not delete-protected",
"path", cfg.Agents.CuratedPath)
} else {
log.Info("curated agents are delete-protected",
"count", len(curated), "ids", definition.SortedIDs(curated))
}
s.definitions = s.definitions.WithCuratedAgents(curated)
// The shared limiter, built only when the routes that use it exist. The
// existing in-process login limiter is untouched: it guards a different
// thing (failed password attempts) with a different model (count failures,
// reset on success), and replacing it is not this change's business.
if cfg.OAuth.Enabled() {
s.limiter = ratelimit.New(database.Pool)
}
mux := http.NewServeMux()
mux.HandleFunc("GET /health", s.handleHealth)
s.endpoints = s.routeAuth(mux) + s.routeResources(mux) + s.routeMe(mux) +
s.routeDefinitions(mux) + s.routeWorkflows(mux) + s.routeOwliver(mux) +
s.routeRuns(mux) + s.routeVersion(mux) + s.routeTools(mux) +
// The MCP surface and the OAuth server behind it. Both return 0 and
// register nothing when OAUTH_ISSUER and MCP_RESOURCE are unset, which
// is every deployment that has not asked for them.
s.routeOAuth(mux) + s.routeMCP(mux)
handler := jsonErrors(mux)
// Authentication sits where devOrgMiddleware used to, so every route below
// it — including the mux's own 404 — is behind the allowlist.
handler = s.authenticate(handler)
handler = recoverer(log)(handler)
// CORS sits outside the recoverer so a preflight is answered without
// touching the router, and inside the logger so refused origins are still
// visible in the log. With no allowlist configured it is not installed at
// all, which is the same-origin default.
if len(cfg.HTTP.CORSOrigins) > 0 {
handler = cors(cfg.HTTP.CORSOrigins)(handler)
}
handler = requestLogger(log)(handler)
s.http = &http.Server{
Addr: net.JoinHostPort(cfg.HTTP.Host, strconv.Itoa(cfg.HTTP.Port)),
Handler: handler,
ReadTimeout: cfg.HTTP.ReadTimeout,
WriteTimeout: cfg.HTTP.WriteTimeout,
IdleTimeout: cfg.HTTP.IdleTimeout,
}
return s, nil
}
// Sessions exposes the session manager, so the process can sweep expired rows
// and tests can drive the clock.
func (s *Server) Sessions() *auth.Manager { return s.sessions }
// Handler exposes the routed handler so tests can drive it without a listener.
func (s *Server) Handler() http.Handler { return s.http.Handler }
// Endpoints is how many routes were registered.
func (s *Server) Endpoints() int { return s.endpoints }
// Addr is the address the server listens on.
func (s *Server) Addr() string { return s.http.Addr }
// Start blocks until the server stops accepting connections.
func (s *Server) Start() error {
err := s.http.ListenAndServe()
if errors.Is(err, http.ErrServerClosed) {
return nil
}
return err
}
// Shutdown drains in-flight requests, then gives up after the configured grace.
func (s *Server) Shutdown(ctx context.Context) error {
ctx, cancel := context.WithTimeout(ctx, s.cfg.HTTP.ShutdownTimeout)
defer cancel()
return s.http.Shutdown(ctx)
}
// healthResponse is the entire public /health body: one field, deliberately.
//
// /health is unauthenticated and reachable by anyone who can reach the port,
// so it is treated as a public document rather than as an operator's console.
// Everything an unauthenticated caller legitimately needs is the answer to
// "should traffic be sent here", and that fits in a status string plus the
// HTTP status code.
//
// What used to be here and is now deliberately absent: the PostgreSQL version,
// the database name, the schema name, the applied migration version, the table
// count, the connection error text, the deployment environment and the process
// uptime. Individually each is small; together they are a free reconnaissance
// report — the server version to look up known CVEs against, the migration
// version to date the deployment, the table count and error text to infer
// shape and topology. None of it is diagnostic to anyone who could not already
// read it from the database directly.
//
// The check itself is unchanged. db.Check still runs on every request and
// still decides the answer; its full detail now goes to the server log, where
// the operator is, instead of into the response, where the internet is. See
// logHealth.
type healthResponse struct {
Status string `json:"status"`
}
// routeTools lists the tools an agent author may choose from.
//
// The frontend's agent editor had no tools field at all, so an authored agent
// carried none and could talk without being able to look anything up. Serving
// the catalogue rather than hard-coding it in the UI keeps one list: a tool
// added or renamed here cannot leave a stale copy behind in a form.
func (s *Server) routeTools(mux *http.ServeMux) int {
mux.HandleFunc("GET /api/v1/tools", func(w http.ResponseWriter, r *http.Request) {
writeJSON(w, http.StatusOK, envelope{Data: s.toolCatalogue})
})
return 1
}
// routeVersion exposes the build identifier to an authenticated caller.
//
// Under /api/v1 rather than on /health deliberately. /health is public, and it
// already withholds its detail from the internet for the reason given above; a
// build identifier is exactly the kind of thing that tells an unauthenticated
// reader which source to go and read. An operator has a session, so this is
// where an operator can reach it and a stranger cannot.
func (s *Server) routeVersion(mux *http.ServeMux) int {
mux.HandleFunc("GET /api/v1/version", func(w http.ResponseWriter, r *http.Request) {
writeJSON(w, http.StatusOK, envelope{Data: map[string]any{
"version": s.version,
"env": s.cfg.AppEnv,
"endpoints": s.endpoints,
}})
})
return 1
}
// handleHealth reports whether this instance should be sent traffic.
//
// 200 "ok" serving normally
// 200 "degraded" the process is healthy, the schema is not: unmigrated,
// or a migration left the version dirty. Still 200,
// because the fault is the database's and taking the
// instance out of rotation would not fix it.
// 503 "unavailable" the database is unreachable, so a load balancer can act
// on the status code alone without parsing the body.
//
// The three status words are a coarse operational signal, not infrastructure
// detail: they say what a caller should do, and nothing about what is running.
func (s *Server) handleHealth(w http.ResponseWriter, r *http.Request) {
ctx, cancel := context.WithTimeout(r.Context(), 5*time.Second)
defer cancel()
health := s.db.Check(ctx)
status, code := "ok", http.StatusOK
switch {
case !health.Reachable:
status, code = "unavailable", http.StatusServiceUnavailable
case health.MigrationDirty, !health.SchemaPresent:
status = "degraded"
}
s.logHealth(status, health)
w.Header().Set("Content-Type", "application/json; charset=utf-8")
w.Header().Set("Cache-Control", "no-store")
w.WriteHeader(code)
enc := json.NewEncoder(w)
enc.SetIndent("", " ")
_ = enc.Encode(healthResponse{Status: status})
}
// logHealth writes the detail the response body used to carry.
//
// This is the "internally" half of the change: nothing was deleted from
// db.Check, and nothing it learns is thrown away — the audience moved from the
// response to the log, which is already authenticated by virtue of being on
// the host.
//
// A load balancer polls this endpoint every few seconds, so a healthy check
// logs at debug and a bad one at warn. Anything other than "ok" is worth
// seeing without turning debug on.
func (s *Server) logHealth(status string, h db.Health) {
attrs := []any{
"status", status,
"env", s.cfg.AppEnv,
"uptime_seconds", int64(time.Since(s.started).Seconds()),
"reachable", h.Reachable,
"schema", h.Schema,
"schema_present", h.SchemaPresent,
"table_count", h.TableCount,
"migration_dirty", h.MigrationDirty,
"latency_ms", h.LatencyMS,
}
if h.Database != "" {
attrs = append(attrs, "database", h.Database, "postgres_version", h.Version)
}
if h.AppliedMigration != nil {
attrs = append(attrs, "applied_migration", *h.AppliedMigration)
}
if h.Error != "" {
attrs = append(attrs, "error", h.Error)
}
if status == "ok" {
s.log.Debug("health", attrs...)
return
}
s.log.Warn("health", attrs...)
}
type statusRecorder struct {
http.ResponseWriter
code int
}
func (r *statusRecorder) WriteHeader(code int) {
r.code = code
r.ResponseWriter.WriteHeader(code)
}
// Flush forwards to the writer underneath.
//
// A wrapper that embeds http.ResponseWriter inherits Write and WriteHeader and
// SILENTLY DROPS every optional interface the real writer implements — Flusher
// among them. Nothing errors: the handler simply asks "can this flush?", is
// told no, and takes whatever fallback it has.
//
// That is exactly how it presented. The SSE endpoint answered ordinary JSON,
// correctly and completely, with no error anywhere — because two middlewares
// deep the writer had stopped being a Flusher and the streaming path politely
// declined to stream.
func (r *statusRecorder) Flush() {
if f, ok := r.ResponseWriter.(http.Flusher); ok {
f.Flush()
}
}
func requestLogger(log *slog.Logger) func(http.Handler) http.Handler {
return func(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
started := time.Now()
rec := &statusRecorder{ResponseWriter: w, code: http.StatusOK}
next.ServeHTTP(rec, r)
log.Info("request",
"method", r.Method, "path", r.URL.Path,
"status", rec.code, "duration_ms", time.Since(started).Milliseconds())
})
}
}
// jsonErrors converts net/http's own plain-text 404 and 405 replies into the
// documented error envelope.
//
// ServeMux writes those itself, before any handler of ours runs, so a client
// that hit a wrong path or method would otherwise get "404 page not found" in
// text/plain while every other response is JSON. Only the mux's own replies are
// rewritten: anything that set a content type has already answered properly.
func jsonErrors(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
iw := &interceptor{ResponseWriter: w}
next.ServeHTTP(iw, r)
if !iw.rewritten || iw.wrote {
return
}
errCode, message := "not_found", "resource not found"
if iw.code == http.StatusMethodNotAllowed {
errCode = "method_not_allowed"
message = r.Method + " is not supported for this resource"
}
writeJSON(w, iw.code, errorEnvelope{Error: errorBody{
Code: errCode, Message: message, Details: map[string]string{},
}})
})
}
// interceptor defers the mux's plain-text 404/405 body so it can be replaced.
type interceptor struct {
http.ResponseWriter
code int
rewritten bool // this is a mux-generated 404/405 we intend to replace
wrote bool // a body already went to the client
}
func (i *interceptor) WriteHeader(code int) {
i.code = code
if code == http.StatusNotFound || code == http.StatusMethodNotAllowed {
if i.Header().Get("Content-Type") != "application/json; charset=utf-8" {
i.rewritten = true
return // hold the header back; jsonErrors writes its own
}
}
i.ResponseWriter.WriteHeader(code)
}
func (i *interceptor) Write(b []byte) (int, error) {
if i.rewritten {
return len(b), nil // swallow the mux's plain-text body
}
i.wrote = true
return i.ResponseWriter.Write(b)
}
// Flush forwards to the writer underneath. See statusRecorder.Flush.
func (i *interceptor) Flush() {
if f, ok := i.ResponseWriter.(http.Flusher); ok {
f.Flush()
}
}
// recoverer turns a panic into a logged 500 rather than a dropped connection.
func recoverer(log *slog.Logger) func(http.Handler) http.Handler {
return func(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
defer func() {
if v := recover(); v != nil {
log.Error("panic", "value", v, "path", r.URL.Path)
writeJSON(w, http.StatusInternalServerError, errorEnvelope{Error: errorBody{
Code: "internal", Message: "internal error", Details: map[string]string{},
}})
}
}()
next.ServeHTTP(w, r)
})
}
}