| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103 |
- // package dict_topic_serve -- словарь топиков обработчиков запросов.
- package dict_topic_serve
- import (
- "context"
- "sync"
- "time"
- "gitp78su.ipnodns.ru/svi/kern/v4/lev0"
- "gitp78su.ipnodns.ru/svi/kern/v4/lev2/kern_ctx"
- "gitp78su.ipnodns.ru/svi/kern/v4/lev2/lti/bus_ent"
- "gitp78su.ipnodns.ru/svi/kern/v4/lev2/lti/bus_ent/topic"
- "gitp78su.ipnodns.ru/svi/kern/v4/lev2/lti/bus_spec"
- )
- // dictServe -- потокобезопасный словарь обработчиков запросов.
- //
- // Допускается только один обработчик запросов на один топик.
- type dictServe struct {
- sync.RWMutex
- kCtx *kern_ctx.KernCtx
- dictServe map[bus_ent.ATopic]bus_spec.IBusHandlerServe
- }
- // NewDictServe -- возвращает потокобезопасный словарь обработчиков запросов.
- func NewDictServe() *dictServe {
- sf := &dictServe{
- kCtx: kern_ctx.GetKernCtx(),
- dictServe: map[bus_ent.ATopic]bus_spec.IBusHandlerServe{},
- }
- return sf
- }
- // Register -- регистрирует обработчик запросов.
- func (sf *dictServe) Register(handler bus_spec.IBusHandlerServe) *lev0.Result[lev0.Bool] {
- sf.Lock()
- defer sf.Unlock()
- if handler == nil {
- err := lev0.LetErr("dictServe.Register(): IBusHandlerSubscribe==nil")
- return lev0.ResErr[lev0.Bool](err)
- }
- topic := handler.Topic()
- isTwinRegister := sf.register(handler)
- if isTwinRegister {
- err := lev0.LetErr("dictServe.Register(): handler of topic (%v) already register", topic.Get())
- return lev0.ResErr[lev0.Bool](err)
- }
- return lev0.ResOk(lev0.Bool(true))
- }
- // Unregister -- удаляет обработчик запросов из словаря.
- func (sf *dictServe) Unregister(handler bus_spec.IBusHandlerServe) {
- sf.Lock()
- defer sf.Unlock()
- if handler == nil {
- err := lev0.LetErr("dictServe.Unregister(): IBusHandlerSubscribe==nil")
- sf.kCtx.Log().Err(err)
- return
- }
- delete(sf.dictServe, handler.Topic().Get())
- }
- // SendRequest -- вызывает обработчик при поступлении запроса.
- func (sf *dictServe) SendRequest(topic *topic.LTopic, binReq []byte) *lev0.Result[[]byte] {
- sf.RLock()
- defer sf.RUnlock()
- handler, isOk := sf.dictServe[topic.Get()]
- if !isOk {
- err := lev0.LetErr("dictServe.SendRequest(): handler for topic (%v) not exists", topic.Get())
- return lev0.ResErr[[]byte](err)
- }
- var (
- chRes = make(chan *lev0.Result[[]byte], 2)
- )
- ctx, fnCancel := context.WithTimeout(sf.kCtx.SelfCtx(), time.Millisecond*time.Duration(TimeoutDefault))
- defer fnCancel()
- fnCall := func() {
- defer close(chRes)
- res := handler.FnBack(binReq)
- chRes <- res
- }
- go fnCall()
- select {
- case <-ctx.Done():
- err := lev0.LetErr("dictServe.SendRequest(): in call for topic (%v), err=\n\t%w", topic.Get(), ctx.Err())
- return lev0.ResErr[[]byte](err)
- case res := <-chRes:
- return res
- }
- }
- var TimeoutDefault = 15000
- // регистрирует обработчик запросов.
- func (sf *dictServe) register(handler bus_spec.IBusHandlerServe) bool {
- topic := handler.Topic()
- _, isOk := sf.dictServe[topic.Get()]
- if isOk {
- return true
- }
- sf.dictServe[topic.Get()] = handler
- return false
- }
|