Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
141 changes: 125 additions & 16 deletions kagenti-operator/internal/keycloak/audience.go
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,8 @@ type protocolMapperRep struct {
Config map[string]string `json:"config"`
}

const oidcAudienceMapper = "oidc-audience-mapper"

// AudienceScopeName derives the realm client-scope name from CLIENT_NAME (same as Python).
func AudienceScopeName(clientName string) string {
return "agent-" + strings.ReplaceAll(clientName, "/", "-") + "-aud"
Expand Down Expand Up @@ -89,9 +91,9 @@ func (a *Admin) getOrCreateAudienceClientScope(ctx context.Context, token, realm
return "", err
}
if scopeID != "" {
if err := a.ensureAudienceMapper(ctx, token, realm, scopeID, scopeName, audience); err != nil {
return "", fmt.Errorf("ensure audience mapper for existing scope %q: %w", scopeName, err)
}
// Scope already exists — do NOT touch its mappers here.
// verifyAudienceMapper handles mapper verification via GET+PUT (never POST for existing scopes).
// This matches the Python AuthBridge sidecar which only adds mappers during initial creation.
return scopeID, nil
}

Expand All @@ -110,7 +112,10 @@ func (a *Admin) getOrCreateAudienceClientScope(ctx context.Context, token, realm
return "", fmt.Errorf("create client scope %q returned empty id", scopeName)
}
if err := a.ensureAudienceMapper(ctx, token, realm, scopeID, scopeName, audience); err != nil {
return "", fmt.Errorf("ensure audience mapper for new scope %q: %w", scopeName, err)
// Mapper creation failed for a brand-new scope. This can happen if createClientScope
// hit a 409 race and another reconcile already created the mapper. Non-fatal —
// verifyAudienceMapper will repair below in this reconcile.
return scopeID, nil
}
return scopeID, nil
}
Expand Down Expand Up @@ -183,7 +188,7 @@ func (a *Admin) ensureAudienceMapper(ctx context.Context, token, realm, scopeID,
mapper := protocolMapperRep{
Name: scopeName,
Protocol: "openid-connect",
ProtocolMapper: "oidc-audience-mapper",
ProtocolMapper: oidcAudienceMapper,
ConsentRequired: false,
Config: map[string]string{
"included.custom.audience": audience,
Expand Down Expand Up @@ -254,28 +259,125 @@ func (a *Admin) listAudienceMappers(ctx context.Context, token, realm, scopeID s

// updateAudienceMapperIfNeeded fetches the existing mapper for the scope and updates
// its included.custom.audience if it differs from the desired value.
// Returns an error if no matching mapper is found — this treats "no match" as a real
// failure (e.g. Keycloak race or name mismatch) rather than silently ignoring it.
// If a mapper with the correct name exists but has the wrong ProtocolMapper type
// (e.g. corrupted state), it deletes the stale mapper and re-creates via best-effort POST.
// If no mapper is found at all (ghost 409 from Keycloak's name index), returns nil
// to allow verifyAudienceMapper to retry on the next reconcile.
func (a *Admin) updateAudienceMapperIfNeeded(ctx context.Context, token, realm, scopeID, scopeName, audience string) error {
mappers, err := a.listAudienceMappers(ctx, token, realm, scopeID)
if err != nil {
return err
}

for i := range mappers {
if mappers[i].Name != scopeName || mappers[i].ProtocolMapper != "oidc-audience-mapper" {
if mappers[i].Name != scopeName {
continue
}
if mappers[i].ProtocolMapper != oidcAudienceMapper {
if err := a.deleteMapper(ctx, token, realm, scopeID, mappers[i].ID); err != nil {
return fmt.Errorf("delete stale mapper %q (type %q) for scope %q: %w",
mappers[i].Name, mappers[i].ProtocolMapper, scopeName, err)
}
return a.createAudienceMapperBestEffort(ctx, token, realm, scopeID, scopeName, audience)
}
if mappers[i].Config == nil {
continue
}
if mappers[i].Config["included.custom.audience"] == audience {
return nil // already correct
return nil
}
mappers[i].Config["included.custom.audience"] = audience
return a.putAudienceMapper(ctx, token, realm, scopeID, mappers[i])
}
return fmt.Errorf("no matching audience mapper found for scope %q (scopeID %s)", scopeName, scopeID)
// No mapper found despite 409 — Keycloak's internal name index is stale.
// Return nil; verifyAudienceMapper will retry on next reconcile.
return nil
}

func (a *Admin) deleteMapper(ctx context.Context, token, realm, scopeID, mapperID string) error {
base := trimBaseURL(a.BaseURL)
endpoint := base + "/admin/realms/" + url.PathEscape(realm) + "/client-scopes/" + url.PathEscape(scopeID) + "/protocol-mappers/models/" + url.PathEscape(mapperID)
req, err := http.NewRequestWithContext(ctx, http.MethodDelete, endpoint, nil)
if err != nil {
return err
}
req.Header.Set("Authorization", "Bearer "+token)

resp, err := a.httpc().Do(req)
if err != nil {
return err
}
defer func() { _ = resp.Body.Close() }()
if resp.StatusCode == http.StatusNoContent || (resp.StatusCode >= 200 && resp.StatusCode < 300) {
return nil
}
body, _ := io.ReadAll(resp.Body)
return fmt.Errorf("keycloak delete mapper: status %d: %s", resp.StatusCode, truncate(body, 256))
}

// createAudienceMapperBestEffort posts a new audience mapper. On 409 (conflict), it verifies
// the mapper actually exists with the correct audience — if it does, that's fine (another
// reconcile got there first). If 409 but no mapper is visible (Keycloak ghost-conflict),
// it returns nil and the reconciler will retry on the next pass.
func (a *Admin) createAudienceMapperBestEffort(ctx context.Context, token, realm, scopeID, scopeName, audience string) error {
mapper := protocolMapperRep{
Name: scopeName,
Protocol: "openid-connect",
ProtocolMapper: oidcAudienceMapper,
ConsentRequired: false,
Config: map[string]string{
"included.custom.audience": audience,
"id.token.claim": "false",
"access.token.claim": "true",
"userinfo.token.claim": "false",
},
}
payload, err := json.Marshal(mapper)
if err != nil {
return err
}
base := trimBaseURL(a.BaseURL)
endpoint := base + "/admin/realms/" + url.PathEscape(realm) + "/client-scopes/" + url.PathEscape(scopeID) + "/protocol-mappers/models"
req, err := http.NewRequestWithContext(ctx, http.MethodPost, endpoint, bytes.NewReader(payload))
if err != nil {
return err
}
req.Header.Set("Authorization", "Bearer "+token)
req.Header.Set("Content-Type", "application/json")

resp, err := a.httpc().Do(req)
if err != nil {
return err
}
defer func() { _ = resp.Body.Close() }()
body, _ := io.ReadAll(resp.Body)
if resp.StatusCode == http.StatusCreated || (resp.StatusCode >= 200 && resp.StatusCode < 300) {
return nil
}
if resp.StatusCode == http.StatusConflict {
// 409 can mean: (a) another reconcile created it already, or (b) Keycloak ghost index.
// Verify via GET — if the mapper exists with correct audience, success.
mappers, err := a.listAudienceMappers(ctx, token, realm, scopeID)
if err != nil {
return nil
}
for i := range mappers {
if mappers[i].Name == scopeName && mappers[i].ProtocolMapper == oidcAudienceMapper {
if mappers[i].Config != nil && mappers[i].Config["included.custom.audience"] == audience {
return nil
}
// Mapper exists but wrong audience — update it.
if mappers[i].Config == nil {
mappers[i].Config = make(map[string]string)
}
mappers[i].Config["included.custom.audience"] = audience
return a.putAudienceMapper(ctx, token, realm, scopeID, mappers[i])
}
}
// Ghost-409: mapper not visible. Return nil; next reconcile will retry.
return nil
}
return fmt.Errorf("keycloak create mapper best-effort: status %d: %s", resp.StatusCode, truncate(body, 256))
}

func (a *Admin) putAudienceMapper(ctx context.Context, token, realm, scopeID string, mapper protocolMapperRep) error {
Expand Down Expand Up @@ -306,20 +408,26 @@ func (a *Admin) putAudienceMapper(ctx context.Context, token, realm, scopeID str

// verifyAudienceMapper is a defense-in-depth check that runs on every reconcile.
// It GETs the mappers for a scope and ensures the oidc-audience-mapper exists with the
// correct audience. If the mapper is missing (e.g. due to a prior transient failure),
// it re-creates it. If the audience is stale, it updates it.
// Cost: one extra GET per reconcile per audience-enabled scope; accepted tradeoff for
// catching scopes left broken by prior transient failures.
// correct audience. If the audience is stale, it PUTs an update. If a mapper with the
// correct name but wrong type exists, it deletes and re-creates. If the mapper is missing
// entirely, it attempts a POST but treats 409 as success (avoids ghost-409 cascades).
func (a *Admin) verifyAudienceMapper(ctx context.Context, token, realm, scopeID, scopeName, audience string) error {
mappers, err := a.listAudienceMappers(ctx, token, realm, scopeID)
if err != nil {
return err
}

for i := range mappers {
if mappers[i].Name != scopeName || mappers[i].ProtocolMapper != "oidc-audience-mapper" {
if mappers[i].Name != scopeName {
continue
}
if mappers[i].ProtocolMapper != oidcAudienceMapper {
if err := a.deleteMapper(ctx, token, realm, scopeID, mappers[i].ID); err != nil {
return fmt.Errorf("delete stale mapper %q (type %q): %w",
mappers[i].Name, mappers[i].ProtocolMapper, err)
}
return a.createAudienceMapperBestEffort(ctx, token, realm, scopeID, scopeName, audience)
}
if mappers[i].Config != nil && mappers[i].Config["included.custom.audience"] == audience {
return nil
}
Expand All @@ -329,7 +437,8 @@ func (a *Admin) verifyAudienceMapper(ctx context.Context, token, realm, scopeID,
mappers[i].Config["included.custom.audience"] = audience
return a.putAudienceMapper(ctx, token, realm, scopeID, mappers[i])
}
return a.ensureAudienceMapper(ctx, token, realm, scopeID, scopeName, audience)
// Mapper not found — create it. Treat 409 as success (Keycloak name-index ghost).
return a.createAudienceMapperBestEffort(ctx, token, realm, scopeID, scopeName, audience)
}

func (a *Admin) putRealmDefaultDefaultClientScope(ctx context.Context, token, realm, scopeID string) error {
Expand Down
126 changes: 113 additions & 13 deletions kagenti-operator/internal/keycloak/audience_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -154,8 +154,8 @@ func TestEnsureAudienceScope_UpdatesStaleMapper(t *testing.T) {
if err != nil {
t.Fatal(err)
}
if getMapperCalls != 2 {
t.Fatalf("expected 2 GET mapper calls (update + verify), got %d", getMapperCalls)
if getMapperCalls != 1 {
t.Fatalf("expected 1 GET mapper call (verify only, no ensureAudienceMapper for existing scopes), got %d", getMapperCalls)
}
if putMapperCalls != 1 {
t.Fatalf("expected 1 PUT mapper call, got %d", putMapperCalls)
Expand Down Expand Up @@ -231,10 +231,12 @@ func TestEnsureAudienceScope_SkipsUpdateWhenCorrect(t *testing.T) {
}
}

// TestEnsureAudienceScope_MapperFailurePropagated verifies that when the mapper POST
// returns a server error (e.g. 500), the error propagates to EnsureAudienceScope
// instead of being silently swallowed (regression test for #348).
func TestEnsureAudienceScope_MapperFailurePropagated(t *testing.T) {
// TestEnsureAudienceScope_MapperFailureForNewScope verifies that when a new scope is created
// but the initial mapper POST fails (500), verifyAudienceMapper repairs the missing mapper
// via createAudienceMapperBestEffort. The initial failure is non-fatal (matches Python sidecar
// which swallows mapper creation exceptions).
func TestEnsureAudienceScope_MapperFailureForNewScope(t *testing.T) {
var postMapperCalls int
var srv *httptest.Server
srv = httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
path := r.URL.Path
Expand All @@ -252,10 +254,24 @@ func TestEnsureAudienceScope_MapperFailurePropagated(t *testing.T) {
w.Header().Set("Location", srv.URL+"/admin/realms/kagenti/client-scopes/new-scope-id")
w.WriteHeader(http.StatusCreated)

// Mapper POST returns 500 (server error)
// Mapper POST: first call (from ensureAudienceMapper) returns 500, second call
// (from verifyAudienceMapper → createAudienceMapperBestEffort) succeeds
case strings.Contains(path, "/client-scopes/new-scope-id/protocol-mappers/models") && r.Method == http.MethodPost:
w.WriteHeader(http.StatusInternalServerError)
_, _ = w.Write([]byte(`{"error":"internal"}`))
postMapperCalls++
if postMapperCalls == 1 {
w.WriteHeader(http.StatusInternalServerError)
_, _ = w.Write([]byte(`{"error":"internal"}`))
} else {
w.WriteHeader(http.StatusCreated)
}

// GET mappers (from verifyAudienceMapper): returns empty list to trigger re-creation
case strings.Contains(path, "/client-scopes/new-scope-id/protocol-mappers/models") && r.Method == http.MethodGet:
_ = json.NewEncoder(w).Encode([]protocolMapperRep{})

// Realm default scope
case path == "/admin/realms/kagenti/default-default-client-scopes/new-scope-id" && r.Method == http.MethodPut:
w.WriteHeader(http.StatusNoContent)

default:
t.Fatalf("unexpected %s %s", r.Method, path)
Expand All @@ -275,11 +291,11 @@ func TestEnsureAudienceScope_MapperFailurePropagated(t *testing.T) {
AudienceClientID: "spiffe://example.org/ns/ns/sa/wl",
AudienceScopeEnabled: true,
})
if err == nil {
t.Fatal("expected error when mapper POST fails, got nil")
if err != nil {
t.Fatalf("expected success (verifyAudienceMapper repairs), got: %s", err)
}
if !strings.Contains(err.Error(), "ensure audience mapper") {
t.Fatalf("expected error to contain 'ensure audience mapper', got: %s", err.Error())
if postMapperCalls != 2 {
t.Fatalf("expected 2 POST mapper calls (initial fail + verify repair), got %d", postMapperCalls)
}
}

Expand Down Expand Up @@ -353,6 +369,90 @@ func TestEnsureAudienceScope_VerifyRecreatesMissingMapper(t *testing.T) {
}
}

// TestEnsureAudienceScope_DeletesCorruptedMapper verifies that when a mapper exists with
// the correct name but the wrong ProtocolMapper type (corrupted state from issue #358),
// the operator deletes the stale mapper and re-creates the correct oidc-audience-mapper.
func TestEnsureAudienceScope_DeletesCorruptedMapper(t *testing.T) {
var deleteMapperCalls, recreatePostCalls int
spiffeURI := "spiffe://example.org/ns/ns/sa/wl"

srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
path := r.URL.Path
switch {
case path == testMasterRealmTokenPath:
w.Header().Set("Content-Type", "application/json")
_ = json.NewEncoder(w).Encode(map[string]string{"access_token": "tok"})

// Scope already exists
case path == "/admin/realms/kagenti/client-scopes" && r.Method == http.MethodGet:
_ = json.NewEncoder(w).Encode([]clientScopeListItem{{ID: "scope-123", Name: "agent-ns-wl-aud"}})

// ensureAudienceMapper POST — 409 conflict (name collision with corrupted mapper)
case strings.Contains(path, "/client-scopes/scope-123/protocol-mappers/models") && r.Method == http.MethodPost:
if deleteMapperCalls > 0 {
// After deletion, re-create succeeds
recreatePostCalls++
w.WriteHeader(http.StatusCreated)
} else {
w.WriteHeader(http.StatusConflict)
}

// GET mappers — returns mapper with wrong type (corrupted)
case strings.Contains(path, "/client-scopes/scope-123/protocol-mappers/models") && r.Method == http.MethodGet:
if deleteMapperCalls > 0 {
// After delete+recreate, verify sees the correct mapper
_ = json.NewEncoder(w).Encode([]protocolMapperRep{{
ID: "mapper-new", Name: "agent-ns-wl-aud", Protocol: "openid-connect",
ProtocolMapper: "oidc-audience-mapper",
Config: map[string]string{"included.custom.audience": spiffeURI},
}})
} else {
// Corrupted: same name, wrong ProtocolMapper type
_ = json.NewEncoder(w).Encode([]protocolMapperRep{{
ID: "mapper-corrupted", Name: "agent-ns-wl-aud", Protocol: "openid-connect",
ProtocolMapper: "oidc-usermodel-attribute-mapper", // wrong type!
Config: map[string]string{"claim.name": "audience"},
}})
}

// DELETE the corrupted mapper
case strings.Contains(path, "/protocol-mappers/models/mapper-corrupted") && r.Method == http.MethodDelete:
deleteMapperCalls++
w.WriteHeader(http.StatusNoContent)

// Realm default scope
case path == "/admin/realms/kagenti/default-default-client-scopes/scope-123" && r.Method == http.MethodPut:
w.WriteHeader(http.StatusNoContent)

default:
t.Fatalf("unexpected %s %s", r.Method, path)
}
}))
defer srv.Close()

a := Admin{BaseURL: srv.URL, HTTPClient: srv.Client()}
token, err := a.PasswordGrantToken(context.Background(), "u", "p")
if err != nil {
t.Fatal(err)
}

err = a.EnsureAudienceScope(context.Background(), token, AudienceParams{
Realm: "kagenti",
ClientName: "ns/wl",
AudienceClientID: spiffeURI,
AudienceScopeEnabled: true,
})
if err != nil {
t.Fatal(err)
}
if deleteMapperCalls != 1 {
t.Fatalf("expected 1 DELETE for corrupted mapper, got %d", deleteMapperCalls)
}
if recreatePostCalls != 1 {
t.Fatalf("expected 1 POST to recreate correct mapper, got %d", recreatePostCalls)
}
}

func TestEnsureAudienceScope_Disabled(t *testing.T) {
a := Admin{}
err := a.EnsureAudienceScope(context.Background(), "t", AudienceParams{AudienceScopeEnabled: false})
Expand Down
Loading