| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235 |
- // package kbus_http -- шина сообщений поверх HTTP.
- package kbus_http
- import (
- "net/http"
- "sync"
- "github.com/gofiber/fiber/v3"
- "gitp78su.ipnodns.ru/svi/kern/v4/lev0"
- "gitp78su.ipnodns.ru/svi/kern/v4/lev1"
- "gitp78su.ipnodns.ru/svi/kern/v4/lev2/kern_ctx"
- "gitp78su.ipnodns.ru/svi/kern/v4/lev2/kserv_http"
- "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 *lev1.LogBuf
- }
- var (
- Bus_ *kBusHttp
- block sync.Mutex
- )
- // GetKernelBusHttp -- возвращает шину HTTP.
- func GetKernelBusHttp() bus_spec.IKernelBus {
- block.Lock()
- defer block.Unlock()
- if Bus_ != nil {
- return Bus_
- }
- paramLogBuf := &lev1.LogBufParam{
- IsTerm_: lev0.MutSafeBool(true),
- Prefix_: lev0.LetTxt("kBusHttp"),
- }
- log := lev1.NewLogBuf(paramLogBuf)
- log.Debug("GetKernelBusHttp(): new")
- kCtx := kern_ctx.GetKernCtx()
- 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(lev0.LetTxt("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_: lev0.LetTxt("kernelBusHttp.postSub(): in parse request, err=\n\t%v\n", err),
- Uuid_: req.Uuid_,
- }
- resp.SelfCheck()
- ctx.Response().SetStatusCode(http.StatusBadRequest)
- err := lev0.LetErr("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 := topic.LetTopic(req.Topic_)
- handler := mock_hand_sub_http.NewMockHandSubHttp(topic, req.WebHook_)
- resp := &dot.SubscribeResp{
- Status_: lev0.LetTxt("ok"),
- Uuid_: req.Uuid_,
- Name_: handler.Name().Get(),
- }
- res := sf.Subscribe(handler)
- if res.IsErr() {
- resp.Status_ = lev0.LetTxt(
- "kernelBusHttp.processSubscribe(): err=\n\t%v",
- res.Err())
- return resp
- }
- return resp
- }
- // Входящая публикация.
- func (sf *kBusHttp) postPublish(ctx fiber.Ctx) error {
- sf.log.Debug("postPublish()")
- sf.postCommon(ctx)
- req := &dot.PublishReq{}
- err := ctx.Bind().Body(req)
- if err != nil {
- resp := &dot.PublishResp{
- Status_: lev0.LetTxt("kernelBusHttp.postPublish(): in parse request, err=\n\t%v\n", err),
- Uuid_: req.Uuid_,
- }
- resp.SelfCheck()
- ctx.Response().SetStatusCode(http.StatusBadRequest)
- err := lev0.LetErr("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)
- }
- // Входящая отписка от топика по HTTP.
- func (sf *kBusHttp) postUnsub(ctx fiber.Ctx) error {
- sf.log.Debug("postUnsub()")
- sf.postCommon(ctx)
- req := &dot.UnsubReq{}
- err := ctx.Bind().Body(req)
- if err != nil {
- resp := &dot.ServeResp{
- Status_: lev0.LetTxt("kernelBusHttp.postSendRequest(): err=\n\t%v", err),
- Uuid_: req.Uuid_,
- }
- resp.SelfCheck()
- ctx.Response().SetStatusCode(http.StatusBadRequest)
- err := lev0.LetErr("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) postCommon(ctx fiber.Ctx) {
- ctx.Set("Content-type", "text/html; charset=utf8")
- ctx.Set("Content-type", "text/json")
- ctx.Set("Cache-Control", "no-cache")
- }
- // Выполняет процесс публикации.
- func (sf *kBusHttp) processPublish(req *dot.PublishReq) *dot.PublishResp {
- req.SelfCheck()
- topic := topic.LetTopic(req.Topic_)
- res := sf.Publish(topic, req.BinMsg_)
- resp := &dot.PublishResp{
- Status_: lev0.LetTxt("ok"),
- Uuid_: req.Uuid_,
- }
- if res.IsErr() {
- resp.Status_ = lev0.LetTxt("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_: lev0.LetTxt("kernelBusHttp.postSendRequest(): err=\n\t%v", err),
- Uuid_: req.Uuid_,
- }
- resp.SelfCheck()
- ctx.Response().SetStatusCode(http.StatusBadRequest)
- err := lev0.LetErr("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 := topic.LetTopic(req.Topic_)
- res := sf.SendRequest(topic, req.BinReq_)
- resp := &dot.ServeResp{
- Status_: lev0.LetTxt("ok"),
- Uuid_: req.Uuid_,
- }
- if res.IsErr() {
- resp.Status_ = lev0.LetTxt("kernelBusHttp.processSendRequest(): err=\n\t%v", res.Err())
- return resp
- }
- resp.BinResp_ = res.Ok()
- return resp
- }
- // Процесс отписки от топика.
- func (sf *kBusHttp) processUnsubRequest(req *dot.UnsubReq) *dot.UnsubResp {
- req.SelfCheck()
- resp := &dot.UnsubResp{
- Status_: lev0.LetTxt("ok"),
- Uuid_: req.Uuid_,
- }
- optHandler := sf.KCtx_.Get(lev0.LetTxt(req.HandlerName_))
- if optHandler.IsNone() {
- resp.Status_ = lev0.LetTxt("kernelBusHttp.processUnsubRequest(): not get handler(%v) from kernel ctx",
- req.HandlerName_)
- return resp
- }
- if optHandler == nil {
- resp.Status_ = lev0.LetTxt("kernelBusHttp.processUnsubRequest(): handler(%v) not exists", req.HandlerName_)
- return resp
- }
- hand := optHandler.Some().Val().(bus_spec.IBusHandlerSubscribe)
- sf.Unsubscribe(hand)
- return resp
- }
|