mock

package
v0.1.35 Latest Latest
Warning

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

Go to latest
Published: Aug 9, 2026 License: Apache-2.0 Imports: 10 Imported by: 0

Documentation

Index

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 (kp *KafkaDataSourceMock) CreateConnector(connect, name string, config map[string]string) (api.Connector, error)

func (*KafkaDataSourceMock) CreateTopic added in v0.1.34

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

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

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

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 (kp *KafkaDataSourceMock) UpdateConnectorConfig(connect, name string, config map[string]string) (api.Connector, error)

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)

Jump to

Keyboard shortcuts

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