From 1b83936b3b587563546a802bda7ad6d670277560 Mon Sep 17 00:00:00 2001 From: Christer Edvartsen Date: Thu, 10 Sep 2026 14:10:58 +0200 Subject: [PATCH 1/2] fix(kafka): align revoke grant arguments --- internal/kafka/command/flag/flag.go | 8 +- internal/kafka/command/list_grants_test.go | 4 +- internal/kafka/command/revoke_grant.go | 51 +++++--- internal/kafka/kafka.go | 52 ++++++++ internal/naisapi/gql/generated.go | 144 +++++++++++++++++++++ schema.graphql | 69 ++++++++++ 6 files changed, 305 insertions(+), 23 deletions(-) diff --git a/internal/kafka/command/flag/flag.go b/internal/kafka/command/flag/flag.go index 1af6e773..c26bcf9c 100644 --- a/internal/kafka/command/flag/flag.go +++ b/internal/kafka/command/flag/flag.go @@ -34,14 +34,14 @@ func (a *KafkaTopicGrantAccess) AutoComplete(context.Context, *naistrix.Argument return []string{"read", "write", "readwrite"}, "Available access levels." } -func (a *KafkaTopicGrantAccess) Validate() error { +func (a KafkaTopicGrantAccess) Validate() error { valid := []string{"read", "write", "readwrite"} - if a == nil || *a == "" { + if a == "" { return naistrix.Errorf("access level is required, must be one of: %s", strings.Join(valid, ", ")) } - if !slices.Contains(valid, string(*a)) { - return naistrix.Errorf("invalid access level: %q, must be one of: %s", *a, strings.Join(valid, ", ")) + if !slices.Contains(valid, string(a)) { + return naistrix.Errorf("invalid access level: %q, must be one of: %s", a, strings.Join(valid, ", ")) } return nil diff --git a/internal/kafka/command/list_grants_test.go b/internal/kafka/command/list_grants_test.go index 68e036ea..4f3823c1 100644 --- a/internal/kafka/command/list_grants_test.go +++ b/internal/kafka/command/list_grants_test.go @@ -58,8 +58,8 @@ func TestRevokeGrantCommand(t *testing.T) { if command.Name != "revoke-grant" { t.Errorf("Name = %q, want %q", command.Name, "revoke-grant") } - if len(command.Args) != 3 || command.Args[0].Name != "topic" || command.Args[1].Name != "username" || command.Args[2].Name != "access" { - t.Errorf("Args = %#v, want topic, username, and access arguments", command.Args) + if len(command.Args) != 3 || command.Args[0].Name != "username" || command.Args[1].Name != "topic" || command.Args[2].Name != "access" { + t.Errorf("Args = %#v, want username, topic, and access arguments", command.Args) } if command.ValidateFunc == nil { t.Error("ValidateFunc is nil") diff --git a/internal/kafka/command/revoke_grant.go b/internal/kafka/command/revoke_grant.go index 42dc3cd1..abb3a69d 100644 --- a/internal/kafka/command/revoke_grant.go +++ b/internal/kafka/command/revoke_grant.go @@ -18,19 +18,21 @@ func revokeGrant(parentFlags *flag.Kafka) *naistrix.Command { Title: "Revoke a user's service-user access to a Kafka topic.", Description: "Removes an ACL entry for a user on a Kafka topic with the specified access level.", Args: []naistrix.Argument{ - {Name: "topic"}, {Name: "username"}, + {Name: "topic"}, {Name: "access"}, }, - ValidateFunc: validation.RequireTeamAndEnvironment(parentFlags), + ValidateFunc: naistrix.ValidateFuncs( + validation.RequireTeamAndEnvironment(parentFlags), + func(_ context.Context, args *naistrix.Arguments) error { + return flag.KafkaTopicGrantAccess(args.Get("access")).Validate() + }, + ), AutoCompleteFunc: autoCompleteKafkaGrantArguments(parentFlags), RunFunc: func(ctx context.Context, args *naistrix.Arguments, out *naistrix.OutputWriter) error { - topicName := args.Get("topic") subject := kafkaApplicationName(args.Get("username")) + topicName := args.Get("topic") access := flag.KafkaTopicGrantAccess(args.Get("access")) - if err := access.Validate(); err != nil { - return err - } grant := gql.KafkaTopicGrantInput{ Subject: subject, TeamName: parentFlags.Team, @@ -51,12 +53,7 @@ func revokeGrant(parentFlags *flag.Kafka) *naistrix.Command { } func autoCompleteKafkaGrantArguments(flags *flag.Kafka) naistrix.AutoCompleteFunc { - topic := autoCompleteKafkaTopicName(flags, 0) - return func(ctx context.Context, args *naistrix.Arguments, toComplete string) ([]string, string) { - if args.Len() == 0 { - return topic(ctx, args, toComplete) - } if args.Len() > 2 { return nil, "" } @@ -65,12 +62,12 @@ func autoCompleteKafkaGrantArguments(flags *flag.Kafka) naistrix.AutoCompleteFun return nil, "Please provide team and environment to auto-complete Kafka grants." } - grants, err := kafka.GetKafkaTopicGrants(ctx, args.Get("topic"), flags.Team, flags.Environment) + grants, err := kafka.GetTeamKafkaTopicGrants(ctx, flags.Team, flags.Environment) if err != nil { - return nil, "Unable to fetch Kafka topic grants." + return nil, "Unable to fetch Kafka grants." } - if args.Len() == 1 { + if args.Len() == 0 { subjects := make([]string, 0, len(grants)) seen := make(map[string]struct{}) for _, grant := range grants { @@ -82,15 +79,35 @@ func autoCompleteKafkaGrantArguments(flags *flag.Kafka) naistrix.AutoCompleteFun } sort.Strings(subjects) if len(subjects) == 0 { - return nil, "No access grants found for this Kafka topic." + return nil, "No Kafka grants found in the selected environment." } - return subjects, "Select a subject with access to this Kafka topic." + return subjects, "Select a subject with a Kafka grant." } subject := kafkaApplicationName(args.Get("username")) + if args.Len() == 1 { + topics := make([]string, 0, len(grants)) + seen := make(map[string]struct{}) + for _, grant := range grants { + if grant.WorkloadName != subject { + continue + } + if _, ok := seen[grant.TopicName]; ok { + continue + } + seen[grant.TopicName] = struct{}{} + topics = append(topics, grant.TopicName) + } + sort.Strings(topics) + if len(topics) == 0 { + return nil, "No Kafka grants found for this subject." + } + return topics, "Select a Kafka topic with a grant for this subject." + } + accesses := make([]string, 0, len(grants)) for _, grant := range grants { - if grant.WorkloadName == subject { + if grant.WorkloadName == subject && grant.TopicName == args.Get("topic") { accesses = append(accesses, strings.ToLower(grant.Access)) } } diff --git a/internal/kafka/kafka.go b/internal/kafka/kafka.go index 89d14a7b..9a683237 100644 --- a/internal/kafka/kafka.go +++ b/internal/kafka/kafka.go @@ -20,6 +20,11 @@ type Grant struct { Access string `heading:"Access level" json:"access"` } +type TopicGrant struct { + TopicName string `json:"topicName"` + Grant +} + func GetTeamTopics(ctx context.Context, team string, environment string, labels []gql.LabelFilter) ([]Topic, error) { _ = `# @genqlient query GetTeamKafkaTopics($team: Slug!, $filter: KafkaTopicFilter) { @@ -188,3 +193,50 @@ func GetKafkaTopicGrants(ctx context.Context, topicName, teamSlug string, enviro return ret, nil } + +func GetTeamKafkaTopicGrants(ctx context.Context, teamSlug string, environmentName flags.Environment) ([]TopicGrant, error) { + _ = `# @genqlient + query GetTeamKafkaTopicGrants($teamSlug: Slug!, $environmentName: String!) { + team(slug: $teamSlug) { + kafkaTopics(first: 1000, filter: { environments: [$environmentName] }) { + nodes { + name + acl(first: 1000, filter: { team: $teamSlug }) { + nodes { + workloadName + teamName + access + } + } + } + } + } + } + ` + + client, err := naisapi.GraphqlClient(ctx) + if err != nil { + return nil, err + } + + resp, err := gql.GetTeamKafkaTopicGrants(ctx, client, teamSlug, string(environmentName)) + if err != nil { + return nil, err + } + + var ret []TopicGrant + for _, topic := range resp.Team.KafkaTopics.Nodes { + for _, grant := range topic.Acl.Nodes { + ret = append(ret, TopicGrant{ + TopicName: topic.Name, + Grant: Grant{ + WorkloadName: grant.WorkloadName, + TeamName: grant.TeamName, + Access: string(grant.Access), + }, + }) + } + } + + return ret, nil +} diff --git a/internal/naisapi/gql/generated.go b/internal/naisapi/gql/generated.go index 468d7a21..57d6de69 100644 --- a/internal/naisapi/gql/generated.go +++ b/internal/naisapi/gql/generated.go @@ -29455,6 +29455,91 @@ func (v *GetTeamJobsTeamJobsJobConnectionNodesJobTeamEnvironmentEnvironment) Get return v.Name } +// GetTeamKafkaTopicGrantsResponse is returned by GetTeamKafkaTopicGrants on success. +type GetTeamKafkaTopicGrantsResponse struct { + // Get a team by its slug. + Team GetTeamKafkaTopicGrantsTeam `json:"team"` +} + +// GetTeam returns GetTeamKafkaTopicGrantsResponse.Team, and is useful for accessing the field via an interface. +func (v *GetTeamKafkaTopicGrantsResponse) GetTeam() GetTeamKafkaTopicGrantsTeam { return v.Team } + +// GetTeamKafkaTopicGrantsTeam includes the requested fields of the GraphQL type Team. +// The GraphQL type's documentation follows. +// +// The team type represents a team on the [Nais platform](https://nais.io/). +// +// Learn more about what Nais teams are and what they can be used for in the [official Nais documentation](https://docs.nais.io/explanations/team/). +// +// External resources (e.g. entraIDGroupID, gitHubTeamSlug) are managed by [Nais API reconcilers](https://github.com/nais/api-reconcilers). +type GetTeamKafkaTopicGrantsTeam struct { + // Kafka topics owned by the team. + KafkaTopics GetTeamKafkaTopicGrantsTeamKafkaTopicsKafkaTopicConnection `json:"kafkaTopics"` +} + +// GetKafkaTopics returns GetTeamKafkaTopicGrantsTeam.KafkaTopics, and is useful for accessing the field via an interface. +func (v *GetTeamKafkaTopicGrantsTeam) GetKafkaTopics() GetTeamKafkaTopicGrantsTeamKafkaTopicsKafkaTopicConnection { + return v.KafkaTopics +} + +// GetTeamKafkaTopicGrantsTeamKafkaTopicsKafkaTopicConnection includes the requested fields of the GraphQL type KafkaTopicConnection. +type GetTeamKafkaTopicGrantsTeamKafkaTopicsKafkaTopicConnection struct { + Nodes []GetTeamKafkaTopicGrantsTeamKafkaTopicsKafkaTopicConnectionNodesKafkaTopic `json:"nodes"` +} + +// GetNodes returns GetTeamKafkaTopicGrantsTeamKafkaTopicsKafkaTopicConnection.Nodes, and is useful for accessing the field via an interface. +func (v *GetTeamKafkaTopicGrantsTeamKafkaTopicsKafkaTopicConnection) GetNodes() []GetTeamKafkaTopicGrantsTeamKafkaTopicsKafkaTopicConnectionNodesKafkaTopic { + return v.Nodes +} + +// GetTeamKafkaTopicGrantsTeamKafkaTopicsKafkaTopicConnectionNodesKafkaTopic includes the requested fields of the GraphQL type KafkaTopic. +type GetTeamKafkaTopicGrantsTeamKafkaTopicsKafkaTopicConnectionNodesKafkaTopic struct { + Name string `json:"name"` + Acl GetTeamKafkaTopicGrantsTeamKafkaTopicsKafkaTopicConnectionNodesKafkaTopicAclKafkaTopicAclConnection `json:"acl"` +} + +// GetName returns GetTeamKafkaTopicGrantsTeamKafkaTopicsKafkaTopicConnectionNodesKafkaTopic.Name, and is useful for accessing the field via an interface. +func (v *GetTeamKafkaTopicGrantsTeamKafkaTopicsKafkaTopicConnectionNodesKafkaTopic) GetName() string { + return v.Name +} + +// GetAcl returns GetTeamKafkaTopicGrantsTeamKafkaTopicsKafkaTopicConnectionNodesKafkaTopic.Acl, and is useful for accessing the field via an interface. +func (v *GetTeamKafkaTopicGrantsTeamKafkaTopicsKafkaTopicConnectionNodesKafkaTopic) GetAcl() GetTeamKafkaTopicGrantsTeamKafkaTopicsKafkaTopicConnectionNodesKafkaTopicAclKafkaTopicAclConnection { + return v.Acl +} + +// GetTeamKafkaTopicGrantsTeamKafkaTopicsKafkaTopicConnectionNodesKafkaTopicAclKafkaTopicAclConnection includes the requested fields of the GraphQL type KafkaTopicAclConnection. +type GetTeamKafkaTopicGrantsTeamKafkaTopicsKafkaTopicConnectionNodesKafkaTopicAclKafkaTopicAclConnection struct { + Nodes []GetTeamKafkaTopicGrantsTeamKafkaTopicsKafkaTopicConnectionNodesKafkaTopicAclKafkaTopicAclConnectionNodesKafkaTopicAcl `json:"nodes"` +} + +// GetNodes returns GetTeamKafkaTopicGrantsTeamKafkaTopicsKafkaTopicConnectionNodesKafkaTopicAclKafkaTopicAclConnection.Nodes, and is useful for accessing the field via an interface. +func (v *GetTeamKafkaTopicGrantsTeamKafkaTopicsKafkaTopicConnectionNodesKafkaTopicAclKafkaTopicAclConnection) GetNodes() []GetTeamKafkaTopicGrantsTeamKafkaTopicsKafkaTopicConnectionNodesKafkaTopicAclKafkaTopicAclConnectionNodesKafkaTopicAcl { + return v.Nodes +} + +// GetTeamKafkaTopicGrantsTeamKafkaTopicsKafkaTopicConnectionNodesKafkaTopicAclKafkaTopicAclConnectionNodesKafkaTopicAcl includes the requested fields of the GraphQL type KafkaTopicAcl. +type GetTeamKafkaTopicGrantsTeamKafkaTopicsKafkaTopicConnectionNodesKafkaTopicAclKafkaTopicAclConnectionNodesKafkaTopicAcl struct { + WorkloadName string `json:"workloadName"` + TeamName string `json:"teamName"` + Access KafkaTopicGrantAccess `json:"access"` +} + +// GetWorkloadName returns GetTeamKafkaTopicGrantsTeamKafkaTopicsKafkaTopicConnectionNodesKafkaTopicAclKafkaTopicAclConnectionNodesKafkaTopicAcl.WorkloadName, and is useful for accessing the field via an interface. +func (v *GetTeamKafkaTopicGrantsTeamKafkaTopicsKafkaTopicConnectionNodesKafkaTopicAclKafkaTopicAclConnectionNodesKafkaTopicAcl) GetWorkloadName() string { + return v.WorkloadName +} + +// GetTeamName returns GetTeamKafkaTopicGrantsTeamKafkaTopicsKafkaTopicConnectionNodesKafkaTopicAclKafkaTopicAclConnectionNodesKafkaTopicAcl.TeamName, and is useful for accessing the field via an interface. +func (v *GetTeamKafkaTopicGrantsTeamKafkaTopicsKafkaTopicConnectionNodesKafkaTopicAclKafkaTopicAclConnectionNodesKafkaTopicAcl) GetTeamName() string { + return v.TeamName +} + +// GetAccess returns GetTeamKafkaTopicGrantsTeamKafkaTopicsKafkaTopicConnectionNodesKafkaTopicAclKafkaTopicAclConnectionNodesKafkaTopicAcl.Access, and is useful for accessing the field via an interface. +func (v *GetTeamKafkaTopicGrantsTeamKafkaTopicsKafkaTopicConnectionNodesKafkaTopicAclKafkaTopicAclConnectionNodesKafkaTopicAcl) GetAccess() KafkaTopicGrantAccess { + return v.Access +} + // GetTeamKafkaTopicsResponse is returned by GetTeamKafkaTopics on success. type GetTeamKafkaTopicsResponse struct { // Get a team by its slug. @@ -34337,6 +34422,18 @@ func (v *__GetTeamJobsInput) GetOrderBy() *JobOrder { return v.OrderBy } // GetFilter returns __GetTeamJobsInput.Filter, and is useful for accessing the field via an interface. func (v *__GetTeamJobsInput) GetFilter() *TeamJobsFilter { return v.Filter } +// __GetTeamKafkaTopicGrantsInput is used internally by genqlient +type __GetTeamKafkaTopicGrantsInput struct { + TeamSlug string `json:"teamSlug"` + EnvironmentName string `json:"environmentName"` +} + +// GetTeamSlug returns __GetTeamKafkaTopicGrantsInput.TeamSlug, and is useful for accessing the field via an interface. +func (v *__GetTeamKafkaTopicGrantsInput) GetTeamSlug() string { return v.TeamSlug } + +// GetEnvironmentName returns __GetTeamKafkaTopicGrantsInput.EnvironmentName, and is useful for accessing the field via an interface. +func (v *__GetTeamKafkaTopicGrantsInput) GetEnvironmentName() string { return v.EnvironmentName } + // __GetTeamKafkaTopicsInput is used internally by genqlient type __GetTeamKafkaTopicsInput struct { Team string `json:"team"` @@ -37413,6 +37510,53 @@ func GetTeamJobs( return data_, err_ } +// The query executed by GetTeamKafkaTopicGrants. +const GetTeamKafkaTopicGrants_Operation = ` +query GetTeamKafkaTopicGrants ($teamSlug: Slug!, $environmentName: String!) { + team(slug: $teamSlug) { + kafkaTopics(first: 1000, filter: {environments:[$environmentName]}) { + nodes { + name + acl(first: 1000, filter: {team:$teamSlug}) { + nodes { + workloadName + teamName + access + } + } + } + } + } +} +` + +func GetTeamKafkaTopicGrants( + ctx_ context.Context, + client_ graphql.Client, + teamSlug string, + environmentName string, +) (data_ *GetTeamKafkaTopicGrantsResponse, err_ error) { + req_ := &graphql.Request{ + OpName: "GetTeamKafkaTopicGrants", + Query: GetTeamKafkaTopicGrants_Operation, + Variables: &__GetTeamKafkaTopicGrantsInput{ + TeamSlug: teamSlug, + EnvironmentName: environmentName, + }, + } + + data_ = &GetTeamKafkaTopicGrantsResponse{} + resp_ := &graphql.Response{Data: data_} + + err_ = client_.MakeRequest( + ctx_, + req_, + resp_, + ) + + return data_, err_ +} + // The query executed by GetTeamKafkaTopics. const GetTeamKafkaTopics_Operation = ` query GetTeamKafkaTopics ($team: Slug!, $filter: KafkaTopicFilter) { diff --git a/schema.graphql b/schema.graphql index 24695ab9..017f979f 100644 --- a/schema.graphql +++ b/schema.graphql @@ -123,6 +123,10 @@ Filter for Kafka credential creation events. """ KAFKA_CREDENTIALS_CREATED """ +Filter for Kafka topic update events. +""" + KAFKA_TOPIC_UPDATED +""" OpenSearch was created. """ OPENSEARCH_CREATED @@ -5454,6 +5458,71 @@ enum KafkaTopicOrderField { ENVIRONMENT } +type KafkaTopicUpdatedActivityLogEntry implements ActivityLogEntry & Node{ +""" +ID of the entry. +""" + id: ID! +""" +The identity of the actor who performed the action. +""" + actor: String! +""" +Creation time of the entry. +""" + createdAt: Time! +""" +Message that summarizes the entry. +""" + message: String! +""" +Type of the resource that was affected by the action. +""" + resourceType: ActivityLogEntryResourceType! +""" +Name of the resource that was affected by the action. +""" + resourceName: String! +""" +The team slug that the entry belongs to. +""" + teamSlug: Slug! +""" +The environment name that the entry belongs to. +""" + environmentName: String +""" +Data associated with the update. +""" + data: KafkaTopicUpdatedActivityLogEntryData! +} + +type KafkaTopicUpdatedActivityLogEntryData { +""" +Grants added to the Kafka topic. +""" + addedGrants: [KafkaTopicUpdatedActivityLogEntryDataGrant!]! +""" +Grants revoked from the Kafka topic. +""" + revokedGrants: [KafkaTopicUpdatedActivityLogEntryDataGrant!]! +} + +type KafkaTopicUpdatedActivityLogEntryDataGrant { +""" +Subject affected by the grant. +""" + subject: String! +""" +Team affected by the grant. +""" + teamName: String! +""" +Access level affected by the grant. +""" + access: KafkaTopicGrantAccess! +} + """ A shared facet item representing a label value distribution (key-value pairs). """ From dc35225fc146e724c6843831838f24f71f7740a5 Mon Sep 17 00:00:00 2001 From: Christer Edvartsen Date: Thu, 10 Sep 2026 14:23:20 +0200 Subject: [PATCH 2/2] fix(kafka): scope revoke access completion --- internal/kafka/command/revoke_grant.go | 71 +++++++++++++++----------- 1 file changed, 40 insertions(+), 31 deletions(-) diff --git a/internal/kafka/command/revoke_grant.go b/internal/kafka/command/revoke_grant.go index abb3a69d..7b0ab458 100644 --- a/internal/kafka/command/revoke_grant.go +++ b/internal/kafka/command/revoke_grant.go @@ -62,30 +62,31 @@ func autoCompleteKafkaGrantArguments(flags *flag.Kafka) naistrix.AutoCompleteFun return nil, "Please provide team and environment to auto-complete Kafka grants." } - grants, err := kafka.GetTeamKafkaTopicGrants(ctx, flags.Team, flags.Environment) - if err != nil { - return nil, "Unable to fetch Kafka grants." - } + switch args.Len() { + case 0, 1: + grants, err := kafka.GetTeamKafkaTopicGrants(ctx, flags.Team, flags.Environment) + if err != nil { + return nil, "Unable to fetch Kafka grants." + } - if args.Len() == 0 { - subjects := make([]string, 0, len(grants)) - seen := make(map[string]struct{}) - for _, grant := range grants { - if _, ok := seen[grant.WorkloadName]; ok { - continue + if args.Len() == 0 { + subjects := make([]string, 0, len(grants)) + seen := make(map[string]struct{}) + for _, grant := range grants { + if _, ok := seen[grant.WorkloadName]; ok { + continue + } + seen[grant.WorkloadName] = struct{}{} + subjects = append(subjects, grant.WorkloadName) } - seen[grant.WorkloadName] = struct{}{} - subjects = append(subjects, grant.WorkloadName) - } - sort.Strings(subjects) - if len(subjects) == 0 { - return nil, "No Kafka grants found in the selected environment." + sort.Strings(subjects) + if len(subjects) == 0 { + return nil, "No Kafka grants found in the selected environment." + } + return subjects, "Select a subject with a Kafka grant." } - return subjects, "Select a subject with a Kafka grant." - } - subject := kafkaApplicationName(args.Get("username")) - if args.Len() == 1 { + subject := kafkaApplicationName(args.Get("username")) topics := make([]string, 0, len(grants)) seen := make(map[string]struct{}) for _, grant := range grants { @@ -103,19 +104,27 @@ func autoCompleteKafkaGrantArguments(flags *flag.Kafka) naistrix.AutoCompleteFun return nil, "No Kafka grants found for this subject." } return topics, "Select a Kafka topic with a grant for this subject." - } + case 2: + grants, err := kafka.GetKafkaTopicGrants(ctx, args.Get("topic"), flags.Team, flags.Environment) + if err != nil { + return nil, "Unable to fetch Kafka topic grants." + } - accesses := make([]string, 0, len(grants)) - for _, grant := range grants { - if grant.WorkloadName == subject && grant.TopicName == args.Get("topic") { - accesses = append(accesses, strings.ToLower(grant.Access)) + subject := kafkaApplicationName(args.Get("username")) + accesses := make([]string, 0, len(grants)) + for _, grant := range grants { + if grant.WorkloadName == subject { + accesses = append(accesses, strings.ToLower(grant.Access)) + } + } + sort.Strings(accesses) + if len(accesses) == 0 { + return nil, "No access grants found for this subject on the Kafka topic." } - } - sort.Strings(accesses) - if len(accesses) == 0 { - return nil, "No access grants found for this subject on the Kafka topic." - } - return accesses, "Select an access level to revoke." + return accesses, "Select an access level to revoke." + default: + return nil, "" + } } }