// package kbus_http -- шина сообщений поверх HTTP. package kbus_http import ( "fmt" "net/http" "sync" "github.com/gofiber/fiber/v3" "gitp78su.ipnodns.ru/svi/kern/v4/lev0/defs" mL1 "gitp78su.ipnodns.ru/svi/kern/v4/lev1" "gitp78su.ipnodns.ru/svi/kern/v4/lev1/comp_spec" "gitp78su.ipnodns.ru/svi/kern/v4/lev1/safe_bool" "gitp78su.ipnodns.ru/svi/kern/v4/lev2/kctx" "gitp78su.ipnodns.ru/svi/kern/v4/lev2/kserv_http" "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_mock/mock_hand_sub_http" "gitp78su.ipnodns.ru/svi/kern/v4/lev2/lti/bus_spec" "gitp78su.ipnodns.ru/svi/kern/v4/lev2/lti/dot" "gitp78su.ipnodns.ru/svi/kern/v4/lev2/lti/kbus_base" ) // kBusHttp -- шина данных поверх HTTP. type kBusHttp struct { *kbus_base.KBusBase log comp_spec.ILogBuf } var ( Bus_ *kBusHttp block sync.Mutex ) // GetKernelBusHttp -- возвращает шину HTTP. func GetKernelBusHttp() bus_spec.IKernelBus { block.Lock() defer block.Unlock() if Bus_ != nil { return Bus_ } paramLogBuf := &mL1.LogBufParam{ IsTerm_: safe_bool.NewSafeBool(true), Prefix_: "kBusHttp", } log := mL1.NewLogBuf(paramLogBuf) log.Debug("GetKernelBusHttp(): new") kCtx := kctx.GetKernelCtx() kBus := kbus_base.GetKernelBusBase() sf := &kBusHttp{ KBusBase: kBus, log: log, } kServHttp := kserv_http.GetKernelServHttp() fibApp := kServHttp.Fiber() fibApp.Post("/bus/sub", sf.postSub) // Топик подписки, IN fibApp.Post("/bus/unsub", sf.postUnsub) // Топик отписки, IN fibApp.Post("/bus/request", sf.postSendRequest) // Топик входящих запросов, IN fibApp.Post("/bus/pub", sf.postPublish) // Топик публикаций подписки, IN kCtx.Set("kernBus", sf, "GetKernelBusHttp(): http data bus") Bus_ = sf return Bus_ } // Входящий запрос HTTP на подписку. func (sf *kBusHttp) postSub(ctx fiber.Ctx) error { sf.log.Debug("postSub()") ctx.Set("Content-type", "text/html; charset=utf8") ctx.Set("Content-type", "text/json") ctx.Set("Cache-Control", "no-cache") sf.log.Debug("postSub()") req := &dot.SubscribeReq{} err := ctx.Bind().Body(req) if err != nil { resp := &dot.SubscribeResp{ Status_: fmt.Sprintf("kernelBusHttp.postSub(): in parse request, err=\n\t%v\n", err), Uuid_: req.Uuid_, } resp.SelfCheck() ctx.Response().SetStatusCode(http.StatusBadRequest) err := defs.NewErr("kBusHttp.postSub(): in body parser, status=%q", resp.Status_) sf.log.Err(err) return ctx.JSON(resp) } resp := sf.processSubscribe(req) resp.SelfCheck() return ctx.JSON(resp) } // Процесс подписки веб-хука. func (sf *kBusHttp) processSubscribe(req *dot.SubscribeReq) *dot.SubscribeResp { req.SelfCheck() topic := bus_ent.ITopicFix(topic.NewTopic(req.Topic_)) handler := mock_hand_sub_http.NewMockHandSubHttp(topic, req.WebHook_) resp := &dot.SubscribeResp{ Status_: "ok", Uuid_: req.Uuid_, Name_: handler.Name().Get(), } res := sf.Subscribe(handler) if res.IsErr() { resp.Status_ = fmt.Sprintf("kernelBusHttp.processSubscribe(): err=\n\t%v", res.Err()) return resp } return resp } // Входящая публикация. func (sf *kBusHttp) postPublish(ctx fiber.Ctx) error { sf.log.Debug("postPublish()") ctx.Set("Content-type", "text/html; charset=utf8") ctx.Set("Content-type", "text/json") ctx.Set("Cache-Control", "no-cache") req := &dot.PublishReq{} err := ctx.Bind().Body(req) if err != nil { resp := &dot.PublishResp{ Status_: fmt.Sprintf("kernelBusHttp.postPublish(): in parse request, err=\n\t%v\n", err), Uuid_: req.Uuid_, } resp.SelfCheck() ctx.Response().SetStatusCode(http.StatusBadRequest) err := defs.NewErr("postPublish(): in body parser, status=%v", resp.Status_) sf.log.Err(err) return ctx.JSON(resp) } resp := sf.processPublish(req) resp.SelfCheck() return ctx.JSON(resp) } // Выполняет процесс публикации. func (sf *kBusHttp) processPublish(req *dot.PublishReq) *dot.PublishResp { req.SelfCheck() topic := bus_ent.ITopicFix(topic.NewTopic(req.Topic_)) res := sf.Publish(topic, req.BinMsg_) resp := &dot.PublishResp{ Status_: "ok", Uuid_: req.Uuid_, } if res.IsErr() { resp.Status_ = fmt.Sprintf("kernelBusHttp.processPublish(): err=\n\t%v", res.Err()) return resp } return resp } // Входящий запрос. func (sf *kBusHttp) postSendRequest(ctx fiber.Ctx) error { sf.log.Debug("postSendRequest()") ctx.Set("Content-type", "text/html; charset=utf8") ctx.Set("Content-type", "text/json") ctx.Set("Cache-Control", "no-cache") req := &dot.ServeReq{} err := ctx.Bind().Body(req) if err != nil { resp := &dot.ServeResp{ Status_: defs.Str("kernelBusHttp.postSendRequest(): err=\n\t%v", err), Uuid_: req.Uuid_, } resp.SelfCheck() ctx.Response().SetStatusCode(http.StatusBadRequest) err := defs.NewErr("kBusHttp.postSendRequest(): in body parser, status=%v", resp.Status_) sf.log.Err(err) return ctx.JSON(resp) } resp := sf.processSendRequest(req) resp.SelfCheck() return ctx.JSON(resp) } // Обрабатывает входящий запрос. func (sf *kBusHttp) processSendRequest(req *dot.ServeReq) *dot.ServeResp { req.SelfCheck() topic := bus_ent.ITopicFix(topic.NewTopic(req.Topic_)) res := sf.SendRequest(topic, req.BinReq_) resp := &dot.ServeResp{ Status_: "ok", Uuid_: req.Uuid_, } if res.IsErr() { resp.Status_ = defs.Str("kernelBusHttp.processSendRequest(): err=\n\t%v", res.Err()) return resp } resp.BinResp_ = res.Ok() return resp } // Входящая отписка от топика по HTTP. func (sf *kBusHttp) postUnsub(ctx fiber.Ctx) error { sf.log.Debug("postUnsub()") ctx.Set("Content-type", "text/html; charset=utf8") ctx.Set("Content-type", "text/json") ctx.Set("Cache-Control", "no-cache") req := &dot.UnsubReq{} err := ctx.Bind().Body(req) if err != nil { resp := &dot.ServeResp{ Status_: defs.Str("kernelBusHttp.postSendRequest(): err=\n\t%v", err), Uuid_: req.Uuid_, } resp.SelfCheck() ctx.Response().SetStatusCode(http.StatusBadRequest) err := defs.NewErr("kBusHttp.postUnsub(): in body parser, status=%q", resp.Status_) sf.log.Err(err) return ctx.JSON(resp) } resp := sf.processUnsubRequest(req) resp.SelfCheck() return ctx.JSON(resp) } // Процесс отписки от топика. func (sf *kBusHttp) processUnsubRequest(req *dot.UnsubReq) *dot.UnsubResp { req.SelfCheck() resp := &dot.UnsubResp{ Status_: "ok", Uuid_: req.Uuid_, } optHandler := sf.KCtx_.Get(defs.AStr(req.HandlerName_)) if optHandler.IsNone() { resp.Status_ = defs.Str("kernelBusHttp.processUnsubRequest(): not get handler(%v) from kernel ctx", req.HandlerName_) return resp } if optHandler == nil { resp.Status_ = defs.Str("kernelBusHttp.processUnsubRequest(): handler(%v) not exists", req.HandlerName_) return resp } hand := optHandler.Some().Val().(bus_spec.IBusHandlerSubscribe) sf.Unsubscribe(hand) return resp }