Documentation
¶
Index ¶
- Variables
- func DoConsume(ctx context.Context, topic string, consumeFlags api.ConsumeFlags, ...)
- func DoConsumeWithConfig(ctx context.Context, topic string, consumeFlags api.ConsumeFlags, ...)
- func DoConsumeWithDeps(ctx context.Context, topic string, consumeFlags api.ConsumeFlags, ...)
- func InitFromConfig(cfgPath string) error
- func InitTUIWriters()
- func PrepareOAuthDeviceFlow(cfgPath string, w io.Writer) error
- func SetOverrides(brokers []string, schemaRegistry, cluster string, verboseLogging bool)
- type AlterConfigCall
- type AlterQuotaCall
- type ClusterAdminInterface
- type ConfigManager
- type ConfigProviderInterface
- type ConsumeConfig
- type ConsumerInterface
- type CreatePartitionsCall
- type CreateTopicCall
- type DefaultConfigManager
- type DefaultConfigProvider
- type DefaultConsumer
- type DefaultKafkaClientFactory
- type DeleteOffsetCall
- type DeleteRecordsCall
- type DescribeQuotasCall
- type KafkaClientFactory
- type KafkaDataSourceKaf
- func (kp KafkaDataSourceKaf) AlterBrokerConfig(brokerID int32, key, value string) error
- func (kp KafkaDataSourceKaf) AlterClientQuotas(entity api.ClientQuotaEntity, quotas map[string]float64) error
- func (kp KafkaDataSourceKaf) AlterReplicaLogDir(brokerID int32, topic string, partition int32, logDir string) error
- func (kp KafkaDataSourceKaf) CancelTopicAnalysis(topicName string) error
- func (kp KafkaDataSourceKaf) ChangeReplicationFactor(name string, newFactor int16) error
- func (kp KafkaDataSourceKaf) CheckSchemaCompatibility(subject, schemaText, schemaType string) (bool, []string, error)
- func (kp KafkaDataSourceKaf) ConsumeTopic(ctx context.Context, topicName string, flags api.ConsumeFlags, ...) error
- func (kp KafkaDataSourceKaf) CreateACL(entry api.ACLEntry) error
- func (kp KafkaDataSourceKaf) CreateConnector(connect, name string, config map[string]string) (api.Connector, error)
- func (kp KafkaDataSourceKaf) CreateTopic(name string, numPartitions int32, replicationFactor int16, ...) error
- func (kp KafkaDataSourceKaf) DecodeMessage(_ context.Context, msg api.Message) (api.Message, error)
- func (kp KafkaDataSourceKaf) DeleteACL(entry api.ACLEntry) error
- func (kp KafkaDataSourceKaf) DeleteConnector(connect, name string) error
- func (kp KafkaDataSourceKaf) DeleteConsumerGroup(groupID string) error
- func (kp KafkaDataSourceKaf) DeleteConsumerGroupOffsets(groupID string, topic string) error
- func (kp KafkaDataSourceKaf) DeleteSchemaVersion(subject string, version int, permanent bool) error
- func (kp KafkaDataSourceKaf) DeleteSubject(subject string, permanent bool) ([]int, error)
- func (kp KafkaDataSourceKaf) DeleteTopic(name string) error
- func (kp KafkaDataSourceKaf) ExecuteKsql(ctx context.Context, sql string, props map[string]string) (<-chan api.KsqlResultTable, error)
- func (kp KafkaDataSourceKaf) GetACLs() ([]api.ACLEntry, error)
- func (kp KafkaDataSourceKaf) GetACLsFiltered(filter api.ACLFilter) ([]api.ACLEntry, error)
- func (kp KafkaDataSourceKaf) GetBrokerConfig(brokerID int32) ([]api.BrokerConfigEntry, error)
- func (kp KafkaDataSourceKaf) GetBrokerLogDirs(brokerIDs []int32) (map[int32][]api.BrokerLogDir, error)
- func (kp KafkaDataSourceKaf) GetBrokerMetrics(brokerID int32) (string, error)
- func (kp KafkaDataSourceKaf) GetBrokerStats() (map[int32]api.BrokerStats, api.BrokerSummary, error)
- func (kp KafkaDataSourceKaf) GetBrokers() ([]api.BrokerInfo, error)
- func (kp KafkaDataSourceKaf) GetClientQuotas() ([]api.ClientQuotaEntry, error)
- func (kp KafkaDataSourceKaf) GetClusterCapabilities(_ context.Context, clusterName string) ([]api.Capability, error)
- func (kp KafkaDataSourceKaf) GetClusterDetails(clusterName string) (api.ClusterInfo, error)
- func (kp KafkaDataSourceKaf) GetClusterStatistics(_ context.Context, clusterName string) (api.ClusterStatistics, error)
- func (kp KafkaDataSourceKaf) GetConnectClusters(withStats bool) ([]api.ConnectCluster, error)
- func (kp KafkaDataSourceKaf) GetConnectorDetails(connect, name string) (api.ConnectorDetails, error)
- func (kp KafkaDataSourceKaf) GetConnectorNames(connect string) ([]string, error)
- func (kp KafkaDataSourceKaf) GetConnectorPlugins(connect string) ([]api.ConnectorPlugin, error)
- func (kp KafkaDataSourceKaf) GetConnectors() ([]api.Connector, error)
- func (kp KafkaDataSourceKaf) GetConsumerGroupDetail(groupID string) (api.ConsumerGroupDetail, error)
- func (kp KafkaDataSourceKaf) GetConsumerGroupDetails(groupIDs []string) ([]api.ConsumerGroup, error)
- func (kp KafkaDataSourceKaf) GetConsumerGroups() ([]api.ConsumerGroup, error)
- func (kp KafkaDataSourceKaf) GetConsumerGroupsForTopic(topic string) ([]api.ConsumerGroup, error)
- func (kp KafkaDataSourceKaf) GetContext() string
- func (kp KafkaDataSourceKaf) GetContexts() ([]string, error)
- func (kp KafkaDataSourceKaf) GetGlobalCompatibility() (api.CompatibilityLevel, error)
- func (kp KafkaDataSourceKaf) GetMessageSchemaInfo(keySchemaID, valueSchemaID string) (*api.MessageSchemaInfo, error)
- func (kp KafkaDataSourceKaf) GetSchemaContent(subject string, version int) (string, error)
- func (kp KafkaDataSourceKaf) GetSchemaDetails(subjects []string) ([]api.Schema, error)
- func (kp KafkaDataSourceKaf) GetSchemaVersions(subject string) ([]api.SchemaVersion, error)
- func (kp KafkaDataSourceKaf) GetSchemas() ([]api.Schema, error)
- func (kp KafkaDataSourceKaf) GetSubjectCompatibility(subject string) (api.CompatibilityLevel, bool, error)
- func (kp KafkaDataSourceKaf) GetTopicAnalysis(topicName string) (*api.TopicAnalysis, error)
- func (kp KafkaDataSourceKaf) GetTopicConfig(topicName string) ([]api.TopicConfigEntry, error)
- func (kp KafkaDataSourceKaf) GetTopicDetails(topicName string) (api.TopicDetails, error)
- func (kp KafkaDataSourceKaf) GetTopicHealth(topicNames []string) (map[string]api.TopicHealth, error)
- func (kp KafkaDataSourceKaf) GetTopicMessageCounts(topics map[string]int32) (map[string]int64, error)
- func (kp KafkaDataSourceKaf) GetTopicNames() ([]string, error)
- func (kp KafkaDataSourceKaf) GetTopicSizes(topicNames []string) (map[string]int64, error)
- func (kp KafkaDataSourceKaf) GetTopics() (map[string]api.Topic, error)
- func (kp KafkaDataSourceKaf) IncreasePartitions(name string, totalCount int32) error
- func (kp *KafkaDataSourceKaf) Init(cfgOption string)
- func (kp KafkaDataSourceKaf) IsTopicDeletionEnabled() (bool, error)
- func (kp KafkaDataSourceKaf) ListKsqlStreams() ([]api.KsqlStream, error)
- func (kp KafkaDataSourceKaf) ListKsqlTables() ([]api.KsqlTable, error)
- func (kp KafkaDataSourceKaf) ListSerdes() []string
- func (kp KafkaDataSourceKaf) PauseConnector(connect, name string) error
- func (kp KafkaDataSourceKaf) ProduceMessage(ctx context.Context, topic string, rec api.ProduceRecord) error
- func (kp KafkaDataSourceKaf) PurgeTopicMessages(name string, partition int32) error
- func (kp KafkaDataSourceKaf) RecreateTopic(name string) error
- func (kp KafkaDataSourceKaf) RegisterSchema(subject, schemaText, schemaType string) (api.Schema, error)
- func (kp *KafkaDataSourceKaf) Reload(effective appconfig.Config) error
- func (kp KafkaDataSourceKaf) ResetConnectorOffsets(connect, name string) error
- func (kp KafkaDataSourceKaf) ResetConsumerGroupOffsets(ctx context.Context, req api.OffsetResetRequest) error
- func (kp KafkaDataSourceKaf) RestartConnector(connect, name string) error
- func (kp KafkaDataSourceKaf) RestartConnectorTask(connect, name string, taskID int) error
- func (kp KafkaDataSourceKaf) ResumeConnector(connect, name string) error
- func (kp KafkaDataSourceKaf) SetContext(contextName string) error
- func (kp KafkaDataSourceKaf) SetGlobalCompatibility(level api.CompatibilityLevel) error
- func (kp KafkaDataSourceKaf) SetSubjectCompatibility(subject string, level api.CompatibilityLevel) error
- func (kp KafkaDataSourceKaf) StartTopicAnalysis(ctx context.Context, topicName string) error
- func (kp KafkaDataSourceKaf) StopConnector(connect, name string) error
- func (kp KafkaDataSourceKaf) UpdateConnectorConfig(connect, name string, config map[string]string) (api.Connector, error)
- func (kp KafkaDataSourceKaf) UpdateTopicConfig(name string, entries map[string]*string) error
- func (kp KafkaDataSourceKaf) ValidateCandidate(ctx context.Context, candidate appconfig.Config) api.ValidationReport
- func (kp KafkaDataSourceKaf) ValidateClusterConnection(ctx context.Context, clusterName string) ([]api.ValidationResult, error)
- func (kp KafkaDataSourceKaf) ValidateConnectorConfig(connect, pluginClass string, config map[string]string) (api.ConnectorValidationResult, error)
- type MockClusterAdmin
- func (m *MockClusterAdmin) AlterClientQuotas(entity []sarama.QuotaEntityComponent, op sarama.ClientQuotasOp, ...) error
- func (m *MockClusterAdmin) AlterPartitionReassignments(topic string, assignment [][]int32) error
- func (m *MockClusterAdmin) Close() error
- func (m *MockClusterAdmin) CreateACLs(resourceACLs []*sarama.ResourceAcls) error
- func (m *MockClusterAdmin) CreatePartitions(topic string, count int32, assignment [][]int32, validateOnly bool) error
- func (m *MockClusterAdmin) CreateTopic(topic string, detail *sarama.TopicDetail, validateOnly bool) error
- func (m *MockClusterAdmin) DeleteACL(filter sarama.AclFilter, validateOnly bool) ([]sarama.MatchingAcl, error)
- func (m *MockClusterAdmin) DeleteConsumerGroup(group string) error
- func (m *MockClusterAdmin) DeleteConsumerGroupOffset(group string, topic string, partition int32) error
- func (m *MockClusterAdmin) DeleteRecords(topic string, partitionOffsets map[int32]int64) error
- func (m *MockClusterAdmin) DeleteTopic(topic string) error
- func (m *MockClusterAdmin) DescribeClientQuotas(components []sarama.QuotaFilterComponent, strict bool) ([]sarama.DescribeClientQuotasEntry, error)
- func (m *MockClusterAdmin) DescribeCluster() ([]*sarama.Broker, int32, error)
- func (m *MockClusterAdmin) DescribeConfig(resource sarama.ConfigResource) ([]sarama.ConfigEntry, error)
- func (m *MockClusterAdmin) DescribeConsumerGroups(groups []string) ([]*sarama.GroupDescription, error)
- func (m *MockClusterAdmin) DescribeLogDirs(brokers []int32) (map[int32][]sarama.DescribeLogDirsResponseDirMetadata, error)
- func (m *MockClusterAdmin) DescribeTopics(topics []string) ([]*sarama.TopicMetadata, error)
- func (m *MockClusterAdmin) IncrementalAlterConfig(resourceType sarama.ConfigResourceType, name string, ...) error
- func (m *MockClusterAdmin) ListAcls(filter sarama.AclFilter) ([]sarama.ResourceAcls, error)
- func (m *MockClusterAdmin) ListConsumerGroupOffsets(group string, topicPartitions map[string][]int32) (*sarama.OffsetFetchResponse, error)
- func (m *MockClusterAdmin) ListConsumerGroups() (map[string]string, error)
- func (m *MockClusterAdmin) ListTopics() (map[string]sarama.TopicDetail, error)
- type MockConfigManager
- type MockKafkaClientFactory
- type ReassignmentCall
- type XDGSCRAMClient
Constants ¶
This section is empty.
Variables ¶
var SHA256 scram.HashGeneratorFcn = func() hash.Hash { return sha256.New() }
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
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
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
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
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
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 (*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 (*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
DeleteOffsetCall captures the arguments of a DeleteConsumerGroupOffset call.
type DeleteRecordsCall ¶ added in v0.1.34
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 (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
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
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 (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 (*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 (*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
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
Source Files
¶
- acls.go
- broker.go
- broker_stats.go
- cluster_stats.go
- connect.go
- consume.go
- consume_interfaces.go
- consumer_groups.go
- datasource_kaf.go
- deviceflow.go
- interfaces.go
- ksql_client.go
- ksql_execute.go
- ksql_listings.go
- ksql_query_stream.go
- ksql_response.go
- ksql_statement.go
- mocks.go
- msk.go
- oauth.go
- offset_reset.go
- produce.go
- quotas.go
- reassignment.go
- schema_registry.go
- scram_client.go
- serde.go
- tokencache.go
- topic_admin.go
- topic_analysis.go
- topic_health.go
- validate.go