Compare commits

..
Author SHA1 Message Date
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
mathiasandClaude Opus 4.8 6606b38a76 feat(gitea): IssueTracker client + inject into brain server (#52)
Implements the IssueTracker port as a real Gitea REST client (#49c) — the
new outbound dependency the brain server gains for capture.

- CreateIssue / CommentIssue / CloseIssue(+optional closing comment) over
  the Gitea API. Owner is the const "mathias", never caller-supplied, so
  a caller cannot redirect a write to another owner's repo.
- Token read once at construction (BRAIN_GITEA_TOKEN), held in the struct,
  travels only in the Authorization header — never logged or in argv.
  Error messages carry status + truncated body, never the token
  (regression-tested). gitea.New returns nil when URL or token is unset,
  so missing config = tracker disabled via one nil check.
- Injected into the MCP server behind the capture.IssueTracker interface
  via WithIssueTracker (constructor injection, swappable/testable); main
  wires it from BRAIN_GITEA_URL (default https://git.d-ma.be) +
  BRAIN_GITEA_TOKEN. Consumed by the capture use-case in #53.

Tests use httptest transports: create (owner+auth header asserted),
comment, close with/without comment, error path that proves the token
never leaks into an error string.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-22 23:27:29 +02:00
mathiasandClaude Opus 4.8 4cfc98de56 refactor(capture): CloseIssue carries a closing comment (#52)
#52's IssueTracker spec is CloseIssue(repo, number, comment). Refine the
#51 port signature to match and have the service pass the ticket body as
the closing comment (empty ⇒ close only). Keeps the close-with-comment
flow first-class rather than forcing two separate ticket items.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-22 23:27:29 +02:00
21 changed files with 1973 additions and 17 deletions
+82
View File
@@ -8,6 +8,7 @@ import (
"net/http" "net/http"
"net/url" "net/url"
"os" "os"
"path/filepath"
"strconv" "strconv"
"strings" "strings"
"time" "time"
@@ -15,8 +16,13 @@ import (
chassisauth "gitea.d-ma.be/mathias/mcp-chassis/auth" chassisauth "gitea.d-ma.be/mathias/mcp-chassis/auth"
"github.com/mathiasbq/hyperguild/ingestion/internal/api" "github.com/mathiasbq/hyperguild/ingestion/internal/api"
"github.com/mathiasbq/hyperguild/ingestion/internal/audit"
"github.com/mathiasbq/hyperguild/ingestion/internal/capture"
"github.com/mathiasbq/hyperguild/ingestion/internal/capturehttp"
"github.com/mathiasbq/hyperguild/ingestion/internal/classification"
"github.com/mathiasbq/hyperguild/ingestion/internal/claudewatcher" "github.com/mathiasbq/hyperguild/ingestion/internal/claudewatcher"
"github.com/mathiasbq/hyperguild/ingestion/internal/embed" "github.com/mathiasbq/hyperguild/ingestion/internal/embed"
"github.com/mathiasbq/hyperguild/ingestion/internal/gitea"
"github.com/mathiasbq/hyperguild/ingestion/internal/graphstore" "github.com/mathiasbq/hyperguild/ingestion/internal/graphstore"
"github.com/mathiasbq/hyperguild/ingestion/internal/graphsync" "github.com/mathiasbq/hyperguild/ingestion/internal/graphsync"
"github.com/mathiasbq/hyperguild/ingestion/internal/llm" "github.com/mathiasbq/hyperguild/ingestion/internal/llm"
@@ -118,6 +124,47 @@ func envInt(key string, fallback int) int {
return fallback return fallback
} }
// buildAuditSink selects the capture audit sink. When BRAIN_LOKI_URL is
// set it builds the classification-aware DegradingSink (loki central +
// durable file buffer + optional ntfy) and starts the reconcile loop;
// otherwise it falls back to a plain slog sink. The buffer lives under the
// brain dir so it survives process restarts.
func buildAuditSink(ctx context.Context, brainDir string, logger *slog.Logger) capture.AuditSink {
lokiURL := os.Getenv("BRAIN_LOKI_URL")
central := audit.NewLokiCentral(lokiURL)
if central == nil {
logger.Info("capture audit: slog sink (BRAIN_LOKI_URL unset)")
return audit.NewSlogSink(logger)
}
buffer, err := audit.NewFileBuffer(filepath.Join(brainDir, ".audit-buffer", "capture.jsonl"))
if err != nil {
logger.Error("capture audit buffer init", "err", err)
os.Exit(1)
}
// Keep notifier as a nil interface (not a typed-nil) when unconfigured
// so DegradingSink/Reconcile skip it cleanly.
var notifier audit.Notifier
if n := audit.NewNtfyNotifier(os.Getenv("BRAIN_NTFY_URL"), os.Getenv("BRAIN_NTFY_TOKEN")); n != nil {
notifier = n
}
reconcileInterval := time.Duration(envInt("BRAIN_AUDIT_RECONCILE_INTERVAL", 60)) * time.Second
audit.StartReconcile(ctx, central, buffer, notifier, reconcileInterval)
logger.Info("capture audit: loki+buffer sink", "loki", lokiURL, "reconcile_s", int(reconcileInterval.Seconds()))
return audit.NewDegradingSink(central, buffer, notifier)
}
// splitList parses a comma-separated env value into a trimmed,
// empty-free slice. Used for the capture sovereign-principal allowlist.
func splitList(v string) []string {
var out []string
for _, p := range strings.Split(v, ",") {
if p = strings.TrimSpace(p); p != "" {
out = append(out, p)
}
}
return out
}
// systemHostname returns os.Hostname() with a "unknown" fallback so the // systemHostname returns os.Hostname() with a "unknown" fallback so the
// caller never has to handle the rare error path. // caller never has to handle the rare error path.
func systemHostname() string { func systemHostname() string {
@@ -175,6 +222,15 @@ func main() {
logger.Info("brain reranker configured", "url", rerankURL, "model", rerankModel) logger.Info("brain reranker configured", "url", rerankURL, "model", rerankModel)
} }
// Gitea ticket tracker for the capture capability (#52). Token via env
// only — never logged or in argv. Both vars must be set to enable it;
// gitea.New returns nil otherwise, leaving ticket integration off.
giteaURL := envOr("BRAIN_GITEA_URL", "https://git.d-ma.be")
if tracker := gitea.New(giteaURL, os.Getenv("BRAIN_GITEA_TOKEN")); tracker != nil {
mcpSrv = mcpSrv.WithIssueTracker(tracker)
logger.Info("brain gitea tracker configured", "url", giteaURL)
}
// Hybrid retrieval (pgvector + nomic-embed-text). Both env vars must // Hybrid retrieval (pgvector + nomic-embed-text). Both env vars must
// be set together for the path to wire on; otherwise BM25-only. // be set together for the path to wire on; otherwise BM25-only.
var vectorStore *vectorstore.PGStore var vectorStore *vectorstore.PGStore
@@ -339,6 +395,32 @@ func main() {
mux.Handle("/mcp", chassisauth.BearerMiddleware(mcpToken, jwtValidator, "brain", resourceMetadataURL, mcpSrv)) mux.Handle("/mcp", chassisauth.BearerMiddleware(mcpToken, jwtValidator, "brain", resourceMetadataURL, mcpSrv))
// POST /capture (#53/#54): the uniform capture REST door. Needs a ticket
// tracker to file action items, so it only mounts when Gitea is
// configured. It reuses the MCP server's graph-wired brain store (one
// implementation), the classification tags for the I1 gate, and a
// classification-aware audit sink (loki + durable buffer + ntfy when
// BRAIN_LOKI_URL is set, else a plain slog sink). The handler does its
// own auth (static + JWT) because it needs the principal to derive the
// trust-zone origin — the chassis middleware hides it.
if tracker := mcpSrv.IssueTracker(); tracker != nil {
classCfg, cerr := classification.Load(brainDir)
if cerr != nil {
logger.Error("load classification config", "err", cerr)
os.Exit(1)
}
auditSink := buildAuditSink(ctx, brainDir, logger)
captureSvc := capture.NewService(
mcpSrv.BrainStore(), tracker, nil, classCfg, auditSink)
sovereign := splitList(os.Getenv("BRAIN_CAPTURE_SOVEREIGN_PRINCIPALS"))
captureH := capturehttp.New(captureSvc, jwtValidator, mcpToken, "local-cli",
capturehttp.NewOriginResolver(sovereign))
mux.Handle("POST /capture", captureH)
logger.Info("capture endpoint enabled", "sovereign_principals", len(sovereign))
} else {
logger.Info("capture endpoint disabled (BRAIN_GITEA_TOKEN unset)")
}
// Opt-in OAuth 2.0 client_credentials flow for claude.ai's custom-MCP // Opt-in OAuth 2.0 client_credentials flow for claude.ai's custom-MCP
// integration UI, which has no static-Bearer field. Setting both // integration UI, which has no static-Bearer field. Setting both
// OAUTH_CLIENT_ID and OAUTH_CLIENT_SECRET enables the token exchange; // OAUTH_CLIENT_ID and OAUTH_CLIENT_SECRET enables the token exchange;
+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. // (stricter wins), best-effort orchestration, and the partial receipt.
package capture 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. // CaptureContext is the per-session metadata accompanying a capture.
// //
// Classification is the caller-declared sensitivity (model C, spec §4.1): // Classification is the caller-declared sensitivity (model C, spec §4.1):
// the server independently derives the target's classification and gates // the server independently derives the target's classification and gates
// on the stricter of the two. Principal is server-derived from the // on the stricter of the two. Principal and Origin are server-derived from
// authenticated identity (#53 populates it); it is never caller-asserted. // the authenticated identity (the REST adapter populates them); they are
// Harness is descriptive telemetry only — never a gate input. // never caller-asserted. Harness is descriptive telemetry only — never a
// gate input.
type CaptureContext struct { type CaptureContext struct {
Harness string Harness string
SessionRef string SessionRef string
@@ -29,6 +60,7 @@ type CaptureContext struct {
Actor string Actor string
Classification string // caller-declared level token ("" = unspecified) Classification string // caller-declared level token ("" = unspecified)
Principal string // server-derived (auth); audit identity 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 // Insight is one piece of session knowledge bound for the brain. A
@@ -108,4 +140,9 @@ type CaptureReceipt struct {
Errors []ItemError `json:"errors"` Errors []ItemError `json:"errors"`
EffectiveClassification string `json:"effective_classification,omitempty"` EffectiveClassification string `json:"effective_classification,omitempty"`
DryRun bool `json:"dry_run"` 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"`
} }
+29 -5
View File
@@ -63,7 +63,9 @@ type IssueRef struct {
// scopes to owner "mathias"; the port deliberately omits owner. // scopes to owner "mathias"; the port deliberately omits owner.
type IssueTracker interface { type IssueTracker interface {
CreateIssue(ctx context.Context, repo, title, body string) (IssueRef, error) CreateIssue(ctx context.Context, repo, title, body string) (IssueRef, error)
CloseIssue(ctx context.Context, repo string, number int) (IssueRef, error) // CloseIssue closes an issue, optionally posting a closing comment
// first (empty comment ⇒ close only).
CloseIssue(ctx context.Context, repo string, number int, comment string) (IssueRef, error)
CommentIssue(ctx context.Context, repo string, number int, body string) (IssueRef, error) CommentIssue(ctx context.Context, repo string, number int, body string) (IssueRef, error)
} }
@@ -94,9 +96,31 @@ type AuditEntry struct {
SecurityEvents []string SecurityEvents []string
} }
// AuditSink records the audit entry. The classification-aware // AuditOutcome is how a capture's audit record was (or will be) persisted.
// degradation/refusal policy (confidential fails closed, internal type AuditOutcome int
// degrades) is the caller's concern in #54; this port just records.
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 { 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
} }
+77 -5
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} 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 // Capture runs the use-case: validate (fail-closed), resolve effective
// classification (stricter of declared vs target-derived), then persist // classification (stricter of declared vs target-derived), then persist
// insights → tickets → summary best-effort, emit an audit record, and // 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) 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{ receipt := CaptureReceipt{
Errors: []ItemError{}, Errors: []ItemError{},
EffectiveClassification: effective.String(), EffectiveClassification: effective.String(),
@@ -78,6 +137,16 @@ func (s *Service) Capture(ctx context.Context, in CaptureInput) (CaptureReceipt,
return receipt, nil 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 var landed []string
for i, ins := range in.Insights { 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. // I5: persist the request-level audit record of exactly what landed,
// Best-effort here; the classification-aware refusal/degradation // using the outcome reserved before the writes. AuditBuffered surfaces
// policy is #54. // the degraded (locally-buffered) state on the receipt.
if err := s.audit.Record(ctx, AuditEntry{ if err := s.audit.Record(ctx, AuditEntry{
Timestamp: s.now().UTC(), Timestamp: s.now().UTC(),
Principal: in.Context.Principal, Principal: in.Context.Principal,
@@ -122,9 +191,12 @@ func (s *Service) Capture(ctx context.Context, in CaptureInput) (CaptureReceipt,
EffectiveClassification: effective.String(), EffectiveClassification: effective.String(),
Items: landed, Items: landed,
SecurityEvents: securityEvents, SecurityEvents: securityEvents,
}); err != nil { }, outcome); err != nil {
receipt.Errors = append(receipt.Errors, ItemError{Item: "audit", Error: err.Error()}) receipt.Errors = append(receipt.Errors, ItemError{Item: "audit", Error: err.Error()})
} }
if outcome == AuditBuffered {
receipt.AuditBuffered = true
}
return receipt, nil return receipt, nil
} }
@@ -222,7 +294,7 @@ func (s *Service) persistTicket(ctx context.Context, tk Ticket) (TicketResult, e
case "create": case "create":
ref, err = s.issues.CreateIssue(ctx, tk.Repo, tk.Title, tk.Body) ref, err = s.issues.CreateIssue(ctx, tk.Repo, tk.Title, tk.Body)
case "close": case "close":
ref, err = s.issues.CloseIssue(ctx, tk.Repo, tk.Number) ref, err = s.issues.CloseIssue(ctx, tk.Repo, tk.Number, tk.Body)
case "comment": case "comment":
ref, err = s.issues.CommentIssue(ctx, tk.Repo, tk.Number, tk.Body) ref, err = s.issues.CommentIssue(ctx, tk.Repo, tk.Number, tk.Body)
} }
+143 -3
View File
@@ -69,7 +69,7 @@ func (f *fakeTracker) CreateIssue(_ context.Context, repo, title, _ string) (Iss
return IssueRef{Repo: repo, Number: 100 + len(f.created), URL: "https://git/" + repo + "/issues/x"}, nil return IssueRef{Repo: repo, Number: 100 + len(f.created), URL: "https://git/" + repo + "/issues/x"}, nil
} }
func (f *fakeTracker) CloseIssue(_ context.Context, repo string, number int) (IssueRef, error) { func (f *fakeTracker) CloseIssue(_ context.Context, repo string, number int, _ string) (IssueRef, error) {
if f.err != nil { if f.err != nil {
return IssueRef{}, f.err return IssueRef{}, f.err
} }
@@ -113,10 +113,19 @@ func (p fakePolicy) Derive(t classification.Target) classification.Level {
type fakeAudit struct { type fakeAudit struct {
entries []AuditEntry 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 { if f.err != nil {
return f.err return f.err
} }
@@ -329,3 +338,134 @@ func TestCaptureSummaryPathAndFidelity(t *testing.T) {
assert.Contains(t, sw.paths[0], "session-wrap") assert.Contains(t, sw.paths[0], "session-wrap")
assert.Contains(t, sw.content[0], "fidelity: transcript-parse", "fidelity stamped in frontmatter") 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)
}
+215
View File
@@ -0,0 +1,215 @@
// 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"
"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 := h.authenticate(r)
if !ok {
http.Error(w, "unauthorized", http.StatusUnauthorized)
return
}
var req request
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "invalid JSON"})
return
}
in := req.toInput()
// 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 wins, then
// JWT) but returns the resolved principal and whether the static path was
// taken — the chassis middleware hides both, and capture needs them to
// derive the origin.
func (h *Handler) authenticate(r *http.Request) (principal string, viaStatic, ok bool) {
raw, found := strings.CutPrefix(r.Header.Get("Authorization"), "Bearer ")
if !found || raw == "" {
return "", false, false
}
if h.staticToken != "" && subtle.ConstantTimeCompare([]byte(raw), []byte(h.staticToken)) == 1 {
return h.staticPrincipal, true, true
}
if h.validator != nil {
if sub, err := h.validator.Validate(r.Context(), raw); err == nil && sub != "" {
return sub, false, true
}
}
return "", false, false
}
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
}
+129
View File
@@ -0,0 +1,129 @@
// Package gitea implements capture.IssueTracker against a Gitea instance
// over its REST API. It is the new outbound dependency the brain server
// gains for the capture capability (#49c/#52): the server otherwise does
// brain-local file ops only.
//
// Owner is hard-coded to the operator and never taken from caller input.
// The API token is read once at construction, held in the struct, and
// never logged or placed in argv — it travels only in the Authorization
// header of outbound requests (AGENTS.md secret-handling).
package gitea
import (
"bytes"
"context"
"encoding/json"
"fmt"
"io"
"net/http"
"strings"
"time"
"github.com/mathiasbq/hyperguild/ingestion/internal/capture"
)
// owner is the fixed repository owner for every ticket operation. It is a
// constant, not a parameter, so a caller can never redirect a write to
// another owner's repo.
const owner = "mathias"
// Client is a Gitea REST API IssueTracker.
type Client struct {
baseURL string
token string
http *http.Client
}
// New constructs a Client. It returns nil when either baseURL or token is
// empty, so callers can treat missing config as "tracker disabled" with a
// single nil check (mirrors embed.New).
func New(baseURL, token string) *Client {
if baseURL == "" || token == "" {
return nil
}
return &Client{
baseURL: strings.TrimRight(baseURL, "/"),
token: token,
http: &http.Client{Timeout: 15 * time.Second},
}
}
// issueResponse is the subset of a Gitea issue/comment payload we read.
type issueResponse struct {
Number int `json:"number"`
HTMLURL string `json:"html_url"`
}
// CreateIssue opens a new issue under the fixed owner.
func (c *Client) CreateIssue(ctx context.Context, repo, title, body string) (capture.IssueRef, error) {
var out issueResponse
if err := c.do(ctx, http.MethodPost,
fmt.Sprintf("/api/v1/repos/%s/%s/issues", owner, repo),
map[string]any{"title": title, "body": body}, &out); err != nil {
return capture.IssueRef{}, err
}
return capture.IssueRef{Repo: repo, Number: out.Number, URL: out.HTMLURL}, nil
}
// CommentIssue posts a comment on an existing issue.
func (c *Client) CommentIssue(ctx context.Context, repo string, number int, body string) (capture.IssueRef, error) {
var out issueResponse
if err := c.do(ctx, http.MethodPost,
fmt.Sprintf("/api/v1/repos/%s/%s/issues/%d/comments", owner, repo, number),
map[string]any{"body": body}, &out); err != nil {
return capture.IssueRef{}, err
}
return capture.IssueRef{Repo: repo, Number: number, URL: out.HTMLURL}, nil
}
// CloseIssue closes an issue, first posting a closing comment when one is
// given (empty comment ⇒ close only).
func (c *Client) CloseIssue(ctx context.Context, repo string, number int, comment string) (capture.IssueRef, error) {
if strings.TrimSpace(comment) != "" {
if _, err := c.CommentIssue(ctx, repo, number, comment); err != nil {
return capture.IssueRef{}, err
}
}
var out issueResponse
if err := c.do(ctx, http.MethodPatch,
fmt.Sprintf("/api/v1/repos/%s/%s/issues/%d", owner, repo, number),
map[string]any{"state": "closed"}, &out); err != nil {
return capture.IssueRef{}, err
}
return capture.IssueRef{Repo: repo, Number: number, URL: out.HTMLURL}, nil
}
// do performs a JSON request against the Gitea API and decodes the
// response into out. Errors carry the status and a truncated body for
// diagnosis but never the token.
func (c *Client) do(ctx context.Context, method, path string, payload any, out *issueResponse) error {
reqBody, err := json.Marshal(payload)
if err != nil {
return fmt.Errorf("marshal request: %w", err)
}
req, err := http.NewRequestWithContext(ctx, method, c.baseURL+path, bytes.NewReader(reqBody))
if err != nil {
return err
}
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Accept", "application/json")
// Gitea's token scheme. Held here only; never logged.
req.Header.Set("Authorization", "token "+c.token)
resp, err := c.http.Do(req)
if err != nil {
return fmt.Errorf("gitea %s %s: %w", method, path, err)
}
defer func() { _ = resp.Body.Close() }()
respBody, _ := io.ReadAll(io.LimitReader(resp.Body, 4096))
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
return fmt.Errorf("gitea %s %s: status %d: %s", method, path, resp.StatusCode, strings.TrimSpace(string(respBody)))
}
if out != nil && len(respBody) > 0 {
if err := json.Unmarshal(respBody, out); err != nil {
return fmt.Errorf("gitea %s %s: decode response: %w", method, path, err)
}
}
return nil
}
+117
View File
@@ -0,0 +1,117 @@
package gitea_test
import (
"context"
"encoding/json"
"io"
"net/http"
"net/http/httptest"
"strings"
"testing"
"github.com/mathiasbq/hyperguild/ingestion/internal/gitea"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
const testToken = "super-secret-token-value"
func TestNewNilWhenUnconfigured(t *testing.T) {
assert.Nil(t, gitea.New("", testToken))
assert.Nil(t, gitea.New("https://git.example", ""))
}
func TestCreateIssueForcesOwnerAndAuth(t *testing.T) {
var gotPath, gotAuth, gotBody string
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
gotPath = r.URL.Path
gotAuth = r.Header.Get("Authorization")
b, _ := io.ReadAll(r.Body)
gotBody = string(b)
assert.Equal(t, http.MethodPost, r.Method)
w.WriteHeader(http.StatusCreated)
_ = json.NewEncoder(w).Encode(map[string]any{"number": 42, "html_url": "https://git.d-ma.be/mathias/hyperguild/issues/42"})
}))
defer srv.Close()
c := gitea.New(srv.URL, testToken)
require.NotNil(t, c)
ref, err := c.CreateIssue(context.Background(), "hyperguild", "Do the thing", "details")
require.NoError(t, err)
assert.Equal(t, "/api/v1/repos/mathias/hyperguild/issues", gotPath, "owner forced to mathias")
assert.Equal(t, "token "+testToken, gotAuth)
assert.Contains(t, gotBody, "Do the thing")
assert.Equal(t, "hyperguild", ref.Repo)
assert.Equal(t, 42, ref.Number)
assert.Contains(t, ref.URL, "/issues/42")
}
func TestCommentIssue(t *testing.T) {
var gotPath string
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
gotPath = r.URL.Path
w.WriteHeader(http.StatusCreated)
_ = json.NewEncoder(w).Encode(map[string]any{"html_url": "https://git/c/1"})
}))
defer srv.Close()
ref, err := gitea.New(srv.URL, testToken).CommentIssue(context.Background(), "hyperguild", 7, "a comment")
require.NoError(t, err)
assert.Equal(t, "/api/v1/repos/mathias/hyperguild/issues/7/comments", gotPath)
assert.Equal(t, 7, ref.Number)
}
func TestCloseIssueWithComment(t *testing.T) {
var paths []string
var states []string
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
paths = append(paths, r.Method+" "+r.URL.Path)
if r.Method == http.MethodPatch {
var body map[string]any
b, _ := io.ReadAll(r.Body)
_ = json.Unmarshal(b, &body)
states = append(states, body["state"].(string))
}
w.WriteHeader(http.StatusOK)
_ = json.NewEncoder(w).Encode(map[string]any{"number": 9, "html_url": "https://git/i/9"})
}))
defer srv.Close()
ref, err := gitea.New(srv.URL, testToken).CloseIssue(context.Background(), "hyperguild", 9, "closing because done")
require.NoError(t, err)
assert.Equal(t, 9, ref.Number)
// Comment posted first, then state PATCHed to closed.
assert.Contains(t, paths, "POST /api/v1/repos/mathias/hyperguild/issues/9/comments")
assert.Contains(t, paths, "PATCH /api/v1/repos/mathias/hyperguild/issues/9")
assert.Equal(t, []string{"closed"}, states)
}
func TestCloseIssueNoComment(t *testing.T) {
var commented bool
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if strings.HasSuffix(r.URL.Path, "/comments") {
commented = true
}
w.WriteHeader(http.StatusOK)
_ = json.NewEncoder(w).Encode(map[string]any{"number": 3, "html_url": "https://git/i/3"})
}))
defer srv.Close()
_, err := gitea.New(srv.URL, testToken).CloseIssue(context.Background(), "hyperguild", 3, "")
require.NoError(t, err)
assert.False(t, commented, "empty comment ⇒ no comment POST")
}
func TestErrorPathDoesNotLeakToken(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusInternalServerError)
_, _ = w.Write([]byte("boom"))
}))
defer srv.Close()
_, err := gitea.New(srv.URL, testToken).CreateIssue(context.Background(), "hyperguild", "t", "b")
require.Error(t, err)
assert.NotContains(t, err.Error(), testToken, "token must never appear in an error message")
assert.Contains(t, err.Error(), "500")
}
+23
View File
@@ -11,6 +11,7 @@ import (
"net/http" "net/http"
"github.com/mathiasbq/hyperguild/ingestion/internal/brainstore" "github.com/mathiasbq/hyperguild/ingestion/internal/brainstore"
"github.com/mathiasbq/hyperguild/ingestion/internal/capture"
"github.com/mathiasbq/hyperguild/ingestion/internal/graphstore" "github.com/mathiasbq/hyperguild/ingestion/internal/graphstore"
"github.com/mathiasbq/hyperguild/ingestion/internal/graphsync" "github.com/mathiasbq/hyperguild/ingestion/internal/graphsync"
"github.com/mathiasbq/hyperguild/ingestion/internal/pipeline" "github.com/mathiasbq/hyperguild/ingestion/internal/pipeline"
@@ -48,6 +49,7 @@ type Server struct {
embedder search.Embedder // nil = BM25-only retrieval embedder search.Embedder // nil = BM25-only retrieval
graph graphsync.Store // nil = brain_graph and GraphRAG augmentation disabled graph graphsync.Store // nil = brain_graph and GraphRAG augmentation disabled
store *brainstore.Store // shared brain write/update/get impl (also used by capture) 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)
} }
// NewServer constructs a Server bound to brainDir. pipelineCfg supplies the // NewServer constructs a Server bound to brainDir. pipelineCfg supplies the
@@ -100,6 +102,27 @@ func (s *Server) WithGraph(g *graphstore.PGStore) *Server {
return s return s
} }
// WithIssueTracker injects the Gitea ticket tracker behind the
// capture.IssueTracker interface. nil leaves ticket integration off. The
// use-case (capture) consumes this in #53; it is wired here so the
// dependency is constructed once and stays swappable/testable.
func (s *Server) WithIssueTracker(t capture.IssueTracker) *Server {
s.tracker = t
return s
}
// IssueTracker returns the injected ticket tracker (nil when unconfigured).
func (s *Server) IssueTracker() capture.IssueTracker {
return s.tracker
}
// BrainStore returns the shared brain store (graph-wired once WithGraph
// has run), so the capture use-case writes through the exact same
// implementation as the MCP handlers.
func (s *Server) BrainStore() *brainstore.Store {
return s.store
}
func (s *Server) ServeHTTP(w http.ResponseWriter, r *http.Request) { func (s *Server) ServeHTTP(w http.ResponseWriter, r *http.Request) {
// MCP streamable HTTP: GET establishes the SSE stream for server-to-client events. // MCP streamable HTTP: GET establishes the SSE stream for server-to-client events.
if r.Method == http.MethodGet { if r.Method == http.MethodGet {
+21
View File
@@ -2,12 +2,14 @@ package mcp_test
import ( import (
"bytes" "bytes"
"context"
"encoding/json" "encoding/json"
"net/http" "net/http"
"net/http/httptest" "net/http/httptest"
"strings" "strings"
"testing" "testing"
"github.com/mathiasbq/hyperguild/ingestion/internal/capture"
"github.com/mathiasbq/hyperguild/ingestion/internal/mcp" "github.com/mathiasbq/hyperguild/ingestion/internal/mcp"
"github.com/stretchr/testify/assert" "github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require" "github.com/stretchr/testify/require"
@@ -93,3 +95,22 @@ func TestServerUnknownMethodReturnsError(t *testing.T) {
assert.Equal(t, float64(-32601), errObj["code"]) assert.Equal(t, float64(-32601), errObj["code"])
assert.Contains(t, errObj["message"].(string), "unknown/method") assert.Contains(t, errObj["message"].(string), "unknown/method")
} }
type stubTracker struct{}
func (stubTracker) CreateIssue(context.Context, string, string, string) (capture.IssueRef, error) {
return capture.IssueRef{}, nil
}
func (stubTracker) CloseIssue(context.Context, string, int, string) (capture.IssueRef, error) {
return capture.IssueRef{}, nil
}
func (stubTracker) CommentIssue(context.Context, string, int, string) (capture.IssueRef, error) {
return capture.IssueRef{}, nil
}
func TestWithIssueTrackerInjects(t *testing.T) {
srv := mcp.NewServer(t.TempDir(), nil, nil, nil)
assert.Nil(t, srv.IssueTracker(), "tracker is off by default")
srv = srv.WithIssueTracker(stubTracker{})
assert.NotNil(t, srv.IssueTracker(), "tracker injected behind the interface")
}