kbus_http.go 6.9 KB

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