kafds

package
v0.1.37 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 1, 2026 License: Apache-2.0 Imports: 40 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

View Source
var SHA256 scram.HashGeneratorFcn = func() hash.Hash { return sha256.New() }
View Source
var SHA512 scram.HashGeneratorFcn = func() hash.Hash { return sha512.New() }

Functions

func DoConsume

func DoConsume(ctx context.Context, topic string, consumeFlags api.ConsumeFlags, handleMessage api.MessageHandlerFunc, onError func(err any))

func DoConsumeWithConfig added in v0.1.34

func DoConsumeWithConfig(ctx context.Context, topic string, consumeFlags api.ConsumeFlags, handleMessage api.MessageHandlerFunc, onError func(err any), configProvider ConfigProviderInterface, consumer ConsumerInterface, config *ConsumeConfig)

func DoConsumeWithDeps added in v0.1.34

func DoConsumeWithDeps(ctx context.Context, topic string, consumeFlags api.ConsumeFlags, handleMessage api.MessageHandlerFunc, onError func(err any), configProvider ConfigProviderInterface, consumer ConsumerInterface)

func InitFromConfig added in v0.1.34

func InitFromConfig(cfgPath string) error

InitFromConfig reads the kaf config file at cfgPath (pass "" to use the default ~/.kaf/config) and sets the active cluster, like Init does but without the CLI overrides; use it in examples and standalone programs.

func InitTUIWriters added in v0.1.34

func InitTUIWriters()

InitTUIWriters routes the sarama Kafka client logger to the structured logger so nothing corrupts the TUI. Call this once before starting tea.NewProgram.

func PrepareOAuthDeviceFlow added in v0.1.34

func PrepareOAuthDeviceFlow(cfgPath string, w io.Writer) error

PrepareOAuthDeviceFlow runs the interactive OAuth2 device-code grant for the active cluster when configured, BEFORE the TUI redirects stdout (AA-13). It loads the kaf config at cfgPath (pass "" for the default), validates the OAUTHBEARER credential combination, and — when device flow applies and no usable/refreshable cached token exists — displays the verification URL and user code on w and caches the resulting token. It is a no-op for non-device clusters and returns a descriptive error for invalid configuration.

func SetOverrides added in v0.1.34

func SetOverrides(brokers []string, schemaRegistry, cluster string, verboseLogging bool)

SetOverrides applies CLI overrides before Init/onInit runs. Empty/nil values leave the corresponding config value untouched. Kafui's own CLI calls this.

Types

type AlterConfigCall added in v0.1.34

type AlterConfigCall struct {
	Name  string
	Key   string
	Value string
}

AlterConfigCall captures the arguments of an IncrementalAlterConfig invocation.

type AlterQuotaCall added in v0.1.34

type AlterQuotaCall struct {
	Entity []sarama.QuotaEntityComponent
	Op     sarama.ClientQuotasOp
}

AlterQuotaCall captures the arguments of an AlterClientQuotas call.

type ClusterAdminInterface added in v0.1.34

type ClusterAdminInterface interface {
	ListTopics() (map[string]sarama.TopicDetail, error)
	ListConsumerGroups() (map[string]string, error)
	DescribeConsumerGroups(groups []string) ([]*sarama.GroupDescription, error)
	ListAcls(filter sarama.AclFilter) ([]sarama.ResourceAcls, error)
	// CreateACLs creates one or more ACL bindings.
	CreateACLs(resourceACLs []*sarama.ResourceAcls) error
	// DeleteACL deletes ACLs matching the filter, returning the bindings that
	// were removed (empty when nothing matched).
	DeleteACL(filter sarama.AclFilter, validateOnly bool) ([]sarama.MatchingAcl, error)
	// DescribeClientQuotas returns the client quotas matching the components.
	DescribeClientQuotas(components []sarama.QuotaFilterComponent, strict bool) ([]sarama.DescribeClientQuotasEntry, error)
	// AlterClientQuotas applies a single set/remove op to the entity's quotas.
	AlterClientQuotas(entity []sarama.QuotaEntityComponent, op sarama.ClientQuotasOp, validateOnly bool) error
	// DescribeCluster returns the online brokers and the active controller ID.
	DescribeCluster() (brokers []*sarama.Broker, controllerID int32, err error)
	// DescribeConfig returns the config entries for a resource (e.g. a broker).
	DescribeConfig(resource sarama.ConfigResource) ([]sarama.ConfigEntry, error)
	// IncrementalAlterConfig incrementally updates config entries, preserving
	// other dynamic configs. sarama v1.45.1's ClusterAdmin implements this
	// natively, so no AlterConfig fallback is required.
	IncrementalAlterConfig(resourceType sarama.ConfigResourceType, name string, entries map[string]sarama.IncrementalAlterConfigsEntry, validateOnly bool) error
	// DescribeLogDirs returns log-directory metadata for the given broker IDs.
	DescribeLogDirs(brokers []int32) (map[int32][]sarama.DescribeLogDirsResponseDirMetadata, error)
	// DescribeTopics returns full partition metadata for the given topics.
	DescribeTopics(topics []string) ([]*sarama.TopicMetadata, error)
	// ListConsumerGroupOffsets returns the committed offsets of a group. Pass a
	// nil topicPartitions map to fetch all committed offsets for the group.
	ListConsumerGroupOffsets(group string, topicPartitions map[string][]int32) (*sarama.OffsetFetchResponse, error)
	// DeleteConsumerGroup deletes a consumer group.
	DeleteConsumerGroup(group string) error
	// DeleteConsumerGroupOffset deletes the committed offset of a single
	// group/topic/partition.
	DeleteConsumerGroupOffset(group string, topic string, partition int32) error

	// --- Topic administration (TP-1). All pass-throughs on sarama.ClusterAdmin. ---
	CreateTopic(topic string, detail *sarama.TopicDetail, validateOnly bool) error
	DeleteTopic(topic string) error
	CreatePartitions(topic string, count int32, assignment [][]int32, validateOnly bool) error
	DeleteRecords(topic string, partitionOffsets map[int32]int64) error
	AlterPartitionReassignments(topic string, assignment [][]int32) error

	Close() error
}

