diff --git a/go-api/internal/memory/memory.go b/go-api/internal/memory/memory.go index 7a810bc..8f7c3ca 100644 --- a/go-api/internal/memory/memory.go +++ b/go-api/internal/memory/memory.go @@ -38,6 +38,7 @@ import ( "time" "github.com/krow/krow-backend/go-api/internal/authctx" + "github.com/krow/krow-backend/go-api/internal/knowledge" "github.com/krow/krow-backend/go-api/internal/repo" ) @@ -141,11 +142,15 @@ func (w Write) Validate() error { return nil } -// Embedder turns text into a comparable vector. The knowledge package's -// embedder satisfies this; memory does not define its own, so a deployment -// cannot end up with two embedding models and vectors that cannot be compared. +// Embedder turns text into a comparable vector. +// +// It is knowledge.Embedder by structure rather than by import: memory does not +// define its own embedding, because two embedding models in one deployment +// produce vectors that cannot be compared and the failure is silent — a recall +// that returns nothing rather than an error. Declaring the shape here keeps +// the dependency one-way while making it impossible to pass a different one. type Embedder interface { - Embed(ctx context.Context, texts []string, kind string) ([][]float32, error) + Embed(ctx context.Context, texts []string, kind knowledge.Kind) ([][]float32, error) Model() string } @@ -183,7 +188,7 @@ func (s *Store) Remember(ctx context.Context, who authctx.Identity, w Write) (st var vector []float32 model := "" if s.embedder != nil { - vectors, err := s.embedder.Embed(ctx, []string{w.Text}, "document") + vectors, err := s.embedder.Embed(ctx, []string{w.Text}, knowledge.KindDocument) // Degraded, not failed: a memory that is stored but not yet searchable // is recoverable by re-embedding, and losing it is not. if err == nil && len(vectors) == 1 && len(vectors[0]) > 0 { @@ -274,7 +279,7 @@ func (s *Store) Recall(ctx context.Context, who authctx.Identity, question strin } if s.embedder != nil && strings.TrimSpace(question) != "" { - vectors, err := s.embedder.Embed(ctx, []string{question}, "query") + vectors, err := s.embedder.Embed(ctx, []string{question}, knowledge.KindQuery) if err == nil && len(vectors) == 1 && len(vectors[0]) > 0 { rows, err := s.query(ctx, ` SELECT id, subject_type, coalesce(subject_id::text, ''), text, author, diff --git a/go-api/internal/runtime/loop.go b/go-api/internal/runtime/loop.go index 418101a..e8f553b 100644 --- a/go-api/internal/runtime/loop.go +++ b/go-api/internal/runtime/loop.go @@ -51,6 +51,11 @@ type ModelExecutor struct { // agent built so far answers from the operational tables through tools, and // none of them needs a corpus. retriever Retriever + + // memory is the long-term store, or nil for a deployment that does not + // remember. Nil is a supported state: the loop behaves exactly as it did + // before memory existed. + memory Recaller } // Retriever is what the loop needs from the knowledge layer. @@ -333,19 +338,37 @@ func (m *ModelExecutor) executeRun( // never on the question, so without this a greeting was handed eight policy // chunks AHEAD of the word "hi" — which is both the dominant cost of the // turn and the reason the answer came back as an operational briefing. - conversation := []gateway.Message{{Role: gateway.RoleUser, Text: question}} + /* Evidence blocks, in the order the model should meet them: what this + workspace remembered, then what the corpus says, then the question. + Both are fenced and labelled as data; neither is an instruction. */ + var blocks []string + + // Memory, before retrieval. It is the smaller block and the more general + // one — a standing preference frames how the documents should be read, and + // a reader meeting it first is not being told a conclusion, only a + // context. Skipped for smalltalk on the same terms as retrieval: a + // greeting does not need remembering, and paying for it is how "hi" came + // to cost six thousand tokens. + if !smalltalk { + if block := m.recall(runCtx, rec, input, question); block != "" { + blocks = append(blocks, block) + } + } + if !smalltalk { if block, retrieved := m.retrieve(runCtx, rec, agent, input, question); block != "" { - conversation = []gateway.Message{{ - Role: gateway.RoleUser, - // Context first, question second. A model reads the question last - // and answers it, rather than treating the evidence as the prompt. - Text: block + "\n\n" + question, - }} + blocks = append(blocks, block) rec.Retrieval(retrieved) } } + // Context first, question second. A model reads the question last and + // answers it, rather than treating the evidence as the prompt. + conversation := []gateway.Message{{ + Role: gateway.RoleUser, + Text: strings.TrimSpace(strings.Join(append(blocks, question), "\n\n")), + }} + // Recorded once rather than on each remaining step, so a long run does not // fill its trajectory with the same note. var toolsWithheld bool diff --git a/go-api/internal/runtime/recall.go b/go-api/internal/runtime/recall.go new file mode 100644 index 0000000..96c4912 --- /dev/null +++ b/go-api/internal/runtime/recall.go @@ -0,0 +1,69 @@ +package runtime + +import ( + "context" + "strconv" + + "github.com/krow/krow-backend/go-api/internal/authctx" + "github.com/krow/krow-backend/go-api/internal/memory" +) + +// Recaller is the memory store's read half, as the loop needs it. +// +// An interface rather than the struct so the runtime does not depend on +// memory's internals and a test can drive the loop without a database — the +// same shape Retriever has, for the same reason. +type Recaller interface { + Recall(ctx context.Context, who authctx.Identity, question string, limit int) ([]memory.Record, string, error) +} + +// WithMemory gives the executor a long-term memory to read. +// +// Optional, like the retriever. Nil means a deployment that has not migrated +// 000017, or has chosen not to remember, and the loop behaves exactly as it +// did before memory existed. +func (m *ModelExecutor) WithMemory(r Recaller) *ModelExecutor { + m.memory = r + return m +} + +// recall returns the memory block for this question, or "". +// +// FAILS QUIET, LOUDLY RECORDED. A memory store that is unreachable must not +// take the run with it: the answer without memory is worse, not wrong, and the +// alternative is an outage in the knowledge layer becoming an outage in the +// product. The trajectory says what happened, so an answer that reads thin is +// explainable afterwards rather than mysterious. +// +// The degraded note from the store — "no embedder is configured; these +// memories are the most recent rather than the most relevant" — is recorded +// too. A reader comparing two answers needs to know which one got relevance +// and which got recency. +func (m *ModelExecutor) recall(ctx context.Context, rec *Recorder, input ExecutionInput, question string) string { + if m.memory == nil { + return "" + } + + records, degraded, err := m.memory.Recall(ctx, input.Identity, question, memory.DefaultRecall) + if err != nil { + rec.Error("memory.failed", err.Error()) + return "" + } + if degraded != "" { + rec.Error("memory.degraded", degraded) + } + if len(records) == 0 { + return "" + } + + rec.Error("runtime.memory_recalled", + pluralMemories(len(records))+" carried into this run") + return memory.Render(records) +} + +func pluralMemories(n int) string { + if n == 1 { + return "1 memory" + } + return strconv.Itoa(n) + " memories" +} diff --git a/go-api/internal/runtime/recall_test.go b/go-api/internal/runtime/recall_test.go new file mode 100644 index 0000000..8f2a655 --- /dev/null +++ b/go-api/internal/runtime/recall_test.go @@ -0,0 +1,113 @@ +package runtime + +import ( + "context" + "errors" + "strings" + "testing" + + "github.com/krow/krow-backend/go-api/internal/authctx" + "github.com/krow/krow-backend/go-api/internal/memory" +) + +type fakeMemory struct { + records []memory.Record + degraded string + err error + asked string +} + +func (f *fakeMemory) Recall(_ context.Context, _ authctx.Identity, question string, _ int) ([]memory.Record, string, error) { + f.asked = question + return f.records, f.degraded, f.err +} + +func TestARecalledMemoryReachesThePrompt(t *testing.T) { + gw := &fakeGateway{text: "answered"} + mem := &fakeMemory{records: []memory.Record{ + {SubjectType: memory.SubjectWorkspace, Author: memory.AuthorModel, Text: "Thursdays are short-staffed."}, + }} + exec := NewModelExecutor(gw, &MemorySink{}, nil).WithMemory(mem) + + if _, err := exec.ExecuteAgent(context.Background(), testAgent(), testInput("who is free?")); err != nil { + t.Fatalf("run failed: %v", err) + } + sent := gw.lastReq.Messages[0].Text + if !strings.Contains(sent, "Thursdays are short-staffed.") { + t.Errorf("the memory did not reach the model:\n%s", sent) + } + if !strings.Contains(sent, "") { + t.Error("the memory was not fenced") + } +} + +// The question goes last. A model reads the last thing and answers it; put the +// evidence after it and the evidence becomes the prompt. +func TestTheQuestionStaysLastWhenMemoryIsCarried(t *testing.T) { + gw := &fakeGateway{text: "answered"} + mem := &fakeMemory{records: []memory.Record{ + {SubjectType: memory.SubjectWorkspace, Author: memory.AuthorModel, Text: "A remembered thing."}, + }} + exec := NewModelExecutor(gw, &MemorySink{}, nil).WithMemory(mem) + exec.ExecuteAgent(context.Background(), testAgent(), testInput("who is free?")) + + sent := gw.lastReq.Messages[0].Text + if !strings.HasSuffix(strings.TrimSpace(sent), "who is free?") { + t.Errorf("the question is not last:\n%s", sent) + } +} + +// A store that is down must not take the run with it: an answer without +// memory is worse, not wrong. +func TestAMemoryFailureDoesNotFailTheRun(t *testing.T) { + gw := &fakeGateway{text: "answered anyway"} + mem := &fakeMemory{err: errors.New("the memory store is unreachable")} + sink := &MemorySink{} + exec := NewModelExecutor(gw, sink, nil).WithMemory(mem) + + res, err := exec.ExecuteAgent(context.Background(), testAgent(), testInput("who is free?")) + if err != nil { + t.Fatalf("a memory failure took the run with it: %v", err) + } + if res.Termination != TerminationCompleted { + t.Errorf("Termination = %q, want Completed", res.Termination) + } + + // ...and it is recorded, so a thin answer is explainable afterwards. + var noted bool + for _, e := range sink.Last().Entries { + if strings.Contains(e.Name, "memory") || strings.Contains(e.Text, "memory store") { + noted = true + } + } + if !noted { + t.Error("the memory failure left no trace in the trajectory") + } +} + +// Greetings skip memory for the same reason they skip retrieval: nobody needs +// remembering to say good morning, and paying for it is how "hi" came to cost +// six thousand tokens. +func TestSmalltalkCarriesNoMemory(t *testing.T) { + gw := &fakeGateway{text: "Good morning."} + mem := &fakeMemory{records: []memory.Record{ + {SubjectType: memory.SubjectWorkspace, Author: memory.AuthorModel, Text: "A remembered thing."}, + }} + exec := NewModelExecutor(gw, &MemorySink{}, nil).WithMemory(mem) + exec.ExecuteAgent(context.Background(), testAgent(), testInput("good morning")) + + if strings.Contains(gw.lastReq.Messages[0].Text, "") { + t.Errorf("a greeting carried memory:\n%s", gw.lastReq.Messages[0].Text) + } +} + +// With no store the loop is exactly what it was. +func TestWithoutAMemoryStoreNothingChanges(t *testing.T) { + gw := &fakeGateway{text: "answered"} + exec := NewModelExecutor(gw, &MemorySink{}, nil) + exec.ExecuteAgent(context.Background(), testAgent(), testInput("who is free?")) + + if gw.lastReq.Messages[0].Text != "who is free?" { + t.Errorf("the question was altered with no memory configured:\n%q", gw.lastReq.Messages[0].Text) + } +} diff --git a/go-api/internal/runtime/wire.go b/go-api/internal/runtime/wire.go index 31c3f45..a81dd06 100644 --- a/go-api/internal/runtime/wire.go +++ b/go-api/internal/runtime/wire.go @@ -4,6 +4,7 @@ import ( "github.com/krow/krow-backend/go-api/internal/config" "github.com/krow/krow-backend/go-api/internal/gateway" "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/repo" "github.com/krow/krow-backend/go-api/internal/tools" ) @@ -26,8 +27,16 @@ func NewModelEngine(db repo.Querier, cfg config.Config) *Engine { // alternative provider needed an edit to this file to be reachable. gw := gateway.New(gateway.FromConfig(cfg.Model)) retriever := knowledge.NewRetriever(db, NewEmbedder(cfg)) - exec := NewModelExecutor(gw, NewPostgresSink(db), DefaultTools(db, retriever)). + + // Long-term memory shares the embedder with retrieval, deliberately: two + // embedding models in one deployment produce vectors that cannot be + // compared, and the failure is silent — a recall that quietly returns + // nothing rather than an error. + memories := memory.New(db, NewEmbedder(cfg)) + + exec := NewModelExecutor(gw, NewPostgresSink(db), DefaultToolsWithMemory(db, retriever, memories)). WithRetriever(retriever). + WithMemory(memories). // Without this, a spec's `subagents:` parses, loads, and is then // dropped — which is how krow-workforce-agent came to declare five // subagents and answer every question by itself. The resolver is the @@ -118,6 +127,16 @@ func NewEmbedder(cfg config.Config) knowledge.Embedder { // service that booted without a capability its specs name would fail one run // at a time instead of once, loudly, at startup. func DefaultTools(db repo.Querier, retriever *knowledge.Retriever) *tools.Registry { + return DefaultToolsWithMemory(db, retriever, nil) +} + +// DefaultToolsWithMemory is DefaultTools with long-term memory available. +// +// A separate constructor rather than a nil check inside the old one, because a +// deployment that has not migrated 000017 must not offer a tool whose every +// call would fail against a table that is not there. Passing nil registers the +// catalogue exactly as it was. +func DefaultToolsWithMemory(db repo.Querier, retriever *knowledge.Retriever, memories tools.MemoryWriter) *tools.Registry { // The confirmation store is Postgres-backed, not in-process. A pending // write is asked about in one request and approved in another, and nothing // guarantees those two reach the same replica — an in-memory store would @@ -162,5 +181,12 @@ func DefaultTools(db repo.Querier, retriever *knowledge.Retriever) *tools.Regist } { reg.MustRegister(t) } + + // Memory, only where there is somewhere to put it. It is a confirmed + // write like assign_worker: it stores personal data that shapes later + // hiring answers, so a person sees the sentence before it is kept. + if memories != nil { + reg.MustRegister(tools.Remember(memories)) + } return reg }