kbus_http.go 7.1 KB

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