Documentation
¶
Index ¶
- type ConsumerCallback
- type ConsumerCallback_Consume_Call
- func (_c *ConsumerCallback_Consume_Call[M]) Return(_a0 bool, _a1 error) *ConsumerCallback_Consume_Call[M]
- func (_c *ConsumerCallback_Consume_Call[M]) Run(run func(ctx context.Context, model M, attributes map[string]string)) *ConsumerCallback_Consume_Call[M]
- func (_c *ConsumerCallback_Consume_Call[M]) RunAndReturn(run func(context.Context, M, map[string]string) (bool, error)) *ConsumerCallback_Consume_Call[M]
- type ConsumerCallback_Expecter
- type Input
- type Input_Expecter
- type Input_IsHealthy_Call
- type Input_Run_Call
- type Input_Stop_Call
- type KinsumerAutoscaleOrchestrator
- type KinsumerAutoscaleOrchestrator_Expecter
- func (_e *KinsumerAutoscaleOrchestrator_Expecter) GetCurrentTaskCount(ctx interface{}) *KinsumerAutoscaleOrchestrator_GetCurrentTaskCount_Call
- func (_e *KinsumerAutoscaleOrchestrator_Expecter) UpdateTaskCount(ctx interface{}, taskCount interface{}) *KinsumerAutoscaleOrchestrator_UpdateTaskCount_Call
- type KinsumerAutoscaleOrchestrator_GetCurrentTaskCount_Call
- func (_c *KinsumerAutoscaleOrchestrator_GetCurrentTaskCount_Call) Return(_a0 int32, _a1 error) *KinsumerAutoscaleOrchestrator_GetCurrentTaskCount_Call
- func (_c *KinsumerAutoscaleOrchestrator_GetCurrentTaskCount_Call) Run(run func(ctx context.Context)) *KinsumerAutoscaleOrchestrator_GetCurrentTaskCount_Call
- func (_c *KinsumerAutoscaleOrchestrator_GetCurrentTaskCount_Call) RunAndReturn(run func(context.Context) (int32, error)) *KinsumerAutoscaleOrchestrator_GetCurrentTaskCount_Call
- type KinsumerAutoscaleOrchestrator_UpdateTaskCount_Call
- func (_c *KinsumerAutoscaleOrchestrator_UpdateTaskCount_Call) Return(_a0 error) *KinsumerAutoscaleOrchestrator_UpdateTaskCount_Call
- func (_c *KinsumerAutoscaleOrchestrator_UpdateTaskCount_Call) Run(run func(ctx context.Context, taskCount int32)) *KinsumerAutoscaleOrchestrator_UpdateTaskCount_Call
- func (_c *KinsumerAutoscaleOrchestrator_UpdateTaskCount_Call) RunAndReturn(run func(context.Context, int32) error) *KinsumerAutoscaleOrchestrator_UpdateTaskCount_Call
- type MessageEncoder
- func (_m *MessageEncoder) Decode(ctx context.Context, msg *stream.Message, out interface{}) (context.Context, map[string]string, error)
- func (_m *MessageEncoder) EXPECT() *MessageEncoder_Expecter
- func (_m *MessageEncoder) Encode(ctx context.Context, data interface{}, attributeSets ...map[string]string) (*stream.Message, error)
- type MessageEncoder_Decode_Call
- func (_c *MessageEncoder_Decode_Call) Return(_a0 context.Context, _a1 map[string]string, _a2 error) *MessageEncoder_Decode_Call
- func (_c *MessageEncoder_Decode_Call) Run(run func(ctx context.Context, msg *stream.Message, out interface{})) *MessageEncoder_Decode_Call
- func (_c *MessageEncoder_Decode_Call) RunAndReturn(...) *MessageEncoder_Decode_Call
- type MessageEncoder_Encode_Call
- type MessageEncoder_Expecter
- type Output
- type Output_Expecter
- type Output_WriteOne_Call
- func (_c *Output_WriteOne_Call) Return(_a0 error) *Output_WriteOne_Call
- func (_c *Output_WriteOne_Call) Run(run func(ctx context.Context, msg stream.WritableMessage)) *Output_WriteOne_Call
- func (_c *Output_WriteOne_Call) RunAndReturn(run func(context.Context, stream.WritableMessage) error) *Output_WriteOne_Call
- type Output_Write_Call
- func (_c *Output_Write_Call) Return(_a0 error) *Output_Write_Call
- func (_c *Output_Write_Call) Run(run func(ctx context.Context, batch []stream.WritableMessage)) *Output_Write_Call
- func (_c *Output_Write_Call) RunAndReturn(run func(context.Context, []stream.WritableMessage) error) *Output_Write_Call
- type PartitionerRand
- type PartitionerRand_Expecter
- type PartitionerRand_Intn_Call
- type Producer
- type ProducerDaemonAggregator
- type ProducerDaemonAggregator_Expecter
- type ProducerDaemonAggregator_Flush_Call
- func (_c *ProducerDaemonAggregator_Flush_Call) Return(_a0 []stream.AggregateFlush, _a1 error) *ProducerDaemonAggregator_Flush_Call
- func (_c *ProducerDaemonAggregator_Flush_Call) Run(run func()) *ProducerDaemonAggregator_Flush_Call
- func (_c *ProducerDaemonAggregator_Flush_Call) RunAndReturn(run func() ([]stream.AggregateFlush, error)) *ProducerDaemonAggregator_Flush_Call
- type ProducerDaemonAggregator_Write_Call
- func (_c *ProducerDaemonAggregator_Write_Call) Return(_a0 []stream.AggregateFlush, _a1 error) *ProducerDaemonAggregator_Write_Call
- func (_c *ProducerDaemonAggregator_Write_Call) Run(run func(ctx context.Context, msg *stream.Message)) *ProducerDaemonAggregator_Write_Call
- func (_c *ProducerDaemonAggregator_Write_Call) RunAndReturn(run func(context.Context, *stream.Message) ([]stream.AggregateFlush, error)) *ProducerDaemonAggregator_Write_Call
- type ProducerDaemonBatcher
- type ProducerDaemonBatcher_Append_Call
- func (_c *ProducerDaemonBatcher_Append_Call) Return(_a0 []stream.WritableMessage, _a1 error) *ProducerDaemonBatcher_Append_Call
- func (_c *ProducerDaemonBatcher_Append_Call) Run(run func(msg *stream.Message)) *ProducerDaemonBatcher_Append_Call
- func (_c *ProducerDaemonBatcher_Append_Call) RunAndReturn(run func(*stream.Message) ([]stream.WritableMessage, error)) *ProducerDaemonBatcher_Append_Call
- type ProducerDaemonBatcher_Expecter
- type ProducerDaemonBatcher_Flush_Call
- func (_c *ProducerDaemonBatcher_Flush_Call) Return(_a0 []stream.WritableMessage) *ProducerDaemonBatcher_Flush_Call
- func (_c *ProducerDaemonBatcher_Flush_Call) Run(run func()) *ProducerDaemonBatcher_Flush_Call
- func (_c *ProducerDaemonBatcher_Flush_Call) RunAndReturn(run func() []stream.WritableMessage) *ProducerDaemonBatcher_Flush_Call
- type Producer_Expecter
- type Producer_WriteOne_Call
- type Producer_Write_Call
- type RetryHandler
- type RetryHandler_Expecter
- type RetryHandler_Put_Call
- func (_c *RetryHandler_Put_Call) Return(_a0 error) *RetryHandler_Put_Call
- func (_c *RetryHandler_Put_Call) Run(run func(ctx context.Context, msg *stream.Message)) *RetryHandler_Put_Call
- func (_c *RetryHandler_Put_Call) RunAndReturn(run func(context.Context, *stream.Message) error) *RetryHandler_Put_Call
- type RunnableCallback
- type RunnableCallback_Expecter
- type RunnableCallback_Run_Call
- type RunnableConsumerCallback
- type RunnableConsumerCallback_Consume_Call
- func (_c *RunnableConsumerCallback_Consume_Call[M]) Return(_a0 bool, _a1 error) *RunnableConsumerCallback_Consume_Call[M]
- func (_c *RunnableConsumerCallback_Consume_Call[M]) Run(run func(ctx context.Context, model M, attributes map[string]string)) *RunnableConsumerCallback_Consume_Call[M]
- func (_c *RunnableConsumerCallback_Consume_Call[M]) RunAndReturn(run func(context.Context, M, map[string]string) (bool, error)) *RunnableConsumerCallback_Consume_Call[M]
- type RunnableConsumerCallback_Expecter
- type RunnableConsumerCallback_Run_Call
- func (_c *RunnableConsumerCallback_Run_Call[M]) Return(_a0 error) *RunnableConsumerCallback_Run_Call[M]
- func (_c *RunnableConsumerCallback_Run_Call[M]) Run(run func(ctx context.Context)) *RunnableConsumerCallback_Run_Call[M]
- func (_c *RunnableConsumerCallback_Run_Call[M]) RunAndReturn(run func(context.Context) error) *RunnableConsumerCallback_Run_Call[M]
- type RunnableUntypedConsumerCallback
- func (_m *RunnableUntypedConsumerCallback) Consume(ctx context.Context, model interface{}, attributes map[string]string) (bool, error)
- func (_m *RunnableUntypedConsumerCallback) EXPECT() *RunnableUntypedConsumerCallback_Expecter
- func (_m *RunnableUntypedConsumerCallback) GetModel(attributes map[string]string) (interface{}, error)
- func (_m *RunnableUntypedConsumerCallback) Run(ctx context.Context) error
- type RunnableUntypedConsumerCallback_Consume_Call
- func (_c *RunnableUntypedConsumerCallback_Consume_Call) Return(_a0 bool, _a1 error) *RunnableUntypedConsumerCallback_Consume_Call
- func (_c *RunnableUntypedConsumerCallback_Consume_Call) Run(run func(ctx context.Context, model interface{}, attributes map[string]string)) *RunnableUntypedConsumerCallback_Consume_Call
- func (_c *RunnableUntypedConsumerCallback_Consume_Call) RunAndReturn(run func(context.Context, interface{}, map[string]string) (bool, error)) *RunnableUntypedConsumerCallback_Consume_Call
- type RunnableUntypedConsumerCallback_Expecter
- func (_e *RunnableUntypedConsumerCallback_Expecter) Consume(ctx interface{}, model interface{}, attributes interface{}) *RunnableUntypedConsumerCallback_Consume_Call
- func (_e *RunnableUntypedConsumerCallback_Expecter) GetModel(attributes interface{}) *RunnableUntypedConsumerCallback_GetModel_Call
- func (_e *RunnableUntypedConsumerCallback_Expecter) Run(ctx interface{}) *RunnableUntypedConsumerCallback_Run_Call
- type RunnableUntypedConsumerCallback_GetModel_Call
- func (_c *RunnableUntypedConsumerCallback_GetModel_Call) Return(_a0 interface{}, _a1 error) *RunnableUntypedConsumerCallback_GetModel_Call
- func (_c *RunnableUntypedConsumerCallback_GetModel_Call) Run(run func(attributes map[string]string)) *RunnableUntypedConsumerCallback_GetModel_Call
- func (_c *RunnableUntypedConsumerCallback_GetModel_Call) RunAndReturn(run func(map[string]string) (interface{}, error)) *RunnableUntypedConsumerCallback_GetModel_Call
- type RunnableUntypedConsumerCallback_Run_Call
- func (_c *RunnableUntypedConsumerCallback_Run_Call) Return(_a0 error) *RunnableUntypedConsumerCallback_Run_Call
- func (_c *RunnableUntypedConsumerCallback_Run_Call) Run(run func(ctx context.Context)) *RunnableUntypedConsumerCallback_Run_Call
- func (_c *RunnableUntypedConsumerCallback_Run_Call) RunAndReturn(run func(context.Context) error) *RunnableUntypedConsumerCallback_Run_Call
- type SchemaRegistryAwareInput
- func (_m *SchemaRegistryAwareInput) EXPECT() *SchemaRegistryAwareInput_Expecter
- func (_m *SchemaRegistryAwareInput) InitSchemaRegistry(ctx context.Context, settings stream.SchemaSettingsWithEncoding) (stream.MessageBodyEncoder, error)
- func (_m *SchemaRegistryAwareInput) IsHealthy() bool
- func (_m *SchemaRegistryAwareInput) Run(ctx context.Context, process stream.InputProcess) error
- func (_m *SchemaRegistryAwareInput) Stop(ctx context.Context)
- type SchemaRegistryAwareInput_Expecter
- func (_e *SchemaRegistryAwareInput_Expecter) InitSchemaRegistry(ctx interface{}, settings interface{}) *SchemaRegistryAwareInput_InitSchemaRegistry_Call
- func (_e *SchemaRegistryAwareInput_Expecter) IsHealthy() *SchemaRegistryAwareInput_IsHealthy_Call
- func (_e *SchemaRegistryAwareInput_Expecter) Run(ctx interface{}, process interface{}) *SchemaRegistryAwareInput_Run_Call
- func (_e *SchemaRegistryAwareInput_Expecter) Stop(ctx interface{}) *SchemaRegistryAwareInput_Stop_Call
- type SchemaRegistryAwareInput_InitSchemaRegistry_Call
- func (_c *SchemaRegistryAwareInput_InitSchemaRegistry_Call) Return(_a0 stream.MessageBodyEncoder, _a1 error) *SchemaRegistryAwareInput_InitSchemaRegistry_Call
- func (_c *SchemaRegistryAwareInput_InitSchemaRegistry_Call) Run(run func(ctx context.Context, settings stream.SchemaSettingsWithEncoding)) *SchemaRegistryAwareInput_InitSchemaRegistry_Call
- func (_c *SchemaRegistryAwareInput_InitSchemaRegistry_Call) RunAndReturn(...) *SchemaRegistryAwareInput_InitSchemaRegistry_Call
- type SchemaRegistryAwareInput_IsHealthy_Call
- func (_c *SchemaRegistryAwareInput_IsHealthy_Call) Return(_a0 bool) *SchemaRegistryAwareInput_IsHealthy_Call
- func (_c *SchemaRegistryAwareInput_IsHealthy_Call) Run(run func()) *SchemaRegistryAwareInput_IsHealthy_Call
- func (_c *SchemaRegistryAwareInput_IsHealthy_Call) RunAndReturn(run func() bool) *SchemaRegistryAwareInput_IsHealthy_Call
- type SchemaRegistryAwareInput_Run_Call
- func (_c *SchemaRegistryAwareInput_Run_Call) Return(_a0 error) *SchemaRegistryAwareInput_Run_Call
- func (_c *SchemaRegistryAwareInput_Run_Call) Run(run func(ctx context.Context, process stream.InputProcess)) *SchemaRegistryAwareInput_Run_Call
- func (_c *SchemaRegistryAwareInput_Run_Call) RunAndReturn(run func(context.Context, stream.InputProcess) error) *SchemaRegistryAwareInput_Run_Call
- type SchemaRegistryAwareInput_Stop_Call
- func (_c *SchemaRegistryAwareInput_Stop_Call) Return() *SchemaRegistryAwareInput_Stop_Call
- func (_c *SchemaRegistryAwareInput_Stop_Call) Run(run func(ctx context.Context)) *SchemaRegistryAwareInput_Stop_Call
- func (_c *SchemaRegistryAwareInput_Stop_Call) RunAndReturn(run func(context.Context)) *SchemaRegistryAwareInput_Stop_Call
- type SchemaRegistryAwareOutput
- func (_m *SchemaRegistryAwareOutput) EXPECT() *SchemaRegistryAwareOutput_Expecter
- func (_m *SchemaRegistryAwareOutput) InitSchemaRegistry(ctx context.Context, settings stream.SchemaSettingsWithEncoding) (stream.MessageBodyEncoder, error)
- func (_m *SchemaRegistryAwareOutput) Write(ctx context.Context, batch []stream.WritableMessage) error
- func (_m *SchemaRegistryAwareOutput) WriteOne(ctx context.Context, msg stream.WritableMessage) error
- type SchemaRegistryAwareOutput_Expecter
- func (_e *SchemaRegistryAwareOutput_Expecter) InitSchemaRegistry(ctx interface{}, settings interface{}) *SchemaRegistryAwareOutput_InitSchemaRegistry_Call
- func (_e *SchemaRegistryAwareOutput_Expecter) Write(ctx interface{}, batch interface{}) *SchemaRegistryAwareOutput_Write_Call
- func (_e *SchemaRegistryAwareOutput_Expecter) WriteOne(ctx interface{}, msg interface{}) *SchemaRegistryAwareOutput_WriteOne_Call
- type SchemaRegistryAwareOutput_InitSchemaRegistry_Call
- func (_c *SchemaRegistryAwareOutput_InitSchemaRegistry_Call) Return(_a0 stream.MessageBodyEncoder, _a1 error) *SchemaRegistryAwareOutput_InitSchemaRegistry_Call
- func (_c *SchemaRegistryAwareOutput_InitSchemaRegistry_Call) Run(run func(ctx context.Context, settings stream.SchemaSettingsWithEncoding)) *SchemaRegistryAwareOutput_InitSchemaRegistry_Call
- func (_c *SchemaRegistryAwareOutput_InitSchemaRegistry_Call) RunAndReturn(...) *SchemaRegistryAwareOutput_InitSchemaRegistry_Call
- type SchemaRegistryAwareOutput_WriteOne_Call
- func (_c *SchemaRegistryAwareOutput_WriteOne_Call) Return(_a0 error) *SchemaRegistryAwareOutput_WriteOne_Call
- func (_c *SchemaRegistryAwareOutput_WriteOne_Call) Run(run func(ctx context.Context, msg stream.WritableMessage)) *SchemaRegistryAwareOutput_WriteOne_Call
- func (_c *SchemaRegistryAwareOutput_WriteOne_Call) RunAndReturn(run func(context.Context, stream.WritableMessage) error) *SchemaRegistryAwareOutput_WriteOne_Call
- type SchemaRegistryAwareOutput_Write_Call
- func (_c *SchemaRegistryAwareOutput_Write_Call) Return(_a0 error) *SchemaRegistryAwareOutput_Write_Call
- func (_c *SchemaRegistryAwareOutput_Write_Call) Run(run func(ctx context.Context, batch []stream.WritableMessage)) *SchemaRegistryAwareOutput_Write_Call
- func (_c *SchemaRegistryAwareOutput_Write_Call) RunAndReturn(run func(context.Context, []stream.WritableMessage) error) *SchemaRegistryAwareOutput_Write_Call
- type SchemaSettingsAwareCallback
- type SchemaSettingsAwareCallback_Expecter
- type SchemaSettingsAwareCallback_GetSchemaSettings_Call
- func (_c *SchemaSettingsAwareCallback_GetSchemaSettings_Call) Return(_a0 *stream.SchemaSettings, _a1 error) *SchemaSettingsAwareCallback_GetSchemaSettings_Call
- func (_c *SchemaSettingsAwareCallback_GetSchemaSettings_Call) Run(run func()) *SchemaSettingsAwareCallback_GetSchemaSettings_Call
- func (_c *SchemaSettingsAwareCallback_GetSchemaSettings_Call) RunAndReturn(run func() (*stream.SchemaSettings, error)) *SchemaSettingsAwareCallback_GetSchemaSettings_Call
- type UntypedConsumerCallback
- type UntypedConsumerCallback_Consume_Call
- func (_c *UntypedConsumerCallback_Consume_Call) Return(_a0 bool, _a1 error) *UntypedConsumerCallback_Consume_Call
- func (_c *UntypedConsumerCallback_Consume_Call) Run(run func(ctx context.Context, model interface{}, attributes map[string]string)) *UntypedConsumerCallback_Consume_Call
- func (_c *UntypedConsumerCallback_Consume_Call) RunAndReturn(run func(context.Context, interface{}, map[string]string) (bool, error)) *UntypedConsumerCallback_Consume_Call
- type UntypedConsumerCallback_Expecter
- type UntypedConsumerCallback_GetModel_Call
- func (_c *UntypedConsumerCallback_GetModel_Call) Return(_a0 interface{}, _a1 error) *UntypedConsumerCallback_GetModel_Call
- func (_c *UntypedConsumerCallback_GetModel_Call) Run(run func(attributes map[string]string)) *UntypedConsumerCallback_GetModel_Call
- func (_c *UntypedConsumerCallback_GetModel_Call) RunAndReturn(run func(map[string]string) (interface{}, error)) *UntypedConsumerCallback_GetModel_Call
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
This section is empty.
Types ¶
type ConsumerCallback ¶
ConsumerCallback is an autogenerated mock type for the ConsumerCallback type
func NewConsumerCallback ¶
func NewConsumerCallback[M interface{}](t interface {
mock.TestingT
Cleanup(func())
}) *ConsumerCallback[M]
NewConsumerCallback creates a new instance of ConsumerCallback. It also registers a testing interface on the mock and a cleanup function to assert the mocks expectations. The first argument is typically a *testing.T value.
func (*ConsumerCallback[M]) Consume ¶
func (_m *ConsumerCallback[M]) Consume(ctx context.Context, model M, attributes map[string]string) (bool, error)
Consume provides a mock function with given fields: ctx, model, attributes
func (*ConsumerCallback[M]) EXPECT ¶
func (_m *ConsumerCallback[M]) EXPECT() *ConsumerCallback_Expecter[M]
type ConsumerCallback_Consume_Call ¶
ConsumerCallback_Consume_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'Consume'
func (*ConsumerCallback_Consume_Call[M]) Return ¶
func (_c *ConsumerCallback_Consume_Call[M]) Return(_a0 bool, _a1 error) *ConsumerCallback_Consume_Call[M]
func (*ConsumerCallback_Consume_Call[M]) Run ¶
func (_c *ConsumerCallback_Consume_Call[M]) Run(run func(ctx context.Context, model M, attributes map[string]string)) *ConsumerCallback_Consume_Call[M]
func (*ConsumerCallback_Consume_Call[M]) RunAndReturn ¶
func (_c *ConsumerCallback_Consume_Call[M]) RunAndReturn(run func(context.Context, M, map[string]string) (bool, error)) *ConsumerCallback_Consume_Call[M]
type ConsumerCallback_Expecter ¶
type ConsumerCallback_Expecter[M interface{}] struct {
// contains filtered or unexported fields
}
func (*ConsumerCallback_Expecter[M]) Consume ¶
func (_e *ConsumerCallback_Expecter[M]) Consume(ctx interface{}, model interface{}, attributes interface{}) *ConsumerCallback_Consume_Call[M]
Consume is a helper method to define mock.On call
- ctx context.Context
- model M
- attributes map[string]string
type Input ¶
Input is an autogenerated mock type for the Input type
func NewInput ¶
NewInput creates a new instance of Input. It also registers a testing interface on the mock and a cleanup function to assert the mocks expectations. The first argument is typically a *testing.T value.
func (*Input) EXPECT ¶
func (_m *Input) EXPECT() *Input_Expecter
type Input_Expecter ¶
type Input_Expecter struct {
// contains filtered or unexported fields
}
func (*Input_Expecter) IsHealthy ¶ added in v0.40.17
func (_e *Input_Expecter) IsHealthy() *Input_IsHealthy_Call
IsHealthy is a helper method to define mock.On call
func (*Input_Expecter) Run ¶
func (_e *Input_Expecter) Run(ctx interface{}, process interface{}) *Input_Run_Call
Run is a helper method to define mock.On call
- ctx context.Context
- process stream.InputProcess
func (*Input_Expecter) Stop ¶
func (_e *Input_Expecter) Stop(ctx interface{}) *Input_Stop_Call
Stop is a helper method to define mock.On call
- ctx context.Context
type Input_IsHealthy_Call ¶ added in v0.40.17
Input_IsHealthy_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'IsHealthy'
func (*Input_IsHealthy_Call) Return ¶ added in v0.40.17
func (_c *Input_IsHealthy_Call) Return(_a0 bool) *Input_IsHealthy_Call
func (*Input_IsHealthy_Call) Run ¶ added in v0.40.17
func (_c *Input_IsHealthy_Call) Run(run func()) *Input_IsHealthy_Call
func (*Input_IsHealthy_Call) RunAndReturn ¶ added in v0.40.17
func (_c *Input_IsHealthy_Call) RunAndReturn(run func() bool) *Input_IsHealthy_Call
type Input_Run_Call ¶
Input_Run_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'Run'
func (*Input_Run_Call) Return ¶
func (_c *Input_Run_Call) Return(_a0 error) *Input_Run_Call
func (*Input_Run_Call) Run ¶
func (_c *Input_Run_Call) Run(run func(ctx context.Context, process stream.InputProcess)) *Input_Run_Call
func (*Input_Run_Call) RunAndReturn ¶
func (_c *Input_Run_Call) RunAndReturn(run func(context.Context, stream.InputProcess) error) *Input_Run_Call
type Input_Stop_Call ¶
Input_Stop_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'Stop'
func (*Input_Stop_Call) Return ¶
func (_c *Input_Stop_Call) Return() *Input_Stop_Call
func (*Input_Stop_Call) Run ¶
func (_c *Input_Stop_Call) Run(run func(ctx context.Context)) *Input_Stop_Call
func (*Input_Stop_Call) RunAndReturn ¶
func (_c *Input_Stop_Call) RunAndReturn(run func(context.Context)) *Input_Stop_Call
type KinsumerAutoscaleOrchestrator ¶ added in v0.34.1
KinsumerAutoscaleOrchestrator is an autogenerated mock type for the KinsumerAutoscaleOrchestrator type
func NewKinsumerAutoscaleOrchestrator ¶ added in v0.34.1
func NewKinsumerAutoscaleOrchestrator(t interface {
mock.TestingT
Cleanup(func())
}) *KinsumerAutoscaleOrchestrator
NewKinsumerAutoscaleOrchestrator creates a new instance of KinsumerAutoscaleOrchestrator. It also registers a testing interface on the mock and a cleanup function to assert the mocks expectations. The first argument is typically a *testing.T value.
func (*KinsumerAutoscaleOrchestrator) EXPECT ¶ added in v0.34.1
func (_m *KinsumerAutoscaleOrchestrator) EXPECT() *KinsumerAutoscaleOrchestrator_Expecter
func (*KinsumerAutoscaleOrchestrator) GetCurrentTaskCount ¶ added in v0.34.1
func (_m *KinsumerAutoscaleOrchestrator) GetCurrentTaskCount(ctx context.Context) (int32, error)
GetCurrentTaskCount provides a mock function with given fields: ctx
func (*KinsumerAutoscaleOrchestrator) UpdateTaskCount ¶ added in v0.34.1
func (_m *KinsumerAutoscaleOrchestrator) UpdateTaskCount(ctx context.Context, taskCount int32) error
UpdateTaskCount provides a mock function with given fields: ctx, taskCount
type KinsumerAutoscaleOrchestrator_Expecter ¶ added in v0.34.1
type KinsumerAutoscaleOrchestrator_Expecter struct {
// contains filtered or unexported fields
}
func (*KinsumerAutoscaleOrchestrator_Expecter) GetCurrentTaskCount ¶ added in v0.34.1
func (_e *KinsumerAutoscaleOrchestrator_Expecter) GetCurrentTaskCount(ctx interface{}) *KinsumerAutoscaleOrchestrator_GetCurrentTaskCount_Call
GetCurrentTaskCount is a helper method to define mock.On call
- ctx context.Context
func (*KinsumerAutoscaleOrchestrator_Expecter) UpdateTaskCount ¶ added in v0.34.1
func (_e *KinsumerAutoscaleOrchestrator_Expecter) UpdateTaskCount(ctx interface{}, taskCount interface{}) *KinsumerAutoscaleOrchestrator_UpdateTaskCount_Call
UpdateTaskCount is a helper method to define mock.On call
- ctx context.Context
- taskCount int32
type KinsumerAutoscaleOrchestrator_GetCurrentTaskCount_Call ¶ added in v0.34.1
KinsumerAutoscaleOrchestrator_GetCurrentTaskCount_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'GetCurrentTaskCount'
func (*KinsumerAutoscaleOrchestrator_GetCurrentTaskCount_Call) Return ¶ added in v0.34.1
func (_c *KinsumerAutoscaleOrchestrator_GetCurrentTaskCount_Call) Return(_a0 int32, _a1 error) *KinsumerAutoscaleOrchestrator_GetCurrentTaskCount_Call
func (*KinsumerAutoscaleOrchestrator_GetCurrentTaskCount_Call) Run ¶ added in v0.34.1
func (_c *KinsumerAutoscaleOrchestrator_GetCurrentTaskCount_Call) Run(run func(ctx context.Context)) *KinsumerAutoscaleOrchestrator_GetCurrentTaskCount_Call
func (*KinsumerAutoscaleOrchestrator_GetCurrentTaskCount_Call) RunAndReturn ¶ added in v0.34.1
func (_c *KinsumerAutoscaleOrchestrator_GetCurrentTaskCount_Call) RunAndReturn(run func(context.Context) (int32, error)) *KinsumerAutoscaleOrchestrator_GetCurrentTaskCount_Call
type KinsumerAutoscaleOrchestrator_UpdateTaskCount_Call ¶ added in v0.34.1
KinsumerAutoscaleOrchestrator_UpdateTaskCount_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'UpdateTaskCount'
func (*KinsumerAutoscaleOrchestrator_UpdateTaskCount_Call) Run ¶ added in v0.34.1
func (_c *KinsumerAutoscaleOrchestrator_UpdateTaskCount_Call) Run(run func(ctx context.Context, taskCount int32)) *KinsumerAutoscaleOrchestrator_UpdateTaskCount_Call
func (*KinsumerAutoscaleOrchestrator_UpdateTaskCount_Call) RunAndReturn ¶ added in v0.34.1
func (_c *KinsumerAutoscaleOrchestrator_UpdateTaskCount_Call) RunAndReturn(run func(context.Context, int32) error) *KinsumerAutoscaleOrchestrator_UpdateTaskCount_Call
type MessageEncoder ¶
MessageEncoder is an autogenerated mock type for the MessageEncoder type
func NewMessageEncoder ¶
func NewMessageEncoder(t interface {
mock.TestingT
Cleanup(func())
}) *MessageEncoder
NewMessageEncoder creates a new instance of MessageEncoder. It also registers a testing interface on the mock and a cleanup function to assert the mocks expectations. The first argument is typically a *testing.T value.
func (*MessageEncoder) Decode ¶
func (_m *MessageEncoder) Decode(ctx context.Context, msg *stream.Message, out interface{}) (context.Context, map[string]string, error)
Decode provides a mock function with given fields: ctx, msg, out
func (*MessageEncoder) EXPECT ¶
func (_m *MessageEncoder) EXPECT() *MessageEncoder_Expecter
type MessageEncoder_Decode_Call ¶
MessageEncoder_Decode_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'Decode'
func (*MessageEncoder_Decode_Call) Return ¶
func (_c *MessageEncoder_Decode_Call) Return(_a0 context.Context, _a1 map[string]string, _a2 error) *MessageEncoder_Decode_Call
func (*MessageEncoder_Decode_Call) Run ¶
func (_c *MessageEncoder_Decode_Call) Run(run func(ctx context.Context, msg *stream.Message, out interface{})) *MessageEncoder_Decode_Call
func (*MessageEncoder_Decode_Call) RunAndReturn ¶
type MessageEncoder_Encode_Call ¶
MessageEncoder_Encode_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'Encode'
func (*MessageEncoder_Encode_Call) Return ¶
func (_c *MessageEncoder_Encode_Call) Return(_a0 *stream.Message, _a1 error) *MessageEncoder_Encode_Call
func (*MessageEncoder_Encode_Call) Run ¶
func (_c *MessageEncoder_Encode_Call) Run(run func(ctx context.Context, data interface{}, attributeSets ...map[string]string)) *MessageEncoder_Encode_Call
func (*MessageEncoder_Encode_Call) RunAndReturn ¶
func (_c *MessageEncoder_Encode_Call) RunAndReturn(run func(context.Context, interface{}, ...map[string]string) (*stream.Message, error)) *MessageEncoder_Encode_Call
type MessageEncoder_Expecter ¶
type MessageEncoder_Expecter struct {
// contains filtered or unexported fields
}
func (*MessageEncoder_Expecter) Decode ¶
func (_e *MessageEncoder_Expecter) Decode(ctx interface{}, msg interface{}, out interface{}) *MessageEncoder_Decode_Call
Decode is a helper method to define mock.On call
- ctx context.Context
- msg *stream.Message
- out interface{}
func (*MessageEncoder_Expecter) Encode ¶
func (_e *MessageEncoder_Expecter) Encode(ctx interface{}, data interface{}, attributeSets ...interface{}) *MessageEncoder_Encode_Call
Encode is a helper method to define mock.On call
- ctx context.Context
- data interface{}
- attributeSets ...map[string]string
type Output ¶
Output is an autogenerated mock type for the Output type
func NewOutput ¶
NewOutput creates a new instance of Output. It also registers a testing interface on the mock and a cleanup function to assert the mocks expectations. The first argument is typically a *testing.T value.
func (*Output) EXPECT ¶
func (_m *Output) EXPECT() *Output_Expecter
type Output_Expecter ¶
type Output_Expecter struct {
// contains filtered or unexported fields
}
func (*Output_Expecter) Write ¶
func (_e *Output_Expecter) Write(ctx interface{}, batch interface{}) *Output_Write_Call
Write is a helper method to define mock.On call
- ctx context.Context
- batch []stream.WritableMessage
func (*Output_Expecter) WriteOne ¶
func (_e *Output_Expecter) WriteOne(ctx interface{}, msg interface{}) *Output_WriteOne_Call
WriteOne is a helper method to define mock.On call
- ctx context.Context
- msg stream.WritableMessage
type Output_WriteOne_Call ¶
Output_WriteOne_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'WriteOne'
func (*Output_WriteOne_Call) Return ¶
func (_c *Output_WriteOne_Call) Return(_a0 error) *Output_WriteOne_Call
func (*Output_WriteOne_Call) Run ¶
func (_c *Output_WriteOne_Call) Run(run func(ctx context.Context, msg stream.WritableMessage)) *Output_WriteOne_Call
func (*Output_WriteOne_Call) RunAndReturn ¶
func (_c *Output_WriteOne_Call) RunAndReturn(run func(context.Context, stream.WritableMessage) error) *Output_WriteOne_Call
type Output_Write_Call ¶
Output_Write_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'Write'
func (*Output_Write_Call) Return ¶
func (_c *Output_Write_Call) Return(_a0 error) *Output_Write_Call
func (*Output_Write_Call) Run ¶
func (_c *Output_Write_Call) Run(run func(ctx context.Context, batch []stream.WritableMessage)) *Output_Write_Call
func (*Output_Write_Call) RunAndReturn ¶
func (_c *Output_Write_Call) RunAndReturn(run func(context.Context, []stream.WritableMessage) error) *Output_Write_Call
type PartitionerRand ¶
PartitionerRand is an autogenerated mock type for the PartitionerRand type
func NewPartitionerRand ¶
func NewPartitionerRand(t interface {
mock.TestingT
Cleanup(func())
}) *PartitionerRand
NewPartitionerRand creates a new instance of PartitionerRand. It also registers a testing interface on the mock and a cleanup function to assert the mocks expectations. The first argument is typically a *testing.T value.
func (*PartitionerRand) EXPECT ¶
func (_m *PartitionerRand) EXPECT() *PartitionerRand_Expecter
func (*PartitionerRand) Intn ¶
func (_m *PartitionerRand) Intn(n int) int
Intn provides a mock function with given fields: n
type PartitionerRand_Expecter ¶
type PartitionerRand_Expecter struct {
// contains filtered or unexported fields
}
func (*PartitionerRand_Expecter) Intn ¶
func (_e *PartitionerRand_Expecter) Intn(n interface{}) *PartitionerRand_Intn_Call
Intn is a helper method to define mock.On call
- n int
type PartitionerRand_Intn_Call ¶
PartitionerRand_Intn_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'Intn'
func (*PartitionerRand_Intn_Call) Return ¶
func (_c *PartitionerRand_Intn_Call) Return(_a0 int) *PartitionerRand_Intn_Call
func (*PartitionerRand_Intn_Call) Run ¶
func (_c *PartitionerRand_Intn_Call) Run(run func(n int)) *PartitionerRand_Intn_Call
func (*PartitionerRand_Intn_Call) RunAndReturn ¶
func (_c *PartitionerRand_Intn_Call) RunAndReturn(run func(int) int) *PartitionerRand_Intn_Call
type Producer ¶
Producer is an autogenerated mock type for the Producer type
func NewProducer ¶
NewProducer creates a new instance of Producer. It also registers a testing interface on the mock and a cleanup function to assert the mocks expectations. The first argument is typically a *testing.T value.
func (*Producer) EXPECT ¶
func (_m *Producer) EXPECT() *Producer_Expecter
type ProducerDaemonAggregator ¶
ProducerDaemonAggregator is an autogenerated mock type for the ProducerDaemonAggregator type
func NewProducerDaemonAggregator ¶
func NewProducerDaemonAggregator(t interface {
mock.TestingT
Cleanup(func())
}) *ProducerDaemonAggregator
NewProducerDaemonAggregator creates a new instance of ProducerDaemonAggregator. It also registers a testing interface on the mock and a cleanup function to assert the mocks expectations. The first argument is typically a *testing.T value.
func (*ProducerDaemonAggregator) EXPECT ¶
func (_m *ProducerDaemonAggregator) EXPECT() *ProducerDaemonAggregator_Expecter
func (*ProducerDaemonAggregator) Flush ¶
func (_m *ProducerDaemonAggregator) Flush() ([]stream.AggregateFlush, error)
Flush provides a mock function with no fields
func (*ProducerDaemonAggregator) Write ¶
func (_m *ProducerDaemonAggregator) Write(ctx context.Context, msg *stream.Message) ([]stream.AggregateFlush, error)
Write provides a mock function with given fields: ctx, msg
type ProducerDaemonAggregator_Expecter ¶
type ProducerDaemonAggregator_Expecter struct {
// contains filtered or unexported fields
}
func (*ProducerDaemonAggregator_Expecter) Flush ¶
func (_e *ProducerDaemonAggregator_Expecter) Flush() *ProducerDaemonAggregator_Flush_Call
Flush is a helper method to define mock.On call
func (*ProducerDaemonAggregator_Expecter) Write ¶
func (_e *ProducerDaemonAggregator_Expecter) Write(ctx interface{}, msg interface{}) *ProducerDaemonAggregator_Write_Call
Write is a helper method to define mock.On call
- ctx context.Context
- msg *stream.Message
type ProducerDaemonAggregator_Flush_Call ¶
ProducerDaemonAggregator_Flush_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'Flush'
func (*ProducerDaemonAggregator_Flush_Call) Return ¶
func (_c *ProducerDaemonAggregator_Flush_Call) Return(_a0 []stream.AggregateFlush, _a1 error) *ProducerDaemonAggregator_Flush_Call
func (*ProducerDaemonAggregator_Flush_Call) Run ¶
func (_c *ProducerDaemonAggregator_Flush_Call) Run(run func()) *ProducerDaemonAggregator_Flush_Call
func (*ProducerDaemonAggregator_Flush_Call) RunAndReturn ¶
func (_c *ProducerDaemonAggregator_Flush_Call) RunAndReturn(run func() ([]stream.AggregateFlush, error)) *ProducerDaemonAggregator_Flush_Call
type ProducerDaemonAggregator_Write_Call ¶
ProducerDaemonAggregator_Write_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'Write'
func (*ProducerDaemonAggregator_Write_Call) Return ¶
func (_c *ProducerDaemonAggregator_Write_Call) Return(_a0 []stream.AggregateFlush, _a1 error) *ProducerDaemonAggregator_Write_Call
func (*ProducerDaemonAggregator_Write_Call) Run ¶
func (_c *ProducerDaemonAggregator_Write_Call) Run(run func(ctx context.Context, msg *stream.Message)) *ProducerDaemonAggregator_Write_Call
func (*ProducerDaemonAggregator_Write_Call) RunAndReturn ¶
func (_c *ProducerDaemonAggregator_Write_Call) RunAndReturn(run func(context.Context, *stream.Message) ([]stream.AggregateFlush, error)) *ProducerDaemonAggregator_Write_Call
type ProducerDaemonBatcher ¶
ProducerDaemonBatcher is an autogenerated mock type for the ProducerDaemonBatcher type
func NewProducerDaemonBatcher ¶
func NewProducerDaemonBatcher(t interface {
mock.TestingT
Cleanup(func())
}) *ProducerDaemonBatcher
NewProducerDaemonBatcher creates a new instance of ProducerDaemonBatcher. It also registers a testing interface on the mock and a cleanup function to assert the mocks expectations. The first argument is typically a *testing.T value.
func (*ProducerDaemonBatcher) Append ¶
func (_m *ProducerDaemonBatcher) Append(msg *stream.Message) ([]stream.WritableMessage, error)
Append provides a mock function with given fields: msg
func (*ProducerDaemonBatcher) EXPECT ¶
func (_m *ProducerDaemonBatcher) EXPECT() *ProducerDaemonBatcher_Expecter
func (*ProducerDaemonBatcher) Flush ¶
func (_m *ProducerDaemonBatcher) Flush() []stream.WritableMessage
Flush provides a mock function with no fields
type ProducerDaemonBatcher_Append_Call ¶
ProducerDaemonBatcher_Append_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'Append'
func (*ProducerDaemonBatcher_Append_Call) Return ¶
func (_c *ProducerDaemonBatcher_Append_Call) Return(_a0 []stream.WritableMessage, _a1 error) *ProducerDaemonBatcher_Append_Call
func (*ProducerDaemonBatcher_Append_Call) Run ¶
func (_c *ProducerDaemonBatcher_Append_Call) Run(run func(msg *stream.Message)) *ProducerDaemonBatcher_Append_Call
func (*ProducerDaemonBatcher_Append_Call) RunAndReturn ¶
func (_c *ProducerDaemonBatcher_Append_Call) RunAndReturn(run func(*stream.Message) ([]stream.WritableMessage, error)) *ProducerDaemonBatcher_Append_Call
type ProducerDaemonBatcher_Expecter ¶
type ProducerDaemonBatcher_Expecter struct {
// contains filtered or unexported fields
}
func (*ProducerDaemonBatcher_Expecter) Append ¶
func (_e *ProducerDaemonBatcher_Expecter) Append(msg interface{}) *ProducerDaemonBatcher_Append_Call
Append is a helper method to define mock.On call
- msg *stream.Message
func (*ProducerDaemonBatcher_Expecter) Flush ¶
func (_e *ProducerDaemonBatcher_Expecter) Flush() *ProducerDaemonBatcher_Flush_Call
Flush is a helper method to define mock.On call
type ProducerDaemonBatcher_Flush_Call ¶
ProducerDaemonBatcher_Flush_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'Flush'
func (*ProducerDaemonBatcher_Flush_Call) Return ¶
func (_c *ProducerDaemonBatcher_Flush_Call) Return(_a0 []stream.WritableMessage) *ProducerDaemonBatcher_Flush_Call
func (*ProducerDaemonBatcher_Flush_Call) Run ¶
func (_c *ProducerDaemonBatcher_Flush_Call) Run(run func()) *ProducerDaemonBatcher_Flush_Call
func (*ProducerDaemonBatcher_Flush_Call) RunAndReturn ¶
func (_c *ProducerDaemonBatcher_Flush_Call) RunAndReturn(run func() []stream.WritableMessage) *ProducerDaemonBatcher_Flush_Call
type Producer_Expecter ¶
type Producer_Expecter struct {
// contains filtered or unexported fields
}
func (*Producer_Expecter) Write ¶
func (_e *Producer_Expecter) Write(ctx interface{}, models interface{}, attributeSets ...interface{}) *Producer_Write_Call
Write is a helper method to define mock.On call
- ctx context.Context
- models interface{}
- attributeSets ...map[string]string
func (*Producer_Expecter) WriteOne ¶
func (_e *Producer_Expecter) WriteOne(ctx interface{}, model interface{}, attributeSets ...interface{}) *Producer_WriteOne_Call
WriteOne is a helper method to define mock.On call
- ctx context.Context
- model interface{}
- attributeSets ...map[string]string
type Producer_WriteOne_Call ¶
Producer_WriteOne_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'WriteOne'
func (*Producer_WriteOne_Call) Return ¶
func (_c *Producer_WriteOne_Call) Return(_a0 error) *Producer_WriteOne_Call
func (*Producer_WriteOne_Call) Run ¶
func (_c *Producer_WriteOne_Call) Run(run func(ctx context.Context, model interface{}, attributeSets ...map[string]string)) *Producer_WriteOne_Call
func (*Producer_WriteOne_Call) RunAndReturn ¶
func (_c *Producer_WriteOne_Call) RunAndReturn(run func(context.Context, interface{}, ...map[string]string) error) *Producer_WriteOne_Call
type Producer_Write_Call ¶
Producer_Write_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'Write'
func (*Producer_Write_Call) Return ¶
func (_c *Producer_Write_Call) Return(_a0 error) *Producer_Write_Call
func (*Producer_Write_Call) Run ¶
func (_c *Producer_Write_Call) Run(run func(ctx context.Context, models interface{}, attributeSets ...map[string]string)) *Producer_Write_Call
func (*Producer_Write_Call) RunAndReturn ¶
func (_c *Producer_Write_Call) RunAndReturn(run func(context.Context, interface{}, ...map[string]string) error) *Producer_Write_Call
type RetryHandler ¶
RetryHandler is an autogenerated mock type for the RetryHandler type
func NewRetryHandler ¶
func NewRetryHandler(t interface {
mock.TestingT
Cleanup(func())
}) *RetryHandler
NewRetryHandler creates a new instance of RetryHandler. It also registers a testing interface on the mock and a cleanup function to assert the mocks expectations. The first argument is typically a *testing.T value.
func (*RetryHandler) EXPECT ¶
func (_m *RetryHandler) EXPECT() *RetryHandler_Expecter
type RetryHandler_Expecter ¶
type RetryHandler_Expecter struct {
// contains filtered or unexported fields
}
func (*RetryHandler_Expecter) Put ¶
func (_e *RetryHandler_Expecter) Put(ctx interface{}, msg interface{}) *RetryHandler_Put_Call
Put is a helper method to define mock.On call
- ctx context.Context
- msg *stream.Message
type RetryHandler_Put_Call ¶
RetryHandler_Put_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'Put'
func (*RetryHandler_Put_Call) Return ¶
func (_c *RetryHandler_Put_Call) Return(_a0 error) *RetryHandler_Put_Call
func (*RetryHandler_Put_Call) Run ¶
func (_c *RetryHandler_Put_Call) Run(run func(ctx context.Context, msg *stream.Message)) *RetryHandler_Put_Call
func (*RetryHandler_Put_Call) RunAndReturn ¶
func (_c *RetryHandler_Put_Call) RunAndReturn(run func(context.Context, *stream.Message) error) *RetryHandler_Put_Call
type RunnableCallback ¶
RunnableCallback is an autogenerated mock type for the RunnableCallback type
func NewRunnableCallback ¶
func NewRunnableCallback(t interface {
mock.TestingT
Cleanup(func())
}) *RunnableCallback
NewRunnableCallback creates a new instance of RunnableCallback. It also registers a testing interface on the mock and a cleanup function to assert the mocks expectations. The first argument is typically a *testing.T value.
func (*RunnableCallback) EXPECT ¶
func (_m *RunnableCallback) EXPECT() *RunnableCallback_Expecter
type RunnableCallback_Expecter ¶
type RunnableCallback_Expecter struct {
// contains filtered or unexported fields
}
func (*RunnableCallback_Expecter) Run ¶
func (_e *RunnableCallback_Expecter) Run(ctx interface{}) *RunnableCallback_Run_Call
Run is a helper method to define mock.On call
- ctx context.Context
type RunnableCallback_Run_Call ¶
RunnableCallback_Run_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'Run'
func (*RunnableCallback_Run_Call) Return ¶
func (_c *RunnableCallback_Run_Call) Return(_a0 error) *RunnableCallback_Run_Call
func (*RunnableCallback_Run_Call) Run ¶
func (_c *RunnableCallback_Run_Call) Run(run func(ctx context.Context)) *RunnableCallback_Run_Call
func (*RunnableCallback_Run_Call) RunAndReturn ¶
func (_c *RunnableCallback_Run_Call) RunAndReturn(run func(context.Context) error) *RunnableCallback_Run_Call
type RunnableConsumerCallback ¶
RunnableConsumerCallback is an autogenerated mock type for the RunnableConsumerCallback type
func NewRunnableConsumerCallback ¶
func NewRunnableConsumerCallback[M interface{}](t interface {
mock.TestingT
Cleanup(func())
}) *RunnableConsumerCallback[M]
NewRunnableConsumerCallback creates a new instance of RunnableConsumerCallback. It also registers a testing interface on the mock and a cleanup function to assert the mocks expectations. The first argument is typically a *testing.T value.
func (*RunnableConsumerCallback[M]) Consume ¶
func (_m *RunnableConsumerCallback[M]) Consume(ctx context.Context, model M, attributes map[string]string) (bool, error)
Consume provides a mock function with given fields: ctx, model, attributes
func (*RunnableConsumerCallback[M]) EXPECT ¶
func (_m *RunnableConsumerCallback[M]) EXPECT() *RunnableConsumerCallback_Expecter[M]
type RunnableConsumerCallback_Consume_Call ¶
RunnableConsumerCallback_Consume_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'Consume'
func (*RunnableConsumerCallback_Consume_Call[M]) Return ¶
func (_c *RunnableConsumerCallback_Consume_Call[M]) Return(_a0 bool, _a1 error) *RunnableConsumerCallback_Consume_Call[M]
func (*RunnableConsumerCallback_Consume_Call[M]) Run ¶
func (_c *RunnableConsumerCallback_Consume_Call[M]) Run(run func(ctx context.Context, model M, attributes map[string]string)) *RunnableConsumerCallback_Consume_Call[M]
func (*RunnableConsumerCallback_Consume_Call[M]) RunAndReturn ¶
func (_c *RunnableConsumerCallback_Consume_Call[M]) RunAndReturn(run func(context.Context, M, map[string]string) (bool, error)) *RunnableConsumerCallback_Consume_Call[M]
type RunnableConsumerCallback_Expecter ¶
type RunnableConsumerCallback_Expecter[M interface{}] struct {
// contains filtered or unexported fields
}
func (*RunnableConsumerCallback_Expecter[M]) Consume ¶
func (_e *RunnableConsumerCallback_Expecter[M]) Consume(ctx interface{}, model interface{}, attributes interface{}) *RunnableConsumerCallback_Consume_Call[M]
Consume is a helper method to define mock.On call
- ctx context.Context
- model M
- attributes map[string]string
func (*RunnableConsumerCallback_Expecter[M]) Run ¶
func (_e *RunnableConsumerCallback_Expecter[M]) Run(ctx interface{}) *RunnableConsumerCallback_Run_Call[M]
Run is a helper method to define mock.On call
- ctx context.Context
type RunnableConsumerCallback_Run_Call ¶
RunnableConsumerCallback_Run_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'Run'
func (*RunnableConsumerCallback_Run_Call[M]) Return ¶
func (_c *RunnableConsumerCallback_Run_Call[M]) Return(_a0 error) *RunnableConsumerCallback_Run_Call[M]
func (*RunnableConsumerCallback_Run_Call[M]) Run ¶
func (_c *RunnableConsumerCallback_Run_Call[M]) Run(run func(ctx context.Context)) *RunnableConsumerCallback_Run_Call[M]
func (*RunnableConsumerCallback_Run_Call[M]) RunAndReturn ¶
func (_c *RunnableConsumerCallback_Run_Call[M]) RunAndReturn(run func(context.Context) error) *RunnableConsumerCallback_Run_Call[M]
type RunnableUntypedConsumerCallback ¶ added in v0.41.0
RunnableUntypedConsumerCallback is an autogenerated mock type for the RunnableUntypedConsumerCallback type
func NewRunnableUntypedConsumerCallback ¶ added in v0.41.0
func NewRunnableUntypedConsumerCallback(t interface {
mock.TestingT
Cleanup(func())
}) *RunnableUntypedConsumerCallback
NewRunnableUntypedConsumerCallback creates a new instance of RunnableUntypedConsumerCallback. It also registers a testing interface on the mock and a cleanup function to assert the mocks expectations. The first argument is typically a *testing.T value.
func (*RunnableUntypedConsumerCallback) Consume ¶ added in v0.41.0
func (_m *RunnableUntypedConsumerCallback) Consume(ctx context.Context, model interface{}, attributes map[string]string) (bool, error)
Consume provides a mock function with given fields: ctx, model, attributes
func (*RunnableUntypedConsumerCallback) EXPECT ¶ added in v0.41.0
func (_m *RunnableUntypedConsumerCallback) EXPECT() *RunnableUntypedConsumerCallback_Expecter
type RunnableUntypedConsumerCallback_Consume_Call ¶ added in v0.41.0
RunnableUntypedConsumerCallback_Consume_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'Consume'
func (*RunnableUntypedConsumerCallback_Consume_Call) Return ¶ added in v0.41.0
func (_c *RunnableUntypedConsumerCallback_Consume_Call) Return(_a0 bool, _a1 error) *RunnableUntypedConsumerCallback_Consume_Call
func (*RunnableUntypedConsumerCallback_Consume_Call) Run ¶ added in v0.41.0
func (_c *RunnableUntypedConsumerCallback_Consume_Call) Run(run func(ctx context.Context, model interface{}, attributes map[string]string)) *RunnableUntypedConsumerCallback_Consume_Call
func (*RunnableUntypedConsumerCallback_Consume_Call) RunAndReturn ¶ added in v0.41.0
func (_c *RunnableUntypedConsumerCallback_Consume_Call) RunAndReturn(run func(context.Context, interface{}, map[string]string) (bool, error)) *RunnableUntypedConsumerCallback_Consume_Call
type RunnableUntypedConsumerCallback_Expecter ¶ added in v0.41.0
type RunnableUntypedConsumerCallback_Expecter struct {
// contains filtered or unexported fields
}
func (*RunnableUntypedConsumerCallback_Expecter) Consume ¶ added in v0.41.0
func (_e *RunnableUntypedConsumerCallback_Expecter) Consume(ctx interface{}, model interface{}, attributes interface{}) *RunnableUntypedConsumerCallback_Consume_Call
Consume is a helper method to define mock.On call
- ctx context.Context
- model interface{}
- attributes map[string]string
func (*RunnableUntypedConsumerCallback_Expecter) GetModel ¶ added in v0.41.0
func (_e *RunnableUntypedConsumerCallback_Expecter) GetModel(attributes interface{}) *RunnableUntypedConsumerCallback_GetModel_Call
GetModel is a helper method to define mock.On call
- attributes map[string]string
func (*RunnableUntypedConsumerCallback_Expecter) Run ¶ added in v0.41.0
func (_e *RunnableUntypedConsumerCallback_Expecter) Run(ctx interface{}) *RunnableUntypedConsumerCallback_Run_Call
Run is a helper method to define mock.On call
- ctx context.Context
type RunnableUntypedConsumerCallback_GetModel_Call ¶ added in v0.41.0
RunnableUntypedConsumerCallback_GetModel_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'GetModel'
func (*RunnableUntypedConsumerCallback_GetModel_Call) Return ¶ added in v0.41.0
func (_c *RunnableUntypedConsumerCallback_GetModel_Call) Return(_a0 interface{}, _a1 error) *RunnableUntypedConsumerCallback_GetModel_Call
func (*RunnableUntypedConsumerCallback_GetModel_Call) Run ¶ added in v0.41.0
func (_c *RunnableUntypedConsumerCallback_GetModel_Call) Run(run func(attributes map[string]string)) *RunnableUntypedConsumerCallback_GetModel_Call
func (*RunnableUntypedConsumerCallback_GetModel_Call) RunAndReturn ¶ added in v0.41.0
func (_c *RunnableUntypedConsumerCallback_GetModel_Call) RunAndReturn(run func(map[string]string) (interface{}, error)) *RunnableUntypedConsumerCallback_GetModel_Call
type RunnableUntypedConsumerCallback_Run_Call ¶ added in v0.41.0
RunnableUntypedConsumerCallback_Run_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'Run'
func (*RunnableUntypedConsumerCallback_Run_Call) Return ¶ added in v0.41.0
func (_c *RunnableUntypedConsumerCallback_Run_Call) Return(_a0 error) *RunnableUntypedConsumerCallback_Run_Call
func (*RunnableUntypedConsumerCallback_Run_Call) Run ¶ added in v0.41.0
func (_c *RunnableUntypedConsumerCallback_Run_Call) Run(run func(ctx context.Context)) *RunnableUntypedConsumerCallback_Run_Call
func (*RunnableUntypedConsumerCallback_Run_Call) RunAndReturn ¶ added in v0.41.0
func (_c *RunnableUntypedConsumerCallback_Run_Call) RunAndReturn(run func(context.Context) error) *RunnableUntypedConsumerCallback_Run_Call
type SchemaRegistryAwareInput ¶ added in v0.51.0
SchemaRegistryAwareInput is an autogenerated mock type for the SchemaRegistryAwareInput type
func NewSchemaRegistryAwareInput ¶ added in v0.51.0
func NewSchemaRegistryAwareInput(t interface {
mock.TestingT
Cleanup(func())
}) *SchemaRegistryAwareInput
NewSchemaRegistryAwareInput creates a new instance of SchemaRegistryAwareInput. It also registers a testing interface on the mock and a cleanup function to assert the mocks expectations. The first argument is typically a *testing.T value.
func (*SchemaRegistryAwareInput) EXPECT ¶ added in v0.51.0
func (_m *SchemaRegistryAwareInput) EXPECT() *SchemaRegistryAwareInput_Expecter
func (*SchemaRegistryAwareInput) InitSchemaRegistry ¶ added in v0.51.0
func (_m *SchemaRegistryAwareInput) InitSchemaRegistry(ctx context.Context, settings stream.SchemaSettingsWithEncoding) (stream.MessageBodyEncoder, error)
InitSchemaRegistry provides a mock function with given fields: ctx, settings
func (*SchemaRegistryAwareInput) IsHealthy ¶ added in v0.51.0
func (_m *SchemaRegistryAwareInput) IsHealthy() bool
IsHealthy provides a mock function with no fields
func (*SchemaRegistryAwareInput) Run ¶ added in v0.51.0
func (_m *SchemaRegistryAwareInput) Run(ctx context.Context, process stream.InputProcess) error
Run provides a mock function with given fields: ctx, process
func (*SchemaRegistryAwareInput) Stop ¶ added in v0.51.0
func (_m *SchemaRegistryAwareInput) Stop(ctx context.Context)
Stop provides a mock function with given fields: ctx
type SchemaRegistryAwareInput_Expecter ¶ added in v0.51.0
type SchemaRegistryAwareInput_Expecter struct {
// contains filtered or unexported fields
}
func (*SchemaRegistryAwareInput_Expecter) InitSchemaRegistry ¶ added in v0.51.0
func (_e *SchemaRegistryAwareInput_Expecter) InitSchemaRegistry(ctx interface{}, settings interface{}) *SchemaRegistryAwareInput_InitSchemaRegistry_Call
InitSchemaRegistry is a helper method to define mock.On call
- ctx context.Context
- settings stream.SchemaSettingsWithEncoding
func (*SchemaRegistryAwareInput_Expecter) IsHealthy ¶ added in v0.51.0
func (_e *SchemaRegistryAwareInput_Expecter) IsHealthy() *SchemaRegistryAwareInput_IsHealthy_Call
IsHealthy is a helper method to define mock.On call
func (*SchemaRegistryAwareInput_Expecter) Run ¶ added in v0.51.0
func (_e *SchemaRegistryAwareInput_Expecter) Run(ctx interface{}, process interface{}) *SchemaRegistryAwareInput_Run_Call
Run is a helper method to define mock.On call
- ctx context.Context
- process stream.InputProcess
func (*SchemaRegistryAwareInput_Expecter) Stop ¶ added in v0.51.0
func (_e *SchemaRegistryAwareInput_Expecter) Stop(ctx interface{}) *SchemaRegistryAwareInput_Stop_Call
Stop is a helper method to define mock.On call
- ctx context.Context
type SchemaRegistryAwareInput_InitSchemaRegistry_Call ¶ added in v0.51.0
SchemaRegistryAwareInput_InitSchemaRegistry_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'InitSchemaRegistry'
func (*SchemaRegistryAwareInput_InitSchemaRegistry_Call) Return ¶ added in v0.51.0
func (_c *SchemaRegistryAwareInput_InitSchemaRegistry_Call) Return(_a0 stream.MessageBodyEncoder, _a1 error) *SchemaRegistryAwareInput_InitSchemaRegistry_Call
func (*SchemaRegistryAwareInput_InitSchemaRegistry_Call) Run ¶ added in v0.51.0
func (_c *SchemaRegistryAwareInput_InitSchemaRegistry_Call) Run(run func(ctx context.Context, settings stream.SchemaSettingsWithEncoding)) *SchemaRegistryAwareInput_InitSchemaRegistry_Call
func (*SchemaRegistryAwareInput_InitSchemaRegistry_Call) RunAndReturn ¶ added in v0.51.0
func (_c *SchemaRegistryAwareInput_InitSchemaRegistry_Call) RunAndReturn(run func(context.Context, stream.SchemaSettingsWithEncoding) (stream.MessageBodyEncoder, error)) *SchemaRegistryAwareInput_InitSchemaRegistry_Call
type SchemaRegistryAwareInput_IsHealthy_Call ¶ added in v0.51.0
SchemaRegistryAwareInput_IsHealthy_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'IsHealthy'
func (*SchemaRegistryAwareInput_IsHealthy_Call) Return ¶ added in v0.51.0
func (_c *SchemaRegistryAwareInput_IsHealthy_Call) Return(_a0 bool) *SchemaRegistryAwareInput_IsHealthy_Call
func (*SchemaRegistryAwareInput_IsHealthy_Call) Run ¶ added in v0.51.0
func (_c *SchemaRegistryAwareInput_IsHealthy_Call) Run(run func()) *SchemaRegistryAwareInput_IsHealthy_Call
func (*SchemaRegistryAwareInput_IsHealthy_Call) RunAndReturn ¶ added in v0.51.0
func (_c *SchemaRegistryAwareInput_IsHealthy_Call) RunAndReturn(run func() bool) *SchemaRegistryAwareInput_IsHealthy_Call
type SchemaRegistryAwareInput_Run_Call ¶ added in v0.51.0
SchemaRegistryAwareInput_Run_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'Run'
func (*SchemaRegistryAwareInput_Run_Call) Return ¶ added in v0.51.0
func (_c *SchemaRegistryAwareInput_Run_Call) Return(_a0 error) *SchemaRegistryAwareInput_Run_Call
func (*SchemaRegistryAwareInput_Run_Call) Run ¶ added in v0.51.0
func (_c *SchemaRegistryAwareInput_Run_Call) Run(run func(ctx context.Context, process stream.InputProcess)) *SchemaRegistryAwareInput_Run_Call
func (*SchemaRegistryAwareInput_Run_Call) RunAndReturn ¶ added in v0.51.0
func (_c *SchemaRegistryAwareInput_Run_Call) RunAndReturn(run func(context.Context, stream.InputProcess) error) *SchemaRegistryAwareInput_Run_Call
type SchemaRegistryAwareInput_Stop_Call ¶ added in v0.51.0
SchemaRegistryAwareInput_Stop_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'Stop'
func (*SchemaRegistryAwareInput_Stop_Call) Return ¶ added in v0.51.0
func (_c *SchemaRegistryAwareInput_Stop_Call) Return() *SchemaRegistryAwareInput_Stop_Call
func (*SchemaRegistryAwareInput_Stop_Call) Run ¶ added in v0.51.0
func (_c *SchemaRegistryAwareInput_Stop_Call) Run(run func(ctx context.Context)) *SchemaRegistryAwareInput_Stop_Call
func (*SchemaRegistryAwareInput_Stop_Call) RunAndReturn ¶ added in v0.51.0
func (_c *SchemaRegistryAwareInput_Stop_Call) RunAndReturn(run func(context.Context)) *SchemaRegistryAwareInput_Stop_Call
type SchemaRegistryAwareOutput ¶ added in v0.51.0
SchemaRegistryAwareOutput is an autogenerated mock type for the SchemaRegistryAwareOutput type
func NewSchemaRegistryAwareOutput ¶ added in v0.51.0
func NewSchemaRegistryAwareOutput(t interface {
mock.TestingT
Cleanup(func())
}) *SchemaRegistryAwareOutput
NewSchemaRegistryAwareOutput creates a new instance of SchemaRegistryAwareOutput. It also registers a testing interface on the mock and a cleanup function to assert the mocks expectations. The first argument is typically a *testing.T value.
func (*SchemaRegistryAwareOutput) EXPECT ¶ added in v0.51.0
func (_m *SchemaRegistryAwareOutput) EXPECT() *SchemaRegistryAwareOutput_Expecter
func (*SchemaRegistryAwareOutput) InitSchemaRegistry ¶ added in v0.51.0
func (_m *SchemaRegistryAwareOutput) InitSchemaRegistry(ctx context.Context, settings stream.SchemaSettingsWithEncoding) (stream.MessageBodyEncoder, error)
InitSchemaRegistry provides a mock function with given fields: ctx, settings
func (*SchemaRegistryAwareOutput) Write ¶ added in v0.51.0
func (_m *SchemaRegistryAwareOutput) Write(ctx context.Context, batch []stream.WritableMessage) error
Write provides a mock function with given fields: ctx, batch
func (*SchemaRegistryAwareOutput) WriteOne ¶ added in v0.51.0
func (_m *SchemaRegistryAwareOutput) WriteOne(ctx context.Context, msg stream.WritableMessage) error
WriteOne provides a mock function with given fields: ctx, msg
type SchemaRegistryAwareOutput_Expecter ¶ added in v0.51.0
type SchemaRegistryAwareOutput_Expecter struct {
// contains filtered or unexported fields
}
func (*SchemaRegistryAwareOutput_Expecter) InitSchemaRegistry ¶ added in v0.51.0
func (_e *SchemaRegistryAwareOutput_Expecter) InitSchemaRegistry(ctx interface{}, settings interface{}) *SchemaRegistryAwareOutput_InitSchemaRegistry_Call
InitSchemaRegistry is a helper method to define mock.On call
- ctx context.Context
- settings stream.SchemaSettingsWithEncoding
func (*SchemaRegistryAwareOutput_Expecter) Write ¶ added in v0.51.0
func (_e *SchemaRegistryAwareOutput_Expecter) Write(ctx interface{}, batch interface{}) *SchemaRegistryAwareOutput_Write_Call
Write is a helper method to define mock.On call
- ctx context.Context
- batch []stream.WritableMessage
func (*SchemaRegistryAwareOutput_Expecter) WriteOne ¶ added in v0.51.0
func (_e *SchemaRegistryAwareOutput_Expecter) WriteOne(ctx interface{}, msg interface{}) *SchemaRegistryAwareOutput_WriteOne_Call
WriteOne is a helper method to define mock.On call
- ctx context.Context
- msg stream.WritableMessage
type SchemaRegistryAwareOutput_InitSchemaRegistry_Call ¶ added in v0.51.0
SchemaRegistryAwareOutput_InitSchemaRegistry_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'InitSchemaRegistry'
func (*SchemaRegistryAwareOutput_InitSchemaRegistry_Call) Return ¶ added in v0.51.0
func (_c *SchemaRegistryAwareOutput_InitSchemaRegistry_Call) Return(_a0 stream.MessageBodyEncoder, _a1 error) *SchemaRegistryAwareOutput_InitSchemaRegistry_Call
func (*SchemaRegistryAwareOutput_InitSchemaRegistry_Call) Run ¶ added in v0.51.0
func (_c *SchemaRegistryAwareOutput_InitSchemaRegistry_Call) Run(run func(ctx context.Context, settings stream.SchemaSettingsWithEncoding)) *SchemaRegistryAwareOutput_InitSchemaRegistry_Call
func (*SchemaRegistryAwareOutput_InitSchemaRegistry_Call) RunAndReturn ¶ added in v0.51.0
func (_c *SchemaRegistryAwareOutput_InitSchemaRegistry_Call) RunAndReturn(run func(context.Context, stream.SchemaSettingsWithEncoding) (stream.MessageBodyEncoder, error)) *SchemaRegistryAwareOutput_InitSchemaRegistry_Call
type SchemaRegistryAwareOutput_WriteOne_Call ¶ added in v0.51.0
SchemaRegistryAwareOutput_WriteOne_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'WriteOne'
func (*SchemaRegistryAwareOutput_WriteOne_Call) Return ¶ added in v0.51.0
func (_c *SchemaRegistryAwareOutput_WriteOne_Call) Return(_a0 error) *SchemaRegistryAwareOutput_WriteOne_Call
func (*SchemaRegistryAwareOutput_WriteOne_Call) Run ¶ added in v0.51.0
func (_c *SchemaRegistryAwareOutput_WriteOne_Call) Run(run func(ctx context.Context, msg stream.WritableMessage)) *SchemaRegistryAwareOutput_WriteOne_Call
func (*SchemaRegistryAwareOutput_WriteOne_Call) RunAndReturn ¶ added in v0.51.0
func (_c *SchemaRegistryAwareOutput_WriteOne_Call) RunAndReturn(run func(context.Context, stream.WritableMessage) error) *SchemaRegistryAwareOutput_WriteOne_Call
type SchemaRegistryAwareOutput_Write_Call ¶ added in v0.51.0
SchemaRegistryAwareOutput_Write_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'Write'
func (*SchemaRegistryAwareOutput_Write_Call) Return ¶ added in v0.51.0
func (_c *SchemaRegistryAwareOutput_Write_Call) Return(_a0 error) *SchemaRegistryAwareOutput_Write_Call
func (*SchemaRegistryAwareOutput_Write_Call) Run ¶ added in v0.51.0
func (_c *SchemaRegistryAwareOutput_Write_Call) Run(run func(ctx context.Context, batch []stream.WritableMessage)) *SchemaRegistryAwareOutput_Write_Call
func (*SchemaRegistryAwareOutput_Write_Call) RunAndReturn ¶ added in v0.51.0
func (_c *SchemaRegistryAwareOutput_Write_Call) RunAndReturn(run func(context.Context, []stream.WritableMessage) error) *SchemaRegistryAwareOutput_Write_Call
type SchemaSettingsAwareCallback ¶ added in v0.51.0
SchemaSettingsAwareCallback is an autogenerated mock type for the SchemaSettingsAwareCallback type
func NewSchemaSettingsAwareCallback ¶ added in v0.51.0
func NewSchemaSettingsAwareCallback(t interface {
mock.TestingT
Cleanup(func())
}) *SchemaSettingsAwareCallback
NewSchemaSettingsAwareCallback creates a new instance of SchemaSettingsAwareCallback. It also registers a testing interface on the mock and a cleanup function to assert the mocks expectations. The first argument is typically a *testing.T value.
func (*SchemaSettingsAwareCallback) EXPECT ¶ added in v0.51.0
func (_m *SchemaSettingsAwareCallback) EXPECT() *SchemaSettingsAwareCallback_Expecter
func (*SchemaSettingsAwareCallback) GetSchemaSettings ¶ added in v0.51.0
func (_m *SchemaSettingsAwareCallback) GetSchemaSettings() (*stream.SchemaSettings, error)
GetSchemaSettings provides a mock function with no fields
type SchemaSettingsAwareCallback_Expecter ¶ added in v0.51.0
type SchemaSettingsAwareCallback_Expecter struct {
// contains filtered or unexported fields
}
func (*SchemaSettingsAwareCallback_Expecter) GetSchemaSettings ¶ added in v0.51.0
func (_e *SchemaSettingsAwareCallback_Expecter) GetSchemaSettings() *SchemaSettingsAwareCallback_GetSchemaSettings_Call
GetSchemaSettings is a helper method to define mock.On call
type SchemaSettingsAwareCallback_GetSchemaSettings_Call ¶ added in v0.51.0
SchemaSettingsAwareCallback_GetSchemaSettings_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'GetSchemaSettings'
func (*SchemaSettingsAwareCallback_GetSchemaSettings_Call) Return ¶ added in v0.51.0
func (_c *SchemaSettingsAwareCallback_GetSchemaSettings_Call) Return(_a0 *stream.SchemaSettings, _a1 error) *SchemaSettingsAwareCallback_GetSchemaSettings_Call
func (*SchemaSettingsAwareCallback_GetSchemaSettings_Call) Run ¶ added in v0.51.0
func (_c *SchemaSettingsAwareCallback_GetSchemaSettings_Call) Run(run func()) *SchemaSettingsAwareCallback_GetSchemaSettings_Call
func (*SchemaSettingsAwareCallback_GetSchemaSettings_Call) RunAndReturn ¶ added in v0.51.0
func (_c *SchemaSettingsAwareCallback_GetSchemaSettings_Call) RunAndReturn(run func() (*stream.SchemaSettings, error)) *SchemaSettingsAwareCallback_GetSchemaSettings_Call
type UntypedConsumerCallback ¶ added in v0.41.0
UntypedConsumerCallback is an autogenerated mock type for the UntypedConsumerCallback type
func NewUntypedConsumerCallback ¶ added in v0.41.0
func NewUntypedConsumerCallback(t interface {
mock.TestingT
Cleanup(func())
}) *UntypedConsumerCallback
NewUntypedConsumerCallback creates a new instance of UntypedConsumerCallback. It also registers a testing interface on the mock and a cleanup function to assert the mocks expectations. The first argument is typically a *testing.T value.
func (*UntypedConsumerCallback) Consume ¶ added in v0.41.0
func (_m *UntypedConsumerCallback) Consume(ctx context.Context, model interface{}, attributes map[string]string) (bool, error)
Consume provides a mock function with given fields: ctx, model, attributes
func (*UntypedConsumerCallback) EXPECT ¶ added in v0.41.0
func (_m *UntypedConsumerCallback) EXPECT() *UntypedConsumerCallback_Expecter
type UntypedConsumerCallback_Consume_Call ¶ added in v0.41.0
UntypedConsumerCallback_Consume_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'Consume'
func (*UntypedConsumerCallback_Consume_Call) Return ¶ added in v0.41.0
func (_c *UntypedConsumerCallback_Consume_Call) Return(_a0 bool, _a1 error) *UntypedConsumerCallback_Consume_Call
func (*UntypedConsumerCallback_Consume_Call) Run ¶ added in v0.41.0
func (_c *UntypedConsumerCallback_Consume_Call) Run(run func(ctx context.Context, model interface{}, attributes map[string]string)) *UntypedConsumerCallback_Consume_Call
func (*UntypedConsumerCallback_Consume_Call) RunAndReturn ¶ added in v0.41.0
func (_c *UntypedConsumerCallback_Consume_Call) RunAndReturn(run func(context.Context, interface{}, map[string]string) (bool, error)) *UntypedConsumerCallback_Consume_Call
type UntypedConsumerCallback_Expecter ¶ added in v0.41.0
type UntypedConsumerCallback_Expecter struct {
// contains filtered or unexported fields
}
func (*UntypedConsumerCallback_Expecter) Consume ¶ added in v0.41.0
func (_e *UntypedConsumerCallback_Expecter) Consume(ctx interface{}, model interface{}, attributes interface{}) *UntypedConsumerCallback_Consume_Call
Consume is a helper method to define mock.On call
- ctx context.Context
- model interface{}
- attributes map[string]string
func (*UntypedConsumerCallback_Expecter) GetModel ¶ added in v0.41.0
func (_e *UntypedConsumerCallback_Expecter) GetModel(attributes interface{}) *UntypedConsumerCallback_GetModel_Call
GetModel is a helper method to define mock.On call
- attributes map[string]string
type UntypedConsumerCallback_GetModel_Call ¶ added in v0.41.0
UntypedConsumerCallback_GetModel_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'GetModel'
func (*UntypedConsumerCallback_GetModel_Call) Return ¶ added in v0.41.0
func (_c *UntypedConsumerCallback_GetModel_Call) Return(_a0 interface{}, _a1 error) *UntypedConsumerCallback_GetModel_Call
func (*UntypedConsumerCallback_GetModel_Call) Run ¶ added in v0.41.0
func (_c *UntypedConsumerCallback_GetModel_Call) Run(run func(attributes map[string]string)) *UntypedConsumerCallback_GetModel_Call
func (*UntypedConsumerCallback_GetModel_Call) RunAndReturn ¶ added in v0.41.0
func (_c *UntypedConsumerCallback_GetModel_Call) RunAndReturn(run func(map[string]string) (interface{}, error)) *UntypedConsumerCallback_GetModel_Call
Source Files
¶
- ConsumerCallback.go
- Input.go
- KinsumerAutoscaleOrchestrator.go
- MessageEncoder.go
- Output.go
- PartitionerRand.go
- Producer.go
- ProducerDaemonAggregator.go
- ProducerDaemonBatcher.go
- RetryHandler.go
- RunnableCallback.go
- RunnableConsumerCallback.go
- RunnableUntypedConsumerCallback.go
- SchemaRegistryAwareInput.go
- SchemaRegistryAwareOutput.go
- SchemaSettingsAwareCallback.go
- UntypedConsumerCallback.go