From 43f92e310202e09829b58055da7bc0c4e197c97f Mon Sep 17 00:00:00 2001 From: Mathias Date: Mon, 22 Jun 2026 23:21:43 +0200 Subject: [PATCH 1/2] feat(capture): CaptureService use-case + ports + entities (#51) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The Clean-Architecture core of the capture capability (#49b). Pure orchestration over ports — no HTTP, no live Gitea, no audit I/O — fully unit-tested against fakes before any adapter exists. - Ports: BrainStore (#45 write/update/get), IssueTracker, SummaryWriter, ClassificationPolicy (satisfied by #50's classification.Config), AuditSink. Entities: Insight, Ticket, Summary, CaptureContext, CaptureInput, CaptureReceipt. - CaptureService.Capture: validate-before-write (fail-closed), resolve effective classification (stricter of declared vs target-derived; under-declaration logged as a security event), orchestrate insights (write/supersede) → tickets → summary best-effort, emit a request-level audit record of exactly what landed, return a partial-aware receipt. - dry_run short-circuits after validation, writes nothing (not even audit). Out of scope here, layered on later: the I1 origin sovereignty gate (#53, needs the server-derived principal) and the classification-aware audit degradation/refusal (#54). "Effective" is folded into the service as classification.Stricter rather than a port method — the stricter-wins rule is use-case policy. Co-Authored-By: Claude Opus 4.8 (1M context) --- ingestion/internal/capture/entities.go | 111 +++++++ ingestion/internal/capture/ports.go | 102 +++++++ ingestion/internal/capture/service.go | 316 ++++++++++++++++++++ ingestion/internal/capture/service_test.go | 331 +++++++++++++++++++++ 4 files changed, 860 insertions(+) create mode 100644 ingestion/internal/capture/entities.go create mode 100644 ingestion/internal/capture/ports.go create mode 100644 ingestion/internal/capture/service.go create mode 100644 ingestion/internal/capture/service_test.go diff --git a/ingestion/internal/capture/entities.go b/ingestion/internal/capture/entities.go new file mode 100644 index 0000000..05f295d --- /dev/null +++ b/ingestion/internal/capture/entities.go @@ -0,0 +1,111 @@ +// Package capture is the Clean-Architecture use-case for the uniform +// capture capability (issue #49/#51): persist a finished session's +// valuable output — insights → brain, action items → Gitea tickets, +// optional summary → ai-sessions — with one invocation, identical core +// behaviour across every harness. +// +// This package is pure orchestration. It depends only on ports +// (interfaces) and plain entities — no HTTP, no live Gitea, no embedding +// or audit I/O. The real adapters are wired in #52 (Gitea tracker), #53 +// (REST + I1 origin gate), and #54/#55 (audit path + relay). The I1 +// sovereignty refusal and the classification-aware audit degradation are +// deliberately NOT here — those need the server-derived principal origin +// (#53) and the loki/buffer machinery (#54). What lives here is everything +// testable against fakes: validation, effective-classification resolution +// (stricter wins), best-effort orchestration, and the partial receipt. +package capture + +// CaptureContext is the per-session metadata accompanying a capture. +// +// Classification is the caller-declared sensitivity (model C, spec §4.1): +// the server independently derives the target's classification and gates +// on the stricter of the two. Principal is server-derived from the +// authenticated identity (#53 populates it); it is never caller-asserted. +// Harness is descriptive telemetry only — never a gate input. +type CaptureContext struct { + Harness string + SessionRef string + Fidelity string + Actor string + Classification string // caller-declared level token ("" = unspecified) + Principal string // server-derived (auth); audit identity +} + +// Insight is one piece of session knowledge bound for the brain. A +// non-empty SupersedeSlug routes to Update (revise in place); otherwise +// Write (create). +type Insight struct { + Text string + Wing string + Hall string + SupersedeSlug string +} + +// Ticket is one action item bound for a Gitea repo. Owner is always the +// operator (set by the tracker adapter), never carried here. +type Ticket struct { + Repo string + Action string // create | close | comment + Number int // required for close/comment + Title string // required for create + Body string +} + +// Summary is an optional session summary bound for ai-sessions. +type Summary struct { + Title string + Body string + ReposTouched []string +} + +// CaptureInput is the whole capture request. +type CaptureInput struct { + Context CaptureContext + Insights []Insight + Tickets []Ticket + Summary *Summary + DryRun bool +} + +// InsightResult is the per-insight outcome in the receipt. +type InsightResult struct { + ID string `json:"id,omitempty"` + Path string `json:"path,omitempty"` + ContentHash string `json:"content_hash,omitempty"` + Superseded bool `json:"superseded"` + OK bool `json:"ok"` +} + +// TicketResult is the per-ticket outcome in the receipt. +type TicketResult struct { + Repo string `json:"repo"` + Number int `json:"number,omitempty"` + Action string `json:"action"` + URL string `json:"url,omitempty"` + OK bool `json:"ok"` +} + +// SummaryResult is the summary outcome in the receipt. +type SummaryResult struct { + Path string `json:"path,omitempty"` + OK bool `json:"ok"` +} + +// ItemError pins a failure to a specific request item for the partial +// receipt. Item is a stable locator like "insight[1]" or "ticket[0]". +type ItemError struct { + Item string `json:"item"` + Error string `json:"error"` +} + +// CaptureReceipt is the structured, partial-aware result. Per-item ok +// flags plus a flat Errors list make partial success explicit; the +// caller never has to infer what landed. +type CaptureReceipt struct { + Insights []InsightResult `json:"insights"` + Tickets []TicketResult `json:"tickets"` + Summary *SummaryResult `json:"summary,omitempty"` + Errors []ItemError `json:"errors"` + EffectiveClassification string `json:"effective_classification,omitempty"` + DryRun bool `json:"dry_run"` +} diff --git a/ingestion/internal/capture/ports.go b/ingestion/internal/capture/ports.go new file mode 100644 index 0000000..58a7d52 --- /dev/null +++ b/ingestion/internal/capture/ports.go @@ -0,0 +1,102 @@ +package capture + +import ( + "context" + "time" + + "github.com/mathiasbq/hyperguild/ingestion/internal/classification" +) + +// Ref is the read-after-write handle returned by a brain write/update — +// the #45 contract. ContentHash lets the caller confirm what landed +// without a re-query; for an Update, Superseded is true. +type Ref struct { + ID string + Path string + ContentHash string + Superseded bool +} + +// StoredNote is a brain note fetched by Get: the read-after-write +// confirmation primitive (a direct fetch, never a semantic query). +type StoredNote struct { + ID string + Path string + ContentHash string + Frontmatter map[string]string + Body string +} + +// Note is the brain-write payload. It carries both the wing/hall taxonomy +// and the legacy type/domain fields so a single BrainStore serves both +// capture insights and the existing MCP brain_write surface. Reason is +// the supersede rationale, used only by Update. +type Note struct { + Content string + Filename string + Wing string + Hall string + Type string + Domain string + Reason string +} + +// BrainStore is the brain persistence port — the shared implementation of +// the #45 write/update/get verbs that both the MCP handlers and capture +// call, so there is one implementation, not two. The read-after-write + +// staleness discipline lives behind this interface so no caller carries +// the rule. +type BrainStore interface { + Write(ctx context.Context, n Note) (Ref, error) + Update(ctx context.Context, slug string, n Note) (Ref, error) + Get(ctx context.Context, id string) (StoredNote, error) +} + +// IssueRef identifies a ticket touched by the tracker. +type IssueRef struct { + Repo string + Number int + URL string +} + +// IssueTracker is the Gitea ticket port. The implementation (#52) always +// scopes to owner "mathias"; the port deliberately omits owner. +type IssueTracker interface { + CreateIssue(ctx context.Context, repo, title, body string) (IssueRef, error) + CloseIssue(ctx context.Context, repo string, number int) (IssueRef, error) + CommentIssue(ctx context.Context, repo string, number int, body string) (IssueRef, error) +} + +// SummaryWriter is the ai-sessions summary port. +type SummaryWriter interface { + WriteFile(ctx context.Context, repo, path, content string) error +} + +// ClassificationPolicy derives a target's sensitivity (model C). The +// "stricter wins" combination of declared vs derived is use-case policy +// and lives in the service, so the port stays minimal. Satisfied by +// classification.Config (#50). +type ClassificationPolicy interface { + Derive(target classification.Target) classification.Level +} + +// AuditEntry is the request-level audit record (I5): who/what captured +// what, when, via which principal. SecurityEvents carries anomalies such +// as a caller under-declaring sensitivity relative to the target floor. +type AuditEntry struct { + Timestamp time.Time + Principal string + Actor string + Harness string + SessionRef string + EffectiveClassification string + Items []string + SecurityEvents []string +} + +// AuditSink records the audit entry. The classification-aware +// degradation/refusal policy (confidential fails closed, internal +// degrades) is the caller's concern in #54; this port just records. +type AuditSink interface { + Record(ctx context.Context, e AuditEntry) error +} diff --git a/ingestion/internal/capture/service.go b/ingestion/internal/capture/service.go new file mode 100644 index 0000000..132e537 --- /dev/null +++ b/ingestion/internal/capture/service.go @@ -0,0 +1,316 @@ +package capture + +import ( + "context" + "crypto/sha256" + "encoding/hex" + "fmt" + "strings" + "time" + + "github.com/mathiasbq/hyperguild/ingestion/internal/brain" + "github.com/mathiasbq/hyperguild/ingestion/internal/classification" +) + +// Service is the CaptureSession use-case. It depends only on ports. +type Service struct { + brain BrainStore + issues IssueTracker + summaries SummaryWriter + policy ClassificationPolicy + audit AuditSink + + // now is the clock, injectable for deterministic summary paths and + // audit timestamps in tests. + now func() time.Time +} + +// NewService constructs a Service from its ports. summaries may be nil +// when no summary persistence is wired; a CaptureInput with a Summary +// then fails that item rather than panicking. +func NewService(b BrainStore, tr IssueTracker, sw SummaryWriter, p ClassificationPolicy, a AuditSink) *Service { + return &Service{brain: b, issues: tr, summaries: sw, policy: p, audit: a, now: time.Now} +} + +var validActions = map[string]bool{"create": true, "close": true, "comment": true} + +// Capture runs the use-case: validate (fail-closed), resolve effective +// classification (stricter of declared vs target-derived), then persist +// insights → tickets → summary best-effort, emit an audit record, and +// return a partial-aware receipt. +// +// A validation failure returns a non-nil error with nothing written. A +// per-item execution failure is recorded in the receipt (no rollback); +// the call still returns a nil error so the caller gets the partial +// receipt. The I1 origin gate and audit-down degradation are layered on +// by #53/#54 around this core. +func (s *Service) Capture(ctx context.Context, in CaptureInput) (CaptureReceipt, error) { + if err := s.validate(in); err != nil { + return CaptureReceipt{}, err + } + + declared := classification.Public // unspecified ⇒ lowest ⇒ target floor governs + if in.Context.Classification != "" { + // Already validated parseable. + declared, _ = classification.ParseLevel(in.Context.Classification) + } + + effective, securityEvents := s.resolveClassification(declared, in) + + receipt := CaptureReceipt{ + Errors: []ItemError{}, + EffectiveClassification: effective.String(), + DryRun: in.DryRun, + } + + if in.DryRun { + // Would-be receipt: mark planned items ok, write nothing (not even + // audit — dry_run touches nothing). + for range in.Insights { + receipt.Insights = append(receipt.Insights, InsightResult{OK: true}) + } + for _, tk := range in.Tickets { + receipt.Tickets = append(receipt.Tickets, TicketResult{Repo: tk.Repo, Action: tk.Action, Number: tk.Number, OK: true}) + } + if in.Summary != nil { + receipt.Summary = &SummaryResult{Path: s.summaryPath(in.Context, in.Summary), OK: true} + } + return receipt, nil + } + + var landed []string + + for i, ins := range in.Insights { + res, item, err := s.persistInsight(ctx, ins) + receipt.Insights = append(receipt.Insights, res) + if err != nil { + receipt.Errors = append(receipt.Errors, ItemError{Item: fmt.Sprintf("insight[%d]", i), Error: err.Error()}) + continue + } + landed = append(landed, item) + } + + for i, tk := range in.Tickets { + res, err := s.persistTicket(ctx, tk) + receipt.Tickets = append(receipt.Tickets, res) + if err != nil { + receipt.Errors = append(receipt.Errors, ItemError{Item: fmt.Sprintf("ticket[%d]", i), Error: err.Error()}) + continue + } + landed = append(landed, fmt.Sprintf("ticket:%s#%d", tk.Repo, res.Number)) + } + + if in.Summary != nil { + res, err := s.persistSummary(ctx, in.Context, in.Summary) + receipt.Summary = &res + if err != nil { + receipt.Errors = append(receipt.Errors, ItemError{Item: "summary", Error: err.Error()}) + } else { + landed = append(landed, "summary:"+res.Path) + } + } + + // I5: emit a request-level audit record of exactly what landed. + // Best-effort here; the classification-aware refusal/degradation + // policy is #54. + if err := s.audit.Record(ctx, AuditEntry{ + Timestamp: s.now().UTC(), + Principal: in.Context.Principal, + Actor: in.Context.Actor, + Harness: in.Context.Harness, + SessionRef: in.Context.SessionRef, + EffectiveClassification: effective.String(), + Items: landed, + SecurityEvents: securityEvents, + }); err != nil { + receipt.Errors = append(receipt.Errors, ItemError{Item: "audit", Error: err.Error()}) + } + + return receipt, nil +} + +// validate enforces fail-closed structural validity over the whole +// request before any write. A bad declared classification, an invalid +// wing/hall, an empty insight, or a malformed ticket aborts the capture +// with nothing written. +func (s *Service) validate(in CaptureInput) error { + if in.Context.Classification != "" { + if _, err := classification.ParseLevel(in.Context.Classification); err != nil { + return fmt.Errorf("context.classification: %w", err) + } + } + for i, ins := range in.Insights { + if strings.TrimSpace(ins.Text) == "" { + return fmt.Errorf("insight[%d]: text is required", i) + } + if strings.TrimSpace(ins.Wing) == "" { + return fmt.Errorf("insight[%d]: wing is required", i) + } + if !brain.IsValidHall(ins.Hall) { + return fmt.Errorf("insight[%d]: invalid hall %q", i, ins.Hall) + } + } + for i, tk := range in.Tickets { + if strings.TrimSpace(tk.Repo) == "" { + return fmt.Errorf("ticket[%d]: repo is required", i) + } + if !validActions[tk.Action] { + return fmt.Errorf("ticket[%d]: invalid action %q (want create/close/comment)", i, tk.Action) + } + if tk.Action == "create" && strings.TrimSpace(tk.Title) == "" { + return fmt.Errorf("ticket[%d]: create requires a title", i) + } + if (tk.Action == "close" || tk.Action == "comment") && tk.Number <= 0 { + return fmt.Errorf("ticket[%d]: %s requires an issue number", i, tk.Action) + } + } + return nil +} + +// resolveClassification computes the effective level (stricter of +// declared and every target's derived level) and collects a security +// event whenever the caller under-declared relative to a target floor. +func (s *Service) resolveClassification(declared classification.Level, in CaptureInput) (classification.Level, []string) { + effective := declared + var events []string + consider := func(kind classification.TargetKind, name string) { + derived := s.policy.Derive(classification.Target{Kind: kind, Name: name}) + effective = classification.Stricter(effective, derived) + if declared < derived { + events = append(events, fmt.Sprintf("classification under-declared: declared=%s target=%s(%s) derived=%s", + declared, name, kindString(kind), derived)) + } + } + for _, ins := range in.Insights { + consider(classification.WingTarget, ins.Wing) + } + for _, tk := range in.Tickets { + consider(classification.RepoTarget, tk.Repo) + } + if in.Summary != nil { + for _, repo := range in.Summary.ReposTouched { + consider(classification.RepoTarget, repo) + } + } + return effective, events +} + +func (s *Service) persistInsight(ctx context.Context, ins Insight) (InsightResult, string, error) { + note := Note{Content: ins.Text, Wing: ins.Wing, Hall: ins.Hall, Filename: brain.Sanitise(firstLine(ins.Text))} + var ref Ref + var err error + if ins.SupersedeSlug != "" { + note.Reason = "superseded via capture" + ref, err = s.brain.Update(ctx, ins.SupersedeSlug, note) + } else { + ref, err = s.brain.Write(ctx, note) + } + if err != nil { + return InsightResult{OK: false, Superseded: ins.SupersedeSlug != ""}, "", err + } + return InsightResult{ + ID: ref.ID, Path: ref.Path, ContentHash: ref.ContentHash, + Superseded: ref.Superseded, OK: true, + }, "insight:" + ref.ID, nil +} + +func (s *Service) persistTicket(ctx context.Context, tk Ticket) (TicketResult, error) { + res := TicketResult{Repo: tk.Repo, Action: tk.Action, Number: tk.Number} + var ref IssueRef + var err error + switch tk.Action { + case "create": + ref, err = s.issues.CreateIssue(ctx, tk.Repo, tk.Title, tk.Body) + case "close": + ref, err = s.issues.CloseIssue(ctx, tk.Repo, tk.Number) + case "comment": + ref, err = s.issues.CommentIssue(ctx, tk.Repo, tk.Number, tk.Body) + } + if err != nil { + return res, err + } + if ref.Number != 0 { + res.Number = ref.Number + } + res.URL = ref.URL + res.OK = true + return res, nil +} + +func (s *Service) persistSummary(ctx context.Context, c CaptureContext, sum *Summary) (SummaryResult, error) { + if s.summaries == nil { + return SummaryResult{OK: false}, fmt.Errorf("no summary writer configured") + } + path := s.summaryPath(c, sum) + content := s.renderSummary(c, sum) + repo := "ai-sessions" + if err := s.summaries.WriteFile(ctx, repo, path, content); err != nil { + return SummaryResult{Path: path, OK: false}, err + } + return SummaryResult{Path: path, OK: true}, nil +} + +// summaryPath builds summaries///--.md. +// The ref8 disambiguator is derived from the session_ref (or the title +// when no ref is present) so distinct sessions never collide. +func (s *Service) summaryPath(c CaptureContext, sum *Summary) string { + t := s.now().UTC() + slug := brain.Sanitise(sum.Title) + if slug == "" { + slug = "summary" + } + seed := c.SessionRef + if seed == "" { + seed = sum.Title + sum.Body + } + sum8 := shortHash(seed) + return fmt.Sprintf("summaries/%s/%s/%s-%s-%s.md", + brain.Sanitise(c.Harness), t.Format("2006-01"), t.Format("2006-01-02"), slug, sum8) +} + +// renderSummary stamps fidelity + session metadata into frontmatter so the +// richer-fidelity-supersedes-thinner collision rule has the data it needs. +func (s *Service) renderSummary(c CaptureContext, sum *Summary) string { + var b strings.Builder + b.WriteString("---\n") + fmt.Fprintf(&b, "title: %s\n", sum.Title) + fmt.Fprintf(&b, "harness: %s\n", c.Harness) + if c.SessionRef != "" { + fmt.Fprintf(&b, "session_ref: %s\n", c.SessionRef) + } + fmt.Fprintf(&b, "fidelity: %s\n", c.Fidelity) + fmt.Fprintf(&b, "captured_at: %s\n", s.now().UTC().Format(time.RFC3339)) + if len(sum.ReposTouched) > 0 { + fmt.Fprintf(&b, "repos_touched: [%s]\n", strings.Join(sum.ReposTouched, ", ")) + } + b.WriteString("---\n\n") + b.WriteString(sum.Body) + if !strings.HasSuffix(sum.Body, "\n") { + b.WriteByte('\n') + } + return b.String() +} + +func kindString(k classification.TargetKind) string { + if k == classification.RepoTarget { + return "repo" + } + return "wing" +} + +func firstLine(s string) string { + s = strings.TrimSpace(s) + if i := strings.IndexByte(s, '\n'); i >= 0 { + s = s[:i] + } + s = strings.TrimLeft(s, "# ") + if len(s) > 60 { + s = s[:60] + } + return s +} + +func shortHash(s string) string { + sum := sha256.Sum256([]byte(s)) + return hex.EncodeToString(sum[:])[:8] +} diff --git a/ingestion/internal/capture/service_test.go b/ingestion/internal/capture/service_test.go new file mode 100644 index 0000000..c1608c5 --- /dev/null +++ b/ingestion/internal/capture/service_test.go @@ -0,0 +1,331 @@ +package capture + +import ( + "context" + "errors" + "strings" + "testing" + "time" + + "github.com/mathiasbq/hyperguild/ingestion/internal/classification" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// --- fakes --- + +type fakeBrain struct { + writes []Note + updates []Note + gets []string + failOn func(Note) error // nil = always succeed + hashSeq int +} + +func (f *fakeBrain) ref(prefix string, n Note, superseded bool) Ref { + f.hashSeq++ + path := "wiki/" + n.Wing + "/" + n.Hall + "/" + n.Filename + ".md" + return Ref{ID: path, Path: path, ContentHash: prefix + string(rune('0'+f.hashSeq)), Superseded: superseded} +} + +func (f *fakeBrain) Write(_ context.Context, n Note) (Ref, error) { + if f.failOn != nil { + if err := f.failOn(n); err != nil { + return Ref{}, err + } + } + f.writes = append(f.writes, n) + return f.ref("w", n, false), nil +} + +func (f *fakeBrain) Update(_ context.Context, slug string, n Note) (Ref, error) { + if f.failOn != nil { + if err := f.failOn(n); err != nil { + return Ref{}, err + } + } + n.Filename = slug + f.updates = append(f.updates, n) + return f.ref("u", n, true), nil +} + +func (f *fakeBrain) Get(_ context.Context, id string) (StoredNote, error) { + f.gets = append(f.gets, id) + return StoredNote{ID: id, Path: id}, nil +} + +type fakeTracker struct { + created []string + closed []int + comments []int + err error +} + +func (f *fakeTracker) CreateIssue(_ context.Context, repo, title, _ string) (IssueRef, error) { + if f.err != nil { + return IssueRef{}, f.err + } + f.created = append(f.created, repo+":"+title) + return IssueRef{Repo: repo, Number: 100 + len(f.created), URL: "https://git/" + repo + "/issues/x"}, nil +} + +func (f *fakeTracker) CloseIssue(_ context.Context, repo string, number int) (IssueRef, error) { + if f.err != nil { + return IssueRef{}, f.err + } + f.closed = append(f.closed, number) + return IssueRef{Repo: repo, Number: number}, nil +} + +func (f *fakeTracker) CommentIssue(_ context.Context, repo string, number int, _ string) (IssueRef, error) { + if f.err != nil { + return IssueRef{}, f.err + } + f.comments = append(f.comments, number) + return IssueRef{Repo: repo, Number: number}, nil +} + +type fakeSummary struct { + paths []string + content []string + err error +} + +func (f *fakeSummary) WriteFile(_ context.Context, _, path, content string) error { + if f.err != nil { + return f.err + } + f.paths = append(f.paths, path) + f.content = append(f.content, content) + return nil +} + +// fakePolicy derives from an explicit map; default Internal so tests pin +// behaviour without depending on the real defaulting. +type fakePolicy struct{ tags map[string]classification.Level } + +func (p fakePolicy) Derive(t classification.Target) classification.Level { + if lvl, ok := p.tags[t.Name]; ok { + return lvl + } + return classification.Internal +} + +type fakeAudit struct { + entries []AuditEntry + err error +} + +func (f *fakeAudit) Record(_ context.Context, e AuditEntry) error { + if f.err != nil { + return f.err + } + f.entries = append(f.entries, e) + return nil +} + +// --- helpers --- + +func newSvc(b BrainStore, tr IssueTracker, sw SummaryWriter, p ClassificationPolicy, a AuditSink) *Service { + s := NewService(b, tr, sw, p, a) + s.now = func() time.Time { return time.Date(2026, 6, 22, 12, 0, 0, 0, time.UTC) } + return s +} + +func baseCtx() CaptureContext { + return CaptureContext{Harness: "claude-code", Actor: "mathias", Principal: "mathias", Classification: "internal"} +} + +// --- scenarios --- + +func TestCaptureHappyPath(t *testing.T) { + b := &fakeBrain{} + tr := &fakeTracker{} + au := &fakeAudit{} + svc := newSvc(b, tr, nil, fakePolicy{}, au) + + rec, err := svc.Capture(context.Background(), CaptureInput{ + Context: baseCtx(), + Insights: []Insight{ + {Text: "a", Wing: "hyperguild", Hall: "decisions", SupersedeSlug: ""}, + {Text: "b", Wing: "hyperguild", Hall: "facts"}, + }, + Tickets: []Ticket{{Repo: "hyperguild", Action: "create", Title: "do x", Body: "y"}}, + }) + require.NoError(t, err) + require.Len(t, rec.Insights, 2) + for _, r := range rec.Insights { + assert.True(t, r.OK) + assert.NotEmpty(t, r.ContentHash, "read-after-write hash returned") + } + require.Len(t, rec.Tickets, 1) + assert.True(t, rec.Tickets[0].OK) + assert.Equal(t, 2, len(b.writes)) + assert.Empty(t, rec.Errors) + // Audit emitted naming principal/harness + items that landed. + require.Len(t, au.entries, 1) + assert.Equal(t, "mathias", au.entries[0].Principal) + assert.Equal(t, "claude-code", au.entries[0].Harness) + assert.Len(t, au.entries[0].Items, 3) +} + +func TestCaptureSupersedeNotDuplicate(t *testing.T) { + b := &fakeBrain{} + svc := newSvc(b, &fakeTracker{}, nil, fakePolicy{}, &fakeAudit{}) + + rec, err := svc.Capture(context.Background(), CaptureInput{ + Context: baseCtx(), + Insights: []Insight{{Text: "revised", Wing: "hyperguild", Hall: "facts", SupersedeSlug: "prior-note"}}, + }) + require.NoError(t, err) + assert.Empty(t, b.writes, "supersede must not create") + require.Len(t, b.updates, 1) + assert.Equal(t, "prior-note", b.updates[0].Filename) + assert.True(t, rec.Insights[0].Superseded) +} + +func TestCaptureValidationFailClosed(t *testing.T) { + b := &fakeBrain{} + tr := &fakeTracker{} + au := &fakeAudit{} + svc := newSvc(b, tr, nil, fakePolicy{}, au) + + _, err := svc.Capture(context.Background(), CaptureInput{ + Context: baseCtx(), + Insights: []Insight{ + {Text: "ok", Wing: "hyperguild", Hall: "facts"}, + {Text: "bad", Wing: "hyperguild", Hall: "garbage-hall"}, // invalid hall + }, + Tickets: []Ticket{{Repo: "hyperguild", Action: "create", Title: "t"}}, + }) + require.Error(t, err) + // Nothing written anywhere. + assert.Empty(t, b.writes) + assert.Empty(t, b.updates) + assert.Empty(t, tr.created) + assert.Empty(t, au.entries) +} + +func TestCaptureValidationRejectsBadTicket(t *testing.T) { + svc := newSvc(&fakeBrain{}, &fakeTracker{}, nil, fakePolicy{}, &fakeAudit{}) + _, err := svc.Capture(context.Background(), CaptureInput{ + Context: baseCtx(), + Tickets: []Ticket{{Repo: "hyperguild", Action: "frobnicate"}}, // bad action + }) + require.Error(t, err) + + _, err = svc.Capture(context.Background(), CaptureInput{ + Context: baseCtx(), + Tickets: []Ticket{{Repo: "hyperguild", Action: "close"}}, // close needs number + }) + require.Error(t, err) +} + +func TestCapturePartialFailureBestEffort(t *testing.T) { + b := &fakeBrain{failOn: func(n Note) error { + if strings.Contains(n.Content, "FAIL") { + return errors.New("disk full") + } + return nil + }} + tr := &fakeTracker{} + au := &fakeAudit{} + svc := newSvc(b, tr, nil, fakePolicy{}, au) + + rec, err := svc.Capture(context.Background(), CaptureInput{ + Context: baseCtx(), + Insights: []Insight{ + {Text: "good one", Wing: "hyperguild", Hall: "facts"}, + {Text: "FAIL here", Wing: "hyperguild", Hall: "facts"}, + }, + Tickets: []Ticket{{Repo: "hyperguild", Action: "create", Title: "t"}}, + }) + require.NoError(t, err, "partial failure is not a request-level error") + assert.True(t, rec.Insights[0].OK) + assert.False(t, rec.Insights[1].OK) + assert.True(t, rec.Tickets[0].OK, "ticket still persisted; no rollback") + require.Len(t, rec.Errors, 1) + assert.Equal(t, "insight[1]", rec.Errors[0].Item) + // Audit reflects exactly what landed: 1 insight + 1 ticket. + require.Len(t, au.entries, 1) + assert.Len(t, au.entries[0].Items, 2) +} + +func TestCaptureDryRunWritesNothing(t *testing.T) { + b := &fakeBrain{} + tr := &fakeTracker{} + au := &fakeAudit{} + svc := newSvc(b, tr, nil, fakePolicy{}, au) + + rec, err := svc.Capture(context.Background(), CaptureInput{ + Context: baseCtx(), + DryRun: true, + Insights: []Insight{{Text: "a", Wing: "hyperguild", Hall: "facts"}}, + Tickets: []Ticket{{Repo: "hyperguild", Action: "create", Title: "t"}}, + }) + require.NoError(t, err) + assert.True(t, rec.DryRun) + assert.Len(t, rec.Insights, 1) + assert.True(t, rec.Insights[0].OK, "would-be receipt marks planned items ok") + // Nothing written anywhere, including audit. + assert.Empty(t, b.writes) + assert.Empty(t, tr.created) + assert.Empty(t, au.entries) +} + +func TestCaptureStricterClassificationWins(t *testing.T) { + // Caller declares internal; target wing tagged confidential → effective confidential + security event. + b := &fakeBrain{} + au := &fakeAudit{} + pol := fakePolicy{tags: map[string]classification.Level{"client-seb": classification.Confidential}} + svc := newSvc(b, &fakeTracker{}, nil, pol, au) + + ctx := baseCtx() + ctx.Classification = "internal" + rec, err := svc.Capture(context.Background(), CaptureInput{ + Context: ctx, + Insights: []Insight{{Text: "x", Wing: "client-seb", Hall: "facts"}}, + }) + require.NoError(t, err) + assert.Equal(t, "confidential", rec.EffectiveClassification) + require.Len(t, au.entries, 1) + assert.NotEmpty(t, au.entries[0].SecurityEvents, "under-declaration logged as security event") + assert.Equal(t, "confidential", au.entries[0].EffectiveClassification) +} + +func TestCaptureCallerRaisingSensitivityHonoured(t *testing.T) { + // Caller declares confidential; target internal → effective confidential, NOT a security event. + au := &fakeAudit{} + pol := fakePolicy{tags: map[string]classification.Level{"hyperguild": classification.Internal}} + svc := newSvc(&fakeBrain{}, &fakeTracker{}, nil, pol, au) + + ctx := baseCtx() + ctx.Classification = "confidential" + rec, err := svc.Capture(context.Background(), CaptureInput{ + Context: ctx, + Insights: []Insight{{Text: "x", Wing: "hyperguild", Hall: "facts"}}, + }) + require.NoError(t, err) + assert.Equal(t, "confidential", rec.EffectiveClassification) + assert.Empty(t, au.entries[0].SecurityEvents, "raising sensitivity is honoured, not flagged") +} + +func TestCaptureSummaryPathAndFidelity(t *testing.T) { + sw := &fakeSummary{} + svc := newSvc(&fakeBrain{}, &fakeTracker{}, sw, fakePolicy{}, &fakeAudit{}) + + ctx := baseCtx() + ctx.Fidelity = "transcript-parse" + ctx.SessionRef = "abc123def456" + rec, err := svc.Capture(context.Background(), CaptureInput{ + Context: ctx, + Summary: &Summary{Title: "Session Wrap", Body: "did stuff", ReposTouched: []string{"hyperguild"}}, + }) + require.NoError(t, err) + require.NotNil(t, rec.Summary) + assert.True(t, rec.Summary.OK) + require.Len(t, sw.paths, 1) + assert.True(t, strings.HasPrefix(sw.paths[0], "summaries/claude-code/2026-06/"), "path: %s", sw.paths[0]) + assert.Contains(t, sw.paths[0], "session-wrap") + assert.Contains(t, sw.content[0], "fidelity: transcript-parse", "fidelity stamped in frontmatter") +} From 0ac165cca37c31cd1aa7935c270a0fcf6d7a9864 Mon Sep 17 00:00:00 2001 From: Mathias Date: Mon, 22 Jun 2026 23:21:43 +0200 Subject: [PATCH 2/2] feat(brainstore): shared BrainStore impl; re-point MCP handlers (#51) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Extracts the #45 write/update/get logic + the wiki upkeep that must accompany a write (wing _index rebuild, cross-wing auto-tunnel, graph re-index) into a single concrete brainstore.Store implementing capture.BrainStore. The MCP brain_write/brain_update/brain_get handlers are re-pointed at it, so there is one implementation, not two — the DRY payoff #51 is named for. capture and MCP now share the exact same brain write path and read-after-write contract. The Server gains a *brainstore.Store, constructed in NewServer and given the graph store in WithGraph. Embedding refresh stays out-of-band (mtime-driven vectorstore.Sync), unchanged. Existing MCP brain_update/ brain_get/brain_write tests pass unmodified — behaviour and the {id, path, content_hash} response contract are preserved. Co-Authored-By: Claude Opus 4.8 (1M context) --- ingestion/internal/brainstore/store.go | 129 ++++++++++++++++++++ ingestion/internal/brainstore/store_test.go | 85 +++++++++++++ ingestion/internal/mcp/handlers.go | 81 ++++-------- ingestion/internal/mcp/server.go | 12 +- 4 files changed, 248 insertions(+), 59 deletions(-) create mode 100644 ingestion/internal/brainstore/store.go create mode 100644 ingestion/internal/brainstore/store_test.go diff --git a/ingestion/internal/brainstore/store.go b/ingestion/internal/brainstore/store.go new file mode 100644 index 0000000..3666217 --- /dev/null +++ b/ingestion/internal/brainstore/store.go @@ -0,0 +1,129 @@ +// Package brainstore is the concrete BrainStore: the single shared +// implementation of the #45 write/update/get verbs, used by BOTH the MCP +// handlers and the capture use-case so there is one implementation, not +// two (the Clean-Architecture / DRY payoff of #51). +// +// It composes the file-level primitives in package api (WriteNote, +// UpdateNote, ReadNote — the read-after-write contract) with the wiki +// upkeep that must accompany a write: wing _index rebuild, cross-wing +// auto-tunnel, and graph re-index. Embedding refresh is intentionally +// out-of-band (mtime-driven vectorstore.Sync) and not triggered here — +// see the brain note on out-of-band sync. +package brainstore + +import ( + "context" + "log/slog" + "strings" + + "github.com/mathiasbq/hyperguild/ingestion/internal/api" + "github.com/mathiasbq/hyperguild/ingestion/internal/brain" + "github.com/mathiasbq/hyperguild/ingestion/internal/capture" + "github.com/mathiasbq/hyperguild/ingestion/internal/graphsync" +) + +// Store implements capture.BrainStore against a brain directory on disk, +// optionally re-indexing each write into the knowledge graph. +type Store struct { + brainDir string + graph graphsync.Store // nil = graph re-index disabled +} + +// New constructs a Store bound to brainDir with graph indexing disabled. +func New(brainDir string) *Store { + return &Store{brainDir: brainDir} +} + +// WithGraph enables graph re-index on every write/update. nil disables it. +func (s *Store) WithGraph(g graphsync.Store) *Store { + s.graph = g + return s +} + +// Write creates a brain note and returns its read-after-write handle. +func (s *Store) Write(ctx context.Context, n capture.Note) (capture.Ref, error) { + relPath, err := api.WriteNote(s.brainDir, api.WriteNoteOptions{ + Content: n.Content, + Filename: n.Filename, + Type: n.Type, + Domain: n.Domain, + Wing: n.Wing, + Hall: n.Hall, + }) + if err != nil { + return capture.Ref{}, err + } + s.wikiUpkeep(relPath, n.Wing, n.Content) + s.indexInGraph(ctx, "brain_write", relPath) + + _, _, hash, _ := api.ReadNote(s.brainDir, relPath) + return capture.Ref{ID: relPath, Path: relPath, ContentHash: hash}, nil +} + +// Update supersedes an existing note in place. slug may be a bare slug +// (resolved against n.Wing/n.Hall) or a full brain-relative path (when it +// contains a slash). It never creates — a missing target is an error. +func (s *Store) Update(ctx context.Context, slug string, n capture.Note) (capture.Ref, error) { + opts := api.UpdateNoteOptions{Content: n.Content, Reason: n.Reason} + if strings.Contains(slug, "/") { + opts.Path = slug + } else { + opts.Wing, opts.Hall, opts.Slug = n.Wing, n.Hall, slug + } + + relPath, hash, _, err := api.UpdateNote(s.brainDir, opts) + if err != nil { + return capture.Ref{}, err + } + if wing := wingFromRelPath(relPath); wing != "" { + s.wikiUpkeep(relPath, wing, n.Content) + } + s.indexInGraph(ctx, "brain_update", relPath) + + return capture.Ref{ID: relPath, Path: relPath, ContentHash: hash, Superseded: true}, nil +} + +// Get fetches a note by id/path — the read-after-write confirmation +// primitive (a direct fetch, never a semantic query). +func (s *Store) Get(_ context.Context, id string) (capture.StoredNote, error) { + fm, body, hash, err := api.ReadNote(s.brainDir, id) + if err != nil { + return capture.StoredNote{}, err + } + return capture.StoredNote{ID: id, Path: id, ContentHash: hash, Frontmatter: fm, Body: body}, nil +} + +// wikiUpkeep rebuilds the wing _index and re-tunnels cross-wing matches +// when a note lands in the structured wiki. Both are best-effort: the +// note is already written, so a failure here is logged, not propagated. +func (s *Store) wikiUpkeep(relPath, wing, content string) { + if wing == "" { + return + } + if err := brain.BuildWingIndex(s.brainDir, wing); err != nil { + slog.Warn("brainstore: auto-index failed", "wing", wing, "err", err) + } + if err := brain.AutoTunnel(s.brainDir, relPath, content); err != nil { + slog.Warn("brainstore: auto-tunnel failed", "src", relPath, "err", err) + } +} + +// indexInGraph re-indexes a written doc into the graph, best-effort. +func (s *Store) indexInGraph(ctx context.Context, op, relPath string) { + if s.graph == nil || relPath == "" { + return + } + if err := graphsync.IndexDoc(ctx, s.graph, s.brainDir, relPath); err != nil { + slog.Warn(op+": graph index failed", "path", relPath, "err", err) + } +} + +// wingFromRelPath extracts the wing from a structured wiki path +// (wiki///.md). Returns "" for legacy/non-wiki paths. +func wingFromRelPath(relPath string) string { + parts := strings.Split(relPath, "/") + if len(parts) >= 4 && parts[0] == "wiki" { + return parts[1] + } + return "" +} diff --git a/ingestion/internal/brainstore/store_test.go b/ingestion/internal/brainstore/store_test.go new file mode 100644 index 0000000..e2a444b --- /dev/null +++ b/ingestion/internal/brainstore/store_test.go @@ -0,0 +1,85 @@ +package brainstore_test + +import ( + "context" + "os" + "path/filepath" + "testing" + + "github.com/mathiasbq/hyperguild/ingestion/internal/brainstore" + "github.com/mathiasbq/hyperguild/ingestion/internal/capture" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestStoreWriteReturnsHandle(t *testing.T) { + dir := t.TempDir() + s := brainstore.New(dir) + + ref, err := s.Write(context.Background(), capture.Note{ + Content: "# X\n\nbody\n", Filename: "x", Wing: "a", Hall: "facts", + }) + require.NoError(t, err) + assert.Equal(t, "wiki/a/facts/x.md", ref.Path) + assert.Equal(t, ref.Path, ref.ID) + assert.NotEmpty(t, ref.ContentHash) + assert.False(t, ref.Superseded) + + _, err = os.Stat(filepath.Join(dir, "wiki/a/facts/x.md")) + require.NoError(t, err) +} + +func TestStoreUpdateSupersedes(t *testing.T) { + dir := t.TempDir() + s := brainstore.New(dir) + _, err := s.Write(context.Background(), capture.Note{ + Content: "old\n", Filename: "n", Wing: "a", Hall: "facts", + }) + require.NoError(t, err) + + ref, err := s.Update(context.Background(), "n", capture.Note{ + Content: "new\n", Wing: "a", Hall: "facts", Reason: "changed", + }) + require.NoError(t, err) + assert.True(t, ref.Superseded) + assert.Equal(t, "wiki/a/facts/n.md", ref.Path) + + got, _ := os.ReadFile(filepath.Join(dir, "wiki/a/facts/n.md")) + assert.Contains(t, string(got), "new") + assert.Contains(t, string(got), "supersede_reason: changed") +} + +func TestStoreUpdateByFullPath(t *testing.T) { + dir := t.TempDir() + s := brainstore.New(dir) + _, err := s.Write(context.Background(), capture.Note{Content: "old\n", Filename: "n", Wing: "a", Hall: "facts"}) + require.NoError(t, err) + + ref, err := s.Update(context.Background(), "wiki/a/facts/n.md", capture.Note{Content: "fresh\n"}) + require.NoError(t, err) + assert.Equal(t, "wiki/a/facts/n.md", ref.Path) +} + +func TestStoreUpdateMissingErrors(t *testing.T) { + dir := t.TempDir() + s := brainstore.New(dir) + _, err := s.Update(context.Background(), "ghost", capture.Note{Content: "x\n", Wing: "a", Hall: "facts"}) + require.Error(t, err) + _, statErr := os.Stat(filepath.Join(dir, "wiki/a/facts/ghost.md")) + assert.True(t, os.IsNotExist(statErr), "update must not create") +} + +func TestStoreGetRoundTripsHash(t *testing.T) { + dir := t.TempDir() + s := brainstore.New(dir) + ref, err := s.Write(context.Background(), capture.Note{ + Content: "# Body\n\ntext\n", Filename: "n", Wing: "a", Hall: "facts", + }) + require.NoError(t, err) + + note, err := s.Get(context.Background(), ref.ID) + require.NoError(t, err) + assert.Equal(t, ref.ContentHash, note.ContentHash, "write→get hash round-trips") + assert.Equal(t, "a", note.Frontmatter["wing"]) + assert.Contains(t, note.Body, "# Body") +} diff --git a/ingestion/internal/mcp/handlers.go b/ingestion/internal/mcp/handlers.go index 8966b9b..12a9193 100644 --- a/ingestion/internal/mcp/handlers.go +++ b/ingestion/internal/mcp/handlers.go @@ -9,8 +9,8 @@ import ( "strings" "time" - "github.com/mathiasbq/hyperguild/ingestion/internal/api" "github.com/mathiasbq/hyperguild/ingestion/internal/brain" + "github.com/mathiasbq/hyperguild/ingestion/internal/capture" "github.com/mathiasbq/hyperguild/ingestion/internal/extract" "github.com/mathiasbq/hyperguild/ingestion/internal/graphsync" "github.com/mathiasbq/hyperguild/ingestion/internal/pipeline" @@ -219,7 +219,11 @@ func (s *Server) brainWrite(ctx context.Context, args json.RawMessage) (json.Raw if err := json.Unmarshal(args, &a); err != nil { return nil, fmt.Errorf("parse args: %w", err) } - relPath, err := api.WriteNote(s.brainDir, api.WriteNoteOptions{ + // Delegate to the shared BrainStore so write+index+tunnel+graph live in + // one implementation (capture uses the same store). The read-after-write + // handle {id, path, content_hash} comes back from the store; path is kept + // for backward compatibility. + ref, err := s.store.Write(ctx, capture.Note{ Content: a.Content, Filename: a.Filename, Type: a.Type, @@ -230,22 +234,7 @@ func (s *Server) brainWrite(ctx context.Context, args json.RawMessage) (json.Raw if err != nil { return nil, err } - // Auto-regenerate the wing _index.md when the write landed in the - // structured wiki, and auto-tunnel cross-wing matches. Both are - // best-effort: the note is already written. - if a.Wing != "" && a.Hall != "" { - if err := brain.BuildWingIndex(s.brainDir, a.Wing); err != nil { - slog.Warn("brain_write: auto-index failed", "wing", a.Wing, "err", err) - } - if err := brain.AutoTunnel(s.brainDir, relPath, a.Content); err != nil { - slog.Warn("brain_write: auto-tunnel failed", "src", relPath, "err", err) - } - } - s.indexInGraph(ctx, "brain_write", relPath) - // Read-after-write handle: id == relPath, content_hash == sha256 of - // the bytes just written. path is kept for backward compatibility. - _, _, hash, _ := api.ReadNote(s.brainDir, relPath) - return json.Marshal(map[string]string{"id": relPath, "path": relPath, "content_hash": hash}) + return json.Marshal(map[string]string{"id": ref.ID, "path": ref.Path, "content_hash": ref.ContentHash}) } type brainUpdateArgs struct { @@ -274,50 +263,26 @@ func (s *Server) brainUpdate(ctx context.Context, args json.RawMessage) (json.Ra return nil, fmt.Errorf("content is required") } - opts := api.UpdateNoteOptions{Content: a.Content, Reason: a.Reason} - switch { - case a.Path != "": - opts.Path = a.Path - case strings.Contains(a.Slug, "/"): - // slug carries a full path (issue #45: "slug ... OR full path"). - opts.Path = a.Slug - default: - opts.Wing, opts.Hall, opts.Slug = a.Wing, a.Hall, a.Slug + // path takes precedence over slug; the store treats any slug containing + // a slash as a full brain-relative path (issue #45: "slug ... OR path"). + slug := a.Slug + if a.Path != "" { + slug = a.Path } - - relPath, hash, _, err := api.UpdateNote(s.brainDir, opts) + ref, err := s.store.Update(ctx, slug, capture.Note{ + Content: a.Content, + Wing: a.Wing, + Hall: a.Hall, + Reason: a.Reason, + }) if err != nil { return nil, err } - - // Best-effort wiki upkeep, mirroring brain_write: rebuild the wing - // _index and re-tunnel cross-wing matches against the new body. Both - // are idempotent and never block — the note is already superseded. - if wing := wingFromRelPath(relPath); wing != "" { - if err := brain.BuildWingIndex(s.brainDir, wing); err != nil { - slog.Warn("brain_update: auto-index failed", "wing", wing, "err", err) - } - if err := brain.AutoTunnel(s.brainDir, relPath, a.Content); err != nil { - slog.Warn("brain_update: auto-tunnel failed", "src", relPath, "err", err) - } - } - s.indexInGraph(ctx, "brain_update", relPath) - return json.Marshal(map[string]any{ - "id": relPath, "path": relPath, "content_hash": hash, "superseded": true, + "id": ref.ID, "path": ref.Path, "content_hash": ref.ContentHash, "superseded": ref.Superseded, }) } -// wingFromRelPath extracts the wing segment from a structured wiki path -// (wiki///.md). Returns "" for legacy/non-wiki paths. -func wingFromRelPath(relPath string) string { - parts := strings.Split(relPath, "/") - if len(parts) >= 4 && parts[0] == "wiki" { - return parts[1] - } - return "" -} - type brainGetArgs struct { ID string `json:"id,omitempty"` Path string `json:"path,omitempty"` @@ -327,7 +292,7 @@ type brainGetArgs struct { // path — the de-facto handle). Read-only; the create-path read-after- // write primitive that lets callers confirm a write landed without a // lexical re-query. -func (s *Server) brainGet(_ context.Context, args json.RawMessage) (json.RawMessage, error) { +func (s *Server) brainGet(ctx context.Context, args json.RawMessage) (json.RawMessage, error) { var a brainGetArgs if err := json.Unmarshal(args, &a); err != nil { return nil, fmt.Errorf("parse args: %w", err) @@ -339,13 +304,13 @@ func (s *Server) brainGet(_ context.Context, args json.RawMessage) (json.RawMess if target == "" { return nil, fmt.Errorf("id or path is required") } - fm, body, hash, err := api.ReadNote(s.brainDir, target) + note, err := s.store.Get(ctx, target) if err != nil { return nil, err } return json.Marshal(map[string]any{ - "id": target, "path": target, "content_hash": hash, - "frontmatter": fm, "body": body, + "id": note.ID, "path": note.Path, "content_hash": note.ContentHash, + "frontmatter": note.Frontmatter, "body": note.Body, }) } diff --git a/ingestion/internal/mcp/server.go b/ingestion/internal/mcp/server.go index 0f9051d..b072db5 100644 --- a/ingestion/internal/mcp/server.go +++ b/ingestion/internal/mcp/server.go @@ -10,6 +10,7 @@ import ( "fmt" "net/http" + "github.com/mathiasbq/hyperguild/ingestion/internal/brainstore" "github.com/mathiasbq/hyperguild/ingestion/internal/graphstore" "github.com/mathiasbq/hyperguild/ingestion/internal/graphsync" "github.com/mathiasbq/hyperguild/ingestion/internal/pipeline" @@ -46,6 +47,7 @@ type Server struct { vector search.VectorSearcher // nil = BM25-only retrieval embedder search.Embedder // nil = BM25-only retrieval graph graphsync.Store // nil = brain_graph and GraphRAG augmentation disabled + store *brainstore.Store // shared brain write/update/get impl (also used by capture) } // NewServer constructs a Server bound to brainDir. pipelineCfg supplies the @@ -56,7 +58,13 @@ func NewServer(brainDir string, pipelineCfg *pipeline.Config, llm pipeline.Compl if pipelineCfg != nil { cfg = *pipelineCfg } - return &Server{brainDir: brainDir, pipeline: cfg, llm: llm, answerLLM: answerLLM} + return &Server{ + brainDir: brainDir, + pipeline: cfg, + llm: llm, + answerLLM: answerLLM, + store: brainstore.New(brainDir), + } } // WithReranker installs an opt-in cross-encoder reranker. When set, @@ -84,9 +92,11 @@ func (s *Server) WithHybridRetrieval(v search.VectorSearcher, e search.Embedder) func (s *Server) WithGraph(g *graphstore.PGStore) *Server { if g == nil { s.graph = nil + s.store.WithGraph(nil) return s } s.graph = g + s.store.WithGraph(g) return s }