diff --git a/README.md b/README.md index b14e4e8..e485f2d 100644 --- a/README.md +++ b/README.md @@ -225,10 +225,11 @@ kubeconfig, endpoint, or token through the Sith collector contract. Only normali freshness-bounded, and stored behind forced RLS. Failed refreshes retain the last snapshot as explicitly stale evidence and record only a closed failure category. Concurrent refresh requests are authorized independently and then coalesced only -within the same validated workspace. The shared refresh runs on a detached, locally traced context, -so one canceled caller cannot cancel its peers and no request credential, context value, or trace -identity crosses caller boundaries; completed and panicking refresh flights are removed. Within one -refresh, transport and validation use a bounded worker pool: `CollectorConfig` defaults to four +within the same validated workspace. The shared refresh runs on a locally traced context detached +from individual callers but bounded by the hub lifecycle, so one canceled caller cannot cancel its +peers, shutdown cancels outstanding work, and no request credential, context value, or trace identity +crosses caller boundaries; completed and panicking refresh flights are removed. Within one refresh, +transport and validation use a bounded worker pool: `CollectorConfig` defaults to four concurrent spokes and accepts only explicit limits from 1 through 64. Persistence and coverage mutation remain serialized, so one store failure cancels the remaining workers before returning; per-spoke deadlines and sorted stale/unreachable coverage remain unchanged. The pinned diff --git a/internal/hubdb/policy_audit_test.go b/internal/hubdb/policy_audit_test.go index e3e877a..eed9fa7 100644 --- a/internal/hubdb/policy_audit_test.go +++ b/internal/hubdb/policy_audit_test.go @@ -231,14 +231,21 @@ func FuzzPolicyAuditEntryHashUsesLengthFraming(f *testing.F) { f.Add("user:alice", "phase-1-read") f.Add("a", "bc") f.Fuzz(func(t *testing.T, left, right string) { - if right == "" || len(left)+len(right) > 512 { + if len(right) < 2 || len(left)+len(right) > 256 { t.Skip() } first := policyAuditTestEntry() first.actor, first.reasonCode = left, right second := policyAuditTestEntry() - second.actor, second.reasonCode = left+right, "" - if bytes.Equal(policyAuditEntryHash(first), policyAuditEntryHash(second)) { + second.actor, second.reasonCode = left+right[:1], right[1:] + if validatePolicyAuditEntry(first) != nil || validatePolicyAuditEntry(second) != nil { + t.Skip() + } + firstHash, secondHash := policyAuditEntryHash(first), policyAuditEntryHash(second) + if len(firstHash) != 32 || len(secondHash) != 32 { + t.Fatal("valid audit entries did not produce SHA-256 digests") + } + if bytes.Equal(firstHash, secondHash) { t.Fatal("length-delimited audit fields produced an ambiguous digest") } }) diff --git a/internal/hubfleet/collector.go b/internal/hubfleet/collector.go index e5bccbf..4bab20f 100644 --- a/internal/hubfleet/collector.go +++ b/internal/hubfleet/collector.go @@ -121,6 +121,7 @@ type Store interface { // CollectorConfig defines bounded collection behavior. type CollectorConfig struct { + LifecycleContext context.Context Store Store Transport Transport PEP *pep.Enforcer @@ -149,8 +150,8 @@ type Collector struct { // NewCollector constructs a fail-closed collector with bounded per-spoke work. func NewCollector(config CollectorConfig) (*Collector, error) { - if config.Store == nil || config.Transport == nil || config.PEP == nil { - return nil, fmt.Errorf("new spoke collector: store, transport, and policy enforcer are required") + if config.LifecycleContext == nil || config.Store == nil || config.Transport == nil || config.PEP == nil { + return nil, fmt.Errorf("new spoke collector: lifecycle context, store, transport, and policy enforcer are required") } if config.SpokeTimeout == 0 { config.SpokeTimeout = defaultSpokeTimeout @@ -183,7 +184,7 @@ func NewCollector(config CollectorConfig) (*Collector, error) { store: config.Store, transport: config.Transport, pep: config.PEP, - refreshes: newRefreshCoordinator(), + refreshes: newRefreshCoordinator(config.LifecycleContext), observer: config.Observer, tracer: config.TraceObserver, spokeTimeout: config.SpokeTimeout, diff --git a/internal/hubfleet/collector_concurrency_test.go b/internal/hubfleet/collector_concurrency_test.go index cfb0930..974cc15 100644 --- a/internal/hubfleet/collector_concurrency_test.go +++ b/internal/hubfleet/collector_concurrency_test.go @@ -25,8 +25,9 @@ func TestCollectorBoundsParallelSpokesAndSortsCoverage(t *testing.T) { var active atomic.Int64 var maximum atomic.Int64 collector, err := NewCollector(CollectorConfig{ - Store: store, - PEP: testReadPEP(t), + LifecycleContext: t.Context(), + Store: store, + PEP: testReadPEP(t), Transport: transportFunc(func(_ context.Context, _ tenancy.WorkspaceID, spoke Spoke) (Snapshot, error) { current := active.Add(1) defer active.Add(-1) @@ -84,8 +85,9 @@ func TestCollectorCollapsesSerialTimeoutsIntoParallelWave(t *testing.T) { store := newMemoryStore(testSpokes("spoke-a", "spoke-b", "spoke-c", "spoke-d")) collector, err := NewCollector(CollectorConfig{ - Store: store, - PEP: testReadPEP(t), + LifecycleContext: t.Context(), + Store: store, + PEP: testReadPEP(t), Transport: transportFunc(func(ctx context.Context, _ tenancy.WorkspaceID, _ Spoke) (Snapshot, error) { <-ctx.Done() return Snapshot{}, ctx.Err() @@ -120,8 +122,9 @@ func TestCollectorPersistsHealthyPeerWhileAnotherWorkerIsBlocked(t *testing.T) { releaseBlocked := make(chan struct{}) blockedStarted := make(chan struct{}) collector, err := NewCollector(CollectorConfig{ - Store: store, - PEP: testReadPEP(t), + LifecycleContext: t.Context(), + Store: store, + PEP: testReadPEP(t), Transport: transportFunc(func(_ context.Context, _ tenancy.WorkspaceID, spoke Spoke) (Snapshot, error) { if spoke.ID == "spoke-blocked" { close(blockedStarted) @@ -164,9 +167,10 @@ func TestCollectWorkspaceCancellationStopsAdmissionAndWorkers(t *testing.T) { var active atomic.Int64 observer := &recordingSnapshotObserver{} collector, err := NewCollector(CollectorConfig{ - Store: store, - PEP: testReadPEP(t), - Observer: observer, + LifecycleContext: t.Context(), + Store: store, + PEP: testReadPEP(t), + Observer: observer, Transport: transportFunc(func(ctx context.Context, _ tenancy.WorkspaceID, spoke Spoke) (Snapshot, error) { active.Add(1) defer active.Add(-1) @@ -235,8 +239,9 @@ func TestCollectorStoreFailureCancelsWorkersAndLaterAdmission(t *testing.T) { secondWorkerStarted := make(chan struct{}) blockedCanceled := make(chan struct{}) collector, err := NewCollector(CollectorConfig{ - Store: store, - PEP: testReadPEP(t), + LifecycleContext: t.Context(), + Store: store, + PEP: testReadPEP(t), Transport: transportFunc(func(ctx context.Context, _ tenancy.WorkspaceID, spoke Spoke) (Snapshot, error) { started <- spoke.ID if spoke.ID == "spoke-a" { @@ -290,8 +295,9 @@ func TestCollectorRecoversPanickingTransportWorkerAndLaterRefresh(t *testing.T) var panicFlight atomic.Bool panicFlight.Store(true) collector, err := NewCollector(CollectorConfig{ - Store: store, - PEP: testReadPEP(t), + LifecycleContext: t.Context(), + Store: store, + PEP: testReadPEP(t), Transport: transportFunc(func(ctx context.Context, _ tenancy.WorkspaceID, spoke Spoke) (Snapshot, error) { if panicFlight.Load() { switch spoke.ID { @@ -338,8 +344,9 @@ func TestCollectorJoinsWorkersWhenStorePanicsAndLaterRefresh(t *testing.T) { secondWorkerStarted := make(chan struct{}) blockedCanceled := make(chan struct{}) collector, err := NewCollector(CollectorConfig{ - Store: store, - PEP: testReadPEP(t), + LifecycleContext: t.Context(), + Store: store, + PEP: testReadPEP(t), Transport: transportFunc(func(ctx context.Context, _ tenancy.WorkspaceID, spoke Spoke) (Snapshot, error) { if store.panicNext.Load() { switch spoke.ID { diff --git a/internal/hubfleet/collector_test.go b/internal/hubfleet/collector_test.go index ba90583..bae7c14 100644 --- a/internal/hubfleet/collector_test.go +++ b/internal/hubfleet/collector_test.go @@ -29,8 +29,9 @@ func TestCollectorIsolatesSpokeFailuresAndRetainsStaleSnapshots(t *testing.T) { failures: make(map[string]FailureKind), } collector, err := NewCollector(CollectorConfig{ - Store: store, - PEP: testReadPEP(t), + LifecycleContext: t.Context(), + Store: store, + PEP: testReadPEP(t), Transport: transportFunc(func(_ context.Context, workspaceID tenancy.WorkspaceID, spoke Spoke) (Snapshot, error) { if workspaceID != "workspace-a" { return Snapshot{}, errors.New("wrong workspace") @@ -77,8 +78,9 @@ func TestCollectorBoundsEveryTransportCall(t *testing.T) { failures: make(map[string]FailureKind), } collector, err := NewCollector(CollectorConfig{ - Store: store, - PEP: testReadPEP(t), + LifecycleContext: t.Context(), + Store: store, + PEP: testReadPEP(t), Transport: transportFunc(func(ctx context.Context, _ tenancy.WorkspaceID, _ Spoke) (Snapshot, error) { deadline, exists := ctx.Deadline() if !exists || time.Until(deadline) > 1100*time.Millisecond { @@ -170,8 +172,9 @@ func TestCollectorCopiesTransportOwnedSnapshot(t *testing.T) { failures: make(map[string]FailureKind), } collector, err := NewCollector(CollectorConfig{ - Store: store, - PEP: testReadPEP(t), + LifecycleContext: t.Context(), + Store: store, + PEP: testReadPEP(t), Transport: transportFunc(func(context.Context, tenancy.WorkspaceID, Spoke) (Snapshot, error) { return transportSnapshot, nil }), @@ -195,20 +198,28 @@ func TestCollectorAndSourceRejectUnsafeConfiguration(t *testing.T) { t.Parallel() store := &memoryStore{snapshots: make(map[string]Snapshot), failures: make(map[string]FailureKind)} - collector, err := NewCollector(CollectorConfig{Store: store, Transport: transportFunc(func(context.Context, tenancy.WorkspaceID, Spoke) (Snapshot, error) { + transport := transportFunc(func(context.Context, tenancy.WorkspaceID, Spoke) (Snapshot, error) { return Snapshot{}, nil - }), PEP: testReadPEP(t)}) + }) + collector, err := NewCollector(CollectorConfig{ + LifecycleContext: t.Context(), Store: store, Transport: transport, PEP: testReadPEP(t), + }) if err != nil || collector.maxConcurrentSpokes != defaultSpokeConcurrency { t.Fatalf("default worker configuration = %d/%v", collector.maxConcurrentSpokes, err) } - if _, err := NewCollector(CollectorConfig{Store: store, Transport: transportFunc(func(context.Context, tenancy.WorkspaceID, Spoke) (Snapshot, error) { - return Snapshot{}, nil - }), PEP: testReadPEP(t), SpokeTimeout: 999 * time.Millisecond}); err == nil { + if _, err := NewCollector(CollectorConfig{Store: store, Transport: transport, PEP: testReadPEP(t)}); err == nil { + t.Fatal("collector without lifecycle context unexpectedly accepted") + } + if _, err := NewCollector(CollectorConfig{ + LifecycleContext: t.Context(), Store: store, Transport: transport, + PEP: testReadPEP(t), SpokeTimeout: 999 * time.Millisecond, + }); err == nil { t.Fatal("sub-second collector timeout unexpectedly accepted") } - if _, err := NewCollector(CollectorConfig{Store: store, Transport: transportFunc(func(context.Context, tenancy.WorkspaceID, Spoke) (Snapshot, error) { - return Snapshot{}, nil - }), PEP: testReadPEP(t), MaxConcurrentSpokes: maximumSpokeConcurrency + 1}); err == nil { + if _, err := NewCollector(CollectorConfig{ + LifecycleContext: t.Context(), Store: store, Transport: transport, + PEP: testReadPEP(t), MaxConcurrentSpokes: maximumSpokeConcurrency + 1, + }); err == nil { t.Fatal("oversized collector worker pool unexpectedly accepted") } if _, err := NewSource(SourceConfig{Reader: fleetReaderFunc(func(context.Context, tenancy.Scope, time.Duration, time.Time) (fleet.FleetResult, error) { diff --git a/internal/hubfleet/metrics_test.go b/internal/hubfleet/metrics_test.go index 8a6e5d2..d47d2c6 100644 --- a/internal/hubfleet/metrics_test.go +++ b/internal/hubfleet/metrics_test.go @@ -47,7 +47,8 @@ func TestCollectorObservesClosedSnapshotOutcomes(t *testing.T) { spokes: []Spoke{{ID: "spoke-a", ManagedClusterRef: "ocm/spoke-a"}}, snapshots: make(map[string]Snapshot), failures: make(map[string]FailureKind), } collector, err := NewCollector(CollectorConfig{ - Store: store, Transport: test.transport, PEP: testReadPEP(t), Observer: observer, Now: func() time.Time { return now }, + LifecycleContext: t.Context(), Store: store, Transport: test.transport, + PEP: testReadPEP(t), Observer: observer, Now: func() time.Time { return now }, }) if err != nil { t.Fatal(err) @@ -65,6 +66,7 @@ func TestCollectorObservesClosedSnapshotOutcomes(t *testing.T) { func TestCollectorRecoversFromPanickingSnapshotObserver(t *testing.T) { now := time.Date(2026, time.July, 12, 20, 0, 0, 0, time.UTC) collector, err := NewCollector(CollectorConfig{ + LifecycleContext: t.Context(), Store: &memoryStore{ spokes: []Spoke{{ID: "spoke-a", ManagedClusterRef: "ocm/spoke-a"}}, snapshots: make(map[string]Snapshot), failures: make(map[string]FailureKind), }, diff --git a/internal/hubfleet/policy_test.go b/internal/hubfleet/policy_test.go index 1b3eddd..822121e 100644 --- a/internal/hubfleet/policy_test.go +++ b/internal/hubfleet/policy_test.go @@ -20,7 +20,7 @@ func TestHubReadEntrypointsStopBeforeDependenciesWhenPolicyRefuses(t *testing.T) store := &memoryStore{snapshots: make(map[string]Snapshot), failures: make(map[string]FailureKind)} collector, err := NewCollector(CollectorConfig{ - Store: store, PEP: refusal.enforcer(t), + LifecycleContext: t.Context(), Store: store, PEP: refusal.enforcer(t), Transport: transportFunc(func(context.Context, tenancy.WorkspaceID, Spoke) (Snapshot, error) { return Snapshot{}, errors.New("transport must not run") }), @@ -118,7 +118,7 @@ func (refusal *policyRefusal) enforcer(t *testing.T) *pep.Enforcer { func TestHubReadConstructorsRequirePolicyEnforcer(t *testing.T) { store := &memoryStore{snapshots: make(map[string]Snapshot), failures: make(map[string]FailureKind)} - if _, err := NewCollector(CollectorConfig{Store: store, Transport: transportFunc(func(context.Context, tenancy.WorkspaceID, Spoke) (Snapshot, error) { + if _, err := NewCollector(CollectorConfig{LifecycleContext: t.Context(), Store: store, Transport: transportFunc(func(context.Context, tenancy.WorkspaceID, Spoke) (Snapshot, error) { return Snapshot{}, nil })}); err == nil { t.Fatal("NewCollector() accepted no policy enforcer") diff --git a/internal/hubfleet/refresh_coordinator.go b/internal/hubfleet/refresh_coordinator.go index ea7d2b7..f3937ab 100644 --- a/internal/hubfleet/refresh_coordinator.go +++ b/internal/hubfleet/refresh_coordinator.go @@ -7,6 +7,7 @@ import ( "errors" "fmt" "sync" + "time" "github.com/ArdurAI/sith/internal/fleet" "github.com/ArdurAI/sith/internal/tenancy" @@ -22,14 +23,28 @@ type refreshFlight struct { } type refreshCoordinator struct { - mu sync.Mutex - flights map[tenancy.WorkspaceID]*refreshFlight + lifecycle context.Context + mu sync.Mutex + flights map[tenancy.WorkspaceID]*refreshFlight } -func newRefreshCoordinator() *refreshCoordinator { - return &refreshCoordinator{flights: make(map[tenancy.WorkspaceID]*refreshFlight)} +func newRefreshCoordinator(lifecycle context.Context) *refreshCoordinator { + return &refreshCoordinator{ + lifecycle: valueFreeContext{Context: lifecycle}, + flights: make(map[tenancy.WorkspaceID]*refreshFlight), + } +} + +// valueFreeContext preserves only lifecycle cancellation and deadlines. Refresh work must not +// inherit request-scoped credentials, authorization state, or trace identity from its owner. +type valueFreeContext struct { + context.Context } +func (ctx valueFreeContext) Value(any) any { return nil } + +func (ctx valueFreeContext) Deadline() (time.Time, bool) { return ctx.Context.Deadline() } + func (coordinator *refreshCoordinator) collect( ctx context.Context, scope tenancy.Scope, @@ -48,7 +63,7 @@ func (coordinator *refreshCoordinator) collect( if !exists { flight = &refreshFlight{done: make(chan struct{})} coordinator.flights[workspaceID] = flight - //nolint:gosec // A flight must outlive one caller; run supplies a fresh value-free context. + //nolint:gosec // A flight outlives one caller but remains bounded by the collector lifecycle. go coordinator.run(workspaceID, flight, scope, collect) } coordinator.mu.Unlock() @@ -72,7 +87,7 @@ func (coordinator *refreshCoordinator) run( coordinator.finish(workspaceID, flight, coverage, err) }() - flightContext, _, traceErr := tracing.Ensure(context.Background()) + flightContext, _, traceErr := tracing.Ensure(coordinator.lifecycle) if traceErr != nil { err = fmt.Errorf("collect spoke snapshots: establish refresh flight trace: %w", traceErr) return diff --git a/internal/hubfleet/refresh_coordinator_test.go b/internal/hubfleet/refresh_coordinator_test.go index 0b6d257..d18a1ba 100644 --- a/internal/hubfleet/refresh_coordinator_test.go +++ b/internal/hubfleet/refresh_coordinator_test.go @@ -5,6 +5,7 @@ package hubfleet import ( "context" "errors" + "fmt" "sync" "sync/atomic" "testing" @@ -39,8 +40,9 @@ func TestCollectorCoalescesAuthorizedRequestsForOneWorkspace(t *testing.T) { t.Fatal(err) } collector, err := NewCollector(CollectorConfig{ - Store: store, - PEP: enforcer, + LifecycleContext: t.Context(), + Store: store, + PEP: enforcer, Transport: transportFunc(func(context.Context, tenancy.WorkspaceID, Spoke) (Snapshot, error) { select { case <-started: @@ -133,8 +135,9 @@ func TestCollectorRefusesCallerBeforeJoiningWorkspaceFlight(t *testing.T) { t.Fatal(err) } collector, err := NewCollector(CollectorConfig{ - Store: store, - PEP: enforcer, + LifecycleContext: t.Context(), + Store: store, + PEP: enforcer, Transport: transportFunc(func(context.Context, tenancy.WorkspaceID, Spoke) (Snapshot, error) { close(started) <-release @@ -170,7 +173,7 @@ func TestCollectorRefusesCallerBeforeJoiningWorkspaceFlight(t *testing.T) { func TestRefreshCoordinatorKeepsWorkspacesIndependent(t *testing.T) { t.Parallel() - coordinator := newRefreshCoordinator() + coordinator := newRefreshCoordinator(t.Context()) releaseA := make(chan struct{}) startedA := make(chan struct{}) startedB := make(chan struct{}) @@ -215,7 +218,7 @@ func TestRefreshCoordinatorKeepsWorkspacesIndependent(t *testing.T) { func TestRefreshCoordinatorDetachesLeaderAndWaiterCancellation(t *testing.T) { t.Parallel() - coordinator := newRefreshCoordinator() + coordinator := newRefreshCoordinator(t.Context()) started := make(chan struct{}) release := make(chan struct{}) var calls atomic.Int32 @@ -280,7 +283,7 @@ func TestRefreshCoordinatorDetachesLeaderAndWaiterCancellation(t *testing.T) { func TestRefreshCoordinatorSharesFailuresAndRecoversAfterPanic(t *testing.T) { t.Parallel() - coordinator := newRefreshCoordinator() + coordinator := newRefreshCoordinator(t.Context()) sharedFailure := errors.New("closed store failure") started := make(chan struct{}) release := make(chan struct{}) @@ -341,6 +344,38 @@ func TestRefreshCoordinatorSharesFailuresAndRecoversAfterPanic(t *testing.T) { } } +func TestRefreshCoordinatorStopsFlightsAtLifecycleBoundaryWithoutValues(t *testing.T) { + t.Parallel() + + type lifecycleValueKey struct{} + lifecycle, cancelLifecycle := context.WithCancel( + context.WithValue(context.Background(), lifecycleValueKey{}, "must-not-cross-refresh-boundary"), + ) + defer cancelLifecycle() + coordinator := newRefreshCoordinator(lifecycle) + started := make(chan struct{}) + result := make(chan error, 1) + go func() { + _, err := coordinator.collect(context.Background(), readerScope(t, "workspace-a"), func(ctx context.Context, _ tenancy.Scope) (fleet.Coverage, error) { + close(started) + if value := ctx.Value(lifecycleValueKey{}); value != nil { + return fleet.Coverage{}, fmt.Errorf("lifecycle value leaked into refresh: %v", value) + } + <-ctx.Done() + return fleet.Coverage{}, ctx.Err() + }) + result <- err + }() + <-started + cancelLifecycle() + if err := <-result; !errors.Is(err, context.Canceled) { + t.Fatalf("lifecycle cancellation error = %v, want context canceled", err) + } + if coordinator.activeFlights() != 0 { + t.Fatalf("lifecycle cancellation retained %d refresh flights", coordinator.activeFlights()) + } +} + func eventually(t *testing.T, condition func() bool) { t.Helper() deadline := time.Now().Add(5 * time.Second) diff --git a/internal/hubfleet/tracing_test.go b/internal/hubfleet/tracing_test.go index 3b06341..a047554 100644 --- a/internal/hubfleet/tracing_test.go +++ b/internal/hubfleet/tracing_test.go @@ -35,7 +35,8 @@ func TestCollectorSeparatesCallerAuthorizationTraceFromRefreshFlight(t *testing. spokes: []Spoke{{ID: "spoke-a", ManagedClusterRef: "ocm/spoke-a"}}, snapshots: make(map[string]Snapshot), failures: make(map[string]FailureKind), } collector, err := NewCollector(CollectorConfig{ - Store: store, + LifecycleContext: t.Context(), + Store: store, Transport: transportFunc(func(ctx context.Context, _ tenancy.WorkspaceID, spoke Spoke) (Snapshot, error) { callerValueLeaked = ctx.Value(callerValueKey{}) != nil var ok bool @@ -77,7 +78,8 @@ func TestCollectorSurvivesPanickingTraceObserver(t *testing.T) { t.Fatal(err) } collector, err := NewCollector(CollectorConfig{ - Store: &memoryStore{spokes: []Spoke{{ID: "spoke-a", ManagedClusterRef: "ocm/spoke-a"}}, snapshots: make(map[string]Snapshot), failures: make(map[string]FailureKind)}, + LifecycleContext: t.Context(), + Store: &memoryStore{spokes: []Spoke{{ID: "spoke-a", ManagedClusterRef: "ocm/spoke-a"}}, snapshots: make(map[string]Snapshot), failures: make(map[string]FailureKind)}, Transport: transportFunc(func(context.Context, tenancy.WorkspaceID, Spoke) (Snapshot, error) { return validSnapshot("spoke-a", now), nil }), diff --git a/internal/hubruntime/config.go b/internal/hubruntime/config.go index 709c934..84c7cf6 100644 --- a/internal/hubruntime/config.go +++ b/internal/hubruntime/config.go @@ -150,7 +150,10 @@ func NewFromEnvironment(ctx context.Context, logger *slog.Logger) (*Runtime, err _ = processAuthObserver.Close() database.Close() } - collector, err := hubfleet.NewCollector(hubfleet.CollectorConfig{Store: database, Transport: transport, PEP: enforcer, Observer: metrics, TraceObserver: tracer}) + collector, err := hubfleet.NewCollector(hubfleet.CollectorConfig{ + LifecycleContext: ctx, Store: database, Transport: transport, PEP: enforcer, + Observer: metrics, TraceObserver: tracer, + }) if err != nil { cleanup() return nil, fmt.Errorf("construct hub runtime: collector configuration is invalid") diff --git a/internal/hubruntime/ocm_integration_test.go b/internal/hubruntime/ocm_integration_test.go index dcde282..1cdeabc 100644 --- a/internal/hubruntime/ocm_integration_test.go +++ b/internal/hubruntime/ocm_integration_test.go @@ -66,7 +66,9 @@ func TestHubRuntimeDirectClusterProxyM0(t *testing.T) { if err != nil { t.Fatal(err) } - collector, err := hubfleet.NewCollector(hubfleet.CollectorConfig{Store: store, Transport: transport, PEP: enforcer}) + collector, err := hubfleet.NewCollector(hubfleet.CollectorConfig{ + LifecycleContext: ctx, Store: store, Transport: transport, PEP: enforcer, + }) if err != nil { t.Fatal(err) } diff --git a/internal/pep/audit_test.go b/internal/pep/audit_test.go index d9340ce..c7e5c87 100644 --- a/internal/pep/audit_test.go +++ b/internal/pep/audit_test.go @@ -190,6 +190,16 @@ func TestAuditEventRejectsUnsafeInvalidRequestSentinel(t *testing.T) { } } +func TestAuditEventRejectsMalformedUTF8Actor(t *testing.T) { + t.Parallel() + + event := policyAuditEvent(time.Now().UTC(), VerdictAllow, "phase-1-read") + event.Actor = string([]byte{'u', 's', 'e', 'r', ':', 0x80}) + if err := event.Validate(); err == nil { + t.Fatal("AuditEvent.Validate() accepted a malformed UTF-8 actor") + } +} + func TestAuditEventRejectsImpossibleProposalRoleOutcome(t *testing.T) { base := policyAuditEvent(time.Now().UTC(), VerdictDeny, "role-denied") base.Action = tenancy.ActionProposeIntent diff --git a/internal/pep/pep.go b/internal/pep/pep.go index 68f3cb8..e56f709 100644 --- a/internal/pep/pep.go +++ b/internal/pep/pep.go @@ -11,6 +11,7 @@ import ( "strings" "time" "unicode" + "unicode/utf8" "github.com/ArdurAI/sith/internal/fleet" "github.com/ArdurAI/sith/internal/intent" @@ -459,7 +460,7 @@ func (enforcer *Enforcer) record(ctx context.Context, request Request, decision } func validateSafeText(name, value string, maximum int) error { - if value == "" || strings.TrimSpace(value) != value || len(value) > maximum { + if value == "" || !utf8.ValidString(value) || strings.TrimSpace(value) != value || len(value) > maximum { return fmt.Errorf("%s is invalid", name) } for _, character := range value { diff --git a/internal/tenancy/model.go b/internal/tenancy/model.go index 5b93643..2d6e0f3 100644 --- a/internal/tenancy/model.go +++ b/internal/tenancy/model.go @@ -6,6 +6,7 @@ import ( "fmt" "strings" "unicode" + "unicode/utf8" ) const maxIdentityLength = 256 @@ -110,7 +111,7 @@ func (role Role) Allows(action Action) bool { } func validateIdentity(name, value string) error { - if value == "" || strings.TrimSpace(value) != value { + if value == "" || !utf8.ValidString(value) || strings.TrimSpace(value) != value { return fmt.Errorf("%s must be a non-empty, trimmed value", name) } if len(value) > maxIdentityLength { @@ -125,7 +126,7 @@ func validateIdentity(name, value string) error { } func validateDisplayName(name, value string) error { - if value == "" || strings.TrimSpace(value) != value { + if value == "" || !utf8.ValidString(value) || strings.TrimSpace(value) != value { return fmt.Errorf("%s must be a non-empty, trimmed value", name) } if len(value) > maxIdentityLength { diff --git a/internal/tenancy/model_test.go b/internal/tenancy/model_test.go index 8355cb7..19c0847 100644 --- a/internal/tenancy/model_test.go +++ b/internal/tenancy/model_test.go @@ -86,6 +86,15 @@ func TestTenancyModelsRejectAmbiguousIdentity(t *testing.T) { {name: "padded tenant key", run: func() error { return (Workspace{ID: "a", Name: "A", TenantKey: " tenant-a"}).Validate() }}, {name: "unknown role", run: func() error { return (Membership{WorkspaceID: "a", Subject: "alice", Role: "owner"}).Validate() }}, {name: "control character", run: func() error { return (Membership{WorkspaceID: "a\n", Subject: "alice", Role: RoleReader}).Validate() }}, + {name: "malformed UTF-8 workspace ID", run: func() error { + return ValidateWorkspaceID(WorkspaceID(string([]byte{'a', 0x80}))) + }}, + {name: "malformed UTF-8 display name", run: func() error { + return (Workspace{ID: "a", Name: string([]byte{'A', 0x80}), TenantKey: "tenant-a"}).Validate() + }}, + {name: "malformed UTF-8 subject", run: func() error { + return (Membership{WorkspaceID: "a", Subject: string([]byte{'a', 0x80}), Role: RoleReader}).Validate() + }}, } for _, test := range tests { t.Run(test.name, func(t *testing.T) { diff --git a/sessions/2026-07-21-deep-quality-audit.md b/sessions/2026-07-21-deep-quality-audit.md new file mode 100644 index 0000000..844fa40 --- /dev/null +++ b/sessions/2026-07-21-deep-quality-audit.md @@ -0,0 +1,66 @@ +# Session — 2026-07-21 — deep quality audit + +**Builder:** Gnani Rahul Nutakki · **Branch:** `gnanirahulnutakki/audit-hardening-20260721` +**Scope:** repository-wide deep audit and open-issue completion · **Status:** in progress + +## [G] Goal + +Audit exact `dev` with current primary-source tooling, fix every confirmed defect in narrow signed +changes, then complete every independently buildable open issue without weakening fail-closed +boundaries. + +## [S] Scope + +- This checkpoint covers malformed UTF-8 rejection at the tenancy and policy-audit boundary, + audit-hash fuzz-oracle correctness, and lifecycle ownership for coalesced spoke refreshes. +- Release-policy drift, dependency/toolchain updates, desktop-build version drift, approval expiry, + and remaining open-issue slices stay separate so each can be reviewed and reverted independently. +- No credential, authorization, schema, cloud resource, external package, or API surface changes. + +## [A] Action + +- Repaired the audit-hash fuzz target so it compares two valid records with the same unframed bytes + split at different actor/reason boundaries instead of comparing validation-error sentinels. +- Added explicit UTF-8 validation for policy actors and tenancy identities. This aligns the durable + writer with the portable audit verifier and fails before malformed bytes reach hashing or storage. +- Made the collector lifecycle context mandatory. Shared refreshes remain detached from individual + callers but now inherit hub shutdown cancellation and deadlines through a value-free boundary. +- Updated every collector construction site and the README's current refresh-lifecycle contract. + +## [T] Proof + +- `go test -race -count=10 ./internal/hubfleet ./internal/pep ./internal/tenancy ./internal/hubdb` + passed. +- `FuzzPolicyAuditEntryHashUsesLengthFraming` passed 250,000 executions after the regression fix; + all 23 repository fuzz targets separately passed 50,000 executions each. +- The first full local CI run exposed a stale release-policy assertion and a hosted coverage gap. + The separately signed prerequisite landed through PR #294, and exact merge SHA `f0e071a` passed + expanded CI, reproducible release, real kind, and all three CodeQL analyses. +- After rebasing onto that verified SHA, complete local CI passed formatting, vet, zero-issue lint, + `govulncheck`, every race test and shell policy, nine alert rules, performance, binary end-to-end, + and production build. +- `make e2e-isolation` passed PostgreSQL 18.4 forced-RLS coverage and both workspace-isolation + fuzzers at 50,000 executions each. +- `make release-check` passed dual four-platform reproducibility, SPDX SBOMs, checksums, Homebrew + formula generation, and release-derived amd64/arm64 OCI verification. +- The Kubernetes 1.36.1 kind gate passed fleet fan-out, OCI image, and Argo Application projection + under the race detector in 260.143 seconds; teardown left no kind cluster or Buildx builder. +- `README.md` was reviewed and updated because refresh lifetime is now bounded by hub shutdown. + +## [S] Security, reliability, and cost + +Malformed identities now fail consistently across online and offline audit paths. Detached work no +longer survives process cancellation, while lifecycle values and request credentials remain outside +the refresh context. The change adds no cloud resource, API call, storage, telemetry cardinality, +or recurring cost. + +## [N] Next + +Push the rebased signed branch, obtain exact-head hosted CI and CodeQL proof, merge without rewriting +the signed commit, and verify the exact post-merge `dev` SHA. + +## [C] Checkpoint #1 + +The validation alignment, lifecycle boundary, adversarial tests, fuzz repair, construction-site +migration, README correction, prerequisite merge, and complete local gate matrix are frozen for +signed review. Hosted exact-head and exact post-merge `dev` proof remain mandatory. diff --git a/tests/e2e/kind_read_federation_test.go b/tests/e2e/kind_read_federation_test.go index 3c13df8..21b0ce1 100644 --- a/tests/e2e/kind_read_federation_test.go +++ b/tests/e2e/kind_read_federation_test.go @@ -39,9 +39,10 @@ func exerciseReadFederationSnapshots( failures: make(map[string]hubfleet.FailureKind), } collector, err := hubfleet.NewCollector(hubfleet.CollectorConfig{ - Store: store, - Transport: kindSnapshotTransport{adapter: adapter}, - PEP: e2eReadPEP(t), + LifecycleContext: ctx, + Store: store, + Transport: kindSnapshotTransport{adapter: adapter}, + PEP: e2eReadPEP(t), }) if err != nil { t.Fatalf("construct read-federation collector: %v", err)