Compare commits

..
Author SHA1 Message Date
mathiasandClaude Opus 4.8 5288554338 fix(gitea): WriteFile creates via POST, updates via PUT (real contents API)
CI / Lint / Test / Vet (pull_request) Successful in 13s
CI / Mirror to GitHub (pull_request) Has been skipped
A live token-scope probe against mathias/ai-sessions revealed gitea's
contents API uses POST to create and PUT (sha required) to update — the
first impl always PUT'd, so creating a new summary 422'd "[SHA]: Required".
The httptest mock had the same wrong assumption. Pick the method by
whether the file exists (GET sha). Token confirmed contents:write
(push:true) by the probe.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-23 17:22:37 +02:00
mathiasandClaude Opus 4.8 2368564523 fix(capture): wire ai-sessions SummaryWriter into the relay (#66)
CI / Lint / Test / Vet (pull_request) Successful in 13s
CI / Mirror to GitHub (pull_request) Has been skipped
The deployed CaptureService was constructed with a nil SummaryWriter, so
a capture carrying a summary block returned the partial-failure
"no summary writer configured" (surfaced in the 2026-06-23 claude.ai
dogfood). The Gitea client already reaches mathias/* over BRAIN_GITEA_TOKEN
and now implements SummaryWriter, so inject it (type-asserted from the
tracker) — no new credential, no manifest change. Session summaries now
write to mathias/ai-sessions at summaries/<harness>/<YYYY-MM>/...

Reuses the existing token deliberately; assumes it carries contents:write
scope on ai-sessions (it already does issue writes for the tracker). If
the token is issue-scoped only, the live write 422s — a token-scope widen,
not a code fix.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-23 10:47:00 +02:00
mathiasandClaude Opus 4.8 0e28b2125b feat(gitea): WriteFile contents-API upsert — satisfies SummaryWriter (#66)
Adds Client.WriteFile (Gitea contents API) so the Gitea client also
implements capture.SummaryWriter. Upserts: a GET resolves the current
blob sha so an existing file is updated (the richer-fidelity-supersedes
rule for re-captured sessions) rather than 422'd. Owner stays the fixed
const; token only in the Authorization header (no leak — regression
tested). Refactors the HTTP path into a shared request() helper so the
contents flow can branch on 404 without it being an error.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-23 10:47:00 +02:00
mathiasandClaude Opus 4.8 06e21c019e docs(capture): as-built implementation report for the #49 epic
CI / Lint / Test / Vet (push) Successful in 12s
CI / Mirror to GitHub (push) Successful in 4s
Records what shipped (v0.11.0): architecture, sub-issue→PR map, I1–I5
compliance, the operational env reference, test coverage, and the
deferred follow-ups. Companion to specs/capture-bdd-spec.md (the design
contract) for onboarding + future audit.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-23 07:41:24 +02:00
mathias 76514215f4 Merge pull request 'feat: capture relay — MCP capture tool for non-library harnesses (#55, capture 49f)' (#61) from feat/capture-mcp-relay into main
CI / Lint / Test / Vet (push) Successful in 12s
CI / Mirror to GitHub (push) Has been skipped
2026-06-23 05:26:55 +00:00
mathiasandClaude Opus 4.8 7cf5bc221d feat(mcp): capture relay tool — MCP door for non-library harnesses (#55)
CI / Lint / Test / Vet (pull_request) Successful in 13s
CI / Mirror to GitHub (pull_request) Has been skipped
Adds the `capture` MCP tool: the #55 relay for harnesses that cannot run
the use-case in-process (claude.ai Chat/Cowork/Design, Crush, Pi, LLM
Council). They reach it through the existing /mcp OAuth connector.

- Thin: forwards to the SAME CaptureService as POST /capture; holds no
  state and retains nothing beyond the I5 audit record. The containment
  properties accepted in infra security-baseline (I2 ledger) hold by
  construction.
- Per-principal: ServeHTTP re-derives the caller's principal from the
  Bearer header (the chassis middleware gates but discards it) and stashes
  it in context; the tool resolves the trust-zone origin from it. A
  caller-asserted harness/origin in the body is ignored — origin is
  server-derived, so the I1 confidential refusal still fires for us-nexus
  callers (claude.ai), and sovereign-allowlisted JWT principals pass.
- Registered only when WithCapture is wired (all three sites: tools(),
  handleCall, package doc); main wires REST + MCP from the same service,
  resolver, and credentials.

Tests: listed-only-when-wired, forwards-via-static-principal,
confidential-via-us-nexus-refused, confidential-via-sovereign-allowed,
unauthenticated-rejected, caller-cannot-forge-origin.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-23 00:21:11 +02:00
mathiasandClaude Opus 4.8 f78a5474a5 refactor(capturehttp): export Authenticate + DecodeRequest for reuse (#55)
Lifts the Bearer principal-derivation and the request→CaptureInput decode
out of the REST handler into exported package funcs, so the MCP capture
tool (#55 relay) reuses the exact same auth precedence and wire shape —
one implementation, not two. No behaviour change to POST /capture.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-23 00:21:11 +02:00
mathias c307b72bd5 Merge pull request 'feat: I5 audit path + classification-aware degradation (#54, capture 49e)' (#60) from feat/capture-audit-degradation into main
CI / Lint / Test / Vet (push) Successful in 13s
CI / Mirror to GitHub (push) Successful in 4s
2026-06-22 21:57:52 +00:00
mathiasandClaude Opus 4.8 38a2e91002 feat(capturehttp): 503 on audit-unavailable; wire degrading sink (#54)
CI / Lint / Test / Vet (pull_request) Successful in 12s
CI / Mirror to GitHub (pull_request) Has been skipped
- Map capture.ErrAuditUnavailable → HTTP 503 (audit substrate down /
  confidential unauditable / floor).
- main: buildAuditSink selects the DegradingSink (loki central + durable
  file buffer under brainDir + optional ntfy) when BRAIN_LOKI_URL is set
  and starts the reconcile loop; else the plain slog sink. Notifier kept
  as a nil interface (not typed-nil) when unconfigured so the sink and
  reconcile skip it cleanly.

Env: BRAIN_LOKI_URL, BRAIN_NTFY_URL, BRAIN_NTFY_TOKEN,
BRAIN_AUDIT_RECONCILE_INTERVAL (default 60s). Buffer at
<brain>/.audit-buffer/capture.jsonl.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-22 23:54:24 +02:00
mathiasandClaude Opus 4.8 77f5e06d6b feat(audit): DegradingSink + durable buffer + loki/ntfy + reconcile (#54)
The classification-aware I5 audit path (§4.4):

- DegradingSink.Reserve: central up → AuditCentral; central down +
  confidential → refuse (no buffer); central down + internal/public +
  buffer writable → AuditBuffered; central down + buffer unwritable →
  floor refuse. Record executes the reserved outcome and, when buffered,
  fires an ntfy alert.
- FileBuffer: durable JSONL buffer that survives process restart; Confirm
  rewrites the file without a record, so a buffered record is cleared ONLY
  after its central write is confirmed.
- LokiCentral: /ready probe + /loki/api/v1/push (full audit entry as the
  structured line). NtfyNotifier: degraded-state alerts; token only in the
  auth header, never logged (regression-tested).
- Reconcile + StartReconcile: replay buffered records to central on
  recovery, confirm-then-clear per record; a failed push keeps the record
  buffered (no loss). SlogSink updated to the two-phase shape (always
  central, never fails) — the default when no loki endpoint is set.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-22 23:54:24 +02:00
mathiasandClaude Opus 4.8 202212e8d5 feat(capture): two-phase classification-aware AuditSink port (#54)
Splits the audit port into Reserve (before any write) + Record (after),
so "confidential + sink-down → refuse before any write" is literally true
even though the audit record — which lists what landed — can only be
written afterwards.

- AuditSink.Reserve(ctx, level) → AuditOutcome | error. The error path
  refuses the capture before writing: confidential + central sink down,
  or the all-tiers floor (nothing can record).
- AuditSink.Record(ctx, entry, outcome) persists per the reserved outcome.
- Service: I5 gate runs after the I1 gate and after the dry-run
  short-circuit (dry-run never probes the sink). AuditBuffered surfaces on
  the receipt. New ErrAuditUnavailable sentinel (→ HTTP 503).

The tier→behaviour decision lives in the sink impl (#54's DegradingSink),
not the service — the service just honours Reserve's verdict.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-22 23:54:24 +02:00
mathias b7938d4636 Merge pull request 'feat: POST /capture REST adapter + OAuth2 + I1 sovereignty gate (#53, capture 49d)' (#59) from feat/capture-rest-i1-gate into main
CI / Lint / Test / Vet (push) Successful in 12s
CI / Mirror to GitHub (push) Successful in 4s
2026-06-22 21:44:14 +00:00
mathiasandClaude Opus 4.8 a1997838b0 feat(capturehttp): POST /capture REST adapter + OAuth2 + origin resolver (#53)
CI / Lint / Test / Vet (pull_request) Successful in 12s
CI / Mirror to GitHub (pull_request) Has been skipped
The HTTP door for the capture capability. Thin: authenticate → derive
trust-zone origin → decode → capture.Service → map receipt to status.

- Auth mirrors the chassis Bearer precedence (static token wins, then Dex
  JWT) but returns the resolved principal + auth path, which the chassis
  middleware hides — capture needs the principal to derive the origin.
  Depends on a small Validator interface (the chassis *JWTValidator
  satisfies it) so the JWT/origin path is testable without a live JWKS.
- OriginResolver maps principal → trust zone: static-token caller and
  allowlisted JWT subjects → sovereign; every other principal → us-nexus
  (fail safe, so the I1 gate refuses confidential by default). Principal
  and origin are server-set on the input, overwriting any body the caller
  sent.
- HTTP status: 200 all-ok / dry-run, 207 partial, 502 all-failed, 403 on
  the I1 refusal, 400 on fail-closed validation.
- Wired in main behind the same static+JWT credentials as /mcp, reusing
  the MCP server's graph-wired brain store (one implementation) and the
  classification tags (#50). Mounts only when a Gitea tracker is
  configured. Sovereign JWT principals via BRAIN_CAPTURE_SOVEREIGN_PRINCIPALS.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-22 23:42:18 +02:00
mathiasandClaude Opus 4.8 77680c7445 feat(audit): minimal slog AuditSink for capture I5 (#53)
Emits the request-level audit record to structured logs (scraped by the
existing alloy/loki substrate) and surfaces security events at warn
level. Placeholder behind the AuditSink interface — the classification-
aware loki+buffer+reconcile sink (confidential fails closed, internal
degrades) lands in #54 and replaces this without touching callers.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-22 23:42:18 +02:00
mathiasandClaude Opus 4.8 d7a842f356 feat(capture): I1 sovereignty gate + server-derived origin (#53)
Adds the trust-zone Origin to CaptureContext and the I1 gate to the
use-case: a confidential effective classification through a us-nexus
origin is refused before ANY write (ErrSovereigntyRefused), and the
refusal is itself audited. A caller-asserted harness label that names a
different zone than the server-derived origin is logged as a security
event — context.Harness is descriptive-only, never a gate input.

The gate triggers only on an explicit ZoneUSNexus, so the unset default
(ZoneUnknown) can never make it fire on caller-controllable input; the
REST adapter always sets a concrete zone.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-22 23:42:18 +02:00
mathias aad90f2dfe Merge pull request 'feat: Gitea IssueTracker client + inject into brain server (#52, capture 49c)' (#58) from feat/capture-gitea-tracker into main
CI / Lint / Test / Vet (push) Successful in 12s
CI / Mirror to GitHub (push) Successful in 4s
2026-06-22 21:32:02 +00:00
mathias 07fca9ee73 Merge pull request 'feat: CaptureService use-case + BrainStore extraction (#51, capture 49b)' (#57) from feat/capture-service into main
CI / Lint / Test / Vet (push) Successful in 12s
CI / Mirror to GitHub (push) Has been cancelled
2026-06-22 21:31:47 +00:00
mathias f6bf9b5f57 Merge pull request 'feat: capture classification taxonomy + per-wing/repo tags (#50, capture 49a)' (#56) from feat/capture-classification into main
CI / Lint / Test / Vet (push) Has been cancelled
CI / Mirror to GitHub (push) Has been cancelled
2026-06-22 21:31:18 +00:00
24 changed files with 2290 additions and 39 deletions
+81
View File
@@ -8,6 +8,7 @@ import (
"net/http"
"net/url"
"os"
"path/filepath"
"strconv"
"strings"
"time"
@@ -15,6 +16,10 @@ import (
chassisauth "gitea.d-ma.be/mathias/mcp-chassis/auth"
"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/embed"
"github.com/mathiasbq/hyperguild/ingestion/internal/gitea"
@@ -119,6 +124,47 @@ func envInt(key string, fallback int) int {
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
// caller never has to handle the rare error path.
func systemHostname() string {
@@ -349,6 +395,41 @@ func main() {
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)
// The Gitea client also satisfies SummaryWriter (#66): session
// summaries are written to mathias/ai-sessions over the same API
// token. nil only if a future tracker impl lacks file writes.
summaryWriter, _ := tracker.(capture.SummaryWriter)
captureSvc := capture.NewService(
mcpSrv.BrainStore(), tracker, summaryWriter, 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
// integration UI, which has no static-Bearer field. Setting both
// OAUTH_CLIENT_ID and OAUTH_CLIENT_SECRET enables the token exchange;
+167
View File
@@ -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]
}
+96
View File
@@ -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)
}
}
+191
View File
@@ -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)
}
+99
View File
@@ -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"))
}
+55
View File
@@ -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
}
+65
View File
@@ -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)
}
}
}
}()
}
+56
View File
@@ -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
}
+41
View File
@@ -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)
})
}
+40 -3
View File
@@ -15,13 +15,44 @@
// (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 is server-derived from the
// authenticated identity (#53 populates it); it is never caller-asserted.
// Harness is descriptive telemetry only — never a gate input.
// 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
@@ -29,6 +60,7 @@ type CaptureContext struct {
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
@@ -108,4 +140,9 @@ type CaptureReceipt struct {
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"`
}
+26 -4
View File
@@ -96,9 +96,31 @@ type AuditEntry struct {
SecurityEvents []string
}
// AuditSink records the audit entry. The classification-aware
// degradation/refusal policy (confidential fails closed, internal
// degrades) is the caller's concern in #54; this port just records.
// 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 {
Record(ctx context.Context, e AuditEntry) error
Reserve(ctx context.Context, level classification.Level) (AuditOutcome, error)
Record(ctx context.Context, e AuditEntry, outcome AuditOutcome) error
}
+76 -4
View File
@@ -34,6 +34,39 @@ func NewService(b BrainStore, tr IssueTracker, sw SummaryWriter, p Classificatio
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
@@ -57,6 +90,32 @@ func (s *Service) Capture(ctx context.Context, in CaptureInput) (CaptureReceipt,
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(),
@@ -78,6 +137,16 @@ func (s *Service) Capture(ctx context.Context, in CaptureInput) (CaptureReceipt,
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 {
@@ -110,9 +179,9 @@ func (s *Service) Capture(ctx context.Context, in CaptureInput) (CaptureReceipt,
}
}
// I5: emit a request-level audit record of exactly what landed.
// Best-effort here; the classification-aware refusal/degradation
// policy is #54.
// 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,
@@ -122,9 +191,12 @@ func (s *Service) Capture(ctx context.Context, in CaptureInput) (CaptureReceipt,
EffectiveClassification: effective.String(),
Items: landed,
SecurityEvents: securityEvents,
}); err != nil {
}, outcome); err != nil {
receipt.Errors = append(receipt.Errors, ItemError{Item: "audit", Error: err.Error()})
}
if outcome == AuditBuffered {
receipt.AuditBuffered = true
}
return receipt, nil
}
+142 -2
View File
@@ -113,10 +113,19 @@ func (p fakePolicy) Derive(t classification.Target) classification.Level {
type fakeAudit struct {
entries []AuditEntry
err error
err error // Record error
reserveErr error // Reserve error (refuse before write)
reserveMode AuditOutcome
}
func (f *fakeAudit) Record(_ context.Context, e AuditEntry) error {
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
}
@@ -329,3 +338,134 @@ func TestCaptureSummaryPathAndFidelity(t *testing.T) {
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)
}
+232
View File
@@ -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)
}
+42
View File
@@ -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
}
+91 -22
View File
@@ -12,6 +12,7 @@ package gitea
import (
"bytes"
"context"
"encoding/base64"
"encoding/json"
"fmt"
"io"
@@ -93,37 +94,105 @@ func (c *Client) CloseIssue(ctx context.Context, repo string, number int, commen
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))
// WriteFile creates or updates a file in repo at path via the Gitea
// contents API — the SummaryWriter port (#66). It upserts: a GET resolves
// the current blob sha (if any) so an existing file is updated rather than
// rejected (the richer-fidelity-supersedes rule for re-captured sessions).
// Owner is the fixed const, like every other call.
func (c *Client) WriteFile(ctx context.Context, repo, path, content string) error {
cpath := fmt.Sprintf("/api/v1/repos/%s/%s/contents/%s", owner, repo, path)
sha, err := c.fileSHA(ctx, cpath)
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)
payload := map[string]any{
"message": "capture: " + path,
"content": base64.StdEncoding.EncodeToString([]byte(content)),
}
// Gitea contents API: POST creates a new file, PUT updates an existing
// one (PUT requires the current sha). Pick by whether the file exists.
method := http.MethodPost
if sha != "" {
method = http.MethodPut
payload["sha"] = sha
}
status, body, err := c.request(ctx, method, cpath, payload)
if err != nil {
return fmt.Errorf("gitea %s %s: %w", method, path, err)
return err
}
if status < 200 || status >= 300 {
return fmt.Errorf("gitea %s %s: status %d: %s", method, cpath, status, strings.TrimSpace(string(body)))
}
return nil
}
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)))
// fileSHA returns the current blob sha for a contents path, or "" when the
// file does not exist (404). Any other non-2xx is an error.
func (c *Client) fileSHA(ctx context.Context, cpath string) (string, error) {
status, body, err := c.request(ctx, http.MethodGet, cpath, nil)
if err != nil {
return "", err
}
if out != nil && len(respBody) > 0 {
if err := json.Unmarshal(respBody, out); err != nil {
if status == http.StatusNotFound {
return "", nil
}
if status < 200 || status >= 300 {
return "", fmt.Errorf("gitea GET %s: status %d: %s", cpath, status, strings.TrimSpace(string(body)))
}
var meta struct {
SHA string `json:"sha"`
}
if err := json.Unmarshal(body, &meta); err != nil {
return "", fmt.Errorf("gitea GET %s: decode: %w", cpath, err)
}
return meta.SHA, nil
}
// do performs a JSON request against the Gitea API and decodes a 2xx
// 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 {
status, body, err := c.request(ctx, method, path, payload)
if err != nil {
return err
}
if status < 200 || status >= 300 {
return fmt.Errorf("gitea %s %s: status %d: %s", method, path, status, strings.TrimSpace(string(body)))
}
if out != nil && len(body) > 0 {
if err := json.Unmarshal(body, out); err != nil {
return fmt.Errorf("gitea %s %s: decode response: %w", method, path, err)
}
}
return nil
}
// request is the shared HTTP path: marshals an optional JSON payload,
// attaches auth (token only ever in the header), and returns the status +
// body so callers can branch on status (e.g. 404) without it being an
// error. Never logs the token.
func (c *Client) request(ctx context.Context, method, path string, payload any) (int, []byte, error) {
var reader io.Reader
if payload != nil {
reqBody, err := json.Marshal(payload)
if err != nil {
return 0, nil, fmt.Errorf("marshal request: %w", err)
}
reader = bytes.NewReader(reqBody)
}
req, err := http.NewRequestWithContext(ctx, method, c.baseURL+path, reader)
if err != nil {
return 0, nil, err
}
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Accept", "application/json")
req.Header.Set("Authorization", "token "+c.token)
resp, err := c.http.Do(req)
if err != nil {
return 0, nil, fmt.Errorf("gitea %s %s: %w", method, path, err)
}
defer func() { _ = resp.Body.Close() }()
body, _ := io.ReadAll(io.LimitReader(resp.Body, 8192))
return resp.StatusCode, body, nil
}
+66
View File
@@ -115,3 +115,69 @@ func TestErrorPathDoesNotLeakToken(t *testing.T) {
assert.NotContains(t, err.Error(), testToken, "token must never appear in an error message")
assert.Contains(t, err.Error(), "500")
}
func TestWriteFileCreatesNewFile(t *testing.T) {
var getPath, postPath, postBody string
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
switch r.Method {
case http.MethodGet:
getPath = r.URL.Path
w.WriteHeader(http.StatusNotFound) // file does not exist yet
case http.MethodPost: // gitea contents API: POST = create
postPath = r.URL.Path
b, _ := io.ReadAll(r.Body)
postBody = string(b)
w.WriteHeader(http.StatusCreated)
_ = json.NewEncoder(w).Encode(map[string]any{"content": map[string]any{"html_url": "https://git/x"}})
default:
t.Errorf("create must POST, got %s", r.Method)
}
}))
defer srv.Close()
err := gitea.New(srv.URL, testToken).WriteFile(context.Background(),
"ai-sessions", "summaries/claude-code/2026-06/2026-06-23-x-abcd1234.md", "# Summary\n\nbody\n")
require.NoError(t, err)
assert.Equal(t, "/api/v1/repos/mathias/ai-sessions/contents/summaries/claude-code/2026-06/2026-06-23-x-abcd1234.md", getPath)
assert.Equal(t, getPath, postPath)
// base64 of the content, no sha on create.
assert.Contains(t, postBody, "IyBTdW1tYXJ5") // base64("# Summary")
assert.NotContains(t, postBody, `"sha"`)
}
func TestWriteFileUpdatesExisting(t *testing.T) {
var putBody string
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
switch r.Method {
case http.MethodGet:
w.WriteHeader(http.StatusOK)
_ = json.NewEncoder(w).Encode(map[string]any{"sha": "deadbeef"})
case http.MethodPut:
b, _ := io.ReadAll(r.Body)
putBody = string(b)
w.WriteHeader(http.StatusOK)
_ = json.NewEncoder(w).Encode(map[string]any{"content": map[string]any{"html_url": "https://git/x"}})
}
}))
defer srv.Close()
err := gitea.New(srv.URL, testToken).WriteFile(context.Background(), "ai-sessions", "p/x.md", "new")
require.NoError(t, err)
assert.Contains(t, putBody, `"sha":"deadbeef"`, "existing file → update with sha")
}
func TestWriteFileErrorNoTokenLeak(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.Method == http.MethodGet {
w.WriteHeader(http.StatusNotFound)
return
}
w.WriteHeader(http.StatusUnprocessableEntity)
_, _ = w.Write([]byte("bad"))
}))
defer srv.Close()
err := gitea.New(srv.URL, testToken).WriteFile(context.Background(), "ai-sessions", "p/x.md", "x")
require.Error(t, err)
assert.NotContains(t, err.Error(), testToken)
assert.Contains(t, err.Error(), "422")
}
+8 -1
View File
@@ -38,7 +38,7 @@ func (s *Server) tools() []map[string]any {
return b
}
return []map[string]any{
tools := []map[string]any{
{
"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).",
@@ -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 {
+60 -2
View File
@@ -1,7 +1,8 @@
// Package mcp implements an MCP HTTP handler for the ingestion service.
// Exposed tools: brain_query, brain_write, brain_update, brain_get,
// 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
import (
@@ -12,6 +13,7 @@ import (
"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/graphsync"
"github.com/mathiasbq/hyperguild/ingestion/internal/pipeline"
@@ -50,6 +52,19 @@ type Server struct {
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
@@ -116,6 +131,36 @@ 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
}
func (s *Server) ServeHTTP(w http.ResponseWriter, r *http.Request) {
// MCP streamable HTTP: GET establishes the SSE stream for server-to-client events.
if r.Method == http.MethodGet {
@@ -165,7 +210,18 @@ func (s *Server) ServeHTTP(w http.ResponseWriter, r *http.Request) {
rpcErr = &rpcError{Code: -32602, Message: "invalid params"}
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 {
rpcErr = &rpcError{Code: -32000, Message: err.Error()}
break
@@ -207,6 +263,8 @@ func (s *Server) handleCall(ctx context.Context, name string, args json.RawMessa
return s.brainUpdate(ctx, args)
case "brain_get":
return s.brainGet(ctx, args)
case "capture":
return s.brainCapture(ctx, args)
case "brain_index":
return s.brainIndex(ctx, args)
case "brain_tunnel":
+105
View File
@@ -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")
}
+116
View File
@@ -0,0 +1,116 @@
# Capture capability — implementation report (as-built)
**Status:** Shipped 2026-06-23, tagged `v0.11.0`. Epic hyperguild #49 (sub-issues #50#55) closed.
**Spec:** `specs/capture-bdd-spec.md` (the design contract this implements).
**Governed by:** `infra/docs/architecture/01-invariants.md` (I1I5) + the I2 acceptance ledger entry in `infra/docs/security-baseline.md`.
This document records what was actually built, where it lives, how it maps to the spec, and what was deferred — for onboarding and future audit. It does not restate the design rationale (see the spec and the linked brain entries).
---
## 1. Outcome
One uniform capture capability — insights → brain, action items → Gitea tickets, optional summary → ai-sessions — reachable identically from every harness:
- **In-process / direct-REST harnesses** (Claude Code CLI, Agentsquad, claude.ai Code, headless): `POST /capture` on the brain server.
- **MCP-native harnesses** (claude.ai Chat/Cowork/Design, Crush, Pi, LLM Council): the `capture` MCP tool, reached over the existing `/mcp` OAuth connector.
Both doors call the **same** `CaptureService`; only the transport and credential assembly differ. The persistence behaviour (validation, classification, I1 gate, orchestration, I5 audit, partial receipt) is written once.
---
## 2. Architecture (as-built)
```
POST /capture (REST) capture MCP tool
capturehttp.Handler mcp.Server.brainCapture
\ /
\ (auth → principal → /
\ origin; decode) /
v v
capture.CaptureService (use-case, pure)
┌───────────────┬───────────────┬──────────────┬───────────────┐
BrainStore IssueTracker SummaryWriter ClassificationPolicy AuditSink
brainstore. gitea.Client (nil today) classification.Config audit.Degrading
Store (REST) Sink / SlogSink
│ │
api.WriteNote/UpdateNote/ReadNote (#45) LokiCentral + FileBuffer
+ wing index + auto-tunnel + graph re-index + NtfyNotifier + Reconcile
```
- **`internal/capture/`** — the use-case + ports + entities. Pure; no I/O. Owns validation (fail-closed), effective-classification resolution (stricter wins), the **I1 sovereignty gate**, best-effort orchestration, the **two-phase I5 audit** (Reserve before writes / Record after), and the partial-aware receipt.
- **`internal/brainstore/`** — concrete `BrainStore` wrapping the #45 `api` primitives + wiki upkeep (wing `_index`, auto-tunnel, graph re-index). The MCP `brain_write`/`brain_update`/`brain_get` handlers were re-pointed at it: one implementation, not two.
- **`internal/classification/`** — `public < internal < confidential` taxonomy + per-wing/repo tags from an optional `classification.yaml`; fail-safe to confidential.
- **`internal/gitea/`** — `IssueTracker` over the Gitea REST API; owner forced to `mathias`; token only in the Authorization header.
- **`internal/capturehttp/`** — the REST adapter + the shared `Authenticate` / `DecodeRequest` / `OriginResolver` (also used by the MCP tool).
- **`internal/audit/`** — `SlogSink` (default) and the `DegradingSink` (loki + durable `FileBuffer` + `NtfyNotifier` + `Reconcile`).
- **`internal/mcp/`** — the `capture` relay tool + principal threading (re-derives the caller's principal from the Bearer header the chassis middleware discards).
---
## 3. Sub-issue → PR map
| Sub | Issue | PR(s) | Delivered |
|-----|-------|-------|-----------|
| 49a | #50 | #56 | classification taxonomy + per-wing/repo tags (fail-safe to confidential) |
| 49b | #51 | #57 | `CaptureService` use-case + ports + entities; `BrainStore` extraction (MCP re-pointed) |
| 49c | #52 | #58 | Gitea `IssueTracker` (owner forced mathias; token never logged) |
| 49d | #53 | #59 | `POST /capture` REST + OAuth2 + I1 sovereignty gate (server-derived origin) |
| 49e | #54 | #60 | I5 audit path + classification-aware degradation (loki + buffer + reconcile) |
| 49f | #55 | #61, infra #151 (ledger), #152 (deploy) | MCP `capture` relay tool + I2 ledger + I3 deploy |
Predecessor: #45 (`brain_update`/`brain_get` verbs, PR #46) — the read-after-write contract capture reuses.
---
## 4. Invariant compliance
| Inv | How satisfied |
|-----|---------------|
| **I1** sovereign containment | Effective classification = stricter(caller-declared, target-derived #50). Origin is **server-derived from the authenticated principal**, never `context.harness`. Confidential + us-nexus origin → refused before any write; refusal audited. Asserted-vs-derived mismatch → security event. |
| **I2** deliberate acceptance | The relay's cross-harness reach is recorded in `infra/docs/security-baseline.md` with six containment properties + Revisit-if, **merged before relay code shipped** (infra #151). |
| **I3** GitOps reconcilability | Env + `gitea-api-token` ExternalSecret under `infra/k3s/apps/supervisor/`, Flux-reconciled; image bumped by CD. No untracked runtime. |
| **I5** auditability | Every capture emits a request-level audit record. Classification-aware degradation: confidential + sink-down → hard-refuse; internal/public + sink-down → durable local buffer + ntfy + reconcile-on-recovery; floor → refuse if nothing can record. |
---
## 5. Operational reference (env)
Set on the `ingestion` deployment (`infra/k3s/apps/supervisor/ingestion-deployment.yaml`):
| Env | Purpose | Notes |
|-----|---------|-------|
| `BRAIN_GITEA_URL` / `BRAIN_GITEA_TOKEN` | enables the `IssueTracker` → gates `/capture` + the MCP tool | token from 1P `DMABE_GITEA_API_TOKEN` via ESO; unset ⇒ capture disabled |
| `BRAIN_LOKI_URL` | activates the `DegradingSink` | unset ⇒ `SlogSink` (audit to stdout → alloy → loki; no refuse/buffer semantics) |
| `BRAIN_NTFY_URL` / `BRAIN_NTFY_TOKEN` | degraded-state alerts | optional |
| `BRAIN_CAPTURE_SOVEREIGN_PRINCIPALS` | JWT subjects treated as sovereign-soil | comma-separated; static-token caller is always sovereign; unknown JWT ⇒ us-nexus (fail safe) |
| `BRAIN_AUDIT_RECONCILE_INTERVAL` | buffer→loki replay tick | default 60s |
The audit buffer lives at `<brain>/.audit-buffer/capture.jsonl` on the brain hostPath (nodeSelector-pinned to koala) — durable across restart without a separate PV.
---
## 6. Tests
67 test functions across the six packages. Coverage maps to the spec's Gherkin: happy path, supersede-not-duplicate, fail-closed validation, partial-failure receipt, dry-run, stricter-classification-wins, I1 confidential-via-us-nexus-refused / via-sovereign-allowed / asserted-label-ignored / caller-cannot-forge-origin, I5 confidential-refuse / internal-buffer / floor-refuse / reconcile / buffer-survives-restart, and the MCP relay tool (forwards, preserves principal, unauth rejected). `task check` green.
---
## 7. Deferred (not in this epic)
Tracked here so they aren't lost; file as issues when picked up:
- **SKILL veneer** — the `close-session` SKILL becomes the claude.ai trigger/harvest layer that calls capture.
- **Per-harness token provisioning** for Crush / Pi / LLM Council (claude.ai is done via the existing `/mcp` connector).
- **Harvest adapters** — transcript-parse vs chat-memory-reconstruct vs agent-runlog, each assembling capture args at its own fidelity.
- **`SummaryWriter` impl** — ai-sessions summary persistence (the port + path logic exist; the concrete writer is nil today, so a request with a summary fails that one item).
---
## 8. Brain learnings
- `wiki/hyperguild/decisions/capture-classification-taxonomy`
- `wiki/hyperguild/decisions/gate-on-server-derived-signals-fail-safe`
- `wiki/hyperguild/decisions/two-phase-reserve-record-audit-gate`
- `wiki/hyperguild/failures/mcp-bearer-middleware-discards-principal`
- `wiki/hyperguild/facts/brain-mcp-embeddings-out-of-band-sync` (from #45, the predecessor)