70 lines
1.2 KiB
Go
70 lines
1.2 KiB
Go
|
package consumer
|
||
|
|
||
|
import (
|
||
|
"fmt"
|
||
|
"strings"
|
||
|
"sync/atomic"
|
||
|
)
|
||
|
|
||
|
type ConsumerHandler interface {
|
||
|
GetName() string
|
||
|
GetConsumerName() string
|
||
|
|
||
|
ConsumerCtx() ConsumerCtx
|
||
|
|
||
|
Init(consumerCtx ConsumerCtx) error
|
||
|
OnStart(consumerCtx ConsumerCtx) error
|
||
|
OnStop(consumerCtx ConsumerCtx)
|
||
|
Destroy(consumerCtx ConsumerCtx)
|
||
|
|
||
|
Validate() error
|
||
|
}
|
||
|
|
||
|
type ConsumerHandlers struct {
|
||
|
Name string `json:"name,omitempty"`
|
||
|
ConsumerName string `json:"consumerName,omitempty"`
|
||
|
|
||
|
validated atomic.Value
|
||
|
}
|
||
|
|
||
|
func (ch *ConsumerHandlers) ConsumerCtx() ConsumerCtx {
|
||
|
return NewConsumerCtx(nil)
|
||
|
}
|
||
|
|
||
|
func (ch *ConsumerHandlers) Init(consumerCtx ConsumerCtx) error {
|
||
|
return nil
|
||
|
}
|
||
|
|
||
|
func (ch *ConsumerHandlers) OnStart(consumerCtx ConsumerCtx) error {
|
||
|
return nil
|
||
|
}
|
||
|
|
||
|
func (ch *ConsumerHandlers) OnStop(consumerCtx ConsumerCtx) {
|
||
|
|
||
|
}
|
||
|
|
||
|
func (ch *ConsumerHandlers) Destroy(consumerCtx ConsumerCtx) {
|
||
|
|
||
|
}
|
||
|
|
||
|
func (ch *ConsumerHandlers) GetName() string {
|
||
|
return ch.Name
|
||
|
}
|
||
|
|
||
|
func (ch *ConsumerHandlers) GetConsumerName() string {
|
||
|
return ch.ConsumerName
|
||
|
}
|
||
|
|
||
|
func (ch *ConsumerHandlers) Validate() error {
|
||
|
if nil != ch.validated.Load() {
|
||
|
return nil
|
||
|
}
|
||
|
ch.validated.Store(true)
|
||
|
|
||
|
if "" == strings.TrimSpace(ch.ConsumerName) {
|
||
|
return fmt.Errorf("ConsumerName is not valid")
|
||
|
}
|
||
|
|
||
|
return nil
|
||
|
}
|