Documentation
¶
Index ¶
- type KafkaDataSourceMock
- func (kp *KafkaDataSourceMock) AlterBrokerConfig(brokerID int32, key, value string) error
- func (kp *KafkaDataSourceMock) AlterClientQuotas(entity api.ClientQuotaEntity, quotas map[string]float64) error
- func (kp *KafkaDataSourceMock) AlterReplicaLogDir(brokerID int32, topic string, partition int32, logDir string) error
- func (kp *KafkaDataSourceMock) CancelTopicAnalysis(topicName string) error
- func (kp *KafkaDataSourceMock) ChangeReplicationFactor(name string, newFactor int16) error
- func (kp *KafkaDataSourceMock) CheckSchemaCompatibility(subject, schemaText, schemaType string) (bool, []string, error)
- func (kp *KafkaDataSourceMock) ConsumeTopic(ctx context.Context, topicName string, flags api.ConsumeFlags, ...) error
- func (kp *KafkaDataSourceMock) CreateACL(entry api.ACLEntry) error
- func (kp *KafkaDataSourceMock) CreateConnector(connect, name string, config map[string]string) (api.Connector, error)
- func (kp *KafkaDataSourceMock) CreateTopic(name string, numPartitions int32, replicationFactor int16, ...) error
- func (kp *KafkaDataSourceMock) DecodeMessage(_ context.Context, msg api.Message) (api.Message, error)
- func (kp *KafkaDataSourceMock) DeleteACL(entry api.ACLEntry) error
- func (kp *KafkaDataSourceMock) DeleteConnector(connect, name string) error
- func (kp *KafkaDataSourceMock) DeleteConsumerGroup(groupID string) error
- func (kp *KafkaDataSourceMock) DeleteConsumerGroupOffsets(groupID string, topic string) error
- func (kp *KafkaDataSourceMock) DeleteSchemaVersion(subject string, version int, permanent bool) error
- func (kp *KafkaDataSourceMock) DeleteSubject(subject string, permanent bool) ([]int, error)
- func (kp *KafkaDataSourceMock) DeleteTopic(name string) error
- func (kp *KafkaDataSourceMock) ExecuteKsql(ctx context.Context, sql string, props map[string]string) (<-chan api.KsqlResultTable, error)
- func (kp *KafkaDataSourceMock) GetACLs() ([]api.ACLEntry, error)
- func (kp *KafkaDataSourceMock) GetACLsFiltered(filter api.ACLFilter) ([]api.ACLEntry, error)
- func (kp *KafkaDataSourceMock) GetBrokerConfig(brokerID int32) ([]api.BrokerConfigEntry, error)
- func (kp *KafkaDataSourceMock) GetBrokerLogDirs(brokerIDs []int32) (map[int32][]api.BrokerLogDir, error)
- func (kp *KafkaDataSourceMock) GetBrokerMetrics(brokerID int32) (string, error)
- func (kp *KafkaDataSourceMock) GetBrokerStats() (map[int32]api.BrokerStats, api.BrokerSummary, error)
- func (kp *KafkaDataSourceMock) GetBrokers() ([]api.BrokerInfo, error)
- func (kp *KafkaDataSourceMock) GetClientQuotas() ([]api.ClientQuotaEntry, error)
- func (kp *KafkaDataSourceMock) GetClusterCapabilities(_ context.Context, clusterName string) ([]api.Capability, error)
- func (kp *KafkaDataSourceMock) GetClusterDetails(clusterName string) (api.ClusterInfo, error)
- func (kp *KafkaDataSourceMock) GetClusterStatistics(_ context.Context, clusterName string) (api.ClusterStatistics, error)
- func (kp *KafkaDataSourceMock) GetConnectClusters(withStats bool) ([]api.ConnectCluster, error)
- func (kp *KafkaDataSourceMock) GetConnectorDetails(connect, name string) (api.ConnectorDetails, error)
- func (kp *KafkaDataSourceMock) GetConnectorNames(connect string) ([]string, error)
- func (kp *KafkaDataSourceMock) GetConnectorPlugins(connect string) ([]api.ConnectorPlugin, error)
- func (kp *KafkaDataSourceMock) GetConnectors() ([]api.Connector, error)
- func (kp *KafkaDataSourceMock) GetConsumerGroupDetail(groupID string) (api.ConsumerGroupDetail, error)
- func (kp *KafkaDataSourceMock) GetConsumerGroupDetails(groupIDs []string) ([]api.ConsumerGroup, error)
- func (kp *KafkaDataSourceMock) GetConsumerGroups() ([]api.ConsumerGroup, error)
- func (kp *KafkaDataSourceMock) GetConsumerGroupsForTopic(topic string) ([]api.ConsumerGroup, error)
- func (kp *KafkaDataSourceMock) GetContext() string
- func (kp *KafkaDataSourceMock) GetContexts() ([]string, error)
- func (kp *KafkaDataSourceMock) GetGlobalCompatibility() (api.CompatibilityLevel, error)
- func (kp *KafkaDataSourceMock) GetMessageSchemaInfo(keySchemaID, valueSchemaID string) (*api.MessageSchemaInfo, error)
- func (kp *KafkaDataSourceMock) GetSchemaContent(subject string, version int) (string, error)
- func (kp *KafkaDataSourceMock) GetSchemaDetails(subjects []string) ([]api.Schema, error)
- func (kp *KafkaDataSourceMock) GetSchemaVersions(subject string) ([]api.SchemaVersion, error)
- func (kp *KafkaDataSourceMock) GetSchemas() ([]api.Schema, error)
- func (kp *KafkaDataSourceMock) GetSubjectCompatibility(subject string) (api.CompatibilityLevel, bool, error)
- func (kp *KafkaDataSourceMock) GetTopicAnalysis(topicName string) (*api.TopicAnalysis, error)
- func (kp *KafkaDataSourceMock) GetTopicConfig(topicName string) ([]api.TopicConfigEntry, error)
- func (kp *KafkaDataSourceMock) GetTopicDetails(topicName string) (api.TopicDetails, error)
- func (kp *KafkaDataSourceMock) GetTopicMessageCounts(topics map[string]int32) (map[string]int64, error)
- func (kp *KafkaDataSourceMock) GetTopicNames() ([]string, error)
- func (kp *KafkaDataSourceMock) GetTopicSizes(topicNames []string) (map[string]int64, error)
- func (kp *KafkaDataSourceMock) GetTopics() (map[string]api.Topic, error)
- func (kp *KafkaDataSourceMock) IncreasePartitions(name string, totalCount int32) error
- func (kp *KafkaDataSourceMock) Init(cfgOption string)
- func (kp *KafkaDataSourceMock) IsTopicDeletionEnabled() (bool, error)
- func (kp *KafkaDataSourceMock) ListKsqlStreams() ([]api.KsqlStream, error)
- func (kp *KafkaDataSourceMock) ListKsqlTables() ([]api.KsqlTable, error)
- func (kp *KafkaDataSourceMock) ListSerdes() []string
- func (kp *KafkaDataSourceMock) PauseConnector(connect, name string) error
- func (kp *KafkaDataSourceMock) ProduceMessage(ctx context.Context, topic string, rec api.ProduceRecord) error
- func (kp *KafkaDataSourceMock) PurgeTopicMessages(name string, partition int32) error
- func (kp *KafkaDataSourceMock) RecreateTopic(name string) error
- func (kp *KafkaDataSourceMock) RegisterSchema(subject, schemaText, schemaType string) (api.Schema, error)
- func (kp *KafkaDataSourceMock) ResetConnectorOffsets(connect, name string) error
- func (kp *KafkaDataSourceMock) ResetConsumerGroupOffsets(ctx context.Context, req api.OffsetResetRequest) error
- func (kp *KafkaDataSourceMock) RestartConnector(connect, name string) error
- func (kp *KafkaDataSourceMock) RestartConnectorTask(connect, name string, taskID int) error
- func (kp *KafkaDataSourceMock) ResumeConnector(connect, name string) error
- func (kp *KafkaDataSourceMock) SetContext(contextName string) error
- func (kp *KafkaDataSourceMock) SetDeletionDisabled(disabled bool)
- func (kp *KafkaDataSourceMock) SetGlobalCompatibility(level api.CompatibilityLevel) error
- func (kp *KafkaDataSourceMock) SetSubjectCompatibility(subject string, level api.CompatibilityLevel) error
- func (kp *KafkaDataSourceMock) StartTopicAnalysis(_ context.Context, topicName string) error
- func (kp *KafkaDataSourceMock) StopConnector(connect, name string) error
- func (kp *KafkaDataSourceMock) UpdateConnectorConfig(connect, name string, config map[string]string) (api.Connector, error)
- func (kp *KafkaDataSourceMock) UpdateTopicConfig(name string, entries map[string]*string) error
- func (kp *KafkaDataSourceMock) ValidateClusterConnection(_ context.Context, clusterName string) ([]api.ValidationResult, error)
- func (kp *KafkaDataSourceMock) ValidateConnectorConfig(connect, pluginClass string, config map[string]string) (api.ConnectorValidationResult, error)
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
This section is empty.
Types ¶
type KafkaDataSourceMock ¶
type KafkaDataSourceMock struct {
// contains filtered or unexported fields
}
func (*KafkaDataSourceMock) AlterBrokerConfig ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) AlterBrokerConfig(brokerID int32, key, value string) error
AlterBrokerConfig implements api.KafkaDataSource, mutating in-memory config. A value of "invalid" is rejected to exercise the InvalidConfigError path.
func (*KafkaDataSourceMock) AlterClientQuotas ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) AlterClientQuotas(entity api.ClientQuotaEntity, quotas map[string]float64) error
AlterClientQuotas implements api.KafkaDataSource with replace/delete semantics: the submitted map fully replaces the entity's properties; an empty or nil map deletes the entity.
func (*KafkaDataSourceMock) AlterReplicaLogDir ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) AlterReplicaLogDir(brokerID int32, topic string, partition int32, logDir string) error
AlterReplicaLogDir implements api.KafkaDataSource, moving a partition's log to a different directory on the broker.
func (*KafkaDataSourceMock) CancelTopicAnalysis ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) CancelTopicAnalysis(topicName string) error
func (*KafkaDataSourceMock) ChangeReplicationFactor ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) ChangeReplicationFactor(name string, newFactor int16) error
func (*KafkaDataSourceMock) CheckSchemaCompatibility ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) CheckSchemaCompatibility(subject, schemaText, schemaType string) (bool, []string, error)
CheckSchemaCompatibility returns incompatible when the candidate contains the magic "INCOMPATIBLE" marker, compatible otherwise.
func (*KafkaDataSourceMock) ConsumeTopic ¶
func (kp *KafkaDataSourceMock) ConsumeTopic(ctx context.Context, topicName string, flags api.ConsumeFlags, handleMessage api.MessageHandlerFunc, onError func(err any)) error
func (*KafkaDataSourceMock) CreateACL ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) CreateACL(entry api.ACLEntry) error
CreateACL implements api.KafkaDataSource. It validates the entry and appends it (defaulting host and pattern type).
func (*KafkaDataSourceMock) CreateConnector ¶ added in v0.1.34
func (*KafkaDataSourceMock) CreateTopic ¶ added in v0.1.34
func (*KafkaDataSourceMock) DecodeMessage ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) DecodeMessage(_ context.Context, msg api.Message) (api.Message, error)
DecodeMessage implements api.KafkaDataSource. Mock messages already have decoded Key/Value; RawKey/RawValue are treated as plain text.
func (*KafkaDataSourceMock) DeleteACL ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) DeleteACL(entry api.ACLEntry) error
DeleteACL implements api.KafkaDataSource. It removes the binding matching the full definition or returns an ACLNotFoundError.
func (*KafkaDataSourceMock) DeleteConnector ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) DeleteConnector(connect, name string) error
func (*KafkaDataSourceMock) DeleteConsumerGroup ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) DeleteConsumerGroup(groupID string) error
DeleteConsumerGroup implements api.KafkaDataSource (CG-6).
func (*KafkaDataSourceMock) DeleteConsumerGroupOffsets ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) DeleteConsumerGroupOffsets(groupID string, topic string) error
DeleteConsumerGroupOffsets implements api.KafkaDataSource (CG-6). Only the named topic's committed offsets are removed.
func (*KafkaDataSourceMock) DeleteSchemaVersion ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) DeleteSchemaVersion(subject string, version int, permanent bool) error
DeleteSchemaVersion removes a single version (version=-1 targets the latest).
func (*KafkaDataSourceMock) DeleteSubject ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) DeleteSubject(subject string, permanent bool) ([]int, error)
DeleteSubject removes all versions of a subject, returning the deleted numbers.
func (*KafkaDataSourceMock) DeleteTopic ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) DeleteTopic(name string) error
func (*KafkaDataSourceMock) ExecuteKsql ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) ExecuteKsql(ctx context.Context, sql string, props map[string]string) (<-chan api.KsqlResultTable, error)
ExecuteKsql mirrors the real datasource contract: validation runs first (sharing the semantics of the real classifier), SHOW/LIST/DESCRIBE return canned tables, DDL returns a success table, and SELECT streams a schema table followed by a ticking row every ~300 ms until ctx is cancelled.
func (*KafkaDataSourceMock) GetACLs ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) GetACLs() ([]api.ACLEntry, error)
GetACLs implements api.KafkaDataSource (the match-any case).
func (*KafkaDataSourceMock) GetACLsFiltered ¶ added in v0.1.34
GetACLsFiltered implements api.KafkaDataSource.
func (*KafkaDataSourceMock) GetBrokerConfig ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) GetBrokerConfig(brokerID int32) ([]api.BrokerConfigEntry, error)
GetBrokerConfig implements api.KafkaDataSource.
func (*KafkaDataSourceMock) GetBrokerLogDirs ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) GetBrokerLogDirs(brokerIDs []int32) (map[int32][]api.BrokerLogDir, error)
GetBrokerLogDirs implements api.KafkaDataSource, honouring the filter/all-brokers semantics (empty = all, unknown IDs dropped).
func (*KafkaDataSourceMock) GetBrokerMetrics ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) GetBrokerMetrics(brokerID int32) (string, error)
GetBrokerMetrics implements api.KafkaDataSource, returning a sample JSON snapshot.
func (*KafkaDataSourceMock) GetBrokerStats ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) GetBrokerStats() (map[int32]api.BrokerStats, api.BrokerSummary, error)
GetBrokerStats implements api.KafkaDataSource with fixed, deterministic values (including one replica skew >= 20% for styling tests).
func (*KafkaDataSourceMock) GetBrokers ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) GetBrokers() ([]api.BrokerInfo, error)
GetBrokers implements api.KafkaDataSource.
func (*KafkaDataSourceMock) GetClientQuotas ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) GetClientQuotas() ([]api.ClientQuotaEntry, error)
GetClientQuotas implements api.KafkaDataSource.
func (*KafkaDataSourceMock) GetClusterCapabilities ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) GetClusterCapabilities(_ context.Context, clusterName string) ([]api.Capability, error)
GetClusterCapabilities implements api.KafkaDataSource.
func (*KafkaDataSourceMock) GetClusterDetails ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) GetClusterDetails(clusterName string) (api.ClusterInfo, error)
GetClusterDetails returns mock configuration details for the named cluster.
func (*KafkaDataSourceMock) GetClusterStatistics ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) GetClusterStatistics(_ context.Context, clusterName string) (api.ClusterStatistics, error)
GetClusterStatistics implements api.KafkaDataSource.
func (*KafkaDataSourceMock) GetConnectClusters ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) GetConnectClusters(withStats bool) ([]api.ConnectCluster, error)
func (*KafkaDataSourceMock) GetConnectorDetails ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) GetConnectorDetails(connect, name string) (api.ConnectorDetails, error)
func (*KafkaDataSourceMock) GetConnectorNames ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) GetConnectorNames(connect string) ([]string, error)
func (*KafkaDataSourceMock) GetConnectorPlugins ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) GetConnectorPlugins(connect string) ([]api.ConnectorPlugin, error)
func (*KafkaDataSourceMock) GetConnectors ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) GetConnectors() ([]api.Connector, error)
func (*KafkaDataSourceMock) GetConsumerGroupDetail ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) GetConsumerGroupDetail(groupID string) (api.ConsumerGroupDetail, error)
GetConsumerGroupDetail implements api.KafkaDataSource (CG-3).
func (*KafkaDataSourceMock) GetConsumerGroupDetails ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) GetConsumerGroupDetails(groupIDs []string) ([]api.ConsumerGroup, error)
GetConsumerGroupDetails implements api.KafkaDataSource (CG-4).
func (*KafkaDataSourceMock) GetConsumerGroups ¶
func (kp *KafkaDataSourceMock) GetConsumerGroups() ([]api.ConsumerGroup, error)
GetConsumerGroups retrieves consumer groups for the current context
func (*KafkaDataSourceMock) GetConsumerGroupsForTopic ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) GetConsumerGroupsForTopic(topic string) ([]api.ConsumerGroup, error)
GetConsumerGroupsForTopic implements api.KafkaDataSource (CG-5).
func (*KafkaDataSourceMock) GetContext ¶
func (kp *KafkaDataSourceMock) GetContext() string
func (*KafkaDataSourceMock) GetContexts ¶
func (kp *KafkaDataSourceMock) GetContexts() ([]string, error)
GetContexts retrieves a list of Kafka contexts
func (*KafkaDataSourceMock) GetGlobalCompatibility ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) GetGlobalCompatibility() (api.CompatibilityLevel, error)
GetGlobalCompatibility returns the mock global compatibility level.
func (*KafkaDataSourceMock) GetMessageSchemaInfo ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) GetMessageSchemaInfo(keySchemaID, valueSchemaID string) (*api.MessageSchemaInfo, error)
GetMessageSchemaInfo implements api.KafkaDataSource
func (*KafkaDataSourceMock) GetSchemaContent ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) GetSchemaContent(subject string, version int) (string, error)
GetSchemaContent returns a version's schema text (latest when version <= 0).
func (*KafkaDataSourceMock) GetSchemaDetails ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) GetSchemaDetails(subjects []string) ([]api.Schema, error)
GetSchemaDetails returns latest-version metadata plus effective compatibility.
func (*KafkaDataSourceMock) GetSchemaVersions ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) GetSchemaVersions(subject string) ([]api.SchemaVersion, error)
GetSchemaVersions lists all versions of a subject (ascending).
func (*KafkaDataSourceMock) GetSchemas ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) GetSchemas() ([]api.Schema, error)
GetSchemas returns all subject names currently registered.
func (*KafkaDataSourceMock) GetSubjectCompatibility ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) GetSubjectCompatibility(subject string) (api.CompatibilityLevel, bool, error)
GetSubjectCompatibility returns a subject's effective level with a fallback flag.
func (*KafkaDataSourceMock) GetTopicAnalysis ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) GetTopicAnalysis(topicName string) (*api.TopicAnalysis, error)
func (*KafkaDataSourceMock) GetTopicConfig ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) GetTopicConfig(topicName string) ([]api.TopicConfigEntry, error)
GetTopicConfig returns a realistic config set with the topic's own overrides layered over cluster defaults, plus one sensitive entry.
func (*KafkaDataSourceMock) GetTopicDetails ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) GetTopicDetails(topicName string) (api.TopicDetails, error)
GetTopicDetails builds multi-partition fixtures including one under-replicated partition.
func (*KafkaDataSourceMock) GetTopicMessageCounts ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) GetTopicMessageCounts(topics map[string]int32) (map[string]int64, error)
GetTopicMessageCounts returns simulated message counts for the given topics. Counts grow with elapsed time so successive calls yield increasing values, letting the background collector derive message-in rates from the deltas.
func (*KafkaDataSourceMock) GetTopicNames ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) GetTopicNames() ([]string, error)
GetTopicNames returns only topic names (mock version, same data as GetTopics but names only).
func (*KafkaDataSourceMock) GetTopicSizes ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) GetTopicSizes(topicNames []string) (map[string]int64, error)
GetTopicSizes returns deterministic sizes (1 KiB per message) for known topics.
func (*KafkaDataSourceMock) GetTopics ¶
func (kp *KafkaDataSourceMock) GetTopics() (map[string]api.Topic, error)
GetTopics retrieves a list of Kafka topics for the current context
func (*KafkaDataSourceMock) IncreasePartitions ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) IncreasePartitions(name string, totalCount int32) error
func (*KafkaDataSourceMock) Init ¶
func (kp *KafkaDataSourceMock) Init(cfgOption string)
func (*KafkaDataSourceMock) IsTopicDeletionEnabled ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) IsTopicDeletionEnabled() (bool, error)
func (*KafkaDataSourceMock) ListKsqlStreams ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) ListKsqlStreams() ([]api.KsqlStream, error)
func (*KafkaDataSourceMock) ListKsqlTables ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) ListKsqlTables() ([]api.KsqlTable, error)
func (*KafkaDataSourceMock) ListSerdes ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) ListSerdes() []string
ListSerdes returns a plausible static list of serde names. (MSG-18)
func (*KafkaDataSourceMock) PauseConnector ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) PauseConnector(connect, name string) error
func (*KafkaDataSourceMock) ProduceMessage ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) ProduceMessage(ctx context.Context, topic string, rec api.ProduceRecord) error
ProduceMessage appends a record to the in-memory store so it is browsable via ConsumeTopic (MSG-30).
func (*KafkaDataSourceMock) PurgeTopicMessages ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) PurgeTopicMessages(name string, partition int32) error
func (*KafkaDataSourceMock) RecreateTopic ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) RecreateTopic(name string) error
func (*KafkaDataSourceMock) RegisterSchema ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) RegisterSchema(subject, schemaText, schemaType string) (api.Schema, error)
RegisterSchema appends a new version (creating the subject when new).
func (*KafkaDataSourceMock) ResetConnectorOffsets ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) ResetConnectorOffsets(connect, name string) error
func (*KafkaDataSourceMock) ResetConsumerGroupOffsets ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) ResetConsumerGroupOffsets(ctx context.Context, req api.OffsetResetRequest) error
ResetConsumerGroupOffsets implements api.KafkaDataSource (CG-7, CG-8).
func (*KafkaDataSourceMock) RestartConnector ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) RestartConnector(connect, name string) error
func (*KafkaDataSourceMock) RestartConnectorTask ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) RestartConnectorTask(connect, name string, taskID int) error
func (*KafkaDataSourceMock) ResumeConnector ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) ResumeConnector(connect, name string) error
func (*KafkaDataSourceMock) SetContext ¶
func (kp *KafkaDataSourceMock) SetContext(contextName string) error
SetContext implements api.KafkaDataSource.
func (*KafkaDataSourceMock) SetDeletionDisabled ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) SetDeletionDisabled(disabled bool)
SetDeletionDisabled toggles the simulated delete.topic.enable=false state for UI tests.
func (*KafkaDataSourceMock) SetGlobalCompatibility ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) SetGlobalCompatibility(level api.CompatibilityLevel) error
SetGlobalCompatibility validates and updates the global level.
func (*KafkaDataSourceMock) SetSubjectCompatibility ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) SetSubjectCompatibility(subject string, level api.CompatibilityLevel) error
SetSubjectCompatibility validates and updates a subject's level.
func (*KafkaDataSourceMock) StartTopicAnalysis ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) StartTopicAnalysis(_ context.Context, topicName string) error
StartTopicAnalysis simulates a fast scan over a sample of generated messages and stores a completed result.
func (*KafkaDataSourceMock) StopConnector ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) StopConnector(connect, name string) error
func (*KafkaDataSourceMock) UpdateConnectorConfig ¶ added in v0.1.34
func (*KafkaDataSourceMock) UpdateTopicConfig ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) UpdateTopicConfig(name string, entries map[string]*string) error
func (*KafkaDataSourceMock) ValidateClusterConnection ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) ValidateClusterConnection(_ context.Context, clusterName string) ([]api.ValidationResult, error)
ValidateClusterConnection implements api.KafkaDataSource.
func (*KafkaDataSourceMock) ValidateConnectorConfig ¶ added in v0.1.34
func (kp *KafkaDataSourceMock) ValidateConnectorConfig(connect, pluginClass string, config map[string]string) (api.ConnectorValidationResult, error)