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///--.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] }