kbus_http.go 6.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223
  1. // package kbus_http -- шина сообщений поверх HTTP.
  2. package kbus_http
  3. import (
  4. "fmt"
  5. "net/http"
  6. "sync"
  7. "github.com/gofiber/fiber/v3"
  8. mKs "gitp78su.ipnodns.ru/svi/kern/v4/lev0/core_spec"
  9. mL1 "gitp78su.ipnodns.ru/svi/kern/v4/lev1"
  10. "gitp78su.ipnodns.ru/svi/kern/v4/lev1/safe_bool"
  11. "gitp78su.ipnodns.ru/svi/kern/v4/lev2/kctx"
  12. "gitp78su.ipnodns.ru/svi/kern/v4/lev2/kserv_http"
  13. "gitp78su.ipnodns.ru/svi/kern/v4/lev2/lti/bus_mock/mock_hand_sub_http"
  14. "gitp78su.ipnodns.ru/svi/kern/v4/lev2/lti/dot"
  15. "gitp78su.ipnodns.ru/svi/kern/v4/lev2/lti/kbus_base"
  16. )
  17. // kBusHttp -- шина данных поверх HTTP.
  18. type kBusHttp struct {
  19. *kbus_base.KBusBase
  20. log mKs.ILogBuf
  21. }
  22. var (
  23. Bus_ *kBusHttp
  24. block sync.Mutex
  25. )
  26. // GetKernelBusHttp -- возвращает шину HTTP.
  27. func GetKernelBusHttp() mKs.IKernelBus {
  28. block.Lock()
  29. defer block.Unlock()
  30. if Bus_ != nil {
  31. return Bus_
  32. }
  33. paramLogBuf := &mL1.LogBufParam{
  34. IsTerm_: safe_bool.NewSafeBool(true),
  35. Prefix_: "kBusHttp",
  36. }
  37. log := mL1.NewLogBuf(paramLogBuf)
  38. log.Debug("GetKernelBusHttp(): new")
  39. kCtx := kctx.GetKernelCtx()
  40. kBus := kbus_base.GetKernelBusBase()
  41. sf := &kBusHttp{
  42. KBusBase: kBus,
  43. log: log,
  44. }
  45. kServHttp := kserv_http.GetKernelServHttp()
  46. fibApp := kServHttp.Fiber()
  47. fibApp.Post("/bus/sub", sf.postSub) // Топик подписки, IN
  48. fibApp.Post("/bus/unsub", sf.postUnsub) // Топик отписки, IN
  49. fibApp.Post("/bus/request", sf.postSendRequest) // Топик входящих запросов, IN
  50. fibApp.Post("/bus/pub", sf.postPublish) // Топик публикаций подписки, IN
  51. kCtx.Set("kernBus", sf, "GetKernelBusHttp(): http data bus")
  52. Bus_ = sf
  53. return Bus_
  54. }
  55. // Входящий запрос HTTP на подписку.
  56. func (sf *kBusHttp) postSub(ctx fiber.Ctx) error {
  57. sf.log.Debug("postSub()")
  58. ctx.Set("Content-type", "text/html; charset=utf8")
  59. ctx.Set("Content-type", "text/json")
  60. ctx.Set("Cache-Control", "no-cache")
  61. sf.log.Debug("postSub()")
  62. req := &dict_topic_serve.SubscribeReq{}
  63. err := ctx.Bind().Body(req)
  64. if err != nil {
  65. resp := &dot.SubscribeResp{
  66. Status_: fmt.Sprintf("kernelBusHttp.postSub(): in parse request, err=\n\t%v\n", err),
  67. Uuid_: req.Uuid_,
  68. }
  69. resp.SelfCheck()
  70. ctx.Response().SetStatusCode(http.StatusBadRequest)
  71. sf.log.Err("postSub(): in body parser, status=%q", resp.Status_)
  72. return ctx.JSON(resp)
  73. }
  74. resp := sf.processSubscribe(req)
  75. resp.SelfCheck()
  76. return ctx.JSON(resp)
  77. }
  78. // Процесс подписки веб-хука.
  79. func (sf *kBusHttp) processSubscribe(req *dot.SubscribeReq) *dot.SubscribeResp {
  80. req.SelfCheck()
  81. handler := mock_hand_sub_http.NewMockHandSubHttp(req.Topic_, req.WebHook_)
  82. resp := &dot.SubscribeResp{
  83. Status_: "ok",
  84. Uuid_: req.Uuid_,
  85. Name_: handler.Name(),
  86. }
  87. res := sf.Subscribe(handler)
  88. if res.IsErr() {
  89. resp.Status_ = fmt.Sprintf("kernelBusHttp.processSubscribe(): err=\n\t%v", res.Err())
  90. return resp
  91. }
  92. return resp
  93. }
  94. // Входящая публикация.
  95. func (sf *kBusHttp) postPublish(ctx fiber.Ctx) error {
  96. sf.log.Debug("postPublish()")
  97. ctx.Set("Content-type", "text/html; charset=utf8")
  98. ctx.Set("Content-type", "text/json")
  99. ctx.Set("Cache-Control", "no-cache")
  100. req := &dict_topic_serve.PublishReq{}
  101. err := ctx.Bind().Body(req)
  102. if err != nil {
  103. resp := &dot.PublishResp{
  104. Status_: fmt.Sprintf("kernelBusHttp.postPublish(): in parse request, err=\n\t%v\n", err),
  105. Uuid_: req.Uuid_,
  106. }
  107. resp.SelfCheck()
  108. ctx.Response().SetStatusCode(http.StatusBadRequest)
  109. sf.log.Err("postPublish(): in body parser, status=%v", resp.Status_)
  110. return ctx.JSON(resp)
  111. }
  112. resp := sf.processPublish(req)
  113. resp.SelfCheck()
  114. return ctx.JSON(resp)
  115. }
  116. // Выполняет процесс публикации.
  117. func (sf *kBusHttp) processPublish(req *dot.PublishReq) *dot.PublishResp {
  118. req.SelfCheck()
  119. res := sf.Publish(req.Topic_, req.BinMsg_)
  120. resp := &dot.PublishResp{
  121. Status_: "ok",
  122. Uuid_: req.Uuid_,
  123. }
  124. if res.IsErr() {
  125. resp.Status_ = fmt.Sprintf("kernelBusHttp.processPublish(): err=\n\t%v", res.Err())
  126. return resp
  127. }
  128. return resp
  129. }
  130. // Входящий запрос.
  131. func (sf *kBusHttp) postSendRequest(ctx fiber.Ctx) error {
  132. sf.log.Debug("postSendRequest()")
  133. ctx.Set("Content-type", "text/html; charset=utf8")
  134. ctx.Set("Content-type", "text/json")
  135. ctx.Set("Cache-Control", "no-cache")
  136. req := &dot.ServeReq{}
  137. err := ctx.Bind().Body(req)
  138. if err != nil {
  139. resp := &dot.ServeResp{
  140. Status_: fmt.Sprintf("kernelBusHttp.postSendRequest(): err=\n\t%v", err),
  141. Uuid_: req.Uuid_,
  142. }
  143. resp.SelfCheck()
  144. ctx.Response().SetStatusCode(http.StatusBadRequest)
  145. sf.log.Err("postSendRequest(): in body parser, status=%v", resp.Status_)
  146. return ctx.JSON(resp)
  147. }
  148. resp := sf.processSendRequest(req)
  149. resp.SelfCheck()
  150. return ctx.JSON(resp)
  151. }
  152. // Обрабатывает входящий запрос.
  153. func (sf *kBusHttp) processSendRequest(req *dot.ServeReq) *dot.ServeResp {
  154. req.SelfCheck()
  155. res := sf.SendRequest(req.Topic_, req.BinReq_)
  156. resp := &dot.ServeResp{
  157. Status_: "ok",
  158. Uuid_: req.Uuid_,
  159. }
  160. if res.IsErr() {
  161. resp.Status_ = fmt.Sprintf("kernelBusHttp.processSendRequest(): err=\n\t%v", res.Err())
  162. return resp
  163. }
  164. resp.BinResp_ = res.Ok()
  165. return resp
  166. }
  167. // Входящая отписка от топика по HTTP.
  168. func (sf *kBusHttp) postUnsub(ctx fiber.Ctx) error {
  169. sf.log.Debug("postUnsub()")
  170. ctx.Set("Content-type", "text/html; charset=utf8")
  171. ctx.Set("Content-type", "text/json")
  172. ctx.Set("Cache-Control", "no-cache")
  173. req := &dot.UnsubReq{}
  174. err := ctx.Bind().Body(req)
  175. if err != nil {
  176. resp := &dot.ServeResp{
  177. Status_: fmt.Sprintf("kernelBusHttp.postSendRequest(): err=\n\t%v", err),
  178. Uuid_: req.Uuid_,
  179. }
  180. resp.SelfCheck()
  181. ctx.Response().SetStatusCode(http.StatusBadRequest)
  182. sf.log.Err("postUnsub(): in body_parser, status=%q", resp.Status_)
  183. return ctx.JSON(resp)
  184. }
  185. resp := sf.processUnsubRequest(req)
  186. resp.SelfCheck()
  187. return ctx.JSON(resp)
  188. }
  189. // Процесс отписки от топика.
  190. func (sf *kBusHttp) processUnsubRequest(req *dot.UnsubReq) *dot.UnsubResp {
  191. req.SelfCheck()
  192. resp := &dot.UnsubResp{
  193. Status_: "ok",
  194. Uuid_: req.Uuid_,
  195. }
  196. optHandler := sf.KCtx_.Get(req.Name_.Get())
  197. if optHandler.IsNone() {
  198. resp.Status_ = fmt.Sprintf("kernelBusHttp.processUnsubRequest(): not get handler(%v) from kernel ctx",
  199. req.Name_)
  200. return resp
  201. }
  202. if optHandler == nil {
  203. resp.Status_ = fmt.Sprintf("kernelBusHttp.processUnsubRequest(): handler(%v) not exists", req.Name_)
  204. return resp
  205. }
  206. hand := optHandler.Some().Val().(mKs.IBusHandlerSubscribe)
  207. sf.Unsubscribe(hand)
  208. return resp
  209. }