kbus_http.go 6.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235
  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"
  8. "gitp78su.ipnodns.ru/svi/kern/v4/lev1"
  9. "gitp78su.ipnodns.ru/svi/kern/v4/lev2/kern_ctx"
  10. "gitp78su.ipnodns.ru/svi/kern/v4/lev2/kserv_http"
  11. "gitp78su.ipnodns.ru/svi/kern/v4/lev2/lti/bus_ent/topic"
  12. "gitp78su.ipnodns.ru/svi/kern/v4/lev2/lti/bus_mock/mock_hand_sub_http"
  13. "gitp78su.ipnodns.ru/svi/kern/v4/lev2/lti/bus_spec"
  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 *lev1.LogBuf
  21. }
  22. var (
  23. Bus_ *kBusHttp
  24. block sync.Mutex
  25. )
  26. // GetKernelBusHttp -- возвращает шину HTTP.
  27. func GetKernelBusHttp() bus_spec.IKernelBus {
  28. block.Lock()
  29. defer block.Unlock()
  30. if Bus_ != nil {
  31. return Bus_
  32. }
  33. paramLogBuf := &lev1.LogBufParam{
  34. IsTerm_: lev0.MutSafeBool(true),
  35. Prefix_: lev0.LetTxt("kBusHttp"),
  36. }
  37. log := lev1.NewLogBuf(paramLogBuf)
  38. log.Debug("GetKernelBusHttp(): new")
  39. kCtx := kern_ctx.GetKernCtx()
  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(lev0.LetTxt("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 := &dot.SubscribeReq{}
  63. err := ctx.Bind().Body(req)
  64. if err != nil {
  65. resp := &dot.SubscribeResp{
  66. Status_: lev0.LetTxt("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. err := lev0.LetErr("kBusHttp.postSub(): in body parser, status=%q", resp.Status_)
  72. sf.log.Err(err)
  73. return ctx.JSON(resp)
  74. }
  75. resp := sf.processSubscribe(req)
  76. resp.SelfCheck()
  77. return ctx.JSON(resp)
  78. }
  79. // Процесс подписки веб-хука.
  80. func (sf *kBusHttp) processSubscribe(req *dot.SubscribeReq) *dot.SubscribeResp {
  81. req.SelfCheck()
  82. topic := topic.LetTopic(req.Topic_)
  83. handler := mock_hand_sub_http.NewMockHandSubHttp(topic, req.WebHook_)
  84. resp := &dot.SubscribeResp{
  85. Status_: lev0.LetTxt("ok"),
  86. Uuid_: req.Uuid_,
  87. Name_: handler.Name().Get(),
  88. }
  89. res := sf.Subscribe(handler)
  90. if res.IsErr() {
  91. resp.Status_ = lev0.LetTxt(
  92. "kernelBusHttp.processSubscribe(): err=\n\t%v",
  93. res.Err())
  94. return resp
  95. }
  96. return resp
  97. }
  98. // Входящая публикация.
  99. func (sf *kBusHttp) postPublish(ctx fiber.Ctx) error {
  100. sf.log.Debug("postPublish()")
  101. sf.postCommon(ctx)
  102. req := &dot.PublishReq{}
  103. err := ctx.Bind().Body(req)
  104. if err != nil {
  105. resp := &dot.PublishResp{
  106. Status_: lev0.LetTxt("kernelBusHttp.postPublish(): in parse request, err=\n\t%v\n", err),
  107. Uuid_: req.Uuid_,
  108. }
  109. resp.SelfCheck()
  110. ctx.Response().SetStatusCode(http.StatusBadRequest)
  111. err := lev0.LetErr("postPublish(): in body parser, status=%v", resp.Status_)
  112. sf.log.Err(err)
  113. return ctx.JSON(resp)
  114. }
  115. resp := sf.processPublish(req)
  116. resp.SelfCheck()
  117. return ctx.JSON(resp)
  118. }
  119. // Входящая отписка от топика по HTTP.
  120. func (sf *kBusHttp) postUnsub(ctx fiber.Ctx) error {
  121. sf.log.Debug("postUnsub()")
  122. sf.postCommon(ctx)
  123. req := &dot.UnsubReq{}
  124. err := ctx.Bind().Body(req)
  125. if err != nil {
  126. resp := &dot.ServeResp{
  127. Status_: lev0.LetTxt("kernelBusHttp.postSendRequest(): err=\n\t%v", err),
  128. Uuid_: req.Uuid_,
  129. }
  130. resp.SelfCheck()
  131. ctx.Response().SetStatusCode(http.StatusBadRequest)
  132. err := lev0.LetErr("kBusHttp.postUnsub(): in body parser, status=%q", resp.Status_)
  133. sf.log.Err(err)
  134. return ctx.JSON(resp)
  135. }
  136. resp := sf.processUnsubRequest(req)
  137. resp.SelfCheck()
  138. return ctx.JSON(resp)
  139. }
  140. // Общая часть процесса публикации и отписки.
  141. func (sf *kBusHttp) postCommon(ctx fiber.Ctx) {
  142. ctx.Set("Content-type", "text/html; charset=utf8")
  143. ctx.Set("Content-type", "text/json")
  144. ctx.Set("Cache-Control", "no-cache")
  145. }
  146. // Выполняет процесс публикации.
  147. func (sf *kBusHttp) processPublish(req *dot.PublishReq) *dot.PublishResp {
  148. req.SelfCheck()
  149. topic := topic.LetTopic(req.Topic_)
  150. res := sf.Publish(topic, req.BinMsg_)
  151. resp := &dot.PublishResp{
  152. Status_: lev0.LetTxt("ok"),
  153. Uuid_: req.Uuid_,
  154. }
  155. if res.IsErr() {
  156. resp.Status_ = lev0.LetTxt("kernelBusHttp.processPublish(): err=\n\t%v", res.Err())
  157. return resp
  158. }
  159. return resp
  160. }
  161. // Входящий запрос.
  162. func (sf *kBusHttp) postSendRequest(ctx fiber.Ctx) error {
  163. sf.log.Debug("postSendRequest()")
  164. ctx.Set("Content-type", "text/html; charset=utf8")
  165. ctx.Set("Content-type", "text/json")
  166. ctx.Set("Cache-Control", "no-cache")
  167. req := &dot.ServeReq{}
  168. err := ctx.Bind().Body(req)
  169. if err != nil {
  170. resp := &dot.ServeResp{
  171. Status_: lev0.LetTxt("kernelBusHttp.postSendRequest(): err=\n\t%v", err),
  172. Uuid_: req.Uuid_,
  173. }
  174. resp.SelfCheck()
  175. ctx.Response().SetStatusCode(http.StatusBadRequest)
  176. err := lev0.LetErr("kBusHttp.postSendRequest(): in body parser, status=%v", resp.Status_)
  177. sf.log.Err(err)
  178. return ctx.JSON(resp)
  179. }
  180. resp := sf.processSendRequest(req)
  181. resp.SelfCheck()
  182. return ctx.JSON(resp)
  183. }
  184. // Обрабатывает входящий запрос.
  185. func (sf *kBusHttp) processSendRequest(req *dot.ServeReq) *dot.ServeResp {
  186. req.SelfCheck()
  187. topic := topic.LetTopic(req.Topic_)
  188. res := sf.SendRequest(topic, req.BinReq_)
  189. resp := &dot.ServeResp{
  190. Status_: lev0.LetTxt("ok"),
  191. Uuid_: req.Uuid_,
  192. }
  193. if res.IsErr() {
  194. resp.Status_ = lev0.LetTxt("kernelBusHttp.processSendRequest(): err=\n\t%v", res.Err())
  195. return resp
  196. }
  197. resp.BinResp_ = res.Ok()
  198. return resp
  199. }
  200. // Процесс отписки от топика.
  201. func (sf *kBusHttp) processUnsubRequest(req *dot.UnsubReq) *dot.UnsubResp {
  202. req.SelfCheck()
  203. resp := &dot.UnsubResp{
  204. Status_: lev0.LetTxt("ok"),
  205. Uuid_: req.Uuid_,
  206. }
  207. optHandler := sf.KCtx_.Get(lev0.LetTxt(req.HandlerName_))
  208. if optHandler.IsNone() {
  209. resp.Status_ = lev0.LetTxt("kernelBusHttp.processUnsubRequest(): not get handler(%v) from kernel ctx",
  210. req.HandlerName_)
  211. return resp
  212. }
  213. if optHandler == nil {
  214. resp.Status_ = lev0.LetTxt("kernelBusHttp.processUnsubRequest(): handler(%v) not exists", req.HandlerName_)
  215. return resp
  216. }
  217. hand := optHandler.Some().Val().(bus_spec.IBusHandlerSubscribe)
  218. sf.Unsubscribe(hand)
  219. return resp
  220. }