diff --git a/README.md b/README.md index d39b0f3..03da84d 100644 --- a/README.md +++ b/README.md @@ -49,8 +49,9 @@ pilotctl publish ticker.btcusd --data '...' | File | What it does | |---|---| -| `eventstream.go` | Wire format: `Event{Topic, Payload}`, length-prefixed framing. `WriteEvent` / `ReadEvent`. | +| `eventstream.go` | Wire format: `Event{Topic, Payload}`, subscription policy interface, and length-prefixed framing. `WriteEvent` / `ReadEvent`. | | `client.go` | Subscriber side: `Client.Dial`, `.Subscribe(topic)`, `.Read`, `.Close`. | +| `governed.go` | Signed publication envelope, broker-side verifier, and enforceable topic/payload constraints. | | `server.go` | Publisher side: `Server` accepts inbound stream connections and broadcasts. | | `service.go` | `*Service` — `coreapi.Service` adapter, binds port 1002. Build tag `!no_eventstream`. | | `service_disabled.go` | Stub when `-tags no_eventstream` is set. | @@ -62,6 +63,58 @@ pilotctl publish ticker.btcusd --data '...' |---|---| | `no_eventstream` | Compiles a no-op stub service. | +## Governed publication + +An enterprise broker can call `SetGovernedPublication` before `Start` with a +`DecisionEventVerifier` (or an equivalent local verifier). Its +`GovernedTopic` transport envelope binds a publication topic and exact bytes to +a signed `decision.Intent` and `decision.Decision`. The broker verifies the +tenant authority state, local deterministic ceiling, exact broker resource, +and understood constraints before it forwards the original event. Subscribers +receive the original topic and payload, never the envelope. + +After all publishers have been upgraded, set `require` to `true`; unsigned +legacy publications are then dropped at the broker. This is intentionally a +publication control: subscription access remains governed by the existing +`TopicPolicy`. A workflow-approved publication uses the same short-lived +execution Decision as an ordinary allowed publication, rather than a reusable +workflow token. + +For a typed-disclosure profile, use `PublishGovernedWithDisclosure` and bind +the canonical `decision.DisclosureBinding` hash into the signed Intent. A +`DecisionEventVerifier` with `RequireDisclosure` rejects otherwise-valid +governed publications that do not carry matching content metadata. +With a receipt recorder configured, typed publications require V2 +disclosure-evidence support before broker fanout; the signed receipt contains +only the disclosure hash, not event plaintext. + +For auditable enterprise publishing, also call `SetGovernedReceiptRecorder` +before `Start` and require it. The broker appends evidence for the exact signed +Intent and Decision before fanout; a recorder failure denies the publication, +so subscribers never receive an unreceipted governed event. + +### Local content inspection + +Call `SetGovernedContentInspector` before `Start` to inspect verified governed +event bytes locally before broker fan-out. `RequireGovernedContentInspection` +makes a missing hook a startup failure, while a detector error rejects only the +publication. `decision.PresidioInspector` provides a bounded text/structured- +text adapter for a tenant-local OSS Presidio service; unsupported binary types +are rejected rather than bypassing inspection. This hook is not invoked by the +central decision authority. + +Typed disclosure binding V2 can carry a signed `retention_class`, which a +policy may restrict before fan-out. This binds metadata for downstream +retention operations; it is not a substitute for a retention executor. + +### Per-agent publication quotas + +`SetGovernedTransferQuota` applies a bounded local byte/action budget to each +verified publisher `Intent.AgentID`. The budget is charged after signature and +local-policy verification, never from a network address or caller-supplied +identity, and covers admitted attempts that later fail local inspection or +receipt recording. + ## License AGPL-3.0-or-later. See [LICENSE](LICENSE). diff --git a/client.go b/client.go index 0ef9670..f4f4f4c 100644 --- a/client.go +++ b/client.go @@ -3,6 +3,7 @@ package eventstream import ( + "github.com/pilot-protocol/common/decision" "github.com/pilot-protocol/common/driver" "github.com/pilot-protocol/common/protocol" ) @@ -35,6 +36,35 @@ func (c *Client) Publish(topic string, payload []byte) error { return WriteEvent(c.conn, &Event{Topic: topic, Payload: payload}) } +// PublishGoverned publishes an event with its exact signed intent and +// authority decision. A broker configured to require governed publications +// verifies the envelope before forwarding the inner event to subscribers. +func (c *Client) PublishGoverned(event *Event, intent decision.Intent, result decision.Decision) error { + governed, err := NewGovernedEvent(event, intent, result) + if err != nil { + return err + } + envelope, err := EncodeGovernedEvent(governed) + if err != nil { + return err + } + return WriteEvent(c.conn, envelope) +} + +// PublishGovernedWithDisclosure publishes a signed event whose Intent binds +// typed disclosure metadata. The broker can require this form per topic. +func (c *Client) PublishGovernedWithDisclosure(event *Event, intent decision.Intent, result decision.Decision, disclosure decision.DisclosureBinding) error { + governed, err := NewGovernedEventWithDisclosure(event, intent, result, disclosure) + if err != nil { + return err + } + envelope, err := EncodeGovernedEvent(governed) + if err != nil { + return err + } + return WriteEvent(c.conn, envelope) +} + // Recv waits for the next event from the broker. func (c *Client) Recv() (*Event, error) { return ReadEvent(c.conn) diff --git a/eventstream.go b/eventstream.go index ee8ae80..0b00c8b 100644 --- a/eventstream.go +++ b/eventstream.go @@ -7,6 +7,8 @@ import ( "fmt" "io" "unicode/utf8" + + "github.com/pilot-protocol/common/coreapi" ) // Event is a typed message published to the event stream. @@ -16,6 +18,14 @@ type Event struct { Payload []byte } +// TopicPolicy is the authorization gate for topic subscription. +// Implementations check whether a peer may subscribe to a named topic. +// The enabled service uses an allow-all policy unless its caller installs a +// stricter implementation. +type TopicPolicy interface { + AllowSubscribe(remoteAddr coreapi.Addr, topic string) bool +} + // WriteEvent writes an event to a writer. func WriteEvent(w io.Writer, e *Event) error { topic := []byte(e.Topic) diff --git a/go.mod b/go.mod index 1ffb668..7cbec6e 100644 --- a/go.mod +++ b/go.mod @@ -2,4 +2,4 @@ module github.com/pilot-protocol/eventstream go 1.25.12 -require github.com/pilot-protocol/common v0.5.11 +require github.com/pilot-protocol/common v0.5.12 diff --git a/go.sum b/go.sum index 94c6484..7b880c7 100644 --- a/go.sum +++ b/go.sum @@ -1,2 +1,2 @@ -github.com/pilot-protocol/common v0.5.11 h1:gaPOT2v3/FUAx61lqPw7yRsrzmDSRWL4+LRuQh8mrqs= -github.com/pilot-protocol/common v0.5.11/go.mod h1:Ybc6f1A37s3ShoEh1nBMVL9DPyYlxvkqPTvtbxaNWg4= +github.com/pilot-protocol/common v0.5.12 h1:ZQ7v8oX0VYtEcluraQZvqrYPMvRIxjixSZtNz8Xo5Uc= +github.com/pilot-protocol/common v0.5.12/go.mod h1:Ybc6f1A37s3ShoEh1nBMVL9DPyYlxvkqPTvtbxaNWg4= diff --git a/governed.go b/governed.go new file mode 100644 index 0000000..8c2471c --- /dev/null +++ b/governed.go @@ -0,0 +1,287 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +package eventstream + +import ( + "bytes" + "context" + "crypto/sha256" + "encoding/binary" + "encoding/json" + "fmt" + "io" + "strconv" + "strings" + "unicode/utf8" + + "github.com/pilot-protocol/common/coreapi" + "github.com/pilot-protocol/common/decision" +) + +// GovernedTopic is reserved for a signed publication envelope. The broker +// unwraps it only after local verification and forwards the original topic and +// payload to subscribers; subscribers never receive this transport topic. +const GovernedTopic = "\x00pilot.governed.v1" + +const governedEventVersion uint16 = 1 + +// GovernedEvent transports one published Event with its exact signed intent +// and authority decision. Its payload binding includes both topic and bytes, +// so neither can be changed after authorization. +type GovernedEvent struct { + Version uint16 `json:"version"` + Topic string `json:"topic"` + Payload []byte `json:"payload"` + Disclosure *decision.DisclosureBinding `json:"disclosure,omitempty"` + Intent decision.Intent `json:"intent"` + Decision decision.Decision `json:"decision"` +} + +// GovernedEventVerifier verifies a governed publication against the local +// broker's authority state before it is distributed to subscribers. +type GovernedEventVerifier interface { + VerifyGovernedEvent(context.Context, coreapi.Addr, GovernedEvent) error +} + +// GovernedReceiptRecorder durably records a broker-side enforcement receipt +// for a verified governed publication before it is fanned out to subscribers. +// Required enterprise deployments treat a recorder failure as a publication +// denial, so an event is never delivered without its evidence. +type GovernedReceiptRecorder interface { + RecordGovernedReceipt(context.Context, decision.Intent, decision.Decision) error +} + +// GovernedDisclosureReceiptRecorder records V2 evidence for typed governed +// publications. Brokers fail closed instead of recording a V1 receipt that +// omits disclosure proof. +type GovernedDisclosureReceiptRecorder interface { + RecordGovernedDisclosureReceipt(context.Context, decision.Intent, decision.Decision, decision.DisclosureBinding) error +} + +func recordGovernedReceipt(ctx context.Context, recorder GovernedReceiptRecorder, intent decision.Intent, result decision.Decision, disclosure *decision.DisclosureBinding) error { + if recorder == nil { + return fmt.Errorf("eventstream: governed receipt recorder is not configured") + } + if disclosure == nil { + return recorder.RecordGovernedReceipt(ctx, intent, result) + } + typed, supported := recorder.(GovernedDisclosureReceiptRecorder) + if !supported { + return fmt.Errorf("eventstream: governed receipt recorder does not support disclosure evidence") + } + return typed.RecordGovernedDisclosureReceipt(ctx, intent, result, *disclosure) +} + +// DecisionEventVerifier is the reference broker-side verifier. Resource must +// return the exact local resource for the incoming publication (for example, +// "eventstream:alerts"). It prevents a valid decision for one topic from +// being replayed to another topic at the same broker. +type DecisionEventVerifier struct { + Enforcer *decision.Enforcer + Resource func(coreapi.Addr, *Event) string + RequireDisclosure bool +} + +func NewGovernedEvent(event *Event, intent decision.Intent, result decision.Decision) (GovernedEvent, error) { + if event == nil { + return GovernedEvent{}, fmt.Errorf("eventstream: governed event is required") + } + governed := GovernedEvent{ + Version: governedEventVersion, Topic: event.Topic, Payload: append([]byte(nil), event.Payload...), Intent: intent, Decision: result, + } + if err := governed.Validate(); err != nil { + return GovernedEvent{}, err + } + return governed, nil +} + +// NewGovernedEventWithDisclosure creates a governed event whose Intent binds +// canonical disclosure metadata. The caller is responsible for requesting a +// Decision for that exact disclosure-bound Intent before publication. +func NewGovernedEventWithDisclosure(event *Event, intent decision.Intent, result decision.Decision, disclosure decision.DisclosureBinding) (GovernedEvent, error) { + if event == nil { + return GovernedEvent{}, fmt.Errorf("eventstream: governed event is required") + } + disclosure.Labels = append([]string(nil), disclosure.Labels...) + governed := GovernedEvent{ + Version: governedEventVersion, Topic: event.Topic, Payload: append([]byte(nil), event.Payload...), Disclosure: &disclosure, Intent: intent, Decision: result, + } + if err := governed.Validate(); err != nil { + return GovernedEvent{}, err + } + return governed, nil +} + +func (event GovernedEvent) Validate() error { + if event.Version != governedEventVersion || !validGovernedTopic(event.Topic) { + return fmt.Errorf("eventstream: invalid governed topic") + } + if len(event.Payload) > 1<<24 { + return fmt.Errorf("eventstream: governed payload exceeds maximum event size") + } + if err := event.Intent.Validate(); err != nil { + return fmt.Errorf("eventstream: invalid governed intent: %w", err) + } + if event.Intent.Signature == "" || event.Decision.Signature == "" { + return fmt.Errorf("eventstream: governed intent and decision must be signed") + } + if event.Intent.Action != "event.publish" { + return fmt.Errorf("eventstream: governed intent action must be event.publish") + } + if err := event.verifyPayloadBinding(); err != nil { + return err + } + if err := event.Decision.Validate(); err != nil { + return fmt.Errorf("eventstream: invalid governed decision: %w", err) + } + return nil +} + +func (event GovernedEvent) verifyPayloadBinding() error { + if event.Disclosure == nil { + if event.Intent.PayloadHash != GovernedEventPayloadHash(event.Topic, event.Payload) { + return fmt.Errorf("eventstream: governed intent payload binding mismatch") + } + return nil + } + if event.Disclosure.ContentHash != decision.HashPayload(event.Payload) || event.Disclosure.DeclaredBytes != uint64(len(event.Payload)) || event.Disclosure.Filename != "" || event.Disclosure.TransferID != "" { + return fmt.Errorf("eventstream: governed disclosure does not match event") + } + if err := event.Disclosure.VerifyIntent(event.Intent); err != nil { + return fmt.Errorf("eventstream: governed disclosure intent binding: %w", err) + } + return nil +} + +func (event GovernedEvent) Event() *Event { + return &Event{Topic: event.Topic, Payload: append([]byte(nil), event.Payload...)} +} + +func EncodeGovernedEvent(event GovernedEvent) (*Event, error) { + if err := event.Validate(); err != nil { + return nil, err + } + body, err := json.Marshal(event) + if err != nil { + return nil, fmt.Errorf("eventstream: encode governed event: %w", err) + } + if len(body) > 1<<24 { + return nil, fmt.Errorf("eventstream: governed envelope exceeds maximum event size") + } + return &Event{Topic: GovernedTopic, Payload: body}, nil +} + +func DecodeGovernedEvent(event *Event) (GovernedEvent, error) { + if event == nil || event.Topic != GovernedTopic || len(event.Payload) > 1<<24 { + return GovernedEvent{}, fmt.Errorf("eventstream: invalid governed envelope") + } + decoder := json.NewDecoder(bytes.NewReader(event.Payload)) + decoder.DisallowUnknownFields() + var governed GovernedEvent + if err := decoder.Decode(&governed); err != nil { + return GovernedEvent{}, fmt.Errorf("eventstream: decode governed envelope: %w", err) + } + var trailing any + if err := decoder.Decode(&trailing); err != io.EOF { + return GovernedEvent{}, fmt.Errorf("eventstream: trailing governed envelope data") + } + if err := governed.Validate(); err != nil { + return GovernedEvent{}, err + } + return governed, nil +} + +func (verifier DecisionEventVerifier) VerifyGovernedEvent(ctx context.Context, remote coreapi.Addr, governed GovernedEvent) error { + if verifier.Enforcer == nil || verifier.Resource == nil { + return fmt.Errorf("eventstream: decision event verifier is not initialized") + } + if err := governed.Validate(); err != nil { + return err + } + if verifier.RequireDisclosure && governed.Disclosure == nil { + return fmt.Errorf("eventstream: governed disclosure is required") + } + event := governed.Event() + resource := verifier.Resource(remote, event) + if resource == "" || governed.Intent.Resource != resource { + return fmt.Errorf("eventstream: governed intent resource binding mismatch") + } + var verifyErr error + if governed.Disclosure != nil { + verifyErr = verifier.Enforcer.VerifyDisclosure(ctx, governed.Intent, governed.Decision, *governed.Disclosure) + } else { + verifyErr = verifier.Enforcer.Verify(ctx, governed.Intent, governed.Decision) + } + if verifyErr != nil { + return fmt.Errorf("eventstream: verify governed decision: %w", verifyErr) + } + switch governed.Decision.Outcome { + case decision.Allow: + return nil + case decision.Constrain: + return enforceEventConstraints(governed.Decision.Constraints, remote, event) + default: + return fmt.Errorf("eventstream: governed decision outcome %q cannot permit publication", governed.Decision.Outcome) + } +} + +// GovernedEventPayloadHash returns the exact payload binding required by a +// GovernedTopic Intent. It includes topic and bytes, so callers must use it +// when creating the Intent before PublishGoverned. +func GovernedEventPayloadHash(topic string, payload []byte) string { + var header [8]byte + binary.BigEndian.PutUint32(header[:4], uint32(len(topic))) + binary.BigEndian.PutUint32(header[4:8], uint32(len(payload))) + hash := sha256.New() + _, _ = hash.Write([]byte("pilot-eventstream-governed-event-v1\x00")) + _, _ = hash.Write(header[:]) + _, _ = hash.Write([]byte(topic)) + _, _ = hash.Write(payload) + return decision.HashPayload(hash.Sum(nil)) +} + +func validGovernedTopic(topic string) bool { + return topic != "" && topic != GovernedTopic && len(topic) <= 1024 && utf8.ValidString(topic) +} + +func enforceEventConstraints(constraints []decision.Constraint, remote coreapi.Addr, event *Event) error { + attributes := map[string]string{ + "publisher": remote.String(), "topic": event.Topic, "bytes": strconv.Itoa(len(event.Payload)), + } + for _, constraint := range constraints { + actual, found := attributes[constraint.Key] + if !found { + return fmt.Errorf("eventstream: constraint %q has no enforceable event attribute", constraint.Key) + } + switch constraint.Operator { + case "eq": + if actual != constraint.Value { + return fmt.Errorf("eventstream: constraint %s rejected", constraint.Key) + } + case "one_of": + matched := false + for _, allowed := range strings.Split(constraint.Value, ",") { + if actual == strings.TrimSpace(allowed) { + matched = true + break + } + } + if !matched { + return fmt.Errorf("eventstream: constraint %s rejected", constraint.Key) + } + case "max", "min": + value, valueErr := strconv.ParseUint(actual, 10, 64) + limit, limitErr := strconv.ParseUint(constraint.Value, 10, 64) + if valueErr != nil || limitErr != nil || (constraint.Operator == "max" && value > limit) || (constraint.Operator == "min" && value < limit) { + return fmt.Errorf("eventstream: numeric constraint %s rejected", constraint.Key) + } + case "require": + if constraint.Value != "" && actual != constraint.Value { + return fmt.Errorf("eventstream: required constraint %s rejected", constraint.Key) + } + default: + return fmt.Errorf("eventstream: constraint operator %q is not enforceable", constraint.Operator) + } + } + return nil +} diff --git a/governed_replay.go b/governed_replay.go new file mode 100644 index 0000000..4500004 --- /dev/null +++ b/governed_replay.go @@ -0,0 +1,61 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +package eventstream + +import ( + "fmt" + "sync" + "time" + + "github.com/pilot-protocol/common/decision" +) + +// maxGovernedReplayEntries bounds the receiver-side replay cache; entries +// expire at the intent's own ExpiresAt (<= 5-minute MaxIntentTTL) so the cache +// drains under legitimate signed traffic. On overflow a publication is refused +// (fail closed) rather than admitted un-deduplicated. +const maxGovernedReplayEntries = 1 << 20 + +// governedReplayGuard rejects a second fan-out of the same signed governed +// event. Without it a verified GovernedEvent is a bearer capability: any peer +// that observed one could re-publish the exact bytes within the intent TTL, +// producing duplicate authorized publications and re-charging the signing +// agent's quota. Dedup is keyed on the signature-authenticated +// (tenant, agent, intent id); a legitimate re-publish must carry a fresh +// intent. +type governedReplayGuard struct { + mu sync.Mutex + seen map[string]int64 + now func() time.Time +} + +func newGovernedReplayGuard() *governedReplayGuard { + return &governedReplayGuard{seen: make(map[string]int64), now: time.Now} +} + +func (g *governedReplayGuard) admit(intent decision.Intent) error { + if intent.ID == "" { + return nil + } + now := g.now().Unix() + key := intent.TenantID + "\x1f" + intent.AgentID + "\x1f" + intent.ID + g.mu.Lock() + defer g.mu.Unlock() + for k, exp := range g.seen { + if exp <= now { + delete(g.seen, k) + } + } + if exp, ok := g.seen[key]; ok && exp > now { + return fmt.Errorf("governed event already published (replay rejected)") + } + if len(g.seen) >= maxGovernedReplayEntries { + return fmt.Errorf("governed replay cache saturated") + } + expiresAt := intent.ExpiresAt + if expiresAt <= now { + expiresAt = now + int64(decision.MaxIntentTTL/time.Second) + } + g.seen[key] = expiresAt + return nil +} diff --git a/governed_test.go b/governed_test.go new file mode 100644 index 0000000..0505994 --- /dev/null +++ b/governed_test.go @@ -0,0 +1,475 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +//go:build !no_eventstream +// +build !no_eventstream + +package eventstream + +import ( + "bytes" + "context" + "crypto/ed25519" + "crypto/rand" + "encoding/json" + "errors" + "io" + "strings" + "testing" + "time" + + "github.com/pilot-protocol/common/coreapi" + "github.com/pilot-protocol/common/decision" +) + +type governedEventTestTrust struct { + intentKey ed25519.PublicKey + decisionKey ed25519.PublicKey +} + +func (trust governedEventTestTrust) IntentKey(context.Context, string, string, string) (ed25519.PublicKey, error) { + return trust.intentKey, nil +} + +func (trust governedEventTestTrust) DecisionKey(context.Context, string, string) (ed25519.PublicKey, error) { + return trust.decisionKey, nil +} + +func (governedEventTestTrust) MinimumState(context.Context, string) (uint64, uint64, error) { + return 7, 3, nil +} + +type governedEventTestCeiling struct{} + +func (governedEventTestCeiling) Check(context.Context, decision.Intent, decision.Decision) error { + return nil +} + +func (governedEventTestCeiling) CheckDisclosure(context.Context, decision.Intent, decision.Decision, decision.DisclosureBinding) error { + return nil +} + +type governedEventReceiptRecorder struct { + calls int + intent decision.Intent + decision decision.Decision + err error +} + +type eventContentInspectorFunc func(context.Context, decision.Intent, *decision.DisclosureBinding, string, string, io.Reader) error + +func (inspect eventContentInspectorFunc) InspectDisclosureContent(ctx context.Context, intent decision.Intent, disclosure *decision.DisclosureBinding, contentType, filename string, content io.Reader) error { + return inspect(ctx, intent, disclosure, contentType, filename, content) +} + +func TestGovernedPublicationContentInspectionRunsBeforeFanout(t *testing.T) { + governed, verifier := newGovernedEventForTest(t, &Event{Topic: "alerts", Payload: []byte("classified")}, decision.Allow, nil) + envelope, err := EncodeGovernedEvent(governed) + if err != nil { + t.Fatal(err) + } + senderServer, senderClient := newPipeStreamPair() + defer senderClient.Close() + broker := newBroker(nil, defaultAllowPolicy{}) + broker.governedVerifier = verifier + // Reuses one envelope to re-run inspection with a failing scanner; the + // replay guard (tested separately) would reject the second identical + // delivery first. + broker.replay = nil + var observed []byte + broker.contentInspector = eventContentInspectorFunc(func(_ context.Context, intent decision.Intent, disclosure *decision.DisclosureBinding, contentType, filename string, content io.Reader) error { + if intent.ID != governed.Intent.ID || disclosure != nil || contentType != "application/octet-stream" || filename != "" { + t.Fatalf("inspection metadata intent=%+v disclosure=%+v content_type=%q filename=%q", intent, disclosure, contentType, filename) + } + var readErr error + observed, readErr = io.ReadAll(content) + return readErr + }) + published, err := broker.governPublication(newSubscriber(senderServer), envelope) + if err != nil || published == nil || !bytes.Equal(observed, governed.Payload) { + t.Fatalf("inspection published=%+v err=%v payload=%q", published, err, observed) + } + broker.contentInspector = eventContentInspectorFunc(func(context.Context, decision.Intent, *decision.DisclosureBinding, string, string, io.Reader) error { + return errors.New("scanner unavailable") + }) + if _, err := broker.governPublication(newSubscriber(senderServer), envelope); err == nil || err.Error() != "eventstream: governed content inspection rejected" { + t.Fatalf("inspection failure leaked or was accepted: %v", err) + } +} + +func TestGovernedPublicationQuotaUsesSignedAgentIdentity(t *testing.T) { + governed, verifier := newGovernedEventForTest(t, &Event{Topic: "alerts", Payload: []byte("12345")}, decision.Allow, nil) + envelope, err := EncodeGovernedEvent(governed) + if err != nil { + t.Fatal(err) + } + limiter, err := decision.NewTransferQuotaLimiter(decision.TransferQuotaConfig{Window: time.Minute, MaxBytes: 5, MaxSenders: 2}) + if err != nil { + t.Fatal(err) + } + senderServer, senderClient := newPipeStreamPair() + defer senderClient.Close() + broker := newBroker(nil, defaultAllowPolicy{}) + broker.governedVerifier = verifier + broker.governedTransferQuota = limiter + // This test reuses one signed envelope to exercise quota accumulation for + // the same agent; the receiver-side replay guard (covered by + // TestGovernedPublicationReplayGuard) would otherwise reject the second + // delivery of the identical intent first. + broker.replay = nil + if _, err := broker.governPublication(newSubscriber(senderServer), envelope); err != nil { + t.Fatal(err) + } + if _, err := broker.governPublication(newSubscriber(senderServer), envelope); err == nil || err.Error() != "eventstream: governed transfer quota rejected" { + t.Fatalf("quota error=%v", err) + } +} + +func (recorder *governedEventReceiptRecorder) RecordGovernedReceipt(_ context.Context, intent decision.Intent, result decision.Decision) error { + recorder.calls++ + recorder.intent, recorder.decision = intent, result + return recorder.err +} + +type legacyGovernedEventReceiptRecorder struct{} + +func (legacyGovernedEventReceiptRecorder) RecordGovernedReceipt(context.Context, decision.Intent, decision.Decision) error { + return nil +} + +type disclosureGovernedEventReceiptRecorder struct { + governedEventReceiptRecorder + disclosure decision.DisclosureBinding +} + +func (recorder *disclosureGovernedEventReceiptRecorder) RecordGovernedDisclosureReceipt(_ context.Context, intent decision.Intent, result decision.Decision, disclosure decision.DisclosureBinding) error { + recorder.calls++ + recorder.intent, recorder.decision, recorder.disclosure = intent, result, disclosure + return recorder.err +} + +func TestDisclosurePublicationReceiptRecorderRequiresV2Evidence(t *testing.T) { + disclosure := decision.DisclosureBinding{Version: decision.DisclosureBindingVersion} + if err := recordGovernedReceipt(context.Background(), legacyGovernedEventReceiptRecorder{}, decision.Intent{ID: "intent"}, decision.Decision{ID: "decision"}, &disclosure); err == nil || !strings.Contains(err.Error(), "does not support disclosure") { + t.Fatalf("legacy disclosure recorder err=%v", err) + } + recorder := &disclosureGovernedEventReceiptRecorder{} + if err := recordGovernedReceipt(context.Background(), recorder, decision.Intent{ID: "intent"}, decision.Decision{ID: "decision"}, &disclosure); err != nil { + t.Fatal(err) + } + if recorder.calls != 1 || recorder.disclosure.Version != decision.DisclosureBindingVersion { + t.Fatalf("disclosure recorder=%+v", recorder) + } +} + +type verifierStub struct{} + +func (verifierStub) VerifyGovernedEvent(context.Context, coreapi.Addr, GovernedEvent) error { + return nil +} + +func newGovernedEventForTest(t *testing.T, event *Event, outcome decision.Outcome, constraints []decision.Constraint) (GovernedEvent, DecisionEventVerifier) { + t.Helper() + intentPublic, intentPrivate, err := ed25519.GenerateKey(rand.Reader) + if err != nil { + t.Fatalf("generate intent key: %v", err) + } + decisionPublic, decisionPrivate, err := ed25519.GenerateKey(rand.Reader) + if err != nil { + t.Fatalf("generate decision key: %v", err) + } + now := time.Now().UTC().Truncate(time.Second) + nonce, err := decision.NewNonce() + if err != nil { + t.Fatalf("nonce: %v", err) + } + intent := decision.Intent{ + Version: decision.SchemaVersion, ID: "event-intent", TenantID: "tenant-a", AgentID: "publisher-a", + Action: "event.publish", Resource: "eventstream:" + event.Topic, PayloadHash: GovernedEventPayloadHash(event.Topic, event.Payload), + Risk: decision.RiskMedium, IssuedAt: now.Unix(), ExpiresAt: now.Add(2 * time.Minute).Unix(), Nonce: nonce, KeyID: "publisher-key", + } + if err := intent.Sign(intentPrivate); err != nil { + t.Fatalf("sign intent: %v", err) + } + intentHash, err := intent.Hash() + if err != nil { + t.Fatalf("hash intent: %v", err) + } + result := decision.Decision{ + Version: decision.SchemaVersion, ID: "event-decision", IntentHash: intentHash, TenantID: intent.TenantID, AgentID: intent.AgentID, + Outcome: outcome, Constraints: constraints, PolicyRevision: 7, RevocationEpoch: 3, ProviderID: "authority-a", + IssuedAt: now.Unix(), ExpiresAt: now.Add(90 * time.Second).Unix(), KeyID: "authority-key", + } + if err := result.Sign(decisionPrivate); err != nil { + t.Fatalf("sign decision: %v", err) + } + governed, err := NewGovernedEvent(event, intent, result) + if err != nil { + t.Fatalf("new governed event: %v", err) + } + verifier := DecisionEventVerifier{ + Enforcer: &decision.Enforcer{ + Trust: governedEventTestTrust{intentKey: intentPublic, decisionKey: decisionPublic}, Ceiling: governedEventTestCeiling{}, Now: func() time.Time { return now }, + }, + Resource: func(_ coreapi.Addr, event *Event) string { return "eventstream:" + event.Topic }, + } + return governed, verifier +} + +func TestGovernedEventRoundTripBindsTopicAndPayload(t *testing.T) { + governed, _ := newGovernedEventForTest(t, &Event{Topic: "alerts", Payload: []byte("approved alert")}, decision.Allow, nil) + envelope, err := EncodeGovernedEvent(governed) + if err != nil { + t.Fatalf("encode governed event: %v", err) + } + decoded, err := DecodeGovernedEvent(envelope) + if err != nil { + t.Fatalf("decode governed event: %v", err) + } + if got := decoded.Event(); got.Topic != "alerts" || string(got.Payload) != "approved alert" { + t.Fatalf("decoded event = %#v", got) + } + + decoded.Topic = "other-alerts" + tampered, err := json.Marshal(decoded) + if err != nil { + t.Fatalf("marshal tampered event: %v", err) + } + if _, err := DecodeGovernedEvent(&Event{Topic: GovernedTopic, Payload: tampered}); err == nil || !strings.Contains(err.Error(), "payload binding") { + t.Fatalf("tampered event error = %v, want payload binding failure", err) + } + + invalidTopic := governed + invalidTopic.Topic = string([]byte{0xff}) + if err := invalidTopic.Validate(); err == nil || !strings.Contains(err.Error(), "invalid governed topic") { + t.Fatalf("invalid topic error = %v, want UTF-8 validation failure", err) + } +} + +func TestGovernedEventDisclosureBindingAndRequiredProfile(t *testing.T) { + event := &Event{Topic: "finance.alerts", Payload: []byte(`{"amount":42}`)} + governed, verifier := newGovernedEventForTest(t, event, decision.Allow, nil) + strict := verifier + strict.RequireDisclosure = true + if err := strict.VerifyGovernedEvent(context.Background(), coreapi.Addr{}, governed); err == nil || !strings.Contains(err.Error(), "disclosure is required") { + t.Fatalf("missing disclosure error=%v", err) + } + binding := decision.DisclosureBinding{ + Version: decision.DisclosureBindingVersion, ContentHash: decision.HashPayload(event.Payload), DeclaredBytes: uint64(len(event.Payload)), + ContentType: "application/json", Labels: []string{"finance", "pii"}, Recipient: "broker:finance", Purpose: "alert-publication", Residency: "eu-west-1", + } + hash, err := binding.Hash() + if err != nil { + t.Fatal(err) + } + governed.Intent.PayloadHash = hash + governed.Intent.Audience = binding.Recipient + governed.Intent.Purpose = binding.Purpose + governed.Intent.Signature = "transport-test-signature" + governed.Decision.Signature = "transport-test-signature" + governed, err = NewGovernedEventWithDisclosure(event, governed.Intent, governed.Decision, binding) + if err != nil { + t.Fatalf("valid disclosure envelope: %v", err) + } + tampered := governed + tampered.Disclosure = &binding + tampered.Disclosure.Labels = []string{"finance"} + if err := tampered.Validate(); err == nil || !strings.Contains(err.Error(), "disclosure intent binding") { + t.Fatalf("disclosure mutation error=%v", err) + } +} + +func TestDecisionEventVerifierBindsBrokerAndEnforcesConstraints(t *testing.T) { + governed, verifier := newGovernedEventForTest(t, &Event{Topic: "alerts", Payload: []byte("hello")}, decision.Constrain, []decision.Constraint{ + {Key: "topic", Operator: "eq", Value: "alerts"}, + {Key: "bytes", Operator: "max", Value: "5"}, + }) + if err := verifier.VerifyGovernedEvent(context.Background(), coreapi.Addr{}, governed); err != nil { + t.Fatalf("verify governed event: %v", err) + } + + wrongBroker := verifier + wrongBroker.Resource = func(_ coreapi.Addr, _ *Event) string { return "eventstream:other" } + if err := wrongBroker.VerifyGovernedEvent(context.Background(), coreapi.Addr{}, governed); err == nil || !strings.Contains(err.Error(), "resource binding") { + t.Fatalf("wrong broker error = %v, want resource-binding failure", err) + } + + tooLarge, constrainedVerifier := newGovernedEventForTest(t, &Event{Topic: "alerts", Payload: []byte("too large")}, decision.Constrain, []decision.Constraint{{Key: "bytes", Operator: "max", Value: "3"}}) + if err := constrainedVerifier.VerifyGovernedEvent(context.Background(), coreapi.Addr{}, tooLarge); err == nil || !strings.Contains(err.Error(), "numeric constraint") { + t.Fatalf("oversize event error = %v, want constraint failure", err) + } +} + +func TestBrokerRequireGovernedUnwrapsVerifiedAndRejectsLegacy(t *testing.T) { + governed, verifier := newGovernedEventForTest(t, &Event{Topic: "alerts", Payload: []byte("approved")}, decision.Allow, nil) + envelope, err := EncodeGovernedEvent(governed) + if err != nil { + t.Fatalf("encode governed event: %v", err) + } + server, peer := newPipeStreamPair() + defer server.Close() + defer peer.Close() + broker := newBroker(nil, defaultAllowPolicy{}) + broker.requireGoverned = true + broker.governedVerifier = verifier + sender := newSubscriber(server) + + got, err := broker.governPublication(sender, envelope) + if err != nil { + t.Fatalf("governed publication: %v", err) + } + if got.Topic != "alerts" || string(got.Payload) != "approved" { + t.Fatalf("unwrapped publication = %#v", got) + } + if _, err := broker.governPublication(sender, &Event{Topic: "alerts", Payload: []byte("legacy")}); err == nil || !strings.Contains(err.Error(), "unsigned legacy") { + t.Fatalf("legacy publication error = %v, want governed rejection", err) + } +} + +func TestBrokerHandleConnRequireGovernedDropsLegacyPublication(t *testing.T) { + bus := &stubEventBus{} + governed, verifier := newGovernedEventForTest(t, &Event{Topic: "alerts", Payload: []byte("approved")}, decision.Allow, nil) + if _, err := EncodeGovernedEvent(governed); err != nil { + t.Fatalf("governed fixture must encode: %v", err) + } + broker := newBroker(bus, defaultAllowPolicy{}) + broker.requireGoverned = true + broker.governedVerifier = verifier + server, client := newPipeStreamPair() + sender := newSubscriber(server) + done := make(chan struct{}) + go func() { + defer close(done) + broker.handleConn(sender) + }() + + if err := WriteEvent(client, &Event{Topic: "publisher"}); err != nil { + t.Fatalf("write subscription: %v", err) + } + if err := WriteEvent(client, &Event{Topic: "alerts", Payload: []byte("legacy")}); err != nil { + t.Fatalf("write legacy publication: %v", err) + } + deadline := time.Now().Add(2 * time.Second) + for time.Now().Before(deadline) && !bus.seen("pubsub.publish_denied") { + time.Sleep(5 * time.Millisecond) + } + if !bus.seen("pubsub.publish_denied") { + t.Fatalf("expected pubsub.publish_denied, got %v", bus.topics) + } + _ = client.Close() + select { + case <-done: + case <-time.After(2 * time.Second): + t.Fatal("broker did not return after client close") + } +} + +func TestBrokerRecordsGovernedReceiptBeforeFanout(t *testing.T) { + bus := &stubEventBus{} + governed, verifier := newGovernedEventForTest(t, &Event{Topic: "alerts", Payload: []byte("approved")}, decision.Allow, nil) + envelope, err := EncodeGovernedEvent(governed) + if err != nil { + t.Fatal(err) + } + recorder := &governedEventReceiptRecorder{} + broker := newBroker(bus, defaultAllowPolicy{}) + broker.requireGoverned, broker.governedVerifier = true, verifier + broker.requireReceipts, broker.receiptRecorder = true, recorder + + senderServer, senderClient := newPipeStreamPair() + defer senderClient.Close() + sender := newSubscriber(senderServer) + receiverServer, receiverClient := newPipeStreamPair() + defer receiverClient.Close() + if !broker.addSub("alerts", newSubscriber(receiverServer)) { + t.Fatal("add receiver") + } + received := make(chan *Event, 1) + go func() { + event, _ := ReadEvent(receiverClient) + received <- event + }() + done := make(chan struct{}) + go func() { + defer close(done) + broker.handleConn(sender) + }() + if err := WriteEvent(senderClient, &Event{Topic: "publisher"}); err != nil { + t.Fatal(err) + } + if err := WriteEvent(senderClient, envelope); err != nil { + t.Fatal(err) + } + select { + case event := <-received: + if event == nil || event.Topic != "alerts" || string(event.Payload) != "approved" { + t.Fatalf("fanout event=%+v", event) + } + case <-time.After(2 * time.Second): + t.Fatal("governed event was not fanned out") + } + if recorder.calls != 1 || recorder.intent.ID != governed.Intent.ID || recorder.decision.ID != governed.Decision.ID { + t.Fatalf("receipt recorder=%+v", recorder) + } + _ = senderClient.Close() + select { + case <-done: + case <-time.After(2 * time.Second): + t.Fatal("broker did not stop after publisher close") + } +} + +func TestBrokerReceiptFailurePreventsGovernedFanout(t *testing.T) { + bus := &stubEventBus{} + governed, verifier := newGovernedEventForTest(t, &Event{Topic: "alerts", Payload: []byte("approved")}, decision.Allow, nil) + envelope, err := EncodeGovernedEvent(governed) + if err != nil { + t.Fatal(err) + } + broker := newBroker(bus, defaultAllowPolicy{}) + broker.requireGoverned, broker.governedVerifier = true, verifier + broker.requireReceipts, broker.receiptRecorder = true, &governedEventReceiptRecorder{err: errors.New("journal unavailable")} + senderServer, senderClient := newPipeStreamPair() + sender := newSubscriber(senderServer) + done := make(chan struct{}) + go func() { + defer close(done) + broker.handleConn(sender) + }() + if err := WriteEvent(senderClient, &Event{Topic: "publisher"}); err != nil { + t.Fatal(err) + } + if err := WriteEvent(senderClient, envelope); err != nil { + t.Fatal(err) + } + deadline := time.Now().Add(2 * time.Second) + for time.Now().Before(deadline) && !bus.seen("pubsub.publish_denied") { + time.Sleep(5 * time.Millisecond) + } + if !bus.seen("pubsub.publish_denied") || bus.seen("pubsub.published") { + t.Fatal("receipt failure did not deny publication before fanout") + } + _ = senderClient.Close() + select { + case <-done: + case <-time.After(2 * time.Second): + t.Fatal("broker did not stop after publisher close") + } +} + +func TestServiceRequireGovernedNeedsVerifier(t *testing.T) { + service := NewService() + service.SetGovernedPublication(nil, true) + if err := service.Start(context.Background(), coreapi.Deps{}); err == nil || !strings.Contains(err.Error(), "no verifier") { + t.Fatalf("start error = %v, want missing verifier", err) + } + service = NewService() + service.SetGovernedPublication(verifierStub{}, true) + service.SetGovernedReceiptRecorder(nil, true) + if err := service.Start(context.Background(), coreapi.Deps{}); err == nil || !strings.Contains(err.Error(), "receipt recorder") { + t.Fatalf("start error = %v, want missing receipt recorder", err) + } +} + +var _ decision.TrustStore = governedEventTestTrust{} +var _ decision.AuthorityCeiling = governedEventTestCeiling{} +var _ GovernedReceiptRecorder = (*governedEventReceiptRecorder)(nil) diff --git a/service.go b/service.go index c2bdd7e..8d0c9a3 100644 --- a/service.go +++ b/service.go @@ -6,13 +6,16 @@ package eventstream import ( + "bytes" "context" "fmt" "log/slog" "sync" "sync/atomic" "time" + "github.com/pilot-protocol/common/coreapi" + "github.com/pilot-protocol/common/decision" "github.com/pilot-protocol/common/protocol" ) @@ -44,19 +47,6 @@ const ( // bound, growing the subscriber slice linearly. const maxSubsPerTopic = 1000 -// TopicPolicy is the authorization gate for topic subscription. -// Implementations check whether a peer (identified by its 48-bit -// virtual address) may subscribe to a given topic. -// -// The defaultAllowPolicy permits all subscriptions (backward -// compatible). Daemon-level callers that need topic-level access -// control provide a restricted implementation via SetTopicPolicy. -type TopicPolicy interface { - // AllowSubscribe returns true if the peer at remoteAddr is - // permitted to subscribe to the named topic. - AllowSubscribe(remoteAddr coreapi.Addr, topic string) bool -} - // defaultAllowPolicy permits all subscriptions. Used when no // TopicPolicy has been set. type defaultAllowPolicy struct{} @@ -67,12 +57,19 @@ func (defaultAllowPolicy) AllowSubscribe(_ coreapi.Addr, _ string) bool { return // cmd/pilotctl _daemon-run construct it via NewService and register // via daemon.RegisterPlugin. type Service struct { - listener coreapi.Listener - deps coreapi.Deps - cancel context.CancelFunc - done chan struct{} - broker *broker - topicPolicy TopicPolicy + listener coreapi.Listener + deps coreapi.Deps + cancel context.CancelFunc + done chan struct{} + broker *broker + topicPolicy TopicPolicy + governedVerifier GovernedEventVerifier + requireGoverned bool + receiptRecorder GovernedReceiptRecorder + requireReceipts bool + contentInspector decision.DisclosureContentInspector + requireContentInspection bool + governedTransferQuota *decision.TransferQuotaLimiter } // NewService returns a Service ready for daemon.RegisterPlugin. @@ -89,6 +86,45 @@ func (s *Service) SetTopicPolicy(p TopicPolicy) { } } +// SetGovernedPublication configures the broker-side authorization gate for +// published events. When require is true every publication must be a signed +// governed envelope. Call before Start; a required gate without a verifier +// makes Start fail closed. +func (s *Service) SetGovernedPublication(verifier GovernedEventVerifier, require bool) { + s.governedVerifier = verifier + s.requireGoverned = require +} + +// SetGovernedReceiptRecorder attaches a durable evidence recorder to governed +// publications. When require is true the broker refuses to start without both +// a required governed gate and this recorder. +func (s *Service) SetGovernedReceiptRecorder(recorder GovernedReceiptRecorder, require bool) { + s.receiptRecorder = recorder + s.requireReceipts = require +} + +// SetGovernedContentInspector installs a receiver-local DLP hook. It runs +// after a signed governed envelope has passed its local policy ceiling and +// before the event reaches subscribers. The hook receives a reader over the +// event bytes; it never runs in the central decision authority. +func (s *Service) SetGovernedContentInspector(inspector decision.DisclosureContentInspector) { + s.contentInspector = inspector +} + +// RequireGovernedContentInspection makes startup fail closed when a local DLP +// hook is absent. Call this from the enterprise attachment after an operator +// has supplied the tenant-local inspector. +func (s *Service) RequireGovernedContentInspection(require bool) { + s.requireContentInspection = require +} + +// SetGovernedTransferQuota installs a receiver-local per-agent admission +// quota. It is charged only after a signed governed event verifies, using the +// signed Intent.AgentID rather than a transport address. +func (s *Service) SetGovernedTransferQuota(limiter *decision.TransferQuotaLimiter) { + s.governedTransferQuota = limiter +} + func (s *Service) Name() string { return "eventstream" } // Order: 120 — after the trust subsystem (50) and other application @@ -96,6 +132,18 @@ func (s *Service) Name() string { return "eventstream" } func (s *Service) Order() int { return 120 } func (s *Service) Start(ctx context.Context, deps coreapi.Deps) error { + if s.requireGoverned && s.governedVerifier == nil { + return fmt.Errorf("eventstream: governed publication is required but no verifier is configured") + } + if s.requireReceipts && (!s.requireGoverned || s.receiptRecorder == nil) { + return fmt.Errorf("eventstream: governed receipts require a governed broker and receipt recorder") + } + if s.requireContentInspection && (!s.requireGoverned || s.contentInspector == nil) { + return fmt.Errorf("eventstream: required content inspection needs a governed broker and local inspector") + } + if s.governedTransferQuota != nil && !s.requireGoverned { + return fmt.Errorf("eventstream: governed transfer quota requires a governed broker") + } s.deps = deps ln, err := deps.Streams.Listen(protocol.PortEventStream) if err != nil { @@ -103,6 +151,13 @@ func (s *Service) Start(ctx context.Context, deps coreapi.Deps) error { } s.listener = ln s.broker = newBroker(deps.Events, s.topicPolicy) + s.broker.governedVerifier = s.governedVerifier + s.broker.requireGoverned = s.requireGoverned + s.broker.receiptRecorder = s.receiptRecorder + s.broker.requireReceipts = s.requireReceipts + s.broker.contentInspector = s.contentInspector + s.broker.requireContentInspection = s.requireContentInspection + s.broker.governedTransferQuota = s.governedTransferQuota runCtx, cancel := context.WithCancel(ctx) s.cancel = cancel @@ -185,12 +240,20 @@ func (s *subscriber) remote() string { // Service; one instance per daemon. Tracks subscribers by topic plus // per-publisher rate buckets. type broker struct { - mu sync.RWMutex - subs map[string][]*subscriber - rateMu sync.Mutex - rate map[*subscriber]*publishBucket - events coreapi.EventBus // for pubsub.* observability events - topicPolicy TopicPolicy + mu sync.RWMutex + subs map[string][]*subscriber + rateMu sync.Mutex + rate map[*subscriber]*publishBucket + events coreapi.EventBus // for pubsub.* observability events + topicPolicy TopicPolicy + governedVerifier GovernedEventVerifier + requireGoverned bool + receiptRecorder GovernedReceiptRecorder + requireReceipts bool + contentInspector decision.DisclosureContentInspector + requireContentInspection bool + governedTransferQuota *decision.TransferQuotaLimiter + replay *governedReplayGuard } func newBroker(events coreapi.EventBus, policy TopicPolicy) *broker { @@ -198,6 +261,7 @@ func newBroker(events coreapi.EventBus, policy TopicPolicy) *broker { subs: make(map[string][]*subscriber), events: events, topicPolicy: policy, + replay: newGovernedReplayGuard(), } } @@ -323,8 +387,87 @@ func (b *broker) handleConn(sub *subscriber) { } continue } - b.publish(evt, sub) + published, governErr := b.governPublication(sub, evt) + if governErr != nil { + slog.Warn("eventstream broker: publication denied", "remote", sub.remote(), "err", governErr) + if b.events != nil { + b.events.Publish("pubsub.publish_denied", map[string]any{ + "topic": evt.Topic, "from": sub.remote(), + }) + } + continue + } + if evt.Topic == GovernedTopic && b.receiptRecorder != nil { + governed, decodeErr := DecodeGovernedEvent(evt) + if decodeErr != nil { + // governPublication already successfully decoded and verified this + // envelope. Keep a defensive failure boundary here in case this + // code changes independently in the future. + slog.Error("eventstream broker: verified envelope could not be reconstructed for receipt", "remote", sub.remote(), "err", decodeErr) + continue + } + if receiptErr := recordGovernedReceipt(context.Background(), b.receiptRecorder, governed.Intent, governed.Decision, governed.Disclosure); receiptErr != nil { + slog.Warn("eventstream broker: publication receipt failed", "remote", sub.remote(), "err", receiptErr) + if b.events != nil { + b.events.Publish("pubsub.publish_denied", map[string]any{ + "topic": evt.Topic, "from": sub.remote(), "reason": "receipt_failed", + }) + } + continue + } + } else if evt.Topic == GovernedTopic && b.requireReceipts { + // Start rejects this configuration, but preserve a per-publication + // fail-closed check for direct broker construction in tests or embedders. + slog.Warn("eventstream broker: publication receipt recorder is unavailable", "remote", sub.remote()) + continue + } + b.publish(published, sub) + } +} + +// governPublication unwraps and verifies the reserved envelope before fanout. +// Even while legacy publishing remains allowed, a peer cannot use the reserved +// topic to bypass verification or expose envelope contents to subscribers. +func (b *broker) governPublication(sub *subscriber, event *Event) (*Event, error) { + if event.Topic != GovernedTopic { + if b.requireGoverned { + return nil, fmt.Errorf("unsigned legacy publication rejected by governed broker") + } + return event, nil + } + if b.governedVerifier == nil { + return nil, fmt.Errorf("governed publication received but no verifier is configured") + } + governed, err := DecodeGovernedEvent(event) + if err != nil { + return nil, err + } + if err := b.governedVerifier.VerifyGovernedEvent(context.Background(), sub.conn.RemoteAddr(), governed); err != nil { + return nil, err + } + if b.replay != nil { + if err := b.replay.admit(governed.Intent); err != nil { + slog.Warn("eventstream governed publication replay rejected", "agent_id", governed.Intent.AgentID, "intent_id", governed.Intent.ID, "error", err) + return nil, fmt.Errorf("eventstream: %w", err) + } + } + if b.governedTransferQuota != nil { + if err := b.governedTransferQuota.Allow(governed.Intent.AgentID, uint64(len(governed.Payload))); err != nil { + slog.Warn("eventstream governed transfer quota rejected", "agent_id", governed.Intent.AgentID, "bytes", len(governed.Payload), "error", err) + return nil, fmt.Errorf("eventstream: governed transfer quota rejected") + } + } + if b.contentInspector != nil { + contentType := "application/octet-stream" + if governed.Disclosure != nil { + contentType = governed.Disclosure.ContentType + } + if err := b.contentInspector.InspectDisclosureContent(context.Background(), governed.Intent, governed.Disclosure, contentType, "", bytes.NewReader(governed.Payload)); err != nil { + slog.Warn("eventstream governed content inspection rejected", "action", governed.Intent.Action, "resource", governed.Intent.Resource, "error", err) + return nil, fmt.Errorf("eventstream: governed content inspection rejected") + } } + return governed.Event(), nil } func (b *broker) addSub(topic string, sub *subscriber) bool { diff --git a/service_disabled.go b/service_disabled.go index d28d427..2e07237 100644 --- a/service_disabled.go +++ b/service_disabled.go @@ -14,6 +14,7 @@ import ( "context" "github.com/pilot-protocol/common/coreapi" + "github.com/pilot-protocol/common/decision" ) // Service is a no-op replacement for the real plugin Service. @@ -23,6 +24,18 @@ type Service struct{} // the real NewService. func NewService() *Service { return &Service{} } +// SetTopicPolicy, SetGovernedPublication, and SetGovernedReceiptRecorder preserve the enabled Service API +// in a no_eventstream build. Neither has an effect while the plugin is off. +func (s *Service) SetTopicPolicy(_ TopicPolicy) {} +func (s *Service) SetGovernedPublication(_ GovernedEventVerifier, _ bool) {} +func (s *Service) SetGovernedReceiptRecorder(_ GovernedReceiptRecorder, _ bool) {} + +func (s *Service) SetGovernedContentInspector(_ decision.DisclosureContentInspector) {} + +func (s *Service) RequireGovernedContentInspection(_ bool) {} + +func (s *Service) SetGovernedTransferQuota(_ *decision.TransferQuotaLimiter) {} + func (s *Service) Name() string { return "eventstream-disabled" } func (s *Service) Order() int { return 120 } func (s *Service) Start(_ context.Context, _ coreapi.Deps) error { return nil } diff --git a/zz_governed_replay_test.go b/zz_governed_replay_test.go new file mode 100644 index 0000000..df5667d --- /dev/null +++ b/zz_governed_replay_test.go @@ -0,0 +1,35 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +package eventstream + +import ( + "testing" + "time" + + "github.com/pilot-protocol/common/decision" +) + +// TestGovernedPublicationReplayGuard pins SECURITY_REVIEW_v1.14 finding M5 for +// the event broker: a verified governed intent may fan out at most once; a +// replay within its validity window is rejected, and distinct intents pass. +func TestGovernedPublicationReplayGuard(t *testing.T) { + base := time.Unix(1_700_000_000, 0) + clock := base + g := newGovernedReplayGuard() + g.now = func() time.Time { return clock } + + intent := decision.Intent{ID: "evt-1", TenantID: "t", AgentID: "pub", ExpiresAt: base.Add(2 * time.Minute).Unix()} + if err := g.admit(intent); err != nil { + t.Fatalf("first publication rejected: %v", err) + } + if err := g.admit(intent); err == nil { + t.Fatal("replayed governed publication was ACCEPTED") + } + if err := g.admit(decision.Intent{ID: "evt-2", TenantID: "t", AgentID: "pub", ExpiresAt: clock.Add(2 * time.Minute).Unix()}); err != nil { + t.Fatalf("distinct intent rejected: %v", err) + } + clock = base.Add(3 * time.Minute) + if err := g.admit(decision.Intent{ID: "evt-1", TenantID: "t", AgentID: "pub", ExpiresAt: clock.Add(2 * time.Minute).Unix()}); err != nil { + t.Fatalf("post-expiry reuse rejected: %v", err) + } +}