Compare commits

...
Author SHA1 Message Date
mathiasandClaude Opus 4.8 9f8fb9c138 docs(close-session): refresh classification gate post-#67 (#62)
CI / Lint / Test / Vet (pull_request) Successful in 13s
CI / Mirror to GitHub (pull_request) Has been skipped
The capture-routing in 723dab5 (#62 minimal) carried pre-#67 prose: it
warned "no populated classification.yaml yet" and steered repos_touched
away from brain/ai-sessions as if they'd escalate to confidential. #67
tagged the homelab repos internal, so that guidance is now wrong and
over-restrictive — listing the central repos is fine, and the summary's
fixed ai-sessions target no longer escalates the call. Point at
classification.yaml as the source of truth; reserve the refusal warning
for client-*/untagged material. Completes #62's veneer.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-26 16:25:57 +02:00
mathias d39a18dd69 Merge pull request 'fix: wire ai-sessions SummaryWriter into the capture relay (#66)' (#68) from fix/capture-summarywriter into main
CI / Lint / Test / Vet (push) Successful in 12s
CI / Mirror to GitHub (push) Successful in 4s
2026-06-23 15:22:42 +00:00
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
mathias 723dab51ae feat(close-session): route closeout through the capture tool (#62 minimal)
CI / Lint / Test / Vet (push) Successful in 13s
CI / Mirror to GitHub (push) Successful in 3s
Collapse Phases 4 (summary commit) + 5 (brain note) into a single
capture-driven Phase 4. The skill now assembles one brain:capture payload
(insights + tickets + summary + context) and dry-runs-then-executes;
capture owns the brain/gitea/ai-sessions writes, the I1 gate, the I5
audit record, and the supersession/read-after-write discipline server-side.

Adds the classification-gate lesson learned dogfooding capture this
session: effective classification = strictest across ALL targets
(wings, ticket repos, repos_touched); untagged repos (brain, ai-sessions)
fail-safe to confidential and get refused via claude.ai, so repos_touched
must stay to internal-default repos. Legacy inline path kept only as the
capture-unreachable fallback. Staleness prose removed (server owns it now).

Partial of #62; harness-token + harvest-adapter + fidelity-supersession
work remain.
2026-06-23 06:27:57 +00: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
24 changed files with 1614 additions and 124 deletions
+49 -8
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,11 +16,11 @@ 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/claudewatcher"
"github.com/mathiasbq/hyperguild/ingestion/internal/audit" "github.com/mathiasbq/hyperguild/ingestion/internal/audit"
"github.com/mathiasbq/hyperguild/ingestion/internal/capture" "github.com/mathiasbq/hyperguild/ingestion/internal/capture"
"github.com/mathiasbq/hyperguild/ingestion/internal/capturehttp" "github.com/mathiasbq/hyperguild/ingestion/internal/capturehttp"
"github.com/mathiasbq/hyperguild/ingestion/internal/classification" "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/embed"
"github.com/mathiasbq/hyperguild/ingestion/internal/gitea" "github.com/mathiasbq/hyperguild/ingestion/internal/gitea"
"github.com/mathiasbq/hyperguild/ingestion/internal/graphstore" "github.com/mathiasbq/hyperguild/ingestion/internal/graphstore"
@@ -123,6 +124,35 @@ 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, // splitList parses a comma-separated env value into a trimmed,
// empty-free slice. Used for the capture sovereign-principal allowlist. // empty-free slice. Used for the capture sovereign-principal allowlist.
func splitList(v string) []string { func splitList(v string) []string {
@@ -365,11 +395,12 @@ 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): the uniform capture REST door. Needs a ticket // POST /capture (#53/#54): the uniform capture REST door. Needs a ticket
// tracker to file action items, so it only mounts when Gitea is // tracker to file action items, so it only mounts when Gitea is
// configured. It reuses the MCP server's graph-wired brain store (one // configured. It reuses the MCP server's graph-wired brain store (one
// implementation), the classification tags for the I1 gate, and a slog // implementation), the classification tags for the I1 gate, and a
// audit sink (the loki+buffer sink lands in #54). The handler does its // 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 // own auth (static + JWT) because it needs the principal to derive the
// trust-zone origin — the chassis middleware hides it. // trust-zone origin — the chassis middleware hides it.
if tracker := mcpSrv.IssueTracker(); tracker != nil { if tracker := mcpSrv.IssueTracker(); tracker != nil {
@@ -378,13 +409,23 @@ func main() {
logger.Error("load classification config", "err", cerr) logger.Error("load classification config", "err", cerr)
os.Exit(1) 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( captureSvc := capture.NewService(
mcpSrv.BrainStore(), tracker, nil, classCfg, audit.NewSlogSink(logger)) mcpSrv.BrainStore(), tracker, summaryWriter, classCfg, auditSink)
sovereign := splitList(os.Getenv("BRAIN_CAPTURE_SOVEREIGN_PRINCIPALS")) sovereign := splitList(os.Getenv("BRAIN_CAPTURE_SOVEREIGN_PRINCIPALS"))
captureH := capturehttp.New(captureSvc, jwtValidator, mcpToken, "local-cli", resolver := capturehttp.NewOriginResolver(sovereign)
capturehttp.NewOriginResolver(sovereign)) captureH := capturehttp.New(captureSvc, jwtValidator, mcpToken, "local-cli", resolver)
mux.Handle("POST /capture", captureH) mux.Handle("POST /capture", captureH)
logger.Info("capture endpoint enabled", "sovereign_principals", len(sovereign)) // 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 { } else {
logger.Info("capture endpoint disabled (BRAIN_GITEA_TOKEN unset)") logger.Info("capture endpoint disabled (BRAIN_GITEA_TOKEN unset)")
} }
+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)
}
}
}
}()
}
+12 -4
View File
@@ -11,11 +11,14 @@ import (
"log/slog" "log/slog"
"github.com/mathiasbq/hyperguild/ingestion/internal/capture" "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, so it // SlogSink records audit entries to an slog.Logger. It never fails and is
// does not exercise the I5 floor (refuse-if-unauditable) — that is #54's // always centrally available, so its Reserve always grants AuditCentral —
// loki+buffer sink. A nil logger falls back to slog.Default(). // 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 { type SlogSink struct {
logger *slog.Logger logger *slog.Logger
} }
@@ -28,10 +31,15 @@ func NewSlogSink(logger *slog.Logger) *SlogSink {
return &SlogSink{logger: logger} 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 // Record emits the audit entry at info level. Security events, when
// present, are logged at warn level so they surface independently of the // present, are logged at warn level so they surface independently of the
// routine audit stream. // routine audit stream.
func (s *SlogSink) Record(_ context.Context, e capture.AuditEntry) error { func (s *SlogSink) Record(_ context.Context, e capture.AuditEntry, _ capture.AuditOutcome) error {
s.logger.Info("capture audit", s.logger.Info("capture audit",
"principal", e.Principal, "principal", e.Principal,
"actor", e.Actor, "actor", e.Actor,
+2 -2
View File
@@ -22,7 +22,7 @@ func TestSlogSinkRecordsEntryAndSecurityEvents(t *testing.T) {
EffectiveClassification: "confidential", EffectiveClassification: "confidential",
Items: []string{"insight:wiki/a/facts/x.md"}, Items: []string{"insight:wiki/a/facts/x.md"},
SecurityEvents: []string{"asserted-vs-derived origin mismatch"}, SecurityEvents: []string{"asserted-vs-derived origin mismatch"},
}) }, capture.AuditCentral)
require.NoError(t, err) require.NoError(t, err)
out := buf.String() out := buf.String()
@@ -36,6 +36,6 @@ func TestSlogSinkRecordsEntryAndSecurityEvents(t *testing.T) {
func TestSlogSinkNilLoggerDefaults(t *testing.T) { func TestSlogSinkNilLoggerDefaults(t *testing.T) {
// nil logger must not panic. // nil logger must not panic.
require.NotPanics(t, func() { require.NotPanics(t, func() {
_ = audit.NewSlogSink(nil).Record(context.Background(), capture.AuditEntry{}) _ = audit.NewSlogSink(nil).Record(context.Background(), capture.AuditEntry{}, capture.AuditCentral)
}) })
} }
+5
View File
@@ -140,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"`
} }
+26 -4
View File
@@ -96,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
} }
+24 -5
View File
@@ -39,6 +39,12 @@ var validActions = map[string]bool{"create": true, "close": true, "comment": tru
// REST adapter maps it to HTTP 403. Callers test with errors.Is. // REST adapter maps it to HTTP 403. Callers test with errors.Is.
var ErrSovereigntyRefused = fmt.Errorf("capture refused by I1 sovereignty gate") 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 // assertedZoneMismatch returns a security-event string when the caller's
// harness label asserts a trust zone that contradicts the server-derived // harness label asserts a trust zone that contradicts the server-derived
// origin. A harness label that names no zone (the normal case, e.g. // origin. A harness label that names no zone (the normal case, e.g.
@@ -105,7 +111,7 @@ func (s *Service) Capture(ctx context.Context, in CaptureInput) (CaptureReceipt,
EffectiveClassification: effective.String(), EffectiveClassification: effective.String(),
Items: nil, // refused before any write Items: nil, // refused before any write
SecurityEvents: append(securityEvents, "I1 refusal: confidential capture via us-nexus origin"), SecurityEvents: append(securityEvents, "I1 refusal: confidential capture via us-nexus origin"),
}) }, AuditCentral)
return CaptureReceipt{}, fmt.Errorf("%w: effective classification confidential through %s origin", return CaptureReceipt{}, fmt.Errorf("%w: effective classification confidential through %s origin",
ErrSovereigntyRefused, in.Context.Origin) ErrSovereigntyRefused, in.Context.Origin)
} }
@@ -131,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 {
@@ -163,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,
@@ -175,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
} }
+64 -3
View File
@@ -112,11 +112,20 @@ 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
} }
@@ -408,3 +417,55 @@ func TestCaptureInternalViaUSNexusAllowed(t *testing.T) {
require.NoError(t, err) require.NoError(t, err)
assert.True(t, rec.Insights[0].OK) 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)
}
+35 -14
View File
@@ -11,6 +11,7 @@ import (
"crypto/subtle" "crypto/subtle"
"encoding/json" "encoding/json"
"errors" "errors"
"io"
"net/http" "net/http"
"strings" "strings"
@@ -90,19 +91,22 @@ type summaryBody struct {
// ServeHTTP authenticates, derives origin, runs the use-case, and maps the // ServeHTTP authenticates, derives origin, runs the use-case, and maps the
// result to an HTTP status. // result to an HTTP status.
func (h *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) { func (h *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
principal, viaStatic, ok := h.authenticate(r) principal, viaStatic, ok := Authenticate(r, h.staticToken, h.staticPrincipal, h.validator)
if !ok { if !ok {
http.Error(w, "unauthorized", http.StatusUnauthorized) http.Error(w, "unauthorized", http.StatusUnauthorized)
return return
} }
var req request body, err := io.ReadAll(r.Body)
if err := json.NewDecoder(r.Body).Decode(&req); err != nil { 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"}) writeJSON(w, http.StatusBadRequest, map[string]string{"error": "invalid JSON"})
return return
} }
in := req.toInput()
// Principal and origin are server-derived — overwrite anything the // Principal and origin are server-derived — overwrite anything the
// caller may have tried to put in the body. // caller may have tried to put in the body.
in.Context.Principal = principal in.Context.Principal = principal
@@ -113,6 +117,10 @@ func (h *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
case errors.Is(err, capture.ErrSovereigntyRefused): case errors.Is(err, capture.ErrSovereigntyRefused):
writeJSON(w, http.StatusForbidden, map[string]string{"error": err.Error()}) writeJSON(w, http.StatusForbidden, map[string]string{"error": err.Error()})
return 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: case err != nil:
// Pre-write validation failure (fail-closed). // Pre-write validation failure (fail-closed).
writeJSON(w, http.StatusBadRequest, map[string]string{"error": err.Error()}) writeJSON(w, http.StatusBadRequest, map[string]string{"error": err.Error()})
@@ -121,26 +129,39 @@ func (h *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
writeJSON(w, statusFor(rec), rec) writeJSON(w, statusFor(rec), rec)
} }
// authenticate mirrors the chassis Bearer precedence (static wins, then // Authenticate mirrors the chassis Bearer precedence (static token wins,
// JWT) but returns the resolved principal and whether the static path was // then Dex JWT) and returns the resolved principal plus whether the static
// taken — the chassis middleware hides both, and capture needs them to // path was taken — the chassis middleware hides both, and capture (REST or
// derive the origin. // MCP) needs them to derive the trust-zone origin. ok is false when no
func (h *Handler) authenticate(r *http.Request) (principal string, viaStatic, ok bool) { // 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 ") raw, found := strings.CutPrefix(r.Header.Get("Authorization"), "Bearer ")
if !found || raw == "" { if !found || raw == "" {
return "", false, false return "", false, false
} }
if h.staticToken != "" && subtle.ConstantTimeCompare([]byte(raw), []byte(h.staticToken)) == 1 { if staticToken != "" && subtle.ConstantTimeCompare([]byte(raw), []byte(staticToken)) == 1 {
return h.staticPrincipal, true, true return staticPrincipal, true, true
} }
if h.validator != nil { if validator != nil {
if sub, err := h.validator.Validate(r.Context(), raw); err == nil && sub != "" { if sub, err := validator.Validate(r.Context(), raw); err == nil && sub != "" {
return sub, false, true return sub, false, true
} }
} }
return "", false, false 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 { func (b request) toInput() capture.CaptureInput {
in := capture.CaptureInput{ in := capture.CaptureInput{
Context: capture.CaptureContext{ Context: capture.CaptureContext{
@@ -176,3 +176,23 @@ func TestCallerCannotForgeOrigin(t *testing.T) {
}) })
assert.Equal(t, http.StatusForbidden, rr.Code) 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)
}
+91 -22
View File
@@ -12,6 +12,7 @@ package gitea
import ( import (
"bytes" "bytes"
"context" "context"
"encoding/base64"
"encoding/json" "encoding/json"
"fmt" "fmt"
"io" "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 return capture.IssueRef{Repo: repo, Number: number, URL: out.HTMLURL}, nil
} }
// do performs a JSON request against the Gitea API and decodes the // WriteFile creates or updates a file in repo at path via the Gitea
// response into out. Errors carry the status and a truncated body for // contents API — the SummaryWriter port (#66). It upserts: a GET resolves
// diagnosis but never the token. // the current blob sha (if any) so an existing file is updated rather than
func (c *Client) do(ctx context.Context, method, path string, payload any, out *issueResponse) error { // rejected (the richer-fidelity-supersedes rule for re-captured sessions).
reqBody, err := json.Marshal(payload) // Owner is the fixed const, like every other call.
if err != nil { func (c *Client) WriteFile(ctx context.Context, repo, path, content string) error {
return fmt.Errorf("marshal request: %w", err) cpath := fmt.Sprintf("/api/v1/repos/%s/%s/contents/%s", owner, repo, path)
} sha, err := c.fileSHA(ctx, cpath)
req, err := http.NewRequestWithContext(ctx, method, c.baseURL+path, bytes.NewReader(reqBody))
if err != nil { if err != nil {
return err return err
} }
req.Header.Set("Content-Type", "application/json") payload := map[string]any{
req.Header.Set("Accept", "application/json") "message": "capture: " + path,
// Gitea's token scheme. Held here only; never logged. "content": base64.StdEncoding.EncodeToString([]byte(content)),
req.Header.Set("Authorization", "token "+c.token) }
// Gitea contents API: POST creates a new file, PUT updates an existing
resp, err := c.http.Do(req) // 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 { if err != nil {
return fmt.Errorf("gitea %s %s: %w", method, path, err) return err
} }
defer func() { _ = resp.Body.Close() }() if status < 200 || status >= 300 {
return fmt.Errorf("gitea %s %s: status %d: %s", method, cpath, status, strings.TrimSpace(string(body)))
}
return nil
}
respBody, _ := io.ReadAll(io.LimitReader(resp.Body, 4096)) // fileSHA returns the current blob sha for a contents path, or "" when the
if resp.StatusCode < 200 || resp.StatusCode >= 300 { // file does not exist (404). Any other non-2xx is an error.
return fmt.Errorf("gitea %s %s: status %d: %s", method, path, resp.StatusCode, strings.TrimSpace(string(respBody))) 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 status == http.StatusNotFound {
if err := json.Unmarshal(respBody, out); err != nil { 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 fmt.Errorf("gitea %s %s: decode response: %w", method, path, err)
} }
} }
return nil 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.NotContains(t, err.Error(), testToken, "token must never appear in an error message")
assert.Contains(t, err.Error(), "500") 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 b
} }
return []map[string]any{ tools := []map[string]any{
{ {
"name": "brain_query", "name": "brain_query",
"description": "BM25 full-text search across brain/knowledge/ and brain/wiki/ markdown files. Optionally scope by wing (topic domain) and hall (memory type).", "description": "BM25 full-text search across brain/knowledge/ and brain/wiki/ markdown files. Optionally scope by wing (topic domain) and hall (memory type).",
@@ -171,6 +171,13 @@ func (s *Server) tools() []map[string]any {
}), }),
}, },
} }
// The capture relay tool (#55) is advertised only when wired via
// WithCapture — MCP-native harnesses (claude.ai, Crush, Pi, LLM Council)
// reach capture through it.
if s.capture != nil {
tools = append(tools, captureToolDescriptor())
}
return tools
} }
type brainQueryArgs struct { type brainQueryArgs struct {
+53 -2
View File
@@ -1,7 +1,8 @@
// Package mcp implements an MCP HTTP handler for the ingestion service. // Package mcp implements an MCP HTTP handler for the ingestion service.
// Exposed tools: brain_query, brain_write, brain_update, brain_get, // Exposed tools: brain_query, brain_write, brain_update, brain_get,
// brain_index, brain_tunnel, brain_ingest, brain_ingest_raw, // brain_index, brain_tunnel, brain_ingest, brain_ingest_raw,
// brain_answer, brain_classify, brain_graph, brain_context, session_log. // brain_answer, brain_classify, brain_graph, brain_context, session_log,
// and capture (the #55 relay tool, registered only when WithCapture is set).
package mcp package mcp
import ( import (
@@ -12,6 +13,7 @@ import (
"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/capture"
"github.com/mathiasbq/hyperguild/ingestion/internal/capturehttp"
"github.com/mathiasbq/hyperguild/ingestion/internal/graphstore" "github.com/mathiasbq/hyperguild/ingestion/internal/graphstore"
"github.com/mathiasbq/hyperguild/ingestion/internal/graphsync" "github.com/mathiasbq/hyperguild/ingestion/internal/graphsync"
"github.com/mathiasbq/hyperguild/ingestion/internal/pipeline" "github.com/mathiasbq/hyperguild/ingestion/internal/pipeline"
@@ -50,6 +52,19 @@ type Server struct {
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) tracker capture.IssueTracker // nil = no Gitea ticket integration; wired for capture (#53)
capture *captureDeps // nil = capture MCP tool disabled (#55 relay)
}
// captureDeps holds what the MCP `capture` tool (the #55 relay door for
// MCP-native harnesses like claude.ai) needs: the use-case, the auth bits
// to re-derive the caller's principal from the Bearer header (the chassis
// middleware gates but discards the principal), and the origin resolver.
type captureDeps struct {
svc *capture.Service
validator capturehttp.Validator
staticToken string
staticPrincipal string
resolver capturehttp.OriginResolver
} }
// NewServer constructs a Server bound to brainDir. pipelineCfg supplies the // NewServer constructs a Server bound to brainDir. pipelineCfg supplies the
@@ -123,6 +138,29 @@ func (s *Server) BrainStore() *brainstore.Store {
return s.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) { 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 {
@@ -172,7 +210,18 @@ func (s *Server) ServeHTTP(w http.ResponseWriter, r *http.Request) {
rpcErr = &rpcError{Code: -32602, Message: "invalid params"} rpcErr = &rpcError{Code: -32602, Message: "invalid params"}
break break
} }
out, err := s.handleCall(r.Context(), p.Name, p.Arguments) // Re-derive the authenticated principal from the Bearer header so
// the capture tool can compute the trust-zone origin. The request
// is already gated by BearerMiddleware; this only recovers the
// identity that middleware discards.
ctx := r.Context()
if s.capture != nil {
if principal, viaStatic, ok := capturehttp.Authenticate(
r, s.capture.staticToken, s.capture.staticPrincipal, s.capture.validator); ok {
ctx = withPrincipal(ctx, principal, viaStatic)
}
}
out, err := s.handleCall(ctx, p.Name, p.Arguments)
if err != nil { if err != nil {
rpcErr = &rpcError{Code: -32000, Message: err.Error()} rpcErr = &rpcError{Code: -32000, Message: err.Error()}
break break
@@ -214,6 +263,8 @@ func (s *Server) handleCall(ctx context.Context, name string, args json.RawMessa
return s.brainUpdate(ctx, args) return s.brainUpdate(ctx, args)
case "brain_get": case "brain_get":
return s.brainGet(ctx, args) return s.brainGet(ctx, args)
case "capture":
return s.brainCapture(ctx, args)
case "brain_index": case "brain_index":
return s.brainIndex(ctx, args) return s.brainIndex(ctx, args)
case "brain_tunnel": case "brain_tunnel":
+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")
}
+29 -59
View File
@@ -7,14 +7,14 @@ description: Disciplined end-of-session closeout for a Claude.ai chat before arc
Capture a finishing Claude.ai work session into durable storage before the chat is archived and its context is lost. The goal is simple and load-bearing: **after this runs, a fresh session (or another agent) can reconstruct what was decided, what was shipped, and what is still open — without the original chat.** Capture a finishing Claude.ai work session into durable storage before the chat is archived and its context is lost. The goal is simple and load-bearing: **after this runs, a fresh session (or another agent) can reconstruct what was decided, what was shipped, and what is still open — without the original chat.**
This skill is **batch**: one session in, findings out, done. It does not loop or re-read its own fresh output semantically (see Phase 5). Run the phases in order. Stop at any confirmation gate that says STOP. This skill is **batch**: one session in, findings out, done. Run the phases in order. Stop at any confirmation gate that says STOP.
## Operating constraints (read first) ## Operating constraints (read first)
- **Gitea owner is always `mathias`.** Never guess another owner. - **Gitea owner is always `mathias`.** Never guess another owner.
- **Ground-truth at HEAD before acting.** Issue bodies and doc references rot — stale hostnames, retired services, moved endpoints. Before closing/commenting on any issue, `gitea:issue_get` it fresh. Before asserting an infra fact, verify it; do not copy it from memory or from a stale issue body. - **Ground-truth at HEAD before acting.** Issue bodies and doc references rot — stale hostnames, retired services, moved endpoints. Before closing/commenting on any issue, `gitea:issue_get` it fresh. Before asserting an infra fact, verify it; do not copy it from memory or from a stale issue body.
- **Current infra truths** (verify rather than trust, but these are the known-good baseline): Gitea is `git.d-ma.be` (not `gitea.d-ma.be`). LiteLLM is `http://koala:30401/v1/` (public `https://llm-api.d-ma.be`); piguard runs NGINX Proxy Manager only — never reference `piguard:4000` or `koala:4000`. Identity provider is Authentik (Dex migration complete). - **Current infra truths** (verify rather than trust, but these are the known-good baseline): Gitea is `git.d-ma.be` (not `gitea.d-ma.be`). LiteLLM is `http://koala:30401/v1/` (public `https://llm-api.d-ma.be`); piguard runs NGINX Proxy Manager only — never reference `piguard:4000` or `koala:4000`. Identity provider is Authentik (Dex migration complete).
- **Side-effects need a confirmation gate.** Closing issues, committing files, and writing to the brain are all real writes. Surface exactly what will happen and get a clear yes before doing it. Reads are free; writes are gated. - **Side-effects need a confirmation gate.** Closing issues and capturing to brain/Gitea/ai-sessions are real writes. Surface exactly what will happen and get a clear yes before doing it. Reads are free; writes are gated.
- **Never fabricate.** If the session didn't produce a decision worth persisting, say so and skip that write. An empty-but-honest closeout beats an invented one. - **Never fabricate.** If the session didn't produce a decision worth persisting, say so and skip that write. An empty-but-honest closeout beats an invented one.
## Phase 1 — Harvest ## Phase 1 — Harvest
@@ -34,76 +34,46 @@ For every repo touched this session, get its true current state before proposing
Do not write anything in this phase. This is the read pass. Do not write anything in this phase. This is the read pass.
## Phase 3 — Confirm and act on issue changes ## Phase 3 — Plan the issue changes
Present a single consolidated plan of issue actions: which to close (with closing comment), which to file (discovered-but-deferred work — token-budget gaps, recorded limitations, v2 follow-ups), which to comment on. Include the exact title/body for any new issue and the closing rationale for any close. Decide the issue actions: which to close (with closing comment), which to file (discovered-but-deferred work — token-budget gaps, recorded limitations, v2 follow-ups), which to comment on. Include the exact title/body for any new issue and the closing rationale for any close.
**GATE — STOP and get explicit confirmation before any issue write.** Issue closes and new issues are side-effects. Once confirmed, execute them (`gitea:issue_close`, `gitea:issue_create`, `gitea:issue_comment`, all owner `mathias`), correcting any rotted references you found in Phase 2 as you go. These actions are **carried into the Phase 4 capture call** as `tickets[]` rather than executed here with direct `gitea:issue_*` calls — routing them through capture puts each one into the I5 audit record. (Closing an issue that needs a separate explanatory comment first is the one case to do directly; otherwise prefer the capture path.)
## Phase 4 — Commit the canonical session summary ## Phase 4 — Capture (one uniform call)
Write one summary file to `mathias/ai-sessions`, committed directly to `main` via `gitea:file_write_branch` (no PR — this repo is solo and unprotected; if branch protection is ever added, fall back to a branch + PR). Persist the session via a **single `capture` call** (the `brain:capture` MCP tool, live on the Claude.ai connector). Capture owns the writes server-side — insights → brain, action items → Gitea tickets, summary → ai-sessions — plus the I1 sovereignty gate, the I5 audit record, and the supersession/read-after-write discipline. The skill's job is to *assemble the payload*, not to write each store itself. Do NOT fall back to separate `gitea:file_write_branch` + `brain_write` steps unless `capture` is unreachable (see fallback below).
**Path:** `summaries/claudeai/<YYYY-MM>/<YYYY-MM-DD>-<topic-slug>-<chatid8>.md` **Assemble one payload:**
where `<chatid8>` is the first 8 chars of the chat's UUID if known, else a short stable slug. `claudeai` has no host segment — Claude.ai is Anthropic-side, not a homelab host.
**Frontmatter — the REDUCED live-capture schema.** A live close-session capture cannot populate the batch-export telemetry (token counts, message counts, duration_ms, permission_mode) — those only exist in the account export pipeline. Write only what's truthfully known, and mark fidelity so a reader (or the batch pipeline) can tell a live capture from an export: - **`insights[]`** — the generalizable learnings from Phase 1 (decisions/failures worth re-reading). Each: `{text, wing, hall}`; add `supersede_slug` to revise a prior note in place instead of creating a duplicate. `hall` ∈ facts/decisions/failures/hypotheses/sources.
- **`tickets[]`** — the issue actions from Phase 3: `{repo, action, ...}` where action ∈ create/close/comment. Owner is always `mathias` (server-forced).
- **`summary`** — `{title, body, repos_touched}`. Capture writes it to `ai-sessions` and stamps `fidelity` in frontmatter. Body stays reconstructable: one-paragraph summary, decisions, key artifacts, open threads.
- **`context`** — `{harness: "claudeai-chat", session_ref: <chatid8-or-slug>, fidelity: "live-capture", actor: "mathias", classification: <see gate below>}`.
```yaml **THE CLASSIFICATION GATE (read before calling — this is where capture refuses).**
--- Capture computes an **effective classification = the strictest across EVERY target it touches** (each insight's `wing`, each ticket's `repo`, and every entry in `summary.repos_touched`), then refuses if that effective level is `confidential` and the origin is us-nexus (claude.ai is us-nexus). Levels come from `classification.yaml` at the brain root (source of truth, #67), with the code defaults as the floor: `hyperguild`/`homelab` → internal; `client-*` → confidential; **anything untagged → confidential (fail-safe)**.
title: "<concise session title>" - **Tagged `internal` today** (safe through claude.ai): wings `hyperguild`, `homelab`; repos `brain`, `ai-sessions`, `infra`, `hyperguild`, `homelab`, `tapir`, `agentsquad`, `jepa-fx-risk`, `swedsl`. Treat `classification.yaml` as authoritative — this list is a hint, not gospel.
client: "claudeai" - Declare `context.classification: "internal"` for normal homelab work.
interface: "claudeai-chat" - `summary.repos_touched`, insight `wing`s, and ticket `repo`s are classification INPUTS, not free-form metadata — every target must resolve `internal` or the whole capture escalates to `confidential` and the gate refuses via claude.ai. Listing the central homelab repos (incl. `brain`/`ai-sessions`) is now fine; they're tagged. The summary always lands in `ai-sessions` (internal), so the summary path itself never escalates.
date: "<YYYY-MM-DD>" - If a session genuinely touched **`client-*` or otherwise-untagged** material, it cannot be captured through claude.ai — note that in the verdict rather than trying to force it.
repos_touched: [<repo slugs>]
topic_tags: [<tags>]
outcome: "<shipped|in-progress|abandoned>"
fidelity: "live-capture" # NOT an export; reconstructed live from chat
captured_by: "close-session-skill"
---
```
Do not invent the export-only fields. `fidelity: live-capture` is the honest signal; if the batch export later produces a richer summary for the same session, the export is source of truth and supersedes this. **GATE — dry-run first, then execute.**
1. Call `capture` with `dry_run: true`. It validates the whole payload and returns the would-be receipt + `effective_classification`, writing nothing.
2. **STOP. Show the dry-run receipt** (effective classification, the insights/tickets/summary that would land) and get explicit confirmation.
3. On confirmation, call `capture` again with `dry_run: false`. Read the returned receipt: it is partial-aware (`errors[]`, per-item `ok`). Report exactly what landed.
**Body** (keep it reconstructable, not exhaustive): If `capture` is **unreachable** (tool not on the connector — e.g. a session that started before a deploy; a tool-list refresh usually fixes it): say so. Only then fall back to the legacy inline path (`gitea:file_write_branch` summary + `brain_write`/`brain_update` + `brain_get` confirm), and note in the verdict that the I5 audit record was NOT produced.
```markdown
## One-paragraph summary
## Decisions
## Key artifacts
## Open threads
```
**GATE — STOP, show the full file (path + frontmatter + body), get explicit confirmation before committing.** ## Phase 5 — Verdict
## Phase 5 — Brain orientation note (the durable "where we are" record)
Write one brain note so a fresh session can orient without the chat. This uses the `brain_update`/`brain_get` verbs (live since 2026-06).
**Target:** `wing: <domain>` (the project/topic domain, e.g. `hyperguild`, `jepa-fx`), `hall: decisions`. The note is a knowledge-type record (a decision/orientation), grouped by knowledge-type, not by interface surface.
**Batch read-after-write discipline (important — do these in order, do not interleave):**
1. **Read first, before any write.** Check whether an orientation note already exists for this wing/topic. Do your "does this already exist / what should I supersede" reads NOW, up front. BM25/keyword search and `brain_get` are immediate; semantic/vector search may lag up to ~5 min after a write, so never rely on a semantic query to find something you wrote earlier in this same run.
2. **Write or supersede:**
- **New note** → `brain_write` (wing, hall: decisions). Returns `{id, path, content_hash}`.
- **Superseding a prior orientation note** → `brain_update` (slug or path, wing, hall, content, reason). Whole-note replace; stamps `supersedes`/`updated_at`; returns `{id, path, content_hash, superseded}`. Use this instead of a second `brain_write` to the same slug — blind re-write creates duplicates/contradictions, which is the exact failure brain_update exists to prevent.
3. **Confirm it landed** via `brain_get(id)` and check the returned `content_hash` matches what the write returned. This is the read-after-write confirmation — do it with `brain_get`, never a semantic query.
**RULE: no semantic/vector brain query after the first `brain_update` in this run.** The batch shape makes this natural — read up front, write, confirm by id. If you ever find the skill wanting to semantic-search a just-superseded note, stop and flag it (that's the signal the staleness window matters and needs the synchronous-reembed follow-up).
**GATE — STOP, show the note (target wing/hall, new-vs-supersede, full content), get explicit confirmation before the brain write.**
After the note lands, if it relates to a note in another wing, create the cross-link inline with `brain_tunnel(source, target)` (idempotent; both paths brain-relative, must be in different wings). Optionally append a `session_log` entry (`session_id`, `skill: close-session`, `phase`, `final_status`) for telemetry. Both are now callable directly from Claude.ai — no Claude Code/Crush handoff needed.
## Phase 6 — Verdict
Deliver a final "safe to archive" verdict in the chat. Either: Deliver a final "safe to archive" verdict in the chat. Either:
- **SAFE TO ARCHIVE** — list what landed (issues closed/filed with numbers, summary path, brain note id, any tunnels) so the trail is auditable. Then list anything still in the user's queue (e.g. a PR awaiting their merge, a decision owed next session). - **SAFE TO ARCHIVE** — list what landed from the capture receipt (issues closed/filed with numbers, summary path, brain note ids/paths) so the trail is auditable. Then list anything still in the user's queue (e.g. a PR awaiting their merge, a decision owed next session).
- **NOT YET** — name the specific gate that wasn't passed or the write that failed, and what to do about it. - **NOT YET** — name the specific gate that wasn't passed, the capture refusal reason, or the per-item error from the receipt, and what to do about it.
Never claim safe-to-archive if any gated write was declined or errored. The verdict is the skill's contract: if it says safe, the session can be lost without losing the work. Never claim safe-to-archive if the capture refused, any receipt item errored, or a gated confirmation was declined. The verdict is the skill's contract: if it says safe, the session can be lost without losing the work.
## Why the gates and the batch discipline matter ## Why the gate and the single-call shape matter
The whole point is durability across a context reset. Every gate is a place where a wrong write would silently corrupt the record (close the wrong issue, overwrite a good brain note, commit a half-truth). The batch read-discipline in Phase 5 exists because the brain's vector index refreshes out-of-band: write-then-semantically-reread in the same run can read stale, so the skill front-loads reads and confirms writes by id. Get those right and the skill does what it promises — nothing important is lost when the chat goes away. The whole point is durability across a context reset. The capture call is the one place a wrong payload would silently corrupt the record (close the wrong issue, escalate to a refusal, commit a half-truth), which is why it is dry-run-then-confirm. Routing everything through one `capture` keeps the supersession discipline, the read-after-write confirmation, and the I5 audit trail server-side — the skill never has to carry those rules itself, and every closeout is uniformly audited. Get the payload and the classification right and the skill does what it promises — nothing important is lost when the chat goes away.
+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)