diff --git a/weed/iam/integration/iam_manager.go b/weed/iam/integration/iam_manager.go index f9fc01d57..48e146aa8 100644 --- a/weed/iam/integration/iam_manager.go +++ b/weed/iam/integration/iam_manager.go @@ -32,6 +32,7 @@ type IAMManager struct { roleStore RoleStore userStore UserStore oidcProviderStore OIDCProviderStore + oidcAuditSink OIDCProviderAuditSink revocationStore SessionRevocationStore filerAddressProvider func() string // Function to get current filer address initialized bool @@ -39,6 +40,31 @@ type IAMManager struct { runtimePolicyNames map[string]struct{} } +// SetOIDCProviderAuditSink configures the lifecycle event sink. When nil +// (default), GlogAuditSink is used so events still surface in logs. +func (m *IAMManager) SetOIDCProviderAuditSink(sink OIDCProviderAuditSink) { + m.oidcAuditSink = sink +} + +// emitOIDCAudit logs a lifecycle event. Errors are swallowed: an audit +// failure must never block an IAM mutation that has already succeeded. +func (m *IAMManager) emitOIDCAudit(ctx context.Context, eventType OIDCProviderAuditEventType, arn, url string, detail map[string]string) { + sink := m.oidcAuditSink + if sink == nil { + sink = GlogAuditSink{} + } + event := &OIDCProviderAuditEvent{ + Type: eventType, + ARN: arn, + URL: url, + Detail: detail, + OccurredAt: time.Now().UTC(), + } + if err := sink.Emit(ctx, event); err != nil { + glog.Warningf("OIDC audit emit %s for %s: %v", eventType, arn, err) + } +} + // SetSessionRevocationStore configures the per-session revocation list. When // nil, RevokeSession returns an error and IsSessionRevoked is a no-op (every // session is considered live until natural expiry). Operators who want @@ -138,6 +164,7 @@ func (m *IAMManager) CreateOIDCProvider(ctx context.Context, rec *OIDCProviderRe return err } m.refreshOIDCProvidersBestEffort(ctx, "CreateOIDCProvider", rec.ARN) + m.emitOIDCAudit(ctx, OIDCAuditEventCreated, rec.ARN, rec.URL, nil) return nil } @@ -150,6 +177,7 @@ func (m *IAMManager) DeleteOIDCProvider(ctx context.Context, arn string) error { return err } m.refreshOIDCProvidersBestEffort(ctx, "DeleteOIDCProvider", arn) + m.emitOIDCAudit(ctx, OIDCAuditEventDeleted, arn, "", nil) return nil } @@ -180,6 +208,7 @@ func (m *IAMManager) AddClientIDToOIDCProvider(ctx context.Context, arn, clientI return err } m.refreshOIDCProvidersBestEffort(ctx, "AddClientIDToOIDCProvider", arn) + m.emitOIDCAudit(ctx, OIDCAuditEventClientIDAdded, rec.ARN, rec.URL, map[string]string{"clientId": clientID}) return nil } @@ -208,6 +237,7 @@ func (m *IAMManager) RemoveClientIDFromOIDCProvider(ctx context.Context, arn, cl return err } m.refreshOIDCProvidersBestEffort(ctx, "RemoveClientIDFromOIDCProvider", arn) + m.emitOIDCAudit(ctx, OIDCAuditEventClientIDRemoved, rec.ARN, rec.URL, map[string]string{"clientId": clientID}) return nil } @@ -235,6 +265,7 @@ func (m *IAMManager) UpdateOIDCProviderThumbprints(ctx context.Context, arn stri return err } m.refreshOIDCProvidersBestEffort(ctx, "UpdateOIDCProviderThumbprints", arn) + m.emitOIDCAudit(ctx, OIDCAuditEventThumbprintsSet, rec.ARN, rec.URL, map[string]string{"count": fmt.Sprintf("%d", len(thumbprints))}) return nil } @@ -254,7 +285,11 @@ func (m *IAMManager) TagOIDCProvider(ctx context.Context, arn string, tags map[s rec.Tags[k] = v } rec.UpdatedAt = time.Now().UTC() - return m.oidcProviderStore.StoreProvider(ctx, m.getFilerAddress(), rec) + if err := m.oidcProviderStore.StoreProvider(ctx, m.getFilerAddress(), rec); err != nil { + return err + } + m.emitOIDCAudit(ctx, OIDCAuditEventTagsAdded, rec.ARN, rec.URL, map[string]string{"count": fmt.Sprintf("%d", len(tags))}) + return nil } // UntagOIDCProvider removes the named tags from the provider's tag set. @@ -270,7 +305,11 @@ func (m *IAMManager) UntagOIDCProvider(ctx context.Context, arn string, keys []s delete(rec.Tags, k) } rec.UpdatedAt = time.Now().UTC() - return m.oidcProviderStore.StoreProvider(ctx, m.getFilerAddress(), rec) + if err := m.oidcProviderStore.StoreProvider(ctx, m.getFilerAddress(), rec); err != nil { + return err + } + m.emitOIDCAudit(ctx, OIDCAuditEventTagsRemoved, rec.ARN, rec.URL, map[string]string{"count": fmt.Sprintf("%d", len(keys))}) + return nil } // validateOIDCProviderRecord enforces the invariants AWS imposes on the diff --git a/weed/iam/integration/oidc_provider_audit.go b/weed/iam/integration/oidc_provider_audit.go new file mode 100644 index 000000000..8769ab371 --- /dev/null +++ b/weed/iam/integration/oidc_provider_audit.go @@ -0,0 +1,173 @@ +package integration + +import ( + "context" + "crypto/sha1" + "encoding/hex" + "encoding/json" + "fmt" + "strings" + "sync" + "time" + + "github.com/seaweedfs/seaweedfs/weed/glog" + "github.com/seaweedfs/seaweedfs/weed/pb" + "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb" + "google.golang.org/grpc" +) + +// OIDCProviderAuditEventType enumerates the lifecycle events emitted when an +// IAM-managed OIDC provider record is mutated. Use-events (every successful +// token validation) are intentionally not emitted by the lifecycle path: +// they're too hot to stream into a filer-backed sink without spikes. A +// future PR can add a sampled-use sink behind an explicit opt-in flag. +type OIDCProviderAuditEventType string + +const ( + OIDCAuditEventCreated OIDCProviderAuditEventType = "Create" + OIDCAuditEventDeleted OIDCProviderAuditEventType = "Delete" + OIDCAuditEventClientIDAdded OIDCProviderAuditEventType = "AddClientID" + OIDCAuditEventClientIDRemoved OIDCProviderAuditEventType = "RemoveClientID" + OIDCAuditEventThumbprintsSet OIDCProviderAuditEventType = "UpdateThumbprints" + OIDCAuditEventTagsAdded OIDCProviderAuditEventType = "Tag" + OIDCAuditEventTagsRemoved OIDCProviderAuditEventType = "Untag" +) + +// OIDCProviderAuditEvent is the payload emitted for each lifecycle event. +type OIDCProviderAuditEvent struct { + Type OIDCProviderAuditEventType `json:"type"` + ARN string `json:"arn"` + URL string `json:"url,omitempty"` + Detail map[string]string `json:"detail,omitempty"` + OccurredAt time.Time `json:"occurredAt"` +} + +// OIDCProviderAuditSink consumes lifecycle events. Implementations must be +// safe for concurrent use; emit happens with the IAM manager's mutation lock +// held, so a slow sink slows the IAM API. +type OIDCProviderAuditSink interface { + Emit(ctx context.Context, event *OIDCProviderAuditEvent) error +} + +// GlogAuditSink writes one structured log line per event. Always-on default +// when an explicit sink isn't configured — events still appear in stdout/the +// log aggregator, just not in a queryable record. +type GlogAuditSink struct{} + +func (GlogAuditSink) Emit(_ context.Context, event *OIDCProviderAuditEvent) error { + if event == nil { + return nil + } + data, err := json.Marshal(event) + if err != nil { + glog.V(0).Infof("oidc-audit: %s arn=%s err-marshal=%v", event.Type, event.ARN, err) + return nil + } + glog.V(0).Infof("oidc-audit: %s", string(data)) + return nil +} + +// MemoryAuditSink keeps events in process memory for tests and short-lived +// inspection. Not durable; drop entries by reading from Events() and +// discarding. +type MemoryAuditSink struct { + mu sync.Mutex + events []*OIDCProviderAuditEvent +} + +func NewMemoryAuditSink() *MemoryAuditSink { + return &MemoryAuditSink{} +} + +func (m *MemoryAuditSink) Emit(_ context.Context, event *OIDCProviderAuditEvent) error { + m.mu.Lock() + defer m.mu.Unlock() + cp := *event + m.events = append(m.events, &cp) + return nil +} + +// Events returns a copy of the captured event log so tests can assert. +func (m *MemoryAuditSink) Events() []*OIDCProviderAuditEvent { + m.mu.Lock() + defer m.mu.Unlock() + out := make([]*OIDCProviderAuditEvent, len(m.events)) + copy(out, m.events) + return out +} + +// FilerAuditSink appends event records as separate files under a filer +// directory. Each event is its own file so concurrent writers don't conflict; +// the filename is `--.json`. For higher volumes, +// switch to an append-only journal — the contract is just `Emit`. +type FilerAuditSink struct { + grpcDialOption grpc.DialOption + basePath string + filerAddressProvider func() string +} + +// NewFilerAuditSink returns a filer-backed sink. Default basePath +// `/etc/iam/audit/oidc-providers` keeps audit records out of the active +// IAM data directories. +func NewFilerAuditSink(config map[string]interface{}, filerAddressProvider func() string) *FilerAuditSink { + sink := &FilerAuditSink{ + basePath: "/etc/iam/audit/oidc-providers", + filerAddressProvider: filerAddressProvider, + } + if config != nil { + if bp, ok := config["basePath"].(string); ok && bp != "" { + sink.basePath = strings.TrimSuffix(bp, "/") + } + } + return sink +} + +func (f *FilerAuditSink) resolveFilerAddress() string { + if f.filerAddressProvider != nil { + return f.filerAddressProvider() + } + return "" +} + +func (f *FilerAuditSink) Emit(ctx context.Context, event *OIDCProviderAuditEvent) error { + addr := f.resolveFilerAddress() + if addr == "" { + return fmt.Errorf("filer address not available") + } + if event == nil { + return nil + } + if event.OccurredAt.IsZero() { + event.OccurredAt = time.Now().UTC() + } + data, err := json.MarshalIndent(event, "", " ") + if err != nil { + return fmt.Errorf("marshal audit event: %v", err) + } + // Filename: --.json. Two events at the same + // nano against the same ARN+type are vanishingly rare, but two events + // at the same nano against different ARNs are perfectly plausible + // (audit-via-batch script). The ARN hash gives us an effectively- + // collision-free name without leaking the ARN into the filer path. + arnSum := sha1.Sum([]byte(event.ARN)) + name := fmt.Sprintf("%d-%s-%s.json", event.OccurredAt.UnixNano(), event.Type, hex.EncodeToString(arnSum[:8])) + occurredAt := event.OccurredAt.Unix() + return pb.WithGrpcFilerClient(false, 0, pb.ServerAddress(addr), f.grpcDialOption, func(client filer_pb.SeaweedFilerClient) error { + _, err := client.CreateEntry(ctx, &filer_pb.CreateEntryRequest{ + Directory: f.basePath, + Entry: &filer_pb.Entry{ + Name: name, + IsDirectory: false, + Attributes: &filer_pb.FuseAttributes{ + // Reflect the event's occurrence time so file metadata + // matches the audit record itself, not the filer write. + Mtime: occurredAt, + Crtime: occurredAt, + FileMode: uint32(0o600), + }, + Content: data, + }, + }) + return err + }) +} diff --git a/weed/iam/integration/oidc_provider_audit_test.go b/weed/iam/integration/oidc_provider_audit_test.go new file mode 100644 index 000000000..d2516e425 --- /dev/null +++ b/weed/iam/integration/oidc_provider_audit_test.go @@ -0,0 +1,105 @@ +package integration + +import ( + "context" + "testing" + "time" + + "github.com/seaweedfs/seaweedfs/weed/iam/policy" + "github.com/seaweedfs/seaweedfs/weed/iam/sts" +) + +func newAuditableManager(t *testing.T) (*IAMManager, *MemoryAuditSink) { + t.Helper() + mgr := NewIAMManager() + sink := NewMemoryAuditSink() + cfg := &IAMConfig{ + STS: &sts.STSConfig{ + TokenDuration: sts.FlexibleDuration{Duration: time.Hour}, + MaxSessionLength: sts.FlexibleDuration{Duration: 12 * time.Hour}, + Issuer: "test-sts", + SigningKey: []byte("test-signing-key-32-characters-long"), + AccountId: "111122223333", + }, + Policy: &policy.PolicyEngineConfig{DefaultEffect: "Deny", StoreType: "memory"}, + Roles: &RoleStoreConfig{StoreType: "memory"}, + } + if err := mgr.Initialize(cfg, func() string { return "localhost:8888" }); err != nil { + t.Fatalf("Initialize: %v", err) + } + mgr.SetOIDCProviderAuditSink(sink) + return mgr, sink +} + +func TestOIDCAuditLifecycleEventsAreEmitted(t *testing.T) { + mgr, sink := newAuditableManager(t) + ctx := context.Background() + + rec := &OIDCProviderRecord{ + AccountID: "111122223333", + ARN: "arn:aws:iam::111122223333:oidc-provider/idp.example", + URL: "https://idp.example", + ClientIDs: []string{"x"}, + } + if err := mgr.CreateOIDCProvider(ctx, rec); err != nil { + t.Fatalf("Create: %v", err) + } + if err := mgr.AddClientIDToOIDCProvider(ctx, rec.ARN, "y"); err != nil { + t.Fatalf("AddClientID: %v", err) + } + if err := mgr.RemoveClientIDFromOIDCProvider(ctx, rec.ARN, "y"); err != nil { + t.Fatalf("RemoveClientID: %v", err) + } + if err := mgr.UpdateOIDCProviderThumbprints(ctx, rec.ARN, []string{"0000000000000000000000000000000000000000"}); err != nil { + t.Fatalf("UpdateThumbprints: %v", err) + } + if err := mgr.TagOIDCProvider(ctx, rec.ARN, map[string]string{"team": "infra"}); err != nil { + t.Fatalf("Tag: %v", err) + } + if err := mgr.UntagOIDCProvider(ctx, rec.ARN, []string{"team"}); err != nil { + t.Fatalf("Untag: %v", err) + } + if err := mgr.DeleteOIDCProvider(ctx, rec.ARN); err != nil { + t.Fatalf("Delete: %v", err) + } + + want := []OIDCProviderAuditEventType{ + OIDCAuditEventCreated, + OIDCAuditEventClientIDAdded, + OIDCAuditEventClientIDRemoved, + OIDCAuditEventThumbprintsSet, + OIDCAuditEventTagsAdded, + OIDCAuditEventTagsRemoved, + OIDCAuditEventDeleted, + } + events := sink.Events() + if len(events) != len(want) { + t.Fatalf("expected %d events, got %d (%+v)", len(want), len(events), events) + } + for i, e := range events { + if e.Type != want[i] { + t.Errorf("event %d: type=%s want=%s", i, e.Type, want[i]) + } + if e.ARN != rec.ARN { + t.Errorf("event %d: ARN=%s want=%s", i, e.ARN, rec.ARN) + } + if e.OccurredAt.IsZero() { + t.Errorf("event %d has zero OccurredAt", i) + } + } +} + +func TestOIDCAuditDefaultsToGlogSink(t *testing.T) { + // No SetOIDCProviderAuditSink call — should default to glog and not panic. + mgr, _ := newAuditableManager(t) + mgr.SetOIDCProviderAuditSink(nil) + rec := &OIDCProviderRecord{ + AccountID: "111122223333", + ARN: "arn:aws:iam::111122223333:oidc-provider/idp.example", + URL: "https://idp.example", + ClientIDs: []string{"x"}, + } + if err := mgr.CreateOIDCProvider(context.Background(), rec); err != nil { + t.Fatalf("Create with default sink: %v", err) + } +}