ClusterAdminInterface wraps the methods we actually use from sarama.ClusterAdmin

type ConfigManager added in v0.1.34

type ConfigManager interface {
	ReadConfig(configFile string) (config.Config, error)
	GetActiveCluster(cfg config.Config) *config.Cluster
}

ConfigManager interface for configuration operations

type ConfigProviderInterface added in v0.1.34

type ConfigProviderInterface interface {
	GetConsumerConfig() (*sarama.Config, error)
	GetClientFromConfig(config *sarama.Config) (sarama.Client, error)
}

ConfigProviderInterface provides configuration for consumers

type ConsumeConfig added in v0.1.34

type ConsumeConfig struct {
	OffsetFlag        string
	GroupFlag         string
	GroupCommitFlag   bool
	Follow            bool
	Tail              int32
	FlagPartitions    []int32
	LimitMessagesFlag int64

	// Typed seek model (MSG-1..4). When Seek is set it drives per-partition
	// offset resolution instead of OffsetFlag.
	Seek          api.SeekMode
	SeekOffset    *int64
	SeekTimestamp *time.Time

	// Resource controls (MSG-10). TailRatePerSec throttles follow delivery;
	// MaxBytesPerSec throttles browse fetches. Zero disables each.
	TailRatePerSec int
	MaxBytesPerSec int

	// OnEvent, when set, receives browse phase/statistics events (MSG-7).
	OnEvent func(api.BrowseEvent)
}

ConsumeConfig holds all configuration for consuming messages

func DefaultConsumeConfig added in v0.1.34

func DefaultConsumeConfig() *ConsumeConfig

DefaultConsumeConfig returns a default configuration

type ConsumerInterface added in v0.1.34

type ConsumerInterface interface {
	CreateConsumerGroupFromClient(group string, client sarama.Client) (sarama.ConsumerGroup, error)
}

ConsumerInterface creates consumer groups; replaceable for testing.

type CreatePartitionsCall added in v0.1.34

type CreatePartitionsCall struct {
	Topic      string
	Count      int32
	Assignment [][]int32
}

CreatePartitionsCall captures the arguments of a CreatePartitions invocation.

type CreateTopicCall added in v0.1.34

type CreateTopicCall struct {
	Topic        string
	Detail       *sarama.TopicDetail
	ValidateOnly bool
}

CreateTopicCall captures the arguments of a CreateTopic invocation.

type DefaultConfigManager added in v0.1.34

type DefaultConfigManager struct{}

DefaultConfigManager implements ConfigManager using real config operations

func (*DefaultConfigManager) GetActiveCluster added in v0.1.34

func (m *DefaultConfigManager) GetActiveCluster(cfg config.Config) *config.Cluster

func (*DefaultConfigManager) ReadConfig added in v0.1.34

func (m *DefaultConfigManager) ReadConfig(configFile string) (config.Config, error)

type DefaultConfigProvider added in v0.1.34

type DefaultConfigProvider struct{}

DefaultConfigProvider implements ConfigProviderInterface

func (*DefaultConfigProvider) GetClientFromConfig added in v0.1.34

func (cp *DefaultConfigProvider) GetClientFromConfig(config *sarama.Config) (sarama.Client, error)

func (*DefaultConfigProvider) GetConsumerConfig added in v0.1.34

func (cp *DefaultConfigProvider) GetConsumerConfig() (*sarama.Config, error)

type DefaultConsumer added in v0.1.34

type DefaultConsumer struct{}

DefaultConsumer implements ConsumerInterface using real Sarama

func (*DefaultConsumer) CreateConsumerGroupFromClient added in v0.1.34

func (c *DefaultConsumer) CreateConsumerGroupFromClient(group string, client sarama.Client) (sarama.ConsumerGroup, error)

type DefaultKafkaClientFactory added in v0.1.34

type DefaultKafkaClientFactory struct{}

DefaultKafkaClientFactory implements KafkaClientFactory using real Sarama clients

func (*DefaultKafkaClientFactory) CreateClient added in v0.1.34

func (f *DefaultKafkaClientFactory) CreateClient(brokers []string, config *sarama.Config) (sarama.Client, error)

func (*DefaultKafkaClientFactory) CreateClusterAdmin added in v0.1.34

func (f *DefaultKafkaClientFactory) CreateClusterAdmin(brokers []string, config *sarama.Config) (ClusterAdminInterface, error)

type DeleteOffsetCall added in v0.1.34

type DeleteOffsetCall struct {
	Group     string
	Topic     string
	Partition int32
}

DeleteOffsetCall captures the arguments of a DeleteConsumerGroupOffset call.

type DeleteRecordsCall added in v0.1.34

type DeleteRecordsCall struct {
	Topic            string
	PartitionOffsets map[int32]int64
}

DeleteRecordsCall captures the arguments of a DeleteRecords invocation.

type DescribeQuotasCall added in v0.1.34

type DescribeQuotasCall struct {
	Components []sarama.QuotaFilterComponent
	Strict     bool
}

DescribeQuotasCall captures the arguments of a DescribeClientQuotas call.

type KafkaClientFactory added in v0.1.34

type KafkaClientFactory interface {
	CreateClusterAdmin(brokers []string, config *sarama.Config) (ClusterAdminInterface, error)
	CreateClient(brokers []string, config *sarama.Config) (sarama.Client, error)
}

KafkaClientFactory interface for creating Kafka clients

type KafkaDataSourceKaf

type KafkaDataSourceKaf struct {
	// contains filtered or unexported fields
}

func NewKafkaDataSourceKaf added in v0.1.34

func NewKafkaDataSourceKaf() *KafkaDataSourceKaf

NewKafkaDataSourceKaf creates a new instance with default dependencies

func NewKafkaDataSourceKafWithDeps added in v0.1.34

