Compare commits

..
Author SHA1 Message Date
mathiasandClaude Opus 4.8 ef11864121 feat(mcp): brain_pending + brain_promote tools + REST routes (#38)
CI / Lint / Test / Vet (pull_request) Successful in 13s
CI / Mirror to GitHub (pull_request) Has been skipped
Wires the curation primitives into the MCP surface (all three sites:
tools() descriptors, handleCall dispatch, package doc) and registers the
GET /pending + POST /promote REST routes. brain_promote additionally
re-indexes the promoted note into the graph (best-effort), matching the
other write paths. brain_pending is the human-review-queue complement to
the agent-facing brain_write.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-26 16:34:15 +02:00
mathiasandClaude Opus 4.8 14b04a25cb feat(brain): ListPending + PromoteNote — raw→wiki curation primitives (#38)
Closes the human curation loop: list brain/raw/ notes awaiting review
(oldest-first, with excerpt) and promote one into brain/wiki/<wing>/<hall>/.

PromoteNote rewrites frontmatter (sets wing/hall/promoted_at, preserves
created_at + custom fields via the existing frontmatter editor), rebuilds
the wing _index, and runs auto-tunnel. Atomic from the caller's view:
hall/wing/slug/collision validation happens before any fs change, and the
source is deleted only after the destination write succeeds (write-then-
delete, never move) — a collision or write failure leaves raw/ intact.
Default slug = filename minus the YYYY-MM-DD- prefix.

Also exposes GET /pending + POST /promote for shell scripts (bad
hall/collision/missing-source → 400, not 500).

Scope note: operates on raw/ (the retrospective-skill review queue) per
this issue's spec. The knowledge/ legacy pile is #22's bulk-migration
concern, not this ongoing queue.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-26 16:34:15 +02:00
mathias 66a9b8e725 Merge pull request 'docs(close-session): refresh classification gate post-#67 (#62)' (#69) from fix/close-session-classification-refresh into main
CI / Lint / Test / Vet (push) Successful in 12s
CI / Mirror to GitHub (push) Successful in 3s
2026-06-26 14:28:48 +00:00
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
mathiasandClaude Opus 4.8 a1997838b0 feat(capturehttp): POST /capture REST adapter + OAuth2 + origin resolver (#53)
CI / Lint / Test / Vet (pull_request) Successful in 12s
CI / Mirror to GitHub (pull_request) Has been skipped
The HTTP door for the capture capability. Thin: authenticate → derive
trust-zone origin → decode → capture.Service → map receipt to status.

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

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

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

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

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-22 23:42:18 +02:00
mathias aad90f2dfe Merge pull request 'feat: Gitea IssueTracker client + inject into brain server (#52, capture 49c)' (#58) from feat/capture-gitea-tracker into main
CI / Lint / Test / Vet (push) Successful in 12s
CI / Mirror to GitHub (push) Successful in 4s
2026-06-22 21:32:02 +00:00
mathias 07fca9ee73 Merge pull request 'feat: CaptureService use-case + BrainStore extraction (#51, capture 49b)' (#57) from feat/capture-service into main
CI / Lint / Test / Vet (push) Successful in 12s
CI / Mirror to GitHub (push) Has been cancelled
2026-06-22 21:31:47 +00:00
mathias f6bf9b5f57 Merge pull request 'feat: capture classification taxonomy + per-wing/repo tags (#50, capture 49a)' (#56) from feat/capture-classification into main
CI / Lint / Test / Vet (push) Has been cancelled
CI / Mirror to GitHub (push) Has been cancelled
2026-06-22 21:31:18 +00:00
mathiasandClaude Opus 4.8 6606b38a76 feat(gitea): IssueTracker client + inject into brain server (#52)
Implements the IssueTracker port as a real Gitea REST client (#49c) — the
new outbound dependency the brain server gains for capture.

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

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

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

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-22 23:27:29 +02:00
mathiasandClaude Opus 4.8 0ac165cca3 feat(brainstore): shared BrainStore impl; re-point MCP handlers (#51)
Extracts the #45 write/update/get logic + the wiki upkeep that must
accompany a write (wing _index rebuild, cross-wing auto-tunnel, graph
re-index) into a single concrete brainstore.Store implementing
capture.BrainStore. The MCP brain_write/brain_update/brain_get handlers
are re-pointed at it, so there is one implementation, not two — the DRY
payoff #51 is named for. capture and MCP now share the exact same brain
write path and read-after-write contract.

The Server gains a *brainstore.Store, constructed in NewServer and given
the graph store in WithGraph. Embedding refresh stays out-of-band
(mtime-driven vectorstore.Sync), unchanged. Existing MCP brain_update/
brain_get/brain_write tests pass unmodified — behaviour and the
{id, path, content_hash} response contract are preserved.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-22 23:21:43 +02:00
mathiasandClaude Opus 4.8 43f92e3102 feat(capture): CaptureService use-case + ports + entities (#51)
The Clean-Architecture core of the capture capability (#49b). Pure
orchestration over ports — no HTTP, no live Gitea, no audit I/O — fully
unit-tested against fakes before any adapter exists.

- Ports: BrainStore (#45 write/update/get), IssueTracker, SummaryWriter,
  ClassificationPolicy (satisfied by #50's classification.Config),
  AuditSink. Entities: Insight, Ticket, Summary, CaptureContext,
  CaptureInput, CaptureReceipt.
- CaptureService.Capture: validate-before-write (fail-closed), resolve
  effective classification (stricter of declared vs target-derived;
  under-declaration logged as a security event), orchestrate insights
  (write/supersede) → tickets → summary best-effort, emit a request-level
  audit record of exactly what landed, return a partial-aware receipt.
- dry_run short-circuits after validation, writes nothing (not even audit).

Out of scope here, layered on later: the I1 origin sovereignty gate (#53,
needs the server-derived principal) and the classification-aware audit
degradation/refusal (#54). "Effective" is folded into the service as
classification.Stricter rather than a port method — the stricter-wins
rule is use-case policy.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-22 23:21:43 +02:00
mathiasandClaude Opus 4.8 98cfae595c feat(classification): data-sensitivity taxonomy + per-wing/repo tags (#50)
CI / Lint / Test / Vet (pull_request) Successful in 13s
CI / Mirror to GitHub (pull_request) Has been skipped
Implements the I1-prerequisite from #49/#50: the classification taxonomy
and the per-wing / per-repo tagging the capture server reads to derive a
target's sensitivity.

- Levels public < internal < confidential, ordered so "stricter wins"
  (spec §4.1 model C) is a plain max via Stricter().
- Tags read from an optional classification.yaml at the brain root
  (wings:/repos: maps). Absent file → defaults-only, not an error.
- Defaulting: client-* → confidential; hyperguild/homelab → internal;
  everything else → confidential. Fail-safe-to-strictest is the
  load-bearing property: a missing tag never silently downgrades.
- Config.Derive(Target) is the function the use-case calls; Wing/Repo
  are the per-kind helpers. ParseLevel rejects unknown tokens; a bad
  level in the config file is a hard load error.

Central classification.yaml (not _index.md frontmatter, not gitea repo
topics): classifying a repo needs no live Gitea client, so #50 has no
dependency on the tracker work (#52); it's auditable in one place; and
it avoids BuildWingIndex clobbering a wing's regenerated _index.md.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-22 22:42:05 +02:00
mathias 2a595b5a92 docs(capture): make Q4 audit-sink-down posture classification-aware (#49)
CI / Lint / Test / Vet (push) Successful in 12s
CI / Mirror to GitHub (push) Successful in 3s
Reconsidered Q4: instead of one global degrade-and-warn, the posture now
inherits from effective classification (Q1):
- confidential + audit-sink-down -> hard-refuse (no buffer; removes the
  buffer-integrity question for confidential data)
- internal/public + audit-sink-down -> degrade-and-warn + durable local
  buffer + ntfy + reconcile-on-recovery
- floor (all tiers): refuse if nothing can record the audit
Updated the I5 Gherkin scenarios + obligations row to match. Couples Q4
to the Q1 classification spine -> one coherent sensitivity model.
2026-06-22 20:20:40 +00:00
mathias b7a2cc5fdf docs(capture): resolve §4 open questions into binding decisions (#49)
CI / Lint / Test / Vet (push) Successful in 14s
CI / Mirror to GitHub (push) Successful in 4s
Q1 classification trust: model (C) — caller declares, server cross-checks
  target tag, stricter wins, mismatch logged. Needs a classification
  taxonomy + per-wing/repo tags (prerequisite, sub-task of #49).
Q2 harness origin: server-derived from authenticated principal;
  context.harness is descriptive-only, never a gate input.
Q3 central relay: ships in v1 (needed for claude.ai/Crush/Pi/LLM Council)
  + I2 security-baseline ledger entry is v1 work.
Q4 audit-sink-down: degrade-and-warn + durable local buffer + reconcile
  on recovery; refuse only if NOTHING can record the audit.

Updated the I1 sovereignty + I5 auditability Gherkin scenarios to match;
added classification-mismatch and origin-spoofing scenarios.
2026-06-22 20:17:07 +00:00
mathias db638cca11 docs(capture): add use-case + BDD spec for the capture capability (#49)
CI / Lint / Test / Vet (push) Successful in 14s
CI / Mirror to GitHub (push) Successful in 3s
Pre-build spec for the uniform cross-harness capture path. Grounds the
design in the homelab invariants (I1-I5, now canonical on infra main):
- I1 sovereignty gate: refuse confidential capture via us-nexus harness
- I2: distributed-library form (no high-degree observer node) is
  admissible; central relay needs a security-baseline ledger entry
- I5: capture is a privileged write path, must emit audit records
Gherkin scenarios cover the happy path, the sovereignty refusal,
supersession + staleness discipline (#45/#47), fail-closed validation,
best-effort partial-failure receipts, dry-run, and fidelity supersession.
Four open questions flagged for review (classification trust is the
highest-risk one).
2026-06-22 17:25:03 +00:00
mathias 38579598e0 feat(skills): add close-session workflow skill
CI / Lint / Test / Vet (push) Successful in 14s
CI / Mirror to GitHub (push) Successful in 4s
Disciplined end-of-session closeout for Claude.ai chats: harvest →
ground-truth gitea → confirm issue actions → commit session summary →
brain orientation note → safe-to-archive verdict.

Finalized against current infra (git.d-ma.be, koala:30401 LiteLLM,
Authentik) and the brain_update/brain_get verbs. Phase 5 uses
supersede-by-slug + brain_get read-after-write with the batch
no-semantic-query-after-supersede discipline. session_log/brain_tunnel
inlined (now callable from claude.ai). Summary frontmatter is the
reduced live-capture schema with fidelity:live-capture to distinguish
from batch-export summaries.
2026-06-22 15:57:02 +00:00
mathias d6fa92b176 Merge pull request 'chore(context): drop unsupported cursor + aider adapters' (#48) from chore/drop-cursor-aider-adapters into main
CI / Lint / Test / Vet (push) Successful in 12s
CI / Mirror to GitHub (push) Successful in 4s
2026-06-22 09:03:59 +00:00
35 changed files with 4712 additions and 62 deletions
+93
View File
@@ -8,6 +8,7 @@ import (
"net/http" "net/http"
"net/url" "net/url"
"os" "os"
"path/filepath"
"strconv" "strconv"
"strings" "strings"
"time" "time"
@@ -15,8 +16,13 @@ import (
chassisauth "gitea.d-ma.be/mathias/mcp-chassis/auth" chassisauth "gitea.d-ma.be/mathias/mcp-chassis/auth"
"github.com/mathiasbq/hyperguild/ingestion/internal/api" "github.com/mathiasbq/hyperguild/ingestion/internal/api"
"github.com/mathiasbq/hyperguild/ingestion/internal/audit"
"github.com/mathiasbq/hyperguild/ingestion/internal/capture"
"github.com/mathiasbq/hyperguild/ingestion/internal/capturehttp"
"github.com/mathiasbq/hyperguild/ingestion/internal/classification"
"github.com/mathiasbq/hyperguild/ingestion/internal/claudewatcher" "github.com/mathiasbq/hyperguild/ingestion/internal/claudewatcher"
"github.com/mathiasbq/hyperguild/ingestion/internal/embed" "github.com/mathiasbq/hyperguild/ingestion/internal/embed"
"github.com/mathiasbq/hyperguild/ingestion/internal/gitea"
"github.com/mathiasbq/hyperguild/ingestion/internal/graphstore" "github.com/mathiasbq/hyperguild/ingestion/internal/graphstore"
"github.com/mathiasbq/hyperguild/ingestion/internal/graphsync" "github.com/mathiasbq/hyperguild/ingestion/internal/graphsync"
"github.com/mathiasbq/hyperguild/ingestion/internal/llm" "github.com/mathiasbq/hyperguild/ingestion/internal/llm"
@@ -118,6 +124,47 @@ func envInt(key string, fallback int) int {
return fallback return fallback
} }
// buildAuditSink selects the capture audit sink. When BRAIN_LOKI_URL is
// set it builds the classification-aware DegradingSink (loki central +
// durable file buffer + optional ntfy) and starts the reconcile loop;
// otherwise it falls back to a plain slog sink. The buffer lives under the
// brain dir so it survives process restarts.
func buildAuditSink(ctx context.Context, brainDir string, logger *slog.Logger) capture.AuditSink {
lokiURL := os.Getenv("BRAIN_LOKI_URL")
central := audit.NewLokiCentral(lokiURL)
if central == nil {
logger.Info("capture audit: slog sink (BRAIN_LOKI_URL unset)")
return audit.NewSlogSink(logger)
}
buffer, err := audit.NewFileBuffer(filepath.Join(brainDir, ".audit-buffer", "capture.jsonl"))
if err != nil {
logger.Error("capture audit buffer init", "err", err)
os.Exit(1)
}
// Keep notifier as a nil interface (not a typed-nil) when unconfigured
// so DegradingSink/Reconcile skip it cleanly.
var notifier audit.Notifier
if n := audit.NewNtfyNotifier(os.Getenv("BRAIN_NTFY_URL"), os.Getenv("BRAIN_NTFY_TOKEN")); n != nil {
notifier = n
}
reconcileInterval := time.Duration(envInt("BRAIN_AUDIT_RECONCILE_INTERVAL", 60)) * time.Second
audit.StartReconcile(ctx, central, buffer, notifier, reconcileInterval)
logger.Info("capture audit: loki+buffer sink", "loki", lokiURL, "reconcile_s", int(reconcileInterval.Seconds()))
return audit.NewDegradingSink(central, buffer, notifier)
}
// splitList parses a comma-separated env value into a trimmed,
// empty-free slice. Used for the capture sovereign-principal allowlist.
func splitList(v string) []string {
var out []string
for _, p := range strings.Split(v, ",") {
if p = strings.TrimSpace(p); p != "" {
out = append(out, p)
}
}
return out
}
// systemHostname returns os.Hostname() with a "unknown" fallback so the // systemHostname returns os.Hostname() with a "unknown" fallback so the
// caller never has to handle the rare error path. // caller never has to handle the rare error path.
func systemHostname() string { func systemHostname() string {
@@ -175,6 +222,15 @@ func main() {
logger.Info("brain reranker configured", "url", rerankURL, "model", rerankModel) logger.Info("brain reranker configured", "url", rerankURL, "model", rerankModel)
} }
// Gitea ticket tracker for the capture capability (#52). Token via env
// only — never logged or in argv. Both vars must be set to enable it;
// gitea.New returns nil otherwise, leaving ticket integration off.
giteaURL := envOr("BRAIN_GITEA_URL", "https://git.d-ma.be")
if tracker := gitea.New(giteaURL, os.Getenv("BRAIN_GITEA_TOKEN")); tracker != nil {
mcpSrv = mcpSrv.WithIssueTracker(tracker)
logger.Info("brain gitea tracker configured", "url", giteaURL)
}
// Hybrid retrieval (pgvector + nomic-embed-text). Both env vars must // Hybrid retrieval (pgvector + nomic-embed-text). Both env vars must
// be set together for the path to wire on; otherwise BM25-only. // be set together for the path to wire on; otherwise BM25-only.
var vectorStore *vectorstore.PGStore var vectorStore *vectorstore.PGStore
@@ -313,6 +369,8 @@ func main() {
mux.HandleFunc("POST /ingest-path", h.IngestPath) mux.HandleFunc("POST /ingest-path", h.IngestPath)
mux.HandleFunc("POST /ingest-raw", h.IngestRaw) mux.HandleFunc("POST /ingest-raw", h.IngestRaw)
mux.HandleFunc("POST /backfill-refs", h.BackfillRefs) mux.HandleFunc("POST /backfill-refs", h.BackfillRefs)
mux.HandleFunc("GET /pending", h.Pending)
mux.HandleFunc("POST /promote", h.Promote)
mux.HandleFunc("POST /backfill-embeddings", h.BackfillEmbeddings) mux.HandleFunc("POST /backfill-embeddings", h.BackfillEmbeddings)
mux.HandleFunc("GET /pass-rate", h.PassRate) mux.HandleFunc("GET /pass-rate", h.PassRate)
jwtValidator, err := chassisauth.NewJWTValidator(ctx, os.Getenv("DEX_ISSUER_URL"), os.Getenv("MCP_AUDIENCE")) jwtValidator, err := chassisauth.NewJWTValidator(ctx, os.Getenv("DEX_ISSUER_URL"), os.Getenv("MCP_AUDIENCE"))
@@ -339,6 +397,41 @@ func main() {
mux.Handle("/mcp", chassisauth.BearerMiddleware(mcpToken, jwtValidator, "brain", resourceMetadataURL, mcpSrv)) mux.Handle("/mcp", chassisauth.BearerMiddleware(mcpToken, jwtValidator, "brain", resourceMetadataURL, mcpSrv))
// POST /capture (#53/#54): the uniform capture REST door. Needs a ticket
// tracker to file action items, so it only mounts when Gitea is
// configured. It reuses the MCP server's graph-wired brain store (one
// implementation), the classification tags for the I1 gate, and a
// classification-aware audit sink (loki + durable buffer + ntfy when
// BRAIN_LOKI_URL is set, else a plain slog sink). The handler does its
// own auth (static + JWT) because it needs the principal to derive the
// trust-zone origin — the chassis middleware hides it.
if tracker := mcpSrv.IssueTracker(); tracker != nil {
classCfg, cerr := classification.Load(brainDir)
if cerr != nil {
logger.Error("load classification config", "err", cerr)
os.Exit(1)
}
auditSink := buildAuditSink(ctx, brainDir, logger)
// The Gitea client also satisfies SummaryWriter (#66): session
// summaries are written to mathias/ai-sessions over the same API
// token. nil only if a future tracker impl lacks file writes.
summaryWriter, _ := tracker.(capture.SummaryWriter)
captureSvc := capture.NewService(
mcpSrv.BrainStore(), tracker, summaryWriter, classCfg, auditSink)
sovereign := splitList(os.Getenv("BRAIN_CAPTURE_SOVEREIGN_PRINCIPALS"))
resolver := capturehttp.NewOriginResolver(sovereign)
captureH := capturehttp.New(captureSvc, jwtValidator, mcpToken, "local-cli", resolver)
mux.Handle("POST /capture", captureH)
// Same use-case behind the MCP `capture` tool (#55 relay) so MCP-native
// harnesses (claude.ai, Crush, Pi, LLM Council) reach capture through
// the existing /mcp OAuth connector. mcpSrv is already wrapped above;
// WithCapture mutates the same instance, so the tool appears live.
mcpSrv.WithCapture(captureSvc, jwtValidator, mcpToken, "local-cli", resolver)
logger.Info("capture enabled (REST + MCP tool)", "sovereign_principals", len(sovereign))
} else {
logger.Info("capture endpoint disabled (BRAIN_GITEA_TOKEN unset)")
}
// Opt-in OAuth 2.0 client_credentials flow for claude.ai's custom-MCP // Opt-in OAuth 2.0 client_credentials flow for claude.ai's custom-MCP
// integration UI, which has no static-Bearer field. Setting both // integration UI, which has no static-Bearer field. Setting both
// OAUTH_CLIENT_ID and OAUTH_CLIENT_SECRET enables the token exchange; // OAUTH_CLIENT_ID and OAUTH_CLIENT_SECRET enables the token exchange;
+34
View File
@@ -483,6 +483,40 @@ func (h *Handler) BackfillRefs(w http.ResponseWriter, r *http.Request) {
writeJSON(w, map[string]int{"updated": n}) writeJSON(w, map[string]int{"updated": n})
} }
// Pending handles GET /pending — list raw/ notes awaiting promotion.
func (h *Handler) Pending(w http.ResponseWriter, _ *http.Request) {
pending, err := ListPending(h.brainDir)
if err != nil {
h.logger.Error("pending failed", "err", err)
writeError(w, http.StatusInternalServerError, "pending error")
return
}
writeJSON(w, map[string]any{"pending": pending})
}
type promoteRequest struct {
Filename string `json:"filename"`
Wing string `json:"wing"`
Hall string `json:"hall"`
Slug string `json:"slug,omitempty"`
}
// Promote handles POST /promote — move a raw/ note into the wiki. A bad
// hall / collision / missing source is a 400 (caller error), not a 500.
func (h *Handler) Promote(w http.ResponseWriter, r *http.Request) {
var req promoteRequest
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
writeError(w, http.StatusBadRequest, "invalid JSON")
return
}
rel, err := PromoteNote(h.brainDir, PromoteOptions(req))
if err != nil {
writeError(w, http.StatusBadRequest, err.Error())
return
}
writeJSON(w, map[string]string{"path": rel})
}
func writeJSON(w http.ResponseWriter, v any) { func writeJSON(w http.ResponseWriter, v any) {
w.Header().Set("Content-Type", "application/json") w.Header().Set("Content-Type", "application/json")
json.NewEncoder(w).Encode(v) //nolint:errcheck json.NewEncoder(w).Encode(v) //nolint:errcheck
+156
View File
@@ -0,0 +1,156 @@
package api
import (
"fmt"
"os"
"path/filepath"
"regexp"
"sort"
"strings"
"time"
"github.com/mathiasbq/hyperguild/ingestion/internal/brain"
)
// PendingNote describes a raw/ note awaiting human promotion to the wiki.
type PendingNote struct {
Filename string `json:"filename"`
CreatedAt string `json:"created_at"`
SizeBytes int64 `json:"size_bytes"`
Excerpt string `json:"excerpt"`
}
// datePrefix matches a leading YYYY-MM-DD- on a raw filename, stripped when
// deriving the default promoted slug.
var datePrefix = regexp.MustCompile(`^\d{4}-\d{2}-\d{2}-`)
// ListPending returns the notes in brain/raw/ awaiting review, oldest-first
// (natural review order). An absent raw/ dir yields an empty slice, not an
// error. Only .md files are listed; tunnel-candidate files are skipped.
func ListPending(brainDir string) ([]PendingNote, error) {
dir := filepath.Join(brainDir, "raw")
entries, err := os.ReadDir(dir)
if err != nil {
if os.IsNotExist(err) {
return []PendingNote{}, nil
}
return nil, fmt.Errorf("read raw dir: %w", err)
}
out := make([]PendingNote, 0, len(entries))
for _, e := range entries {
if e.IsDir() || !strings.HasSuffix(e.Name(), ".md") || strings.HasPrefix(e.Name(), "tunnel-candidates-") {
continue
}
info, statErr := e.Info()
if statErr != nil {
continue
}
raw, readErr := os.ReadFile(filepath.Join(dir, e.Name()))
if readErr != nil {
continue
}
fm, body := parseFrontmatter(string(raw))
created := fm.get("created_at")
if created == "" {
created = info.ModTime().UTC().Format(time.RFC3339)
}
out = append(out, PendingNote{
Filename: e.Name(),
CreatedAt: created,
SizeBytes: info.Size(),
Excerpt: excerpt(body, 200),
})
}
sort.SliceStable(out, func(i, j int) bool { return out[i].CreatedAt < out[j].CreatedAt })
return out, nil
}
// PromoteOptions identifies a raw note to promote and its wiki destination.
type PromoteOptions struct {
Filename string // basename in brain/raw/
Wing string
Hall string
Slug string // optional; defaults to Filename minus date prefix + .md
}
// PromoteNote moves a note from brain/raw/ into the structured wiki: it
// rewrites frontmatter (sets wing/hall/promoted_at, preserves created_at and
// any custom fields), writes to brain/wiki/<wing>/<hall>/<slug>.md, deletes
// the source, then rebuilds the wing index and runs auto-tunnel detection.
//
// It is atomic from the caller's view: validation (hall, wing, slug,
// collision) happens before any filesystem change, and the source is deleted
// only after the destination write succeeds (write-then-delete, never move).
// Returns the promoted note's path relative to brainDir.
func PromoteNote(brainDir string, opts PromoteOptions) (string, error) {
// Validate filename (basename only — no traversal) before touching fs.
base := filepath.Base(opts.Filename)
if base != opts.Filename || base == "." || base == ".." || strings.ContainsAny(opts.Filename, `/\`) {
return "", fmt.Errorf("invalid filename %q", opts.Filename)
}
slug := opts.Slug
if slug == "" {
slug = datePrefix.ReplaceAllString(strings.TrimSuffix(base, ".md"), "")
}
// NotePath validates hall + wing + slug; do this before reading anything.
dest, err := brain.NotePath(brainDir, opts.Wing, opts.Hall, slug)
if err != nil {
return "", err
}
src := filepath.Join(brainDir, "raw", base)
raw, err := os.ReadFile(src)
if err != nil {
if os.IsNotExist(err) {
return "", fmt.Errorf("pending note %q does not exist in raw/", base)
}
return "", fmt.Errorf("read source: %w", err)
}
// Collision: never silently overwrite an existing promoted note.
if _, statErr := os.Stat(dest); statErr == nil {
rel, _ := filepath.Rel(brainDir, dest)
return "", fmt.Errorf("target %s already exists; choose a different slug", filepath.ToSlash(rel))
}
fm, body := parseFrontmatter(string(raw))
now := time.Now().UTC().Format(time.RFC3339)
fm.set("wing", brain.Sanitise(opts.Wing))
fm.set("hall", opts.Hall)
if fm.get("created_at") == "" {
fm.set("created_at", now)
}
fm.set("promoted_at", now)
if err := os.MkdirAll(filepath.Dir(dest), 0o755); err != nil {
return "", fmt.Errorf("create wing dir: %w", err)
}
// Write-then-delete: the source survives any write failure.
if err := os.WriteFile(dest, []byte(fm.render()+body), 0o644); err != nil {
return "", fmt.Errorf("write promoted note: %w", err)
}
if err := os.Remove(src); err != nil {
return "", fmt.Errorf("promoted note written but source removal failed: %w", err)
}
rel, _ := filepath.Rel(brainDir, dest)
relSlash := filepath.ToSlash(rel)
// Best-effort wiki upkeep — the note is already promoted.
_ = brain.BuildWingIndex(brainDir, opts.Wing)
_ = brain.AutoTunnel(brainDir, relSlash, body)
return relSlash, nil
}
// excerpt returns the first n runes of s, trimmed, single-spaced.
func excerpt(s string, n int) string {
s = strings.TrimSpace(s)
r := []rune(s)
if len(r) > n {
r = r[:n]
}
return strings.TrimSpace(string(r))
}
+121
View File
@@ -0,0 +1,121 @@
package api
import (
"os"
"path/filepath"
"testing"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func writeRaw(t *testing.T, brainDir, name, content string) {
t.Helper()
dir := filepath.Join(brainDir, "raw")
require.NoError(t, os.MkdirAll(dir, 0o755))
require.NoError(t, os.WriteFile(filepath.Join(dir, name), []byte(content), 0o644))
}
func TestListPendingEmptyWhenAbsent(t *testing.T) {
got, err := ListPending(t.TempDir())
require.NoError(t, err, "absent raw/ is not an error")
assert.Empty(t, got)
}
func TestListPendingReturnsOldestFirstWithExcerpt(t *testing.T) {
dir := t.TempDir()
writeRaw(t, dir, "2026-06-02-newer.md", "---\ncreated_at: 2026-06-02T00:00:00Z\n---\nNewer body here.\n")
writeRaw(t, dir, "2026-06-01-older.md", "---\ncreated_at: 2026-06-01T00:00:00Z\n---\nOlder body content.\n")
// non-md ignored
writeRaw(t, dir, "notes.txt", "ignore me")
got, err := ListPending(dir)
require.NoError(t, err)
require.Len(t, got, 2)
assert.Equal(t, "2026-06-01-older.md", got[0].Filename, "oldest first")
assert.Equal(t, "2026-06-02-newer.md", got[1].Filename)
assert.Contains(t, got[0].Excerpt, "Older body content")
assert.NotContains(t, got[0].Excerpt, "---", "excerpt is body, not frontmatter")
assert.Positive(t, got[0].SizeBytes)
}
func TestPromoteHappyPath(t *testing.T) {
dir := t.TempDir()
writeRaw(t, dir, "2026-06-01-lejpa-decision.md",
"---\ncreated_at: 2026-06-01T09:00:00Z\ncustom_field: keep-me\n---\n# LeJEPA\n\nbody.\n")
rel, err := PromoteNote(dir, PromoteOptions{
Filename: "2026-06-01-lejpa-decision.md", Wing: "jepa-fx", Hall: "decisions",
})
require.NoError(t, err)
assert.Equal(t, "wiki/jepa-fx/decisions/lejpa-decision.md", rel, "slug defaults to filename minus date prefix")
// Source deleted.
_, statErr := os.Stat(filepath.Join(dir, "raw", "2026-06-01-lejpa-decision.md"))
assert.True(t, os.IsNotExist(statErr), "source removed after promote")
got, err := os.ReadFile(filepath.Join(dir, filepath.FromSlash(rel)))
require.NoError(t, err)
s := string(got)
assert.Contains(t, s, "wing: jepa-fx")
assert.Contains(t, s, "hall: decisions")
assert.Contains(t, s, "created_at: 2026-06-01T09:00:00Z", "original created_at preserved")
assert.Contains(t, s, "promoted_at:")
assert.Contains(t, s, "custom_field: keep-me", "custom frontmatter preserved")
assert.Contains(t, s, "# LeJEPA")
}
func TestPromoteExplicitSlug(t *testing.T) {
dir := t.TempDir()
writeRaw(t, dir, "2026-06-01-x.md", "body\n")
rel, err := PromoteNote(dir, PromoteOptions{Filename: "2026-06-01-x.md", Wing: "a", Hall: "facts", Slug: "custom-slug"})
require.NoError(t, err)
assert.Equal(t, "wiki/a/facts/custom-slug.md", rel)
}
func TestPromoteInvalidHallErrorsBeforeTouchingFS(t *testing.T) {
dir := t.TempDir()
writeRaw(t, dir, "2026-06-01-x.md", "body\n")
_, err := PromoteNote(dir, PromoteOptions{Filename: "2026-06-01-x.md", Wing: "a", Hall: "garbage"})
require.Error(t, err)
// Source untouched.
_, statErr := os.Stat(filepath.Join(dir, "raw", "2026-06-01-x.md"))
assert.NoError(t, statErr, "invalid hall must not delete or move the source")
}
func TestPromoteMissingSourceErrors(t *testing.T) {
_, err := PromoteNote(t.TempDir(), PromoteOptions{Filename: "ghost.md", Wing: "a", Hall: "facts"})
require.Error(t, err)
}
func TestPromoteSlugCollisionNoOverwrite(t *testing.T) {
dir := t.TempDir()
// Pre-existing target.
dest := filepath.Join(dir, "wiki", "a", "facts", "x.md")
require.NoError(t, os.MkdirAll(filepath.Dir(dest), 0o755))
require.NoError(t, os.WriteFile(dest, []byte("EXISTING\n"), 0o644))
writeRaw(t, dir, "2026-06-01-x.md", "NEW\n")
_, err := PromoteNote(dir, PromoteOptions{Filename: "2026-06-01-x.md", Wing: "a", Hall: "facts"})
require.Error(t, err, "collision must error, not overwrite")
got, _ := os.ReadFile(dest)
assert.Equal(t, "EXISTING\n", string(got), "target not overwritten")
_, statErr := os.Stat(filepath.Join(dir, "raw", "2026-06-01-x.md"))
assert.NoError(t, statErr, "source preserved on collision (atomic: no delete without write)")
}
func TestPromoteRejectsTraversalFilename(t *testing.T) {
_, err := PromoteNote(t.TempDir(), PromoteOptions{Filename: "../escape.md", Wing: "a", Hall: "facts"})
require.Error(t, err)
}
func TestPromoteRebuildsWingIndex(t *testing.T) {
dir := t.TempDir()
writeRaw(t, dir, "2026-06-01-x.md", "---\ntitle: X Note\n---\nbody\n")
_, err := PromoteNote(dir, PromoteOptions{Filename: "2026-06-01-x.md", Wing: "a", Hall: "facts"})
require.NoError(t, err)
idx, err := os.ReadFile(filepath.Join(dir, "wiki", "a", "_index.md"))
require.NoError(t, err, "wing _index regenerated")
assert.Contains(t, string(idx), "x", "promoted note appears in the index")
}
+167
View File
@@ -0,0 +1,167 @@
package audit
import (
"bufio"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"fmt"
"os"
"path/filepath"
"sync"
"github.com/mathiasbq/hyperguild/ingestion/internal/capture"
)
// FileBuffer is a durable, restart-surviving audit buffer backed by a
// JSONL file: one {id, entry} record per line. It is the internal/public
// tier fallback when loki is unreachable. Confirm rewrites the file
// without the confirmed record, so a record is cleared only after its
// central write is confirmed.
//
// Access is serialised by a mutex; the buffer is low-throughput (only
// written during a loki outage), so a whole-file rewrite on Confirm is
// acceptable and keeps the on-disk format trivially correct.
type FileBuffer struct {
path string
mu sync.Mutex
}
type bufferLine struct {
ID string `json:"id"`
Entry capture.AuditEntry `json:"entry"`
}
// NewFileBuffer returns a buffer backed by path. The parent directory is
// created if needed. The file itself is created lazily on first Append.
func NewFileBuffer(path string) (*FileBuffer, error) {
if err := os.MkdirAll(filepath.Dir(path), 0o755); err != nil {
return nil, fmt.Errorf("create buffer dir: %w", err)
}
return &FileBuffer{path: path}, nil
}
// Writable reports whether the buffer file can be appended to. It probes
// by opening the file for append (creating it if absent) — the same
// operation Append performs — so Reserve's check matches Append's reality.
func (b *FileBuffer) Writable() error {
b.mu.Lock()
defer b.mu.Unlock()
f, err := os.OpenFile(b.path, os.O_CREATE|os.O_APPEND|os.O_WRONLY, 0o644)
if err != nil {
return err
}
return f.Close()
}
// Append durably writes one audit record. The ID is derived from the
// content + timestamp so it is stable and unique per record.
func (b *FileBuffer) Append(e capture.AuditEntry) error {
b.mu.Lock()
defer b.mu.Unlock()
line := bufferLine{ID: recordID(e), Entry: e}
data, err := json.Marshal(line)
if err != nil {
return fmt.Errorf("marshal buffer line: %w", err)
}
f, err := os.OpenFile(b.path, os.O_CREATE|os.O_APPEND|os.O_WRONLY, 0o644)
if err != nil {
return err
}
defer func() { _ = f.Close() }()
if _, err := f.Write(append(data, '\n')); err != nil {
return err
}
return f.Sync()
}
// Pending reads all buffered records. A missing file means none.
func (b *FileBuffer) Pending() ([]Buffered, error) {
b.mu.Lock()
defer b.mu.Unlock()
return b.readAllLocked()
}
func (b *FileBuffer) readAllLocked() ([]Buffered, error) {
f, err := os.Open(b.path)
if err != nil {
if os.IsNotExist(err) {
return nil, nil
}
return nil, err
}
defer func() { _ = f.Close() }()
var out []Buffered
sc := bufio.NewScanner(f)
sc.Buffer(make([]byte, 0, 64*1024), 1024*1024)
for sc.Scan() {
raw := sc.Bytes()
if len(raw) == 0 {
continue
}
var l bufferLine
if err := json.Unmarshal(raw, &l); err != nil {
return nil, fmt.Errorf("parse buffer line: %w", err)
}
out = append(out, Buffered(l))
}
return out, sc.Err()
}
// Confirm removes a single record after its central write is confirmed, by
// rewriting the file without it. Unknown IDs are a no-op.
func (b *FileBuffer) Confirm(id string) error {
b.mu.Lock()
defer b.mu.Unlock()
all, err := b.readAllLocked()
if err != nil {
return err
}
tmp := b.path + ".tmp"
f, err := os.OpenFile(tmp, os.O_CREATE|os.O_TRUNC|os.O_WRONLY, 0o644)
if err != nil {
return err
}
w := bufio.NewWriter(f)
kept := 0
for _, rec := range all {
if rec.ID == id {
continue
}
data, _ := json.Marshal(bufferLine(rec))
if _, err := w.Write(append(data, '\n')); err != nil {
_ = f.Close()
return err
}
kept++
}
if err := w.Flush(); err != nil {
_ = f.Close()
return err
}
if err := f.Sync(); err != nil {
_ = f.Close()
return err
}
if err := f.Close(); err != nil {
return err
}
// Empty buffer → remove the file entirely so Pending sees nothing.
if kept == 0 {
_ = os.Remove(tmp)
return os.Remove(b.path)
}
return os.Rename(tmp, b.path)
}
// recordID is a stable per-record identifier: sha256 of the principal,
// timestamp, and item list. Distinct captures never collide; the same
// buffered record always hashes the same.
func recordID(e capture.AuditEntry) string {
h := sha256.New()
_, _ = fmt.Fprintf(h, "%s|%s|%v|%s", e.Principal, e.Timestamp.UTC().Format("2006-01-02T15:04:05.000000000Z07:00"), e.Items, e.SessionRef)
return hex.EncodeToString(h.Sum(nil))[:16]
}
+96
View File
@@ -0,0 +1,96 @@
package audit
import (
"context"
"fmt"
"github.com/mathiasbq/hyperguild/ingestion/internal/capture"
"github.com/mathiasbq/hyperguild/ingestion/internal/classification"
)
// Central is the central audit substrate (loki). Ready is a cheap
// reachability probe used by the pre-write reserve; Push writes a record.
type Central interface {
Ready(ctx context.Context) error
Push(ctx context.Context, e capture.AuditEntry) error
}
// Buffer is the durable local fallback for internal/public-tier records
// when the central sink is unreachable. It must survive process restart.
type Buffer interface {
// Writable reports whether the buffer can currently be appended to.
Writable() error
Append(e capture.AuditEntry) error
// Pending returns buffered records awaiting reconciliation, each with a
// stable ID used to Confirm (delete) it after a confirmed central write.
Pending() ([]Buffered, error)
Confirm(id string) error
}
// Buffered is a buffered audit record plus its stable buffer ID.
type Buffered struct {
ID string
Entry capture.AuditEntry
}
// Notifier raises an out-of-band alert (ntfy) about a degraded state.
type Notifier interface {
Notify(ctx context.Context, msg string) error
}
// DegradingSink is the classification-aware AuditSink (§4.4):
//
// - central reachable → AuditCentral (all tiers).
// - central down + confidential → refuse (no buffer): confidential must
// be centrally auditable at write time.
// - central down + internal/public + buffer writable → AuditBuffered.
// - central down + (confidential, or buffer not writable) → refuse (floor).
//
// The decision is made in Reserve, before any write; Record then executes it.
type DegradingSink struct {
central Central
buffer Buffer
notifier Notifier
}
// NewDegradingSink wires the central sink, durable buffer, and notifier.
func NewDegradingSink(central Central, buffer Buffer, notifier Notifier) *DegradingSink {
return &DegradingSink{central: central, buffer: buffer, notifier: notifier}
}
// Reserve decides, before any write, how the capture will be audited — or
// returns an error to refuse it.
func (d *DegradingSink) Reserve(ctx context.Context, level classification.Level) (capture.AuditOutcome, error) {
if err := d.central.Ready(ctx); err == nil {
return capture.AuditCentral, nil
}
// Central sink is down.
if level == classification.Confidential {
return 0, fmt.Errorf("confidential capture requires the central audit sink, which is unreachable")
}
if err := d.buffer.Writable(); err != nil {
// Floor: neither central nor local buffer can record the audit.
return 0, fmt.Errorf("audit floor: central sink down and local buffer unwritable: %w", err)
}
return capture.AuditBuffered, nil
}
// Record persists the entry per the reserved outcome. For AuditBuffered it
// also fires the degraded-state alert.
func (d *DegradingSink) Record(ctx context.Context, e capture.AuditEntry, outcome capture.AuditOutcome) error {
switch outcome {
case capture.AuditBuffered:
if err := d.buffer.Append(e); err != nil {
return fmt.Errorf("buffer audit record: %w", err)
}
// Best-effort alert; the record is already durably buffered.
if d.notifier != nil {
_ = d.notifier.Notify(ctx, fmt.Sprintf(
"capture audit BUFFERED LOCALLY (loki unreachable) — principal=%s class=%s items=%d",
e.Principal, e.EffectiveClassification, len(e.Items)))
}
return nil
default:
return d.central.Push(ctx, e)
}
}
+191
View File
@@ -0,0 +1,191 @@
package audit_test
import (
"context"
"errors"
"path/filepath"
"testing"
"github.com/mathiasbq/hyperguild/ingestion/internal/audit"
"github.com/mathiasbq/hyperguild/ingestion/internal/capture"
"github.com/mathiasbq/hyperguild/ingestion/internal/classification"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
// --- fakes ---
type fakeCentral struct {
down bool
pushed []capture.AuditEntry
pushErr error
}
func (f *fakeCentral) Ready(context.Context) error {
if f.down {
return errors.New("loki down")
}
return nil
}
func (f *fakeCentral) Push(_ context.Context, e capture.AuditEntry) error {
if f.pushErr != nil {
return f.pushErr
}
f.pushed = append(f.pushed, e)
return nil
}
type fakeNotifier struct{ msgs []string }
func (f *fakeNotifier) Notify(_ context.Context, msg string) error {
f.msgs = append(f.msgs, msg)
return nil
}
// unwritableBuffer always reports it cannot be written (floor condition).
type unwritableBuffer struct{}
func (unwritableBuffer) Writable() error { return errors.New("disk full") }
func (unwritableBuffer) Append(capture.AuditEntry) error { return errors.New("disk full") }
func (unwritableBuffer) Pending() ([]audit.Buffered, error) { return nil, nil }
func (unwritableBuffer) Confirm(string) error { return nil }
func newFileBuffer(t *testing.T) *audit.FileBuffer {
t.Helper()
b, err := audit.NewFileBuffer(filepath.Join(t.TempDir(), "audit-buffer.jsonl"))
require.NoError(t, err)
return b
}
func entry(principal string) capture.AuditEntry {
return capture.AuditEntry{Principal: principal, EffectiveClassification: "internal", Items: []string{"insight:x"}}
}
// --- Reserve: classification-aware decision ---
func TestReserveCentralUpGrantsCentral(t *testing.T) {
d := audit.NewDegradingSink(&fakeCentral{}, newFileBuffer(t), &fakeNotifier{})
for _, lvl := range []classification.Level{classification.Public, classification.Internal, classification.Confidential} {
out, err := d.Reserve(context.Background(), lvl)
require.NoError(t, err)
assert.Equal(t, capture.AuditCentral, out)
}
}
func TestReserveConfidentialSinkDownRefuses(t *testing.T) {
d := audit.NewDegradingSink(&fakeCentral{down: true}, newFileBuffer(t), &fakeNotifier{})
_, err := d.Reserve(context.Background(), classification.Confidential)
require.Error(t, err, "confidential + sink down → refuse, no buffer")
}
func TestReserveInternalSinkDownBuffers(t *testing.T) {
d := audit.NewDegradingSink(&fakeCentral{down: true}, newFileBuffer(t), &fakeNotifier{})
out, err := d.Reserve(context.Background(), classification.Internal)
require.NoError(t, err)
assert.Equal(t, capture.AuditBuffered, out)
}
func TestReserveFloorRefusesWhenNothingCanRecord(t *testing.T) {
d := audit.NewDegradingSink(&fakeCentral{down: true}, unwritableBuffer{}, &fakeNotifier{})
_, err := d.Reserve(context.Background(), classification.Internal)
require.Error(t, err, "central down AND buffer unwritable → floor refuse")
}
// --- Record: executes the reserved outcome ---
func TestRecordCentralPushes(t *testing.T) {
c := &fakeCentral{}
d := audit.NewDegradingSink(c, newFileBuffer(t), &fakeNotifier{})
require.NoError(t, d.Record(context.Background(), entry("p"), capture.AuditCentral))
assert.Len(t, c.pushed, 1)
}
func TestRecordBufferedAppendsAndNotifies(t *testing.T) {
buf := newFileBuffer(t)
nt := &fakeNotifier{}
d := audit.NewDegradingSink(&fakeCentral{down: true}, buf, nt)
require.NoError(t, d.Record(context.Background(), entry("p"), capture.AuditBuffered))
pending, err := buf.Pending()
require.NoError(t, err)
assert.Len(t, pending, 1)
assert.NotEmpty(t, nt.msgs, "degraded state alerts via ntfy")
}
// --- FileBuffer durability + Confirm ---
func TestFileBufferSurvivesRestart(t *testing.T) {
path := filepath.Join(t.TempDir(), "buf.jsonl")
b1, err := audit.NewFileBuffer(path)
require.NoError(t, err)
require.NoError(t, b1.Append(entry("p1")))
require.NoError(t, b1.Append(entry("p2")))
// "restart": a fresh FileBuffer over the same file sees the records.
b2, err := audit.NewFileBuffer(path)
require.NoError(t, err)
pending, err := b2.Pending()
require.NoError(t, err)
assert.Len(t, pending, 2)
}
func TestFileBufferConfirmRemovesOnlyThatRecord(t *testing.T) {
buf := newFileBuffer(t)
require.NoError(t, buf.Append(entry("keep")))
require.NoError(t, buf.Append(entry("drop")))
pending, _ := buf.Pending()
require.Len(t, pending, 2)
var dropID string
for _, p := range pending {
if p.Entry.Principal == "drop" {
dropID = p.ID
}
}
require.NoError(t, buf.Confirm(dropID))
after, _ := buf.Pending()
require.Len(t, after, 1)
assert.Equal(t, "keep", after[0].Entry.Principal)
}
// --- Reconcile ---
func TestReconcileReplaysAndClearsOnlyAfterConfirmedWrite(t *testing.T) {
buf := newFileBuffer(t)
require.NoError(t, buf.Append(entry("a")))
require.NoError(t, buf.Append(entry("b")))
c := &fakeCentral{} // up
nt := &fakeNotifier{}
n, err := audit.Reconcile(context.Background(), c, buf, nt)
require.NoError(t, err)
assert.Equal(t, 2, n)
assert.Len(t, c.pushed, 2, "buffered records replayed to central")
pending, _ := buf.Pending()
assert.Empty(t, pending, "buffer cleared after confirmed central writes")
}
func TestReconcileNoopWhenCentralDown(t *testing.T) {
buf := newFileBuffer(t)
require.NoError(t, buf.Append(entry("a")))
n, err := audit.Reconcile(context.Background(), &fakeCentral{down: true}, buf, &fakeNotifier{})
require.NoError(t, err)
assert.Equal(t, 0, n)
pending, _ := buf.Pending()
assert.Len(t, pending, 1, "records stay buffered while central is down")
}
func TestReconcileKeepsRecordWhenPushFails(t *testing.T) {
buf := newFileBuffer(t)
require.NoError(t, buf.Append(entry("a")))
// Ready ok but Push fails → record must remain buffered (not lost).
c := &fakeCentral{pushErr: errors.New("push rejected")}
n, err := audit.Reconcile(context.Background(), c, buf, &fakeNotifier{})
require.NoError(t, err)
assert.Equal(t, 0, n)
pending, _ := buf.Pending()
assert.Len(t, pending, 1)
}
+99
View File
@@ -0,0 +1,99 @@
package audit
import (
"bytes"
"context"
"encoding/json"
"fmt"
"net/http"
"strconv"
"strings"
"time"
"github.com/mathiasbq/hyperguild/ingestion/internal/capture"
)
// LokiCentral pushes capture audit records to a Grafana Loki instance via
// its push API, and probes readiness via /ready. It is the central audit
// substrate behind DegradingSink.
type LokiCentral struct {
baseURL string
labels map[string]string
http *http.Client
}
// NewLokiCentral constructs a LokiCentral for the given base URL (e.g.
// http://loki:3100). Returns nil when baseURL is empty so callers can
// treat missing config as "no central sink" with a single nil check.
func NewLokiCentral(baseURL string) *LokiCentral {
if baseURL == "" {
return nil
}
return &LokiCentral{
baseURL: strings.TrimRight(baseURL, "/"),
labels: map[string]string{"service": "brain-capture", "kind": "audit"},
http: &http.Client{Timeout: 10 * time.Second},
}
}
// Ready probes Loki's readiness endpoint.
func (l *LokiCentral) Ready(ctx context.Context) error {
req, err := http.NewRequestWithContext(ctx, http.MethodGet, l.baseURL+"/ready", nil)
if err != nil {
return err
}
resp, err := l.http.Do(req)
if err != nil {
return fmt.Errorf("loki not ready: %w", err)
}
defer func() { _ = resp.Body.Close() }()
if resp.StatusCode != http.StatusOK {
return fmt.Errorf("loki not ready: status %d", resp.StatusCode)
}
return nil
}
// pushPayload is the Loki push API body: one stream, one entry whose line
// is the JSON-encoded audit record.
type pushPayload struct {
Streams []lokiStream `json:"streams"`
}
type lokiStream struct {
Stream map[string]string `json:"stream"`
Values [][2]string `json:"values"`
}
// Push writes one audit record to Loki as a structured log line.
func (l *LokiCentral) Push(ctx context.Context, e capture.AuditEntry) error {
line, err := json.Marshal(e)
if err != nil {
return fmt.Errorf("marshal audit entry: %w", err)
}
ts := e.Timestamp
if ts.IsZero() {
ts = time.Now()
}
body, err := json.Marshal(pushPayload{Streams: []lokiStream{{
Stream: l.labels,
Values: [][2]string{{strconv.FormatInt(ts.UTC().UnixNano(), 10), string(line)}},
}}})
if err != nil {
return err
}
req, err := http.NewRequestWithContext(ctx, http.MethodPost,
l.baseURL+"/loki/api/v1/push", bytes.NewReader(body))
if err != nil {
return err
}
req.Header.Set("Content-Type", "application/json")
resp, err := l.http.Do(req)
if err != nil {
return fmt.Errorf("loki push: %w", err)
}
defer func() { _ = resp.Body.Close() }()
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
return fmt.Errorf("loki push: status %d", resp.StatusCode)
}
return nil
}
@@ -0,0 +1,86 @@
package audit_test
import (
"context"
"io"
"net/http"
"net/http/httptest"
"strings"
"testing"
"github.com/mathiasbq/hyperguild/ingestion/internal/audit"
"github.com/mathiasbq/hyperguild/ingestion/internal/capture"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func TestLokiReadyAndPush(t *testing.T) {
var pushBody string
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
switch r.URL.Path {
case "/ready":
w.WriteHeader(http.StatusOK)
case "/loki/api/v1/push":
b, _ := io.ReadAll(r.Body)
pushBody = string(b)
w.WriteHeader(http.StatusNoContent)
default:
w.WriteHeader(http.StatusNotFound)
}
}))
defer srv.Close()
c := audit.NewLokiCentral(srv.URL)
require.NotNil(t, c)
require.NoError(t, c.Ready(context.Background()))
err := c.Push(context.Background(), capture.AuditEntry{
Principal: "koala-cli", EffectiveClassification: "internal", Items: []string{"insight:x"},
})
require.NoError(t, err)
assert.Contains(t, pushBody, "streams")
assert.Contains(t, pushBody, "koala-cli", "audit entry serialised into the loki line")
}
func TestLokiReadyFailsWhenDown(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusServiceUnavailable)
}))
defer srv.Close()
require.Error(t, audit.NewLokiCentral(srv.URL).Ready(context.Background()))
}
func TestLokiNilWhenUnconfigured(t *testing.T) {
assert.Nil(t, audit.NewLokiCentral(""))
}
func TestNtfyNotify(t *testing.T) {
var gotBody, gotAuth string
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
b, _ := io.ReadAll(r.Body)
gotBody = string(b)
gotAuth = r.Header.Get("Authorization")
w.WriteHeader(http.StatusOK)
}))
defer srv.Close()
n := audit.NewNtfyNotifier(srv.URL, "ntfy-token")
require.NotNil(t, n)
require.NoError(t, n.Notify(context.Background(), "audit buffered locally"))
assert.Contains(t, gotBody, "audit buffered locally")
assert.Equal(t, "Bearer ntfy-token", gotAuth)
}
func TestNtfyDoesNotLeakTokenOnError(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusInternalServerError)
}))
defer srv.Close()
err := audit.NewNtfyNotifier(srv.URL, "secret-token").Notify(context.Background(), "x")
require.Error(t, err)
assert.False(t, strings.Contains(err.Error(), "secret-token"), "token must not leak into errors")
}
func TestNtfyNilWhenUnconfigured(t *testing.T) {
assert.Nil(t, audit.NewNtfyNotifier("", "tok"))
}
+55
View File
@@ -0,0 +1,55 @@
package audit
import (
"context"
"fmt"
"net/http"
"strings"
"time"
)
// NtfyNotifier posts alerts to an ntfy topic URL. Used to surface a
// degraded audit state (records buffered locally during a loki outage).
type NtfyNotifier struct {
topicURL string
token string
http *http.Client
}
// NewNtfyNotifier constructs a notifier for the given ntfy topic URL
// (e.g. https://ntfy.sh/my-topic). token is an optional bearer for
// protected ntfy instances; it is held here and only sent in the
// Authorization header, never logged. Returns nil when topicURL is empty.
func NewNtfyNotifier(topicURL, token string) *NtfyNotifier {
if topicURL == "" {
return nil
}
return &NtfyNotifier{
topicURL: strings.TrimRight(topicURL, "/"),
token: token,
http: &http.Client{Timeout: 10 * time.Second},
}
}
// Notify posts a message to the ntfy topic.
func (n *NtfyNotifier) Notify(ctx context.Context, msg string) error {
req, err := http.NewRequestWithContext(ctx, http.MethodPost, n.topicURL, strings.NewReader(msg))
if err != nil {
return err
}
req.Header.Set("Title", "brain-capture audit degraded")
req.Header.Set("Priority", "high")
req.Header.Set("Tags", "warning,brain")
if n.token != "" {
req.Header.Set("Authorization", "Bearer "+n.token)
}
resp, err := n.http.Do(req)
if err != nil {
return fmt.Errorf("ntfy notify: %w", err)
}
defer func() { _ = resp.Body.Close() }()
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
return fmt.Errorf("ntfy notify: status %d", resp.StatusCode)
}
return nil
}
+65
View File
@@ -0,0 +1,65 @@
package audit
import (
"context"
"fmt"
"log/slog"
"time"
)
// Reconcile replays locally-buffered audit records to the central sink
// when it is reachable again. A record is removed from the buffer ONLY
// after its central write is confirmed, so a crash mid-reconcile re-plays
// rather than loses. Returns the number of records reconciled.
//
// A no-op (0, nil) when the central sink is still unreachable or the
// buffer is empty.
func Reconcile(ctx context.Context, central Central, buffer Buffer, notifier Notifier) (int, error) {
if err := central.Ready(ctx); err != nil {
return 0, nil // still down; try again next tick
}
pending, err := buffer.Pending()
if err != nil {
return 0, fmt.Errorf("read buffer: %w", err)
}
reconciled := 0
for _, rec := range pending {
if err := central.Push(ctx, rec.Entry); err != nil {
// Central went away mid-drain; stop and keep the rest buffered.
break
}
if err := buffer.Confirm(rec.ID); err != nil {
return reconciled, fmt.Errorf("confirm buffered record %s: %w", rec.ID, err)
}
reconciled++
}
if reconciled > 0 && notifier != nil {
_ = notifier.Notify(ctx, fmt.Sprintf("reconciled %d buffered capture audit record(s) to loki", reconciled))
}
return reconciled, nil
}
// StartReconcile runs Reconcile on a ticker until ctx is cancelled. It is
// the recovery half of the degrade-and-buffer path; pair it with a
// DegradingSink sharing the same buffer + central.
func StartReconcile(ctx context.Context, central Central, buffer Buffer, notifier Notifier, interval time.Duration) {
if interval <= 0 {
interval = time.Minute
}
go func() {
t := time.NewTicker(interval)
defer t.Stop()
for {
select {
case <-ctx.Done():
return
case <-t.C:
if n, err := Reconcile(ctx, central, buffer, notifier); err != nil {
slog.Warn("audit reconcile failed", "err", err)
} else if n > 0 {
slog.Info("audit reconcile", "reconciled", n)
}
}
}
}()
}
+56
View File
@@ -0,0 +1,56 @@
// Package audit provides AuditSink implementations for the capture
// capability (I5). This file ships the minimal slog-backed sink used in
// #53: it emits the request-level audit record to structured logs, which
// the alloy/loki substrate already scrapes. The classification-aware
// degradation/refusal sink (confidential fails closed, internal buffers +
// reconciles) lands in #54 and replaces this behind the same interface.
package audit
import (
"context"
"log/slog"
"github.com/mathiasbq/hyperguild/ingestion/internal/capture"
"github.com/mathiasbq/hyperguild/ingestion/internal/classification"
)
// SlogSink records audit entries to an slog.Logger. It never fails and is
// always centrally available, so its Reserve always grants AuditCentral —
// it does not exercise the I5 degradation/floor. That is DegradingSink's
// job (loki + durable buffer). SlogSink is the default for deployments
// without a loki endpoint configured. A nil logger ⇒ slog.Default().
type SlogSink struct {
logger *slog.Logger
}
// NewSlogSink constructs a SlogSink. nil logger ⇒ slog.Default().
func NewSlogSink(logger *slog.Logger) *SlogSink {
if logger == nil {
logger = slog.Default()
}
return &SlogSink{logger: logger}
}
// Reserve always grants central recording — slog is always available.
func (s *SlogSink) Reserve(_ context.Context, _ classification.Level) (capture.AuditOutcome, error) {
return capture.AuditCentral, nil
}
// Record emits the audit entry at info level. Security events, when
// present, are logged at warn level so they surface independently of the
// routine audit stream.
func (s *SlogSink) Record(_ context.Context, e capture.AuditEntry, _ capture.AuditOutcome) error {
s.logger.Info("capture audit",
"principal", e.Principal,
"actor", e.Actor,
"harness", e.Harness,
"session_ref", e.SessionRef,
"classification", e.EffectiveClassification,
"items", e.Items,
"ts", e.Timestamp,
)
for _, ev := range e.SecurityEvents {
s.logger.Warn("capture security event", "principal", e.Principal, "event", ev)
}
return nil
}
+41
View File
@@ -0,0 +1,41 @@
package audit_test
import (
"bytes"
"context"
"log/slog"
"testing"
"github.com/mathiasbq/hyperguild/ingestion/internal/audit"
"github.com/mathiasbq/hyperguild/ingestion/internal/capture"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func TestSlogSinkRecordsEntryAndSecurityEvents(t *testing.T) {
var buf bytes.Buffer
sink := audit.NewSlogSink(slog.New(slog.NewTextHandler(&buf, nil)))
err := sink.Record(context.Background(), capture.AuditEntry{
Principal: "koala-cli",
Harness: "claude-code",
EffectiveClassification: "confidential",
Items: []string{"insight:wiki/a/facts/x.md"},
SecurityEvents: []string{"asserted-vs-derived origin mismatch"},
}, capture.AuditCentral)
require.NoError(t, err)
out := buf.String()
assert.Contains(t, out, "capture audit")
assert.Contains(t, out, "koala-cli")
assert.Contains(t, out, "confidential")
assert.Contains(t, out, "capture security event")
assert.Contains(t, out, "asserted-vs-derived origin mismatch")
}
func TestSlogSinkNilLoggerDefaults(t *testing.T) {
// nil logger must not panic.
require.NotPanics(t, func() {
_ = audit.NewSlogSink(nil).Record(context.Background(), capture.AuditEntry{}, capture.AuditCentral)
})
}
+129
View File
@@ -0,0 +1,129 @@
// Package brainstore is the concrete BrainStore: the single shared
// implementation of the #45 write/update/get verbs, used by BOTH the MCP
// handlers and the capture use-case so there is one implementation, not
// two (the Clean-Architecture / DRY payoff of #51).
//
// It composes the file-level primitives in package api (WriteNote,
// UpdateNote, ReadNote — the read-after-write contract) with the wiki
// upkeep that must accompany a write: wing _index rebuild, cross-wing
// auto-tunnel, and graph re-index. Embedding refresh is intentionally
// out-of-band (mtime-driven vectorstore.Sync) and not triggered here —
// see the brain note on out-of-band sync.
package brainstore
import (
"context"
"log/slog"
"strings"
"github.com/mathiasbq/hyperguild/ingestion/internal/api"
"github.com/mathiasbq/hyperguild/ingestion/internal/brain"
"github.com/mathiasbq/hyperguild/ingestion/internal/capture"
"github.com/mathiasbq/hyperguild/ingestion/internal/graphsync"
)
// Store implements capture.BrainStore against a brain directory on disk,
// optionally re-indexing each write into the knowledge graph.
type Store struct {
brainDir string
graph graphsync.Store // nil = graph re-index disabled
}
// New constructs a Store bound to brainDir with graph indexing disabled.
func New(brainDir string) *Store {
return &Store{brainDir: brainDir}
}
// WithGraph enables graph re-index on every write/update. nil disables it.
func (s *Store) WithGraph(g graphsync.Store) *Store {
s.graph = g
return s
}
// Write creates a brain note and returns its read-after-write handle.
func (s *Store) Write(ctx context.Context, n capture.Note) (capture.Ref, error) {
relPath, err := api.WriteNote(s.brainDir, api.WriteNoteOptions{
Content: n.Content,
Filename: n.Filename,
Type: n.Type,
Domain: n.Domain,
Wing: n.Wing,
Hall: n.Hall,
})
if err != nil {
return capture.Ref{}, err
}
s.wikiUpkeep(relPath, n.Wing, n.Content)
s.indexInGraph(ctx, "brain_write", relPath)
_, _, hash, _ := api.ReadNote(s.brainDir, relPath)
return capture.Ref{ID: relPath, Path: relPath, ContentHash: hash}, nil
}
// Update supersedes an existing note in place. slug may be a bare slug
// (resolved against n.Wing/n.Hall) or a full brain-relative path (when it
// contains a slash). It never creates — a missing target is an error.
func (s *Store) Update(ctx context.Context, slug string, n capture.Note) (capture.Ref, error) {
opts := api.UpdateNoteOptions{Content: n.Content, Reason: n.Reason}
if strings.Contains(slug, "/") {
opts.Path = slug
} else {
opts.Wing, opts.Hall, opts.Slug = n.Wing, n.Hall, slug
}
relPath, hash, _, err := api.UpdateNote(s.brainDir, opts)
if err != nil {
return capture.Ref{}, err
}
if wing := wingFromRelPath(relPath); wing != "" {
s.wikiUpkeep(relPath, wing, n.Content)
}
s.indexInGraph(ctx, "brain_update", relPath)
return capture.Ref{ID: relPath, Path: relPath, ContentHash: hash, Superseded: true}, nil
}
// Get fetches a note by id/path — the read-after-write confirmation
// primitive (a direct fetch, never a semantic query).
func (s *Store) Get(_ context.Context, id string) (capture.StoredNote, error) {
fm, body, hash, err := api.ReadNote(s.brainDir, id)
if err != nil {
return capture.StoredNote{}, err
}
return capture.StoredNote{ID: id, Path: id, ContentHash: hash, Frontmatter: fm, Body: body}, nil
}
// wikiUpkeep rebuilds the wing _index and re-tunnels cross-wing matches
// when a note lands in the structured wiki. Both are best-effort: the
// note is already written, so a failure here is logged, not propagated.
func (s *Store) wikiUpkeep(relPath, wing, content string) {
if wing == "" {
return
}
if err := brain.BuildWingIndex(s.brainDir, wing); err != nil {
slog.Warn("brainstore: auto-index failed", "wing", wing, "err", err)
}
if err := brain.AutoTunnel(s.brainDir, relPath, content); err != nil {
slog.Warn("brainstore: auto-tunnel failed", "src", relPath, "err", err)
}
}
// indexInGraph re-indexes a written doc into the graph, best-effort.
func (s *Store) indexInGraph(ctx context.Context, op, relPath string) {
if s.graph == nil || relPath == "" {
return
}
if err := graphsync.IndexDoc(ctx, s.graph, s.brainDir, relPath); err != nil {
slog.Warn(op+": graph index failed", "path", relPath, "err", err)
}
}
// wingFromRelPath extracts the wing from a structured wiki path
// (wiki/<wing>/<hall>/<slug>.md). Returns "" for legacy/non-wiki paths.
func wingFromRelPath(relPath string) string {
parts := strings.Split(relPath, "/")
if len(parts) >= 4 && parts[0] == "wiki" {
return parts[1]
}
return ""
}
@@ -0,0 +1,85 @@
package brainstore_test
import (
"context"
"os"
"path/filepath"
"testing"
"github.com/mathiasbq/hyperguild/ingestion/internal/brainstore"
"github.com/mathiasbq/hyperguild/ingestion/internal/capture"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func TestStoreWriteReturnsHandle(t *testing.T) {
dir := t.TempDir()
s := brainstore.New(dir)
ref, err := s.Write(context.Background(), capture.Note{
Content: "# X\n\nbody\n", Filename: "x", Wing: "a", Hall: "facts",
})
require.NoError(t, err)
assert.Equal(t, "wiki/a/facts/x.md", ref.Path)
assert.Equal(t, ref.Path, ref.ID)
assert.NotEmpty(t, ref.ContentHash)
assert.False(t, ref.Superseded)
_, err = os.Stat(filepath.Join(dir, "wiki/a/facts/x.md"))
require.NoError(t, err)
}
func TestStoreUpdateSupersedes(t *testing.T) {
dir := t.TempDir()
s := brainstore.New(dir)
_, err := s.Write(context.Background(), capture.Note{
Content: "old\n", Filename: "n", Wing: "a", Hall: "facts",
})
require.NoError(t, err)
ref, err := s.Update(context.Background(), "n", capture.Note{
Content: "new\n", Wing: "a", Hall: "facts", Reason: "changed",
})
require.NoError(t, err)
assert.True(t, ref.Superseded)
assert.Equal(t, "wiki/a/facts/n.md", ref.Path)
got, _ := os.ReadFile(filepath.Join(dir, "wiki/a/facts/n.md"))
assert.Contains(t, string(got), "new")
assert.Contains(t, string(got), "supersede_reason: changed")
}
func TestStoreUpdateByFullPath(t *testing.T) {
dir := t.TempDir()
s := brainstore.New(dir)
_, err := s.Write(context.Background(), capture.Note{Content: "old\n", Filename: "n", Wing: "a", Hall: "facts"})
require.NoError(t, err)
ref, err := s.Update(context.Background(), "wiki/a/facts/n.md", capture.Note{Content: "fresh\n"})
require.NoError(t, err)
assert.Equal(t, "wiki/a/facts/n.md", ref.Path)
}
func TestStoreUpdateMissingErrors(t *testing.T) {
dir := t.TempDir()
s := brainstore.New(dir)
_, err := s.Update(context.Background(), "ghost", capture.Note{Content: "x\n", Wing: "a", Hall: "facts"})
require.Error(t, err)
_, statErr := os.Stat(filepath.Join(dir, "wiki/a/facts/ghost.md"))
assert.True(t, os.IsNotExist(statErr), "update must not create")
}
func TestStoreGetRoundTripsHash(t *testing.T) {
dir := t.TempDir()
s := brainstore.New(dir)
ref, err := s.Write(context.Background(), capture.Note{
Content: "# Body\n\ntext\n", Filename: "n", Wing: "a", Hall: "facts",
})
require.NoError(t, err)
note, err := s.Get(context.Background(), ref.ID)
require.NoError(t, err)
assert.Equal(t, ref.ContentHash, note.ContentHash, "write→get hash round-trips")
assert.Equal(t, "a", note.Frontmatter["wing"])
assert.Contains(t, note.Body, "# Body")
}
+148
View File
@@ -0,0 +1,148 @@
// Package capture is the Clean-Architecture use-case for the uniform
// capture capability (issue #49/#51): persist a finished session's
// valuable output — insights → brain, action items → Gitea tickets,
// optional summary → ai-sessions — with one invocation, identical core
// behaviour across every harness.
//
// This package is pure orchestration. It depends only on ports
// (interfaces) and plain entities — no HTTP, no live Gitea, no embedding
// or audit I/O. The real adapters are wired in #52 (Gitea tracker), #53
// (REST + I1 origin gate), and #54/#55 (audit path + relay). The I1
// sovereignty refusal and the classification-aware audit degradation are
// deliberately NOT here — those need the server-derived principal origin
// (#53) and the loki/buffer machinery (#54). What lives here is everything
// testable against fakes: validation, effective-classification resolution
// (stricter wins), best-effort orchestration, and the partial receipt.
package capture
// Zone is the trust zone a capture originates from, server-derived from
// the authenticated principal (spec §4.2 / I1). It is NEVER taken from
// caller input — context.Harness is descriptive telemetry only.
type Zone int
const (
// ZoneUnknown means the origin was not set. The REST adapter always
// sets a concrete zone; the service treats Unknown as "not gated" (only
// an explicit ZoneUSNexus triggers the I1 refusal) so the gate can
// never fire on a caller-controllable default.
ZoneUnknown Zone = iota
// ZoneSovereign is sovereign soil (homelab / Tailscale CLI callers).
ZoneSovereign
// ZoneUSNexus is a non-sovereign US-jurisdiction surface (e.g.
// claude.ai). Confidential captures through it are refused (I1).
ZoneUSNexus
)
// String renders the zone for audit/refusal messages.
func (z Zone) String() string {
switch z {
case ZoneSovereign:
return "sovereign-soil"
case ZoneUSNexus:
return "us-nexus"
default:
return "unknown"
}
}
// CaptureContext is the per-session metadata accompanying a capture.
//
// Classification is the caller-declared sensitivity (model C, spec §4.1):
// the server independently derives the target's classification and gates
// on the stricter of the two. Principal and Origin are server-derived from
// the authenticated identity (the REST adapter populates them); they are
// never caller-asserted. Harness is descriptive telemetry only — never a
// gate input.
type CaptureContext struct {
Harness string
SessionRef string
Fidelity string
Actor string
Classification string // caller-declared level token ("" = unspecified)
Principal string // server-derived (auth); audit identity
Origin Zone // server-derived trust zone; the I1 gate input
}
// Insight is one piece of session knowledge bound for the brain. A
// non-empty SupersedeSlug routes to Update (revise in place); otherwise
// Write (create).
type Insight struct {
Text string
Wing string
Hall string
SupersedeSlug string
}
// Ticket is one action item bound for a Gitea repo. Owner is always the
// operator (set by the tracker adapter), never carried here.
type Ticket struct {
Repo string
Action string // create | close | comment
Number int // required for close/comment
Title string // required for create
Body string
}
// Summary is an optional session summary bound for ai-sessions.
type Summary struct {
Title string
Body string
ReposTouched []string
}
// CaptureInput is the whole capture request.
type CaptureInput struct {
Context CaptureContext
Insights []Insight
Tickets []Ticket
Summary *Summary
DryRun bool
}
// InsightResult is the per-insight outcome in the receipt.
type InsightResult struct {
ID string `json:"id,omitempty"`
Path string `json:"path,omitempty"`
ContentHash string `json:"content_hash,omitempty"`
Superseded bool `json:"superseded"`
OK bool `json:"ok"`
}
// TicketResult is the per-ticket outcome in the receipt.
type TicketResult struct {
Repo string `json:"repo"`
Number int `json:"number,omitempty"`
Action string `json:"action"`
URL string `json:"url,omitempty"`
OK bool `json:"ok"`
}
// SummaryResult is the summary outcome in the receipt.
type SummaryResult struct {
Path string `json:"path,omitempty"`
OK bool `json:"ok"`
}
// ItemError pins a failure to a specific request item for the partial
// receipt. Item is a stable locator like "insight[1]" or "ticket[0]".
type ItemError struct {
Item string `json:"item"`
Error string `json:"error"`
}
// CaptureReceipt is the structured, partial-aware result. Per-item ok
// flags plus a flat Errors list make partial success explicit; the
// caller never has to infer what landed.
type CaptureReceipt struct {
Insights []InsightResult `json:"insights"`
Tickets []TicketResult `json:"tickets"`
Summary *SummaryResult `json:"summary,omitempty"`
Errors []ItemError `json:"errors"`
EffectiveClassification string `json:"effective_classification,omitempty"`
DryRun bool `json:"dry_run"`
// AuditBuffered is true when the central audit sink was unreachable and
// this capture's audit record was written to the durable local buffer
// instead (internal/public tier). Surfaces the degraded state to the
// caller per §4.4.
AuditBuffered bool `json:"audit_buffered,omitempty"`
}
+126
View File
@@ -0,0 +1,126 @@
package capture
import (
"context"
"time"
"github.com/mathiasbq/hyperguild/ingestion/internal/classification"
)
// Ref is the read-after-write handle returned by a brain write/update —
// the #45 contract. ContentHash lets the caller confirm what landed
// without a re-query; for an Update, Superseded is true.
type Ref struct {
ID string
Path string
ContentHash string
Superseded bool
}
// StoredNote is a brain note fetched by Get: the read-after-write
// confirmation primitive (a direct fetch, never a semantic query).
type StoredNote struct {
ID string
Path string
ContentHash string
Frontmatter map[string]string
Body string
}
// Note is the brain-write payload. It carries both the wing/hall taxonomy
// and the legacy type/domain fields so a single BrainStore serves both
// capture insights and the existing MCP brain_write surface. Reason is
// the supersede rationale, used only by Update.
type Note struct {
Content string
Filename string
Wing string
Hall string
Type string
Domain string
Reason string
}
// BrainStore is the brain persistence port — the shared implementation of
// the #45 write/update/get verbs that both the MCP handlers and capture
// call, so there is one implementation, not two. The read-after-write +
// staleness discipline lives behind this interface so no caller carries
// the rule.
type BrainStore interface {
Write(ctx context.Context, n Note) (Ref, error)
Update(ctx context.Context, slug string, n Note) (Ref, error)
Get(ctx context.Context, id string) (StoredNote, error)
}
// IssueRef identifies a ticket touched by the tracker.
type IssueRef struct {
Repo string
Number int
URL string
}
// IssueTracker is the Gitea ticket port. The implementation (#52) always
// scopes to owner "mathias"; the port deliberately omits owner.
type IssueTracker interface {
CreateIssue(ctx context.Context, repo, title, body string) (IssueRef, error)
// CloseIssue closes an issue, optionally posting a closing comment
// first (empty comment ⇒ close only).
CloseIssue(ctx context.Context, repo string, number int, comment string) (IssueRef, error)
CommentIssue(ctx context.Context, repo string, number int, body string) (IssueRef, error)
}
// SummaryWriter is the ai-sessions summary port.
type SummaryWriter interface {
WriteFile(ctx context.Context, repo, path, content string) error
}
// ClassificationPolicy derives a target's sensitivity (model C). The
// "stricter wins" combination of declared vs derived is use-case policy
// and lives in the service, so the port stays minimal. Satisfied by
// classification.Config (#50).
type ClassificationPolicy interface {
Derive(target classification.Target) classification.Level
}
// AuditEntry is the request-level audit record (I5): who/what captured
// what, when, via which principal. SecurityEvents carries anomalies such
// as a caller under-declaring sensitivity relative to the target floor.
type AuditEntry struct {
Timestamp time.Time
Principal string
Actor string
Harness string
SessionRef string
EffectiveClassification string
Items []string
SecurityEvents []string
}
// AuditOutcome is how a capture's audit record was (or will be) persisted.
type AuditOutcome int
const (
// AuditCentral means the record goes to the central sink (loki).
AuditCentral AuditOutcome = iota
// AuditBuffered means the central sink was unreachable and the record
// is written to a durable local buffer for later reconciliation
// (internal/public tier only).
AuditBuffered
)
// AuditSink is the two-phase, classification-aware audit port (I5, §4.4).
//
// Reserve runs BEFORE any write and decides whether the capture can be
// audited at its effective classification: it returns the outcome to use,
// or an error to refuse the capture before anything is written
// (confidential + central sink down → refuse; the all-tiers floor when
// nothing can record → refuse). Record runs AFTER the writes and persists
// the final entry per the reserved outcome.
//
// Splitting reserve from record is what lets "confidential + sink-down →
// refuse before any write" be literally true while the record itself
// (which lists what landed) is necessarily written afterwards.
type AuditSink interface {
Reserve(ctx context.Context, level classification.Level) (AuditOutcome, error)
Record(ctx context.Context, e AuditEntry, outcome AuditOutcome) error
}
+388
View File
@@ -0,0 +1,388 @@
package capture
import (
"context"
"crypto/sha256"
"encoding/hex"
"fmt"
"strings"
"time"
"github.com/mathiasbq/hyperguild/ingestion/internal/brain"
"github.com/mathiasbq/hyperguild/ingestion/internal/classification"
)
// Service is the CaptureSession use-case. It depends only on ports.
type Service struct {
brain BrainStore
issues IssueTracker
summaries SummaryWriter
policy ClassificationPolicy
audit AuditSink
// now is the clock, injectable for deterministic summary paths and
// audit timestamps in tests.
now func() time.Time
}
// NewService constructs a Service from its ports. summaries may be nil
// when no summary persistence is wired; a CaptureInput with a Summary
// then fails that item rather than panicking.
func NewService(b BrainStore, tr IssueTracker, sw SummaryWriter, p ClassificationPolicy, a AuditSink) *Service {
return &Service{brain: b, issues: tr, summaries: sw, policy: p, audit: a, now: time.Now}
}
var validActions = map[string]bool{"create": true, "close": true, "comment": true}
// ErrSovereigntyRefused is returned when the I1 gate refuses a capture
// (confidential effective classification through a us-nexus origin). The
// REST adapter maps it to HTTP 403. Callers test with errors.Is.
var ErrSovereigntyRefused = fmt.Errorf("capture refused by I1 sovereignty gate")
// ErrAuditUnavailable is returned when the I5 audit gate refuses a capture
// before any write: a confidential capture whose central audit sink is
// unreachable, or the all-tiers floor where nothing can record the audit.
// The REST adapter maps it to HTTP 503. Callers test with errors.Is.
var ErrAuditUnavailable = fmt.Errorf("capture refused: audit substrate unavailable")
// assertedZoneMismatch returns a security-event string when the caller's
// harness label asserts a trust zone that contradicts the server-derived
// origin. A harness label that names no zone (the normal case, e.g.
// "claude-code") returns "". The label is never used as a gate input —
// this only flags the discrepancy for the audit trail.
func assertedZoneMismatch(harness string, derived Zone) string {
var asserted Zone
switch strings.ToLower(strings.TrimSpace(harness)) {
case "sovereign-soil", "sovereign":
asserted = ZoneSovereign
case "us-nexus", "usnexus":
asserted = ZoneUSNexus
default:
return "" // no zone claim
}
if asserted != derived {
return fmt.Sprintf("asserted-vs-derived origin mismatch: harness asserted %s, principal resolves to %s",
asserted, derived)
}
return ""
}
// Capture runs the use-case: validate (fail-closed), resolve effective
// classification (stricter of declared vs target-derived), then persist
// insights → tickets → summary best-effort, emit an audit record, and
// return a partial-aware receipt.
//
// A validation failure returns a non-nil error with nothing written. A
// per-item execution failure is recorded in the receipt (no rollback);
// the call still returns a nil error so the caller gets the partial
// receipt. The I1 origin gate and audit-down degradation are layered on
// by #53/#54 around this core.
func (s *Service) Capture(ctx context.Context, in CaptureInput) (CaptureReceipt, error) {
if err := s.validate(in); err != nil {
return CaptureReceipt{}, err
}
declared := classification.Public // unspecified ⇒ lowest ⇒ target floor governs
if in.Context.Classification != "" {
// Already validated parseable.
declared, _ = classification.ParseLevel(in.Context.Classification)
}
effective, securityEvents := s.resolveClassification(declared, in)
// Server-derived origin governs the I1 gate; a caller-asserted harness
// label that names a different zone is descriptive-only and logged as a
// security event (spec §4.2: a control keyed on attacker-suppliable
// input is not a control).
if ev := assertedZoneMismatch(in.Context.Harness, in.Context.Origin); ev != "" {
securityEvents = append(securityEvents, ev)
}
// I1 sovereignty gate: a confidential capture through a us-nexus origin
// is refused before ANY write. The refusal itself is audited (best
// effort) — refusals must be reconstructable too.
if effective == classification.Confidential && in.Context.Origin == ZoneUSNexus {
_ = s.audit.Record(ctx, AuditEntry{
Timestamp: s.now().UTC(),
Principal: in.Context.Principal,
Actor: in.Context.Actor,
Harness: in.Context.Harness,
SessionRef: in.Context.SessionRef,
EffectiveClassification: effective.String(),
Items: nil, // refused before any write
SecurityEvents: append(securityEvents, "I1 refusal: confidential capture via us-nexus origin"),
}, AuditCentral)
return CaptureReceipt{}, fmt.Errorf("%w: effective classification confidential through %s origin",
ErrSovereigntyRefused, in.Context.Origin)
}
receipt := CaptureReceipt{
Errors: []ItemError{},
EffectiveClassification: effective.String(),
DryRun: in.DryRun,
}
if in.DryRun {
// Would-be receipt: mark planned items ok, write nothing (not even
// audit — dry_run touches nothing).
for range in.Insights {
receipt.Insights = append(receipt.Insights, InsightResult{OK: true})
}
for _, tk := range in.Tickets {
receipt.Tickets = append(receipt.Tickets, TicketResult{Repo: tk.Repo, Action: tk.Action, Number: tk.Number, OK: true})
}
if in.Summary != nil {
receipt.Summary = &SummaryResult{Path: s.summaryPath(in.Context, in.Summary), OK: true}
}
return receipt, nil
}
// I5 audit gate: decide BEFORE any write whether this capture can be
// audited at its effective classification. Confidential + central sink
// down → refuse here, before writing anything; the all-tiers floor
// (nothing can record) likewise refuses. Internal/public degrade to the
// durable local buffer (signalled by AuditBuffered).
outcome, err := s.audit.Reserve(ctx, effective)
if err != nil {
return CaptureReceipt{}, fmt.Errorf("%w: %v", ErrAuditUnavailable, err)
}
var landed []string
for i, ins := range in.Insights {
res, item, err := s.persistInsight(ctx, ins)
receipt.Insights = append(receipt.Insights, res)
if err != nil {
receipt.Errors = append(receipt.Errors, ItemError{Item: fmt.Sprintf("insight[%d]", i), Error: err.Error()})
continue
}
landed = append(landed, item)
}
for i, tk := range in.Tickets {
res, err := s.persistTicket(ctx, tk)
receipt.Tickets = append(receipt.Tickets, res)
if err != nil {
receipt.Errors = append(receipt.Errors, ItemError{Item: fmt.Sprintf("ticket[%d]", i), Error: err.Error()})
continue
}
landed = append(landed, fmt.Sprintf("ticket:%s#%d", tk.Repo, res.Number))
}
if in.Summary != nil {
res, err := s.persistSummary(ctx, in.Context, in.Summary)
receipt.Summary = &res
if err != nil {
receipt.Errors = append(receipt.Errors, ItemError{Item: "summary", Error: err.Error()})
} else {
landed = append(landed, "summary:"+res.Path)
}
}
// I5: persist the request-level audit record of exactly what landed,
// using the outcome reserved before the writes. AuditBuffered surfaces
// the degraded (locally-buffered) state on the receipt.
if err := s.audit.Record(ctx, AuditEntry{
Timestamp: s.now().UTC(),
Principal: in.Context.Principal,
Actor: in.Context.Actor,
Harness: in.Context.Harness,
SessionRef: in.Context.SessionRef,
EffectiveClassification: effective.String(),
Items: landed,
SecurityEvents: securityEvents,
}, outcome); err != nil {
receipt.Errors = append(receipt.Errors, ItemError{Item: "audit", Error: err.Error()})
}
if outcome == AuditBuffered {
receipt.AuditBuffered = true
}
return receipt, nil
}
// validate enforces fail-closed structural validity over the whole
// request before any write. A bad declared classification, an invalid
// wing/hall, an empty insight, or a malformed ticket aborts the capture
// with nothing written.
func (s *Service) validate(in CaptureInput) error {
if in.Context.Classification != "" {
if _, err := classification.ParseLevel(in.Context.Classification); err != nil {
return fmt.Errorf("context.classification: %w", err)
}
}
for i, ins := range in.Insights {
if strings.TrimSpace(ins.Text) == "" {
return fmt.Errorf("insight[%d]: text is required", i)
}
if strings.TrimSpace(ins.Wing) == "" {
return fmt.Errorf("insight[%d]: wing is required", i)
}
if !brain.IsValidHall(ins.Hall) {
return fmt.Errorf("insight[%d]: invalid hall %q", i, ins.Hall)
}
}
for i, tk := range in.Tickets {
if strings.TrimSpace(tk.Repo) == "" {
return fmt.Errorf("ticket[%d]: repo is required", i)
}
if !validActions[tk.Action] {
return fmt.Errorf("ticket[%d]: invalid action %q (want create/close/comment)", i, tk.Action)
}
if tk.Action == "create" && strings.TrimSpace(tk.Title) == "" {
return fmt.Errorf("ticket[%d]: create requires a title", i)
}
if (tk.Action == "close" || tk.Action == "comment") && tk.Number <= 0 {
return fmt.Errorf("ticket[%d]: %s requires an issue number", i, tk.Action)
}
}
return nil
}
// resolveClassification computes the effective level (stricter of
// declared and every target's derived level) and collects a security
// event whenever the caller under-declared relative to a target floor.
func (s *Service) resolveClassification(declared classification.Level, in CaptureInput) (classification.Level, []string) {
effective := declared
var events []string
consider := func(kind classification.TargetKind, name string) {
derived := s.policy.Derive(classification.Target{Kind: kind, Name: name})
effective = classification.Stricter(effective, derived)
if declared < derived {
events = append(events, fmt.Sprintf("classification under-declared: declared=%s target=%s(%s) derived=%s",
declared, name, kindString(kind), derived))
}
}
for _, ins := range in.Insights {
consider(classification.WingTarget, ins.Wing)
}
for _, tk := range in.Tickets {
consider(classification.RepoTarget, tk.Repo)
}
if in.Summary != nil {
for _, repo := range in.Summary.ReposTouched {
consider(classification.RepoTarget, repo)
}
}
return effective, events
}
func (s *Service) persistInsight(ctx context.Context, ins Insight) (InsightResult, string, error) {
note := Note{Content: ins.Text, Wing: ins.Wing, Hall: ins.Hall, Filename: brain.Sanitise(firstLine(ins.Text))}
var ref Ref
var err error
if ins.SupersedeSlug != "" {
note.Reason = "superseded via capture"
ref, err = s.brain.Update(ctx, ins.SupersedeSlug, note)
} else {
ref, err = s.brain.Write(ctx, note)
}
if err != nil {
return InsightResult{OK: false, Superseded: ins.SupersedeSlug != ""}, "", err
}
return InsightResult{
ID: ref.ID, Path: ref.Path, ContentHash: ref.ContentHash,
Superseded: ref.Superseded, OK: true,
}, "insight:" + ref.ID, nil
}
func (s *Service) persistTicket(ctx context.Context, tk Ticket) (TicketResult, error) {
res := TicketResult{Repo: tk.Repo, Action: tk.Action, Number: tk.Number}
var ref IssueRef
var err error
switch tk.Action {
case "create":
ref, err = s.issues.CreateIssue(ctx, tk.Repo, tk.Title, tk.Body)
case "close":
ref, err = s.issues.CloseIssue(ctx, tk.Repo, tk.Number, tk.Body)
case "comment":
ref, err = s.issues.CommentIssue(ctx, tk.Repo, tk.Number, tk.Body)
}
if err != nil {
return res, err
}
if ref.Number != 0 {
res.Number = ref.Number
}
res.URL = ref.URL
res.OK = true
return res, nil
}
func (s *Service) persistSummary(ctx context.Context, c CaptureContext, sum *Summary) (SummaryResult, error) {
if s.summaries == nil {
return SummaryResult{OK: false}, fmt.Errorf("no summary writer configured")
}
path := s.summaryPath(c, sum)
content := s.renderSummary(c, sum)
repo := "ai-sessions"
if err := s.summaries.WriteFile(ctx, repo, path, content); err != nil {
return SummaryResult{Path: path, OK: false}, err
}
return SummaryResult{Path: path, OK: true}, nil
}
// summaryPath builds summaries/<harness>/<YYYY-MM>/<date>-<slug>-<ref8>.md.
// The ref8 disambiguator is derived from the session_ref (or the title
// when no ref is present) so distinct sessions never collide.
func (s *Service) summaryPath(c CaptureContext, sum *Summary) string {
t := s.now().UTC()
slug := brain.Sanitise(sum.Title)
if slug == "" {
slug = "summary"
}
seed := c.SessionRef
if seed == "" {
seed = sum.Title + sum.Body
}
sum8 := shortHash(seed)
return fmt.Sprintf("summaries/%s/%s/%s-%s-%s.md",
brain.Sanitise(c.Harness), t.Format("2006-01"), t.Format("2006-01-02"), slug, sum8)
}
// renderSummary stamps fidelity + session metadata into frontmatter so the
// richer-fidelity-supersedes-thinner collision rule has the data it needs.
func (s *Service) renderSummary(c CaptureContext, sum *Summary) string {
var b strings.Builder
b.WriteString("---\n")
fmt.Fprintf(&b, "title: %s\n", sum.Title)
fmt.Fprintf(&b, "harness: %s\n", c.Harness)
if c.SessionRef != "" {
fmt.Fprintf(&b, "session_ref: %s\n", c.SessionRef)
}
fmt.Fprintf(&b, "fidelity: %s\n", c.Fidelity)
fmt.Fprintf(&b, "captured_at: %s\n", s.now().UTC().Format(time.RFC3339))
if len(sum.ReposTouched) > 0 {
fmt.Fprintf(&b, "repos_touched: [%s]\n", strings.Join(sum.ReposTouched, ", "))
}
b.WriteString("---\n\n")
b.WriteString(sum.Body)
if !strings.HasSuffix(sum.Body, "\n") {
b.WriteByte('\n')
}
return b.String()
}
func kindString(k classification.TargetKind) string {
if k == classification.RepoTarget {
return "repo"
}
return "wing"
}
func firstLine(s string) string {
s = strings.TrimSpace(s)
if i := strings.IndexByte(s, '\n'); i >= 0 {
s = s[:i]
}
s = strings.TrimLeft(s, "# ")
if len(s) > 60 {
s = s[:60]
}
return s
}
func shortHash(s string) string {
sum := sha256.Sum256([]byte(s))
return hex.EncodeToString(sum[:])[:8]
}
+471
View File
@@ -0,0 +1,471 @@
package capture
import (
"context"
"errors"
"strings"
"testing"
"time"
"github.com/mathiasbq/hyperguild/ingestion/internal/classification"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
// --- fakes ---
type fakeBrain struct {
writes []Note
updates []Note
gets []string
failOn func(Note) error // nil = always succeed
hashSeq int
}
func (f *fakeBrain) ref(prefix string, n Note, superseded bool) Ref {
f.hashSeq++
path := "wiki/" + n.Wing + "/" + n.Hall + "/" + n.Filename + ".md"
return Ref{ID: path, Path: path, ContentHash: prefix + string(rune('0'+f.hashSeq)), Superseded: superseded}
}
func (f *fakeBrain) Write(_ context.Context, n Note) (Ref, error) {
if f.failOn != nil {
if err := f.failOn(n); err != nil {
return Ref{}, err
}
}
f.writes = append(f.writes, n)
return f.ref("w", n, false), nil
}
func (f *fakeBrain) Update(_ context.Context, slug string, n Note) (Ref, error) {
if f.failOn != nil {
if err := f.failOn(n); err != nil {
return Ref{}, err
}
}
n.Filename = slug
f.updates = append(f.updates, n)
return f.ref("u", n, true), nil
}
func (f *fakeBrain) Get(_ context.Context, id string) (StoredNote, error) {
f.gets = append(f.gets, id)
return StoredNote{ID: id, Path: id}, nil
}
type fakeTracker struct {
created []string
closed []int
comments []int
err error
}
func (f *fakeTracker) CreateIssue(_ context.Context, repo, title, _ string) (IssueRef, error) {
if f.err != nil {
return IssueRef{}, f.err
}
f.created = append(f.created, repo+":"+title)
return IssueRef{Repo: repo, Number: 100 + len(f.created), URL: "https://git/" + repo + "/issues/x"}, nil
}
func (f *fakeTracker) CloseIssue(_ context.Context, repo string, number int, _ string) (IssueRef, error) {
if f.err != nil {
return IssueRef{}, f.err
}
f.closed = append(f.closed, number)
return IssueRef{Repo: repo, Number: number}, nil
}
func (f *fakeTracker) CommentIssue(_ context.Context, repo string, number int, _ string) (IssueRef, error) {
if f.err != nil {
return IssueRef{}, f.err
}
f.comments = append(f.comments, number)
return IssueRef{Repo: repo, Number: number}, nil
}
type fakeSummary struct {
paths []string
content []string
err error
}
func (f *fakeSummary) WriteFile(_ context.Context, _, path, content string) error {
if f.err != nil {
return f.err
}
f.paths = append(f.paths, path)
f.content = append(f.content, content)
return nil
}
// fakePolicy derives from an explicit map; default Internal so tests pin
// behaviour without depending on the real defaulting.
type fakePolicy struct{ tags map[string]classification.Level }
func (p fakePolicy) Derive(t classification.Target) classification.Level {
if lvl, ok := p.tags[t.Name]; ok {
return lvl
}
return classification.Internal
}
type fakeAudit struct {
entries []AuditEntry
err error // Record error
reserveErr error // Reserve error (refuse before write)
reserveMode AuditOutcome
}
func (f *fakeAudit) Reserve(_ context.Context, _ classification.Level) (AuditOutcome, error) {
if f.reserveErr != nil {
return 0, f.reserveErr
}
return f.reserveMode, nil
}
func (f *fakeAudit) Record(_ context.Context, e AuditEntry, _ AuditOutcome) error {
if f.err != nil {
return f.err
}
f.entries = append(f.entries, e)
return nil
}
// --- helpers ---
func newSvc(b BrainStore, tr IssueTracker, sw SummaryWriter, p ClassificationPolicy, a AuditSink) *Service {
s := NewService(b, tr, sw, p, a)
s.now = func() time.Time { return time.Date(2026, 6, 22, 12, 0, 0, 0, time.UTC) }
return s
}
func baseCtx() CaptureContext {
return CaptureContext{Harness: "claude-code", Actor: "mathias", Principal: "mathias", Classification: "internal"}
}
// --- scenarios ---
func TestCaptureHappyPath(t *testing.T) {
b := &fakeBrain{}
tr := &fakeTracker{}
au := &fakeAudit{}
svc := newSvc(b, tr, nil, fakePolicy{}, au)
rec, err := svc.Capture(context.Background(), CaptureInput{
Context: baseCtx(),
Insights: []Insight{
{Text: "a", Wing: "hyperguild", Hall: "decisions", SupersedeSlug: ""},
{Text: "b", Wing: "hyperguild", Hall: "facts"},
},
Tickets: []Ticket{{Repo: "hyperguild", Action: "create", Title: "do x", Body: "y"}},
})
require.NoError(t, err)
require.Len(t, rec.Insights, 2)
for _, r := range rec.Insights {
assert.True(t, r.OK)
assert.NotEmpty(t, r.ContentHash, "read-after-write hash returned")
}
require.Len(t, rec.Tickets, 1)
assert.True(t, rec.Tickets[0].OK)
assert.Equal(t, 2, len(b.writes))
assert.Empty(t, rec.Errors)
// Audit emitted naming principal/harness + items that landed.
require.Len(t, au.entries, 1)
assert.Equal(t, "mathias", au.entries[0].Principal)
assert.Equal(t, "claude-code", au.entries[0].Harness)
assert.Len(t, au.entries[0].Items, 3)
}
func TestCaptureSupersedeNotDuplicate(t *testing.T) {
b := &fakeBrain{}
svc := newSvc(b, &fakeTracker{}, nil, fakePolicy{}, &fakeAudit{})
rec, err := svc.Capture(context.Background(), CaptureInput{
Context: baseCtx(),
Insights: []Insight{{Text: "revised", Wing: "hyperguild", Hall: "facts", SupersedeSlug: "prior-note"}},
})
require.NoError(t, err)
assert.Empty(t, b.writes, "supersede must not create")
require.Len(t, b.updates, 1)
assert.Equal(t, "prior-note", b.updates[0].Filename)
assert.True(t, rec.Insights[0].Superseded)
}
func TestCaptureValidationFailClosed(t *testing.T) {
b := &fakeBrain{}
tr := &fakeTracker{}
au := &fakeAudit{}
svc := newSvc(b, tr, nil, fakePolicy{}, au)
_, err := svc.Capture(context.Background(), CaptureInput{
Context: baseCtx(),
Insights: []Insight{
{Text: "ok", Wing: "hyperguild", Hall: "facts"},
{Text: "bad", Wing: "hyperguild", Hall: "garbage-hall"}, // invalid hall
},
Tickets: []Ticket{{Repo: "hyperguild", Action: "create", Title: "t"}},
})
require.Error(t, err)
// Nothing written anywhere.
assert.Empty(t, b.writes)
assert.Empty(t, b.updates)
assert.Empty(t, tr.created)
assert.Empty(t, au.entries)
}
func TestCaptureValidationRejectsBadTicket(t *testing.T) {
svc := newSvc(&fakeBrain{}, &fakeTracker{}, nil, fakePolicy{}, &fakeAudit{})
_, err := svc.Capture(context.Background(), CaptureInput{
Context: baseCtx(),
Tickets: []Ticket{{Repo: "hyperguild", Action: "frobnicate"}}, // bad action
})
require.Error(t, err)
_, err = svc.Capture(context.Background(), CaptureInput{
Context: baseCtx(),
Tickets: []Ticket{{Repo: "hyperguild", Action: "close"}}, // close needs number
})
require.Error(t, err)
}
func TestCapturePartialFailureBestEffort(t *testing.T) {
b := &fakeBrain{failOn: func(n Note) error {
if strings.Contains(n.Content, "FAIL") {
return errors.New("disk full")
}
return nil
}}
tr := &fakeTracker{}
au := &fakeAudit{}
svc := newSvc(b, tr, nil, fakePolicy{}, au)
rec, err := svc.Capture(context.Background(), CaptureInput{
Context: baseCtx(),
Insights: []Insight{
{Text: "good one", Wing: "hyperguild", Hall: "facts"},
{Text: "FAIL here", Wing: "hyperguild", Hall: "facts"},
},
Tickets: []Ticket{{Repo: "hyperguild", Action: "create", Title: "t"}},
})
require.NoError(t, err, "partial failure is not a request-level error")
assert.True(t, rec.Insights[0].OK)
assert.False(t, rec.Insights[1].OK)
assert.True(t, rec.Tickets[0].OK, "ticket still persisted; no rollback")
require.Len(t, rec.Errors, 1)
assert.Equal(t, "insight[1]", rec.Errors[0].Item)
// Audit reflects exactly what landed: 1 insight + 1 ticket.
require.Len(t, au.entries, 1)
assert.Len(t, au.entries[0].Items, 2)
}
func TestCaptureDryRunWritesNothing(t *testing.T) {
b := &fakeBrain{}
tr := &fakeTracker{}
au := &fakeAudit{}
svc := newSvc(b, tr, nil, fakePolicy{}, au)
rec, err := svc.Capture(context.Background(), CaptureInput{
Context: baseCtx(),
DryRun: true,
Insights: []Insight{{Text: "a", Wing: "hyperguild", Hall: "facts"}},
Tickets: []Ticket{{Repo: "hyperguild", Action: "create", Title: "t"}},
})
require.NoError(t, err)
assert.True(t, rec.DryRun)
assert.Len(t, rec.Insights, 1)
assert.True(t, rec.Insights[0].OK, "would-be receipt marks planned items ok")
// Nothing written anywhere, including audit.
assert.Empty(t, b.writes)
assert.Empty(t, tr.created)
assert.Empty(t, au.entries)
}
func TestCaptureStricterClassificationWins(t *testing.T) {
// Caller declares internal; target wing tagged confidential → effective confidential + security event.
b := &fakeBrain{}
au := &fakeAudit{}
pol := fakePolicy{tags: map[string]classification.Level{"client-seb": classification.Confidential}}
svc := newSvc(b, &fakeTracker{}, nil, pol, au)
ctx := baseCtx()
ctx.Classification = "internal"
rec, err := svc.Capture(context.Background(), CaptureInput{
Context: ctx,
Insights: []Insight{{Text: "x", Wing: "client-seb", Hall: "facts"}},
})
require.NoError(t, err)
assert.Equal(t, "confidential", rec.EffectiveClassification)
require.Len(t, au.entries, 1)
assert.NotEmpty(t, au.entries[0].SecurityEvents, "under-declaration logged as security event")
assert.Equal(t, "confidential", au.entries[0].EffectiveClassification)
}
func TestCaptureCallerRaisingSensitivityHonoured(t *testing.T) {
// Caller declares confidential; target internal → effective confidential, NOT a security event.
au := &fakeAudit{}
pol := fakePolicy{tags: map[string]classification.Level{"hyperguild": classification.Internal}}
svc := newSvc(&fakeBrain{}, &fakeTracker{}, nil, pol, au)
ctx := baseCtx()
ctx.Classification = "confidential"
rec, err := svc.Capture(context.Background(), CaptureInput{
Context: ctx,
Insights: []Insight{{Text: "x", Wing: "hyperguild", Hall: "facts"}},
})
require.NoError(t, err)
assert.Equal(t, "confidential", rec.EffectiveClassification)
assert.Empty(t, au.entries[0].SecurityEvents, "raising sensitivity is honoured, not flagged")
}
func TestCaptureSummaryPathAndFidelity(t *testing.T) {
sw := &fakeSummary{}
svc := newSvc(&fakeBrain{}, &fakeTracker{}, sw, fakePolicy{}, &fakeAudit{})
ctx := baseCtx()
ctx.Fidelity = "transcript-parse"
ctx.SessionRef = "abc123def456"
rec, err := svc.Capture(context.Background(), CaptureInput{
Context: ctx,
Summary: &Summary{Title: "Session Wrap", Body: "did stuff", ReposTouched: []string{"hyperguild"}},
})
require.NoError(t, err)
require.NotNil(t, rec.Summary)
assert.True(t, rec.Summary.OK)
require.Len(t, sw.paths, 1)
assert.True(t, strings.HasPrefix(sw.paths[0], "summaries/claude-code/2026-06/"), "path: %s", sw.paths[0])
assert.Contains(t, sw.paths[0], "session-wrap")
assert.Contains(t, sw.content[0], "fidelity: transcript-parse", "fidelity stamped in frontmatter")
}
// --- I1 sovereignty gate (#53) ---
func TestCaptureRefusesConfidentialViaUSNexus(t *testing.T) {
b := &fakeBrain{}
tr := &fakeTracker{}
au := &fakeAudit{}
pol := fakePolicy{tags: map[string]classification.Level{"client-seb": classification.Confidential}}
svc := newSvc(b, tr, nil, pol, au)
ctx := baseCtx()
ctx.Classification = "confidential"
ctx.Origin = ZoneUSNexus
_, err := svc.Capture(context.Background(), CaptureInput{
Context: ctx,
Insights: []Insight{{Text: "x", Wing: "client-seb", Hall: "facts"}},
})
require.Error(t, err)
assert.ErrorIs(t, err, ErrSovereigntyRefused)
// Refused before any write.
assert.Empty(t, b.writes)
assert.Empty(t, tr.created)
// Refusal is audited.
require.Len(t, au.entries, 1)
assert.Empty(t, au.entries[0].Items, "no items landed on refusal")
}
func TestCaptureAllowsConfidentialViaSovereign(t *testing.T) {
b := &fakeBrain{}
pol := fakePolicy{tags: map[string]classification.Level{"client-seb": classification.Confidential}}
svc := newSvc(b, &fakeTracker{}, nil, pol, &fakeAudit{})
ctx := baseCtx()
ctx.Classification = "confidential"
ctx.Origin = ZoneSovereign
rec, err := svc.Capture(context.Background(), CaptureInput{
Context: ctx,
Insights: []Insight{{Text: "x", Wing: "client-seb", Hall: "facts"}},
})
require.NoError(t, err)
assert.True(t, rec.Insights[0].OK)
assert.Len(t, b.writes, 1)
}
func TestCaptureAssertedLabelIgnoredAndLogged(t *testing.T) {
// Caller asserts harness "sovereign-soil" but principal resolves to
// us-nexus; confidential ⇒ refused, and the discrepancy is a security event.
au := &fakeAudit{}
pol := fakePolicy{tags: map[string]classification.Level{"client-seb": classification.Confidential}}
svc := newSvc(&fakeBrain{}, &fakeTracker{}, nil, pol, au)
ctx := baseCtx()
ctx.Harness = "sovereign-soil" // asserted
ctx.Origin = ZoneUSNexus // server-derived
ctx.Classification = "confidential"
_, err := svc.Capture(context.Background(), CaptureInput{
Context: ctx,
Insights: []Insight{{Text: "x", Wing: "client-seb", Hall: "facts"}},
})
require.ErrorIs(t, err, ErrSovereigntyRefused)
require.Len(t, au.entries, 1)
joined := strings.Join(au.entries[0].SecurityEvents, " | ")
assert.Contains(t, joined, "asserted-vs-derived origin mismatch")
assert.Contains(t, joined, "I1 refusal")
}
func TestCaptureInternalViaUSNexusAllowed(t *testing.T) {
// us-nexus origin is fine for non-confidential data.
b := &fakeBrain{}
svc := newSvc(b, &fakeTracker{}, nil, fakePolicy{}, &fakeAudit{})
ctx := baseCtx()
ctx.Origin = ZoneUSNexus // internal classification, so gate doesn't fire
rec, err := svc.Capture(context.Background(), CaptureInput{
Context: ctx,
Insights: []Insight{{Text: "x", Wing: "hyperguild", Hall: "facts"}},
})
require.NoError(t, err)
assert.True(t, rec.Insights[0].OK)
}
// --- I5 audit gate (#54) ---
func TestCaptureRefusesWhenAuditReserveFails(t *testing.T) {
// Reserve refusing (e.g. confidential + central sink down, or the floor)
// aborts the capture before any write.
b := &fakeBrain{}
tr := &fakeTracker{}
au := &fakeAudit{reserveErr: errors.New("central sink unreachable")}
svc := newSvc(b, tr, nil, fakePolicy{}, au)
_, err := svc.Capture(context.Background(), CaptureInput{
Context: baseCtx(),
Insights: []Insight{{Text: "x", Wing: "hyperguild", Hall: "facts"}},
})
require.Error(t, err)
assert.ErrorIs(t, err, ErrAuditUnavailable)
assert.Empty(t, b.writes, "nothing written when audit unavailable")
assert.Empty(t, tr.created)
}
func TestCaptureFlagsLocallyBufferedAudit(t *testing.T) {
// Reserve returns AuditBuffered (internal/public, central down) → capture
// proceeds and the receipt flags the degraded audit state.
b := &fakeBrain{}
au := &fakeAudit{reserveMode: AuditBuffered}
svc := newSvc(b, &fakeTracker{}, nil, fakePolicy{}, au)
rec, err := svc.Capture(context.Background(), CaptureInput{
Context: baseCtx(),
Insights: []Insight{{Text: "x", Wing: "hyperguild", Hall: "facts"}},
})
require.NoError(t, err)
assert.True(t, rec.Insights[0].OK, "capture proceeds on degraded audit")
assert.True(t, rec.AuditBuffered, "receipt flags locally-buffered audit")
require.Len(t, au.entries, 1)
}
func TestCaptureDryRunSkipsAuditGate(t *testing.T) {
// dry_run must not even probe the audit sink (writes nothing anywhere).
au := &fakeAudit{reserveErr: errors.New("would refuse")}
svc := newSvc(&fakeBrain{}, &fakeTracker{}, nil, fakePolicy{}, au)
rec, err := svc.Capture(context.Background(), CaptureInput{
Context: baseCtx(),
DryRun: true,
Insights: []Insight{{Text: "x", Wing: "hyperguild", Hall: "facts"}},
})
require.NoError(t, err, "dry-run does not hit the audit gate")
assert.True(t, rec.DryRun)
assert.Empty(t, au.entries)
}
+232
View File
@@ -0,0 +1,232 @@
// Package capturehttp is the REST adapter for the capture use-case: the
// POST /capture door (#53). It is deliberately thin — authenticate, derive
// the trust-zone origin from the authenticated principal, decode the
// request, call capture.Service, map the receipt to an HTTP status. No
// business logic lives here; the I1 gate, validation, and orchestration
// are all in the use-case.
package capturehttp
import (
"context"
"crypto/subtle"
"encoding/json"
"errors"
"io"
"net/http"
"strings"
"github.com/mathiasbq/hyperguild/ingestion/internal/capture"
)
// Validator validates a Bearer JWT and returns its subject. The chassis
// *auth.JWTValidator satisfies it (including its nil-receiver "disabled"
// behaviour), and tests can substitute a fake without a live JWKS.
type Validator interface {
Validate(ctx context.Context, rawToken string) (string, error)
}
// Handler serves POST /capture.
type Handler struct {
svc *capture.Service
validator Validator // nil ⇒ JWT auth disabled
staticToken string // "" ⇒ static auth disabled
staticPrincipal string // principal name attributed to static-token callers
resolver OriginResolver
}
// New constructs a capture HTTP handler. staticToken callers are
// attributed to staticPrincipal (a sovereign homelab identity); JWT
// callers are attributed to their token subject.
func New(svc *capture.Service, validator Validator, staticToken, staticPrincipal string, resolver OriginResolver) *Handler {
if staticPrincipal == "" {
staticPrincipal = "local-cli"
}
return &Handler{
svc: svc,
validator: validator,
staticToken: staticToken,
staticPrincipal: staticPrincipal,
resolver: resolver,
}
}
// wire types — the POST /capture request body.
type request struct {
Context contextBody `json:"context"`
Insights []insightBody `json:"insights"`
Tickets []ticketBody `json:"tickets"`
Summary *summaryBody `json:"summary,omitempty"`
DryRun bool `json:"dry_run"`
}
type contextBody struct {
Harness string `json:"harness"`
SessionRef string `json:"session_ref"`
Fidelity string `json:"fidelity"`
Actor string `json:"actor"`
Classification string `json:"classification"`
}
type insightBody struct {
Text string `json:"text"`
Wing string `json:"wing"`
Hall string `json:"hall"`
SupersedeSlug string `json:"supersede_slug,omitempty"`
}
type ticketBody struct {
Repo string `json:"repo"`
Action string `json:"action"`
Number int `json:"number,omitempty"`
Title string `json:"title,omitempty"`
Body string `json:"body,omitempty"`
}
type summaryBody struct {
Title string `json:"title"`
Body string `json:"body"`
ReposTouched []string `json:"repos_touched,omitempty"`
}
// ServeHTTP authenticates, derives origin, runs the use-case, and maps the
// result to an HTTP status.
func (h *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
principal, viaStatic, ok := Authenticate(r, h.staticToken, h.staticPrincipal, h.validator)
if !ok {
http.Error(w, "unauthorized", http.StatusUnauthorized)
return
}
body, err := io.ReadAll(r.Body)
if err != nil {
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "read body"})
return
}
in, err := DecodeRequest(body)
if err != nil {
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "invalid JSON"})
return
}
// Principal and origin are server-derived — overwrite anything the
// caller may have tried to put in the body.
in.Context.Principal = principal
in.Context.Origin = h.resolver.Resolve(principal, viaStatic)
rec, err := h.svc.Capture(r.Context(), in)
switch {
case errors.Is(err, capture.ErrSovereigntyRefused):
writeJSON(w, http.StatusForbidden, map[string]string{"error": err.Error()})
return
case errors.Is(err, capture.ErrAuditUnavailable):
// I5 refusal: confidential + audit sink down, or the all-tiers floor.
writeJSON(w, http.StatusServiceUnavailable, map[string]string{"error": err.Error()})
return
case err != nil:
// Pre-write validation failure (fail-closed).
writeJSON(w, http.StatusBadRequest, map[string]string{"error": err.Error()})
return
}
writeJSON(w, statusFor(rec), rec)
}
// Authenticate mirrors the chassis Bearer precedence (static token wins,
// then Dex JWT) and returns the resolved principal plus whether the static
// path was taken — the chassis middleware hides both, and capture (REST or
// MCP) needs them to derive the trust-zone origin. ok is false when no
// credential matched.
func Authenticate(r *http.Request, staticToken, staticPrincipal string, validator Validator) (principal string, viaStatic, ok bool) {
raw, found := strings.CutPrefix(r.Header.Get("Authorization"), "Bearer ")
if !found || raw == "" {
return "", false, false
}
if staticToken != "" && subtle.ConstantTimeCompare([]byte(raw), []byte(staticToken)) == 1 {
return staticPrincipal, true, true
}
if validator != nil {
if sub, err := validator.Validate(r.Context(), raw); err == nil && sub != "" {
return sub, false, true
}
}
return "", false, false
}
// DecodeRequest parses a capture request body into a CaptureInput. Shared
// by the REST adapter and the MCP capture tool so the wire shape has one
// definition. Principal and Origin are NOT set here — the caller sets them
// from the authenticated identity.
func DecodeRequest(data []byte) (capture.CaptureInput, error) {
var b request
if err := json.Unmarshal(data, &b); err != nil {
return capture.CaptureInput{}, err
}
return b.toInput(), nil
}
func (b request) toInput() capture.CaptureInput {
in := capture.CaptureInput{
Context: capture.CaptureContext{
Harness: b.Context.Harness,
SessionRef: b.Context.SessionRef,
Fidelity: b.Context.Fidelity,
Actor: b.Context.Actor,
Classification: b.Context.Classification,
},
DryRun: b.DryRun,
}
for _, i := range b.Insights {
in.Insights = append(in.Insights, capture.Insight{
Text: i.Text, Wing: i.Wing, Hall: i.Hall, SupersedeSlug: i.SupersedeSlug,
})
}
for _, t := range b.Tickets {
in.Tickets = append(in.Tickets, capture.Ticket{
Repo: t.Repo, Action: t.Action, Number: t.Number, Title: t.Title, Body: t.Body,
})
}
if b.Summary != nil {
in.Summary = &capture.Summary{
Title: b.Summary.Title, Body: b.Summary.Body, ReposTouched: b.Summary.ReposTouched,
}
}
return in
}
// statusFor maps a receipt to an HTTP status: 200 all-ok (or dry-run),
// 207 partial, 502 everything-failed.
func statusFor(rec capture.CaptureReceipt) int {
if rec.DryRun {
return http.StatusOK
}
var ok, fail int
for _, i := range rec.Insights {
count(&ok, &fail, i.OK)
}
for _, t := range rec.Tickets {
count(&ok, &fail, t.OK)
}
if rec.Summary != nil {
count(&ok, &fail, rec.Summary.OK)
}
switch {
case fail == 0:
return http.StatusOK
case ok == 0:
return http.StatusBadGateway // every persistence attempt failed
default:
return http.StatusMultiStatus // 207: partial success
}
}
func count(ok, fail *int, isOK bool) {
if isOK {
*ok++
} else {
*fail++
}
}
func writeJSON(w http.ResponseWriter, status int, v any) {
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(status)
_ = json.NewEncoder(w).Encode(v)
}
@@ -0,0 +1,198 @@
package capturehttp_test
import (
"bytes"
"context"
"encoding/json"
"errors"
"net/http"
"net/http/httptest"
"testing"
"github.com/mathiasbq/hyperguild/ingestion/internal/audit"
"github.com/mathiasbq/hyperguild/ingestion/internal/brainstore"
"github.com/mathiasbq/hyperguild/ingestion/internal/capture"
"github.com/mathiasbq/hyperguild/ingestion/internal/capturehttp"
"github.com/mathiasbq/hyperguild/ingestion/internal/classification"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
const staticTok = "static-secret"
// fakeValidator stands in for the chassis JWT validator.
type fakeValidator struct {
subject string
err error
}
func (f fakeValidator) Validate(context.Context, string) (string, error) {
return f.subject, f.err
}
type fakeTracker struct{ failCreate bool }
func (f fakeTracker) CreateIssue(context.Context, string, string, string) (capture.IssueRef, error) {
if f.failCreate {
return capture.IssueRef{}, errors.New("gitea down")
}
return capture.IssueRef{Repo: "hyperguild", Number: 1, URL: "https://git/1"}, nil
}
func (fakeTracker) CloseIssue(context.Context, string, int, string) (capture.IssueRef, error) {
return capture.IssueRef{}, nil
}
func (fakeTracker) CommentIssue(context.Context, string, int, string) (capture.IssueRef, error) {
return capture.IssueRef{}, nil
}
func newHandler(t *testing.T, v capturehttp.Validator, tr capture.IssueTracker, sovereign []string) *capturehttp.Handler {
t.Helper()
cfg, err := classification.Load(t.TempDir())
require.NoError(t, err)
svc := capture.NewService(brainstore.New(t.TempDir()), tr, nil, cfg, audit.NewSlogSink(nil))
return capturehttp.New(svc, v, staticTok, "local-cli", capturehttp.NewOriginResolver(sovereign))
}
func do(t *testing.T, h *capturehttp.Handler, authz string, body any) *httptest.ResponseRecorder {
t.Helper()
b, _ := json.Marshal(body)
req := httptest.NewRequest(http.MethodPost, "/capture", bytes.NewReader(b))
if authz != "" {
req.Header.Set("Authorization", authz)
}
rr := httptest.NewRecorder()
h.ServeHTTP(rr, req)
return rr
}
func internalReq() map[string]any {
return map[string]any{
"context": map[string]any{"harness": "claude-code", "actor": "mathias", "classification": "internal"},
"insights": []map[string]any{{"text": "a fact", "wing": "hyperguild", "hall": "facts"}},
}
}
func TestUnauthorizedWithoutToken(t *testing.T) {
h := newHandler(t, fakeValidator{err: errors.New("no")}, fakeTracker{}, nil)
rr := do(t, h, "", internalReq())
assert.Equal(t, http.StatusUnauthorized, rr.Code)
}
func TestUnauthorizedBadToken(t *testing.T) {
h := newHandler(t, fakeValidator{err: errors.New("bad jwt")}, fakeTracker{}, nil)
rr := do(t, h, "Bearer wrong", internalReq())
assert.Equal(t, http.StatusUnauthorized, rr.Code)
}
func TestHappyPathStaticToken(t *testing.T) {
h := newHandler(t, nil, fakeTracker{}, nil)
rr := do(t, h, "Bearer "+staticTok, map[string]any{
"context": map[string]any{"harness": "claude-code", "actor": "mathias", "classification": "internal"},
"insights": []map[string]any{{"text": "a fact", "wing": "hyperguild", "hall": "facts"}},
"tickets": []map[string]any{{"repo": "hyperguild", "action": "create", "title": "t"}},
})
require.Equal(t, http.StatusOK, rr.Code)
var rec capture.CaptureReceipt
require.NoError(t, json.Unmarshal(rr.Body.Bytes(), &rec))
assert.True(t, rec.Insights[0].OK)
assert.True(t, rec.Tickets[0].OK)
assert.Empty(t, rec.Errors)
}
func TestConfidentialViaUSNexusRefused(t *testing.T) {
// JWT principal not in the sovereign allowlist ⇒ us-nexus; confidential ⇒ 403.
h := newHandler(t, fakeValidator{subject: "claudeai-oauth-client"}, fakeTracker{}, nil)
rr := do(t, h, "Bearer jwt-token", map[string]any{
"context": map[string]any{"harness": "claudeai-chat", "actor": "mathias", "classification": "confidential"},
"insights": []map[string]any{{"text": "secret", "wing": "client-seb", "hall": "facts"}},
})
assert.Equal(t, http.StatusForbidden, rr.Code)
assert.Contains(t, rr.Body.String(), "sovereignty")
}
func TestConfidentialViaSovereignJWTAllowed(t *testing.T) {
// Same confidential payload, but the principal is allowlisted sovereign ⇒ allowed.
h := newHandler(t, fakeValidator{subject: "koala-cli"}, fakeTracker{}, []string{"koala-cli"})
rr := do(t, h, "Bearer jwt-token", map[string]any{
"context": map[string]any{"harness": "claude-code", "actor": "mathias", "classification": "confidential"},
"insights": []map[string]any{{"text": "secret", "wing": "client-seb", "hall": "facts"}},
})
require.Equal(t, http.StatusOK, rr.Code)
}
func TestStaticTokenIsSovereignSoConfidentialAllowed(t *testing.T) {
h := newHandler(t, nil, fakeTracker{}, nil)
rr := do(t, h, "Bearer "+staticTok, map[string]any{
"context": map[string]any{"harness": "claude-code", "actor": "mathias", "classification": "confidential"},
"insights": []map[string]any{{"text": "secret", "wing": "client-seb", "hall": "facts"}},
})
assert.Equal(t, http.StatusOK, rr.Code)
}
func TestValidationRejectedBeforeWrite(t *testing.T) {
h := newHandler(t, nil, fakeTracker{}, nil)
rr := do(t, h, "Bearer "+staticTok, map[string]any{
"context": map[string]any{"actor": "mathias", "classification": "internal"},
"insights": []map[string]any{{"text": "x", "wing": "hyperguild", "hall": "not-a-hall"}},
})
assert.Equal(t, http.StatusBadRequest, rr.Code)
}
func TestPartialFailureIs207(t *testing.T) {
h := newHandler(t, nil, fakeTracker{failCreate: true}, nil)
rr := do(t, h, "Bearer "+staticTok, map[string]any{
"context": map[string]any{"harness": "claude-code", "actor": "mathias", "classification": "internal"},
"insights": []map[string]any{{"text": "ok insight", "wing": "hyperguild", "hall": "facts"}},
"tickets": []map[string]any{{"repo": "hyperguild", "action": "create", "title": "fails"}},
})
assert.Equal(t, http.StatusMultiStatus, rr.Code)
var rec capture.CaptureReceipt
require.NoError(t, json.Unmarshal(rr.Body.Bytes(), &rec))
assert.True(t, rec.Insights[0].OK)
assert.False(t, rec.Tickets[0].OK)
assert.Len(t, rec.Errors, 1)
}
func TestDryRunWritesNothing(t *testing.T) {
h := newHandler(t, nil, fakeTracker{}, nil)
rr := do(t, h, "Bearer "+staticTok, map[string]any{
"context": map[string]any{"harness": "claude-code", "actor": "mathias", "classification": "internal"},
"insights": []map[string]any{{"text": "a", "wing": "hyperguild", "hall": "facts"}},
"dry_run": true,
})
require.Equal(t, http.StatusOK, rr.Code)
var rec capture.CaptureReceipt
require.NoError(t, json.Unmarshal(rr.Body.Bytes(), &rec))
assert.True(t, rec.DryRun)
}
func TestCallerCannotForgeOrigin(t *testing.T) {
// Even if the body tried to assert a sovereign harness, a us-nexus JWT
// principal + confidential ⇒ refused. (Origin is server-derived.)
h := newHandler(t, fakeValidator{subject: "claudeai-oauth-client"}, fakeTracker{}, nil)
rr := do(t, h, "Bearer jwt", map[string]any{
"context": map[string]any{"harness": "sovereign-soil", "actor": "mathias", "classification": "confidential"},
"insights": []map[string]any{{"text": "secret", "wing": "client-seb", "hall": "facts"}},
})
assert.Equal(t, http.StatusForbidden, rr.Code)
}
// refusingAudit refuses at Reserve (e.g. confidential + loki down, or floor).
type refusingAudit struct{}
func (refusingAudit) Reserve(context.Context, classification.Level) (capture.AuditOutcome, error) {
return 0, errors.New("central audit sink unreachable")
}
func (refusingAudit) Record(context.Context, capture.AuditEntry, capture.AuditOutcome) error {
return nil
}
func TestAuditUnavailableIs503(t *testing.T) {
cfg, err := classification.Load(t.TempDir())
require.NoError(t, err)
svc := capture.NewService(brainstore.New(t.TempDir()), fakeTracker{}, nil, cfg, refusingAudit{})
h := capturehttp.New(svc, nil, staticTok, "local-cli", capturehttp.NewOriginResolver(nil))
rr := do(t, h, "Bearer "+staticTok, internalReq())
assert.Equal(t, http.StatusServiceUnavailable, rr.Code)
}
+42
View File
@@ -0,0 +1,42 @@
package capturehttp
import "github.com/mathiasbq/hyperguild/ingestion/internal/capture"
// OriginResolver maps an authenticated principal to its trust zone
// (spec §4.2). The mapping is server-side and never reads caller input.
//
// Rules:
// - The static-token path is a homelab CLI caller on sovereign soil →
// ZoneSovereign.
// - A JWT principal in the sovereign allowlist → ZoneSovereign.
// - Any other JWT principal (e.g. claude.ai's OAuth identity, or any
// unrecognised subject) → ZoneUSNexus.
//
// The default is the strict one: an unknown principal is treated as
// us-nexus so the I1 gate fails safe (refuses confidential), exactly as
// an untagged classification target fails safe to confidential (#50).
type OriginResolver struct {
sovereign map[string]bool
}
// NewOriginResolver builds a resolver whose JWT sovereign principals are
// the given subjects. The static-token caller is always sovereign and
// need not be listed.
func NewOriginResolver(sovereignPrincipals []string) OriginResolver {
m := make(map[string]bool, len(sovereignPrincipals))
for _, p := range sovereignPrincipals {
if p != "" {
m[p] = true
}
}
return OriginResolver{sovereign: m}
}
// Resolve returns the trust zone for a principal. viaStatic is true when
// the static-token auth path was taken.
func (r OriginResolver) Resolve(principal string, viaStatic bool) capture.Zone {
if viaStatic || r.sovereign[principal] {
return capture.ZoneSovereign
}
return capture.ZoneUSNexus
}
@@ -0,0 +1,189 @@
// Package classification defines the data-sensitivity taxonomy and the
// per-wing / per-repo tagging the capture server reads to enforce the I1
// sovereignty gate (issue #50, capture spec §4.1).
//
// The single load-bearing property is fail-safe-to-strictest: a target
// with no explicit tag and no known default classifies as Confidential,
// never as something more permissive. A missing tag must never silently
// downgrade — that would turn the I1 gate into theatre.
//
// Classification is read from an optional classification.yaml at the
// brain root. A central, Flux-reconcilable file is deliberate: it is
// auditable in one place (I2/I5), it does not require a live Gitea client
// to classify a repo (so this package has no dependency on the gitea
// tracker work), and it avoids tagging a wing's _index.md frontmatter —
// which BuildWingIndex regenerates and would clobber.
package classification
import (
"fmt"
"os"
"path/filepath"
"strings"
"gopkg.in/yaml.v3"
)
// Level is a data-sensitivity tier. Higher is stricter, so the "stricter
// wins" rule (spec §4.1 model C) is a plain max.
type Level int
const (
Public Level = iota
Internal
Confidential
)
// String returns the canonical lowercase token for a level.
func (l Level) String() string {
switch l {
case Public:
return "public"
case Internal:
return "internal"
case Confidential:
return "confidential"
default:
return fmt.Sprintf("level(%d)", int(l))
}
}
// ParseLevel parses a level token (case-insensitive, surrounding space
// tolerated). An unknown token is an error — callers must decide what to
// do with bad input rather than have it silently coerced.
func ParseLevel(s string) (Level, error) {
switch strings.ToLower(strings.TrimSpace(s)) {
case "public":
return Public, nil
case "internal":
return Internal, nil
case "confidential":
return Confidential, nil
default:
return Confidential, fmt.Errorf("unknown classification level %q (want public/internal/confidential)", s)
}
}
// Stricter returns the more restrictive of two levels.
func Stricter(a, b Level) Level {
if a > b {
return a
}
return b
}
// TargetKind distinguishes the two kinds of capture destination.
type TargetKind int
const (
WingTarget TargetKind = iota // a brain wing (insights land here)
RepoTarget // a Gitea repo (tickets / summaries land here)
)
// Target names a capture destination to classify.
type Target struct {
Kind TargetKind
Name string
}
// Config holds the explicit per-wing / per-repo classification tags read
// from classification.yaml. Absent entries fall through to the built-in
// defaults in defaultFor. The zero value (no file) is valid and applies
// defaults to everything.
type Config struct {
wings map[string]Level
repos map[string]Level
}
// rawConfig is the on-disk YAML shape: string→string maps, parsed into
// validated levels by Load.
type rawConfig struct {
Wings map[string]string `yaml:"wings"`
Repos map[string]string `yaml:"repos"`
}
// Load reads classification.yaml from brainDir. An absent file is not an
// error — it yields an empty config where every target classifies by the
// built-in defaults. A malformed file, or any unparseable level token in
// it, is a hard error: a classification source the server cannot trust
// must fail loud, not degrade silently.
func Load(brainDir string) (*Config, error) {
cfg := &Config{wings: map[string]Level{}, repos: map[string]Level{}}
data, err := os.ReadFile(filepath.Join(brainDir, "classification.yaml"))
if err != nil {
if os.IsNotExist(err) {
return cfg, nil
}
return nil, fmt.Errorf("read classification.yaml: %w", err)
}
var raw rawConfig
if err := yaml.Unmarshal(data, &raw); err != nil {
return nil, fmt.Errorf("parse classification.yaml: %w", err)
}
for name, lvl := range raw.Wings {
parsed, perr := ParseLevel(lvl)
if perr != nil {
return nil, fmt.Errorf("wing %q: %w", name, perr)
}
cfg.wings[normalise(name)] = parsed
}
for name, lvl := range raw.Repos {
parsed, perr := ParseLevel(lvl)
if perr != nil {
return nil, fmt.Errorf("repo %q: %w", name, perr)
}
cfg.repos[normalise(name)] = parsed
}
return cfg, nil
}
// Derive returns the classification for any target — the function the
// capture use-case calls per item.
func (c *Config) Derive(t Target) Level {
if t.Kind == RepoTarget {
return c.Repo(t.Name)
}
return c.Wing(t.Name)
}
// Wing classifies a brain wing: an explicit tag wins, else defaults.
func (c *Config) Wing(name string) Level {
if lvl, ok := c.wings[normalise(name)]; ok {
return lvl
}
return defaultFor(name)
}
// Repo classifies a Gitea repo: an explicit tag wins, else defaults.
func (c *Config) Repo(name string) Level {
if lvl, ok := c.repos[normalise(name)]; ok {
return lvl
}
return defaultFor(name)
}
// defaultFor applies the built-in defaulting rules when a target has no
// explicit tag:
// - client-* → Confidential (client work is confidential by default)
// - hyperguild / homelab → Internal (the operator's own infra)
// - everything else → Confidential (fail safe to strictest)
func defaultFor(name string) Level {
n := normalise(name)
if strings.HasPrefix(n, "client-") {
return Confidential
}
switch n {
case "hyperguild", "homelab":
return Internal
default:
return Confidential
}
}
// normalise lowercases and trims a wing/repo name so matching and the
// client-* prefix check are case-insensitive.
func normalise(name string) string {
return strings.ToLower(strings.TrimSpace(name))
}
@@ -0,0 +1,112 @@
package classification
import (
"os"
"path/filepath"
"testing"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func TestLevelOrderingAndString(t *testing.T) {
assert.True(t, Public < Internal)
assert.True(t, Internal < Confidential)
assert.Equal(t, "public", Public.String())
assert.Equal(t, "internal", Internal.String())
assert.Equal(t, "confidential", Confidential.String())
}
func TestParseLevel(t *testing.T) {
for s, want := range map[string]Level{
"public": Public, "internal": Internal, "confidential": Confidential,
"PUBLIC": Public, " Confidential ": Confidential,
} {
got, err := ParseLevel(s)
require.NoError(t, err, s)
assert.Equal(t, want, got, s)
}
_, err := ParseLevel("secret")
require.Error(t, err, "unknown level must error, not silently default")
_, err = ParseLevel("")
require.Error(t, err)
}
func TestStricterReturnsMax(t *testing.T) {
assert.Equal(t, Confidential, Stricter(Internal, Confidential))
assert.Equal(t, Confidential, Stricter(Confidential, Public))
assert.Equal(t, Internal, Stricter(Public, Internal))
assert.Equal(t, Public, Stricter(Public, Public))
}
func TestLoadAbsentFileIsDefaultsOnly(t *testing.T) {
cfg, err := Load(t.TempDir())
require.NoError(t, err, "absent classification.yaml must not be an error — defaults apply")
require.NotNil(t, cfg)
// Pure defaulting still works.
assert.Equal(t, Internal, cfg.Wing("hyperguild"))
assert.Equal(t, Confidential, cfg.Wing("anything-unknown"))
}
func TestLoadParsesExplicitTags(t *testing.T) {
dir := t.TempDir()
require.NoError(t, os.WriteFile(filepath.Join(dir, "classification.yaml"), []byte(
"wings:\n research-public: public\n hyperguild: confidential\nrepos:\n infra: internal\n research-public: public\n",
), 0o644))
cfg, err := Load(dir)
require.NoError(t, err)
// Explicit tag wins over the built-in default (hyperguild default is internal).
assert.Equal(t, Confidential, cfg.Wing("hyperguild"))
// Explicit public is honoured.
assert.Equal(t, Public, cfg.Wing("research-public"))
assert.Equal(t, Internal, cfg.Repo("infra"))
assert.Equal(t, Public, cfg.Repo("research-public"))
}
func TestLoadRejectsUnknownLevelInFile(t *testing.T) {
dir := t.TempDir()
require.NoError(t, os.WriteFile(filepath.Join(dir, "classification.yaml"),
[]byte("wings:\n x: top-secret\n"), 0o644))
_, err := Load(dir)
require.Error(t, err, "an unparseable level in the config must fail loud, not be ignored")
}
func TestWingDefaulting(t *testing.T) {
cfg, err := Load(t.TempDir())
require.NoError(t, err)
cases := map[string]Level{
"client-seb": Confidential, // client-* → confidential
"client-mastercard": Confidential,
"hyperguild": Internal,
"homelab": Internal,
"jepa-fx": Confidential, // unknown → fail safe to strictest
"": Confidential, // empty → fail safe
}
for wing, want := range cases {
assert.Equal(t, want, cfg.Wing(wing), "wing %q", wing)
}
}
func TestRepoDefaulting(t *testing.T) {
cfg, err := Load(t.TempDir())
require.NoError(t, err)
assert.Equal(t, Confidential, cfg.Repo("client-seb-pipeline"))
assert.Equal(t, Internal, cfg.Repo("hyperguild"))
assert.Equal(t, Confidential, cfg.Repo("some-unknown-repo"), "untagged repo → confidential (fail safe)")
}
func TestDeriveUnifiedTarget(t *testing.T) {
cfg, err := Load(t.TempDir())
require.NoError(t, err)
assert.Equal(t, Internal, cfg.Derive(Target{Kind: WingTarget, Name: "homelab"}))
assert.Equal(t, Confidential, cfg.Derive(Target{Kind: RepoTarget, Name: "client-x"}))
assert.Equal(t, Confidential, cfg.Derive(Target{Kind: WingTarget, Name: "untagged"}))
}
func TestCaseInsensitiveMatching(t *testing.T) {
cfg, err := Load(t.TempDir())
require.NoError(t, err)
assert.Equal(t, Confidential, cfg.Wing("Client-SEB"), "client- prefix match is case-insensitive")
assert.Equal(t, Internal, cfg.Wing("HyperGuild"))
}
+198
View File
@@ -0,0 +1,198 @@
// Package gitea implements capture.IssueTracker against a Gitea instance
// over its REST API. It is the new outbound dependency the brain server
// gains for the capture capability (#49c/#52): the server otherwise does
// brain-local file ops only.
//
// Owner is hard-coded to the operator and never taken from caller input.
// The API token is read once at construction, held in the struct, and
// never logged or placed in argv — it travels only in the Authorization
// header of outbound requests (AGENTS.md secret-handling).
package gitea
import (
"bytes"
"context"
"encoding/base64"
"encoding/json"
"fmt"
"io"
"net/http"
"strings"
"time"
"github.com/mathiasbq/hyperguild/ingestion/internal/capture"
)
// owner is the fixed repository owner for every ticket operation. It is a
// constant, not a parameter, so a caller can never redirect a write to
// another owner's repo.
const owner = "mathias"
// Client is a Gitea REST API IssueTracker.
type Client struct {
baseURL string
token string
http *http.Client
}
// New constructs a Client. It returns nil when either baseURL or token is
// empty, so callers can treat missing config as "tracker disabled" with a
// single nil check (mirrors embed.New).
func New(baseURL, token string) *Client {
if baseURL == "" || token == "" {
return nil
}
return &Client{
baseURL: strings.TrimRight(baseURL, "/"),
token: token,
http: &http.Client{Timeout: 15 * time.Second},
}
}
// issueResponse is the subset of a Gitea issue/comment payload we read.
type issueResponse struct {
Number int `json:"number"`
HTMLURL string `json:"html_url"`
}
// CreateIssue opens a new issue under the fixed owner.
func (c *Client) CreateIssue(ctx context.Context, repo, title, body string) (capture.IssueRef, error) {
var out issueResponse
if err := c.do(ctx, http.MethodPost,
fmt.Sprintf("/api/v1/repos/%s/%s/issues", owner, repo),
map[string]any{"title": title, "body": body}, &out); err != nil {
return capture.IssueRef{}, err
}
return capture.IssueRef{Repo: repo, Number: out.Number, URL: out.HTMLURL}, nil
}
// CommentIssue posts a comment on an existing issue.
func (c *Client) CommentIssue(ctx context.Context, repo string, number int, body string) (capture.IssueRef, error) {
var out issueResponse
if err := c.do(ctx, http.MethodPost,
fmt.Sprintf("/api/v1/repos/%s/%s/issues/%d/comments", owner, repo, number),
map[string]any{"body": body}, &out); err != nil {
return capture.IssueRef{}, err
}
return capture.IssueRef{Repo: repo, Number: number, URL: out.HTMLURL}, nil
}
// CloseIssue closes an issue, first posting a closing comment when one is
// given (empty comment ⇒ close only).
func (c *Client) CloseIssue(ctx context.Context, repo string, number int, comment string) (capture.IssueRef, error) {
if strings.TrimSpace(comment) != "" {
if _, err := c.CommentIssue(ctx, repo, number, comment); err != nil {
return capture.IssueRef{}, err
}
}
var out issueResponse
if err := c.do(ctx, http.MethodPatch,
fmt.Sprintf("/api/v1/repos/%s/%s/issues/%d", owner, repo, number),
map[string]any{"state": "closed"}, &out); err != nil {
return capture.IssueRef{}, err
}
return capture.IssueRef{Repo: repo, Number: number, URL: out.HTMLURL}, nil
}
// WriteFile creates or updates a file in repo at path via the Gitea
// contents API — the SummaryWriter port (#66). It upserts: a GET resolves
// the current blob sha (if any) so an existing file is updated rather than
// rejected (the richer-fidelity-supersedes rule for re-captured sessions).
// Owner is the fixed const, like every other call.
func (c *Client) WriteFile(ctx context.Context, repo, path, content string) error {
cpath := fmt.Sprintf("/api/v1/repos/%s/%s/contents/%s", owner, repo, path)
sha, err := c.fileSHA(ctx, cpath)
if err != nil {
return err
}
payload := map[string]any{
"message": "capture: " + path,
"content": base64.StdEncoding.EncodeToString([]byte(content)),
}
// Gitea contents API: POST creates a new file, PUT updates an existing
// one (PUT requires the current sha). Pick by whether the file exists.
method := http.MethodPost
if sha != "" {
method = http.MethodPut
payload["sha"] = sha
}
status, body, err := c.request(ctx, method, cpath, payload)
if err != nil {
return err
}
if status < 200 || status >= 300 {
return fmt.Errorf("gitea %s %s: status %d: %s", method, cpath, status, strings.TrimSpace(string(body)))
}
return nil
}
// fileSHA returns the current blob sha for a contents path, or "" when the
// file does not exist (404). Any other non-2xx is an error.
func (c *Client) fileSHA(ctx context.Context, cpath string) (string, error) {
status, body, err := c.request(ctx, http.MethodGet, cpath, nil)
if err != nil {
return "", err
}
if status == http.StatusNotFound {
return "", nil
}
if status < 200 || status >= 300 {
return "", fmt.Errorf("gitea GET %s: status %d: %s", cpath, status, strings.TrimSpace(string(body)))
}
var meta struct {
SHA string `json:"sha"`
}
if err := json.Unmarshal(body, &meta); err != nil {
return "", fmt.Errorf("gitea GET %s: decode: %w", cpath, err)
}
return meta.SHA, nil
}
// do performs a JSON request against the Gitea API and decodes a 2xx
// response into out. Errors carry the status and a truncated body for
// diagnosis but never the token.
func (c *Client) do(ctx context.Context, method, path string, payload any, out *issueResponse) error {
status, body, err := c.request(ctx, method, path, payload)
if err != nil {
return err
}
if status < 200 || status >= 300 {
return fmt.Errorf("gitea %s %s: status %d: %s", method, path, status, strings.TrimSpace(string(body)))
}
if out != nil && len(body) > 0 {
if err := json.Unmarshal(body, out); err != nil {
return fmt.Errorf("gitea %s %s: decode response: %w", method, path, err)
}
}
return nil
}
// request is the shared HTTP path: marshals an optional JSON payload,
// attaches auth (token only ever in the header), and returns the status +
// body so callers can branch on status (e.g. 404) without it being an
// error. Never logs the token.
func (c *Client) request(ctx context.Context, method, path string, payload any) (int, []byte, error) {
var reader io.Reader
if payload != nil {
reqBody, err := json.Marshal(payload)
if err != nil {
return 0, nil, fmt.Errorf("marshal request: %w", err)
}
reader = bytes.NewReader(reqBody)
}
req, err := http.NewRequestWithContext(ctx, method, c.baseURL+path, reader)
if err != nil {
return 0, nil, err
}
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Accept", "application/json")
req.Header.Set("Authorization", "token "+c.token)
resp, err := c.http.Do(req)
if err != nil {
return 0, nil, fmt.Errorf("gitea %s %s: %w", method, path, err)
}
defer func() { _ = resp.Body.Close() }()
body, _ := io.ReadAll(io.LimitReader(resp.Body, 8192))
return resp.StatusCode, body, nil
}
+183
View File
@@ -0,0 +1,183 @@
package gitea_test
import (
"context"
"encoding/json"
"io"
"net/http"
"net/http/httptest"
"strings"
"testing"
"github.com/mathiasbq/hyperguild/ingestion/internal/gitea"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
const testToken = "super-secret-token-value"
func TestNewNilWhenUnconfigured(t *testing.T) {
assert.Nil(t, gitea.New("", testToken))
assert.Nil(t, gitea.New("https://git.example", ""))
}
func TestCreateIssueForcesOwnerAndAuth(t *testing.T) {
var gotPath, gotAuth, gotBody string
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
gotPath = r.URL.Path
gotAuth = r.Header.Get("Authorization")
b, _ := io.ReadAll(r.Body)
gotBody = string(b)
assert.Equal(t, http.MethodPost, r.Method)
w.WriteHeader(http.StatusCreated)
_ = json.NewEncoder(w).Encode(map[string]any{"number": 42, "html_url": "https://git.d-ma.be/mathias/hyperguild/issues/42"})
}))
defer srv.Close()
c := gitea.New(srv.URL, testToken)
require.NotNil(t, c)
ref, err := c.CreateIssue(context.Background(), "hyperguild", "Do the thing", "details")
require.NoError(t, err)
assert.Equal(t, "/api/v1/repos/mathias/hyperguild/issues", gotPath, "owner forced to mathias")
assert.Equal(t, "token "+testToken, gotAuth)
assert.Contains(t, gotBody, "Do the thing")
assert.Equal(t, "hyperguild", ref.Repo)
assert.Equal(t, 42, ref.Number)
assert.Contains(t, ref.URL, "/issues/42")
}
func TestCommentIssue(t *testing.T) {
var gotPath string
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
gotPath = r.URL.Path
w.WriteHeader(http.StatusCreated)
_ = json.NewEncoder(w).Encode(map[string]any{"html_url": "https://git/c/1"})
}))
defer srv.Close()
ref, err := gitea.New(srv.URL, testToken).CommentIssue(context.Background(), "hyperguild", 7, "a comment")
require.NoError(t, err)
assert.Equal(t, "/api/v1/repos/mathias/hyperguild/issues/7/comments", gotPath)
assert.Equal(t, 7, ref.Number)
}
func TestCloseIssueWithComment(t *testing.T) {
var paths []string
var states []string
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
paths = append(paths, r.Method+" "+r.URL.Path)
if r.Method == http.MethodPatch {
var body map[string]any
b, _ := io.ReadAll(r.Body)
_ = json.Unmarshal(b, &body)
states = append(states, body["state"].(string))
}
w.WriteHeader(http.StatusOK)
_ = json.NewEncoder(w).Encode(map[string]any{"number": 9, "html_url": "https://git/i/9"})
}))
defer srv.Close()
ref, err := gitea.New(srv.URL, testToken).CloseIssue(context.Background(), "hyperguild", 9, "closing because done")
require.NoError(t, err)
assert.Equal(t, 9, ref.Number)
// Comment posted first, then state PATCHed to closed.
assert.Contains(t, paths, "POST /api/v1/repos/mathias/hyperguild/issues/9/comments")
assert.Contains(t, paths, "PATCH /api/v1/repos/mathias/hyperguild/issues/9")
assert.Equal(t, []string{"closed"}, states)
}
func TestCloseIssueNoComment(t *testing.T) {
var commented bool
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if strings.HasSuffix(r.URL.Path, "/comments") {
commented = true
}
w.WriteHeader(http.StatusOK)
_ = json.NewEncoder(w).Encode(map[string]any{"number": 3, "html_url": "https://git/i/3"})
}))
defer srv.Close()
_, err := gitea.New(srv.URL, testToken).CloseIssue(context.Background(), "hyperguild", 3, "")
require.NoError(t, err)
assert.False(t, commented, "empty comment ⇒ no comment POST")
}
func TestErrorPathDoesNotLeakToken(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusInternalServerError)
_, _ = w.Write([]byte("boom"))
}))
defer srv.Close()
_, err := gitea.New(srv.URL, testToken).CreateIssue(context.Background(), "hyperguild", "t", "b")
require.Error(t, err)
assert.NotContains(t, err.Error(), testToken, "token must never appear in an error message")
assert.Contains(t, err.Error(), "500")
}
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")
}
+80 -58
View File
@@ -11,6 +11,7 @@ import (
"github.com/mathiasbq/hyperguild/ingestion/internal/api" "github.com/mathiasbq/hyperguild/ingestion/internal/api"
"github.com/mathiasbq/hyperguild/ingestion/internal/brain" "github.com/mathiasbq/hyperguild/ingestion/internal/brain"
"github.com/mathiasbq/hyperguild/ingestion/internal/capture"
"github.com/mathiasbq/hyperguild/ingestion/internal/extract" "github.com/mathiasbq/hyperguild/ingestion/internal/extract"
"github.com/mathiasbq/hyperguild/ingestion/internal/graphsync" "github.com/mathiasbq/hyperguild/ingestion/internal/graphsync"
"github.com/mathiasbq/hyperguild/ingestion/internal/pipeline" "github.com/mathiasbq/hyperguild/ingestion/internal/pipeline"
@@ -38,7 +39,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).",
@@ -81,6 +82,21 @@ func (s *Server) tools() []map[string]any {
"path": str("brain-relative path to the note; equivalent to id"), "path": str("brain-relative path to the note; equivalent to id"),
}), }),
}, },
{
"name": "brain_pending",
"description": "List notes in brain/raw/ awaiting human promotion to the wiki, oldest-first. Returns filename, created_at, size_bytes, excerpt. The human-review queue complement to brain_promote.",
"inputSchema": schema([]string{}, map[string]any{}),
},
{
"name": "brain_promote",
"description": "Promote a brain/raw/ note into brain/wiki/<wing>/<hall>/: rewrites frontmatter (sets wing/hall/promoted_at, preserves created_at + custom fields), deletes the source, rebuilds the wing index, runs auto-tunnel. Errors (without touching the fs) on invalid hall or a slug collision. Returns {path}.",
"inputSchema": schema([]string{"filename", "wing", "hall"}, map[string]any{
"filename": str("basename in brain/raw/, e.g. 2026-06-01-lejpa-decision.md"),
"wing": str("target wing, e.g. jepa-fx"),
"hall": enum("target hall", halls...),
"slug": str("optional target slug; defaults to filename minus date prefix"),
}),
},
{ {
"name": "brain_tunnel", "name": "brain_tunnel",
"description": "Create an explicit bidirectional [[wikilink]] between two notes in different wings. Idempotent.", "description": "Create an explicit bidirectional [[wikilink]] between two notes in different wings. Idempotent.",
@@ -171,6 +187,13 @@ func (s *Server) tools() []map[string]any {
}), }),
}, },
} }
// The capture relay tool (#55) is advertised only when wired via
// WithCapture — MCP-native harnesses (claude.ai, Crush, Pi, LLM Council)
// reach capture through it.
if s.capture != nil {
tools = append(tools, captureToolDescriptor())
}
return tools
} }
type brainQueryArgs struct { type brainQueryArgs struct {
@@ -219,7 +242,11 @@ func (s *Server) brainWrite(ctx context.Context, args json.RawMessage) (json.Raw
if err := json.Unmarshal(args, &a); err != nil { if err := json.Unmarshal(args, &a); err != nil {
return nil, fmt.Errorf("parse args: %w", err) return nil, fmt.Errorf("parse args: %w", err)
} }
relPath, err := api.WriteNote(s.brainDir, api.WriteNoteOptions{ // Delegate to the shared BrainStore so write+index+tunnel+graph live in
// one implementation (capture uses the same store). The read-after-write
// handle {id, path, content_hash} comes back from the store; path is kept
// for backward compatibility.
ref, err := s.store.Write(ctx, capture.Note{
Content: a.Content, Content: a.Content,
Filename: a.Filename, Filename: a.Filename,
Type: a.Type, Type: a.Type,
@@ -230,22 +257,7 @@ func (s *Server) brainWrite(ctx context.Context, args json.RawMessage) (json.Raw
if err != nil { if err != nil {
return nil, err return nil, err
} }
// Auto-regenerate the wing _index.md when the write landed in the return json.Marshal(map[string]string{"id": ref.ID, "path": ref.Path, "content_hash": ref.ContentHash})
// structured wiki, and auto-tunnel cross-wing matches. Both are
// best-effort: the note is already written.
if a.Wing != "" && a.Hall != "" {
if err := brain.BuildWingIndex(s.brainDir, a.Wing); err != nil {
slog.Warn("brain_write: auto-index failed", "wing", a.Wing, "err", err)
}
if err := brain.AutoTunnel(s.brainDir, relPath, a.Content); err != nil {
slog.Warn("brain_write: auto-tunnel failed", "src", relPath, "err", err)
}
}
s.indexInGraph(ctx, "brain_write", relPath)
// Read-after-write handle: id == relPath, content_hash == sha256 of
// the bytes just written. path is kept for backward compatibility.
_, _, hash, _ := api.ReadNote(s.brainDir, relPath)
return json.Marshal(map[string]string{"id": relPath, "path": relPath, "content_hash": hash})
} }
type brainUpdateArgs struct { type brainUpdateArgs struct {
@@ -274,50 +286,26 @@ func (s *Server) brainUpdate(ctx context.Context, args json.RawMessage) (json.Ra
return nil, fmt.Errorf("content is required") return nil, fmt.Errorf("content is required")
} }
opts := api.UpdateNoteOptions{Content: a.Content, Reason: a.Reason} // path takes precedence over slug; the store treats any slug containing
switch { // a slash as a full brain-relative path (issue #45: "slug ... OR path").
case a.Path != "": slug := a.Slug
opts.Path = a.Path if a.Path != "" {
case strings.Contains(a.Slug, "/"): slug = a.Path
// slug carries a full path (issue #45: "slug ... OR full path").
opts.Path = a.Slug
default:
opts.Wing, opts.Hall, opts.Slug = a.Wing, a.Hall, a.Slug
} }
ref, err := s.store.Update(ctx, slug, capture.Note{
relPath, hash, _, err := api.UpdateNote(s.brainDir, opts) Content: a.Content,
Wing: a.Wing,
Hall: a.Hall,
Reason: a.Reason,
})
if err != nil { if err != nil {
return nil, err return nil, err
} }
// Best-effort wiki upkeep, mirroring brain_write: rebuild the wing
// _index and re-tunnel cross-wing matches against the new body. Both
// are idempotent and never block — the note is already superseded.
if wing := wingFromRelPath(relPath); wing != "" {
if err := brain.BuildWingIndex(s.brainDir, wing); err != nil {
slog.Warn("brain_update: auto-index failed", "wing", wing, "err", err)
}
if err := brain.AutoTunnel(s.brainDir, relPath, a.Content); err != nil {
slog.Warn("brain_update: auto-tunnel failed", "src", relPath, "err", err)
}
}
s.indexInGraph(ctx, "brain_update", relPath)
return json.Marshal(map[string]any{ return json.Marshal(map[string]any{
"id": relPath, "path": relPath, "content_hash": hash, "superseded": true, "id": ref.ID, "path": ref.Path, "content_hash": ref.ContentHash, "superseded": ref.Superseded,
}) })
} }
// wingFromRelPath extracts the wing segment from a structured wiki path
// (wiki/<wing>/<hall>/<slug>.md). Returns "" for legacy/non-wiki paths.
func wingFromRelPath(relPath string) string {
parts := strings.Split(relPath, "/")
if len(parts) >= 4 && parts[0] == "wiki" {
return parts[1]
}
return ""
}
type brainGetArgs struct { type brainGetArgs struct {
ID string `json:"id,omitempty"` ID string `json:"id,omitempty"`
Path string `json:"path,omitempty"` Path string `json:"path,omitempty"`
@@ -327,7 +315,7 @@ type brainGetArgs struct {
// path — the de-facto handle). Read-only; the create-path read-after- // path — the de-facto handle). Read-only; the create-path read-after-
// write primitive that lets callers confirm a write landed without a // write primitive that lets callers confirm a write landed without a
// lexical re-query. // lexical re-query.
func (s *Server) brainGet(_ context.Context, args json.RawMessage) (json.RawMessage, error) { func (s *Server) brainGet(ctx context.Context, args json.RawMessage) (json.RawMessage, error) {
var a brainGetArgs var a brainGetArgs
if err := json.Unmarshal(args, &a); err != nil { if err := json.Unmarshal(args, &a); err != nil {
return nil, fmt.Errorf("parse args: %w", err) return nil, fmt.Errorf("parse args: %w", err)
@@ -339,16 +327,50 @@ func (s *Server) brainGet(_ context.Context, args json.RawMessage) (json.RawMess
if target == "" { if target == "" {
return nil, fmt.Errorf("id or path is required") return nil, fmt.Errorf("id or path is required")
} }
fm, body, hash, err := api.ReadNote(s.brainDir, target) note, err := s.store.Get(ctx, target)
if err != nil { if err != nil {
return nil, err return nil, err
} }
return json.Marshal(map[string]any{ return json.Marshal(map[string]any{
"id": target, "path": target, "content_hash": hash, "id": note.ID, "path": note.Path, "content_hash": note.ContentHash,
"frontmatter": fm, "body": body, "frontmatter": note.Frontmatter, "body": note.Body,
}) })
} }
// brainPending lists the raw/ review queue (oldest-first).
func (s *Server) brainPending(_ context.Context, _ json.RawMessage) (json.RawMessage, error) {
pending, err := api.ListPending(s.brainDir)
if err != nil {
return nil, err
}
return json.Marshal(map[string]any{"pending": pending})
}
type brainPromoteArgs struct {
Filename string `json:"filename"`
Wing string `json:"wing"`
Hall string `json:"hall"`
Slug string `json:"slug,omitempty"`
}
// brainPromote moves a raw/ note into the structured wiki (frontmatter
// rewrite + index + auto-tunnel, all owned by api.PromoteNote) and then
// re-indexes it into the graph. The human-facing complement to brain_write.
func (s *Server) brainPromote(ctx context.Context, args json.RawMessage) (json.RawMessage, error) {
var a brainPromoteArgs
if err := json.Unmarshal(args, &a); err != nil {
return nil, fmt.Errorf("parse args: %w", err)
}
relPath, err := api.PromoteNote(s.brainDir, api.PromoteOptions{
Filename: a.Filename, Wing: a.Wing, Hall: a.Hall, Slug: a.Slug,
})
if err != nil {
return nil, err
}
s.indexInGraph(ctx, "brain_promote", relPath)
return json.Marshal(map[string]string{"path": relPath})
}
// indexInGraph is a best-effort wrapper around graphsync.IndexDoc that // indexInGraph is a best-effort wrapper around graphsync.IndexDoc that
// logs failures but never propagates them — the underlying write/ingest // logs failures but never propagates them — the underlying write/ingest
// has already succeeded and the graph is an augmentation, not a // has already succeeded and the graph is an augmentation, not a
+58
View File
@@ -332,3 +332,61 @@ func TestSessionLogRequiresSessionID(t *testing.T) {
resp := toolCall(t, srv, "session_log", map[string]any{"skill": "tdd"}) resp := toolCall(t, srv, "session_log", map[string]any{"skill": "tdd"})
require.NotNil(t, resp["error"]) require.NotNil(t, resp["error"])
} }
func TestBrainPendingListsRaw(t *testing.T) {
brainDir := t.TempDir()
raw := filepath.Join(brainDir, "raw")
require.NoError(t, os.MkdirAll(raw, 0o755))
require.NoError(t, os.WriteFile(filepath.Join(raw, "2026-06-01-x.md"),
[]byte("---\ncreated_at: 2026-06-01T00:00:00Z\n---\npending body\n"), 0o644))
srv := mcp.NewServer(brainDir, nil, nil, nil)
resp := toolCall(t, srv, "brain_pending", map[string]any{})
require.Nil(t, resp["error"])
text := resp["result"].(map[string]any)["content"].([]any)[0].(map[string]any)["text"].(string)
assert.Contains(t, text, "2026-06-01-x.md")
assert.Contains(t, text, "pending body")
}
func TestBrainPendingEmpty(t *testing.T) {
srv := mcp.NewServer(t.TempDir(), nil, nil, nil)
resp := toolCall(t, srv, "brain_pending", map[string]any{})
require.Nil(t, resp["error"])
text := resp["result"].(map[string]any)["content"].([]any)[0].(map[string]any)["text"].(string)
assert.Contains(t, text, `"pending":[]`)
}
func TestBrainPromoteMovesToWiki(t *testing.T) {
brainDir := t.TempDir()
raw := filepath.Join(brainDir, "raw")
require.NoError(t, os.MkdirAll(raw, 0o755))
require.NoError(t, os.WriteFile(filepath.Join(raw, "2026-06-01-decision.md"),
[]byte("---\ncreated_at: 2026-06-01T00:00:00Z\n---\n# D\n\nbody\n"), 0o644))
srv := mcp.NewServer(brainDir, nil, nil, nil)
resp := toolCall(t, srv, "brain_promote", map[string]any{
"filename": "2026-06-01-decision.md", "wing": "jepa-fx", "hall": "decisions",
})
require.Nil(t, resp["error"], "got: %v", resp["error"])
text := resp["result"].(map[string]any)["content"].([]any)[0].(map[string]any)["text"].(string)
assert.Contains(t, text, "wiki/jepa-fx/decisions/decision.md")
_, err := os.Stat(filepath.Join(brainDir, "wiki/jepa-fx/decisions/decision.md"))
require.NoError(t, err)
_, srcErr := os.Stat(filepath.Join(raw, "2026-06-01-decision.md"))
assert.True(t, os.IsNotExist(srcErr), "source removed")
}
func TestBrainPromoteInvalidHallErrors(t *testing.T) {
brainDir := t.TempDir()
raw := filepath.Join(brainDir, "raw")
require.NoError(t, os.MkdirAll(raw, 0o755))
require.NoError(t, os.WriteFile(filepath.Join(raw, "x.md"), []byte("body\n"), 0o644))
srv := mcp.NewServer(brainDir, nil, nil, nil)
resp := toolCall(t, srv, "brain_promote", map[string]any{
"filename": "x.md", "wing": "a", "hall": "garbage",
})
require.NotNil(t, resp["error"])
_, srcErr := os.Stat(filepath.Join(raw, "x.md"))
assert.NoError(t, srcErr, "source untouched on validation error")
}
+93 -4
View File
@@ -1,7 +1,9 @@
// 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_pending, brain_promote, brain_index, brain_tunnel, brain_ingest,
// brain_answer, brain_classify, brain_graph, brain_context, session_log. // brain_ingest_raw, brain_answer, brain_classify, brain_graph,
// brain_context, session_log, and capture (the #55 relay tool, registered
// only when WithCapture is set).
package mcp package mcp
import ( import (
@@ -10,6 +12,9 @@ import (
"fmt" "fmt"
"net/http" "net/http"
"github.com/mathiasbq/hyperguild/ingestion/internal/brainstore"
"github.com/mathiasbq/hyperguild/ingestion/internal/capture"
"github.com/mathiasbq/hyperguild/ingestion/internal/capturehttp"
"github.com/mathiasbq/hyperguild/ingestion/internal/graphstore" "github.com/mathiasbq/hyperguild/ingestion/internal/graphstore"
"github.com/mathiasbq/hyperguild/ingestion/internal/graphsync" "github.com/mathiasbq/hyperguild/ingestion/internal/graphsync"
"github.com/mathiasbq/hyperguild/ingestion/internal/pipeline" "github.com/mathiasbq/hyperguild/ingestion/internal/pipeline"
@@ -46,6 +51,21 @@ type Server struct {
vector search.VectorSearcher // nil = BM25-only retrieval vector search.VectorSearcher // nil = BM25-only retrieval
embedder search.Embedder // nil = BM25-only retrieval embedder search.Embedder // nil = BM25-only retrieval
graph graphsync.Store // nil = brain_graph and GraphRAG augmentation disabled graph graphsync.Store // nil = brain_graph and GraphRAG augmentation disabled
store *brainstore.Store // shared brain write/update/get impl (also used by capture)
tracker capture.IssueTracker // nil = no Gitea ticket integration; wired for capture (#53)
capture *captureDeps // nil = capture MCP tool disabled (#55 relay)
}
// captureDeps holds what the MCP `capture` tool (the #55 relay door for
// MCP-native harnesses like claude.ai) needs: the use-case, the auth bits
// to re-derive the caller's principal from the Bearer header (the chassis
// middleware gates but discards the principal), and the origin resolver.
type captureDeps struct {
svc *capture.Service
validator capturehttp.Validator
staticToken string
staticPrincipal string
resolver capturehttp.OriginResolver
} }
// NewServer constructs a Server bound to brainDir. pipelineCfg supplies the // NewServer constructs a Server bound to brainDir. pipelineCfg supplies the
@@ -56,7 +76,13 @@ func NewServer(brainDir string, pipelineCfg *pipeline.Config, llm pipeline.Compl
if pipelineCfg != nil { if pipelineCfg != nil {
cfg = *pipelineCfg cfg = *pipelineCfg
} }
return &Server{brainDir: brainDir, pipeline: cfg, llm: llm, answerLLM: answerLLM} return &Server{
brainDir: brainDir,
pipeline: cfg,
llm: llm,
answerLLM: answerLLM,
store: brainstore.New(brainDir),
}
} }
// WithReranker installs an opt-in cross-encoder reranker. When set, // WithReranker installs an opt-in cross-encoder reranker. When set,
@@ -84,9 +110,55 @@ func (s *Server) WithHybridRetrieval(v search.VectorSearcher, e search.Embedder)
func (s *Server) WithGraph(g *graphstore.PGStore) *Server { func (s *Server) WithGraph(g *graphstore.PGStore) *Server {
if g == nil { if g == nil {
s.graph = nil s.graph = nil
s.store.WithGraph(nil)
return s return s
} }
s.graph = g s.graph = g
s.store.WithGraph(g)
return s
}
// WithIssueTracker injects the Gitea ticket tracker behind the
// capture.IssueTracker interface. nil leaves ticket integration off. The
// use-case (capture) consumes this in #53; it is wired here so the
// dependency is constructed once and stays swappable/testable.
func (s *Server) WithIssueTracker(t capture.IssueTracker) *Server {
s.tracker = t
return s
}
// IssueTracker returns the injected ticket tracker (nil when unconfigured).
func (s *Server) IssueTracker() capture.IssueTracker {
return s.tracker
}
// BrainStore returns the shared brain store (graph-wired once WithGraph
// has run), so the capture use-case writes through the exact same
// implementation as the MCP handlers.
func (s *Server) BrainStore() *brainstore.Store {
return s.store
}
// WithCapture enables the MCP `capture` tool (#55) — the relay door for
// MCP-native harnesses (claude.ai, Crush, Pi, LLM Council) that cannot run
// the use-case in-process. It forwards to the same CaptureService as
// POST /capture, deriving the caller's principal + origin from the same
// auth credentials that gate /mcp. nil svc leaves the tool unregistered.
func (s *Server) WithCapture(svc *capture.Service, validator capturehttp.Validator, staticToken, staticPrincipal string, resolver capturehttp.OriginResolver) *Server {
if svc == nil {
s.capture = nil
return s
}
if staticPrincipal == "" {
staticPrincipal = "local-cli"
}
s.capture = &captureDeps{
svc: svc,
validator: validator,
staticToken: staticToken,
staticPrincipal: staticPrincipal,
resolver: resolver,
}
return s return s
} }
@@ -139,7 +211,18 @@ func (s *Server) ServeHTTP(w http.ResponseWriter, r *http.Request) {
rpcErr = &rpcError{Code: -32602, Message: "invalid params"} rpcErr = &rpcError{Code: -32602, Message: "invalid params"}
break break
} }
out, err := s.handleCall(r.Context(), p.Name, p.Arguments) // Re-derive the authenticated principal from the Bearer header so
// the capture tool can compute the trust-zone origin. The request
// is already gated by BearerMiddleware; this only recovers the
// identity that middleware discards.
ctx := r.Context()
if s.capture != nil {
if principal, viaStatic, ok := capturehttp.Authenticate(
r, s.capture.staticToken, s.capture.staticPrincipal, s.capture.validator); ok {
ctx = withPrincipal(ctx, principal, viaStatic)
}
}
out, err := s.handleCall(ctx, p.Name, p.Arguments)
if err != nil { if err != nil {
rpcErr = &rpcError{Code: -32000, Message: err.Error()} rpcErr = &rpcError{Code: -32000, Message: err.Error()}
break break
@@ -181,6 +264,12 @@ 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 "brain_pending":
return s.brainPending(ctx, args)
case "brain_promote":
return s.brainPromote(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":
+22
View File
@@ -2,12 +2,14 @@ package mcp_test
import ( import (
"bytes" "bytes"
"context"
"encoding/json" "encoding/json"
"net/http" "net/http"
"net/http/httptest" "net/http/httptest"
"strings" "strings"
"testing" "testing"
"github.com/mathiasbq/hyperguild/ingestion/internal/capture"
"github.com/mathiasbq/hyperguild/ingestion/internal/mcp" "github.com/mathiasbq/hyperguild/ingestion/internal/mcp"
"github.com/stretchr/testify/assert" "github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require" "github.com/stretchr/testify/require"
@@ -56,6 +58,7 @@ func TestServerToolsList(t *testing.T) {
} }
assert.ElementsMatch(t, []string{ assert.ElementsMatch(t, []string{
"brain_query", "brain_write", "brain_update", "brain_get", "brain_query", "brain_write", "brain_update", "brain_get",
"brain_pending", "brain_promote",
"brain_index", "brain_tunnel", "brain_index", "brain_tunnel",
"brain_ingest_raw", "brain_ingest", "brain_ingest_raw", "brain_ingest",
"brain_answer", "brain_classify", "brain_graph", "brain_context", "brain_answer", "brain_classify", "brain_graph", "brain_context",
@@ -93,3 +96,22 @@ func TestServerUnknownMethodReturnsError(t *testing.T) {
assert.Equal(t, float64(-32601), errObj["code"]) assert.Equal(t, float64(-32601), errObj["code"])
assert.Contains(t, errObj["message"].(string), "unknown/method") assert.Contains(t, errObj["message"].(string), "unknown/method")
} }
type stubTracker struct{}
func (stubTracker) CreateIssue(context.Context, string, string, string) (capture.IssueRef, error) {
return capture.IssueRef{}, nil
}
func (stubTracker) CloseIssue(context.Context, string, int, string) (capture.IssueRef, error) {
return capture.IssueRef{}, nil
}
func (stubTracker) CommentIssue(context.Context, string, int, string) (capture.IssueRef, error) {
return capture.IssueRef{}, nil
}
func TestWithIssueTrackerInjects(t *testing.T) {
srv := mcp.NewServer(t.TempDir(), nil, nil, nil)
assert.Nil(t, srv.IssueTracker(), "tracker is off by default")
srv = srv.WithIssueTracker(stubTracker{})
assert.NotNil(t, srv.IssueTracker(), "tracker injected behind the interface")
}
+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")
}
+79
View File
@@ -0,0 +1,79 @@
---
name: close-session
description: Disciplined end-of-session closeout for a Claude.ai chat before archiving it. Harvests the session's decisions, artifacts, and open threads and durably persists them to the brain MCP and the right Gitea repo so nothing is lost when context resets. Use this whenever the user signals they are wrapping up — phrases like "close this out", "let's wrap up", "before I archive", "session retro", "capture this before I go", "did we lose anything", or any end-of-session/handoff cue — even if they don't say the word "close". Also use when the user explicitly asks to retro, archive, or hand off a working session.
---
# close-session
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. Run the phases in order. Stop at any confirmation gate that says STOP.
## Operating constraints (read first)
- **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.
- **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 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.
## Phase 1 — Harvest
Reconstruct what actually happened this session from the conversation itself. Produce, in working memory:
- **Decisions taken** — what was decided and the reasoning, not just the outcome.
- **Artifacts produced** — issues filed/closed, PRs opened/merged, files committed, brain notes written, ADRs. Capture identifiers (issue numbers, PR numbers, paths, commit SHAs) as you go.
- **Open threads** — what was deferred, what's blocked, what the next session should pick up.
- **Generalizable learnings** — reusable patterns or footguns that would bite anyone again (these are brain-worthy; project status is not).
Be honest about fidelity: a long session compresses harder at the start than the end. Flag anything you're reconstructing rather than certain of.
## Phase 2 — Ground-truth Gitea state
For every repo touched this session, get its true current state before proposing any change. `gitea:repo_status` (owner `mathias`) gives branches + open PRs + protection in one call. For each issue you intend to close, comment on, or reference: `gitea:issue_get` it fresh and compare to what the session assumed. Note any drift (closed-already, body rotted, renamed) — you'll surface it in Phase 3.
Do not write anything in this phase. This is the read pass.
## Phase 3 — Plan the issue changes
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.
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 — Capture (one uniform call)
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).
**Assemble one payload:**
- **`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>}`.
**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)**.
- **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.
- Declare `context.classification: "internal"` for normal homelab work.
- `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.
- 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.
**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.
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.
## Phase 5 — Verdict
Deliver a final "safe to archive" verdict in the chat. Either:
- **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, the capture refusal reason, or the per-item error from the receipt, and what to do about it.
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 gate and the single-call shape matter
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.
+248
View File
@@ -0,0 +1,248 @@
# Capture capability — use-case & BDD specification
**Status:** Decisions resolved 2026-06-22 (§4). Ready for implementation scoping. `capture` is a
privileged cross-harness write path touching brain + Gitea + ai-sessions.
**Tracks:** hyperguild #49.
**Governed by:** `infra/docs/architecture/01-invariants.md` (I1I5), the admissibility test in
`00-synthesis-model.md`, and the distributed-consolidation shape mandated by
`brain/wiki/homelab/decisions/no-centralized-cross-harness-observer-2026-06-17.md`.
---
## 1. Use-case (Clean Architecture form)
**Name:** CaptureSession
**Actor:** A harness acting on the user's behalf (claude.ai Chat/Cowork/Code/Design, Claude Code
CLI, Crush, Pi, LLM Council, Agentsquad executor/reviewer) — or the user directly.
**Goal:** Durably persist a finished session's valuable output — insights → brain, action items →
Gitea tickets, optional summary → ai-sessions — with one uniform invocation, identical core
behaviour across harnesses.
**Primary success scenario (essential steps):**
1. Caller assembles capture input (insights, tickets, optional summary) + context (harness,
session_ref, fidelity, actor, **data-classification**).
2. System validates the whole request (fail-closed).
3. System resolves **effective classification** (stricter of caller-declared and target-derived)
and the **server-derived harness origin** (from the authenticated principal). It checks the
**sovereignty gate** (I1): if effective classification is confidential AND the origin is a
non-sovereign (us-nexus) surface, the capture is **refused** before any write.
4. System persists insights (write or supersede), tickets (create/close/comment), summary — each
best-effort, recording per-item outcome.
5. System emits an **audit record** (I5) of who/what captured what, when, via which principal.
6. System returns a structured, partial-aware receipt.
**Architectural shape:** the *logic* is a shared use-case (`CaptureService`), invoked **per-harness
against the caller's own credentials** (distributed consolidation — no high-degree observer node).
A central authenticated relay endpoint exists ONLY as a fallback for harnesses that cannot run the
use-case in-process (Crush/Pi/headless); the relay holds no standing visibility and retains nothing
beyond the I5 audit log.
---
## 2. Invariant obligations (acceptance gates, not nice-to-haves)
| Invariant | Obligation on `capture` |
|---|---|
| **I1 sovereign containment** | A confidential-classified session MUST NOT be captured through a us-nexus harness. Harness origin is **server-derived from the authenticated principal** (not caller-asserted). Classification uses **model (C)**: caller declares, server cross-checks the target's tag, **stricter wins**, mismatch logged. See §4.14.2. |
| **I2 deliberate acceptance** | The *distributed-library* form opens no new acceptance. IF a central relay node is deployed, its cross-harness reach MUST be entered in `infra/docs/security-baseline.md` with Why-accepted / Revisit-if before it ships. |
| **I3 GitOps reconcilability** | IF `capture` runs as a deployed service, its manifest lives under `infra/k3s/apps/**` (sovereign source, Flux-reconciled). No untracked runtime. |
| **I4 decisions captured** | The distributed-vs-central decision and the intent-named-verb pattern are recorded (ADR + brain). |
| **I5 auditability** | Every capture emits a request-level audit record (actor/principal, harness, items written, timestamp) to the alloy/loki substrate. **Classification-aware degradation** (§4.4): confidential + sink-down → hard-refuse; internal/public + sink-down → durable local buffer + ntfy + reconcile. Floor: refuse if nothing can record the audit. |
---
## 3. BDD scenarios (Gherkin)
```gherkin
Feature: Capture session value uniformly across harnesses
As an operator working across many AI harnesses
I want one uniform command to persist insights and file tickets
So that valuable session output is never lost and is always auditable
Background:
Given a brain store, a Gitea issue tracker, and an ai-sessions summary writer
And the caller is authenticated with a principal
And the session context declares a harness, a fidelity, and a data classification
# --- Core happy path ---
Scenario: Capture insights and tickets from a non-confidential session
Given a session classified as "internal"
And the capture input has 2 insights and 1 ticket to create
When capture is invoked
Then both insights are written to the brain and their ids and content hashes are returned
And the ticket is created in the named repo under owner "mathias"
And an audit record is emitted naming the principal, harness, and items written
And the receipt reports every item as ok
# --- I1: sovereignty gate (the load-bearing refusal) ---
# Harness origin is server-derived from the authenticated principal, never from context.harness.
Scenario: Refuse capture of a confidential session through a us-nexus harness
Given a session whose effective classification is "confidential"
And the authenticated principal resolves to a us-nexus harness origin
When capture is invoked
Then the capture is refused before any write
And no insight, ticket, or summary is persisted
And the refusal names the sovereignty invariant as the reason
Scenario: Allow capture of a confidential session through a sovereign harness
Given a session whose effective classification is "confidential"
And the authenticated principal resolves to a sovereign-soil harness origin
When capture is invoked
Then the capture proceeds and persists normally
Scenario: Ignore a caller-asserted harness label and use the server-derived origin
Given the request context asserts harness "sovereign-soil"
But the authenticated principal resolves to a us-nexus origin
And the session classification is "confidential"
When capture is invoked
Then the capture is refused
And the server-derived origin is used, not the asserted label
And the asserted-vs-derived discrepancy is logged as a security event
# --- I1: classification model (C) — stricter of declared vs target-derived wins ---
Scenario: Take the stricter classification when caller and target disagree
Given the caller declares classification "internal"
But the target wing/repo is tagged "confidential"
When capture is invoked
Then the effective classification is "confidential"
And the declared-vs-derived mismatch is logged as a security event
And the I1 gate is evaluated against "confidential"
Scenario: Honour a caller raising sensitivity above the target's tag
Given the caller declares classification "confidential"
And the target wing/repo is tagged "internal"
When capture is invoked
Then the effective classification is "confidential"
And the capture is gated as confidential
# --- Supersession + staleness discipline (reuses #45 / #47 resolution) ---
Scenario: Supersede a prior insight rather than duplicating it
Given an insight whose context names an existing note to supersede
When capture is invoked
Then the existing note is updated in place, not duplicated
And the prior content hash is recorded in the superseding note
And read-after-write confirmation uses a direct fetch, never a semantic query
# --- Validation: fail-closed ---
Scenario: Reject a malformed request before any write
Given a capture input with an invalid wing/hall or unknown repo
When capture is invoked
Then the request is rejected with a validation error
And nothing is written to the brain, Gitea, or ai-sessions
# --- Partial failure: best-effort + honest receipt ---
Scenario: Report partial success when one item fails mid-capture
Given a capture input with 2 insights and 1 ticket
And the second insight write will fail
When capture is invoked
Then the first insight and the ticket are persisted
And the second insight is reported as failed in the receipt
And no rollback is attempted
And the audit record reflects exactly what landed
# --- Dry run ---
Scenario: Preview a capture without writing
Given a valid capture input with dry_run true
When capture is invoked
Then the would-be receipt is returned
And nothing is written anywhere
# --- I5: auditability is classification-aware (confidential fails closed) ---
Scenario: Confidential capture hard-refuses when the central audit sink is down
Given the effective classification is "confidential"
And the central audit substrate (loki) cannot be written to
When capture is invoked
Then the capture is refused before any write
And the reason names the auditability invariant
# Confidential work must be centrally auditable at write time — no buffered exception.
Scenario: Internal capture degrades to a durable local buffer when the sink is down
Given the effective classification is "internal" or "public"
And the central audit substrate (loki) cannot be written to
When capture is invoked
Then the capture proceeds
And the audit record is written to a durable LOCAL fallback buffer
And an ntfy alert is emitted naming the degraded audit state
And the receipt flags that audit was buffered locally, not centrally recorded
Scenario: Locally buffered audit records reconcile to the central sink on recovery
Given internal-tier audit records were buffered locally during a sink outage
When the central audit substrate becomes reachable again
Then the buffered records are replayed to the central sink
And the local buffer is cleared only after confirmed central write
Scenario: Even internal capture refuses if neither sink nor local buffer can be written
Given the effective classification is "internal" or "public"
And neither the central sink nor the local fallback buffer can be written
When capture is invoked
Then the capture is refused
And the reason names the auditability invariant
# Degrade-and-warn has a floor: if NOTHING can record the audit, do not write.
# --- Summary fidelity (collision rule from the retro work) ---
Scenario: A richer-fidelity summary supersedes a thinner one for the same session
Given a summary already exists for session_ref X at fidelity "live-capture"
And a new summary arrives for session_ref X at fidelity "transcript-parse"
When capture is invoked
Then the transcript-parse summary supersedes the live-capture one
And the live-capture summary is not left as a contradicting duplicate
```
---
## 4. Resolved decisions (2026-06-22)
These were open questions at draft; resolved in the 2026-06-22 review session. Recorded here as
binding design decisions for the build.
1. **Classification trust — model (C): caller-declares + server-cross-checks, stricter wins.**
The caller declares `context.classification`; the server **independently derives** the target's
classification (from the target wing/repo's classification tag) and gates on the **stricter of
the two**. The caller can voluntarily *raise* sensitivity but can never *lower* it below the
target's floor. A declared-vs-derived **mismatch is logged as a security event** (I5).
- **Prerequisite (new build work):** a classification taxonomy (e.g. `public` /
`internal` / `confidential`) and a per-wing / per-repo classification tag the server can read.
This must exist before the I1 gate is load-bearing. Tracked as a sub-task of #49.
- **Implemented (#50):** taxonomy `public < internal < confidential` (ordered so "stricter wins"
is `max`) in `ingestion/internal/classification/`. Tags are read from an optional
`classification.yaml` at the brain root (`wings:` / `repos:` maps); absent entries fall to
built-in defaults (`client-*` → confidential; `hyperguild`/`homelab` → internal; everything
else → **confidential, fail-safe**). `Config.Derive(Target)` is the function the use-case
calls. See brain `wiki/hyperguild/decisions/capture-classification-taxonomy`.
- Rationale: composes with decision 2; fails safe; honours a caller flagging something *more*
sensitive than its destination. Pure caller-trust (A) was rejected — it makes the gate theatre.
2. **Sovereign-harness determination — server-derived, not caller-asserted.**
"Is this harness us-nexus / sovereign?" is derived from the **authenticated principal/origin**
(the OAuth2 identity), never from `context.harness`. `context.harness` survives only as a
self-reported label for the audit log — descriptive telemetry, **never a gate input**. A control
keyed on an attacker-suppliable value is not a control.
3. **Central relay — ships in v1, with the I2 ledger entry.**
The relay is required, not optional: claude.ai (Chat/Cowork/Design), Crush, Pi, and LLM Council
cannot run the use-case library in-process, and those are primary day-to-day surfaces. Deferring
the relay would ship a capability that doesn't work from the interfaces actually in use. Because
the relay is a (thin, no-standing-visibility, audit-only-retention) central node, its cross-harness
reach **must be entered in `infra/docs/security-baseline.md`** with Why-accepted / Revisit-if
**before it ships** (I2). That ledger entry is v1 work, not a follow-up.
4. **Audit-sink-down — classification-aware: confidential fails closed, internal/public degrades.**
The posture inherits from the effective classification (decision 1), so there is one coherent
sensitivity model rather than a separate availability policy:
- **Confidential + central audit sink unreachable → hard-refuse.** No buffer, no proceed.
Confidential work must be centrally auditable *at write time*; "buffer and reconcile later"
introduces a buffer-integrity question (can a write tamper with its own pending audit record?)
that must not exist for confidential data. The simplicity of "refuse" is itself the assurance
asset — trivially true, nothing to poke holes in.
- **Internal / public + central sink unreachable → degrade-and-warn** with a durable local buffer
+ ntfy alert + reconcile-on-recovery (the earlier Q4 design, now scoped to lower tiers). Keeps
capture available for your own homelab work during an observability outage; negligible risk
since the buffered record is still durable and the data isn't client-confidential.
- **Floor (all tiers):** if *nothing* — neither central sink nor (for internal/public) the local
buffer — can record the audit, capture **refuses**. No tier writes wholly un-audited.
- Rationale: matches assurance cost to data sensitivity, exactly as the I1/sovereignty model
does for placement. Presentable to a due-diligence client as "audit posture is
classification-aware: confidential fails closed, internal degrades gracefully" — which
demonstrates the judgment, not just a binary. Couples Q4 to Q1's classification machinery
(being built anyway) and removes the buffer-integrity rabbit hole for the only tier where it
mattered.
+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)