Compare commits
18
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
76514215f4 | ||
|
|
7cf5bc221d | ||
|
|
f78a5474a5 | ||
|
|
c307b72bd5 | ||
|
|
38a2e91002 | ||
|
|
77f5e06d6b | ||
|
|
202212e8d5 | ||
|
|
b7938d4636 | ||
|
|
a1997838b0 | ||
|
|
77680c7445 | ||
|
|
d7a842f356 | ||
|
|
aad90f2dfe | ||
|
|
07fca9ee73 | ||
|
|
f6bf9b5f57 | ||
|
|
6606b38a76 | ||
|
|
4cfc98de56 | ||
|
|
0ac165cca3 | ||
|
|
43f92e3102 |
@@ -8,6 +8,7 @@ import (
|
|||||||
"net/http"
|
"net/http"
|
||||||
"net/url"
|
"net/url"
|
||||||
"os"
|
"os"
|
||||||
|
"path/filepath"
|
||||||
"strconv"
|
"strconv"
|
||||||
"strings"
|
"strings"
|
||||||
"time"
|
"time"
|
||||||
@@ -15,8 +16,13 @@ import (
|
|||||||
chassisauth "gitea.d-ma.be/mathias/mcp-chassis/auth"
|
chassisauth "gitea.d-ma.be/mathias/mcp-chassis/auth"
|
||||||
|
|
||||||
"github.com/mathiasbq/hyperguild/ingestion/internal/api"
|
"github.com/mathiasbq/hyperguild/ingestion/internal/api"
|
||||||
|
"github.com/mathiasbq/hyperguild/ingestion/internal/audit"
|
||||||
|
"github.com/mathiasbq/hyperguild/ingestion/internal/capture"
|
||||||
|
"github.com/mathiasbq/hyperguild/ingestion/internal/capturehttp"
|
||||||
|
"github.com/mathiasbq/hyperguild/ingestion/internal/classification"
|
||||||
"github.com/mathiasbq/hyperguild/ingestion/internal/claudewatcher"
|
"github.com/mathiasbq/hyperguild/ingestion/internal/claudewatcher"
|
||||||
"github.com/mathiasbq/hyperguild/ingestion/internal/embed"
|
"github.com/mathiasbq/hyperguild/ingestion/internal/embed"
|
||||||
|
"github.com/mathiasbq/hyperguild/ingestion/internal/gitea"
|
||||||
"github.com/mathiasbq/hyperguild/ingestion/internal/graphstore"
|
"github.com/mathiasbq/hyperguild/ingestion/internal/graphstore"
|
||||||
"github.com/mathiasbq/hyperguild/ingestion/internal/graphsync"
|
"github.com/mathiasbq/hyperguild/ingestion/internal/graphsync"
|
||||||
"github.com/mathiasbq/hyperguild/ingestion/internal/llm"
|
"github.com/mathiasbq/hyperguild/ingestion/internal/llm"
|
||||||
@@ -118,6 +124,47 @@ func envInt(key string, fallback int) int {
|
|||||||
return fallback
|
return fallback
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// buildAuditSink selects the capture audit sink. When BRAIN_LOKI_URL is
|
||||||
|
// set it builds the classification-aware DegradingSink (loki central +
|
||||||
|
// durable file buffer + optional ntfy) and starts the reconcile loop;
|
||||||
|
// otherwise it falls back to a plain slog sink. The buffer lives under the
|
||||||
|
// brain dir so it survives process restarts.
|
||||||
|
func buildAuditSink(ctx context.Context, brainDir string, logger *slog.Logger) capture.AuditSink {
|
||||||
|
lokiURL := os.Getenv("BRAIN_LOKI_URL")
|
||||||
|
central := audit.NewLokiCentral(lokiURL)
|
||||||
|
if central == nil {
|
||||||
|
logger.Info("capture audit: slog sink (BRAIN_LOKI_URL unset)")
|
||||||
|
return audit.NewSlogSink(logger)
|
||||||
|
}
|
||||||
|
buffer, err := audit.NewFileBuffer(filepath.Join(brainDir, ".audit-buffer", "capture.jsonl"))
|
||||||
|
if err != nil {
|
||||||
|
logger.Error("capture audit buffer init", "err", err)
|
||||||
|
os.Exit(1)
|
||||||
|
}
|
||||||
|
// Keep notifier as a nil interface (not a typed-nil) when unconfigured
|
||||||
|
// so DegradingSink/Reconcile skip it cleanly.
|
||||||
|
var notifier audit.Notifier
|
||||||
|
if n := audit.NewNtfyNotifier(os.Getenv("BRAIN_NTFY_URL"), os.Getenv("BRAIN_NTFY_TOKEN")); n != nil {
|
||||||
|
notifier = n
|
||||||
|
}
|
||||||
|
reconcileInterval := time.Duration(envInt("BRAIN_AUDIT_RECONCILE_INTERVAL", 60)) * time.Second
|
||||||
|
audit.StartReconcile(ctx, central, buffer, notifier, reconcileInterval)
|
||||||
|
logger.Info("capture audit: loki+buffer sink", "loki", lokiURL, "reconcile_s", int(reconcileInterval.Seconds()))
|
||||||
|
return audit.NewDegradingSink(central, buffer, notifier)
|
||||||
|
}
|
||||||
|
|
||||||
|
// splitList parses a comma-separated env value into a trimmed,
|
||||||
|
// empty-free slice. Used for the capture sovereign-principal allowlist.
|
||||||
|
func splitList(v string) []string {
|
||||||
|
var out []string
|
||||||
|
for _, p := range strings.Split(v, ",") {
|
||||||
|
if p = strings.TrimSpace(p); p != "" {
|
||||||
|
out = append(out, p)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return out
|
||||||
|
}
|
||||||
|
|
||||||
// systemHostname returns os.Hostname() with a "unknown" fallback so the
|
// systemHostname returns os.Hostname() with a "unknown" fallback so the
|
||||||
// caller never has to handle the rare error path.
|
// caller never has to handle the rare error path.
|
||||||
func systemHostname() string {
|
func systemHostname() string {
|
||||||
@@ -175,6 +222,15 @@ func main() {
|
|||||||
logger.Info("brain reranker configured", "url", rerankURL, "model", rerankModel)
|
logger.Info("brain reranker configured", "url", rerankURL, "model", rerankModel)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Gitea ticket tracker for the capture capability (#52). Token via env
|
||||||
|
// only — never logged or in argv. Both vars must be set to enable it;
|
||||||
|
// gitea.New returns nil otherwise, leaving ticket integration off.
|
||||||
|
giteaURL := envOr("BRAIN_GITEA_URL", "https://git.d-ma.be")
|
||||||
|
if tracker := gitea.New(giteaURL, os.Getenv("BRAIN_GITEA_TOKEN")); tracker != nil {
|
||||||
|
mcpSrv = mcpSrv.WithIssueTracker(tracker)
|
||||||
|
logger.Info("brain gitea tracker configured", "url", giteaURL)
|
||||||
|
}
|
||||||
|
|
||||||
// Hybrid retrieval (pgvector + nomic-embed-text). Both env vars must
|
// Hybrid retrieval (pgvector + nomic-embed-text). Both env vars must
|
||||||
// be set together for the path to wire on; otherwise BM25-only.
|
// be set together for the path to wire on; otherwise BM25-only.
|
||||||
var vectorStore *vectorstore.PGStore
|
var vectorStore *vectorstore.PGStore
|
||||||
@@ -339,6 +395,37 @@ func main() {
|
|||||||
|
|
||||||
mux.Handle("/mcp", chassisauth.BearerMiddleware(mcpToken, jwtValidator, "brain", resourceMetadataURL, mcpSrv))
|
mux.Handle("/mcp", chassisauth.BearerMiddleware(mcpToken, jwtValidator, "brain", resourceMetadataURL, mcpSrv))
|
||||||
|
|
||||||
|
// POST /capture (#53/#54): the uniform capture REST door. Needs a ticket
|
||||||
|
// tracker to file action items, so it only mounts when Gitea is
|
||||||
|
// configured. It reuses the MCP server's graph-wired brain store (one
|
||||||
|
// implementation), the classification tags for the I1 gate, and a
|
||||||
|
// classification-aware audit sink (loki + durable buffer + ntfy when
|
||||||
|
// BRAIN_LOKI_URL is set, else a plain slog sink). The handler does its
|
||||||
|
// own auth (static + JWT) because it needs the principal to derive the
|
||||||
|
// trust-zone origin — the chassis middleware hides it.
|
||||||
|
if tracker := mcpSrv.IssueTracker(); tracker != nil {
|
||||||
|
classCfg, cerr := classification.Load(brainDir)
|
||||||
|
if cerr != nil {
|
||||||
|
logger.Error("load classification config", "err", cerr)
|
||||||
|
os.Exit(1)
|
||||||
|
}
|
||||||
|
auditSink := buildAuditSink(ctx, brainDir, logger)
|
||||||
|
captureSvc := capture.NewService(
|
||||||
|
mcpSrv.BrainStore(), tracker, nil, classCfg, auditSink)
|
||||||
|
sovereign := splitList(os.Getenv("BRAIN_CAPTURE_SOVEREIGN_PRINCIPALS"))
|
||||||
|
resolver := capturehttp.NewOriginResolver(sovereign)
|
||||||
|
captureH := capturehttp.New(captureSvc, jwtValidator, mcpToken, "local-cli", resolver)
|
||||||
|
mux.Handle("POST /capture", captureH)
|
||||||
|
// Same use-case behind the MCP `capture` tool (#55 relay) so MCP-native
|
||||||
|
// harnesses (claude.ai, Crush, Pi, LLM Council) reach capture through
|
||||||
|
// the existing /mcp OAuth connector. mcpSrv is already wrapped above;
|
||||||
|
// WithCapture mutates the same instance, so the tool appears live.
|
||||||
|
mcpSrv.WithCapture(captureSvc, jwtValidator, mcpToken, "local-cli", resolver)
|
||||||
|
logger.Info("capture enabled (REST + MCP tool)", "sovereign_principals", len(sovereign))
|
||||||
|
} else {
|
||||||
|
logger.Info("capture endpoint disabled (BRAIN_GITEA_TOKEN unset)")
|
||||||
|
}
|
||||||
|
|
||||||
// Opt-in OAuth 2.0 client_credentials flow for claude.ai's custom-MCP
|
// Opt-in OAuth 2.0 client_credentials flow for claude.ai's custom-MCP
|
||||||
// integration UI, which has no static-Bearer field. Setting both
|
// integration UI, which has no static-Bearer field. Setting both
|
||||||
// OAUTH_CLIENT_ID and OAUTH_CLIENT_SECRET enables the token exchange;
|
// OAUTH_CLIENT_ID and OAUTH_CLIENT_SECRET enables the token exchange;
|
||||||
|
|||||||
@@ -0,0 +1,167 @@
|
|||||||
|
package audit
|
||||||
|
|
||||||
|
import (
|
||||||
|
"bufio"
|
||||||
|
"crypto/sha256"
|
||||||
|
"encoding/hex"
|
||||||
|
"encoding/json"
|
||||||
|
"fmt"
|
||||||
|
"os"
|
||||||
|
"path/filepath"
|
||||||
|
"sync"
|
||||||
|
|
||||||
|
"github.com/mathiasbq/hyperguild/ingestion/internal/capture"
|
||||||
|
)
|
||||||
|
|
||||||
|
// FileBuffer is a durable, restart-surviving audit buffer backed by a
|
||||||
|
// JSONL file: one {id, entry} record per line. It is the internal/public
|
||||||
|
// tier fallback when loki is unreachable. Confirm rewrites the file
|
||||||
|
// without the confirmed record, so a record is cleared only after its
|
||||||
|
// central write is confirmed.
|
||||||
|
//
|
||||||
|
// Access is serialised by a mutex; the buffer is low-throughput (only
|
||||||
|
// written during a loki outage), so a whole-file rewrite on Confirm is
|
||||||
|
// acceptable and keeps the on-disk format trivially correct.
|
||||||
|
type FileBuffer struct {
|
||||||
|
path string
|
||||||
|
mu sync.Mutex
|
||||||
|
}
|
||||||
|
|
||||||
|
type bufferLine struct {
|
||||||
|
ID string `json:"id"`
|
||||||
|
Entry capture.AuditEntry `json:"entry"`
|
||||||
|
}
|
||||||
|
|
||||||
|
// NewFileBuffer returns a buffer backed by path. The parent directory is
|
||||||
|
// created if needed. The file itself is created lazily on first Append.
|
||||||
|
func NewFileBuffer(path string) (*FileBuffer, error) {
|
||||||
|
if err := os.MkdirAll(filepath.Dir(path), 0o755); err != nil {
|
||||||
|
return nil, fmt.Errorf("create buffer dir: %w", err)
|
||||||
|
}
|
||||||
|
return &FileBuffer{path: path}, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// Writable reports whether the buffer file can be appended to. It probes
|
||||||
|
// by opening the file for append (creating it if absent) — the same
|
||||||
|
// operation Append performs — so Reserve's check matches Append's reality.
|
||||||
|
func (b *FileBuffer) Writable() error {
|
||||||
|
b.mu.Lock()
|
||||||
|
defer b.mu.Unlock()
|
||||||
|
f, err := os.OpenFile(b.path, os.O_CREATE|os.O_APPEND|os.O_WRONLY, 0o644)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
return f.Close()
|
||||||
|
}
|
||||||
|
|
||||||
|
// Append durably writes one audit record. The ID is derived from the
|
||||||
|
// content + timestamp so it is stable and unique per record.
|
||||||
|
func (b *FileBuffer) Append(e capture.AuditEntry) error {
|
||||||
|
b.mu.Lock()
|
||||||
|
defer b.mu.Unlock()
|
||||||
|
|
||||||
|
line := bufferLine{ID: recordID(e), Entry: e}
|
||||||
|
data, err := json.Marshal(line)
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("marshal buffer line: %w", err)
|
||||||
|
}
|
||||||
|
f, err := os.OpenFile(b.path, os.O_CREATE|os.O_APPEND|os.O_WRONLY, 0o644)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
defer func() { _ = f.Close() }()
|
||||||
|
if _, err := f.Write(append(data, '\n')); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
return f.Sync()
|
||||||
|
}
|
||||||
|
|
||||||
|
// Pending reads all buffered records. A missing file means none.
|
||||||
|
func (b *FileBuffer) Pending() ([]Buffered, error) {
|
||||||
|
b.mu.Lock()
|
||||||
|
defer b.mu.Unlock()
|
||||||
|
return b.readAllLocked()
|
||||||
|
}
|
||||||
|
|
||||||
|
func (b *FileBuffer) readAllLocked() ([]Buffered, error) {
|
||||||
|
f, err := os.Open(b.path)
|
||||||
|
if err != nil {
|
||||||
|
if os.IsNotExist(err) {
|
||||||
|
return nil, nil
|
||||||
|
}
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
defer func() { _ = f.Close() }()
|
||||||
|
|
||||||
|
var out []Buffered
|
||||||
|
sc := bufio.NewScanner(f)
|
||||||
|
sc.Buffer(make([]byte, 0, 64*1024), 1024*1024)
|
||||||
|
for sc.Scan() {
|
||||||
|
raw := sc.Bytes()
|
||||||
|
if len(raw) == 0 {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
var l bufferLine
|
||||||
|
if err := json.Unmarshal(raw, &l); err != nil {
|
||||||
|
return nil, fmt.Errorf("parse buffer line: %w", err)
|
||||||
|
}
|
||||||
|
out = append(out, Buffered(l))
|
||||||
|
}
|
||||||
|
return out, sc.Err()
|
||||||
|
}
|
||||||
|
|
||||||
|
// Confirm removes a single record after its central write is confirmed, by
|
||||||
|
// rewriting the file without it. Unknown IDs are a no-op.
|
||||||
|
func (b *FileBuffer) Confirm(id string) error {
|
||||||
|
b.mu.Lock()
|
||||||
|
defer b.mu.Unlock()
|
||||||
|
|
||||||
|
all, err := b.readAllLocked()
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
tmp := b.path + ".tmp"
|
||||||
|
f, err := os.OpenFile(tmp, os.O_CREATE|os.O_TRUNC|os.O_WRONLY, 0o644)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
w := bufio.NewWriter(f)
|
||||||
|
kept := 0
|
||||||
|
for _, rec := range all {
|
||||||
|
if rec.ID == id {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
data, _ := json.Marshal(bufferLine(rec))
|
||||||
|
if _, err := w.Write(append(data, '\n')); err != nil {
|
||||||
|
_ = f.Close()
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
kept++
|
||||||
|
}
|
||||||
|
if err := w.Flush(); err != nil {
|
||||||
|
_ = f.Close()
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if err := f.Sync(); err != nil {
|
||||||
|
_ = f.Close()
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if err := f.Close(); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
// Empty buffer → remove the file entirely so Pending sees nothing.
|
||||||
|
if kept == 0 {
|
||||||
|
_ = os.Remove(tmp)
|
||||||
|
return os.Remove(b.path)
|
||||||
|
}
|
||||||
|
return os.Rename(tmp, b.path)
|
||||||
|
}
|
||||||
|
|
||||||
|
// recordID is a stable per-record identifier: sha256 of the principal,
|
||||||
|
// timestamp, and item list. Distinct captures never collide; the same
|
||||||
|
// buffered record always hashes the same.
|
||||||
|
func recordID(e capture.AuditEntry) string {
|
||||||
|
h := sha256.New()
|
||||||
|
_, _ = fmt.Fprintf(h, "%s|%s|%v|%s", e.Principal, e.Timestamp.UTC().Format("2006-01-02T15:04:05.000000000Z07:00"), e.Items, e.SessionRef)
|
||||||
|
return hex.EncodeToString(h.Sum(nil))[:16]
|
||||||
|
}
|
||||||
@@ -0,0 +1,96 @@
|
|||||||
|
package audit
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"fmt"
|
||||||
|
|
||||||
|
"github.com/mathiasbq/hyperguild/ingestion/internal/capture"
|
||||||
|
"github.com/mathiasbq/hyperguild/ingestion/internal/classification"
|
||||||
|
)
|
||||||
|
|
||||||
|
// Central is the central audit substrate (loki). Ready is a cheap
|
||||||
|
// reachability probe used by the pre-write reserve; Push writes a record.
|
||||||
|
type Central interface {
|
||||||
|
Ready(ctx context.Context) error
|
||||||
|
Push(ctx context.Context, e capture.AuditEntry) error
|
||||||
|
}
|
||||||
|
|
||||||
|
// Buffer is the durable local fallback for internal/public-tier records
|
||||||
|
// when the central sink is unreachable. It must survive process restart.
|
||||||
|
type Buffer interface {
|
||||||
|
// Writable reports whether the buffer can currently be appended to.
|
||||||
|
Writable() error
|
||||||
|
Append(e capture.AuditEntry) error
|
||||||
|
// Pending returns buffered records awaiting reconciliation, each with a
|
||||||
|
// stable ID used to Confirm (delete) it after a confirmed central write.
|
||||||
|
Pending() ([]Buffered, error)
|
||||||
|
Confirm(id string) error
|
||||||
|
}
|
||||||
|
|
||||||
|
// Buffered is a buffered audit record plus its stable buffer ID.
|
||||||
|
type Buffered struct {
|
||||||
|
ID string
|
||||||
|
Entry capture.AuditEntry
|
||||||
|
}
|
||||||
|
|
||||||
|
// Notifier raises an out-of-band alert (ntfy) about a degraded state.
|
||||||
|
type Notifier interface {
|
||||||
|
Notify(ctx context.Context, msg string) error
|
||||||
|
}
|
||||||
|
|
||||||
|
// DegradingSink is the classification-aware AuditSink (§4.4):
|
||||||
|
//
|
||||||
|
// - central reachable → AuditCentral (all tiers).
|
||||||
|
// - central down + confidential → refuse (no buffer): confidential must
|
||||||
|
// be centrally auditable at write time.
|
||||||
|
// - central down + internal/public + buffer writable → AuditBuffered.
|
||||||
|
// - central down + (confidential, or buffer not writable) → refuse (floor).
|
||||||
|
//
|
||||||
|
// The decision is made in Reserve, before any write; Record then executes it.
|
||||||
|
type DegradingSink struct {
|
||||||
|
central Central
|
||||||
|
buffer Buffer
|
||||||
|
notifier Notifier
|
||||||
|
}
|
||||||
|
|
||||||
|
// NewDegradingSink wires the central sink, durable buffer, and notifier.
|
||||||
|
func NewDegradingSink(central Central, buffer Buffer, notifier Notifier) *DegradingSink {
|
||||||
|
return &DegradingSink{central: central, buffer: buffer, notifier: notifier}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Reserve decides, before any write, how the capture will be audited — or
|
||||||
|
// returns an error to refuse it.
|
||||||
|
func (d *DegradingSink) Reserve(ctx context.Context, level classification.Level) (capture.AuditOutcome, error) {
|
||||||
|
if err := d.central.Ready(ctx); err == nil {
|
||||||
|
return capture.AuditCentral, nil
|
||||||
|
}
|
||||||
|
// Central sink is down.
|
||||||
|
if level == classification.Confidential {
|
||||||
|
return 0, fmt.Errorf("confidential capture requires the central audit sink, which is unreachable")
|
||||||
|
}
|
||||||
|
if err := d.buffer.Writable(); err != nil {
|
||||||
|
// Floor: neither central nor local buffer can record the audit.
|
||||||
|
return 0, fmt.Errorf("audit floor: central sink down and local buffer unwritable: %w", err)
|
||||||
|
}
|
||||||
|
return capture.AuditBuffered, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// Record persists the entry per the reserved outcome. For AuditBuffered it
|
||||||
|
// also fires the degraded-state alert.
|
||||||
|
func (d *DegradingSink) Record(ctx context.Context, e capture.AuditEntry, outcome capture.AuditOutcome) error {
|
||||||
|
switch outcome {
|
||||||
|
case capture.AuditBuffered:
|
||||||
|
if err := d.buffer.Append(e); err != nil {
|
||||||
|
return fmt.Errorf("buffer audit record: %w", err)
|
||||||
|
}
|
||||||
|
// Best-effort alert; the record is already durably buffered.
|
||||||
|
if d.notifier != nil {
|
||||||
|
_ = d.notifier.Notify(ctx, fmt.Sprintf(
|
||||||
|
"capture audit BUFFERED LOCALLY (loki unreachable) — principal=%s class=%s items=%d",
|
||||||
|
e.Principal, e.EffectiveClassification, len(e.Items)))
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
default:
|
||||||
|
return d.central.Push(ctx, e)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,191 @@
|
|||||||
|
package audit_test
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"errors"
|
||||||
|
"path/filepath"
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"github.com/mathiasbq/hyperguild/ingestion/internal/audit"
|
||||||
|
"github.com/mathiasbq/hyperguild/ingestion/internal/capture"
|
||||||
|
"github.com/mathiasbq/hyperguild/ingestion/internal/classification"
|
||||||
|
"github.com/stretchr/testify/assert"
|
||||||
|
"github.com/stretchr/testify/require"
|
||||||
|
)
|
||||||
|
|
||||||
|
// --- fakes ---
|
||||||
|
|
||||||
|
type fakeCentral struct {
|
||||||
|
down bool
|
||||||
|
pushed []capture.AuditEntry
|
||||||
|
pushErr error
|
||||||
|
}
|
||||||
|
|
||||||
|
func (f *fakeCentral) Ready(context.Context) error {
|
||||||
|
if f.down {
|
||||||
|
return errors.New("loki down")
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (f *fakeCentral) Push(_ context.Context, e capture.AuditEntry) error {
|
||||||
|
if f.pushErr != nil {
|
||||||
|
return f.pushErr
|
||||||
|
}
|
||||||
|
f.pushed = append(f.pushed, e)
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
type fakeNotifier struct{ msgs []string }
|
||||||
|
|
||||||
|
func (f *fakeNotifier) Notify(_ context.Context, msg string) error {
|
||||||
|
f.msgs = append(f.msgs, msg)
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// unwritableBuffer always reports it cannot be written (floor condition).
|
||||||
|
type unwritableBuffer struct{}
|
||||||
|
|
||||||
|
func (unwritableBuffer) Writable() error { return errors.New("disk full") }
|
||||||
|
func (unwritableBuffer) Append(capture.AuditEntry) error { return errors.New("disk full") }
|
||||||
|
func (unwritableBuffer) Pending() ([]audit.Buffered, error) { return nil, nil }
|
||||||
|
func (unwritableBuffer) Confirm(string) error { return nil }
|
||||||
|
|
||||||
|
func newFileBuffer(t *testing.T) *audit.FileBuffer {
|
||||||
|
t.Helper()
|
||||||
|
b, err := audit.NewFileBuffer(filepath.Join(t.TempDir(), "audit-buffer.jsonl"))
|
||||||
|
require.NoError(t, err)
|
||||||
|
return b
|
||||||
|
}
|
||||||
|
|
||||||
|
func entry(principal string) capture.AuditEntry {
|
||||||
|
return capture.AuditEntry{Principal: principal, EffectiveClassification: "internal", Items: []string{"insight:x"}}
|
||||||
|
}
|
||||||
|
|
||||||
|
// --- Reserve: classification-aware decision ---
|
||||||
|
|
||||||
|
func TestReserveCentralUpGrantsCentral(t *testing.T) {
|
||||||
|
d := audit.NewDegradingSink(&fakeCentral{}, newFileBuffer(t), &fakeNotifier{})
|
||||||
|
for _, lvl := range []classification.Level{classification.Public, classification.Internal, classification.Confidential} {
|
||||||
|
out, err := d.Reserve(context.Background(), lvl)
|
||||||
|
require.NoError(t, err)
|
||||||
|
assert.Equal(t, capture.AuditCentral, out)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestReserveConfidentialSinkDownRefuses(t *testing.T) {
|
||||||
|
d := audit.NewDegradingSink(&fakeCentral{down: true}, newFileBuffer(t), &fakeNotifier{})
|
||||||
|
_, err := d.Reserve(context.Background(), classification.Confidential)
|
||||||
|
require.Error(t, err, "confidential + sink down → refuse, no buffer")
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestReserveInternalSinkDownBuffers(t *testing.T) {
|
||||||
|
d := audit.NewDegradingSink(&fakeCentral{down: true}, newFileBuffer(t), &fakeNotifier{})
|
||||||
|
out, err := d.Reserve(context.Background(), classification.Internal)
|
||||||
|
require.NoError(t, err)
|
||||||
|
assert.Equal(t, capture.AuditBuffered, out)
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestReserveFloorRefusesWhenNothingCanRecord(t *testing.T) {
|
||||||
|
d := audit.NewDegradingSink(&fakeCentral{down: true}, unwritableBuffer{}, &fakeNotifier{})
|
||||||
|
_, err := d.Reserve(context.Background(), classification.Internal)
|
||||||
|
require.Error(t, err, "central down AND buffer unwritable → floor refuse")
|
||||||
|
}
|
||||||
|
|
||||||
|
// --- Record: executes the reserved outcome ---
|
||||||
|
|
||||||
|
func TestRecordCentralPushes(t *testing.T) {
|
||||||
|
c := &fakeCentral{}
|
||||||
|
d := audit.NewDegradingSink(c, newFileBuffer(t), &fakeNotifier{})
|
||||||
|
require.NoError(t, d.Record(context.Background(), entry("p"), capture.AuditCentral))
|
||||||
|
assert.Len(t, c.pushed, 1)
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestRecordBufferedAppendsAndNotifies(t *testing.T) {
|
||||||
|
buf := newFileBuffer(t)
|
||||||
|
nt := &fakeNotifier{}
|
||||||
|
d := audit.NewDegradingSink(&fakeCentral{down: true}, buf, nt)
|
||||||
|
require.NoError(t, d.Record(context.Background(), entry("p"), capture.AuditBuffered))
|
||||||
|
|
||||||
|
pending, err := buf.Pending()
|
||||||
|
require.NoError(t, err)
|
||||||
|
assert.Len(t, pending, 1)
|
||||||
|
assert.NotEmpty(t, nt.msgs, "degraded state alerts via ntfy")
|
||||||
|
}
|
||||||
|
|
||||||
|
// --- FileBuffer durability + Confirm ---
|
||||||
|
|
||||||
|
func TestFileBufferSurvivesRestart(t *testing.T) {
|
||||||
|
path := filepath.Join(t.TempDir(), "buf.jsonl")
|
||||||
|
b1, err := audit.NewFileBuffer(path)
|
||||||
|
require.NoError(t, err)
|
||||||
|
require.NoError(t, b1.Append(entry("p1")))
|
||||||
|
require.NoError(t, b1.Append(entry("p2")))
|
||||||
|
|
||||||
|
// "restart": a fresh FileBuffer over the same file sees the records.
|
||||||
|
b2, err := audit.NewFileBuffer(path)
|
||||||
|
require.NoError(t, err)
|
||||||
|
pending, err := b2.Pending()
|
||||||
|
require.NoError(t, err)
|
||||||
|
assert.Len(t, pending, 2)
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestFileBufferConfirmRemovesOnlyThatRecord(t *testing.T) {
|
||||||
|
buf := newFileBuffer(t)
|
||||||
|
require.NoError(t, buf.Append(entry("keep")))
|
||||||
|
require.NoError(t, buf.Append(entry("drop")))
|
||||||
|
|
||||||
|
pending, _ := buf.Pending()
|
||||||
|
require.Len(t, pending, 2)
|
||||||
|
var dropID string
|
||||||
|
for _, p := range pending {
|
||||||
|
if p.Entry.Principal == "drop" {
|
||||||
|
dropID = p.ID
|
||||||
|
}
|
||||||
|
}
|
||||||
|
require.NoError(t, buf.Confirm(dropID))
|
||||||
|
|
||||||
|
after, _ := buf.Pending()
|
||||||
|
require.Len(t, after, 1)
|
||||||
|
assert.Equal(t, "keep", after[0].Entry.Principal)
|
||||||
|
}
|
||||||
|
|
||||||
|
// --- Reconcile ---
|
||||||
|
|
||||||
|
func TestReconcileReplaysAndClearsOnlyAfterConfirmedWrite(t *testing.T) {
|
||||||
|
buf := newFileBuffer(t)
|
||||||
|
require.NoError(t, buf.Append(entry("a")))
|
||||||
|
require.NoError(t, buf.Append(entry("b")))
|
||||||
|
c := &fakeCentral{} // up
|
||||||
|
nt := &fakeNotifier{}
|
||||||
|
|
||||||
|
n, err := audit.Reconcile(context.Background(), c, buf, nt)
|
||||||
|
require.NoError(t, err)
|
||||||
|
assert.Equal(t, 2, n)
|
||||||
|
assert.Len(t, c.pushed, 2, "buffered records replayed to central")
|
||||||
|
|
||||||
|
pending, _ := buf.Pending()
|
||||||
|
assert.Empty(t, pending, "buffer cleared after confirmed central writes")
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestReconcileNoopWhenCentralDown(t *testing.T) {
|
||||||
|
buf := newFileBuffer(t)
|
||||||
|
require.NoError(t, buf.Append(entry("a")))
|
||||||
|
n, err := audit.Reconcile(context.Background(), &fakeCentral{down: true}, buf, &fakeNotifier{})
|
||||||
|
require.NoError(t, err)
|
||||||
|
assert.Equal(t, 0, n)
|
||||||
|
pending, _ := buf.Pending()
|
||||||
|
assert.Len(t, pending, 1, "records stay buffered while central is down")
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestReconcileKeepsRecordWhenPushFails(t *testing.T) {
|
||||||
|
buf := newFileBuffer(t)
|
||||||
|
require.NoError(t, buf.Append(entry("a")))
|
||||||
|
// Ready ok but Push fails → record must remain buffered (not lost).
|
||||||
|
c := &fakeCentral{pushErr: errors.New("push rejected")}
|
||||||
|
n, err := audit.Reconcile(context.Background(), c, buf, &fakeNotifier{})
|
||||||
|
require.NoError(t, err)
|
||||||
|
assert.Equal(t, 0, n)
|
||||||
|
pending, _ := buf.Pending()
|
||||||
|
assert.Len(t, pending, 1)
|
||||||
|
}
|
||||||
@@ -0,0 +1,99 @@
|
|||||||
|
package audit
|
||||||
|
|
||||||
|
import (
|
||||||
|
"bytes"
|
||||||
|
"context"
|
||||||
|
"encoding/json"
|
||||||
|
"fmt"
|
||||||
|
"net/http"
|
||||||
|
"strconv"
|
||||||
|
"strings"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/mathiasbq/hyperguild/ingestion/internal/capture"
|
||||||
|
)
|
||||||
|
|
||||||
|
// LokiCentral pushes capture audit records to a Grafana Loki instance via
|
||||||
|
// its push API, and probes readiness via /ready. It is the central audit
|
||||||
|
// substrate behind DegradingSink.
|
||||||
|
type LokiCentral struct {
|
||||||
|
baseURL string
|
||||||
|
labels map[string]string
|
||||||
|
http *http.Client
|
||||||
|
}
|
||||||
|
|
||||||
|
// NewLokiCentral constructs a LokiCentral for the given base URL (e.g.
|
||||||
|
// http://loki:3100). Returns nil when baseURL is empty so callers can
|
||||||
|
// treat missing config as "no central sink" with a single nil check.
|
||||||
|
func NewLokiCentral(baseURL string) *LokiCentral {
|
||||||
|
if baseURL == "" {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
return &LokiCentral{
|
||||||
|
baseURL: strings.TrimRight(baseURL, "/"),
|
||||||
|
labels: map[string]string{"service": "brain-capture", "kind": "audit"},
|
||||||
|
http: &http.Client{Timeout: 10 * time.Second},
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Ready probes Loki's readiness endpoint.
|
||||||
|
func (l *LokiCentral) Ready(ctx context.Context) error {
|
||||||
|
req, err := http.NewRequestWithContext(ctx, http.MethodGet, l.baseURL+"/ready", nil)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
resp, err := l.http.Do(req)
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("loki not ready: %w", err)
|
||||||
|
}
|
||||||
|
defer func() { _ = resp.Body.Close() }()
|
||||||
|
if resp.StatusCode != http.StatusOK {
|
||||||
|
return fmt.Errorf("loki not ready: status %d", resp.StatusCode)
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// pushPayload is the Loki push API body: one stream, one entry whose line
|
||||||
|
// is the JSON-encoded audit record.
|
||||||
|
type pushPayload struct {
|
||||||
|
Streams []lokiStream `json:"streams"`
|
||||||
|
}
|
||||||
|
|
||||||
|
type lokiStream struct {
|
||||||
|
Stream map[string]string `json:"stream"`
|
||||||
|
Values [][2]string `json:"values"`
|
||||||
|
}
|
||||||
|
|
||||||
|
// Push writes one audit record to Loki as a structured log line.
|
||||||
|
func (l *LokiCentral) Push(ctx context.Context, e capture.AuditEntry) error {
|
||||||
|
line, err := json.Marshal(e)
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("marshal audit entry: %w", err)
|
||||||
|
}
|
||||||
|
ts := e.Timestamp
|
||||||
|
if ts.IsZero() {
|
||||||
|
ts = time.Now()
|
||||||
|
}
|
||||||
|
body, err := json.Marshal(pushPayload{Streams: []lokiStream{{
|
||||||
|
Stream: l.labels,
|
||||||
|
Values: [][2]string{{strconv.FormatInt(ts.UTC().UnixNano(), 10), string(line)}},
|
||||||
|
}}})
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
req, err := http.NewRequestWithContext(ctx, http.MethodPost,
|
||||||
|
l.baseURL+"/loki/api/v1/push", bytes.NewReader(body))
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
req.Header.Set("Content-Type", "application/json")
|
||||||
|
resp, err := l.http.Do(req)
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("loki push: %w", err)
|
||||||
|
}
|
||||||
|
defer func() { _ = resp.Body.Close() }()
|
||||||
|
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
|
||||||
|
return fmt.Errorf("loki push: status %d", resp.StatusCode)
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
@@ -0,0 +1,86 @@
|
|||||||
|
package audit_test
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"io"
|
||||||
|
"net/http"
|
||||||
|
"net/http/httptest"
|
||||||
|
"strings"
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"github.com/mathiasbq/hyperguild/ingestion/internal/audit"
|
||||||
|
"github.com/mathiasbq/hyperguild/ingestion/internal/capture"
|
||||||
|
"github.com/stretchr/testify/assert"
|
||||||
|
"github.com/stretchr/testify/require"
|
||||||
|
)
|
||||||
|
|
||||||
|
func TestLokiReadyAndPush(t *testing.T) {
|
||||||
|
var pushBody string
|
||||||
|
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||||
|
switch r.URL.Path {
|
||||||
|
case "/ready":
|
||||||
|
w.WriteHeader(http.StatusOK)
|
||||||
|
case "/loki/api/v1/push":
|
||||||
|
b, _ := io.ReadAll(r.Body)
|
||||||
|
pushBody = string(b)
|
||||||
|
w.WriteHeader(http.StatusNoContent)
|
||||||
|
default:
|
||||||
|
w.WriteHeader(http.StatusNotFound)
|
||||||
|
}
|
||||||
|
}))
|
||||||
|
defer srv.Close()
|
||||||
|
|
||||||
|
c := audit.NewLokiCentral(srv.URL)
|
||||||
|
require.NotNil(t, c)
|
||||||
|
require.NoError(t, c.Ready(context.Background()))
|
||||||
|
|
||||||
|
err := c.Push(context.Background(), capture.AuditEntry{
|
||||||
|
Principal: "koala-cli", EffectiveClassification: "internal", Items: []string{"insight:x"},
|
||||||
|
})
|
||||||
|
require.NoError(t, err)
|
||||||
|
assert.Contains(t, pushBody, "streams")
|
||||||
|
assert.Contains(t, pushBody, "koala-cli", "audit entry serialised into the loki line")
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestLokiReadyFailsWhenDown(t *testing.T) {
|
||||||
|
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
|
||||||
|
w.WriteHeader(http.StatusServiceUnavailable)
|
||||||
|
}))
|
||||||
|
defer srv.Close()
|
||||||
|
require.Error(t, audit.NewLokiCentral(srv.URL).Ready(context.Background()))
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestLokiNilWhenUnconfigured(t *testing.T) {
|
||||||
|
assert.Nil(t, audit.NewLokiCentral(""))
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestNtfyNotify(t *testing.T) {
|
||||||
|
var gotBody, gotAuth string
|
||||||
|
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||||
|
b, _ := io.ReadAll(r.Body)
|
||||||
|
gotBody = string(b)
|
||||||
|
gotAuth = r.Header.Get("Authorization")
|
||||||
|
w.WriteHeader(http.StatusOK)
|
||||||
|
}))
|
||||||
|
defer srv.Close()
|
||||||
|
|
||||||
|
n := audit.NewNtfyNotifier(srv.URL, "ntfy-token")
|
||||||
|
require.NotNil(t, n)
|
||||||
|
require.NoError(t, n.Notify(context.Background(), "audit buffered locally"))
|
||||||
|
assert.Contains(t, gotBody, "audit buffered locally")
|
||||||
|
assert.Equal(t, "Bearer ntfy-token", gotAuth)
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestNtfyDoesNotLeakTokenOnError(t *testing.T) {
|
||||||
|
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
|
||||||
|
w.WriteHeader(http.StatusInternalServerError)
|
||||||
|
}))
|
||||||
|
defer srv.Close()
|
||||||
|
err := audit.NewNtfyNotifier(srv.URL, "secret-token").Notify(context.Background(), "x")
|
||||||
|
require.Error(t, err)
|
||||||
|
assert.False(t, strings.Contains(err.Error(), "secret-token"), "token must not leak into errors")
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestNtfyNilWhenUnconfigured(t *testing.T) {
|
||||||
|
assert.Nil(t, audit.NewNtfyNotifier("", "tok"))
|
||||||
|
}
|
||||||
@@ -0,0 +1,55 @@
|
|||||||
|
package audit
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"fmt"
|
||||||
|
"net/http"
|
||||||
|
"strings"
|
||||||
|
"time"
|
||||||
|
)
|
||||||
|
|
||||||
|
// NtfyNotifier posts alerts to an ntfy topic URL. Used to surface a
|
||||||
|
// degraded audit state (records buffered locally during a loki outage).
|
||||||
|
type NtfyNotifier struct {
|
||||||
|
topicURL string
|
||||||
|
token string
|
||||||
|
http *http.Client
|
||||||
|
}
|
||||||
|
|
||||||
|
// NewNtfyNotifier constructs a notifier for the given ntfy topic URL
|
||||||
|
// (e.g. https://ntfy.sh/my-topic). token is an optional bearer for
|
||||||
|
// protected ntfy instances; it is held here and only sent in the
|
||||||
|
// Authorization header, never logged. Returns nil when topicURL is empty.
|
||||||
|
func NewNtfyNotifier(topicURL, token string) *NtfyNotifier {
|
||||||
|
if topicURL == "" {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
return &NtfyNotifier{
|
||||||
|
topicURL: strings.TrimRight(topicURL, "/"),
|
||||||
|
token: token,
|
||||||
|
http: &http.Client{Timeout: 10 * time.Second},
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Notify posts a message to the ntfy topic.
|
||||||
|
func (n *NtfyNotifier) Notify(ctx context.Context, msg string) error {
|
||||||
|
req, err := http.NewRequestWithContext(ctx, http.MethodPost, n.topicURL, strings.NewReader(msg))
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
req.Header.Set("Title", "brain-capture audit degraded")
|
||||||
|
req.Header.Set("Priority", "high")
|
||||||
|
req.Header.Set("Tags", "warning,brain")
|
||||||
|
if n.token != "" {
|
||||||
|
req.Header.Set("Authorization", "Bearer "+n.token)
|
||||||
|
}
|
||||||
|
resp, err := n.http.Do(req)
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("ntfy notify: %w", err)
|
||||||
|
}
|
||||||
|
defer func() { _ = resp.Body.Close() }()
|
||||||
|
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
|
||||||
|
return fmt.Errorf("ntfy notify: status %d", resp.StatusCode)
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
@@ -0,0 +1,65 @@
|
|||||||
|
package audit
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"fmt"
|
||||||
|
"log/slog"
|
||||||
|
"time"
|
||||||
|
)
|
||||||
|
|
||||||
|
// Reconcile replays locally-buffered audit records to the central sink
|
||||||
|
// when it is reachable again. A record is removed from the buffer ONLY
|
||||||
|
// after its central write is confirmed, so a crash mid-reconcile re-plays
|
||||||
|
// rather than loses. Returns the number of records reconciled.
|
||||||
|
//
|
||||||
|
// A no-op (0, nil) when the central sink is still unreachable or the
|
||||||
|
// buffer is empty.
|
||||||
|
func Reconcile(ctx context.Context, central Central, buffer Buffer, notifier Notifier) (int, error) {
|
||||||
|
if err := central.Ready(ctx); err != nil {
|
||||||
|
return 0, nil // still down; try again next tick
|
||||||
|
}
|
||||||
|
pending, err := buffer.Pending()
|
||||||
|
if err != nil {
|
||||||
|
return 0, fmt.Errorf("read buffer: %w", err)
|
||||||
|
}
|
||||||
|
reconciled := 0
|
||||||
|
for _, rec := range pending {
|
||||||
|
if err := central.Push(ctx, rec.Entry); err != nil {
|
||||||
|
// Central went away mid-drain; stop and keep the rest buffered.
|
||||||
|
break
|
||||||
|
}
|
||||||
|
if err := buffer.Confirm(rec.ID); err != nil {
|
||||||
|
return reconciled, fmt.Errorf("confirm buffered record %s: %w", rec.ID, err)
|
||||||
|
}
|
||||||
|
reconciled++
|
||||||
|
}
|
||||||
|
if reconciled > 0 && notifier != nil {
|
||||||
|
_ = notifier.Notify(ctx, fmt.Sprintf("reconciled %d buffered capture audit record(s) to loki", reconciled))
|
||||||
|
}
|
||||||
|
return reconciled, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// StartReconcile runs Reconcile on a ticker until ctx is cancelled. It is
|
||||||
|
// the recovery half of the degrade-and-buffer path; pair it with a
|
||||||
|
// DegradingSink sharing the same buffer + central.
|
||||||
|
func StartReconcile(ctx context.Context, central Central, buffer Buffer, notifier Notifier, interval time.Duration) {
|
||||||
|
if interval <= 0 {
|
||||||
|
interval = time.Minute
|
||||||
|
}
|
||||||
|
go func() {
|
||||||
|
t := time.NewTicker(interval)
|
||||||
|
defer t.Stop()
|
||||||
|
for {
|
||||||
|
select {
|
||||||
|
case <-ctx.Done():
|
||||||
|
return
|
||||||
|
case <-t.C:
|
||||||
|
if n, err := Reconcile(ctx, central, buffer, notifier); err != nil {
|
||||||
|
slog.Warn("audit reconcile failed", "err", err)
|
||||||
|
} else if n > 0 {
|
||||||
|
slog.Info("audit reconcile", "reconciled", n)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
}
|
||||||
@@ -0,0 +1,56 @@
|
|||||||
|
// Package audit provides AuditSink implementations for the capture
|
||||||
|
// capability (I5). This file ships the minimal slog-backed sink used in
|
||||||
|
// #53: it emits the request-level audit record to structured logs, which
|
||||||
|
// the alloy/loki substrate already scrapes. The classification-aware
|
||||||
|
// degradation/refusal sink (confidential fails closed, internal buffers +
|
||||||
|
// reconciles) lands in #54 and replaces this behind the same interface.
|
||||||
|
package audit
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"log/slog"
|
||||||
|
|
||||||
|
"github.com/mathiasbq/hyperguild/ingestion/internal/capture"
|
||||||
|
"github.com/mathiasbq/hyperguild/ingestion/internal/classification"
|
||||||
|
)
|
||||||
|
|
||||||
|
// SlogSink records audit entries to an slog.Logger. It never fails and is
|
||||||
|
// always centrally available, so its Reserve always grants AuditCentral —
|
||||||
|
// it does not exercise the I5 degradation/floor. That is DegradingSink's
|
||||||
|
// job (loki + durable buffer). SlogSink is the default for deployments
|
||||||
|
// without a loki endpoint configured. A nil logger ⇒ slog.Default().
|
||||||
|
type SlogSink struct {
|
||||||
|
logger *slog.Logger
|
||||||
|
}
|
||||||
|
|
||||||
|
// NewSlogSink constructs a SlogSink. nil logger ⇒ slog.Default().
|
||||||
|
func NewSlogSink(logger *slog.Logger) *SlogSink {
|
||||||
|
if logger == nil {
|
||||||
|
logger = slog.Default()
|
||||||
|
}
|
||||||
|
return &SlogSink{logger: logger}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Reserve always grants central recording — slog is always available.
|
||||||
|
func (s *SlogSink) Reserve(_ context.Context, _ classification.Level) (capture.AuditOutcome, error) {
|
||||||
|
return capture.AuditCentral, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// Record emits the audit entry at info level. Security events, when
|
||||||
|
// present, are logged at warn level so they surface independently of the
|
||||||
|
// routine audit stream.
|
||||||
|
func (s *SlogSink) Record(_ context.Context, e capture.AuditEntry, _ capture.AuditOutcome) error {
|
||||||
|
s.logger.Info("capture audit",
|
||||||
|
"principal", e.Principal,
|
||||||
|
"actor", e.Actor,
|
||||||
|
"harness", e.Harness,
|
||||||
|
"session_ref", e.SessionRef,
|
||||||
|
"classification", e.EffectiveClassification,
|
||||||
|
"items", e.Items,
|
||||||
|
"ts", e.Timestamp,
|
||||||
|
)
|
||||||
|
for _, ev := range e.SecurityEvents {
|
||||||
|
s.logger.Warn("capture security event", "principal", e.Principal, "event", ev)
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
@@ -0,0 +1,41 @@
|
|||||||
|
package audit_test
|
||||||
|
|
||||||
|
import (
|
||||||
|
"bytes"
|
||||||
|
"context"
|
||||||
|
"log/slog"
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"github.com/mathiasbq/hyperguild/ingestion/internal/audit"
|
||||||
|
"github.com/mathiasbq/hyperguild/ingestion/internal/capture"
|
||||||
|
"github.com/stretchr/testify/assert"
|
||||||
|
"github.com/stretchr/testify/require"
|
||||||
|
)
|
||||||
|
|
||||||
|
func TestSlogSinkRecordsEntryAndSecurityEvents(t *testing.T) {
|
||||||
|
var buf bytes.Buffer
|
||||||
|
sink := audit.NewSlogSink(slog.New(slog.NewTextHandler(&buf, nil)))
|
||||||
|
|
||||||
|
err := sink.Record(context.Background(), capture.AuditEntry{
|
||||||
|
Principal: "koala-cli",
|
||||||
|
Harness: "claude-code",
|
||||||
|
EffectiveClassification: "confidential",
|
||||||
|
Items: []string{"insight:wiki/a/facts/x.md"},
|
||||||
|
SecurityEvents: []string{"asserted-vs-derived origin mismatch"},
|
||||||
|
}, capture.AuditCentral)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
out := buf.String()
|
||||||
|
assert.Contains(t, out, "capture audit")
|
||||||
|
assert.Contains(t, out, "koala-cli")
|
||||||
|
assert.Contains(t, out, "confidential")
|
||||||
|
assert.Contains(t, out, "capture security event")
|
||||||
|
assert.Contains(t, out, "asserted-vs-derived origin mismatch")
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestSlogSinkNilLoggerDefaults(t *testing.T) {
|
||||||
|
// nil logger must not panic.
|
||||||
|
require.NotPanics(t, func() {
|
||||||
|
_ = audit.NewSlogSink(nil).Record(context.Background(), capture.AuditEntry{}, capture.AuditCentral)
|
||||||
|
})
|
||||||
|
}
|
||||||
@@ -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/<wing>/<hall>/<slug>.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 ""
|
||||||
|
}
|
||||||
@@ -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")
|
||||||
|
}
|
||||||
@@ -0,0 +1,148 @@
|
|||||||
|
// 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
|
||||||
|
|
||||||
|
// Zone is the trust zone a capture originates from, server-derived from
|
||||||
|
// the authenticated principal (spec §4.2 / I1). It is NEVER taken from
|
||||||
|
// caller input — context.Harness is descriptive telemetry only.
|
||||||
|
type Zone int
|
||||||
|
|
||||||
|
const (
|
||||||
|
// ZoneUnknown means the origin was not set. The REST adapter always
|
||||||
|
// sets a concrete zone; the service treats Unknown as "not gated" (only
|
||||||
|
// an explicit ZoneUSNexus triggers the I1 refusal) so the gate can
|
||||||
|
// never fire on a caller-controllable default.
|
||||||
|
ZoneUnknown Zone = iota
|
||||||
|
// ZoneSovereign is sovereign soil (homelab / Tailscale CLI callers).
|
||||||
|
ZoneSovereign
|
||||||
|
// ZoneUSNexus is a non-sovereign US-jurisdiction surface (e.g.
|
||||||
|
// claude.ai). Confidential captures through it are refused (I1).
|
||||||
|
ZoneUSNexus
|
||||||
|
)
|
||||||
|
|
||||||
|
// String renders the zone for audit/refusal messages.
|
||||||
|
func (z Zone) String() string {
|
||||||
|
switch z {
|
||||||
|
case ZoneSovereign:
|
||||||
|
return "sovereign-soil"
|
||||||
|
case ZoneUSNexus:
|
||||||
|
return "us-nexus"
|
||||||
|
default:
|
||||||
|
return "unknown"
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// 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 and Origin are server-derived from
|
||||||
|
// the authenticated identity (the REST adapter populates them); they are
|
||||||
|
// 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
|
||||||
|
Origin Zone // server-derived trust zone; the I1 gate input
|
||||||
|
}
|
||||||
|
|
||||||
|
// 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"`
|
||||||
|
// AuditBuffered is true when the central audit sink was unreachable and
|
||||||
|
// this capture's audit record was written to the durable local buffer
|
||||||
|
// instead (internal/public tier). Surfaces the degraded state to the
|
||||||
|
// caller per §4.4.
|
||||||
|
AuditBuffered bool `json:"audit_buffered,omitempty"`
|
||||||
|
}
|
||||||
@@ -0,0 +1,126 @@
|
|||||||
|
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 closes an issue, optionally posting a closing comment
|
||||||
|
// first (empty comment ⇒ close only).
|
||||||
|
CloseIssue(ctx context.Context, repo string, number int, comment string) (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
|
||||||
|
}
|
||||||
|
|
||||||
|
// AuditOutcome is how a capture's audit record was (or will be) persisted.
|
||||||
|
type AuditOutcome int
|
||||||
|
|
||||||
|
const (
|
||||||
|
// AuditCentral means the record goes to the central sink (loki).
|
||||||
|
AuditCentral AuditOutcome = iota
|
||||||
|
// AuditBuffered means the central sink was unreachable and the record
|
||||||
|
// is written to a durable local buffer for later reconciliation
|
||||||
|
// (internal/public tier only).
|
||||||
|
AuditBuffered
|
||||||
|
)
|
||||||
|
|
||||||
|
// AuditSink is the two-phase, classification-aware audit port (I5, §4.4).
|
||||||
|
//
|
||||||
|
// Reserve runs BEFORE any write and decides whether the capture can be
|
||||||
|
// audited at its effective classification: it returns the outcome to use,
|
||||||
|
// or an error to refuse the capture before anything is written
|
||||||
|
// (confidential + central sink down → refuse; the all-tiers floor when
|
||||||
|
// nothing can record → refuse). Record runs AFTER the writes and persists
|
||||||
|
// the final entry per the reserved outcome.
|
||||||
|
//
|
||||||
|
// Splitting reserve from record is what lets "confidential + sink-down →
|
||||||
|
// refuse before any write" be literally true while the record itself
|
||||||
|
// (which lists what landed) is necessarily written afterwards.
|
||||||
|
type AuditSink interface {
|
||||||
|
Reserve(ctx context.Context, level classification.Level) (AuditOutcome, error)
|
||||||
|
Record(ctx context.Context, e AuditEntry, outcome AuditOutcome) error
|
||||||
|
}
|
||||||
@@ -0,0 +1,388 @@
|
|||||||
|
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}
|
||||||
|
|
||||||
|
// ErrSovereigntyRefused is returned when the I1 gate refuses a capture
|
||||||
|
// (confidential effective classification through a us-nexus origin). The
|
||||||
|
// REST adapter maps it to HTTP 403. Callers test with errors.Is.
|
||||||
|
var ErrSovereigntyRefused = fmt.Errorf("capture refused by I1 sovereignty gate")
|
||||||
|
|
||||||
|
// ErrAuditUnavailable is returned when the I5 audit gate refuses a capture
|
||||||
|
// before any write: a confidential capture whose central audit sink is
|
||||||
|
// unreachable, or the all-tiers floor where nothing can record the audit.
|
||||||
|
// The REST adapter maps it to HTTP 503. Callers test with errors.Is.
|
||||||
|
var ErrAuditUnavailable = fmt.Errorf("capture refused: audit substrate unavailable")
|
||||||
|
|
||||||
|
// assertedZoneMismatch returns a security-event string when the caller's
|
||||||
|
// harness label asserts a trust zone that contradicts the server-derived
|
||||||
|
// origin. A harness label that names no zone (the normal case, e.g.
|
||||||
|
// "claude-code") returns "". The label is never used as a gate input —
|
||||||
|
// this only flags the discrepancy for the audit trail.
|
||||||
|
func assertedZoneMismatch(harness string, derived Zone) string {
|
||||||
|
var asserted Zone
|
||||||
|
switch strings.ToLower(strings.TrimSpace(harness)) {
|
||||||
|
case "sovereign-soil", "sovereign":
|
||||||
|
asserted = ZoneSovereign
|
||||||
|
case "us-nexus", "usnexus":
|
||||||
|
asserted = ZoneUSNexus
|
||||||
|
default:
|
||||||
|
return "" // no zone claim
|
||||||
|
}
|
||||||
|
if asserted != derived {
|
||||||
|
return fmt.Sprintf("asserted-vs-derived origin mismatch: harness asserted %s, principal resolves to %s",
|
||||||
|
asserted, derived)
|
||||||
|
}
|
||||||
|
return ""
|
||||||
|
}
|
||||||
|
|
||||||
|
// 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)
|
||||||
|
|
||||||
|
// Server-derived origin governs the I1 gate; a caller-asserted harness
|
||||||
|
// label that names a different zone is descriptive-only and logged as a
|
||||||
|
// security event (spec §4.2: a control keyed on attacker-suppliable
|
||||||
|
// input is not a control).
|
||||||
|
if ev := assertedZoneMismatch(in.Context.Harness, in.Context.Origin); ev != "" {
|
||||||
|
securityEvents = append(securityEvents, ev)
|
||||||
|
}
|
||||||
|
|
||||||
|
// I1 sovereignty gate: a confidential capture through a us-nexus origin
|
||||||
|
// is refused before ANY write. The refusal itself is audited (best
|
||||||
|
// effort) — refusals must be reconstructable too.
|
||||||
|
if effective == classification.Confidential && in.Context.Origin == ZoneUSNexus {
|
||||||
|
_ = 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: nil, // refused before any write
|
||||||
|
SecurityEvents: append(securityEvents, "I1 refusal: confidential capture via us-nexus origin"),
|
||||||
|
}, AuditCentral)
|
||||||
|
return CaptureReceipt{}, fmt.Errorf("%w: effective classification confidential through %s origin",
|
||||||
|
ErrSovereigntyRefused, in.Context.Origin)
|
||||||
|
}
|
||||||
|
|
||||||
|
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
|
||||||
|
}
|
||||||
|
|
||||||
|
// I5 audit gate: decide BEFORE any write whether this capture can be
|
||||||
|
// audited at its effective classification. Confidential + central sink
|
||||||
|
// down → refuse here, before writing anything; the all-tiers floor
|
||||||
|
// (nothing can record) likewise refuses. Internal/public degrade to the
|
||||||
|
// durable local buffer (signalled by AuditBuffered).
|
||||||
|
outcome, err := s.audit.Reserve(ctx, effective)
|
||||||
|
if err != nil {
|
||||||
|
return CaptureReceipt{}, fmt.Errorf("%w: %v", ErrAuditUnavailable, err)
|
||||||
|
}
|
||||||
|
|
||||||
|
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: persist the request-level audit record of exactly what landed,
|
||||||
|
// using the outcome reserved before the writes. AuditBuffered surfaces
|
||||||
|
// the degraded (locally-buffered) state on the receipt.
|
||||||
|
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,
|
||||||
|
}, outcome); err != nil {
|
||||||
|
receipt.Errors = append(receipt.Errors, ItemError{Item: "audit", Error: err.Error()})
|
||||||
|
}
|
||||||
|
if outcome == AuditBuffered {
|
||||||
|
receipt.AuditBuffered = true
|
||||||
|
}
|
||||||
|
|
||||||
|
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, tk.Body)
|
||||||
|
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/<harness>/<YYYY-MM>/<date>-<slug>-<ref8>.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]
|
||||||
|
}
|
||||||
@@ -0,0 +1,471 @@
|
|||||||
|
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, _ string) (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 // Record error
|
||||||
|
reserveErr error // Reserve error (refuse before write)
|
||||||
|
reserveMode AuditOutcome
|
||||||
|
}
|
||||||
|
|
||||||
|
func (f *fakeAudit) Reserve(_ context.Context, _ classification.Level) (AuditOutcome, error) {
|
||||||
|
if f.reserveErr != nil {
|
||||||
|
return 0, f.reserveErr
|
||||||
|
}
|
||||||
|
return f.reserveMode, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (f *fakeAudit) Record(_ context.Context, e AuditEntry, _ AuditOutcome) 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")
|
||||||
|
}
|
||||||
|
|
||||||
|
// --- I1 sovereignty gate (#53) ---
|
||||||
|
|
||||||
|
func TestCaptureRefusesConfidentialViaUSNexus(t *testing.T) {
|
||||||
|
b := &fakeBrain{}
|
||||||
|
tr := &fakeTracker{}
|
||||||
|
au := &fakeAudit{}
|
||||||
|
pol := fakePolicy{tags: map[string]classification.Level{"client-seb": classification.Confidential}}
|
||||||
|
svc := newSvc(b, tr, nil, pol, au)
|
||||||
|
|
||||||
|
ctx := baseCtx()
|
||||||
|
ctx.Classification = "confidential"
|
||||||
|
ctx.Origin = ZoneUSNexus
|
||||||
|
_, err := svc.Capture(context.Background(), CaptureInput{
|
||||||
|
Context: ctx,
|
||||||
|
Insights: []Insight{{Text: "x", Wing: "client-seb", Hall: "facts"}},
|
||||||
|
})
|
||||||
|
require.Error(t, err)
|
||||||
|
assert.ErrorIs(t, err, ErrSovereigntyRefused)
|
||||||
|
// Refused before any write.
|
||||||
|
assert.Empty(t, b.writes)
|
||||||
|
assert.Empty(t, tr.created)
|
||||||
|
// Refusal is audited.
|
||||||
|
require.Len(t, au.entries, 1)
|
||||||
|
assert.Empty(t, au.entries[0].Items, "no items landed on refusal")
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestCaptureAllowsConfidentialViaSovereign(t *testing.T) {
|
||||||
|
b := &fakeBrain{}
|
||||||
|
pol := fakePolicy{tags: map[string]classification.Level{"client-seb": classification.Confidential}}
|
||||||
|
svc := newSvc(b, &fakeTracker{}, nil, pol, &fakeAudit{})
|
||||||
|
|
||||||
|
ctx := baseCtx()
|
||||||
|
ctx.Classification = "confidential"
|
||||||
|
ctx.Origin = ZoneSovereign
|
||||||
|
rec, err := svc.Capture(context.Background(), CaptureInput{
|
||||||
|
Context: ctx,
|
||||||
|
Insights: []Insight{{Text: "x", Wing: "client-seb", Hall: "facts"}},
|
||||||
|
})
|
||||||
|
require.NoError(t, err)
|
||||||
|
assert.True(t, rec.Insights[0].OK)
|
||||||
|
assert.Len(t, b.writes, 1)
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestCaptureAssertedLabelIgnoredAndLogged(t *testing.T) {
|
||||||
|
// Caller asserts harness "sovereign-soil" but principal resolves to
|
||||||
|
// us-nexus; confidential ⇒ refused, and the discrepancy is a security event.
|
||||||
|
au := &fakeAudit{}
|
||||||
|
pol := fakePolicy{tags: map[string]classification.Level{"client-seb": classification.Confidential}}
|
||||||
|
svc := newSvc(&fakeBrain{}, &fakeTracker{}, nil, pol, au)
|
||||||
|
|
||||||
|
ctx := baseCtx()
|
||||||
|
ctx.Harness = "sovereign-soil" // asserted
|
||||||
|
ctx.Origin = ZoneUSNexus // server-derived
|
||||||
|
ctx.Classification = "confidential"
|
||||||
|
_, err := svc.Capture(context.Background(), CaptureInput{
|
||||||
|
Context: ctx,
|
||||||
|
Insights: []Insight{{Text: "x", Wing: "client-seb", Hall: "facts"}},
|
||||||
|
})
|
||||||
|
require.ErrorIs(t, err, ErrSovereigntyRefused)
|
||||||
|
require.Len(t, au.entries, 1)
|
||||||
|
joined := strings.Join(au.entries[0].SecurityEvents, " | ")
|
||||||
|
assert.Contains(t, joined, "asserted-vs-derived origin mismatch")
|
||||||
|
assert.Contains(t, joined, "I1 refusal")
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestCaptureInternalViaUSNexusAllowed(t *testing.T) {
|
||||||
|
// us-nexus origin is fine for non-confidential data.
|
||||||
|
b := &fakeBrain{}
|
||||||
|
svc := newSvc(b, &fakeTracker{}, nil, fakePolicy{}, &fakeAudit{})
|
||||||
|
ctx := baseCtx()
|
||||||
|
ctx.Origin = ZoneUSNexus // internal classification, so gate doesn't fire
|
||||||
|
rec, err := svc.Capture(context.Background(), CaptureInput{
|
||||||
|
Context: ctx,
|
||||||
|
Insights: []Insight{{Text: "x", Wing: "hyperguild", Hall: "facts"}},
|
||||||
|
})
|
||||||
|
require.NoError(t, err)
|
||||||
|
assert.True(t, rec.Insights[0].OK)
|
||||||
|
}
|
||||||
|
|
||||||
|
// --- I5 audit gate (#54) ---
|
||||||
|
|
||||||
|
func TestCaptureRefusesWhenAuditReserveFails(t *testing.T) {
|
||||||
|
// Reserve refusing (e.g. confidential + central sink down, or the floor)
|
||||||
|
// aborts the capture before any write.
|
||||||
|
b := &fakeBrain{}
|
||||||
|
tr := &fakeTracker{}
|
||||||
|
au := &fakeAudit{reserveErr: errors.New("central sink unreachable")}
|
||||||
|
svc := newSvc(b, tr, nil, fakePolicy{}, au)
|
||||||
|
|
||||||
|
_, err := svc.Capture(context.Background(), CaptureInput{
|
||||||
|
Context: baseCtx(),
|
||||||
|
Insights: []Insight{{Text: "x", Wing: "hyperguild", Hall: "facts"}},
|
||||||
|
})
|
||||||
|
require.Error(t, err)
|
||||||
|
assert.ErrorIs(t, err, ErrAuditUnavailable)
|
||||||
|
assert.Empty(t, b.writes, "nothing written when audit unavailable")
|
||||||
|
assert.Empty(t, tr.created)
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestCaptureFlagsLocallyBufferedAudit(t *testing.T) {
|
||||||
|
// Reserve returns AuditBuffered (internal/public, central down) → capture
|
||||||
|
// proceeds and the receipt flags the degraded audit state.
|
||||||
|
b := &fakeBrain{}
|
||||||
|
au := &fakeAudit{reserveMode: AuditBuffered}
|
||||||
|
svc := newSvc(b, &fakeTracker{}, nil, fakePolicy{}, au)
|
||||||
|
|
||||||
|
rec, err := svc.Capture(context.Background(), CaptureInput{
|
||||||
|
Context: baseCtx(),
|
||||||
|
Insights: []Insight{{Text: "x", Wing: "hyperguild", Hall: "facts"}},
|
||||||
|
})
|
||||||
|
require.NoError(t, err)
|
||||||
|
assert.True(t, rec.Insights[0].OK, "capture proceeds on degraded audit")
|
||||||
|
assert.True(t, rec.AuditBuffered, "receipt flags locally-buffered audit")
|
||||||
|
require.Len(t, au.entries, 1)
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestCaptureDryRunSkipsAuditGate(t *testing.T) {
|
||||||
|
// dry_run must not even probe the audit sink (writes nothing anywhere).
|
||||||
|
au := &fakeAudit{reserveErr: errors.New("would refuse")}
|
||||||
|
svc := newSvc(&fakeBrain{}, &fakeTracker{}, nil, fakePolicy{}, au)
|
||||||
|
|
||||||
|
rec, err := svc.Capture(context.Background(), CaptureInput{
|
||||||
|
Context: baseCtx(),
|
||||||
|
DryRun: true,
|
||||||
|
Insights: []Insight{{Text: "x", Wing: "hyperguild", Hall: "facts"}},
|
||||||
|
})
|
||||||
|
require.NoError(t, err, "dry-run does not hit the audit gate")
|
||||||
|
assert.True(t, rec.DryRun)
|
||||||
|
assert.Empty(t, au.entries)
|
||||||
|
}
|
||||||
@@ -0,0 +1,232 @@
|
|||||||
|
// Package capturehttp is the REST adapter for the capture use-case: the
|
||||||
|
// POST /capture door (#53). It is deliberately thin — authenticate, derive
|
||||||
|
// the trust-zone origin from the authenticated principal, decode the
|
||||||
|
// request, call capture.Service, map the receipt to an HTTP status. No
|
||||||
|
// business logic lives here; the I1 gate, validation, and orchestration
|
||||||
|
// are all in the use-case.
|
||||||
|
package capturehttp
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"crypto/subtle"
|
||||||
|
"encoding/json"
|
||||||
|
"errors"
|
||||||
|
"io"
|
||||||
|
"net/http"
|
||||||
|
"strings"
|
||||||
|
|
||||||
|
"github.com/mathiasbq/hyperguild/ingestion/internal/capture"
|
||||||
|
)
|
||||||
|
|
||||||
|
// Validator validates a Bearer JWT and returns its subject. The chassis
|
||||||
|
// *auth.JWTValidator satisfies it (including its nil-receiver "disabled"
|
||||||
|
// behaviour), and tests can substitute a fake without a live JWKS.
|
||||||
|
type Validator interface {
|
||||||
|
Validate(ctx context.Context, rawToken string) (string, error)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Handler serves POST /capture.
|
||||||
|
type Handler struct {
|
||||||
|
svc *capture.Service
|
||||||
|
validator Validator // nil ⇒ JWT auth disabled
|
||||||
|
staticToken string // "" ⇒ static auth disabled
|
||||||
|
staticPrincipal string // principal name attributed to static-token callers
|
||||||
|
resolver OriginResolver
|
||||||
|
}
|
||||||
|
|
||||||
|
// New constructs a capture HTTP handler. staticToken callers are
|
||||||
|
// attributed to staticPrincipal (a sovereign homelab identity); JWT
|
||||||
|
// callers are attributed to their token subject.
|
||||||
|
func New(svc *capture.Service, validator Validator, staticToken, staticPrincipal string, resolver OriginResolver) *Handler {
|
||||||
|
if staticPrincipal == "" {
|
||||||
|
staticPrincipal = "local-cli"
|
||||||
|
}
|
||||||
|
return &Handler{
|
||||||
|
svc: svc,
|
||||||
|
validator: validator,
|
||||||
|
staticToken: staticToken,
|
||||||
|
staticPrincipal: staticPrincipal,
|
||||||
|
resolver: resolver,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// wire types — the POST /capture request body.
|
||||||
|
type request struct {
|
||||||
|
Context contextBody `json:"context"`
|
||||||
|
Insights []insightBody `json:"insights"`
|
||||||
|
Tickets []ticketBody `json:"tickets"`
|
||||||
|
Summary *summaryBody `json:"summary,omitempty"`
|
||||||
|
DryRun bool `json:"dry_run"`
|
||||||
|
}
|
||||||
|
|
||||||
|
type contextBody struct {
|
||||||
|
Harness string `json:"harness"`
|
||||||
|
SessionRef string `json:"session_ref"`
|
||||||
|
Fidelity string `json:"fidelity"`
|
||||||
|
Actor string `json:"actor"`
|
||||||
|
Classification string `json:"classification"`
|
||||||
|
}
|
||||||
|
|
||||||
|
type insightBody struct {
|
||||||
|
Text string `json:"text"`
|
||||||
|
Wing string `json:"wing"`
|
||||||
|
Hall string `json:"hall"`
|
||||||
|
SupersedeSlug string `json:"supersede_slug,omitempty"`
|
||||||
|
}
|
||||||
|
|
||||||
|
type ticketBody struct {
|
||||||
|
Repo string `json:"repo"`
|
||||||
|
Action string `json:"action"`
|
||||||
|
Number int `json:"number,omitempty"`
|
||||||
|
Title string `json:"title,omitempty"`
|
||||||
|
Body string `json:"body,omitempty"`
|
||||||
|
}
|
||||||
|
|
||||||
|
type summaryBody struct {
|
||||||
|
Title string `json:"title"`
|
||||||
|
Body string `json:"body"`
|
||||||
|
ReposTouched []string `json:"repos_touched,omitempty"`
|
||||||
|
}
|
||||||
|
|
||||||
|
// ServeHTTP authenticates, derives origin, runs the use-case, and maps the
|
||||||
|
// result to an HTTP status.
|
||||||
|
func (h *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
||||||
|
principal, viaStatic, ok := Authenticate(r, h.staticToken, h.staticPrincipal, h.validator)
|
||||||
|
if !ok {
|
||||||
|
http.Error(w, "unauthorized", http.StatusUnauthorized)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
body, err := io.ReadAll(r.Body)
|
||||||
|
if err != nil {
|
||||||
|
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "read body"})
|
||||||
|
return
|
||||||
|
}
|
||||||
|
in, err := DecodeRequest(body)
|
||||||
|
if err != nil {
|
||||||
|
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "invalid JSON"})
|
||||||
|
return
|
||||||
|
}
|
||||||
|
// Principal and origin are server-derived — overwrite anything the
|
||||||
|
// caller may have tried to put in the body.
|
||||||
|
in.Context.Principal = principal
|
||||||
|
in.Context.Origin = h.resolver.Resolve(principal, viaStatic)
|
||||||
|
|
||||||
|
rec, err := h.svc.Capture(r.Context(), in)
|
||||||
|
switch {
|
||||||
|
case errors.Is(err, capture.ErrSovereigntyRefused):
|
||||||
|
writeJSON(w, http.StatusForbidden, map[string]string{"error": err.Error()})
|
||||||
|
return
|
||||||
|
case errors.Is(err, capture.ErrAuditUnavailable):
|
||||||
|
// I5 refusal: confidential + audit sink down, or the all-tiers floor.
|
||||||
|
writeJSON(w, http.StatusServiceUnavailable, map[string]string{"error": err.Error()})
|
||||||
|
return
|
||||||
|
case err != nil:
|
||||||
|
// Pre-write validation failure (fail-closed).
|
||||||
|
writeJSON(w, http.StatusBadRequest, map[string]string{"error": err.Error()})
|
||||||
|
return
|
||||||
|
}
|
||||||
|
writeJSON(w, statusFor(rec), rec)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Authenticate mirrors the chassis Bearer precedence (static token wins,
|
||||||
|
// then Dex JWT) and returns the resolved principal plus whether the static
|
||||||
|
// path was taken — the chassis middleware hides both, and capture (REST or
|
||||||
|
// MCP) needs them to derive the trust-zone origin. ok is false when no
|
||||||
|
// credential matched.
|
||||||
|
func Authenticate(r *http.Request, staticToken, staticPrincipal string, validator Validator) (principal string, viaStatic, ok bool) {
|
||||||
|
raw, found := strings.CutPrefix(r.Header.Get("Authorization"), "Bearer ")
|
||||||
|
if !found || raw == "" {
|
||||||
|
return "", false, false
|
||||||
|
}
|
||||||
|
if staticToken != "" && subtle.ConstantTimeCompare([]byte(raw), []byte(staticToken)) == 1 {
|
||||||
|
return staticPrincipal, true, true
|
||||||
|
}
|
||||||
|
if validator != nil {
|
||||||
|
if sub, err := validator.Validate(r.Context(), raw); err == nil && sub != "" {
|
||||||
|
return sub, false, true
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return "", false, false
|
||||||
|
}
|
||||||
|
|
||||||
|
// DecodeRequest parses a capture request body into a CaptureInput. Shared
|
||||||
|
// by the REST adapter and the MCP capture tool so the wire shape has one
|
||||||
|
// definition. Principal and Origin are NOT set here — the caller sets them
|
||||||
|
// from the authenticated identity.
|
||||||
|
func DecodeRequest(data []byte) (capture.CaptureInput, error) {
|
||||||
|
var b request
|
||||||
|
if err := json.Unmarshal(data, &b); err != nil {
|
||||||
|
return capture.CaptureInput{}, err
|
||||||
|
}
|
||||||
|
return b.toInput(), nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (b request) toInput() capture.CaptureInput {
|
||||||
|
in := capture.CaptureInput{
|
||||||
|
Context: capture.CaptureContext{
|
||||||
|
Harness: b.Context.Harness,
|
||||||
|
SessionRef: b.Context.SessionRef,
|
||||||
|
Fidelity: b.Context.Fidelity,
|
||||||
|
Actor: b.Context.Actor,
|
||||||
|
Classification: b.Context.Classification,
|
||||||
|
},
|
||||||
|
DryRun: b.DryRun,
|
||||||
|
}
|
||||||
|
for _, i := range b.Insights {
|
||||||
|
in.Insights = append(in.Insights, capture.Insight{
|
||||||
|
Text: i.Text, Wing: i.Wing, Hall: i.Hall, SupersedeSlug: i.SupersedeSlug,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
for _, t := range b.Tickets {
|
||||||
|
in.Tickets = append(in.Tickets, capture.Ticket{
|
||||||
|
Repo: t.Repo, Action: t.Action, Number: t.Number, Title: t.Title, Body: t.Body,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
if b.Summary != nil {
|
||||||
|
in.Summary = &capture.Summary{
|
||||||
|
Title: b.Summary.Title, Body: b.Summary.Body, ReposTouched: b.Summary.ReposTouched,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return in
|
||||||
|
}
|
||||||
|
|
||||||
|
// statusFor maps a receipt to an HTTP status: 200 all-ok (or dry-run),
|
||||||
|
// 207 partial, 502 everything-failed.
|
||||||
|
func statusFor(rec capture.CaptureReceipt) int {
|
||||||
|
if rec.DryRun {
|
||||||
|
return http.StatusOK
|
||||||
|
}
|
||||||
|
var ok, fail int
|
||||||
|
for _, i := range rec.Insights {
|
||||||
|
count(&ok, &fail, i.OK)
|
||||||
|
}
|
||||||
|
for _, t := range rec.Tickets {
|
||||||
|
count(&ok, &fail, t.OK)
|
||||||
|
}
|
||||||
|
if rec.Summary != nil {
|
||||||
|
count(&ok, &fail, rec.Summary.OK)
|
||||||
|
}
|
||||||
|
switch {
|
||||||
|
case fail == 0:
|
||||||
|
return http.StatusOK
|
||||||
|
case ok == 0:
|
||||||
|
return http.StatusBadGateway // every persistence attempt failed
|
||||||
|
default:
|
||||||
|
return http.StatusMultiStatus // 207: partial success
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func count(ok, fail *int, isOK bool) {
|
||||||
|
if isOK {
|
||||||
|
*ok++
|
||||||
|
} else {
|
||||||
|
*fail++
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func writeJSON(w http.ResponseWriter, status int, v any) {
|
||||||
|
w.Header().Set("Content-Type", "application/json")
|
||||||
|
w.WriteHeader(status)
|
||||||
|
_ = json.NewEncoder(w).Encode(v)
|
||||||
|
}
|
||||||
@@ -0,0 +1,198 @@
|
|||||||
|
package capturehttp_test
|
||||||
|
|
||||||
|
import (
|
||||||
|
"bytes"
|
||||||
|
"context"
|
||||||
|
"encoding/json"
|
||||||
|
"errors"
|
||||||
|
"net/http"
|
||||||
|
"net/http/httptest"
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"github.com/mathiasbq/hyperguild/ingestion/internal/audit"
|
||||||
|
"github.com/mathiasbq/hyperguild/ingestion/internal/brainstore"
|
||||||
|
"github.com/mathiasbq/hyperguild/ingestion/internal/capture"
|
||||||
|
"github.com/mathiasbq/hyperguild/ingestion/internal/capturehttp"
|
||||||
|
"github.com/mathiasbq/hyperguild/ingestion/internal/classification"
|
||||||
|
"github.com/stretchr/testify/assert"
|
||||||
|
"github.com/stretchr/testify/require"
|
||||||
|
)
|
||||||
|
|
||||||
|
const staticTok = "static-secret"
|
||||||
|
|
||||||
|
// fakeValidator stands in for the chassis JWT validator.
|
||||||
|
type fakeValidator struct {
|
||||||
|
subject string
|
||||||
|
err error
|
||||||
|
}
|
||||||
|
|
||||||
|
func (f fakeValidator) Validate(context.Context, string) (string, error) {
|
||||||
|
return f.subject, f.err
|
||||||
|
}
|
||||||
|
|
||||||
|
type fakeTracker struct{ failCreate bool }
|
||||||
|
|
||||||
|
func (f fakeTracker) CreateIssue(context.Context, string, string, string) (capture.IssueRef, error) {
|
||||||
|
if f.failCreate {
|
||||||
|
return capture.IssueRef{}, errors.New("gitea down")
|
||||||
|
}
|
||||||
|
return capture.IssueRef{Repo: "hyperguild", Number: 1, URL: "https://git/1"}, nil
|
||||||
|
}
|
||||||
|
func (fakeTracker) CloseIssue(context.Context, string, int, string) (capture.IssueRef, error) {
|
||||||
|
return capture.IssueRef{}, nil
|
||||||
|
}
|
||||||
|
func (fakeTracker) CommentIssue(context.Context, string, int, string) (capture.IssueRef, error) {
|
||||||
|
return capture.IssueRef{}, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func newHandler(t *testing.T, v capturehttp.Validator, tr capture.IssueTracker, sovereign []string) *capturehttp.Handler {
|
||||||
|
t.Helper()
|
||||||
|
cfg, err := classification.Load(t.TempDir())
|
||||||
|
require.NoError(t, err)
|
||||||
|
svc := capture.NewService(brainstore.New(t.TempDir()), tr, nil, cfg, audit.NewSlogSink(nil))
|
||||||
|
return capturehttp.New(svc, v, staticTok, "local-cli", capturehttp.NewOriginResolver(sovereign))
|
||||||
|
}
|
||||||
|
|
||||||
|
func do(t *testing.T, h *capturehttp.Handler, authz string, body any) *httptest.ResponseRecorder {
|
||||||
|
t.Helper()
|
||||||
|
b, _ := json.Marshal(body)
|
||||||
|
req := httptest.NewRequest(http.MethodPost, "/capture", bytes.NewReader(b))
|
||||||
|
if authz != "" {
|
||||||
|
req.Header.Set("Authorization", authz)
|
||||||
|
}
|
||||||
|
rr := httptest.NewRecorder()
|
||||||
|
h.ServeHTTP(rr, req)
|
||||||
|
return rr
|
||||||
|
}
|
||||||
|
|
||||||
|
func internalReq() map[string]any {
|
||||||
|
return map[string]any{
|
||||||
|
"context": map[string]any{"harness": "claude-code", "actor": "mathias", "classification": "internal"},
|
||||||
|
"insights": []map[string]any{{"text": "a fact", "wing": "hyperguild", "hall": "facts"}},
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestUnauthorizedWithoutToken(t *testing.T) {
|
||||||
|
h := newHandler(t, fakeValidator{err: errors.New("no")}, fakeTracker{}, nil)
|
||||||
|
rr := do(t, h, "", internalReq())
|
||||||
|
assert.Equal(t, http.StatusUnauthorized, rr.Code)
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestUnauthorizedBadToken(t *testing.T) {
|
||||||
|
h := newHandler(t, fakeValidator{err: errors.New("bad jwt")}, fakeTracker{}, nil)
|
||||||
|
rr := do(t, h, "Bearer wrong", internalReq())
|
||||||
|
assert.Equal(t, http.StatusUnauthorized, rr.Code)
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestHappyPathStaticToken(t *testing.T) {
|
||||||
|
h := newHandler(t, nil, fakeTracker{}, nil)
|
||||||
|
rr := do(t, h, "Bearer "+staticTok, map[string]any{
|
||||||
|
"context": map[string]any{"harness": "claude-code", "actor": "mathias", "classification": "internal"},
|
||||||
|
"insights": []map[string]any{{"text": "a fact", "wing": "hyperguild", "hall": "facts"}},
|
||||||
|
"tickets": []map[string]any{{"repo": "hyperguild", "action": "create", "title": "t"}},
|
||||||
|
})
|
||||||
|
require.Equal(t, http.StatusOK, rr.Code)
|
||||||
|
var rec capture.CaptureReceipt
|
||||||
|
require.NoError(t, json.Unmarshal(rr.Body.Bytes(), &rec))
|
||||||
|
assert.True(t, rec.Insights[0].OK)
|
||||||
|
assert.True(t, rec.Tickets[0].OK)
|
||||||
|
assert.Empty(t, rec.Errors)
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestConfidentialViaUSNexusRefused(t *testing.T) {
|
||||||
|
// JWT principal not in the sovereign allowlist ⇒ us-nexus; confidential ⇒ 403.
|
||||||
|
h := newHandler(t, fakeValidator{subject: "claudeai-oauth-client"}, fakeTracker{}, nil)
|
||||||
|
rr := do(t, h, "Bearer jwt-token", map[string]any{
|
||||||
|
"context": map[string]any{"harness": "claudeai-chat", "actor": "mathias", "classification": "confidential"},
|
||||||
|
"insights": []map[string]any{{"text": "secret", "wing": "client-seb", "hall": "facts"}},
|
||||||
|
})
|
||||||
|
assert.Equal(t, http.StatusForbidden, rr.Code)
|
||||||
|
assert.Contains(t, rr.Body.String(), "sovereignty")
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestConfidentialViaSovereignJWTAllowed(t *testing.T) {
|
||||||
|
// Same confidential payload, but the principal is allowlisted sovereign ⇒ allowed.
|
||||||
|
h := newHandler(t, fakeValidator{subject: "koala-cli"}, fakeTracker{}, []string{"koala-cli"})
|
||||||
|
rr := do(t, h, "Bearer jwt-token", map[string]any{
|
||||||
|
"context": map[string]any{"harness": "claude-code", "actor": "mathias", "classification": "confidential"},
|
||||||
|
"insights": []map[string]any{{"text": "secret", "wing": "client-seb", "hall": "facts"}},
|
||||||
|
})
|
||||||
|
require.Equal(t, http.StatusOK, rr.Code)
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestStaticTokenIsSovereignSoConfidentialAllowed(t *testing.T) {
|
||||||
|
h := newHandler(t, nil, fakeTracker{}, nil)
|
||||||
|
rr := do(t, h, "Bearer "+staticTok, map[string]any{
|
||||||
|
"context": map[string]any{"harness": "claude-code", "actor": "mathias", "classification": "confidential"},
|
||||||
|
"insights": []map[string]any{{"text": "secret", "wing": "client-seb", "hall": "facts"}},
|
||||||
|
})
|
||||||
|
assert.Equal(t, http.StatusOK, rr.Code)
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestValidationRejectedBeforeWrite(t *testing.T) {
|
||||||
|
h := newHandler(t, nil, fakeTracker{}, nil)
|
||||||
|
rr := do(t, h, "Bearer "+staticTok, map[string]any{
|
||||||
|
"context": map[string]any{"actor": "mathias", "classification": "internal"},
|
||||||
|
"insights": []map[string]any{{"text": "x", "wing": "hyperguild", "hall": "not-a-hall"}},
|
||||||
|
})
|
||||||
|
assert.Equal(t, http.StatusBadRequest, rr.Code)
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestPartialFailureIs207(t *testing.T) {
|
||||||
|
h := newHandler(t, nil, fakeTracker{failCreate: true}, nil)
|
||||||
|
rr := do(t, h, "Bearer "+staticTok, map[string]any{
|
||||||
|
"context": map[string]any{"harness": "claude-code", "actor": "mathias", "classification": "internal"},
|
||||||
|
"insights": []map[string]any{{"text": "ok insight", "wing": "hyperguild", "hall": "facts"}},
|
||||||
|
"tickets": []map[string]any{{"repo": "hyperguild", "action": "create", "title": "fails"}},
|
||||||
|
})
|
||||||
|
assert.Equal(t, http.StatusMultiStatus, rr.Code)
|
||||||
|
var rec capture.CaptureReceipt
|
||||||
|
require.NoError(t, json.Unmarshal(rr.Body.Bytes(), &rec))
|
||||||
|
assert.True(t, rec.Insights[0].OK)
|
||||||
|
assert.False(t, rec.Tickets[0].OK)
|
||||||
|
assert.Len(t, rec.Errors, 1)
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestDryRunWritesNothing(t *testing.T) {
|
||||||
|
h := newHandler(t, nil, fakeTracker{}, nil)
|
||||||
|
rr := do(t, h, "Bearer "+staticTok, map[string]any{
|
||||||
|
"context": map[string]any{"harness": "claude-code", "actor": "mathias", "classification": "internal"},
|
||||||
|
"insights": []map[string]any{{"text": "a", "wing": "hyperguild", "hall": "facts"}},
|
||||||
|
"dry_run": true,
|
||||||
|
})
|
||||||
|
require.Equal(t, http.StatusOK, rr.Code)
|
||||||
|
var rec capture.CaptureReceipt
|
||||||
|
require.NoError(t, json.Unmarshal(rr.Body.Bytes(), &rec))
|
||||||
|
assert.True(t, rec.DryRun)
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestCallerCannotForgeOrigin(t *testing.T) {
|
||||||
|
// Even if the body tried to assert a sovereign harness, a us-nexus JWT
|
||||||
|
// principal + confidential ⇒ refused. (Origin is server-derived.)
|
||||||
|
h := newHandler(t, fakeValidator{subject: "claudeai-oauth-client"}, fakeTracker{}, nil)
|
||||||
|
rr := do(t, h, "Bearer jwt", map[string]any{
|
||||||
|
"context": map[string]any{"harness": "sovereign-soil", "actor": "mathias", "classification": "confidential"},
|
||||||
|
"insights": []map[string]any{{"text": "secret", "wing": "client-seb", "hall": "facts"}},
|
||||||
|
})
|
||||||
|
assert.Equal(t, http.StatusForbidden, rr.Code)
|
||||||
|
}
|
||||||
|
|
||||||
|
// refusingAudit refuses at Reserve (e.g. confidential + loki down, or floor).
|
||||||
|
type refusingAudit struct{}
|
||||||
|
|
||||||
|
func (refusingAudit) Reserve(context.Context, classification.Level) (capture.AuditOutcome, error) {
|
||||||
|
return 0, errors.New("central audit sink unreachable")
|
||||||
|
}
|
||||||
|
func (refusingAudit) Record(context.Context, capture.AuditEntry, capture.AuditOutcome) error {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestAuditUnavailableIs503(t *testing.T) {
|
||||||
|
cfg, err := classification.Load(t.TempDir())
|
||||||
|
require.NoError(t, err)
|
||||||
|
svc := capture.NewService(brainstore.New(t.TempDir()), fakeTracker{}, nil, cfg, refusingAudit{})
|
||||||
|
h := capturehttp.New(svc, nil, staticTok, "local-cli", capturehttp.NewOriginResolver(nil))
|
||||||
|
|
||||||
|
rr := do(t, h, "Bearer "+staticTok, internalReq())
|
||||||
|
assert.Equal(t, http.StatusServiceUnavailable, rr.Code)
|
||||||
|
}
|
||||||
@@ -0,0 +1,42 @@
|
|||||||
|
package capturehttp
|
||||||
|
|
||||||
|
import "github.com/mathiasbq/hyperguild/ingestion/internal/capture"
|
||||||
|
|
||||||
|
// OriginResolver maps an authenticated principal to its trust zone
|
||||||
|
// (spec §4.2). The mapping is server-side and never reads caller input.
|
||||||
|
//
|
||||||
|
// Rules:
|
||||||
|
// - The static-token path is a homelab CLI caller on sovereign soil →
|
||||||
|
// ZoneSovereign.
|
||||||
|
// - A JWT principal in the sovereign allowlist → ZoneSovereign.
|
||||||
|
// - Any other JWT principal (e.g. claude.ai's OAuth identity, or any
|
||||||
|
// unrecognised subject) → ZoneUSNexus.
|
||||||
|
//
|
||||||
|
// The default is the strict one: an unknown principal is treated as
|
||||||
|
// us-nexus so the I1 gate fails safe (refuses confidential), exactly as
|
||||||
|
// an untagged classification target fails safe to confidential (#50).
|
||||||
|
type OriginResolver struct {
|
||||||
|
sovereign map[string]bool
|
||||||
|
}
|
||||||
|
|
||||||
|
// NewOriginResolver builds a resolver whose JWT sovereign principals are
|
||||||
|
// the given subjects. The static-token caller is always sovereign and
|
||||||
|
// need not be listed.
|
||||||
|
func NewOriginResolver(sovereignPrincipals []string) OriginResolver {
|
||||||
|
m := make(map[string]bool, len(sovereignPrincipals))
|
||||||
|
for _, p := range sovereignPrincipals {
|
||||||
|
if p != "" {
|
||||||
|
m[p] = true
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return OriginResolver{sovereign: m}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Resolve returns the trust zone for a principal. viaStatic is true when
|
||||||
|
// the static-token auth path was taken.
|
||||||
|
func (r OriginResolver) Resolve(principal string, viaStatic bool) capture.Zone {
|
||||||
|
if viaStatic || r.sovereign[principal] {
|
||||||
|
return capture.ZoneSovereign
|
||||||
|
}
|
||||||
|
return capture.ZoneUSNexus
|
||||||
|
}
|
||||||
@@ -0,0 +1,129 @@
|
|||||||
|
// Package gitea implements capture.IssueTracker against a Gitea instance
|
||||||
|
// over its REST API. It is the new outbound dependency the brain server
|
||||||
|
// gains for the capture capability (#49c/#52): the server otherwise does
|
||||||
|
// brain-local file ops only.
|
||||||
|
//
|
||||||
|
// Owner is hard-coded to the operator and never taken from caller input.
|
||||||
|
// The API token is read once at construction, held in the struct, and
|
||||||
|
// never logged or placed in argv — it travels only in the Authorization
|
||||||
|
// header of outbound requests (AGENTS.md secret-handling).
|
||||||
|
package gitea
|
||||||
|
|
||||||
|
import (
|
||||||
|
"bytes"
|
||||||
|
"context"
|
||||||
|
"encoding/json"
|
||||||
|
"fmt"
|
||||||
|
"io"
|
||||||
|
"net/http"
|
||||||
|
"strings"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/mathiasbq/hyperguild/ingestion/internal/capture"
|
||||||
|
)
|
||||||
|
|
||||||
|
// owner is the fixed repository owner for every ticket operation. It is a
|
||||||
|
// constant, not a parameter, so a caller can never redirect a write to
|
||||||
|
// another owner's repo.
|
||||||
|
const owner = "mathias"
|
||||||
|
|
||||||
|
// Client is a Gitea REST API IssueTracker.
|
||||||
|
type Client struct {
|
||||||
|
baseURL string
|
||||||
|
token string
|
||||||
|
http *http.Client
|
||||||
|
}
|
||||||
|
|
||||||
|
// New constructs a Client. It returns nil when either baseURL or token is
|
||||||
|
// empty, so callers can treat missing config as "tracker disabled" with a
|
||||||
|
// single nil check (mirrors embed.New).
|
||||||
|
func New(baseURL, token string) *Client {
|
||||||
|
if baseURL == "" || token == "" {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
return &Client{
|
||||||
|
baseURL: strings.TrimRight(baseURL, "/"),
|
||||||
|
token: token,
|
||||||
|
http: &http.Client{Timeout: 15 * time.Second},
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// issueResponse is the subset of a Gitea issue/comment payload we read.
|
||||||
|
type issueResponse struct {
|
||||||
|
Number int `json:"number"`
|
||||||
|
HTMLURL string `json:"html_url"`
|
||||||
|
}
|
||||||
|
|
||||||
|
// CreateIssue opens a new issue under the fixed owner.
|
||||||
|
func (c *Client) CreateIssue(ctx context.Context, repo, title, body string) (capture.IssueRef, error) {
|
||||||
|
var out issueResponse
|
||||||
|
if err := c.do(ctx, http.MethodPost,
|
||||||
|
fmt.Sprintf("/api/v1/repos/%s/%s/issues", owner, repo),
|
||||||
|
map[string]any{"title": title, "body": body}, &out); err != nil {
|
||||||
|
return capture.IssueRef{}, err
|
||||||
|
}
|
||||||
|
return capture.IssueRef{Repo: repo, Number: out.Number, URL: out.HTMLURL}, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// CommentIssue posts a comment on an existing issue.
|
||||||
|
func (c *Client) CommentIssue(ctx context.Context, repo string, number int, body string) (capture.IssueRef, error) {
|
||||||
|
var out issueResponse
|
||||||
|
if err := c.do(ctx, http.MethodPost,
|
||||||
|
fmt.Sprintf("/api/v1/repos/%s/%s/issues/%d/comments", owner, repo, number),
|
||||||
|
map[string]any{"body": body}, &out); err != nil {
|
||||||
|
return capture.IssueRef{}, err
|
||||||
|
}
|
||||||
|
return capture.IssueRef{Repo: repo, Number: number, URL: out.HTMLURL}, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// CloseIssue closes an issue, first posting a closing comment when one is
|
||||||
|
// given (empty comment ⇒ close only).
|
||||||
|
func (c *Client) CloseIssue(ctx context.Context, repo string, number int, comment string) (capture.IssueRef, error) {
|
||||||
|
if strings.TrimSpace(comment) != "" {
|
||||||
|
if _, err := c.CommentIssue(ctx, repo, number, comment); err != nil {
|
||||||
|
return capture.IssueRef{}, err
|
||||||
|
}
|
||||||
|
}
|
||||||
|
var out issueResponse
|
||||||
|
if err := c.do(ctx, http.MethodPatch,
|
||||||
|
fmt.Sprintf("/api/v1/repos/%s/%s/issues/%d", owner, repo, number),
|
||||||
|
map[string]any{"state": "closed"}, &out); err != nil {
|
||||||
|
return capture.IssueRef{}, err
|
||||||
|
}
|
||||||
|
return capture.IssueRef{Repo: repo, Number: number, URL: out.HTMLURL}, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// do performs a JSON request against the Gitea API and decodes the
|
||||||
|
// response into out. Errors carry the status and a truncated body for
|
||||||
|
// diagnosis but never the token.
|
||||||
|
func (c *Client) do(ctx context.Context, method, path string, payload any, out *issueResponse) error {
|
||||||
|
reqBody, err := json.Marshal(payload)
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("marshal request: %w", err)
|
||||||
|
}
|
||||||
|
req, err := http.NewRequestWithContext(ctx, method, c.baseURL+path, bytes.NewReader(reqBody))
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
req.Header.Set("Content-Type", "application/json")
|
||||||
|
req.Header.Set("Accept", "application/json")
|
||||||
|
// Gitea's token scheme. Held here only; never logged.
|
||||||
|
req.Header.Set("Authorization", "token "+c.token)
|
||||||
|
|
||||||
|
resp, err := c.http.Do(req)
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("gitea %s %s: %w", method, path, err)
|
||||||
|
}
|
||||||
|
defer func() { _ = resp.Body.Close() }()
|
||||||
|
|
||||||
|
respBody, _ := io.ReadAll(io.LimitReader(resp.Body, 4096))
|
||||||
|
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
|
||||||
|
return fmt.Errorf("gitea %s %s: status %d: %s", method, path, resp.StatusCode, strings.TrimSpace(string(respBody)))
|
||||||
|
}
|
||||||
|
if out != nil && len(respBody) > 0 {
|
||||||
|
if err := json.Unmarshal(respBody, out); err != nil {
|
||||||
|
return fmt.Errorf("gitea %s %s: decode response: %w", method, path, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
@@ -0,0 +1,117 @@
|
|||||||
|
package gitea_test
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"encoding/json"
|
||||||
|
"io"
|
||||||
|
"net/http"
|
||||||
|
"net/http/httptest"
|
||||||
|
"strings"
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"github.com/mathiasbq/hyperguild/ingestion/internal/gitea"
|
||||||
|
"github.com/stretchr/testify/assert"
|
||||||
|
"github.com/stretchr/testify/require"
|
||||||
|
)
|
||||||
|
|
||||||
|
const testToken = "super-secret-token-value"
|
||||||
|
|
||||||
|
func TestNewNilWhenUnconfigured(t *testing.T) {
|
||||||
|
assert.Nil(t, gitea.New("", testToken))
|
||||||
|
assert.Nil(t, gitea.New("https://git.example", ""))
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestCreateIssueForcesOwnerAndAuth(t *testing.T) {
|
||||||
|
var gotPath, gotAuth, gotBody string
|
||||||
|
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||||
|
gotPath = r.URL.Path
|
||||||
|
gotAuth = r.Header.Get("Authorization")
|
||||||
|
b, _ := io.ReadAll(r.Body)
|
||||||
|
gotBody = string(b)
|
||||||
|
assert.Equal(t, http.MethodPost, r.Method)
|
||||||
|
w.WriteHeader(http.StatusCreated)
|
||||||
|
_ = json.NewEncoder(w).Encode(map[string]any{"number": 42, "html_url": "https://git.d-ma.be/mathias/hyperguild/issues/42"})
|
||||||
|
}))
|
||||||
|
defer srv.Close()
|
||||||
|
|
||||||
|
c := gitea.New(srv.URL, testToken)
|
||||||
|
require.NotNil(t, c)
|
||||||
|
ref, err := c.CreateIssue(context.Background(), "hyperguild", "Do the thing", "details")
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
assert.Equal(t, "/api/v1/repos/mathias/hyperguild/issues", gotPath, "owner forced to mathias")
|
||||||
|
assert.Equal(t, "token "+testToken, gotAuth)
|
||||||
|
assert.Contains(t, gotBody, "Do the thing")
|
||||||
|
assert.Equal(t, "hyperguild", ref.Repo)
|
||||||
|
assert.Equal(t, 42, ref.Number)
|
||||||
|
assert.Contains(t, ref.URL, "/issues/42")
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestCommentIssue(t *testing.T) {
|
||||||
|
var gotPath string
|
||||||
|
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||||
|
gotPath = r.URL.Path
|
||||||
|
w.WriteHeader(http.StatusCreated)
|
||||||
|
_ = json.NewEncoder(w).Encode(map[string]any{"html_url": "https://git/c/1"})
|
||||||
|
}))
|
||||||
|
defer srv.Close()
|
||||||
|
|
||||||
|
ref, err := gitea.New(srv.URL, testToken).CommentIssue(context.Background(), "hyperguild", 7, "a comment")
|
||||||
|
require.NoError(t, err)
|
||||||
|
assert.Equal(t, "/api/v1/repos/mathias/hyperguild/issues/7/comments", gotPath)
|
||||||
|
assert.Equal(t, 7, ref.Number)
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestCloseIssueWithComment(t *testing.T) {
|
||||||
|
var paths []string
|
||||||
|
var states []string
|
||||||
|
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||||
|
paths = append(paths, r.Method+" "+r.URL.Path)
|
||||||
|
if r.Method == http.MethodPatch {
|
||||||
|
var body map[string]any
|
||||||
|
b, _ := io.ReadAll(r.Body)
|
||||||
|
_ = json.Unmarshal(b, &body)
|
||||||
|
states = append(states, body["state"].(string))
|
||||||
|
}
|
||||||
|
w.WriteHeader(http.StatusOK)
|
||||||
|
_ = json.NewEncoder(w).Encode(map[string]any{"number": 9, "html_url": "https://git/i/9"})
|
||||||
|
}))
|
||||||
|
defer srv.Close()
|
||||||
|
|
||||||
|
ref, err := gitea.New(srv.URL, testToken).CloseIssue(context.Background(), "hyperguild", 9, "closing because done")
|
||||||
|
require.NoError(t, err)
|
||||||
|
assert.Equal(t, 9, ref.Number)
|
||||||
|
// Comment posted first, then state PATCHed to closed.
|
||||||
|
assert.Contains(t, paths, "POST /api/v1/repos/mathias/hyperguild/issues/9/comments")
|
||||||
|
assert.Contains(t, paths, "PATCH /api/v1/repos/mathias/hyperguild/issues/9")
|
||||||
|
assert.Equal(t, []string{"closed"}, states)
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestCloseIssueNoComment(t *testing.T) {
|
||||||
|
var commented bool
|
||||||
|
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||||
|
if strings.HasSuffix(r.URL.Path, "/comments") {
|
||||||
|
commented = true
|
||||||
|
}
|
||||||
|
w.WriteHeader(http.StatusOK)
|
||||||
|
_ = json.NewEncoder(w).Encode(map[string]any{"number": 3, "html_url": "https://git/i/3"})
|
||||||
|
}))
|
||||||
|
defer srv.Close()
|
||||||
|
|
||||||
|
_, err := gitea.New(srv.URL, testToken).CloseIssue(context.Background(), "hyperguild", 3, "")
|
||||||
|
require.NoError(t, err)
|
||||||
|
assert.False(t, commented, "empty comment ⇒ no comment POST")
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestErrorPathDoesNotLeakToken(t *testing.T) {
|
||||||
|
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
|
||||||
|
w.WriteHeader(http.StatusInternalServerError)
|
||||||
|
_, _ = w.Write([]byte("boom"))
|
||||||
|
}))
|
||||||
|
defer srv.Close()
|
||||||
|
|
||||||
|
_, err := gitea.New(srv.URL, testToken).CreateIssue(context.Background(), "hyperguild", "t", "b")
|
||||||
|
require.Error(t, err)
|
||||||
|
assert.NotContains(t, err.Error(), testToken, "token must never appear in an error message")
|
||||||
|
assert.Contains(t, err.Error(), "500")
|
||||||
|
}
|
||||||
@@ -9,8 +9,8 @@ import (
|
|||||||
"strings"
|
"strings"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/mathiasbq/hyperguild/ingestion/internal/api"
|
|
||||||
"github.com/mathiasbq/hyperguild/ingestion/internal/brain"
|
"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/extract"
|
||||||
"github.com/mathiasbq/hyperguild/ingestion/internal/graphsync"
|
"github.com/mathiasbq/hyperguild/ingestion/internal/graphsync"
|
||||||
"github.com/mathiasbq/hyperguild/ingestion/internal/pipeline"
|
"github.com/mathiasbq/hyperguild/ingestion/internal/pipeline"
|
||||||
@@ -38,7 +38,7 @@ func (s *Server) tools() []map[string]any {
|
|||||||
return b
|
return b
|
||||||
}
|
}
|
||||||
|
|
||||||
return []map[string]any{
|
tools := []map[string]any{
|
||||||
{
|
{
|
||||||
"name": "brain_query",
|
"name": "brain_query",
|
||||||
"description": "BM25 full-text search across brain/knowledge/ and brain/wiki/ markdown files. Optionally scope by wing (topic domain) and hall (memory type).",
|
"description": "BM25 full-text search across brain/knowledge/ and brain/wiki/ markdown files. Optionally scope by wing (topic domain) and hall (memory type).",
|
||||||
@@ -171,6 +171,13 @@ func (s *Server) tools() []map[string]any {
|
|||||||
}),
|
}),
|
||||||
},
|
},
|
||||||
}
|
}
|
||||||
|
// The capture relay tool (#55) is advertised only when wired via
|
||||||
|
// WithCapture — MCP-native harnesses (claude.ai, Crush, Pi, LLM Council)
|
||||||
|
// reach capture through it.
|
||||||
|
if s.capture != nil {
|
||||||
|
tools = append(tools, captureToolDescriptor())
|
||||||
|
}
|
||||||
|
return tools
|
||||||
}
|
}
|
||||||
|
|
||||||
type brainQueryArgs struct {
|
type brainQueryArgs struct {
|
||||||
@@ -219,7 +226,11 @@ func (s *Server) brainWrite(ctx context.Context, args json.RawMessage) (json.Raw
|
|||||||
if err := json.Unmarshal(args, &a); err != nil {
|
if err := json.Unmarshal(args, &a); err != nil {
|
||||||
return nil, fmt.Errorf("parse args: %w", err)
|
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,
|
Content: a.Content,
|
||||||
Filename: a.Filename,
|
Filename: a.Filename,
|
||||||
Type: a.Type,
|
Type: a.Type,
|
||||||
@@ -230,22 +241,7 @@ func (s *Server) brainWrite(ctx context.Context, args json.RawMessage) (json.Raw
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
// Auto-regenerate the wing _index.md when the write landed in the
|
return json.Marshal(map[string]string{"id": ref.ID, "path": ref.Path, "content_hash": ref.ContentHash})
|
||||||
// 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})
|
|
||||||
}
|
}
|
||||||
|
|
||||||
type brainUpdateArgs struct {
|
type brainUpdateArgs struct {
|
||||||
@@ -274,50 +270,26 @@ func (s *Server) brainUpdate(ctx context.Context, args json.RawMessage) (json.Ra
|
|||||||
return nil, fmt.Errorf("content is required")
|
return nil, fmt.Errorf("content is required")
|
||||||
}
|
}
|
||||||
|
|
||||||
opts := api.UpdateNoteOptions{Content: a.Content, Reason: a.Reason}
|
// path takes precedence over slug; the store treats any slug containing
|
||||||
switch {
|
// a slash as a full brain-relative path (issue #45: "slug ... OR path").
|
||||||
case a.Path != "":
|
slug := a.Slug
|
||||||
opts.Path = a.Path
|
if a.Path != "" {
|
||||||
case strings.Contains(a.Slug, "/"):
|
slug = a.Path
|
||||||
// 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
|
|
||||||
}
|
}
|
||||||
|
ref, err := s.store.Update(ctx, slug, capture.Note{
|
||||||
relPath, hash, _, err := api.UpdateNote(s.brainDir, opts)
|
Content: a.Content,
|
||||||
|
Wing: a.Wing,
|
||||||
|
Hall: a.Hall,
|
||||||
|
Reason: a.Reason,
|
||||||
|
})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
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{
|
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/<wing>/<hall>/<slug>.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 {
|
type brainGetArgs struct {
|
||||||
ID string `json:"id,omitempty"`
|
ID string `json:"id,omitempty"`
|
||||||
Path string `json:"path,omitempty"`
|
Path string `json:"path,omitempty"`
|
||||||
@@ -327,7 +299,7 @@ type brainGetArgs struct {
|
|||||||
// path — the de-facto handle). Read-only; the create-path read-after-
|
// path — the de-facto handle). Read-only; the create-path read-after-
|
||||||
// write primitive that lets callers confirm a write landed without a
|
// write primitive that lets callers confirm a write landed without a
|
||||||
// lexical re-query.
|
// 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
|
var a brainGetArgs
|
||||||
if err := json.Unmarshal(args, &a); err != nil {
|
if err := json.Unmarshal(args, &a); err != nil {
|
||||||
return nil, fmt.Errorf("parse args: %w", err)
|
return nil, fmt.Errorf("parse args: %w", err)
|
||||||
@@ -339,13 +311,13 @@ func (s *Server) brainGet(_ context.Context, args json.RawMessage) (json.RawMess
|
|||||||
if target == "" {
|
if target == "" {
|
||||||
return nil, fmt.Errorf("id or path is required")
|
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 {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
return json.Marshal(map[string]any{
|
return json.Marshal(map[string]any{
|
||||||
"id": target, "path": target, "content_hash": hash,
|
"id": note.ID, "path": note.Path, "content_hash": note.ContentHash,
|
||||||
"frontmatter": fm, "body": body,
|
"frontmatter": note.Frontmatter, "body": note.Body,
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -1,7 +1,8 @@
|
|||||||
// Package mcp implements an MCP HTTP handler for the ingestion service.
|
// Package mcp implements an MCP HTTP handler for the ingestion service.
|
||||||
// Exposed tools: brain_query, brain_write, brain_update, brain_get,
|
// Exposed tools: brain_query, brain_write, brain_update, brain_get,
|
||||||
// brain_index, brain_tunnel, brain_ingest, brain_ingest_raw,
|
// brain_index, brain_tunnel, brain_ingest, brain_ingest_raw,
|
||||||
// brain_answer, brain_classify, brain_graph, brain_context, session_log.
|
// brain_answer, brain_classify, brain_graph, brain_context, session_log,
|
||||||
|
// and capture (the #55 relay tool, registered only when WithCapture is set).
|
||||||
package mcp
|
package mcp
|
||||||
|
|
||||||
import (
|
import (
|
||||||
@@ -10,6 +11,9 @@ import (
|
|||||||
"fmt"
|
"fmt"
|
||||||
"net/http"
|
"net/http"
|
||||||
|
|
||||||
|
"github.com/mathiasbq/hyperguild/ingestion/internal/brainstore"
|
||||||
|
"github.com/mathiasbq/hyperguild/ingestion/internal/capture"
|
||||||
|
"github.com/mathiasbq/hyperguild/ingestion/internal/capturehttp"
|
||||||
"github.com/mathiasbq/hyperguild/ingestion/internal/graphstore"
|
"github.com/mathiasbq/hyperguild/ingestion/internal/graphstore"
|
||||||
"github.com/mathiasbq/hyperguild/ingestion/internal/graphsync"
|
"github.com/mathiasbq/hyperguild/ingestion/internal/graphsync"
|
||||||
"github.com/mathiasbq/hyperguild/ingestion/internal/pipeline"
|
"github.com/mathiasbq/hyperguild/ingestion/internal/pipeline"
|
||||||
@@ -46,6 +50,21 @@ type Server struct {
|
|||||||
vector search.VectorSearcher // nil = BM25-only retrieval
|
vector search.VectorSearcher // nil = BM25-only retrieval
|
||||||
embedder search.Embedder // nil = BM25-only retrieval
|
embedder search.Embedder // nil = BM25-only retrieval
|
||||||
graph graphsync.Store // nil = brain_graph and GraphRAG augmentation disabled
|
graph graphsync.Store // nil = brain_graph and GraphRAG augmentation disabled
|
||||||
|
store *brainstore.Store // shared brain write/update/get impl (also used by capture)
|
||||||
|
tracker capture.IssueTracker // nil = no Gitea ticket integration; wired for capture (#53)
|
||||||
|
capture *captureDeps // nil = capture MCP tool disabled (#55 relay)
|
||||||
|
}
|
||||||
|
|
||||||
|
// captureDeps holds what the MCP `capture` tool (the #55 relay door for
|
||||||
|
// MCP-native harnesses like claude.ai) needs: the use-case, the auth bits
|
||||||
|
// to re-derive the caller's principal from the Bearer header (the chassis
|
||||||
|
// middleware gates but discards the principal), and the origin resolver.
|
||||||
|
type captureDeps struct {
|
||||||
|
svc *capture.Service
|
||||||
|
validator capturehttp.Validator
|
||||||
|
staticToken string
|
||||||
|
staticPrincipal string
|
||||||
|
resolver capturehttp.OriginResolver
|
||||||
}
|
}
|
||||||
|
|
||||||
// NewServer constructs a Server bound to brainDir. pipelineCfg supplies the
|
// NewServer constructs a Server bound to brainDir. pipelineCfg supplies the
|
||||||
@@ -56,7 +75,13 @@ func NewServer(brainDir string, pipelineCfg *pipeline.Config, llm pipeline.Compl
|
|||||||
if pipelineCfg != nil {
|
if pipelineCfg != nil {
|
||||||
cfg = *pipelineCfg
|
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,
|
// WithReranker installs an opt-in cross-encoder reranker. When set,
|
||||||
@@ -84,9 +109,55 @@ func (s *Server) WithHybridRetrieval(v search.VectorSearcher, e search.Embedder)
|
|||||||
func (s *Server) WithGraph(g *graphstore.PGStore) *Server {
|
func (s *Server) WithGraph(g *graphstore.PGStore) *Server {
|
||||||
if g == nil {
|
if g == nil {
|
||||||
s.graph = nil
|
s.graph = nil
|
||||||
|
s.store.WithGraph(nil)
|
||||||
return s
|
return s
|
||||||
}
|
}
|
||||||
s.graph = g
|
s.graph = g
|
||||||
|
s.store.WithGraph(g)
|
||||||
|
return s
|
||||||
|
}
|
||||||
|
|
||||||
|
// WithIssueTracker injects the Gitea ticket tracker behind the
|
||||||
|
// capture.IssueTracker interface. nil leaves ticket integration off. The
|
||||||
|
// use-case (capture) consumes this in #53; it is wired here so the
|
||||||
|
// dependency is constructed once and stays swappable/testable.
|
||||||
|
func (s *Server) WithIssueTracker(t capture.IssueTracker) *Server {
|
||||||
|
s.tracker = t
|
||||||
|
return s
|
||||||
|
}
|
||||||
|
|
||||||
|
// IssueTracker returns the injected ticket tracker (nil when unconfigured).
|
||||||
|
func (s *Server) IssueTracker() capture.IssueTracker {
|
||||||
|
return s.tracker
|
||||||
|
}
|
||||||
|
|
||||||
|
// BrainStore returns the shared brain store (graph-wired once WithGraph
|
||||||
|
// has run), so the capture use-case writes through the exact same
|
||||||
|
// implementation as the MCP handlers.
|
||||||
|
func (s *Server) BrainStore() *brainstore.Store {
|
||||||
|
return s.store
|
||||||
|
}
|
||||||
|
|
||||||
|
// WithCapture enables the MCP `capture` tool (#55) — the relay door for
|
||||||
|
// MCP-native harnesses (claude.ai, Crush, Pi, LLM Council) that cannot run
|
||||||
|
// the use-case in-process. It forwards to the same CaptureService as
|
||||||
|
// POST /capture, deriving the caller's principal + origin from the same
|
||||||
|
// auth credentials that gate /mcp. nil svc leaves the tool unregistered.
|
||||||
|
func (s *Server) WithCapture(svc *capture.Service, validator capturehttp.Validator, staticToken, staticPrincipal string, resolver capturehttp.OriginResolver) *Server {
|
||||||
|
if svc == nil {
|
||||||
|
s.capture = nil
|
||||||
|
return s
|
||||||
|
}
|
||||||
|
if staticPrincipal == "" {
|
||||||
|
staticPrincipal = "local-cli"
|
||||||
|
}
|
||||||
|
s.capture = &captureDeps{
|
||||||
|
svc: svc,
|
||||||
|
validator: validator,
|
||||||
|
staticToken: staticToken,
|
||||||
|
staticPrincipal: staticPrincipal,
|
||||||
|
resolver: resolver,
|
||||||
|
}
|
||||||
return s
|
return s
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -139,7 +210,18 @@ func (s *Server) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
|||||||
rpcErr = &rpcError{Code: -32602, Message: "invalid params"}
|
rpcErr = &rpcError{Code: -32602, Message: "invalid params"}
|
||||||
break
|
break
|
||||||
}
|
}
|
||||||
out, err := s.handleCall(r.Context(), p.Name, p.Arguments)
|
// Re-derive the authenticated principal from the Bearer header so
|
||||||
|
// the capture tool can compute the trust-zone origin. The request
|
||||||
|
// is already gated by BearerMiddleware; this only recovers the
|
||||||
|
// identity that middleware discards.
|
||||||
|
ctx := r.Context()
|
||||||
|
if s.capture != nil {
|
||||||
|
if principal, viaStatic, ok := capturehttp.Authenticate(
|
||||||
|
r, s.capture.staticToken, s.capture.staticPrincipal, s.capture.validator); ok {
|
||||||
|
ctx = withPrincipal(ctx, principal, viaStatic)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
out, err := s.handleCall(ctx, p.Name, p.Arguments)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
rpcErr = &rpcError{Code: -32000, Message: err.Error()}
|
rpcErr = &rpcError{Code: -32000, Message: err.Error()}
|
||||||
break
|
break
|
||||||
@@ -181,6 +263,8 @@ func (s *Server) handleCall(ctx context.Context, name string, args json.RawMessa
|
|||||||
return s.brainUpdate(ctx, args)
|
return s.brainUpdate(ctx, args)
|
||||||
case "brain_get":
|
case "brain_get":
|
||||||
return s.brainGet(ctx, args)
|
return s.brainGet(ctx, args)
|
||||||
|
case "capture":
|
||||||
|
return s.brainCapture(ctx, args)
|
||||||
case "brain_index":
|
case "brain_index":
|
||||||
return s.brainIndex(ctx, args)
|
return s.brainIndex(ctx, args)
|
||||||
case "brain_tunnel":
|
case "brain_tunnel":
|
||||||
|
|||||||
@@ -2,12 +2,14 @@ package mcp_test
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"bytes"
|
"bytes"
|
||||||
|
"context"
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
"net/http"
|
"net/http"
|
||||||
"net/http/httptest"
|
"net/http/httptest"
|
||||||
"strings"
|
"strings"
|
||||||
"testing"
|
"testing"
|
||||||
|
|
||||||
|
"github.com/mathiasbq/hyperguild/ingestion/internal/capture"
|
||||||
"github.com/mathiasbq/hyperguild/ingestion/internal/mcp"
|
"github.com/mathiasbq/hyperguild/ingestion/internal/mcp"
|
||||||
"github.com/stretchr/testify/assert"
|
"github.com/stretchr/testify/assert"
|
||||||
"github.com/stretchr/testify/require"
|
"github.com/stretchr/testify/require"
|
||||||
@@ -93,3 +95,22 @@ func TestServerUnknownMethodReturnsError(t *testing.T) {
|
|||||||
assert.Equal(t, float64(-32601), errObj["code"])
|
assert.Equal(t, float64(-32601), errObj["code"])
|
||||||
assert.Contains(t, errObj["message"].(string), "unknown/method")
|
assert.Contains(t, errObj["message"].(string), "unknown/method")
|
||||||
}
|
}
|
||||||
|
|
||||||
|
type stubTracker struct{}
|
||||||
|
|
||||||
|
func (stubTracker) CreateIssue(context.Context, string, string, string) (capture.IssueRef, error) {
|
||||||
|
return capture.IssueRef{}, nil
|
||||||
|
}
|
||||||
|
func (stubTracker) CloseIssue(context.Context, string, int, string) (capture.IssueRef, error) {
|
||||||
|
return capture.IssueRef{}, nil
|
||||||
|
}
|
||||||
|
func (stubTracker) CommentIssue(context.Context, string, int, string) (capture.IssueRef, error) {
|
||||||
|
return capture.IssueRef{}, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestWithIssueTrackerInjects(t *testing.T) {
|
||||||
|
srv := mcp.NewServer(t.TempDir(), nil, nil, nil)
|
||||||
|
assert.Nil(t, srv.IssueTracker(), "tracker is off by default")
|
||||||
|
srv = srv.WithIssueTracker(stubTracker{})
|
||||||
|
assert.NotNil(t, srv.IssueTracker(), "tracker injected behind the interface")
|
||||||
|
}
|
||||||
|
|||||||
@@ -0,0 +1,105 @@
|
|||||||
|
package mcp
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"encoding/json"
|
||||||
|
"fmt"
|
||||||
|
|
||||||
|
"github.com/mathiasbq/hyperguild/ingestion/internal/capturehttp"
|
||||||
|
)
|
||||||
|
|
||||||
|
// principalKey is the context key under which the authenticated principal
|
||||||
|
// (re-derived in ServeHTTP) is stashed for the capture tool.
|
||||||
|
type principalKeyT struct{}
|
||||||
|
|
||||||
|
var principalKey principalKeyT
|
||||||
|
|
||||||
|
type principalInfo struct {
|
||||||
|
principal string
|
||||||
|
viaStatic bool
|
||||||
|
}
|
||||||
|
|
||||||
|
func withPrincipal(ctx context.Context, principal string, viaStatic bool) context.Context {
|
||||||
|
return context.WithValue(ctx, principalKey, principalInfo{principal: principal, viaStatic: viaStatic})
|
||||||
|
}
|
||||||
|
|
||||||
|
// captureToolDescriptor is the tools/list entry for the capture relay.
|
||||||
|
// Appended only when WithCapture has wired the tool.
|
||||||
|
func captureToolDescriptor() map[string]any {
|
||||||
|
str := func(d string) map[string]any { return map[string]any{"type": "string", "description": d} }
|
||||||
|
insightItem := map[string]any{
|
||||||
|
"type": "object",
|
||||||
|
"properties": map[string]any{
|
||||||
|
"text": str("the insight body"), "wing": str("brain wing"),
|
||||||
|
"hall": str("brain hall (facts/decisions/failures/hypotheses/sources)"),
|
||||||
|
"supersede_slug": str("optional: slug of a prior note to revise in place instead of creating"),
|
||||||
|
},
|
||||||
|
"required": []string{"text", "wing", "hall"},
|
||||||
|
}
|
||||||
|
ticketItem := map[string]any{
|
||||||
|
"type": "object",
|
||||||
|
"properties": map[string]any{
|
||||||
|
"repo": str("gitea repo (owner is always mathias)"), "action": str("create|close|comment"),
|
||||||
|
"number": map[string]any{"type": "integer", "description": "issue number (close/comment)"},
|
||||||
|
"title": str("issue title (create)"), "body": str("issue/comment body"),
|
||||||
|
},
|
||||||
|
"required": []string{"repo", "action"},
|
||||||
|
}
|
||||||
|
schema := map[string]any{
|
||||||
|
"type": "object",
|
||||||
|
"properties": map[string]any{
|
||||||
|
"context": map[string]any{
|
||||||
|
"type": "object",
|
||||||
|
"properties": map[string]any{
|
||||||
|
"harness": str("descriptive harness label (telemetry only, never a gate input)"),
|
||||||
|
"session_ref": str("optional session reference"), "fidelity": str("live-capture|transcript-parse|agent-runlog"),
|
||||||
|
"actor": str("acting user/agent"), "classification": str("caller-declared sensitivity: public|internal|confidential"),
|
||||||
|
},
|
||||||
|
},
|
||||||
|
"insights": map[string]any{"type": "array", "items": insightItem},
|
||||||
|
"tickets": map[string]any{"type": "array", "items": ticketItem},
|
||||||
|
"summary": map[string]any{"type": "object", "properties": map[string]any{
|
||||||
|
"title": str("summary title"), "body": str("summary body"),
|
||||||
|
"repos_touched": map[string]any{"type": "array", "items": map[string]any{"type": "string"}},
|
||||||
|
}},
|
||||||
|
"dry_run": map[string]any{"type": "boolean", "description": "validate + return the would-be receipt, write nothing"},
|
||||||
|
},
|
||||||
|
}
|
||||||
|
b, _ := json.Marshal(schema)
|
||||||
|
return map[string]any{
|
||||||
|
"name": "capture",
|
||||||
|
"description": "Persist a session's value uniformly: insights → brain (write or supersede), action items → Gitea tickets, optional summary → ai-sessions. The relay door for MCP-native harnesses. Origin is server-derived from your authenticated identity; confidential captures through a us-nexus surface are refused (I1). Returns a partial-aware receipt.",
|
||||||
|
"inputSchema": json.RawMessage(b),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// brainCapture is the MCP capture tool: the #55 relay for MCP-native
|
||||||
|
// harnesses. It re-uses the same CaptureService, principal-derivation, and
|
||||||
|
// origin resolver as POST /capture — only the transport differs. It holds
|
||||||
|
// no state and retains nothing beyond the I5 audit record.
|
||||||
|
func (s *Server) brainCapture(ctx context.Context, args json.RawMessage) (json.RawMessage, error) {
|
||||||
|
if s.capture == nil {
|
||||||
|
return nil, fmt.Errorf("capture tool not configured")
|
||||||
|
}
|
||||||
|
info, ok := ctx.Value(principalKey).(principalInfo)
|
||||||
|
if !ok || info.principal == "" {
|
||||||
|
// No authenticated principal ⇒ cannot derive origin ⇒ cannot gate.
|
||||||
|
return nil, fmt.Errorf("capture requires an authenticated principal")
|
||||||
|
}
|
||||||
|
|
||||||
|
in, err := capturehttp.DecodeRequest(args)
|
||||||
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("invalid capture request: %w", err)
|
||||||
|
}
|
||||||
|
// Principal and origin are server-derived — never taken from the body.
|
||||||
|
in.Context.Principal = info.principal
|
||||||
|
in.Context.Origin = s.capture.resolver.Resolve(info.principal, info.viaStatic)
|
||||||
|
|
||||||
|
rec, err := s.capture.svc.Capture(ctx, in)
|
||||||
|
if err != nil {
|
||||||
|
// Surface I1/I5 refusals and validation failures verbatim; errors.Is
|
||||||
|
// markers (ErrSovereigntyRefused / ErrAuditUnavailable) ride in the message.
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
return json.Marshal(rec)
|
||||||
|
}
|
||||||
@@ -0,0 +1,150 @@
|
|||||||
|
package mcp_test
|
||||||
|
|
||||||
|
import (
|
||||||
|
"bytes"
|
||||||
|
"context"
|
||||||
|
"encoding/json"
|
||||||
|
"errors"
|
||||||
|
"net/http"
|
||||||
|
"net/http/httptest"
|
||||||
|
"os"
|
||||||
|
"path/filepath"
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"github.com/mathiasbq/hyperguild/ingestion/internal/audit"
|
||||||
|
"github.com/mathiasbq/hyperguild/ingestion/internal/brainstore"
|
||||||
|
"github.com/mathiasbq/hyperguild/ingestion/internal/capture"
|
||||||
|
"github.com/mathiasbq/hyperguild/ingestion/internal/capturehttp"
|
||||||
|
"github.com/mathiasbq/hyperguild/ingestion/internal/classification"
|
||||||
|
"github.com/mathiasbq/hyperguild/ingestion/internal/mcp"
|
||||||
|
"github.com/stretchr/testify/assert"
|
||||||
|
"github.com/stretchr/testify/require"
|
||||||
|
)
|
||||||
|
|
||||||
|
const capStaticTok = "cap-static-tok"
|
||||||
|
|
||||||
|
type capFakeTracker struct{}
|
||||||
|
|
||||||
|
func (capFakeTracker) CreateIssue(context.Context, string, string, string) (capture.IssueRef, error) {
|
||||||
|
return capture.IssueRef{Repo: "hyperguild", Number: 1, URL: "https://git/1"}, nil
|
||||||
|
}
|
||||||
|
func (capFakeTracker) CloseIssue(context.Context, string, int, string) (capture.IssueRef, error) {
|
||||||
|
return capture.IssueRef{}, nil
|
||||||
|
}
|
||||||
|
func (capFakeTracker) CommentIssue(context.Context, string, int, string) (capture.IssueRef, error) {
|
||||||
|
return capture.IssueRef{}, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
type capFakeValidator struct {
|
||||||
|
subject string
|
||||||
|
err error
|
||||||
|
}
|
||||||
|
|
||||||
|
func (v capFakeValidator) Validate(context.Context, string) (string, error) {
|
||||||
|
return v.subject, v.err
|
||||||
|
}
|
||||||
|
|
||||||
|
func captureServer(t *testing.T, validator capturehttp.Validator, sovereign []string) (*mcp.Server, string) {
|
||||||
|
t.Helper()
|
||||||
|
brainDir := t.TempDir()
|
||||||
|
cfg, err := classification.Load(brainDir)
|
||||||
|
require.NoError(t, err)
|
||||||
|
svc := capture.NewService(brainstore.New(brainDir), capFakeTracker{}, nil, cfg, audit.NewSlogSink(nil))
|
||||||
|
srv := mcp.NewServer(brainDir, nil, nil, nil)
|
||||||
|
srv.WithCapture(svc, validator, capStaticTok, "local-cli", capturehttp.NewOriginResolver(sovereign))
|
||||||
|
return srv, brainDir
|
||||||
|
}
|
||||||
|
|
||||||
|
func captureCall(t *testing.T, srv http.Handler, authz string, args map[string]any) map[string]any {
|
||||||
|
t.Helper()
|
||||||
|
body, _ := json.Marshal(map[string]any{
|
||||||
|
"jsonrpc": "2.0", "id": 1, "method": "tools/call",
|
||||||
|
"params": map[string]any{"name": "capture", "arguments": args},
|
||||||
|
})
|
||||||
|
req := httptest.NewRequest(http.MethodPost, "/mcp", bytes.NewReader(body))
|
||||||
|
if authz != "" {
|
||||||
|
req.Header.Set("Authorization", authz)
|
||||||
|
}
|
||||||
|
rr := httptest.NewRecorder()
|
||||||
|
srv.ServeHTTP(rr, req)
|
||||||
|
var resp map[string]any
|
||||||
|
require.NoError(t, json.Unmarshal(rr.Body.Bytes(), &resp))
|
||||||
|
return resp
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestCaptureToolListedWhenWired(t *testing.T) {
|
||||||
|
srv, _ := captureServer(t, nil, nil)
|
||||||
|
body, _ := json.Marshal(map[string]any{"jsonrpc": "2.0", "id": 1, "method": "tools/list"})
|
||||||
|
req := httptest.NewRequest(http.MethodPost, "/mcp", bytes.NewReader(body))
|
||||||
|
rr := httptest.NewRecorder()
|
||||||
|
srv.ServeHTTP(rr, req)
|
||||||
|
assert.Contains(t, rr.Body.String(), `"capture"`)
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestCaptureToolNotListedByDefault(t *testing.T) {
|
||||||
|
srv := mcp.NewServer(t.TempDir(), nil, nil, nil) // no WithCapture
|
||||||
|
body, _ := json.Marshal(map[string]any{"jsonrpc": "2.0", "id": 1, "method": "tools/list"})
|
||||||
|
req := httptest.NewRequest(http.MethodPost, "/mcp", bytes.NewReader(body))
|
||||||
|
rr := httptest.NewRecorder()
|
||||||
|
srv.ServeHTTP(rr, req)
|
||||||
|
assert.NotContains(t, rr.Body.String(), `"capture"`)
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestCaptureToolForwardsViaStaticPrincipal(t *testing.T) {
|
||||||
|
srv, brainDir := captureServer(t, nil, nil)
|
||||||
|
resp := captureCall(t, srv, "Bearer "+capStaticTok, map[string]any{
|
||||||
|
"context": map[string]any{"harness": "claude-code", "actor": "mathias", "classification": "internal"},
|
||||||
|
"insights": []map[string]any{{"text": "a fact", "wing": "hyperguild", "hall": "facts"}},
|
||||||
|
"tickets": []map[string]any{{"repo": "hyperguild", "action": "create", "title": "t"}},
|
||||||
|
})
|
||||||
|
require.Nil(t, resp["error"], "got error: %v", resp["error"])
|
||||||
|
text := resp["result"].(map[string]any)["content"].([]any)[0].(map[string]any)["text"].(string)
|
||||||
|
var rec capture.CaptureReceipt
|
||||||
|
require.NoError(t, json.Unmarshal([]byte(text), &rec))
|
||||||
|
assert.True(t, rec.Insights[0].OK)
|
||||||
|
assert.True(t, rec.Tickets[0].OK)
|
||||||
|
// Forwarded to the real brain store.
|
||||||
|
_, statErr := os.Stat(filepath.Join(brainDir, "wiki/hyperguild/facts"))
|
||||||
|
require.NoError(t, statErr)
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestCaptureToolRefusesConfidentialViaUSNexus(t *testing.T) {
|
||||||
|
// JWT principal not in the sovereign allowlist ⇒ us-nexus origin.
|
||||||
|
srv, _ := captureServer(t, capFakeValidator{subject: "claudeai-oauth"}, nil)
|
||||||
|
resp := captureCall(t, srv, "Bearer jwt-token", map[string]any{
|
||||||
|
"context": map[string]any{"harness": "claudeai-chat", "actor": "mathias", "classification": "confidential"},
|
||||||
|
"insights": []map[string]any{{"text": "secret", "wing": "client-seb", "hall": "facts"}},
|
||||||
|
})
|
||||||
|
require.NotNil(t, resp["error"])
|
||||||
|
assert.Contains(t, resp["error"].(map[string]any)["message"].(string), "sovereignty")
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestCaptureToolAllowsConfidentialViaSovereignJWT(t *testing.T) {
|
||||||
|
srv, _ := captureServer(t, capFakeValidator{subject: "koala-cli"}, []string{"koala-cli"})
|
||||||
|
resp := captureCall(t, srv, "Bearer jwt-token", map[string]any{
|
||||||
|
"context": map[string]any{"harness": "claude-code", "actor": "mathias", "classification": "confidential"},
|
||||||
|
"insights": []map[string]any{{"text": "secret", "wing": "client-seb", "hall": "facts"}},
|
||||||
|
})
|
||||||
|
assert.Nil(t, resp["error"], "sovereign JWT principal should be allowed: %v", resp["error"])
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestCaptureToolRejectsUnauthenticated(t *testing.T) {
|
||||||
|
srv, _ := captureServer(t, capFakeValidator{err: errors.New("no jwt")}, nil)
|
||||||
|
resp := captureCall(t, srv, "", map[string]any{ // no Authorization
|
||||||
|
"context": map[string]any{"harness": "x", "classification": "internal"},
|
||||||
|
"insights": []map[string]any{{"text": "a", "wing": "hyperguild", "hall": "facts"}},
|
||||||
|
})
|
||||||
|
require.NotNil(t, resp["error"])
|
||||||
|
assert.Contains(t, resp["error"].(map[string]any)["message"].(string), "authenticated principal")
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestCaptureToolCallerCannotForgeOrigin(t *testing.T) {
|
||||||
|
// Body asserts sovereign harness, but the us-nexus JWT principal governs.
|
||||||
|
srv, _ := captureServer(t, capFakeValidator{subject: "claudeai-oauth"}, nil)
|
||||||
|
resp := captureCall(t, srv, "Bearer jwt", map[string]any{
|
||||||
|
"context": map[string]any{"harness": "sovereign-soil", "classification": "confidential"},
|
||||||
|
"insights": []map[string]any{{"text": "secret", "wing": "client-seb", "hall": "facts"}},
|
||||||
|
})
|
||||||
|
require.NotNil(t, resp["error"])
|
||||||
|
assert.Contains(t, resp["error"].(map[string]any)["message"].(string), "sovereignty")
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user