func NewKafkaDataSourceKafWithDeps(clientFactory KafkaClientFactory, configManager ConfigManager) *KafkaDataSourceKaf

NewKafkaDataSourceKafWithDeps creates a new instance with custom dependencies for testing

func (KafkaDataSourceKaf) AlterBrokerConfig added in v0.1.34

func (kp KafkaDataSourceKaf) AlterBrokerConfig(brokerID int32, key, value string) error

AlterBrokerConfig implements api.KafkaDataSource using an incremental SET so other dynamic configs are preserved.

func (KafkaDataSourceKaf) AlterClientQuotas added in v0.1.34

func (kp KafkaDataSourceKaf) AlterClientQuotas(entity api.ClientQuotaEntity, quotas map[string]float64) error

AlterClientQuotas implements api.KafkaDataSource with replace semantics: the submitted map becomes the entity's complete property set. Properties present on the entity but missing from the submission are removed; an empty submission removes everything (delete).

func (KafkaDataSourceKaf) AlterReplicaLogDir added in v0.1.34

func (kp KafkaDataSourceKaf) AlterReplicaLogDir(brokerID int32, topic string, partition int32, logDir string) error

AlterReplicaLogDir implements api.KafkaDataSource.

sarama v1.45.1 does not implement AlterReplicaLogDirsRequest (only the API-key constant exists), and ClusterAdmin exposes no helper, so there is no protocol path to perform the move against a real cluster. We therefore return a typed NotSupportedError; the mock datasource implements the full flow for UI work.

func (KafkaDataSourceKaf) CancelTopicAnalysis added in v0.1.34

func (kp KafkaDataSourceKaf) CancelTopicAnalysis(topicName string) error

CancelTopicAnalysis implements api.KafkaDataSource.

func (KafkaDataSourceKaf) ChangeReplicationFactor added in v0.1.34

func (kp KafkaDataSourceKaf) ChangeReplicationFactor(name string, newFactor int16) error

ChangeReplicationFactor implements api.KafkaDataSource by computing a balanced reassignment across online brokers and applying it.

func (KafkaDataSourceKaf) CheckSchemaCompatibility added in v0.1.34

func (kp KafkaDataSourceKaf) CheckSchemaCompatibility(subject, schemaText, schemaType string) (bool, []string, error)

CheckSchemaCompatibility tests a candidate schema against the subject's latest version without registering it (SR-8).

func (KafkaDataSourceKaf) ConsumeTopic

func (kp KafkaDataSourceKaf) ConsumeTopic(ctx context.Context, topicName string, flags api.ConsumeFlags, handleMessage api.MessageHandlerFunc, onError func(err any)) error

func (KafkaDataSourceKaf) CreateACL added in v0.1.34

func (kp KafkaDataSourceKaf) CreateACL(entry api.ACLEntry) error

CreateACL implements api.KafkaDataSource. It validates and maps the binding to sarama enums, then calls admin.CreateACLs.

func (KafkaDataSourceKaf) CreateConnector added in v0.1.34

func (kp KafkaDataSourceKaf) CreateConnector(connect, name string, config map[string]string) (api.Connector, error)

func (KafkaDataSourceKaf) CreateTopic added in v0.1.34

func (kp KafkaDataSourceKaf) CreateTopic(name string, numPartitions int32, replicationFactor int16, configs map[string]*string) error

CreateTopic implements api.KafkaDataSource.

func (KafkaDataSourceKaf) DecodeMessage added in v0.1.34

func (kp KafkaDataSourceKaf) DecodeMessage(_ context.Context, msg api.Message) (api.Message, error)

DecodeMessage decodes Avro-encoded raw bytes stored in msg.RawKey / msg.RawValue into human-readable strings. Messages without raw bytes are returned unchanged. The schema registry client is shared across calls (see cachedSchemaCache).

func (KafkaDataSourceKaf) DeleteACL added in v0.1.34

func (kp KafkaDataSourceKaf) DeleteACL(entry api.ACLEntry) error

DeleteACL implements api.KafkaDataSource. It builds an exact-match filter from the full binding and returns an ACLNotFoundError when nothing matched.

func (KafkaDataSourceKaf) DeleteConnector added in v0.1.34

func (kp KafkaDataSourceKaf) DeleteConnector(connect, name string) error

func (KafkaDataSourceKaf) DeleteConsumerGroup added in v0.1.34

func (kp KafkaDataSourceKaf) DeleteConsumerGroup(groupID string) error

DeleteConsumerGroup implements api.KafkaDataSource (CG-6).

func (KafkaDataSourceKaf) DeleteConsumerGroupOffsets added in v0.1.34

func (kp KafkaDataSourceKaf) DeleteConsumerGroupOffsets(groupID string, topic string) error

DeleteConsumerGroupOffsets implements api.KafkaDataSource (CG-6). It deletes only the named topic's committed offsets, leaving other topics intact.

func (KafkaDataSourceKaf) DeleteSchemaVersion added in v0.1.34

func (kp KafkaDataSourceKaf) DeleteSchemaVersion(subject string, version int, permanent bool) error

DeleteSchemaVersion deletes a single version of a subject. version=-1 targets the registry keyword "latest". permanent=true performs a hard delete (SR-9).

func (KafkaDataSourceKaf) DeleteSubject added in v0.1.34

func (kp KafkaDataSourceKaf) DeleteSubject(subject string, permanent bool) ([]int, error)

