From 0e4214f320ee8289f7eff8b9b831e9b4528d5099 Mon Sep 17 00:00:00 2001 From: Teodor Calin Date: Sun, 2 Aug 2026 19:08:56 +0300 Subject: [PATCH 1/3] dataexchange: WIP governed delivery frames + retention MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Preservation snapshot of uncommitted working-tree work (repo survey 2026-08-02). WIP branch — do NOT push directly; split into reviewed PRs first. Build scratch excluded via .gitignore. Co-Authored-By: Claude Fable 5 --- README.md | 90 +++++- client.go | 110 +++++++ dataexchange.go | 12 + filestream.go | 82 +++++- governed.go | 318 ++++++++++++++++++++ governed_stream.go | 285 ++++++++++++++++++ governed_test.go | 686 ++++++++++++++++++++++++++++++++++++++++++++ retention.go | 295 +++++++++++++++++++ retention_test.go | 171 +++++++++++ service.go | 567 ++++++++++++++++++++++++++++++------ service_disabled.go | 26 +- zz_coverage_test.go | 7 +- zz_more_test.go | 3 + 13 files changed, 2550 insertions(+), 102 deletions(-) create mode 100644 governed.go create mode 100644 governed_stream.go create mode 100644 governed_test.go create mode 100644 retention.go create mode 100644 retention_test.go diff --git a/README.md b/README.md index db93962..f5b20c5 100644 --- a/README.md +++ b/README.md @@ -30,8 +30,9 @@ rt.Register(dataexchange.NewService(dataexchange.ServiceConfig{})) | File | What it does | |---|---| -| `dataexchange.go` | Wire format: `Frame`, `WriteFrame`, `ReadFrame`, `TraceFrame`, `TypeText/Binary/JSON/File/Trace`, `TypeName`. | +| `dataexchange.go` | Wire format: `Frame`, `WriteFrame`, `ReadFrame`, `TraceFrame`, `TypeText/Binary/JSON/File/Trace/Governed`, `TypeName`. | | `client.go` | `Client` — `Dial` and send helpers. | +| `governed.go` | Signed decision envelope, receiver-side verifier, and enforceable transport constraints. | | `server.go` | `Server` — accept loop and handler dispatch. | | `service.go` | `*Service` — `coreapi.Service` adapter. Build tag `!no_dataexchange`. | | `service_disabled.go` | Stub `*Service` for `-tags no_dataexchange` builds. | @@ -45,7 +46,92 @@ rt.Register(dataexchange.NewService(dataexchange.ServiceConfig{})) For `TypeFile` the payload is prefixed with `[2-byte name length][name bytes]`. For `TypeTrace` the payload is `[4-byte inner_type][8-byte sent_at_ns][inner payload]`. -Max frame size: 256 MiB. +Max frame size: 64 MiB by default (configurable at process start with +`PILOT_DATAEXCHANGE_MAX_FRAME` within its documented safe range). + +## Governed delivery + +`TypeGoverned` wraps one text, JSON, binary, or single-frame file delivery in +the sender's signed `decision.Intent` and signed `decision.Decision`. The +intent payload hash binds the frame type, filename, and exact bytes; the +receiver verifies the signatures, tenant authority state, local deterministic +ceiling, exact local destination, and any applicable transport constraints +before it writes to disk. A workflow-approved action uses the same short-lived +execution Decision as an ordinary allowed action—there is no reusable +transport permit. + +Large resumable files use `TypeGovernedFileStream`: the signed envelope binds +the exact `TypeFileStream` INIT (filename, declared length, full SHA-256, +chunk size, and transfer ID). Required receivers admit later chunks only for +that verified transfer on the same connection. Compute the file hash first, +call `BuildStreamInitPayload`, sign an Intent using +`GovernedStreamPayloadHash`, then call `SendGovernedFileStream`. + +For a required typed-disclosure profile, build a `decision.DisclosureBinding` +whose content hash, byte length, filename, and stream transfer ID match the +file, bind its canonical hash in the signed Intent, and use +`SendGovernedWithDisclosure` or `SendGovernedFileStreamWithDisclosure`. A +`DecisionFrameVerifier` with `RequireDisclosure` rejects governed messages, +single-frame files, and resumable stream INITs that omit this evidence. +When a receipt recorder is configured, typed deliveries require its V2 +disclosure-evidence method; the resulting signed receipt binds the canonical +disclosure hash without retaining the file or message body. + +Set `ServiceConfig.GovernedVerifier` to `DecisionFrameVerifier` (or an +equivalent local verifier). Once all senders have been upgraded, set +`ServiceConfig.RequireGoverned` to reject unsigned legacy deliveries. Roll out +in that order: upgraded receiver with verification available, upgraded +senders, then required mode. An older receiver does not understand +`TypeGoverned`; a required receiver intentionally rejects `TypeTrace` and +raw `TypeFileStream` INIT frames. Governed stream INITs are supported. + +For auditable enterprise ingress, configure `GovernedReceiptRecorder` and set +`RequireGovernedReceipts`. The service writes the received message/file first, +then requires the recorder to durably capture the exact signed Intent and +Decision before it emits the success ACK or delivery event. If recording fails, +the staged file is removed and the sender receives an error rather than a +successful but unreceipted delivery. + +### Local content inspection + +`ServiceConfig.GovernedContentInspector` is an optional receiver-local hook +that runs after the signed envelope and local policy ceiling are verified but +before a message/file is released. Set +`RequireGovernedContentInspection` to make startup fail unless the hook is +present; an inspection error removes a staged file or rejects the message. +`decision.PresidioInspector` is the included OSS adapter for bounded text, +JSON, XML, YAML, and form content. It rejects unsupported binary/document +types rather than truncating or silently skipping them. The inspector is local +to the receiver; neither the decision authority nor the sender's authority +receives plaintext for this check. + +Typed disclosure binding V2 adds a tenant-defined `retention_class` to the +same Intent hash. A signed policy can select allowed classes; the receiver sees +the bound metadata. Configure `GovernedRetentionPolicies` to map those classes +to local expiry durations. The service writes an owner-only retention journal +before a governed message/file becomes accepted (and before a streamed file's +final rename), then removes the content after expiry across restarts. An +unknown, V1, or unconfigured class is rejected when retention is enabled. +This is deletion retention, not a legal-hold or WORM-storage implementation. + +### Per-agent transfer quotas + +`ServiceConfig.GovernedTransferQuota` admits a bounded number of bytes and/or +actions for each signed `Intent.AgentID` in a fixed local window. It is charged +only after governed verification, including the declared bytes of a verified +stream INIT; a peer address cannot select or reset another agent's budget. +Quota is deliberately charged for an admitted attempt even if later local DLP +or receipt persistence rejects it, so repeatedly failing submissions cannot +turn the scanner into an unmetered denial-of-service target. + +The v1 action mapping is: + +| Frame | Intent action | +|---|---| +| text | `data.send.text` | +| JSON | `data.send.json` | +| binary | `data.send.binary` | +| file | `file.share` | ## Build tags diff --git a/client.go b/client.go index acddd60..69db791 100644 --- a/client.go +++ b/client.go @@ -3,8 +3,11 @@ package dataexchange import ( + "fmt" + "io" "time" + "github.com/pilot-protocol/common/decision" "github.com/pilot-protocol/common/driver" "github.com/pilot-protocol/common/protocol" ) @@ -43,6 +46,113 @@ func (c *Client) SendFile(name string, data []byte) error { return WriteFrame(c.conn, &Frame{Type: TypeFile, Filename: name, Payload: data}) } +// SendGoverned sends a frame with exact signed intent/decision evidence. +// The remote service must have RequireGoverned enabled to enforce it. +func (c *Client) SendGoverned(frame *Frame, intent decision.Intent, result decision.Decision) error { + governed, err := NewGovernedFrame(frame, intent, result) + if err != nil { + return err + } + envelope, err := EncodeGovernedFrame(governed) + if err != nil { + return err + } + return WriteFrame(c.conn, envelope) +} + +// SendGovernedWithDisclosure sends a single governed message or file with a +// typed disclosure binding. The Intent and Decision must already have been +// obtained for the exact binding. +func (c *Client) SendGovernedWithDisclosure(frame *Frame, intent decision.Intent, result decision.Decision, disclosure decision.DisclosureBinding) error { + governed, err := NewGovernedFrameWithDisclosure(frame, intent, result, disclosure) + if err != nil { + return err + } + envelope, err := EncodeGovernedFrame(governed) + if err != nil { + return err + } + return WriteFrame(c.conn, envelope) +} + +// SendGovernedFileStream sends a large resumable file with a signed authority +// decision bound to its exact INIT metadata. Callers compute the intent payload +// hash with GovernedStreamPayloadHash over the deterministic INIT payload (use +// BuildStreamInitPayload after calculating the file hash), then pass the same +// signed intent and decision here. The receiver accepts chunks only for that +// verified transfer ID and records its receipt on successful completion. +func (c *Client) SendGovernedFileStream(name string, r io.ReadSeeker, size int64, intent decision.Intent, result decision.Decision, stepTimeout time.Duration) (*StreamResult, error) { + return c.SendGovernedFileStreamWithAuthorizer(name, r, size, func(_ []byte) (decision.Intent, decision.Decision, error) { + return intent, result, nil + }, stepTimeout) +} + +// SendGovernedFileStreamWithDisclosure sends a governed resumable file whose +// Intent and authority Decision bind the supplied typed disclosure metadata. +// The metadata must match the final INIT (content hash, size, filename, and +// transfer ID) exactly; callers can derive these with BuildStreamInitPayload. +func (c *Client) SendGovernedFileStreamWithDisclosure(name string, r io.ReadSeeker, size int64, intent decision.Intent, result decision.Decision, disclosure decision.DisclosureBinding, stepTimeout time.Duration) (*StreamResult, error) { + return streamSendWithInit(c.conn, name, r, size, stepTimeout, func(id [transferIDLen]byte, declaredSize uint64, hash [32]byte, chunkSize uint32, filename string) (*Frame, error) { + init := encodeInit(id, declaredSize, hash, chunkSize, filename) + governed, err := NewGovernedStreamInitWithDisclosure(init, intent, result, disclosure) + if err != nil { + return nil, err + } + return EncodeGovernedStreamInit(governed) + }) +} + +// GovernedStreamAuthorizer receives the exact FileStream INIT payload after +// its stable transfer ID and content hash have been calculated. It must return +// a short-lived signed Intent and Decision bound to that exact payload. The +// callback runs before any stream frame is written, so a deny or unavailable +// authority cannot produce a partial transfer. +type GovernedStreamAuthorizer func(initPayload []byte) (decision.Intent, decision.Decision, error) + +// GovernedStreamDisclosureAuthorizer returns typed disclosure evidence for +// the exact INIT. Hosted federation clients use this after uploading the full +// file content and before the first stream byte is released to the peer. +type GovernedStreamDisclosureAuthorizer func(initPayload []byte) (decision.Intent, decision.Decision, decision.DisclosureBinding, error) + +// SendGovernedFileStreamWithAuthorizer obtains an exact signed decision only +// after the resumable stream's INIT bytes are known. This is the preferred +// sender API when the decision comes from an online authority. +func (c *Client) SendGovernedFileStreamWithAuthorizer(name string, r io.ReadSeeker, size int64, authorize GovernedStreamAuthorizer, stepTimeout time.Duration) (*StreamResult, error) { + if authorize == nil { + return nil, fmt.Errorf("dataexchange: governed stream authorizer is required") + } + return streamSendWithInit(c.conn, name, r, size, stepTimeout, func(id [transferIDLen]byte, declaredSize uint64, hash [32]byte, chunkSize uint32, filename string) (*Frame, error) { + init := encodeInit(id, declaredSize, hash, chunkSize, filename) + intent, result, err := authorize(append([]byte(nil), init.Payload...)) + if err != nil { + return nil, err + } + governed, err := NewGovernedStreamInit(init, intent, result) + if err != nil { + return nil, err + } + return EncodeGovernedStreamInit(governed) + }) +} + +func (c *Client) SendGovernedFileStreamWithDisclosureAuthorizer(name string, r io.ReadSeeker, size int64, authorize GovernedStreamDisclosureAuthorizer, stepTimeout time.Duration) (*StreamResult, error) { + if authorize == nil { + return nil, fmt.Errorf("dataexchange: governed stream disclosure authorizer is required") + } + return streamSendWithInit(c.conn, name, r, size, stepTimeout, func(id [transferIDLen]byte, declaredSize uint64, hash [32]byte, chunkSize uint32, filename string) (*Frame, error) { + init := encodeInit(id, declaredSize, hash, chunkSize, filename) + intent, result, disclosure, err := authorize(append([]byte(nil), init.Payload...)) + if err != nil { + return nil, err + } + governed, err := NewGovernedStreamInitWithDisclosure(init, intent, result, disclosure) + if err != nil { + return nil, err + } + return EncodeGovernedStreamInit(governed) + }) +} + // SendTrace wraps data in a TypeTrace frame with the current nanosecond clock. // Returns sentAtNs so the caller can correlate it against the timing ACK. func (c *Client) SendTrace(innerType uint32, data []byte) (sentAtNs int64, err error) { diff --git a/dataexchange.go b/dataexchange.go index d32c158..688b861 100644 --- a/dataexchange.go +++ b/dataexchange.go @@ -33,6 +33,14 @@ const ( // peer that does not understand TypeFileStream never sends INIT-ACK, so // the sender falls back to TypeFile. TypeFileStream uint32 = 7 + // TypeGoverned carries a regular data frame together with the sender's + // signed intent and authority decision. It is opt-in; enterprise receivers + // can require this envelope before persisting a message or file. + TypeGoverned uint32 = 8 + // TypeGovernedFileStream carries the signed authorization evidence for a + // TypeFileStream INIT. Subsequent chunks are accepted only while bound to + // that verified transfer ID on the same connection. + TypeGovernedFileStream uint32 = 9 ) // TraceFrame carries timing metadata around an inner message frame. @@ -253,6 +261,10 @@ func TypeName(t uint32) string { return "TRACE" case TypeFileStream: return "FILESTREAM" + case TypeGoverned: + return "GOVERNED" + case TypeGovernedFileStream: + return "GOVERNED_FILESTREAM" default: return fmt.Sprintf("UNKNOWN(%d)", t) } diff --git a/filestream.go b/filestream.go index df35821..4876b52 100644 --- a/filestream.go +++ b/filestream.go @@ -74,6 +74,19 @@ const ( // fresh connection. var ErrStreamUnsupported = errors.New("dataexchange: peer does not support TypeFileStream") +// BuildStreamInitPayload returns the deterministic INIT payload that a caller +// must bind in a signed file.share Intent before SendGovernedFileStream. The +// transfer ID is derived from the supplied full SHA-256, matching the sender's +// resumable-transfer state machine. +func BuildStreamInitPayload(name string, size int64, fullHash [32]byte) ([]byte, error) { + if size < 0 || !validGovernedFilename(name) { + return nil, fmt.Errorf("dataexchange: governed stream requires a non-negative size and safe filename") + } + var id [transferIDLen]byte + copy(id[:], fullHash[:transferIDLen]) + return append([]byte(nil), encodeInit(id, uint64(size), fullHash, uint32(StreamChunkSize), name).Payload...), nil +} + // --- control-frame codec --------------------------------------------------- func encodeStreamFrame(kind byte, id [transferIDLen]byte, body []byte) *Frame { @@ -116,7 +129,7 @@ func decodeInit(body []byte) (size uint64, hash [32]byte, chunkSize uint32, name copy(hash[:], body[8:40]) chunkSize = binary.BigEndian.Uint32(body[40:44]) nameLen := int(binary.BigEndian.Uint16(body[44:46])) - if 46+nameLen > len(body) { + if 46+nameLen != len(body) { return 0, hash, 0, "", false } name = string(body[46 : 46+nameLen]) @@ -200,6 +213,22 @@ func (c *Client) SendFileStream(name string, r io.ReadSeeker, size int64, stepTi } func streamSend(conn frameRW, name string, r io.ReadSeeker, size int64, stepTimeout time.Duration) (*StreamResult, error) { + return streamSendWithInit(conn, name, r, size, stepTimeout, func(id [transferIDLen]byte, declaredSize uint64, hash [32]byte, chunkSize uint32, filename string) (*Frame, error) { + return encodeInit(id, declaredSize, hash, chunkSize, filename), nil + }) +} + +// streamInitBuilder allows the sender to substitute a governed INIT envelope +// while keeping the exact same chunk/ACK/resume state machine. +type streamInitBuilder func([transferIDLen]byte, uint64, [32]byte, uint32, string) (*Frame, error) + +func streamSendWithInit(conn frameRW, name string, r io.ReadSeeker, size int64, stepTimeout time.Duration, buildInit streamInitBuilder) (*StreamResult, error) { + if buildInit == nil { + return nil, fmt.Errorf("dataexchange: stream INIT builder is required") + } + if size < 0 { + return nil, fmt.Errorf("dataexchange: stream size must not be negative") + } if stepTimeout <= 0 { stepTimeout = streamStepTimeout } @@ -219,7 +248,11 @@ func streamSend(conn frameRW, name string, r io.ReadSeeker, size int64, stepTime } // INIT + negotiate. - if err := WriteFrame(conn, encodeInit(id, uint64(size), fullHash, uint32(StreamChunkSize), name)); err != nil { + init, err := buildInit(id, uint64(size), fullHash, uint32(StreamChunkSize), name) + if err != nil { + return nil, fmt.Errorf("build INIT: %w", err) + } + if err := WriteFrame(conn, init); err != nil { return nil, fmt.Errorf("send INIT: %w", err) } initAck, err := recvFrameTimeout(conn, streamNegTimeout) @@ -390,7 +423,16 @@ func recvFrameTimeout(conn frameRW, d time.Duration) (*Frame, error) { type StreamReceiver struct { receivedDir string onSaved func(name, path string, size int64) - nameSuffix func(base string) string // produces the final unique filename + // onPrepare runs after full-content integrity verification but before the + // atomic rename. It lets a governed service durably record retention work + // for the final path before that path can become visible after a crash. + onPrepare func([transferIDLen]byte, string, string, int64) error + // onCommit runs after integrity verification and atomic rename, but before + // onSaved. A required receipt recorder uses it to make the visible file + // contingent on durable enforcement evidence; a callback error removes the + // just-renamed file and turns COMPLETE into a failure. + onCommit func([transferIDLen]byte, string, string, int64) error + nameSuffix func(base string) string // produces the final unique filename // quotaBytes caps total on-disk bytes under receivedDir (completed files // + .partial fragments). Enforced on INIT (declared size) and on every @@ -423,7 +465,7 @@ type recvTransfer struct { // onSaved (nil ok) fires after a verified file is renamed into place. No disk // quota is enforced — use NewStreamReceiverWithQuota to bound on-disk bytes. func NewStreamReceiver(receivedDir string, nameSuffix func(base string) string, onSaved func(name, path string, size int64)) *StreamReceiver { - return NewStreamReceiverWithQuota(receivedDir, nameSuffix, onSaved, 0) + return NewStreamReceiverWithQuotaAndCommit(receivedDir, nameSuffix, onSaved, nil, 0) } // NewStreamReceiverWithQuota is NewStreamReceiver with a disk quota: @@ -431,6 +473,21 @@ func NewStreamReceiver(receivedDir string, nameSuffix func(base string) string, // plus retained .partial fragments). Enforced on INIT and on every chunk // write so a peer cannot fill the disk mid-stream. Zero ⇒ unlimited. func NewStreamReceiverWithQuota(receivedDir string, nameSuffix func(base string) string, onSaved func(name, path string, size int64), quotaBytes int64) *StreamReceiver { + return NewStreamReceiverWithQuotaAndCommit(receivedDir, nameSuffix, onSaved, nil, quotaBytes) +} + +// NewStreamReceiverWithQuotaAndCommit extends the normal receiver with a +// transactional commit hook. It is intended for governed streams: the hook +// records evidence after the final file is durable but before consumers are +// notified. A hook failure removes the final file and reports COMPLETE failure. +func NewStreamReceiverWithQuotaAndCommit(receivedDir string, nameSuffix func(base string) string, onSaved func(name, path string, size int64), onCommit func([transferIDLen]byte, string, string, int64) error, quotaBytes int64) *StreamReceiver { + return NewStreamReceiverWithQuotaAndPrepareAndCommit(receivedDir, nameSuffix, onSaved, nil, onCommit, quotaBytes) +} + +// NewStreamReceiverWithQuotaAndPrepareAndCommit adds a pre-rename durable +// preparation hook to the governed commit path. A preparation error leaves no +// final file visible; callers may keep the partial for retry/inspection. +func NewStreamReceiverWithQuotaAndPrepareAndCommit(receivedDir string, nameSuffix func(base string) string, onSaved func(name, path string, size int64), onPrepare, onCommit func([transferIDLen]byte, string, string, int64) error, quotaBytes int64) *StreamReceiver { if nameSuffix == nil { nameSuffix = defaultStreamName } @@ -440,6 +497,8 @@ func NewStreamReceiverWithQuota(receivedDir string, nameSuffix func(base string) return &StreamReceiver{ receivedDir: receivedDir, onSaved: onSaved, + onPrepare: onPrepare, + onCommit: onCommit, nameSuffix: nameSuffix, quotaBytes: quotaBytes, transfers: make(map[[transferIDLen]byte]*recvTransfer), @@ -673,6 +732,12 @@ func (sr *StreamReceiver) handleDone(id [transferIDLen]byte) *Frame { _ = t.file.Close() finalName := sr.nameSuffix(t.name) finalPath := filepath.Join(sr.receivedDir, finalName) + if sr.onPrepare != nil { + if err := sr.onPrepare(id, finalName, finalPath, int64(t.size)); err != nil { + sr.forget(id) + return encodeComplete(id, false, "prepare: "+err.Error()) + } + } if err := os.Rename(t.partial, finalPath); err != nil { // Drop the in-memory transfer. t.file was already closed above, so // leaving the entry in place stranded a record holding a closed @@ -684,6 +749,15 @@ func (sr *StreamReceiver) handleDone(id [transferIDLen]byte) *Frame { sr.forget(id) return encodeComplete(id, false, "rename: "+err.Error()) } + if sr.onCommit != nil { + if err := sr.onCommit(id, finalName, finalPath, int64(t.size)); err != nil { + if removeErr := os.Remove(finalPath); removeErr != nil && !os.IsNotExist(removeErr) { + err = fmt.Errorf("%w; remove uncommitted file: %v", err, removeErr) + } + sr.forget(id) + return encodeComplete(id, false, "commit: "+err.Error()) + } + } sr.forget(id) if sr.onSaved != nil { sr.onSaved(finalName, finalPath, int64(t.size)) diff --git a/governed.go b/governed.go new file mode 100644 index 0000000..4eadb6c --- /dev/null +++ b/governed.go @@ -0,0 +1,318 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +package dataexchange + +import ( + "bytes" + "context" + "crypto/sha256" + "encoding/binary" + "encoding/json" + "fmt" + "io" + "path/filepath" + "strconv" + "strings" + "unicode/utf8" + + "github.com/pilot-protocol/common/coreapi" + "github.com/pilot-protocol/common/decision" +) + +const governedFrameVersion uint16 = 1 + +// GovernedFrame transports an ordinary data frame with the exact signed +// authorization evidence for that frame. PayloadHash is computed over frame +// type, filename, and bytes, so a decision for text cannot be replayed as a +// file or against different contents. +type GovernedFrame struct { + Version uint16 `json:"version"` + Type uint32 `json:"type"` + Filename string `json:"filename,omitempty"` + Payload []byte `json:"payload"` + Disclosure *decision.DisclosureBinding `json:"disclosure,omitempty"` + Intent decision.Intent `json:"intent"` + Decision decision.Decision `json:"decision"` +} + +// GovernedFrameVerifier maps a trusted transport peer and governed frame into +// a local enforcement decision. The service invokes it before any disk write. +type GovernedFrameVerifier interface { + VerifyGovernedFrame(context.Context, coreapi.Addr, GovernedFrame) error +} + +// GovernedReceiptRecorder durably records a receiver-side enforcement receipt +// for a verified governed delivery. It receives the exact signed objects that +// authorized the data frame; implementations must fail closed on a recording +// error when their deployment requires evidence before acknowledgement. +type GovernedReceiptRecorder interface { + RecordGovernedReceipt(context.Context, decision.Intent, decision.Decision) error +} + +// GovernedDisclosureReceiptRecorder is the evidence extension for a +// disclosure-bound transport. A configured recorder must implement this when +// it receives typed metadata: silently emitting a V1 receipt would lose the +// evidence a required disclosure profile relies on. +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("dataexchange: governed receipt recorder is not configured") + } + if disclosure == nil { + return recorder.RecordGovernedReceipt(ctx, intent, result) + } + typed, supported := recorder.(GovernedDisclosureReceiptRecorder) + if !supported { + return fmt.Errorf("dataexchange: governed receipt recorder does not support disclosure evidence") + } + return typed.RecordGovernedDisclosureReceipt(ctx, intent, result, *disclosure) +} + +// DecisionFrameVerifier is the reference receiver-side verifier. Resource +// must return the exact local resource identifier expected for this incoming +// frame (for example, "agent:finance/inbox"). This prevents a sender from +// reusing a valid decision for a different local destination. +type DecisionFrameVerifier struct { + Enforcer *decision.Enforcer + Resource func(coreapi.Addr, *Frame) string + RequireDisclosure bool +} + +func NewGovernedFrame(frame *Frame, intent decision.Intent, result decision.Decision) (GovernedFrame, error) { + if frame == nil { + return GovernedFrame{}, fmt.Errorf("dataexchange: governed frame is required") + } + governed := GovernedFrame{ + Version: governedFrameVersion, Type: frame.Type, Filename: frame.Filename, + Payload: append([]byte(nil), frame.Payload...), Intent: intent, Decision: result, + } + if err := governed.Validate(); err != nil { + return GovernedFrame{}, err + } + return governed, nil +} + +// NewGovernedFrameWithDisclosure creates a governed frame whose Intent binds +// a canonical disclosure metadata object. The caller must obtain the Decision +// for the disclosure-bound Intent; this constructor never changes authority. +func NewGovernedFrameWithDisclosure(frame *Frame, intent decision.Intent, result decision.Decision, disclosure decision.DisclosureBinding) (GovernedFrame, error) { + if frame == nil { + return GovernedFrame{}, fmt.Errorf("dataexchange: governed frame is required") + } + disclosure.Labels = append([]string(nil), disclosure.Labels...) + governed := GovernedFrame{ + Version: governedFrameVersion, Type: frame.Type, Filename: frame.Filename, + Payload: append([]byte(nil), frame.Payload...), Disclosure: &disclosure, Intent: intent, Decision: result, + } + if err := governed.Validate(); err != nil { + return GovernedFrame{}, err + } + return governed, nil +} + +func (frame GovernedFrame) Validate() error { + if frame.Version != governedFrameVersion || !governedDataType(frame.Type) { + return fmt.Errorf("dataexchange: invalid governed frame type") + } + if frame.Type == TypeFile { + if !validGovernedFilename(frame.Filename) { + return fmt.Errorf("dataexchange: invalid governed filename") + } + } else if frame.Filename != "" { + return fmt.Errorf("dataexchange: filename is only valid for file frames") + } + if uint64(len(frame.Payload)) > uint64(MaxFrameSize) { + return fmt.Errorf("dataexchange: governed payload exceeds maximum frame size") + } + if err := frame.Intent.Validate(); err != nil { + return fmt.Errorf("dataexchange: invalid governed intent: %w", err) + } + if frame.Intent.Signature == "" || frame.Decision.Signature == "" { + return fmt.Errorf("dataexchange: governed intent and decision must be signed") + } + if frame.Intent.Action != GovernedAction(frame.Type) { + return fmt.Errorf("dataexchange: governed intent action does not match frame type") + } + if err := frame.verifyPayloadBinding(); err != nil { + return err + } + if err := frame.Decision.Validate(); err != nil { + return fmt.Errorf("dataexchange: invalid governed decision: %w", err) + } + return nil +} + +func (frame GovernedFrame) verifyPayloadBinding() error { + if frame.Disclosure == nil { + if frame.Intent.PayloadHash != GovernedPayloadHash(frame.Type, frame.Filename, frame.Payload) { + return fmt.Errorf("dataexchange: governed intent payload binding mismatch") + } + return nil + } + if frame.Disclosure.ContentHash != decision.HashPayload(frame.Payload) || frame.Disclosure.DeclaredBytes != uint64(len(frame.Payload)) || frame.Disclosure.Filename != frame.Filename || frame.Disclosure.TransferID != "" { + return fmt.Errorf("dataexchange: governed disclosure does not match frame") + } + if err := frame.Disclosure.VerifyIntent(frame.Intent); err != nil { + return fmt.Errorf("dataexchange: governed disclosure intent binding: %w", err) + } + return nil +} + +func (frame GovernedFrame) DataFrame() *Frame { + return &Frame{Type: frame.Type, Filename: frame.Filename, Payload: append([]byte(nil), frame.Payload...)} +} + +func EncodeGovernedFrame(frame GovernedFrame) (*Frame, error) { + if err := frame.Validate(); err != nil { + return nil, err + } + body, err := json.Marshal(frame) + if err != nil { + return nil, fmt.Errorf("dataexchange: encode governed frame: %w", err) + } + if uint64(len(body)) > uint64(MaxFrameSize) { + return nil, fmt.Errorf("dataexchange: governed envelope exceeds maximum frame size") + } + return &Frame{Type: TypeGoverned, Payload: body}, nil +} + +func DecodeGovernedFrame(frame *Frame) (GovernedFrame, error) { + if frame == nil || frame.Type != TypeGoverned || uint64(len(frame.Payload)) > uint64(MaxFrameSize) { + return GovernedFrame{}, fmt.Errorf("dataexchange: invalid governed envelope") + } + decoder := json.NewDecoder(bytes.NewReader(frame.Payload)) + decoder.DisallowUnknownFields() + var governed GovernedFrame + if err := decoder.Decode(&governed); err != nil { + return GovernedFrame{}, fmt.Errorf("dataexchange: decode governed envelope: %w", err) + } + var trailing any + if err := decoder.Decode(&trailing); err != io.EOF { + return GovernedFrame{}, fmt.Errorf("dataexchange: trailing governed envelope data") + } + if err := governed.Validate(); err != nil { + return GovernedFrame{}, err + } + return governed, nil +} + +func (verifier DecisionFrameVerifier) VerifyGovernedFrame(ctx context.Context, remote coreapi.Addr, governed GovernedFrame) error { + if verifier.Enforcer == nil || verifier.Resource == nil { + return fmt.Errorf("dataexchange: decision frame verifier is not initialized") + } + if err := governed.Validate(); err != nil { + return err + } + if verifier.RequireDisclosure && governed.Disclosure == nil { + return fmt.Errorf("dataexchange: governed disclosure is required") + } + frame := governed.DataFrame() + resource := verifier.Resource(remote, frame) + if resource == "" || governed.Intent.Resource != resource { + return fmt.Errorf("dataexchange: 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("dataexchange: verify governed decision: %w", verifyErr) + } + switch governed.Decision.Outcome { + case decision.Allow: + return nil + case decision.Constrain: + return enforceFrameConstraints(governed.Decision.Constraints, remote, frame) + default: + return fmt.Errorf("dataexchange: governed decision outcome %q cannot permit delivery", governed.Decision.Outcome) + } +} + +// GovernedAction returns the only signed Intent action that can authorize a +// regular governed data frame of frameType. An empty result means the type is +// not eligible for a governed envelope. +func GovernedAction(frameType uint32) string { + switch frameType { + case TypeText: + return "data.send.text" + case TypeJSON: + return "data.send.json" + case TypeBinary: + return "data.send.binary" + case TypeFile: + return "file.share" + default: + return "" + } +} + +func governedAction(frameType uint32) string { return GovernedAction(frameType) } + +func governedDataType(frameType uint32) bool { return GovernedAction(frameType) != "" } + +// GovernedPayloadHash returns the exact payload binding required by a +// TypeGoverned Intent. It includes the frame kind and filename as well as the +// bytes, so callers must use it when creating the Intent before SendGoverned. +func GovernedPayloadHash(frameType uint32, filename string, payload []byte) string { + var header [12]byte + binary.BigEndian.PutUint32(header[:4], frameType) + binary.BigEndian.PutUint32(header[4:8], uint32(len(filename))) + binary.BigEndian.PutUint32(header[8:12], uint32(len(payload))) + hash := sha256.New() + _, _ = hash.Write([]byte("pilot-dataexchange-governed-frame-v1\x00")) + _, _ = hash.Write(header[:]) + _, _ = hash.Write([]byte(filename)) + _, _ = hash.Write(payload) + return decision.HashPayload(hash.Sum(nil)) +} + +func validGovernedFilename(name string) bool { + return name != "" && utf8.ValidString(name) && len(name) <= maxFilenameLen && filepath.Base(name) == name && !strings.ContainsAny(name, "/\\") +} + +func enforceFrameConstraints(constraints []decision.Constraint, remote coreapi.Addr, frame *Frame) error { + attributes := map[string]string{ + "sender": remote.String(), "frame_type": TypeName(frame.Type), "bytes": strconv.Itoa(len(frame.Payload)), "filename": frame.Filename, + } + for _, constraint := range constraints { + actual, found := attributes[constraint.Key] + if !found { + return fmt.Errorf("dataexchange: constraint %q has no enforceable frame attribute", constraint.Key) + } + switch constraint.Operator { + case "eq": + if actual != constraint.Value { + return fmt.Errorf("dataexchange: 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("dataexchange: 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("dataexchange: numeric constraint %s rejected", constraint.Key) + } + case "require": + if constraint.Value != "" && actual != constraint.Value { + return fmt.Errorf("dataexchange: required constraint %s rejected", constraint.Key) + } + default: + return fmt.Errorf("dataexchange: constraint operator %q is not enforceable", constraint.Operator) + } + } + return nil +} diff --git a/governed_stream.go b/governed_stream.go new file mode 100644 index 0000000..67437ba --- /dev/null +++ b/governed_stream.go @@ -0,0 +1,285 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +package dataexchange + +import ( + "bytes" + "context" + "encoding/binary" + "encoding/hex" + "encoding/json" + "fmt" + "io" + "strconv" + "strings" + + "github.com/pilot-protocol/common/coreapi" + "github.com/pilot-protocol/common/decision" +) + +const governedStreamVersion uint16 = 1 + +// GovernedStreamInit binds a signed authority decision to the exact INIT +// frame for a resumable file stream. The INIT includes the full file hash, +// filename, declared length, and transfer ID; individual chunks are accepted +// only after this envelope has been verified locally. +type GovernedStreamInit struct { + Version uint16 `json:"version"` + InitPayload []byte `json:"init_payload"` + Disclosure *decision.DisclosureBinding `json:"disclosure,omitempty"` + Intent decision.Intent `json:"intent"` + Decision decision.Decision `json:"decision"` +} + +// GovernedStreamVerifier verifies a stream INIT before the receiver creates a +// partial file. It is intentionally distinct from GovernedFrameVerifier so +// embedders that do not support large files need not accidentally authorize +// them. +type GovernedStreamVerifier interface { + VerifyGovernedStreamInit(context.Context, coreapi.Addr, GovernedStreamInit) error +} + +func NewGovernedStreamInit(init *Frame, intent decision.Intent, result decision.Decision) (GovernedStreamInit, error) { + if init == nil || init.Type != TypeFileStream { + return GovernedStreamInit{}, fmt.Errorf("dataexchange: stream INIT frame is required") + } + governed := GovernedStreamInit{ + Version: governedStreamVersion, InitPayload: append([]byte(nil), init.Payload...), Intent: intent, Decision: result, + } + if err := governed.Validate(); err != nil { + return GovernedStreamInit{}, err + } + return governed, nil +} + +// NewGovernedStreamInitWithDisclosure creates a governed resumable transfer +// whose signed Intent binds typed metadata as well as the content hash, +// filename, declared bytes, and stable transfer ID in its INIT frame. +func NewGovernedStreamInitWithDisclosure(init *Frame, intent decision.Intent, result decision.Decision, disclosure decision.DisclosureBinding) (GovernedStreamInit, error) { + if init == nil || init.Type != TypeFileStream { + return GovernedStreamInit{}, fmt.Errorf("dataexchange: stream INIT frame is required") + } + disclosure.Labels = append([]string(nil), disclosure.Labels...) + governed := GovernedStreamInit{ + Version: governedStreamVersion, InitPayload: append([]byte(nil), init.Payload...), Disclosure: &disclosure, Intent: intent, Decision: result, + } + if err := governed.Validate(); err != nil { + return GovernedStreamInit{}, err + } + return governed, nil +} + +func (stream GovernedStreamInit) Validate() error { + if stream.Version != governedStreamVersion || len(stream.InitPayload) == 0 || uint64(len(stream.InitPayload)) > uint64(MaxFrameSize) { + return fmt.Errorf("dataexchange: invalid governed stream INIT") + } + _, id, body, ok := decodeStreamFrame(&Frame{Type: TypeFileStream, Payload: stream.InitPayload}) + if !ok || stream.InitPayload[0] != streamKindInit { + return fmt.Errorf("dataexchange: governed stream must contain a file-stream INIT") + } + size, hash, chunkSize, name, valid := decodeInit(body) + if !valid || !validGovernedFilename(name) || chunkSize == 0 || size > uint64(^uint64(0)>>1) { + return fmt.Errorf("dataexchange: invalid governed stream metadata") + } + if !bytes.Equal(id[:], hash[:transferIDLen]) { + return fmt.Errorf("dataexchange: governed stream transfer ID does not match content hash") + } + if err := stream.Intent.Validate(); err != nil { + return fmt.Errorf("dataexchange: invalid governed stream intent: %w", err) + } + if stream.Intent.Signature == "" || stream.Decision.Signature == "" { + return fmt.Errorf("dataexchange: governed stream intent and decision must be signed") + } + if stream.Intent.Action != "file.share" { + return fmt.Errorf("dataexchange: governed stream intent action mismatch") + } + if stream.Disclosure == nil && stream.Intent.PayloadHash != GovernedStreamPayloadHash(stream.InitPayload) { + return fmt.Errorf("dataexchange: governed stream intent binding mismatch") + } + if stream.Disclosure != nil { + if stream.Disclosure.ContentHash != hex.EncodeToString(hash[:]) || stream.Disclosure.DeclaredBytes != size || stream.Disclosure.Filename != name || stream.Disclosure.TransferID != hex.EncodeToString(id[:]) { + return fmt.Errorf("dataexchange: governed stream disclosure does not match INIT") + } + if err := stream.Disclosure.VerifyIntent(stream.Intent); err != nil { + return fmt.Errorf("dataexchange: governed stream disclosure intent binding: %w", err) + } + } + if err := stream.Decision.Validate(); err != nil { + return fmt.Errorf("dataexchange: invalid governed stream decision: %w", err) + } + return nil +} + +func (stream GovernedStreamInit) InitFrame() *Frame { + return &Frame{Type: TypeFileStream, Payload: append([]byte(nil), stream.InitPayload...)} +} + +func (stream GovernedStreamInit) TransferID() ([transferIDLen]byte, error) { + _, id, _, ok := decodeStreamFrame(stream.InitFrame()) + if !ok { + return [transferIDLen]byte{}, fmt.Errorf("dataexchange: invalid governed stream transfer ID") + } + return id, nil +} + +// DeclaredBytes returns the signed stream length from its verified INIT. It is +// used by receiver-local admission controls before any stream chunks are +// accepted. +func (stream GovernedStreamInit) DeclaredBytes() (uint64, error) { + if err := stream.Validate(); err != nil { + return 0, err + } + _, _, body, ok := decodeStreamFrame(stream.InitFrame()) + if !ok { + return 0, fmt.Errorf("dataexchange: invalid governed stream INIT") + } + size, _, _, _, valid := decodeInit(body) + if !valid { + return 0, fmt.Errorf("dataexchange: invalid governed stream metadata") + } + return size, nil +} + +func EncodeGovernedStreamInit(stream GovernedStreamInit) (*Frame, error) { + if err := stream.Validate(); err != nil { + return nil, err + } + body, err := json.Marshal(stream) + if err != nil { + return nil, fmt.Errorf("dataexchange: encode governed stream: %w", err) + } + if uint64(len(body)) > uint64(MaxFrameSize) { + return nil, fmt.Errorf("dataexchange: governed stream envelope exceeds maximum frame size") + } + return &Frame{Type: TypeGovernedFileStream, Payload: body}, nil +} + +func DecodeGovernedStreamInit(frame *Frame) (GovernedStreamInit, error) { + if frame == nil || frame.Type != TypeGovernedFileStream || uint64(len(frame.Payload)) > uint64(MaxFrameSize) { + return GovernedStreamInit{}, fmt.Errorf("dataexchange: invalid governed stream envelope") + } + decoder := json.NewDecoder(bytes.NewReader(frame.Payload)) + decoder.DisallowUnknownFields() + var stream GovernedStreamInit + if err := decoder.Decode(&stream); err != nil { + return GovernedStreamInit{}, fmt.Errorf("dataexchange: decode governed stream envelope: %w", err) + } + var trailing any + if err := decoder.Decode(&trailing); err != io.EOF { + return GovernedStreamInit{}, fmt.Errorf("dataexchange: trailing governed stream envelope data") + } + if err := stream.Validate(); err != nil { + return GovernedStreamInit{}, err + } + return stream, nil +} + +// GovernedStreamPayloadHash binds the exact bytes of the file-stream INIT +// frame. It therefore covers the filename, size, full content hash, chunk +// size, and transfer ID without having to duplicate the stream wire grammar. +func GovernedStreamPayloadHash(initPayload []byte) string { + var length [4]byte + binary.BigEndian.PutUint32(length[:], uint32(len(initPayload))) + buffer := make([]byte, 0, len("pilot-dataexchange-governed-stream-v1\x00")+len(length)+len(initPayload)) + buffer = append(buffer, []byte("pilot-dataexchange-governed-stream-v1\x00")...) + buffer = append(buffer, length[:]...) + buffer = append(buffer, initPayload...) + return decision.HashPayload(buffer) +} + +// GovernedStreamDisclosureMetadata returns the canonical values a hosted +// federation disclosure must bind for an already-built stream INIT. +func GovernedStreamDisclosureMetadata(initPayload []byte) (transferID string, declaredBytes uint64, contentHash string, filename string, err error) { + _, id, body, ok := decodeStreamFrame(&Frame{Type: TypeFileStream, Payload: initPayload}) + if !ok || len(initPayload) == 0 || initPayload[0] != streamKindInit { + err = fmt.Errorf("dataexchange: invalid file-stream INIT") + return + } + size, hash, _, name, valid := decodeInit(body) + if !valid || !validGovernedFilename(name) || !bytes.Equal(id[:], hash[:transferIDLen]) { + err = fmt.Errorf("dataexchange: invalid file-stream disclosure metadata") + return + } + return hex.EncodeToString(id[:]), size, hex.EncodeToString(hash[:]), name, nil +} + +func (verifier DecisionFrameVerifier) VerifyGovernedStreamInit(ctx context.Context, remote coreapi.Addr, stream GovernedStreamInit) error { + if verifier.Enforcer == nil || verifier.Resource == nil { + return fmt.Errorf("dataexchange: decision stream verifier is not initialized") + } + if err := stream.Validate(); err != nil { + return err + } + if verifier.RequireDisclosure && stream.Disclosure == nil { + return fmt.Errorf("dataexchange: governed stream disclosure is required") + } + init := stream.InitFrame() + _, _, body, _ := decodeStreamFrame(init) + size, _, _, name, _ := decodeInit(body) + resourceFrame := &Frame{Type: TypeFileStream, Filename: name} + resource := verifier.Resource(remote, resourceFrame) + if resource == "" || stream.Intent.Resource != resource { + return fmt.Errorf("dataexchange: governed stream intent resource binding mismatch") + } + var verifyErr error + if stream.Disclosure != nil { + verifyErr = verifier.Enforcer.VerifyDisclosure(ctx, stream.Intent, stream.Decision, *stream.Disclosure) + } else { + verifyErr = verifier.Enforcer.Verify(ctx, stream.Intent, stream.Decision) + } + if verifyErr != nil { + return fmt.Errorf("dataexchange: verify governed stream decision: %w", verifyErr) + } + switch stream.Decision.Outcome { + case decision.Allow: + return nil + case decision.Constrain: + return enforceStreamConstraints(stream.Decision.Constraints, remote, name, size) + default: + return fmt.Errorf("dataexchange: governed stream decision outcome %q cannot permit delivery", stream.Decision.Outcome) + } +} + +func enforceStreamConstraints(constraints []decision.Constraint, remote coreapi.Addr, filename string, size uint64) error { + attributes := map[string]string{ + "sender": remote.String(), "frame_type": TypeName(TypeFileStream), "bytes": strconv.FormatUint(size, 10), "filename": filename, + } + for _, constraint := range constraints { + actual, found := attributes[constraint.Key] + if !found { + return fmt.Errorf("dataexchange: constraint %q has no enforceable stream attribute", constraint.Key) + } + switch constraint.Operator { + case "eq": + if actual != constraint.Value { + return fmt.Errorf("dataexchange: stream 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("dataexchange: stream 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("dataexchange: numeric stream constraint %s rejected", constraint.Key) + } + case "require": + if constraint.Value != "" && actual != constraint.Value { + return fmt.Errorf("dataexchange: required stream constraint %s rejected", constraint.Key) + } + default: + return fmt.Errorf("dataexchange: constraint operator %q is not enforceable for a stream", constraint.Operator) + } + } + return nil +} + +var _ GovernedStreamVerifier = DecisionFrameVerifier{} diff --git a/governed_test.go b/governed_test.go new file mode 100644 index 0000000..a2cef03 --- /dev/null +++ b/governed_test.go @@ -0,0 +1,686 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +//go:build !no_dataexchange +// +build !no_dataexchange + +package dataexchange + +import ( + "bytes" + "context" + "crypto/ed25519" + "crypto/rand" + "crypto/sha256" + "encoding/json" + "errors" + "fmt" + "io" + "os" + "path/filepath" + "strings" + "testing" + "time" + + "github.com/pilot-protocol/common/coreapi" + "github.com/pilot-protocol/common/decision" +) + +type governedTestTrust struct { + intentKey ed25519.PublicKey + decisionKey ed25519.PublicKey +} + +type noWriteStream struct{ writes bytes.Buffer } + +func (stream *noWriteStream) Read([]byte) (int, error) { + return 0, errors.New("read should not occur before authorization") +} +func (stream *noWriteStream) Write(value []byte) (int, error) { + return stream.writes.Write(value) +} +func (stream *noWriteStream) Close() error { return nil } + +func (trust governedTestTrust) IntentKey(context.Context, string, string, string) (ed25519.PublicKey, error) { + return trust.intentKey, nil +} + +func (trust governedTestTrust) DecisionKey(context.Context, string, string) (ed25519.PublicKey, error) { + return trust.decisionKey, nil +} + +func (governedTestTrust) MinimumState(context.Context, string) (uint64, uint64, error) { + return 7, 3, nil +} + +type governedTestCeiling struct{} + +func (governedTestCeiling) Check(context.Context, decision.Intent, decision.Decision) error { + return nil +} + +func (governedTestCeiling) CheckDisclosure(context.Context, decision.Intent, decision.Decision, decision.DisclosureBinding) error { + return nil +} + +type governedReceiptRecorder struct { + calls int + intent decision.Intent + decision decision.Decision + err error +} + +type contentInspectorFunc func(context.Context, decision.Intent, *decision.DisclosureBinding, string, string, io.Reader) error + +func (inspect contentInspectorFunc) 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 TestGovernedContentInspectionIsLocalAndFailClosed(t *testing.T) { + frame := &Frame{Type: TypeFile, Filename: "invoice.pdf", Payload: []byte("classified invoice")} + governed, _ := newGovernedTestFrame(t, frame, decision.Allow, nil) + disclosure := decision.DisclosureBinding{Version: decision.DisclosureBindingVersion, ContentType: "application/pdf", Filename: frame.Filename} + governed.Disclosure = &disclosure + var observed []byte + service := NewService(ServiceConfig{GovernedContentInspector: contentInspectorFunc(func(_ context.Context, intent decision.Intent, got *decision.DisclosureBinding, contentType, filename string, content io.Reader) error { + if intent.ID != governed.Intent.ID || got == nil || contentType != "application/pdf" || filename != frame.Filename { + t.Fatalf("inspection metadata intent=%+v disclosure=%+v content_type=%q filename=%q", intent, got, contentType, filename) + } + var err error + observed, err = io.ReadAll(content) + return err + })}) + if err := service.inspectGovernedFrame(context.Background(), governed); err != nil || string(observed) != string(frame.Payload) { + t.Fatalf("inspection err=%v payload=%q", err, observed) + } + service.cfg.GovernedContentInspector = contentInspectorFunc(func(context.Context, decision.Intent, *decision.DisclosureBinding, string, string, io.Reader) error { + return errors.New("detector unavailable") + }) + if err := service.inspectGovernedFrame(context.Background(), governed); err == nil || err.Error() != "governed content inspection rejected" { + t.Fatalf("inspection failure leaked or was accepted: %v", err) + } + if err := (&Service{cfg: ServiceConfig{RequireGoverned: true, RequireGovernedContentInspection: true}}).validateGovernedConfig(); err == nil || !strings.Contains(err.Error(), "local inspector") { + t.Fatalf("required inspection without local hook err=%v", err) + } +} + +func TestGovernedTransferQuotaChargesOnlyVerifiedAgentIdentity(t *testing.T) { + limiter, err := decision.NewTransferQuotaLimiter(decision.TransferQuotaConfig{Window: time.Minute, MaxBytes: 5, MaxSenders: 2}) + if err != nil { + t.Fatal(err) + } + service := NewService(ServiceConfig{RequireGoverned: true, GovernedTransferQuota: limiter}) + if err := service.validateGovernedConfig(); err != nil { + t.Fatal(err) + } + if err := service.admitGovernedTransfer(decision.Intent{AgentID: "sender-a"}, 5); err != nil { + t.Fatal(err) + } + if err := service.admitGovernedTransfer(decision.Intent{AgentID: "sender-a"}, 1); err == nil || err.Error() != "governed transfer quota rejected" { + t.Fatalf("quota error=%v", err) + } + if err := (&Service{cfg: ServiceConfig{GovernedTransferQuota: limiter}}).validateGovernedConfig(); err == nil || !strings.Contains(err.Error(), "requires a governed receiver") { + t.Fatalf("legacy quota configuration error=%v", err) + } +} + +func (recorder *governedReceiptRecorder) RecordGovernedReceipt(_ context.Context, intent decision.Intent, result decision.Decision) error { + recorder.calls++ + recorder.intent, recorder.decision = intent, result + return recorder.err +} + +type legacyGovernedReceiptRecorder struct{} + +func (legacyGovernedReceiptRecorder) RecordGovernedReceipt(context.Context, decision.Intent, decision.Decision) error { + return nil +} + +type disclosureGovernedReceiptRecorder struct { + governedReceiptRecorder + disclosure decision.DisclosureBinding +} + +func (recorder *disclosureGovernedReceiptRecorder) 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 TestDisclosureReceiptRecorderRequiresV2Evidence(t *testing.T) { + disclosure := decision.DisclosureBinding{Version: decision.DisclosureBindingVersion} + if err := recordGovernedReceipt(context.Background(), legacyGovernedReceiptRecorder{}, 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 := &disclosureGovernedReceiptRecorder{} + 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) + } +} + +func newGovernedTestFrame(t *testing.T, frame *Frame, outcome decision.Outcome, constraints []decision.Constraint) (GovernedFrame, DecisionFrameVerifier) { + 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: "governed-intent", TenantID: "tenant-a", AgentID: "sender-a", + Action: governedAction(frame.Type), Resource: "agent:receiver/inbox", PayloadHash: GovernedPayloadHash(frame.Type, frame.Filename, frame.Payload), + Risk: decision.RiskMedium, IssuedAt: now.Unix(), ExpiresAt: now.Add(2 * time.Minute).Unix(), Nonce: nonce, KeyID: "sender-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: "governed-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 := NewGovernedFrame(frame, intent, result) + if err != nil { + t.Fatalf("new governed frame: %v", err) + } + verifier := DecisionFrameVerifier{ + Enforcer: &decision.Enforcer{ + Trust: governedTestTrust{intentKey: intentPublic, decisionKey: decisionPublic}, Ceiling: governedTestCeiling{}, Now: func() time.Time { return now }, + }, + Resource: func(_ coreapi.Addr, _ *Frame) string { return "agent:receiver/inbox" }, + } + return governed, verifier +} + +func newGovernedStreamForTest(t *testing.T, name string, payload []byte, outcome decision.Outcome, constraints []decision.Constraint) (GovernedStreamInit, DecisionFrameVerifier) { + t.Helper() + intentPublic, intentPrivate, err := ed25519.GenerateKey(rand.Reader) + if err != nil { + t.Fatal(err) + } + decisionPublic, decisionPrivate, err := ed25519.GenerateKey(rand.Reader) + if err != nil { + t.Fatal(err) + } + now := time.Now().UTC().Truncate(time.Second) + fullHash := sha256.Sum256(payload) + initPayload, err := BuildStreamInitPayload(name, int64(len(payload)), fullHash) + if err != nil { + t.Fatal(err) + } + nonce, err := decision.NewNonce() + if err != nil { + t.Fatal(err) + } + intent := decision.Intent{ + Version: decision.SchemaVersion, ID: "governed-stream-intent", TenantID: "tenant-a", AgentID: "sender-a", + Action: "file.share", Resource: "agent:receiver/inbox", PayloadHash: GovernedStreamPayloadHash(initPayload), + Risk: decision.RiskMedium, IssuedAt: now.Unix(), ExpiresAt: now.Add(2 * time.Minute).Unix(), Nonce: nonce, KeyID: "sender-key", + } + if err := intent.Sign(intentPrivate); err != nil { + t.Fatal(err) + } + intentHash, err := intent.Hash() + if err != nil { + t.Fatal(err) + } + result := decision.Decision{ + Version: decision.SchemaVersion, ID: "governed-stream-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.Fatal(err) + } + stream, err := NewGovernedStreamInit(&Frame{Type: TypeFileStream, Payload: initPayload}, intent, result) + if err != nil { + t.Fatal(err) + } + verifier := DecisionFrameVerifier{ + Enforcer: &decision.Enforcer{ + Trust: governedTestTrust{intentKey: intentPublic, decisionKey: decisionPublic}, Ceiling: governedTestCeiling{}, Now: func() time.Time { return now }, + }, + Resource: func(_ coreapi.Addr, _ *Frame) string { return "agent:receiver/inbox" }, + } + return stream, verifier +} + +func TestGovernedStreamAuthorizerReceivesExactInitBeforeAnyWrite(t *testing.T) { + payload := []byte("sensitive export") + stream := &noWriteStream{} + var initPayload []byte + denied := errors.New("policy denied") + _, err := streamSendWithInit(stream, "export.txt", bytes.NewReader(payload), int64(len(payload)), time.Second, func(id [transferIDLen]byte, size uint64, hash [32]byte, chunkSize uint32, name string) (*Frame, error) { + init := encodeInit(id, size, hash, chunkSize, name) + initPayload = append([]byte(nil), init.Payload...) + return nil, denied + }) + if !errors.Is(err, denied) { + t.Fatalf("err=%v, want denied authorization", err) + } + if len(stream.writes.Bytes()) != 0 { + t.Fatalf("authorization rejection wrote %d transport bytes", stream.writes.Len()) + } + fullHash := sha256.Sum256(payload) + want, err := BuildStreamInitPayload("export.txt", int64(len(payload)), fullHash) + if err != nil { + t.Fatal(err) + } + if !bytes.Equal(initPayload, want) { + t.Fatalf("authorizer init does not match governed payload binding") + } +} + +func TestGovernedFrameRoundTripBindsPayloadAndDecision(t *testing.T) { + frame := &Frame{Type: TypeFile, Filename: "report.pdf", Payload: []byte("approved report")} + governed, _ := newGovernedTestFrame(t, frame, decision.Allow, nil) + envelope, err := EncodeGovernedFrame(governed) + if err != nil { + t.Fatalf("encode governed frame: %v", err) + } + decoded, err := DecodeGovernedFrame(envelope) + if err != nil { + t.Fatalf("decode governed frame: %v", err) + } + if got := decoded.DataFrame(); got.Type != frame.Type || got.Filename != frame.Filename || string(got.Payload) != string(frame.Payload) { + t.Fatalf("decoded data frame = %#v, want %#v", got, frame) + } + + decoded.Payload = []byte("substituted report") + tamperedBody, err := json.Marshal(decoded) + if err != nil { + t.Fatalf("marshal tampered frame: %v", err) + } + if _, err := DecodeGovernedFrame(&Frame{Type: TypeGoverned, Payload: tamperedBody}); err == nil || !strings.Contains(err.Error(), "payload binding") { + t.Fatalf("tampered payload error = %v, want payload-binding failure", err) + } + + invalidName := governed + invalidName.Filename = string([]byte{0xff}) + if err := invalidName.Validate(); err == nil || !strings.Contains(err.Error(), "invalid governed filename") { + t.Fatalf("invalid filename error = %v, want UTF-8 validation failure", err) + } +} + +func TestGovernedFrameDisclosureBindingAndRequiredProfile(t *testing.T) { + frame := &Frame{Type: TypeFile, Filename: "report.json", Payload: []byte(`{"amount":42}`)} + governed, verifier := newGovernedTestFrame(t, frame, decision.Allow, nil) + strict := verifier + strict.RequireDisclosure = true + if err := strict.VerifyGovernedFrame(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(frame.Payload), DeclaredBytes: uint64(len(frame.Payload)), + ContentType: "application/json", Labels: []string{"finance", "pii"}, Recipient: "agent:finance", Purpose: "invoice-payment", Residency: "eu-west-1", Filename: frame.Filename, + } + hash, err := binding.Hash() + if err != nil { + t.Fatal(err) + } + // Constructing a fresh signed Intent/Decision is covered by the common + // disclosure tests; this transport test verifies the frame-level mismatch + // and required-profile checks before receiver-side signature resolution. + 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 = NewGovernedFrameWithDisclosure(frame, governed.Intent, governed.Decision, binding) + if err != nil { + t.Fatalf("valid disclosure envelope: %v", err) + } + tampered := governed + tampered.Disclosure = &binding + tampered.Disclosure.Residency = "us-east-1" + if err := tampered.Validate(); err == nil || !strings.Contains(err.Error(), "disclosure intent binding") { + t.Fatalf("disclosure mutation error=%v", err) + } +} + +func TestDecisionFrameVerifierBindsReceiverAndEnforcesConstraints(t *testing.T) { + frame := &Frame{Type: TypeText, Payload: []byte("hello")} + governed, verifier := newGovernedTestFrame(t, frame, decision.Constrain, []decision.Constraint{ + {Key: "frame_type", Operator: "eq", Value: "TEXT"}, + {Key: "bytes", Operator: "max", Value: "5"}, + }) + if err := verifier.VerifyGovernedFrame(context.Background(), coreapi.Addr{}, governed); err != nil { + t.Fatalf("verify governed frame: %v", err) + } + + wrongDestination := verifier + wrongDestination.Resource = func(_ coreapi.Addr, _ *Frame) string { return "agent:other/inbox" } + if err := wrongDestination.VerifyGovernedFrame(context.Background(), coreapi.Addr{}, governed); err == nil || !strings.Contains(err.Error(), "resource binding") { + t.Fatalf("wrong destination error = %v, want resource-binding failure", err) + } + + tooLarge, constrainedVerifier := newGovernedTestFrame(t, &Frame{Type: TypeText, Payload: []byte("too large")}, decision.Constrain, []decision.Constraint{{Key: "bytes", Operator: "max", Value: "3"}}) + if err := constrainedVerifier.VerifyGovernedFrame(context.Background(), coreapi.Addr{}, tooLarge); err == nil || !strings.Contains(err.Error(), "numeric constraint") { + t.Fatalf("oversize payload error = %v, want constraint failure", err) + } +} + +func TestGovernedStreamInitBindsExactMetadataAndConstraints(t *testing.T) { + stream, verifier := newGovernedStreamForTest(t, "large.bin", []byte("governed streaming payload"), decision.Constrain, []decision.Constraint{{Key: "bytes", Operator: "max", Value: "32"}}) + envelope, err := EncodeGovernedStreamInit(stream) + if err != nil { + t.Fatal(err) + } + decoded, err := DecodeGovernedStreamInit(envelope) + if err != nil { + t.Fatal(err) + } + if err := verifier.VerifyGovernedStreamInit(context.Background(), coreapi.Addr{}, decoded); err != nil { + t.Fatalf("verify governed stream: %v", err) + } + tampered := decoded + tampered.InitPayload[len(tampered.InitPayload)-1] ^= 0x01 + body, err := json.Marshal(tampered) + if err != nil { + t.Fatal(err) + } + if _, err := DecodeGovernedStreamInit(&Frame{Type: TypeGovernedFileStream, Payload: body}); err == nil || !strings.Contains(err.Error(), "binding") { + t.Fatalf("tampered stream error=%v", err) + } + tooLarge, constrainedVerifier := newGovernedStreamForTest(t, "large.bin", bytes.Repeat([]byte("x"), 33), decision.Constrain, []decision.Constraint{{Key: "bytes", Operator: "max", Value: "32"}}) + if err := constrainedVerifier.VerifyGovernedStreamInit(context.Background(), coreapi.Addr{}, tooLarge); err == nil || !strings.Contains(err.Error(), "numeric stream constraint") { + t.Fatalf("oversize stream constraint error=%v", err) + } +} + +func TestGovernedStreamDisclosureBindingAndRequiredProfile(t *testing.T) { + stream, verifier := newGovernedStreamForTest(t, "export.txt", []byte("sensitive export"), decision.Allow, nil) + verifier.RequireDisclosure = true + if err := verifier.VerifyGovernedStreamInit(context.Background(), coreapi.Addr{}, stream); err == nil || !strings.Contains(err.Error(), "disclosure is required") { + t.Fatalf("stream without required disclosure err=%v", err) + } + + payload := []byte("sensitive export") + fullHash := sha256.Sum256(payload) + initPayload, err := BuildStreamInitPayload("export.txt", int64(len(payload)), fullHash) + if err != nil { + t.Fatal(err) + } + disclosure := decision.DisclosureBinding{ + Version: decision.DisclosureBindingVersion, ContentHash: decision.HashPayload(payload), DeclaredBytes: uint64(len(payload)), + ContentType: "text/plain", Labels: []string{"confidential", "pii"}, Recipient: "agent:receiver", + Purpose: "customer-support", Residency: "eu-west-1", Filename: "export.txt", TransferID: fmt.Sprintf("%x", fullHash[:transferIDLen]), + } + disclosureHash, err := disclosure.Hash() + if err != nil { + t.Fatal(err) + } + intent := decision.Intent{ + Version: decision.SchemaVersion, ID: "governed-stream-disclosure-intent", TenantID: "tenant-a", AgentID: "sender-a", + Action: "file.share", Resource: "agent:receiver/inbox", Audience: disclosure.Recipient, Purpose: disclosure.Purpose, + PayloadHash: disclosureHash, Risk: decision.RiskHigh, IssuedAt: 1785500000, ExpiresAt: 1785500060, + Nonce: strings.Repeat("a", 32), KeyID: "sender-key", Signature: "test-signature", + } + intentHash, err := intent.Hash() + if err != nil { + t.Fatal(err) + } + result := decision.Decision{ + Version: decision.SchemaVersion, ID: "governed-stream-disclosure-decision", IntentHash: intentHash, TenantID: intent.TenantID, AgentID: intent.AgentID, + Outcome: decision.Allow, PolicyRevision: 7, RevocationEpoch: 3, ProviderID: "authority-a", IssuedAt: 1785500000, ExpiresAt: 1785500060, + KeyID: "authority-key", Signature: "test-signature", + } + governed, err := NewGovernedStreamInitWithDisclosure(&Frame{Type: TypeFileStream, Payload: initPayload}, intent, result, disclosure) + if err != nil { + t.Fatal(err) + } + encoded, err := EncodeGovernedStreamInit(governed) + if err != nil { + t.Fatal(err) + } + decoded, err := DecodeGovernedStreamInit(encoded) + if err != nil || decoded.Disclosure == nil || decoded.Disclosure.TransferID != disclosure.TransferID { + t.Fatalf("decode disclosure stream err=%v decoded=%+v", err, decoded) + } + decoded.Disclosure.Residency = "us-east-1" + if err := decoded.Validate(); err == nil || !strings.Contains(err.Error(), "binding") { + t.Fatalf("residency mutation accepted: %v", err) + } +} + +func TestRequireGovernedServicePersistsOnlyVerifiedFrames(t *testing.T) { + tmp := t.TempDir() + governed, verifier := newGovernedTestFrame(t, &Frame{Type: TypeText, Payload: []byte("approved")}, decision.Allow, nil) + w, r, wait := makeServiceConn(t, ServiceConfig{InboxDir: tmp, RequireGoverned: true, GovernedVerifier: verifier}) + defer wait() + + envelope, err := EncodeGovernedFrame(governed) + if err != nil { + t.Fatalf("encode governed frame: %v", err) + } + if err := WriteFrame(w, envelope); err != nil { + t.Fatalf("write governed frame: %v", err) + } + ack, err := ReadFrame(r) + if err != nil { + t.Fatalf("read governed ack: %v", err) + } + if !strings.Contains(string(ack.Payload), "ACK TEXT 8 bytes") { + t.Fatalf("governed ack = %q", ack.Payload) + } + entries, err := os.ReadDir(tmp) + if err != nil || len(entries) != 1 { + t.Fatalf("inbox after governed delivery: entries=%d err=%v", len(entries), err) + } + + if err := WriteFrame(w, &Frame{Type: TypeText, Payload: []byte("unsigned")}); err != nil { + t.Fatalf("write unsigned frame: %v", err) + } + ack, err = ReadFrame(r) + if err != nil { + t.Fatalf("read unsigned ack: %v", err) + } + if !strings.Contains(string(ack.Payload), "ERR TEXT save failed: unsigned legacy frame rejected") { + t.Fatalf("unsigned ack = %q", ack.Payload) + } + entries, err = os.ReadDir(tmp) + if err != nil || len(entries) != 1 { + t.Fatalf("unsigned delivery changed inbox: entries=%d err=%v", len(entries), err) + } +} + +func TestGovernedReceiptIsRequiredBeforeDeliveryAcknowledgement(t *testing.T) { + governed, verifier := newGovernedTestFrame(t, &Frame{Type: TypeText, Payload: []byte("approved")}, decision.Allow, nil) + envelope, err := EncodeGovernedFrame(governed) + if err != nil { + t.Fatal(err) + } + recorder := &governedReceiptRecorder{} + directory := t.TempDir() + w, r, wait := makeServiceConn(t, ServiceConfig{ + InboxDir: directory, RequireGoverned: true, GovernedVerifier: verifier, + RequireGovernedReceipts: true, GovernedReceiptRecorder: recorder, + }) + defer wait() + if err := WriteFrame(w, envelope); err != nil { + t.Fatal(err) + } + ack, err := ReadFrame(r) + if err != nil || !strings.Contains(string(ack.Payload), "ACK TEXT 8 bytes") { + t.Fatalf("receipt-backed acknowledgement=%q err=%v", ack.Payload, err) + } + if recorder.calls != 1 || recorder.intent.ID != governed.Intent.ID || recorder.decision.ID != governed.Decision.ID { + t.Fatalf("receipt recorder=%+v", recorder) + } + entries, err := os.ReadDir(directory) + if err != nil || len(entries) != 1 { + t.Fatalf("receipt-backed delivery entries=%d err=%v", len(entries), err) + } + + failingDirectory := t.TempDir() + failing := &governedReceiptRecorder{err: os.ErrPermission} + w, r, wait = makeServiceConn(t, ServiceConfig{ + InboxDir: failingDirectory, RequireGoverned: true, GovernedVerifier: verifier, + RequireGovernedReceipts: true, GovernedReceiptRecorder: failing, + }) + defer wait() + if err := WriteFrame(w, envelope); err != nil { + t.Fatal(err) + } + ack, err = ReadFrame(r) + if err != nil || !strings.Contains(string(ack.Payload), "record governed delivery receipt") { + t.Fatalf("receipt failure acknowledgement=%q err=%v", ack.Payload, err) + } + entries, err = os.ReadDir(failingDirectory) + if err != nil || len(entries) != 0 { + t.Fatalf("unreceipted delivery remained on disk: entries=%d err=%v", len(entries), err) + } +} + +func TestGovernedStreamRequiresReceiptBeforeComplete(t *testing.T) { + payload := []byte("a governed resumable file") + stream, verifier := newGovernedStreamForTest(t, "report.bin", payload, decision.Allow, nil) + envelope, err := EncodeGovernedStreamInit(stream) + if err != nil { + t.Fatal(err) + } + id, err := stream.TransferID() + if err != nil { + t.Fatal(err) + } + recorder := &governedReceiptRecorder{} + directory := t.TempDir() + w, r, wait := makeServiceConn(t, ServiceConfig{ + ReceivedDir: directory, RequireGoverned: true, GovernedVerifier: verifier, GovernedStreamVerifier: verifier, + RequireGovernedReceipts: true, GovernedReceiptRecorder: recorder, + }) + defer wait() + if err := WriteFrame(w, envelope); err != nil { + t.Fatal(err) + } + initAck, err := ReadFrame(r) + if err != nil { + t.Fatal(err) + } + kind, gotID, _, ok := decodeStreamFrame(initAck) + if !ok || kind != streamKindInitAck || gotID != id { + t.Fatalf("governed stream init response=%#v", initAck) + } + if err := WriteFrame(w, encodeChunk(id, 0, payload)); err != nil { + t.Fatal(err) + } + if _, err := ReadFrame(r); err != nil { + t.Fatal(err) + } + if err := WriteFrame(w, encodeStreamFrame(streamKindDone, id, nil)); err != nil { + t.Fatal(err) + } + complete, err := ReadFrame(r) + if err != nil { + t.Fatal(err) + } + kind, _, body, ok := decodeStreamFrame(complete) + completeOK, message := decodeComplete(body) + if !ok || kind != streamKindComplete || !completeOK || message != "" { + t.Fatalf("governed stream completion=%#v ok=%v message=%q", complete, completeOK, message) + } + if recorder.calls != 1 || recorder.intent.ID != stream.Intent.ID || recorder.decision.ID != stream.Decision.ID { + t.Fatalf("stream receipt recorder=%+v", recorder) + } + entries, err := os.ReadDir(directory) + if err != nil { + t.Fatal(err) + } + var saved string + for _, entry := range entries { + if !entry.IsDir() { + saved = entry.Name() + } + } + contents, err := os.ReadFile(filepath.Join(directory, saved)) + if err != nil || string(contents) != string(payload) { + t.Fatalf("saved governed stream contents=%q err=%v", contents, err) + } + + failingDirectory := t.TempDir() + w, r, wait = makeServiceConn(t, ServiceConfig{ + ReceivedDir: failingDirectory, RequireGoverned: true, GovernedVerifier: verifier, GovernedStreamVerifier: verifier, + RequireGovernedReceipts: true, GovernedReceiptRecorder: &governedReceiptRecorder{err: os.ErrPermission}, + }) + defer wait() + if err := WriteFrame(w, envelope); err != nil { + t.Fatal(err) + } + if _, err := ReadFrame(r); err != nil { + t.Fatal(err) + } + if err := WriteFrame(w, encodeChunk(id, 0, payload)); err != nil { + t.Fatal(err) + } + if _, err := ReadFrame(r); err != nil { + t.Fatal(err) + } + if err := WriteFrame(w, encodeStreamFrame(streamKindDone, id, nil)); err != nil { + t.Fatal(err) + } + complete, err = ReadFrame(r) + if err != nil { + t.Fatal(err) + } + kind, _, body, ok = decodeStreamFrame(complete) + completeOK, message = decodeComplete(body) + if !ok || kind != streamKindComplete || completeOK || !strings.Contains(message, "record governed stream receipt") { + t.Fatalf("receipt-failed stream completion=%#v ok=%v message=%q", complete, completeOK, message) + } + entries, err = os.ReadDir(failingDirectory) + if err != nil { + t.Fatal(err) + } + for _, entry := range entries { + if !entry.IsDir() { + t.Fatalf("unreceipted stream file remained: %s", entry.Name()) + } + } + + rawDirectory := t.TempDir() + w, r, wait = makeServiceConn(t, ServiceConfig{ + ReceivedDir: rawDirectory, RequireGoverned: true, GovernedVerifier: verifier, GovernedStreamVerifier: verifier, + RequireGovernedReceipts: true, GovernedReceiptRecorder: &governedReceiptRecorder{}, + }) + defer wait() + if err := WriteFrame(w, stream.InitFrame()); err != nil { + t.Fatal(err) + } + rejected, err := ReadFrame(r) + if err != nil { + t.Fatal(err) + } + kind, _, body, ok = decodeStreamFrame(rejected) + completeOK, message = decodeComplete(body) + if !ok || kind != streamKindComplete || completeOK || !strings.Contains(message, "unsigned stream frame") { + t.Fatalf("raw stream rejection=%#v ok=%v message=%q", rejected, completeOK, message) + } +} + +func TestGovernedReceiptRequirementFailsStartupWithoutRecorder(t *testing.T) { + service := NewService(ServiceConfig{RequireGoverned: true, RequireGovernedReceipts: 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 = governedTestTrust{} +var _ decision.AuthorityCeiling = governedTestCeiling{} +var _ GovernedReceiptRecorder = (*governedReceiptRecorder)(nil) diff --git a/retention.go b/retention.go new file mode 100644 index 0000000..4eee828 --- /dev/null +++ b/retention.go @@ -0,0 +1,295 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +package dataexchange + +import ( + "crypto/sha256" + "encoding/hex" + "encoding/json" + "fmt" + "io" + "log/slog" + "os" + "path/filepath" + "strings" + "time" + + "github.com/pilot-protocol/common/decision" +) + +const ( + retentionEntryVersion uint16 = 1 + defaultRetentionSweep = time.Minute + minRetentionDuration = time.Second + maxRetentionDuration = 10 * 365 * 24 * time.Hour + maxRetentionPolicies = 32 +) + +// GovernedRetentionPolicy maps one V2 signed retention class to a local +// retention duration. The class is selected by signed policy and bound to the +// Intent; the duration is attachment-controlled, so a sender cannot extend it. +type GovernedRetentionPolicy struct { + Class string + RetainFor time.Duration +} + +type governedRetentionEntry struct { + Version uint16 `json:"version"` + Root string `json:"root"` + RelativePath string `json:"relative_path"` + RetentionClass string `json:"retention_class"` + DisclosureHash string `json:"disclosure_hash"` + ExpiresAt int64 `json:"expires_at"` +} + +type governedRetentionManager struct { + stateDir string + roots map[string]string + policies map[string]time.Duration + now func() time.Time +} + +type retentionTicket struct { + entryPath string +} + +func newGovernedRetentionManager(stateDir string, roots map[string]string, policies []GovernedRetentionPolicy, now func() time.Time) (*governedRetentionManager, error) { + if len(policies) == 0 || len(policies) > maxRetentionPolicies { + return nil, fmt.Errorf("dataexchange: retention needs 1-%d policies", maxRetentionPolicies) + } + if strings.TrimSpace(stateDir) == "" { + return nil, fmt.Errorf("dataexchange: retention state directory is required") + } + if len(roots) == 0 { + return nil, fmt.Errorf("dataexchange: retention roots are required") + } + cleanRoots := make(map[string]string, len(roots)) + for name, root := range roots { + if name == "" || strings.TrimSpace(root) == "" || !filepath.IsAbs(root) { + return nil, fmt.Errorf("dataexchange: invalid retention root") + } + cleanRoots[name] = filepath.Clean(root) + } + policyMap := make(map[string]time.Duration, len(policies)) + for _, policy := range policies { + if !validRetentionClass(policy.Class) || policy.RetainFor < minRetentionDuration || policy.RetainFor > maxRetentionDuration { + return nil, fmt.Errorf("dataexchange: invalid retention policy") + } + if _, exists := policyMap[policy.Class]; exists { + return nil, fmt.Errorf("dataexchange: duplicate retention class %q", policy.Class) + } + policyMap[policy.Class] = policy.RetainFor + } + if now == nil { + now = time.Now + } + if err := os.MkdirAll(stateDir, 0700); err != nil { + return nil, fmt.Errorf("dataexchange: create retention state directory: %w", err) + } + return &governedRetentionManager{stateDir: filepath.Clean(stateDir), roots: cleanRoots, policies: policyMap, now: now}, nil +} + +func (manager *governedRetentionManager) prepare(disclosure *decision.DisclosureBinding, path string) (retentionTicket, error) { + if manager == nil { + return retentionTicket{}, nil + } + if disclosure == nil || disclosure.Version != decision.DisclosureBindingRetentionVersion || disclosure.RetentionClass == "" { + return retentionTicket{}, fmt.Errorf("dataexchange: governed retention requires a V2 disclosure retention class") + } + retention, exists := manager.policies[disclosure.RetentionClass] + if !exists { + return retentionTicket{}, fmt.Errorf("dataexchange: disclosure retention class is not configured") + } + rootName, relativePath, err := manager.relativePath(path) + if err != nil { + return retentionTicket{}, err + } + disclosureHash, err := disclosure.Hash() + if err != nil { + return retentionTicket{}, err + } + entry := governedRetentionEntry{ + Version: retentionEntryVersion, Root: rootName, RelativePath: relativePath, + RetentionClass: disclosure.RetentionClass, DisclosureHash: disclosureHash, + ExpiresAt: manager.now().UTC().Add(retention).Unix(), + } + encoded, err := json.Marshal(entry) + if err != nil { + return retentionTicket{}, fmt.Errorf("dataexchange: encode retention entry: %w", err) + } + sum := sha256.Sum256([]byte(rootName + "\x00" + relativePath)) + entryPath := filepath.Join(manager.stateDir, hex.EncodeToString(sum[:])+".json") + if err := writeRetentionEntry(entryPath, encoded); err != nil { + return retentionTicket{}, err + } + return retentionTicket{entryPath: entryPath}, nil +} + +func (ticket retentionTicket) rollback() error { + if ticket.entryPath == "" { + return nil + } + if err := os.Remove(ticket.entryPath); err != nil && !os.IsNotExist(err) { + return fmt.Errorf("dataexchange: remove retention entry: %w", err) + } + return nil +} + +// Sweep deletes expired governed content and only then removes its durable +// journal entry. A malformed journal causes an error instead of silently +// disabling a retention obligation after restart. +func (manager *governedRetentionManager) Sweep() error { + if manager == nil { + return nil + } + entries, err := os.ReadDir(manager.stateDir) + if err != nil { + return fmt.Errorf("dataexchange: read retention state: %w", err) + } + now := manager.now().UTC().Unix() + for _, file := range entries { + if file.IsDir() || !strings.HasSuffix(file.Name(), ".json") { + continue + } + entryPath := filepath.Join(manager.stateDir, file.Name()) + entry, err := readRetentionEntry(entryPath) + if err != nil { + return err + } + if err := manager.validateEntry(entry); err != nil { + return err + } + if entry.ExpiresAt > now { + continue + } + path, err := manager.pathFor(entry.Root, entry.RelativePath) + if err != nil { + return err + } + if err := os.Remove(path); err != nil && !os.IsNotExist(err) { + return fmt.Errorf("dataexchange: expire retained content: %w", err) + } + if err := os.Remove(entryPath); err != nil && !os.IsNotExist(err) { + return fmt.Errorf("dataexchange: remove expired retention entry: %w", err) + } + } + return nil +} + +func (manager *governedRetentionManager) run(stop <-chan struct{}, interval time.Duration) { + if interval <= 0 { + interval = defaultRetentionSweep + } + ticker := time.NewTicker(interval) + defer ticker.Stop() + for { + select { + case <-stop: + return + case <-ticker.C: + if err := manager.Sweep(); err != nil { + // Keep retrying. Startup already verifies the journal, and a + // transient removal error must not kill future retention work. + slog.Error("governed retention sweep failed", "error", err) + continue + } + } + } +} + +func (manager *governedRetentionManager) relativePath(path string) (string, string, error) { + cleanPath := filepath.Clean(path) + for name, root := range manager.roots { + relativePath, err := filepath.Rel(root, cleanPath) + if err == nil && relativePath != "." && relativePath != ".." && !strings.HasPrefix(relativePath, ".."+string(filepath.Separator)) && !filepath.IsAbs(relativePath) { + return name, relativePath, nil + } + } + return "", "", fmt.Errorf("dataexchange: retention path is outside managed roots") +} + +func (manager *governedRetentionManager) pathFor(rootName, relativePath string) (string, error) { + root, exists := manager.roots[rootName] + if !exists || relativePath == "" || filepath.IsAbs(relativePath) { + return "", fmt.Errorf("dataexchange: invalid retention entry path") + } + path := filepath.Join(root, relativePath) + _, checkedRelative, err := manager.relativePath(path) + if err != nil || checkedRelative != filepath.Clean(relativePath) { + return "", fmt.Errorf("dataexchange: invalid retention entry path") + } + return path, nil +} + +func (manager *governedRetentionManager) validateEntry(entry governedRetentionEntry) error { + if entry.Version != retentionEntryVersion || !validRetentionClass(entry.RetentionClass) || len(entry.DisclosureHash) != 64 || entry.ExpiresAt <= 0 { + return fmt.Errorf("dataexchange: invalid retention entry") + } + for _, character := range entry.DisclosureHash { + if !(character >= '0' && character <= '9' || character >= 'a' && character <= 'f') { + return fmt.Errorf("dataexchange: invalid retention entry") + } + } + _, err := manager.pathFor(entry.Root, entry.RelativePath) + return err +} + +func writeRetentionEntry(path string, data []byte) error { + temporary, err := os.CreateTemp(filepath.Dir(path), ".retention-*") + if err != nil { + return fmt.Errorf("dataexchange: create retention entry: %w", err) + } + temporaryPath := temporary.Name() + defer os.Remove(temporaryPath) + if err := temporary.Chmod(0600); err != nil { + _ = temporary.Close() + return fmt.Errorf("dataexchange: protect retention entry: %w", err) + } + if _, err := temporary.Write(data); err != nil { + _ = temporary.Close() + return fmt.Errorf("dataexchange: write retention entry: %w", err) + } + if err := temporary.Sync(); err != nil { + _ = temporary.Close() + return fmt.Errorf("dataexchange: sync retention entry: %w", err) + } + if err := temporary.Close(); err != nil { + return fmt.Errorf("dataexchange: close retention entry: %w", err) + } + if err := os.Rename(temporaryPath, path); err != nil { + return fmt.Errorf("dataexchange: publish retention entry: %w", err) + } + return nil +} + +func readRetentionEntry(path string) (governedRetentionEntry, error) { + file, err := os.Open(path) + if err != nil { + return governedRetentionEntry{}, fmt.Errorf("dataexchange: read retention entry: %w", err) + } + defer file.Close() + decoder := json.NewDecoder(io.LimitReader(file, 4<<10)) + decoder.DisallowUnknownFields() + var entry governedRetentionEntry + if err := decoder.Decode(&entry); err != nil { + return governedRetentionEntry{}, fmt.Errorf("dataexchange: decode retention entry: %w", err) + } + var trailing any + if err := decoder.Decode(&trailing); err != io.EOF { + return governedRetentionEntry{}, fmt.Errorf("dataexchange: invalid retention entry") + } + return entry, nil +} + +func validRetentionClass(value string) bool { + if len(value) == 0 || len(value) > 64 || value[0] == '-' || value[len(value)-1] == '-' { + return false + } + for index, character := range value { + if (character >= 'a' && character <= 'z') || (character >= '0' && character <= '9') || (character == '-' && index > 0 && index+1 < len(value)) { + continue + } + return false + } + return true +} diff --git a/retention_test.go b/retention_test.go new file mode 100644 index 0000000..e04f9d3 --- /dev/null +++ b/retention_test.go @@ -0,0 +1,171 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +//go:build !no_dataexchange +// +build !no_dataexchange + +package dataexchange + +import ( + "crypto/sha256" + "os" + "path/filepath" + "testing" + "time" + + "github.com/pilot-protocol/common/decision" + "github.com/pilot-protocol/common/protocol" +) + +func TestGovernedRetentionManagerDurablyExpiresOnlyManagedContent(t *testing.T) { + root := t.TempDir() + inbox := filepath.Join(root, "inbox") + received := filepath.Join(root, "received") + state := filepath.Join(root, "retention") + if err := os.MkdirAll(inbox, 0700); err != nil { + t.Fatal(err) + } + if err := os.MkdirAll(received, 0700); err != nil { + t.Fatal(err) + } + path := filepath.Join(inbox, "message.json") + if err := os.WriteFile(path, []byte("classified"), 0600); err != nil { + t.Fatal(err) + } + now := time.Unix(1785500000, 0).UTC() + manager, err := newGovernedRetentionManager(state, map[string]string{"inbox": inbox, "received": received}, []GovernedRetentionPolicy{{Class: "finance-7y", RetainFor: time.Second}}, func() time.Time { return now }) + if err != nil { + t.Fatal(err) + } + disclosure := decision.DisclosureBinding{ + Version: decision.DisclosureBindingRetentionVersion, ContentHash: decision.HashPayload([]byte("classified")), DeclaredBytes: 10, + ContentType: "application/json", Labels: []string{"finance"}, Recipient: "agent:finance", Purpose: "inbox-delivery", Residency: "eu-west-1", RetentionClass: "finance-7y", + } + ticket, err := manager.prepare(&disclosure, path) + if err != nil || ticket.entryPath == "" { + t.Fatalf("prepare ticket=%+v err=%v", ticket, err) + } + if err := manager.Sweep(); err != nil { + t.Fatal(err) + } + if _, err := os.Stat(path); err != nil { + t.Fatalf("content deleted before expiry: %v", err) + } + now = now.Add(time.Second) + if err := manager.Sweep(); err != nil { + t.Fatal(err) + } + if _, err := os.Stat(path); !os.IsNotExist(err) { + t.Fatalf("expired content remained: %v", err) + } + if _, err := os.Stat(ticket.entryPath); !os.IsNotExist(err) { + t.Fatalf("expired journal remained: %v", err) + } +} + +func TestGovernedRetentionManagerRejectsUnknownClassAndEscape(t *testing.T) { + root := t.TempDir() + inbox := filepath.Join(root, "inbox") + if err := os.MkdirAll(inbox, 0700); err != nil { + t.Fatal(err) + } + manager, err := newGovernedRetentionManager(filepath.Join(root, "retention"), map[string]string{"inbox": inbox}, []GovernedRetentionPolicy{{Class: "finance-7y", RetainFor: time.Hour}}, time.Now) + if err != nil { + t.Fatal(err) + } + disclosure := &decision.DisclosureBinding{Version: decision.DisclosureBindingRetentionVersion, RetentionClass: "finance-30d"} + if _, err := manager.prepare(disclosure, filepath.Join(inbox, "message.json")); err == nil { + t.Fatal("unknown retention class was accepted") + } + disclosure.RetentionClass = "finance-7y" + if _, err := manager.prepare(disclosure, filepath.Join(root, "outside.json")); err == nil { + t.Fatal("retention path escape was accepted") + } +} + +func TestServicePersistsAndExpiresGovernedRetention(t *testing.T) { + root := t.TempDir() + inbox := filepath.Join(root, "inbox") + received := filepath.Join(root, "received") + state := filepath.Join(root, "retention") + now := time.Unix(1785500000, 0).UTC() + service := NewService(ServiceConfig{ + InboxDir: inbox, ReceivedDir: received, RequireGoverned: true, + GovernedRetentionPolicies: []GovernedRetentionPolicy{{Class: "finance-7y", RetainFor: time.Second}}, + RetentionStateDir: state, + }) + if err := service.initializeGovernedRetention(); err != nil { + t.Fatal(err) + } + service.retention.now = func() time.Time { return now } + disclosure := &decision.DisclosureBinding{ + Version: decision.DisclosureBindingRetentionVersion, ContentHash: decision.HashPayload([]byte("classified")), DeclaredBytes: 10, + ContentType: "text/plain", Labels: []string{"finance"}, Recipient: "agent:finance", Purpose: "inbox-delivery", Residency: "eu-west-1", RetentionClass: "finance-7y", + } + if err := service.requireGovernedRetention(disclosure); err != nil { + t.Fatal(err) + } + delivery, err := service.prepareInboxMessage(&Frame{Type: TypeText, Payload: []byte("classified")}, protocol.Addr{Node: 1}, disclosure) + if err != nil { + t.Fatal(err) + } + delivery.commit() + files, err := os.ReadDir(inbox) + if err != nil || len(files) != 1 { + t.Fatalf("retained inbox files=%v err=%v", files, err) + } + now = now.Add(time.Second) + if err := service.retention.Sweep(); err != nil { + t.Fatal(err) + } + files, err = os.ReadDir(inbox) + if err != nil || len(files) != 0 { + t.Fatalf("expired inbox files=%v err=%v", files, err) + } +} + +func TestStreamRetentionPreparationRunsBeforeFinalRename(t *testing.T) { + dir := t.TempDir() + payload := []byte("classified") + hash := sha256.Sum256(payload) + var id [transferIDLen]byte + copy(id[:], hash[:transferIDLen]) + prepared, committed := false, false + receiver := NewStreamReceiverWithQuotaAndPrepareAndCommit( + dir, + func(string) string { return "final.txt" }, + nil, + func(_ [transferIDLen]byte, _ string, path string, _ int64) error { + if _, err := os.Stat(path); !os.IsNotExist(err) { + t.Fatalf("final path was visible before preparation: %v", err) + } + prepared = true + return nil + }, + func(_ [transferIDLen]byte, _ string, path string, _ int64) error { + if !prepared { + t.Fatal("commit ran before preparation") + } + if _, err := os.Stat(path); err != nil { + t.Fatalf("final path was not visible at commit: %v", err) + } + committed = true + return nil + }, + 0, + ) + if response := receiver.HandleFrame(encodeInit(id, uint64(len(payload)), hash, uint32(StreamChunkSize), "source.txt")); response == nil { + t.Fatal("INIT did not produce a response") + } + if response := receiver.HandleFrame(encodeChunk(id, 0, payload)); response == nil { + t.Fatal("CHUNK did not produce a response") + } + response := receiver.HandleFrame(encodeStreamFrame(streamKindDone, id, nil)) + if response == nil { + t.Fatal("DONE did not produce completion") + } + _, _, body, ok := decodeStreamFrame(response) + complete, _ := decodeComplete(body) + if !ok || !complete || !prepared || !committed { + t.Fatalf("stream lifecycle response=%+v prepared=%v committed=%v", response, prepared, committed) + } +} diff --git a/service.go b/service.go index 8f21320..0b42051 100644 --- a/service.go +++ b/service.go @@ -6,6 +6,7 @@ package dataexchange import ( + "bytes" "context" "encoding/base64" "encoding/json" @@ -19,6 +20,7 @@ import ( "unicode/utf8" "github.com/pilot-protocol/common/coreapi" + "github.com/pilot-protocol/common/decision" "github.com/pilot-protocol/common/protocol" ) @@ -55,6 +57,31 @@ type ServiceConfig struct { // instead of pinning the goroutine and its buffers indefinitely. // Zero ⇒ DefaultIdleTimeout; a negative value disables the deadline. IdleTimeout time.Duration + // RequireGoverned rejects legacy frames unless they carry a signed + // governed envelope verified by GovernedVerifier. Enable this only after + // sender rollout; the default preserves compatibility with older peers. + RequireGoverned bool + GovernedVerifier GovernedFrameVerifier + GovernedStreamVerifier GovernedStreamVerifier + RequireGovernedReceipts bool + GovernedReceiptRecorder GovernedReceiptRecorder + // GovernedContentInspector runs only at this receiver, after signed + // governed verification and before a message/file is released. It receives + // a reader, not an exported payload copy, so central authority services do + // not need application plaintext for DLP. + GovernedContentInspector decision.DisclosureContentInspector + RequireGovernedContentInspection bool + // GovernedTransferQuota bounds admitted transfers for each verified + // Intent.AgentID. It never uses an untrusted peer address as the subject. + // Quota is charged at governed-frame verification (or stream INIT) and does + // not silently apply to legacy traffic. + GovernedTransferQuota *decision.TransferQuotaLimiter + // GovernedRetentionPolicies maps signed V2 disclosure retention classes to + // durable local expiry work. When non-empty, governed deliveries without a + // configured V2 retention class are rejected before persistence. + GovernedRetentionPolicies []GovernedRetentionPolicy + RetentionStateDir string + RetentionSweepInterval time.Duration } // Defensible defaults applied when the corresponding ServiceConfig field @@ -123,12 +150,22 @@ const inboxEvictCheckEvery = 64 // Service is the L11 plugin adapter. Daemon (L7) holds it only as // coreapi.Service; cmd/daemon/main.go (L12) constructs it. type Service struct { - cfg ServiceConfig - listener coreapi.Listener - deps coreapi.Deps - cancel context.CancelFunc - done chan struct{} - seq atomic.Uint64 + cfg ServiceConfig + listener coreapi.Listener + deps coreapi.Deps + cancel context.CancelFunc + done chan struct{} + seq atomic.Uint64 + retention *governedRetentionManager +} + +// persistedDelivery separates a disk write from its visible notification so a +// governed receipt can be durably appended before the receiver acknowledges +// or publishes the delivery event. A failed receipt causes the staged file to +// be removed rather than leaving an unreceipted enterprise side effect. +type persistedDelivery struct { + rollback func() error + commit func() } func NewService(cfg ServiceConfig) *Service { @@ -141,7 +178,13 @@ func (s *Service) Name() string { return "dataexchange" } func (s *Service) Order() int { return 110 } func (s *Service) Start(ctx context.Context, deps coreapi.Deps) error { + if err := s.validateGovernedConfig(); err != nil { + return err + } s.deps = deps + if err := s.initializeGovernedRetention(); err != nil { + return err + } ln, err := deps.Streams.Listen(protocol.PortDataExchange) if err != nil { return fmt.Errorf("dataexchange: listen on port %d: %w", protocol.PortDataExchange, err) @@ -152,10 +195,77 @@ func (s *Service) Start(ctx context.Context, deps coreapi.Deps) error { s.cancel = cancel s.done = make(chan struct{}) go s.acceptLoop(runCtx) + if s.retention != nil { + go s.retention.run(runCtx.Done(), s.cfg.RetentionSweepInterval) + } slog.Info("dataexchange service listening", "port", protocol.PortDataExchange) return nil } +func (s *Service) validateGovernedConfig() error { + if s.cfg.RequireGovernedReceipts && (!s.cfg.RequireGoverned || s.cfg.GovernedReceiptRecorder == nil) { + return fmt.Errorf("dataexchange: governed receipts require a governed receiver and receipt recorder") + } + if s.cfg.RequireGovernedContentInspection && (!s.cfg.RequireGoverned || s.cfg.GovernedContentInspector == nil) { + return fmt.Errorf("dataexchange: required content inspection needs a governed receiver and local inspector") + } + if s.cfg.GovernedTransferQuota != nil && !s.cfg.RequireGoverned { + return fmt.Errorf("dataexchange: governed transfer quota requires a governed receiver") + } + if len(s.cfg.GovernedRetentionPolicies) > 0 && !s.cfg.RequireGoverned { + return fmt.Errorf("dataexchange: governed retention requires a governed receiver") + } + return nil +} + +func (s *Service) initializeGovernedRetention() error { + if len(s.cfg.GovernedRetentionPolicies) == 0 { + s.retention = nil + return nil + } + inbox, err := s.inboxDir() + if err != nil { + return fmt.Errorf("dataexchange: retention inbox directory: %w", err) + } + received, err := s.receivedDir() + if err != nil { + return fmt.Errorf("dataexchange: retention received directory: %w", err) + } + stateDir := s.cfg.RetentionStateDir + if stateDir == "" { + stateDir = filepath.Join(filepath.Dir(inbox), "retention") + } + manager, err := newGovernedRetentionManager(stateDir, map[string]string{"inbox": inbox, "received": received}, s.cfg.GovernedRetentionPolicies, nil) + if err != nil { + return err + } + if err := manager.Sweep(); err != nil { + return err + } + s.retention = manager + return nil +} + +func (s *Service) requireGovernedRetention(disclosure *decision.DisclosureBinding) error { + if s.retention == nil { + return nil + } + if disclosure == nil || disclosure.Version != decision.DisclosureBindingRetentionVersion || disclosure.RetentionClass == "" { + return fmt.Errorf("dataexchange: governed retention requires a V2 disclosure retention class") + } + if _, exists := s.retention.policies[disclosure.RetentionClass]; !exists { + return fmt.Errorf("dataexchange: disclosure retention class is not configured") + } + return nil +} + +func (s *Service) prepareGovernedRetention(disclosure *decision.DisclosureBinding, path string) (retentionTicket, error) { + if s.retention == nil { + return retentionTicket{}, nil + } + return s.retention.prepare(disclosure, path) +} + func (s *Service) Stop(ctx context.Context) error { if s.cancel != nil { s.cancel() @@ -166,6 +276,12 @@ func (s *Service) Stop(ctx context.Context) error { if s.done == nil { return nil } + // Prefer a context error that already existed when Stop was called. A + // completed accept loop must not make cancellation semantics depend on + // which ready channel select happens to choose. + if err := ctx.Err(); err != nil { + return err + } select { case <-s.done: case <-ctx.Done(): @@ -202,7 +318,12 @@ func (s *Service) handleConn(ctx context.Context, conn coreapi.Stream) { // Final filenames and the file.received event match saveReceivedFile so // the two transfer paths are indistinguishable to consumers. var sr *StreamReceiver + governedStreams := make(map[[transferIDLen]byte]GovernedStreamInit) + streamRetention := make(map[[transferIDLen]byte]retentionTicket) defer func() { + for _, ticket := range streamRetention { + _ = ticket.rollback() + } if sr != nil { sr.Close() } @@ -222,6 +343,54 @@ func (s *Service) handleConn(ctx context.Context, conn coreapi.Stream) { }) } } + streamOnPrepare := func(id [transferIDLen]byte, _ string, path string, _ int64) error { + governed, governedTransfer := governedStreams[id] + if s.cfg.RequireGoverned && !governedTransfer { + return fmt.Errorf("stream completion is not bound to a governed INIT") + } + if !governedTransfer { + return nil + } + ticket, err := s.prepareGovernedRetention(governed.Disclosure, path) + if err != nil { + return err + } + streamRetention[id] = ticket + return nil + } + streamOnCommit := func(id [transferIDLen]byte, name, path string, _ int64) error { + governed, governedTransfer := governedStreams[id] + if s.cfg.RequireGoverned && !governedTransfer { + return fmt.Errorf("stream completion is not bound to a governed INIT") + } + if !governedTransfer { + return nil + } + defer delete(governedStreams, id) + retention := streamRetention[id] + defer delete(streamRetention, id) + committed := false + defer func() { + if !committed { + _ = retention.rollback() + } + }() + if err := s.inspectGovernedStreamFile(ctx, governed, path, name); err != nil { + return err + } + if s.cfg.GovernedReceiptRecorder == nil { + if s.cfg.RequireGovernedReceipts { + return fmt.Errorf("governed receipt recorder is not configured") + } + committed = true + return nil + } + if err := recordGovernedReceipt(ctx, s.cfg.GovernedReceiptRecorder, governed.Intent, governed.Decision, governed.Disclosure); err != nil { + return fmt.Errorf("record governed stream receipt: %w", err) + } + committed = true + return nil + } idle := s.effectiveIdleTimeout() dl, canDeadline := conn.(readDeadliner) @@ -244,63 +413,177 @@ func (s *Service) handleConn(ctx context.Context, conn coreapi.Stream) { "bytes", len(frame.Payload), "remote", conn.RemoteAddr()) - var saveErr error - var ackFrame *Frame - switch frame.Type { - case TypeFileStream: - // Chunked/resumable transfer. The receiver emits its own - // control responses (INIT-ACK / ACK / COMPLETE), so skip the - // generic per-frame ACK below. - if sr == nil { - dir, derr := s.receivedDir() - if derr != nil { - _ = WriteFrame(conn, &Frame{Type: TypeText, Payload: []byte("ERR received dir: " + derr.Error())}) - return - } - if mderr := os.MkdirAll(dir, 0700); mderr != nil { - _ = WriteFrame(conn, &Frame{Type: TypeText, Payload: []byte("ERR mkdir: " + mderr.Error())}) - return - } - sr = NewStreamReceiverWithQuota(dir, streamNameSuffix, streamOnSaved, s.effectiveReceivedMaxBytes()) + var ( + saveErr error + ackFrame *Frame + governed *GovernedFrame + streamed *GovernedStreamInit + delivery persistedDelivery + ) + if frame.Type == TypeGoverned { + decoded, governErr := DecodeGovernedFrame(frame) + if governErr != nil { + saveErr = governErr + } else if s.cfg.GovernedVerifier == nil { + saveErr = fmt.Errorf("governed frame received but no verifier is configured") + } else if governErr = s.cfg.GovernedVerifier.VerifyGovernedFrame(ctx, conn.RemoteAddr(), decoded); governErr != nil { + saveErr = governErr + } else { + governed = &decoded + frame = governed.DataFrame() } - if resp := sr.HandleFrame(frame); resp != nil { - if werr := WriteFrame(conn, resp); werr != nil { - return + } else if frame.Type == TypeGovernedFileStream { + decoded, governErr := DecodeGovernedStreamInit(frame) + if governErr != nil { + saveErr = governErr + } else if s.cfg.GovernedStreamVerifier == nil { + saveErr = fmt.Errorf("governed stream received but no verifier is configured") + } else if governErr = s.cfg.GovernedStreamVerifier.VerifyGovernedStreamInit(ctx, conn.RemoteAddr(), decoded); governErr != nil { + saveErr = governErr + } else { + streamed = &decoded + frame = streamed.InitFrame() + } + } else if s.cfg.RequireGoverned && frame.Type != TypeFileStream { + saveErr = fmt.Errorf("unsigned legacy frame rejected by governed receiver") + } + if saveErr == nil { + if governed != nil { + saveErr = s.admitGovernedTransfer(governed.Intent, uint64(len(governed.Payload))) + } else if streamed != nil { + declaredBytes, declaredErr := streamed.DeclaredBytes() + if declaredErr != nil { + saveErr = declaredErr + } else { + saveErr = s.admitGovernedTransfer(streamed.Intent, declaredBytes) } } - continue - case TypeFile: - if frame.Filename != "" { - saveErr = s.saveReceivedFile(frame) + } + if saveErr == nil { + if governed != nil { + saveErr = s.requireGovernedRetention(governed.Disclosure) + } else if streamed != nil { + saveErr = s.requireGovernedRetention(streamed.Disclosure) } - case TypeText, TypeJSON, TypeBinary: - saveErr = s.saveInboxMessage(frame, conn.RemoteAddr()) - case TypeTrace: - tf, tferr := ReadTracePayload(frame) - if tferr != nil { - ackFrame = &Frame{ - Type: TypeText, - Payload: []byte(fmt.Sprintf("ERR trace parse: %v", tferr)), + } + if saveErr == nil { + if governed != nil && frame.Type != TypeFileStream { + saveErr = s.inspectGovernedFrame(ctx, *governed) + } + } + if saveErr == nil { + switch frame.Type { + case TypeFileStream: + // Chunked/resumable transfer. The receiver emits its own + // control responses (INIT-ACK / ACK / COMPLETE), so skip the + // generic per-frame ACK below. + if sr == nil { + dir, derr := s.receivedDir() + if derr != nil { + _ = WriteFrame(conn, &Frame{Type: TypeText, Payload: []byte("ERR received dir: " + derr.Error())}) + return + } + if mderr := os.MkdirAll(dir, 0700); mderr != nil { + _ = WriteFrame(conn, &Frame{Type: TypeText, Payload: []byte("ERR mkdir: " + mderr.Error())}) + return + } + sr = NewStreamReceiverWithQuotaAndPrepareAndCommit(dir, streamNameSuffix, streamOnSaved, streamOnPrepare, streamOnCommit, s.effectiveReceivedMaxBytes()) } - } else { - innerFrame := &Frame{Type: tf.InnerType, Payload: tf.Payload} - innerSaveErr := s.saveInboxMessage(innerFrame, conn.RemoteAddr()) - inboxWrittenAtNs := time.Now().UnixNano() - innerAck := fmt.Sprintf("ACK %s %d bytes", TypeName(tf.InnerType), len(tf.Payload)) - if innerSaveErr != nil { - innerAck = fmt.Sprintf("ERR %s save failed: %v", TypeName(tf.InnerType), innerSaveErr) + kind, id, _, validStreamFrame := decodeStreamFrame(frame) + if !validStreamFrame { + _ = WriteFrame(conn, encodeComplete(id, false, "malformed stream frame")) + continue } - ackSentAtNs := time.Now().UnixNano() - timingJSON, _ := json.Marshal(map[string]interface{}{ - "sent_at_ns": tf.SentAtNs, - "received_at_ns": frameReceivedAtNs, - "inbox_written_at_ns": inboxWrittenAtNs, - "ack_sent_at_ns": ackSentAtNs, - "inner_ack": innerAck, - }) - ackFrame = &Frame{Type: TypeJSON, Payload: timingJSON} + if s.cfg.RequireGoverned { + if streamed == nil { + if kind == streamKindInit || !streamBound(governedStreams, id) { + _ = WriteFrame(conn, encodeComplete(id, false, "unsigned stream frame rejected by governed receiver")) + continue + } + } else if kind != streamKindInit { + _ = WriteFrame(conn, encodeComplete(id, false, "governed stream envelope must carry INIT")) + continue + } + } + if resp := sr.HandleFrame(frame); resp != nil { + responseKind, _, responseBody, responseValid := decodeStreamFrame(resp) + if streamed != nil && responseValid && responseKind == streamKindInitAck { + governedStreams[id] = *streamed + } + if responseValid && responseKind == streamKindComplete { + completeOK, _ := decodeComplete(responseBody) + if completeOK { + delete(governedStreams, id) + } + } + if werr := WriteFrame(conn, resp); werr != nil { + return + } + } + if kind == streamKindAbort { + delete(governedStreams, id) + } + continue + case TypeFile: + if frame.Filename != "" { + if governed != nil { + delivery, saveErr = s.prepareReceivedFile(frame, governed.Disclosure) + } else { + saveErr = s.saveReceivedFile(frame) + } + } + case TypeText, TypeJSON, TypeBinary: + if governed != nil { + delivery, saveErr = s.prepareInboxMessage(frame, conn.RemoteAddr(), governed.Disclosure) + } else { + saveErr = s.saveInboxMessage(frame, conn.RemoteAddr()) + } + case TypeTrace: + tf, tferr := ReadTracePayload(frame) + if tferr != nil { + ackFrame = &Frame{ + Type: TypeText, + Payload: []byte(fmt.Sprintf("ERR trace parse: %v", tferr)), + } + } else { + innerFrame := &Frame{Type: tf.InnerType, Payload: tf.Payload} + innerSaveErr := s.saveInboxMessage(innerFrame, conn.RemoteAddr()) + inboxWrittenAtNs := time.Now().UnixNano() + innerAck := fmt.Sprintf("ACK %s %d bytes", TypeName(tf.InnerType), len(tf.Payload)) + if innerSaveErr != nil { + innerAck = fmt.Sprintf("ERR %s save failed: %v", TypeName(tf.InnerType), innerSaveErr) + } + ackSentAtNs := time.Now().UnixNano() + timingJSON, _ := json.Marshal(map[string]interface{}{ + "sent_at_ns": tf.SentAtNs, + "received_at_ns": frameReceivedAtNs, + "inbox_written_at_ns": inboxWrittenAtNs, + "ack_sent_at_ns": ackSentAtNs, + "inner_ack": innerAck, + }) + ackFrame = &Frame{Type: TypeJSON, Payload: timingJSON} + } + default: + saveErr = fmt.Errorf("unsupported frame type %d", frame.Type) + } + } + if saveErr == nil && governed != nil && delivery.commit != nil { + if s.cfg.GovernedReceiptRecorder == nil { + if s.cfg.RequireGovernedReceipts { + saveErr = fmt.Errorf("governed receipt recorder is not configured") + } + } else if receiptErr := recordGovernedReceipt(ctx, s.cfg.GovernedReceiptRecorder, governed.Intent, governed.Decision, governed.Disclosure); receiptErr != nil { + if delivery.rollback != nil { + if rollbackErr := delivery.rollback(); rollbackErr != nil { + receiptErr = fmt.Errorf("%w; remove unreceipted delivery: %v", receiptErr, rollbackErr) + } + } + saveErr = fmt.Errorf("record governed delivery receipt: %w", receiptErr) } } + if saveErr == nil && delivery.commit != nil { + delivery.commit() + } if ackFrame == nil { ackMsg := fmt.Sprintf("ACK %s %d bytes", TypeName(frame.Type), len(frame.Payload)) @@ -322,6 +605,68 @@ func (s *Service) handleConn(ctx context.Context, conn coreapi.Stream) { } } +func (s *Service) admitGovernedTransfer(intent decision.Intent, bytes uint64) error { + if s.cfg.GovernedTransferQuota == nil { + return nil + } + if err := s.cfg.GovernedTransferQuota.Allow(intent.AgentID, bytes); err != nil { + slog.Warn("governed transfer quota rejected", "agent_id", intent.AgentID, "bytes", bytes, "error", err) + return fmt.Errorf("governed transfer quota rejected") + } + return nil +} + +func streamBound(streams map[[transferIDLen]byte]GovernedStreamInit, id [transferIDLen]byte) bool { + _, exists := streams[id] + return exists +} + +func (s *Service) inspectGovernedFrame(ctx context.Context, governed GovernedFrame) error { + if s.cfg.GovernedContentInspector == nil { + return nil + } + contentType := contentTypeForFrame(governed.Type, governed.Disclosure) + if err := s.cfg.GovernedContentInspector.InspectDisclosureContent(ctx, governed.Intent, governed.Disclosure, contentType, governed.Filename, bytes.NewReader(governed.Payload)); err != nil { + slog.Warn("governed content inspection rejected", "action", governed.Intent.Action, "resource", governed.Intent.Resource, "error", err) + return fmt.Errorf("governed content inspection rejected") + } + return nil +} + +func (s *Service) inspectGovernedStreamFile(ctx context.Context, governed GovernedStreamInit, path string, filename string) error { + if s.cfg.GovernedContentInspector == nil { + return nil + } + file, err := os.Open(path) + if err != nil { + return fmt.Errorf("governed content inspection rejected") + } + defer file.Close() + contentType := contentTypeForFrame(TypeFile, governed.Disclosure) + // Stream the verified on-disk body directly. A scanner that cannot process + // the whole file must return an error; truncating the reader would create a + // bypass for sensitive content placed after an arbitrary byte boundary. + if err := s.cfg.GovernedContentInspector.InspectDisclosureContent(ctx, governed.Intent, governed.Disclosure, contentType, filename, file); err != nil { + slog.Warn("governed stream content inspection rejected", "action", governed.Intent.Action, "resource", governed.Intent.Resource, "error", err) + return fmt.Errorf("governed content inspection rejected") + } + return nil +} + +func contentTypeForFrame(frameType uint32, disclosure *decision.DisclosureBinding) string { + if disclosure != nil { + return disclosure.ContentType + } + switch frameType { + case TypeText: + return "text/plain" + case TypeJSON: + return "application/json" + default: + return "application/octet-stream" + } +} + // receivedDir returns the configured received-file directory or the // default ~/.pilot/received. func (s *Service) receivedDir() (string, error) { @@ -349,14 +694,23 @@ func (s *Service) inboxDir() (string, error) { } func (s *Service) saveReceivedFile(frame *Frame) error { + delivery, err := s.prepareReceivedFile(frame, nil) + if err != nil { + return err + } + delivery.commit() + return nil +} + +func (s *Service) prepareReceivedFile(frame *Frame, disclosure *decision.DisclosureBinding) (persistedDelivery, error) { dir, err := s.receivedDir() if err != nil { slog.Warn("save received file: cannot determine dir", "err", err) - return err + return persistedDelivery{}, err } if err := os.MkdirAll(dir, 0700); err != nil { slog.Warn("save received file: mkdir failed", "err", err) - return fmt.Errorf("mkdir: %w", err) + return persistedDelivery{}, fmt.Errorf("mkdir: %w", err) } // Disk quota: reject the file up front if storing it would push the @@ -375,7 +729,7 @@ func (s *Service) saveReceivedFile(frame *Frame) error { "max_bytes": quota, }) } - return fmt.Errorf("received-files quota exceeded: %d + %d > %d", + return persistedDelivery{}, fmt.Errorf("received-files quota exceeded: %d + %d > %d", current, len(frame.Payload), quota) } } @@ -387,28 +741,53 @@ func (s *Service) saveReceivedFile(frame *Frame) error { base := safeName[:len(safeName)-len(ext)] destName := fmt.Sprintf("%s-%s-%06d%s", base, ts, seq, ext) destPath := filepath.Join(dir, destName) + retention, err := s.prepareGovernedRetention(disclosure, destPath) + if err != nil { + return persistedDelivery{}, err + } if err := os.WriteFile(destPath, frame.Payload, 0600); err != nil { _ = os.Remove(destPath) + _ = retention.rollback() slog.Warn("save received file: write failed", "path", destPath, "err", err) - return fmt.Errorf("write: %w", err) - } - slog.Info("file saved", "path", destPath, "bytes", len(frame.Payload)) - if s.deps.Events != nil { - s.deps.Events.Publish("file.received", map[string]any{ - "filename": safeName, "size": len(frame.Payload), "path": destPath, - }) + return persistedDelivery{}, fmt.Errorf("write: %w", err) + } + return persistedDelivery{ + rollback: func() error { + removeErr := os.Remove(destPath) + retentionErr := retention.rollback() + if removeErr != nil && !os.IsNotExist(removeErr) { + return removeErr + } + return retentionErr + }, + commit: func() { + slog.Info("file saved", "path", destPath, "bytes", len(frame.Payload)) + if s.deps.Events != nil { + s.deps.Events.Publish("file.received", map[string]any{ + "filename": safeName, "size": len(frame.Payload), "path": destPath, + }) + } + }, + }, nil +} + +func (s *Service) saveInboxMessage(frame *Frame, from protocol.Addr) error { + delivery, err := s.prepareInboxMessage(frame, from, nil) + if err != nil { + return err } + delivery.commit() return nil } -func (s *Service) saveInboxMessage(frame *Frame, from protocol.Addr) error { +func (s *Service) prepareInboxMessage(frame *Frame, from protocol.Addr, disclosure *decision.DisclosureBinding) (persistedDelivery, error) { dir, err := s.inboxDir() if err != nil { - return err + return persistedDelivery{}, err } if err := os.MkdirAll(dir, 0700); err != nil { - return fmt.Errorf("mkdir: %w", err) + return persistedDelivery{}, fmt.Errorf("mkdir: %w", err) } // Byte-budget check: confirm there is room BEFORE writing. Evict if @@ -435,7 +814,7 @@ func (s *Service) saveInboxMessage(frame *Frame, from protocol.Addr) error { "max_bytes": maxBytes, }) } - return fmt.Errorf("inbox byte budget exceeded: %d + %d > %d", + return persistedDelivery{}, fmt.Errorf("inbox byte budget exceeded: %d + %d > %d", after, estimated, maxBytes) } } @@ -464,29 +843,45 @@ func (s *Service) saveInboxMessage(frame *Frame, from protocol.Addr) error { } data, err := json.Marshal(msg) if err != nil { - return fmt.Errorf("marshal: %w", err) + return persistedDelivery{}, fmt.Errorf("marshal: %w", err) } seq := s.seq.Add(1) filename := fmt.Sprintf("%s-%s-%06d.json", TypeName(frame.Type), ts.Format("20060102-150405.000"), seq) destPath := filepath.Join(dir, filename) - if err := os.WriteFile(destPath, data, 0600); err != nil { - return fmt.Errorf("write: %w", err) - } - slog.Info("inbox message saved", "path", destPath, "type", TypeName(frame.Type), "bytes", len(frame.Payload)) - if s.deps.Events != nil { - s.deps.Events.Publish("message.received", map[string]any{ - "type": TypeName(frame.Type), "from": from.String(), - "size": len(frame.Payload), - }) - } - // Periodic eviction so a misbehaving peer (or sustained inbound - // load) cannot fill the operator's disk. We sample every - // inboxEvictCheckEvery writes — the cap is soft. - if seq%inboxEvictCheckEvery == 0 { - s.evictInboxOverflow(dir) + retention, err := s.prepareGovernedRetention(disclosure, destPath) + if err != nil { + return persistedDelivery{}, err } - return nil + if err := os.WriteFile(destPath, data, 0600); err != nil { + _ = retention.rollback() + return persistedDelivery{}, fmt.Errorf("write: %w", err) + } + return persistedDelivery{ + rollback: func() error { + removeErr := os.Remove(destPath) + retentionErr := retention.rollback() + if removeErr != nil && !os.IsNotExist(removeErr) { + return removeErr + } + return retentionErr + }, + commit: func() { + slog.Info("inbox message saved", "path", destPath, "type", TypeName(frame.Type), "bytes", len(frame.Payload)) + if s.deps.Events != nil { + s.deps.Events.Publish("message.received", map[string]any{ + "type": TypeName(frame.Type), "from": from.String(), + "size": len(frame.Payload), + }) + } + // Periodic eviction so a misbehaving peer (or sustained inbound + // load) cannot fill the operator's disk. We sample every + // inboxEvictCheckEvery writes — the cap is soft. + if seq%inboxEvictCheckEvery == 0 { + s.evictInboxOverflow(dir) + } + }, + }, nil } // evictInboxOverflow trims the inbox to at most cfg.InboxMaxFiles by diff --git a/service_disabled.go b/service_disabled.go index e865588..066004c 100644 --- a/service_disabled.go +++ b/service_disabled.go @@ -15,6 +15,7 @@ import ( "time" "github.com/pilot-protocol/common/coreapi" + "github.com/pilot-protocol/common/decision" ) // ServiceConfig mirrors the real ServiceConfig so cmd/daemon's @@ -22,13 +23,24 @@ import ( // compiles unchanged when the plugin is disabled. Field set kept in sync // with the real ServiceConfig in service.go. type ServiceConfig struct { - ReceivedDir string - InboxDir string - IncludeBase64 bool - InboxMaxFiles int - InboxMaxBytes int64 - ReceivedMaxBytes int64 - IdleTimeout time.Duration + ReceivedDir string + InboxDir string + IncludeBase64 bool + InboxMaxFiles int + InboxMaxBytes int64 + ReceivedMaxBytes int64 + IdleTimeout time.Duration + RequireGoverned bool + GovernedVerifier GovernedFrameVerifier + GovernedStreamVerifier GovernedStreamVerifier + RequireGovernedReceipts bool + GovernedReceiptRecorder GovernedReceiptRecorder + GovernedContentInspector decision.DisclosureContentInspector + RequireGovernedContentInspection bool + GovernedTransferQuota *decision.TransferQuotaLimiter + GovernedRetentionPolicies []GovernedRetentionPolicy + RetentionStateDir string + RetentionSweepInterval time.Duration } // Service is a no-op replacement for the real plugin Service. diff --git a/zz_coverage_test.go b/zz_coverage_test.go index dd7e309..affa4f8 100644 --- a/zz_coverage_test.go +++ b/zz_coverage_test.go @@ -417,8 +417,9 @@ func TestHandleConn_SaveError(t *testing.T) { } } -// TestHandleConn_UnknownType — a frame whose type doesn't match any -// switch arm still produces an ACK (no save, but a default ACK is sent). +// TestHandleConn_UnknownType ensures an unsupported wire type is rejected. +// Silently ACKing it would allow a sender to mistake an unverified governed +// envelope (or a future control frame) for a successfully delivered message. func TestHandleConn_UnknownType(t *testing.T) { t.Parallel() tmp := t.TempDir() @@ -432,7 +433,7 @@ func TestHandleConn_UnknownType(t *testing.T) { if err != nil { t.Fatalf("read ack: %v", err) } - if !bytes.Contains(ack.Payload, []byte("ACK UNKNOWN(999)")) { + if !bytes.Contains(ack.Payload, []byte("ERR UNKNOWN(999) save failed: unsupported frame type 999")) { t.Errorf("ack payload = %q", ack.Payload) } } diff --git a/zz_more_test.go b/zz_more_test.go index d1883ba..e996d1c 100644 --- a/zz_more_test.go +++ b/zz_more_test.go @@ -1,5 +1,8 @@ // SPDX-License-Identifier: AGPL-3.0-or-later +//go:build !no_dataexchange +// +build !no_dataexchange + package dataexchange import ( From 8535fa6cd8919d6d0501903246f9837179d99ed9 Mon Sep 17 00:00:00 2001 From: Teodor Calin Date: Sun, 2 Aug 2026 21:01:44 +0300 Subject: [PATCH 2/3] sec(governed): receiver-side replay dedup for governed transfers/publications MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit M5: a verified governed envelope (frame, stream INIT, or event) was a bearer capability with no receiver-side replay protection — any peer that observed one could re-send the exact bytes within the 5-minute intent TTL, producing duplicate authorized deliveries/publications and re-charging the signing agent's quota (cross-agent grief) plus provenance confusion. Add a bounded per-receiver replay guard keyed on the signature-authenticated (tenant, agent, intent-id), TTL = the intent's own ExpiresAt, checked at the single admit/govern chokepoint before quota and before the side effect. A legitimate retry carries a fresh intent (fresh nonce/id), so this rejects only true replays — full governed stream/publish suites still pass. Regression tests added. SECURITY_REVIEW_v1.14 M5. Co-Authored-By: Claude Fable 5 --- governed_replay.go | 71 ++++++++++++++++++++++++++++++++++++++ service.go | 9 ++++- zz_governed_replay_test.go | 39 +++++++++++++++++++++ 3 files changed, 118 insertions(+), 1 deletion(-) create mode 100644 governed_replay.go create mode 100644 zz_governed_replay_test.go diff --git a/governed_replay.go b/governed_replay.go new file mode 100644 index 0000000..5a15ef9 --- /dev/null +++ b/governed_replay.go @@ -0,0 +1,71 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +package dataexchange + +import ( + "fmt" + "sync" + "time" + + "github.com/pilot-protocol/common/decision" +) + +// maxGovernedReplayEntries bounds the receiver-side replay cache. Entries are +// keyed on the signed intent and expire at the intent's own ExpiresAt (<= the +// 5-minute MaxIntentTTL), so under legitimate signed traffic the cache drains +// continuously. The cap is a backstop against pathological retention; on +// overflow a new transfer is refused (fail closed) rather than admitted +// un-deduplicated. +const maxGovernedReplayEntries = 1 << 20 + +// governedReplayGuard rejects a second delivery of the same signed governed +// intent. A verified governed envelope is otherwise a bearer capability: any +// peer that observed one could re-send the exact bytes within the intent TTL, +// producing duplicate authorized deliveries and re-charging the signing +// agent's quota. Dedup is keyed on (tenant, agent, intent id) — all +// signature-authenticated fields — so a replay cannot dodge it by presenting a +// different transport peer. A legitimate retry must carry a fresh intent +// (fresh nonce/id); reusing a signed intent IS the replay pattern. +type governedReplayGuard struct { + mu sync.Mutex + seen map[string]int64 // key -> intent ExpiresAt (unix seconds) + now func() time.Time +} + +func newGovernedReplayGuard() *governedReplayGuard { + return &governedReplayGuard{seen: make(map[string]int64), now: time.Now} +} + +func governedReplayKey(intent decision.Intent) string { + return intent.TenantID + "\x1f" + intent.AgentID + "\x1f" + intent.ID +} + +// admit records the intent as delivered and returns an error if it was already +// delivered (replay) or if the cache is saturated. A blank intent ID is a +// no-op (upstream freshness/nonce checks still apply). +func (g *governedReplayGuard) admit(intent decision.Intent) error { + if intent.ID == "" { + return nil + } + now := g.now().Unix() + key := governedReplayKey(intent) + 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 intent already delivered (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/service.go b/service.go index 0b42051..d772336 100644 --- a/service.go +++ b/service.go @@ -157,6 +157,7 @@ type Service struct { done chan struct{} seq atomic.Uint64 retention *governedRetentionManager + replay *governedReplayGuard } // persistedDelivery separates a disk write from its visible notification so a @@ -169,7 +170,7 @@ type persistedDelivery struct { } func NewService(cfg ServiceConfig) *Service { - return &Service{cfg: cfg} + return &Service{cfg: cfg, replay: newGovernedReplayGuard()} } func (s *Service) Name() string { return "dataexchange" } @@ -606,6 +607,12 @@ func (s *Service) handleConn(ctx context.Context, conn coreapi.Stream) { } func (s *Service) admitGovernedTransfer(intent decision.Intent, bytes uint64) error { + if s.replay != nil { + if err := s.replay.admit(intent); err != nil { + slog.Warn("governed transfer replay rejected", "agent_id", intent.AgentID, "intent_id", intent.ID, "error", err) + return err + } + } if s.cfg.GovernedTransferQuota == nil { return nil } diff --git a/zz_governed_replay_test.go b/zz_governed_replay_test.go new file mode 100644 index 0000000..7b4890c --- /dev/null +++ b/zz_governed_replay_test.go @@ -0,0 +1,39 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +package dataexchange + +import ( + "testing" + "time" + + "github.com/pilot-protocol/common/decision" +) + +// TestGovernedReplayGuard pins SECURITY_REVIEW_v1.14 finding M5: a verified +// governed intent may be delivered at most once; a replay within the intent's +// validity window is rejected, and the cache drains after expiry. +func TestGovernedReplayGuard(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: "int-1", TenantID: "t", AgentID: "a", ExpiresAt: base.Add(5 * time.Minute).Unix()} + + if err := g.admit(intent); err != nil { + t.Fatalf("first delivery rejected: %v", err) + } + if err := g.admit(intent); err == nil { + t.Fatal("replay within TTL was ACCEPTED") + } + // A distinct intent from the same agent is fine (fresh nonce/id). + if err := g.admit(decision.Intent{ID: "int-2", TenantID: "t", AgentID: "a", ExpiresAt: clock.Add(5 * time.Minute).Unix()}); err != nil { + t.Fatalf("distinct intent rejected: %v", err) + } + // After the intent expires the entry is pruned; the same id no longer + // collides (and upstream freshness would reject the stale intent anyway). + clock = base.Add(6 * time.Minute) + if err := g.admit(decision.Intent{ID: "int-1", TenantID: "t", AgentID: "a", ExpiresAt: clock.Add(5 * time.Minute).Unix()}); err != nil { + t.Fatalf("post-expiry re-use rejected: %v", err) + } +} From 0bac7eddbf55227697160c8f1a4e0285bba8b13a Mon Sep 17 00:00:00 2001 From: Teodor Calin Date: Fri, 7 Aug 2026 02:59:56 +0300 Subject: [PATCH 3/3] build(deps): require public control contracts --- go.mod | 2 +- go.sum | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/go.mod b/go.mod index 672cc90..e151458 100644 --- a/go.mod +++ b/go.mod @@ -3,6 +3,6 @@ module github.com/pilot-protocol/dataexchange go 1.25.12 require ( - github.com/pilot-protocol/common v0.5.11 + github.com/pilot-protocol/common v0.5.12 github.com/pilot-protocol/eventstream v0.2.2 ) diff --git a/go.sum b/go.sum index a6b4ee4..312c552 100644 --- a/go.sum +++ b/go.sum @@ -1,4 +1,4 @@ -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= github.com/pilot-protocol/eventstream v0.2.2 h1:E0IjveK7K+dsIbE/5hD3N821FkHzxVsx1tiAORMzt8k= github.com/pilot-protocol/eventstream v0.2.2/go.mod h1:gUjoMEItW1SRJYEq39VlcIeDe2LcE5B18/4bcaUJNrs=