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>
317 lines
10 KiB
Go
317 lines
10 KiB
Go
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}
|
|
|
|
// 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)
|
|
|
|
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
|
|
}
|
|
|
|
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: emit a request-level audit record of exactly what landed.
|
|
// Best-effort here; the classification-aware refusal/degradation
|
|
// policy is #54.
|
|
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,
|
|
}); err != nil {
|
|
receipt.Errors = append(receipt.Errors, ItemError{Item: "audit", Error: err.Error()})
|
|
}
|
|
|
|
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)
|
|
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]
|
|
}
|