kbus_http.go 7.1 KB

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