DeleteSubject deletes all versions of a subject, returning the deleted version numbers. permanent=true performs a hard delete (requires a prior soft delete; the registry's 40405 error is surfaced) (SR-9).

func (KafkaDataSourceKaf) DeleteTopic added in v0.1.34

func (kp KafkaDataSourceKaf) DeleteTopic(name string) error

DeleteTopic implements api.KafkaDataSource.

func (KafkaDataSourceKaf) ExecuteKsql added in v0.1.34

func (kp KafkaDataSourceKaf) ExecuteKsql(ctx context.Context, sql string, props map[string]string) (<-chan api.KsqlResultTable, error)

ExecuteKsql validates and executes a single ksqlDB statement, delivering all outcomes on the returned channel. Validation (KS-6) runs first; a failure is emitted as a single error table and the channel is closed. Statement-kind input is posted to /ksql and its result tables emitted; SELECT input opens the /query stream and emits a schema table followed by one table per data row. Cancelling ctx closes the underlying HTTP body (terminating the server-side query) and closes the channel. A KsqlNotConfiguredError is returned (with a nil channel) when no endpoint is configured.

func (KafkaDataSourceKaf) GetACLs added in v0.1.34

func (kp KafkaDataSourceKaf) GetACLs() ([]api.ACLEntry, error)

GetACLs implements api.KafkaDataSource (the match-any case of GetACLsFiltered).

func (KafkaDataSourceKaf) GetACLsFiltered added in v0.1.34

func (kp KafkaDataSourceKaf) GetACLsFiltered(filter api.ACLFilter) ([]api.ACLEntry, error)

GetACLsFiltered implements api.KafkaDataSource. Empty filter fields match any value. Results are stably sorted by principal -> resourceType -> resourceName.

func (KafkaDataSourceKaf) GetBrokerConfig added in v0.1.34

func (kp KafkaDataSourceKaf) GetBrokerConfig(brokerID int32) ([]api.BrokerConfigEntry, error)

GetBrokerConfig implements api.KafkaDataSource.

func (KafkaDataSourceKaf) GetBrokerLogDirs added in v0.1.34

func (kp KafkaDataSourceKaf) GetBrokerLogDirs(brokerIDs []int32) (map[int32][]api.BrokerLogDir, error)

GetBrokerLogDirs implements api.KafkaDataSource.

func (KafkaDataSourceKaf) GetBrokerMetrics added in v0.1.34

func (kp KafkaDataSourceKaf) GetBrokerMetrics(brokerID int32) (string, error)

GetBrokerMetrics implements api.KafkaDataSource. Per-broker metrics require the metrics-collection pipeline (feature 12), which does not exist yet.

func (KafkaDataSourceKaf) GetBrokerStats added in v0.1.34

func (kp KafkaDataSourceKaf) GetBrokerStats() (map[int32]api.BrokerStats, api.BrokerSummary, error)

GetBrokerStats implements api.KafkaDataSource. It combines partition distribution (from topic metadata) with disk usage (from log dirs).

func (KafkaDataSourceKaf) GetBrokers added in v0.1.34

func (kp KafkaDataSourceKaf) GetBrokers() ([]api.BrokerInfo, error)

GetBrokers implements api.KafkaDataSource. It describes the cluster and maps each online broker to api.BrokerInfo, marking the active controller.

func (KafkaDataSourceKaf) GetClientQuotas added in v0.1.34

func (kp KafkaDataSourceKaf) GetClientQuotas() ([]api.ClientQuotaEntry, error)

GetClientQuotas implements api.KafkaDataSource. It lists every configured client quota and orders the result by user -> client-id -> ip, with absent identifiers sorted last.

func (KafkaDataSourceKaf) GetClusterCapabilities added in v0.1.34

func (kp KafkaDataSourceKaf) GetClusterCapabilities(_ context.Context, clusterName string) ([]api.Capability, error)

GetClusterCapabilities implements api.KafkaDataSource. Capabilities are derived from configuration (schema registry) plus a best-effort ACL probe on the active cluster. A probe failure removes the capability but never fails the whole call.

func (KafkaDataSourceKaf) GetClusterDetails added in v0.1.34

func (kp KafkaDataSourceKaf) GetClusterDetails(clusterName string) (api.ClusterInfo, error)

GetClusterDetails returns configuration details for the named cluster.

func (KafkaDataSourceKaf) GetClusterStatistics added in v0.1.34

func (kp KafkaDataSourceKaf) GetClusterStatistics(_ context.Context, clusterName string) (api.ClusterStatistics, error)

GetClusterStatistics implements api.KafkaDataSource by reusing the broker statistics collector (DescribeCluster + metadata + DescribeLogDirs) and folding its per-broker/summary view into a cluster-level snapshot.

ponytail: KRaft/ZooKeeper quorum detection is not wrapped by ClusterAdminInterface, so CoordinationType is best-effort "unknown". Byte throughput isn't available from Sarama admin APIs (metrics feature owns it).

Only the active cluster has a connection, so for any other clusterName it returns api.NotSupportedError instead of reporting (and paying for) the active cluster's statistics under the other cluster's name.

func (KafkaDataSourceKaf) GetConnectClusters added in v0.1.34

func (kp KafkaDataSourceKaf) GetConnectClusters(withStats bool) ([]api.ConnectCluster, error)

func (KafkaDataSourceKaf) GetConnectorDetails added in v0.1.34

func (kp KafkaDataSourceKaf) GetConnectorDetails(connect, name string) (api.ConnectorDetails, error)

func (KafkaDataSourceKaf) GetConnectorNames added in v0.1.34

func (kp KafkaDataSourceKaf) GetConnectorNames(connect string) ([]string, error)

func (KafkaDataSourceKaf) GetConnectorPlugins added in v0.1.34

func (kp KafkaDataSourceKaf) GetConnectorPlugins(connect string) ([]api.ConnectorPlugin, error)

func (KafkaDataSourceKaf) GetConnectors added in v0.1.34

func (kp KafkaDataSourceKaf) GetConnectors() ([]api.Connector, error)

func (KafkaDataSourceKaf) GetConsumerGroupDetail added in v0.1.34

func (kp KafkaDataSourceKaf) GetConsumerGroupDetail(groupID string) (api.ConsumerGroupDetail, error)

GetConsumerGroupDetail implements api.KafkaDataSource (CG-3).

func (KafkaDataSourceKaf) GetConsumerGroupDetails added in v0.1.34

func (kp KafkaDataSourceKaf) GetConsumerGroupDetails(groupIDs []string) ([]api.ConsumerGroup, error)

GetConsumerGroupDetails implements api.KafkaDataSource (CG-4).

func (KafkaDataSourceKaf) GetConsumerGroups

func (kp KafkaDataSourceKaf) GetConsumerGroups() ([]api.ConsumerGroup, error)

func (KafkaDataSourceKaf) GetConsumerGroupsForTopic added in v0.1.34

func (kp KafkaDataSourceKaf) GetConsumerGroupsForTopic(topic string) ([]api.ConsumerGroup, error)

GetConsumerGroupsForTopic implements api.KafkaDataSource (CG-5).

func (KafkaDataSourceKaf) GetContext

func (kp KafkaDataSourceKaf) GetContext() string

func (KafkaDataSourceKaf) GetContexts

func (kp KafkaDataSourceKaf) GetContexts() ([]string, error)

GetContexts retrieves a list of Kafka contexts

func (KafkaDataSourceKaf) GetGlobalCompatibility added in v0.1.34

func (kp KafkaDataSourceKaf) GetGlobalCompatibility() (api.CompatibilityLevel, error)

GetGlobalCompatibility returns the registry's global compatibility level (SR-5).

func (KafkaDataSourceKaf) GetMessageSchemaInfo added in v0.1.34

func (kp KafkaDataSourceKaf) GetMessageSchemaInfo(keySchemaID, valueSchemaID string) (*api.MessageSchemaInfo, error)

GetMessageSchemaInfo implements api.KafkaDataSource

func (KafkaDataSourceKaf) GetSchemaContent added in v0.1.34

func (kp KafkaDataSourceKaf) GetSchemaContent(subject string, version int) (string, error)

GetSchemaContent fetches the full schema definition string for the given subject and version. Pass version=0 (or any non-positive value) to fetch the latest version.

func (KafkaDataSourceKaf) GetSchemaDetails added in v0.1.34

func (kp KafkaDataSourceKaf) GetSchemaDetails(subjects []string) ([]api.Schema, error)

GetSchemaDetails fetches the latest version metadata (version, id, schemaType) plus the effective compatibility level for the given subjects using a 20-worker concurrent pool (SR-6).

func (KafkaDataSourceKaf) GetSchemaVersions added in v0.1.34

func (kp KafkaDataSourceKaf) GetSchemaVersions(subject string) ([]api.SchemaVersion, error)

GetSchemaVersions lists all versions of a subject with per-version metadata (id, type), leaving the Schema text empty (SR-4). Versions are returned in ascending order.

func (KafkaDataSourceKaf) GetSchemas added in v0.1.34

func (kp KafkaDataSourceKaf) GetSchemas() ([]api.Schema, error)

GetSchemas returns all registered schema subjects. It performs only a single HTTP request (GET /subjects) so it completes quickly even for large registries. Call GetSchemaDetails to lazily load version/ID/type for a subset of subjects.

func (KafkaDataSourceKaf) GetSubjectCompatibility added in v0.1.34

func (kp KafkaDataSourceKaf) GetSubjectCompatibility(subject string) (api.CompatibilityLevel, bool, error)

GetSubjectCompatibility returns a subject's effective compatibility level, falling back to the global level (isSubjectSpecific=false) when the subject has no own setting (SR-5).

func (KafkaDataSourceKaf) GetTopicAnalysis added in v0.1.34

func (kp KafkaDataSourceKaf) GetTopicAnalysis(topicName string) (*api.TopicAnalysis, error)

GetTopicAnalysis implements api.KafkaDataSource.

func (KafkaDataSourceKaf) GetTopicConfig added in v0.1.34

func (kp KafkaDataSourceKaf) GetTopicConfig(topicName string) ([]api.TopicConfigEntry, error)

GetTopicConfig implements api.KafkaDataSource. On an authorization failure it returns an empty slice (not an error), per spec.

ponytail: sarama's ClusterAdmin.DescribeConfig does not set IncludeSynonyms, so on real clusters config synonyms (and thus derived defaults) may be empty. Default derivation below is best-effort and fully exercised by the mock admin, which can populate synonyms. Wiring a synonym-aware describe would need a new admin method beyond the pass-through interface.

func (KafkaDataSourceKaf) GetTopicDetails added in v0.1.34

func (kp KafkaDataSourceKaf) GetTopicDetails(topicName string) (api.TopicDetails, error)

GetTopicDetails implements api.KafkaDataSource. Offsets are best effort: a partition whose offsets could not be read reports 0..0. Callers that act on offsets (PurgeTopicMessages) fetch them strictly instead.

func (KafkaDataSourceKaf) GetTopicHealth added in v0.1.36

func (kp KafkaDataSourceKaf) GetTopicHealth(topicNames []string) (map[string]api.TopicHealth, error)

GetTopicHealth implements api.KafkaDataSource. One DescribeTopics covers the whole batch: the topics list used to call GetTopicDetails per topic, which opened a fresh cluster admin AND a fresh client per topic and then fetched offsets sequentially per partition — none of which the OSR column uses. On a remote cluster that was the difference between one round trip and hundreds.

func (KafkaDataSourceKaf) GetTopicMessageCounts added in v0.1.34

func (kp KafkaDataSourceKaf) GetTopicMessageCounts(topics map[string]int32) (map[string]int64, error)

GetTopicMessageCounts fetches the approximate message count for each topic by summing (newestOffset - oldestOffset) across all partitions. The offsets of all topics are fetched with one request per leader broker (see fetchOffsetsBatched) on the shared client. Partitions that fail individually are skipped so a partial result is always returned.

func (KafkaDataSourceKaf) GetTopicNames added in v0.1.34

func (kp KafkaDataSourceKaf) GetTopicNames() ([]string, error)

GetTopicNames returns only the topic names using a lightweight Sarama client metadata request. This is faster than GetTopics() because it skips the per-topic DescribeConfigs that ListTopics() adds.

func (KafkaDataSourceKaf) GetTopicSizes added in v0.1.34

func (kp KafkaDataSourceKaf) GetTopicSizes(topicNames []string) (map[string]int64, error)

GetTopicSizes implements api.KafkaDataSource. Sizes count leader replicas only so replicated bytes are not double-counted. Best-effort: topics absent from metadata are omitted. DescribeLogDirs is bounded by describeLogDirsWithTimeout (TP-4/BUG-4) so a broker that doesn't support/answer it (e.g. some managed Kafka offerings) degrades to empty sizes instead of hanging the caller — which otherwise blocks the same tea.Cmd that resolves the OSR column, making both spin forever in the topics table.

Only the brokers leading a partition of the requested topics are asked for their log dirs, since sizes count leader replicas only. Sizing one topic no longer describes every broker in the cluster.

func (KafkaDataSourceKaf) GetTopics

func (kp KafkaDataSourceKaf) GetTopics() (map[string]api.Topic, error)

GetTopics retrieves a list of Kafka topics. It uses ListTopics, which also sends one DescribeConfigs for every topic: the UI reads Topic.ConfigEntries (the "N configs" column of the topic list, and the topic page's config view and cleanup.policy check), so the configs cannot be dropped here. Callers that only need names should use GetTopicNames.

func (KafkaDataSourceKaf) IncreasePartitions added in v0.1.34

func (kp KafkaDataSourceKaf) IncreasePartitions(name string, totalCount int32) error

IncreasePartitions implements api.KafkaDataSource. It rejects a decrease or a no-op before touching the broker.

func (*KafkaDataSourceKaf) Init

func (kp *KafkaDataSourceKaf) Init(cfgOption string)

func (KafkaDataSourceKaf) IsTopicDeletionEnabled added in v0.1.34

func (kp KafkaDataSourceKaf) IsTopicDeletionEnabled() (bool, error)

IsTopicDeletionEnabled implements api.KafkaDataSource. It reads the controller broker's delete.topic.enable config; missing/unparseable defaults to true.

func (KafkaDataSourceKaf) ListKsqlStreams added in v0.1.34

func (kp KafkaDataSourceKaf) ListKsqlStreams() ([]api.KsqlStream, error)

ListKsqlStreams posts LIST STREAMS; to /ksql and maps the response. A response that is not a recognizable streams listing yields a descriptive error.

func (KafkaDataSourceKaf) ListKsqlTables added in v0.1.34

func (kp KafkaDataSourceKaf) ListKsqlTables() ([]api.KsqlTable, error)

ListKsqlTables posts LIST TABLES; to /ksql and maps the response, including the windowed flag. A response that is not a recognizable tables listing yields a descriptive error.

func (KafkaDataSourceKaf) ListSerdes added in v0.1.34

func (kp KafkaDataSourceKaf) ListSerdes() []string

ListSerdes returns the names of serdes available for decoding, driven by the active cluster's registry (built-ins + configured). (MSG-18)

func (KafkaDataSourceKaf) PauseConnector added in v0.1.34

func (kp KafkaDataSourceKaf) PauseConnector(connect, name string) error

func (KafkaDataSourceKaf) ProduceMessage added in v0.1.34

func (kp KafkaDataSourceKaf) ProduceMessage(ctx context.Context, topic string, rec api.ProduceRecord) error

ProduceMessage implements api.KafkaDataSource (MSG-30).

func (KafkaDataSourceKaf) PurgeTopicMessages added in v0.1.34

func (kp KafkaDataSourceKaf) PurgeTopicMessages(name string, partition int32) error

PurgeTopicMessages implements api.KafkaDataSource. partition == -1 purges all partitions to the high-watermark.

func (KafkaDataSourceKaf) RecreateTopic added in v0.1.34

func (kp KafkaDataSourceKaf) RecreateTopic(name string) error

RecreateTopic implements api.KafkaDataSource: snapshot, delete, then recreate with the same partition count / replication factor / non-default configs, retrying while the prior instance is still propagating its deletion.

func (KafkaDataSourceKaf) RegisterSchema added in v0.1.34

func (kp KafkaDataSourceKaf) RegisterSchema(subject, schemaText, schemaType string) (api.Schema, error)

RegisterSchema registers a new schema (new subject or new version) and returns the stored record re-fetched from versions/latest (SR-7).

func (*KafkaDataSourceKaf) Reload added in v0.1.34

func (kp *KafkaDataSourceKaf) Reload(effective appconfig.Config) error

Reload rebuilds the in-memory kaf config and active cluster from the effective kafui configuration, merging fully-kafui-defined clusters into the loaded cluster list (replacing by name or appending), then invalidates caches (mirroring SetContext). It NEVER reads or writes ~/.kaf/config — the merge is entirely in memory. Called after an in-UI config apply to take effect without restarting the process.

func (KafkaDataSourceKaf) ResetConnectorOffsets added in v0.1.34

func (kp KafkaDataSourceKaf) ResetConnectorOffsets(connect, name string) error

func (KafkaDataSourceKaf) ResetConsumerGroupOffsets added in v0.1.34

func (kp KafkaDataSourceKaf) ResetConsumerGroupOffsets(ctx context.Context, req api.OffsetResetRequest) error

ResetConsumerGroupOffsets implements api.KafkaDataSource (CG-7, CG-8).

func (KafkaDataSourceKaf) RestartConnector added in v0.1.34

func (kp KafkaDataSourceKaf) RestartConnector(connect, name string) error

func (KafkaDataSourceKaf) RestartConnectorTask added in v0.1.34

func (kp KafkaDataSourceKaf) RestartConnectorTask(connect, name string, taskID int) error

func (KafkaDataSourceKaf) ResumeConnector added in v0.1.34

func (kp KafkaDataSourceKaf) ResumeConnector(connect, name string) error

func (KafkaDataSourceKaf) SetContext

func (kp KafkaDataSourceKaf) SetContext(contextName string) error

func (KafkaDataSourceKaf) SetGlobalCompatibility added in v0.1.34

func (kp KafkaDataSourceKaf) SetGlobalCompatibility(level api.CompatibilityLevel) error

SetGlobalCompatibility sets the registry's global compatibility level (SR-10).

func (KafkaDataSourceKaf) SetSubjectCompatibility added in v0.1.34

func (kp KafkaDataSourceKaf) SetSubjectCompatibility(subject string, level api.CompatibilityLevel) error

SetSubjectCompatibility sets a subject's compatibility level (SR-10).

func (KafkaDataSourceKaf) StartTopicAnalysis added in v0.1.34

func (kp KafkaDataSourceKaf) StartTopicAnalysis(ctx context.Context, topicName string) error

StartTopicAnalysis implements api.KafkaDataSource. It validates the topic exists (TopicNotFoundError) and captures the current message total for the progress percentage before starting the background scan.

func (KafkaDataSourceKaf) StopConnector added in v0.1.34

func (kp KafkaDataSourceKaf) StopConnector(connect, name string) error

func (KafkaDataSourceKaf) UpdateConnectorConfig added in v0.1.34

func (kp KafkaDataSourceKaf) UpdateConnectorConfig(connect, name string, config map[string]string) (api.Connector, error)

func (KafkaDataSourceKaf) UpdateTopicConfig added in v0.1.34

func (kp KafkaDataSourceKaf) UpdateTopicConfig(name string, entries map[string]*string) error

UpdateTopicConfig implements api.KafkaDataSource via an incremental alter so unrelated dynamic configs are preserved. A nil value deletes the key.

func (KafkaDataSourceKaf) ValidateCandidate added in v0.1.34

func (kp KafkaDataSourceKaf) ValidateCandidate(ctx context.Context, candidate appconfig.Config) api.ValidationReport

ValidateCandidate probes every cluster in a candidate configuration without persisting anything. For each cluster it independently checks the broker connection, the schema registry (when configured) and each configured connect/ksql/metrics endpoint. TLS material is opened and parsed first; a load failure is reported as that cluster's error without attempting a connection. An empty candidate returns an empty report.

func (KafkaDataSourceKaf) ValidateClusterConnection added in v0.1.34

func (kp KafkaDataSourceKaf) ValidateClusterConnection(ctx context.Context, clusterName string) ([]api.ValidationResult, error)

ValidateClusterConnection implements api.KafkaDataSource. It builds a probe-ready extension for the named cluster (empty ⇒ the active cluster) by overlaying the kafui entry on the live kaf cluster, then delegates to the shared connectivity-validation service (AC-11). Nothing is persisted.

func (KafkaDataSourceKaf) ValidateConnectorConfig added in v0.1.34

func (kp KafkaDataSourceKaf) ValidateConnectorConfig(connect, pluginClass string, config map[string]string) (api.ConnectorValidationResult, error)

type MockClusterAdmin added in v0.1.34

type MockClusterAdmin struct {
	ShouldFailListTopics         bool
	ShouldFailListConsumerGroups bool
	ShouldFailDescribeGroups     bool
	MockTopics                   map[string]sarama.TopicDetail
	MockConsumerGroups           map[string]string
	MockGroupDescriptions        []*sarama.GroupDescription
	// ListTopicsCalls counts ListTopics invocations.
	ListTopicsCalls int

	// Broker-management fields (BR-2..BR-7).
	MockBrokers               []*sarama.Broker
	MockControllerID          int32
	ShouldFailDescribeCluster bool
	MockConfigEntries         []sarama.ConfigEntry
	ShouldFailDescribeConfig  bool
	MockLogDirs               map[int32][]sarama.DescribeLogDirsResponseDirMetadata
	ShouldFailDescribeLogDirs bool
	// DescribeLogDirsCalls records the broker IDs of each DescribeLogDirs call.
	DescribeLogDirsCalls     [][]int32
	MockTopicMetadata        []*sarama.TopicMetadata
	ShouldFailDescribeTopics bool
	// AlterConfigErr, when set, is returned from IncrementalAlterConfig.
	AlterConfigErr error
	// IncrementalAlterConfigCalls records (name, key, value) of each SET.
	IncrementalAlterConfigCalls []AlterConfigCall

	// Consumer-group offset/mutation fields (CG-2..CG-8).
	MockGroupOffsets              *sarama.OffsetFetchResponse
	ShouldFailListGroupOffsets    bool
	DeleteConsumerGroupErr        error
	DeleteConsumerGroupCalls      []string
	DeleteConsumerGroupOffsetErr  error
	DeleteConsumerGroupOffsetCall []DeleteOffsetCall

	// Topic-administration fields (TP-1..TP-11).
	CreateTopicErr         error
	CreateTopicCalls       []CreateTopicCall
	DeleteTopicErr         error
	DeleteTopicCalls       []string
	CreatePartitionsErr    error
	CreatePartitionsCalls  []CreatePartitionsCall
	DeleteRecordsErr       error
	DeleteRecordsCalls     []DeleteRecordsCall
	AlterReassignmentErr   error
	AlterReassignmentCalls []ReassignmentCall

	// ACL write fields (AQ-5).
	MockAcls         []sarama.ResourceAcls
	ListAclsErr      error
	CreateACLsErr    error
	CreateACLsCalls  [][]*sarama.ResourceAcls
	DeleteACLErr     error
	DeleteACLCalls   []sarama.AclFilter
	MockMatchingAcls []sarama.MatchingAcl

	// Client-quota fields (AQ-11).
	MockQuotas             []sarama.DescribeClientQuotasEntry
	DescribeQuotasErr      error
	DescribeQuotasCalls    []DescribeQuotasCall
	AlterQuotasErr         error
	AlterClientQuotasCalls []AlterQuotaCall
}

MockClusterAdmin for testing - implements ClusterAdminInterface

func (*MockClusterAdmin) AlterClientQuotas added in v0.1.34

func (m *MockClusterAdmin) AlterClientQuotas(entity []sarama.QuotaEntityComponent, op sarama.ClientQuotasOp, validateOnly bool) error

func (*MockClusterAdmin) AlterPartitionReassignments added in v0.1.34

func (m *MockClusterAdmin) AlterPartitionReassignments(topic string, assignment [][]int32) error

func (*MockClusterAdmin) Close added in v0.1.34

func (m *MockClusterAdmin) Close() error

func (*MockClusterAdmin) CreateACLs added in v0.1.34

func (m *MockClusterAdmin) CreateACLs(resourceACLs []*sarama.ResourceAcls) error

func (*MockClusterAdmin) CreatePartitions added in v0.1.34

func (m *MockClusterAdmin) CreatePartitions(topic string, count int32, assignment [][]int32, validateOnly bool) error

func (*MockClusterAdmin) CreateTopic added in v0.1.34

func (m *MockClusterAdmin) CreateTopic(topic string, detail *sarama.TopicDetail, validateOnly bool) error

func (*MockClusterAdmin) DeleteACL added in v0.1.34

func (m *MockClusterAdmin) DeleteACL(filter sarama.AclFilter, validateOnly bool) ([]sarama.MatchingAcl, error)

func (*MockClusterAdmin) DeleteConsumerGroup added in v0.1.34

func (m *MockClusterAdmin) DeleteConsumerGroup(group string) error

func (*MockClusterAdmin) DeleteConsumerGroupOffset added in v0.1.34

func (m *MockClusterAdmin) DeleteConsumerGroupOffset(group string, topic string, partition int32) error

func (*MockClusterAdmin) DeleteRecords added in v0.1.34

func (m *MockClusterAdmin) DeleteRecords(topic string, partitionOffsets map[int32]int64) error

func (*MockClusterAdmin) DeleteTopic added in v0.1.34

func (m *MockClusterAdmin) DeleteTopic(topic string) error

func (*MockClusterAdmin) DescribeClientQuotas added in v0.1.34

func (m *MockClusterAdmin) DescribeClientQuotas(components []sarama.QuotaFilterComponent, strict bool) ([]sarama.DescribeClientQuotasEntry, error)

func (*MockClusterAdmin) DescribeCluster added in v0.1.34

func (m *MockClusterAdmin) DescribeCluster() ([]*sarama.Broker, int32, error)

func (*MockClusterAdmin) DescribeConfig added in v0.1.34

func (m *MockClusterAdmin) DescribeConfig(resource sarama.ConfigResource) ([]sarama.ConfigEntry, error)

func (*MockClusterAdmin) DescribeConsumerGroups added in v0.1.34

func (m *MockClusterAdmin) DescribeConsumerGroups(groups []string) ([]*sarama.GroupDescription, error)

func (*MockClusterAdmin) DescribeLogDirs added in v0.1.34

func (m *MockClusterAdmin) DescribeLogDirs(brokers []int32) (map[int32][]sarama.DescribeLogDirsResponseDirMetadata, error)

func (*MockClusterAdmin) DescribeTopics added in v0.1.34

func (m *MockClusterAdmin) DescribeTopics(topics []string) ([]*sarama.TopicMetadata, error)

func (*MockClusterAdmin) IncrementalAlterConfig added in v0.1.34

func (m *MockClusterAdmin) IncrementalAlterConfig(resourceType sarama.ConfigResourceType, name string, entries map[string]sarama.IncrementalAlterConfigsEntry, validateOnly bool) error

func (*MockClusterAdmin) ListAcls added in v0.1.34

func (m *MockClusterAdmin) ListAcls(filter sarama.AclFilter) ([]sarama.ResourceAcls, error)

func (*MockClusterAdmin) ListConsumerGroupOffsets added in v0.1.34

func (m *MockClusterAdmin) ListConsumerGroupOffsets(group string, topicPartitions map[string][]int32) (*sarama.OffsetFetchResponse, error)

func (*MockClusterAdmin) ListConsumerGroups added in v0.1.34

func (m *MockClusterAdmin) ListConsumerGroups() (map[string]string, error)

func (*MockClusterAdmin) ListTopics added in v0.1.34

func (m *MockClusterAdmin) ListTopics() (map[string]sarama.TopicDetail, error)

type MockConfigManager added in v0.1.34

type MockConfigManager struct {
	ShouldFailReadConfig bool
	MockConfig           config.Config
	MockActiveCluster    *config.Cluster
	ReadConfigCallCount  int
}

MockConfigManager for testing

func (*MockConfigManager) GetActiveCluster added in v0.1.34

func (m *MockConfigManager) GetActiveCluster(cfg config.Config) *config.Cluster

func (*MockConfigManager) ReadConfig added in v0.1.34

func (m *MockConfigManager) ReadConfig(configFile string) (config.Config, error)

type MockKafkaClientFactory added in v0.1.34

type MockKafkaClientFactory struct {
	ShouldFailClusterAdmin bool
	ShouldFailClient       bool
	MockClusterAdmin       ClusterAdminInterface
	MockClient             sarama.Client
	// CreateClusterAdminCalls counts CreateClusterAdmin invocations.
	CreateClusterAdminCalls int
}

MockKafkaClientFactory for testing

func (*MockKafkaClientFactory) CreateClient added in v0.1.34

func (m *MockKafkaClientFactory) CreateClient(brokers []string, config *sarama.Config) (sarama.Client, error)

func (*MockKafkaClientFactory) CreateClusterAdmin added in v0.1.34

func (m *MockKafkaClientFactory) CreateClusterAdmin(brokers []string, config *sarama.Config) (ClusterAdminInterface, error)

type ReassignmentCall added in v0.1.34

type ReassignmentCall struct {
	Topic      string
	Assignment [][]int32
}

ReassignmentCall captures the arguments of an AlterPartitionReassignments call.

type XDGSCRAMClient

type XDGSCRAMClient struct {
	*scram.Client
	*scram.ClientConversation
	scram.HashGeneratorFcn
}

func (*XDGSCRAMClient) Begin

func (x *XDGSCRAMClient) Begin(userName, password, authzID string) (err error)

func (*XDGSCRAMClient) Done

func (x *XDGSCRAMClient) Done() bool

func (*XDGSCRAMClient) Step

func (x *XDGSCRAMClient) Step(challenge string) (response string, err error)

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL