controller.go 8.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308
  1. package main
  2. import (
  3. "fmt"
  4. "github.com/gogf/gf/encoding/gjson"
  5. "github.com/gogf/gf/os/grpool"
  6. "github.com/gogf/gf/util/guid"
  7. "sparrow/pkg/actors"
  8. "sparrow/pkg/klink"
  9. "sparrow/pkg/protocol"
  10. "sparrow/pkg/queue"
  11. "sparrow/pkg/queue/msgQueue"
  12. "sparrow/pkg/rpcs"
  13. "sparrow/pkg/rule"
  14. "sparrow/pkg/ruleEngine"
  15. "sparrow/pkg/server"
  16. "time"
  17. )
  18. type Controller struct {
  19. producer queue.QueueProducer
  20. timer *rule.Timer
  21. ift *rule.Ifttt
  22. actorContext *ruleEngine.SystemContext
  23. consumer queue.QueueConsumer
  24. cluster *ClusterService
  25. pool *grpool.Pool
  26. }
  27. func NewController(rabbithost string) (*Controller, error) {
  28. admin := msgQueue.NewRabbitMessageQueueAdmin(&msgQueue.RabbitMqSettings{Host: rabbithost}, nil)
  29. producer := msgQueue.NewRabbitMqProducer(admin, "default")
  30. consumer := msgQueue.NewRabbitConsumer(admin, "MAIN")
  31. tp := make([]*queue.TopicPartitionInfo, 0)
  32. tp = append(tp, &queue.TopicPartitionInfo{
  33. Topic: "MAIN",
  34. TenantId: "1ps9djpswi0cds7cofynkso300eql4iu",
  35. Partition: 0,
  36. MyPartition: false,
  37. })
  38. tp = append(tp, &queue.TopicPartitionInfo{
  39. Topic: "MAIN",
  40. TenantId: "1ps9djpswi0cds7cofynkso300eql4sw",
  41. Partition: 0,
  42. MyPartition: false,
  43. })
  44. _ = consumer.SubscribeWithPartitions(tp)
  45. if err := producer.Init(); err != nil {
  46. return nil, err
  47. }
  48. return &Controller{
  49. producer: producer,
  50. consumer: consumer,
  51. cluster: &ClusterService{producer: producer},
  52. pool: grpool.New(),
  53. }, nil
  54. }
  55. func (c *Controller) SetStatus(args rpcs.ArgsSetStatus, reply *rpcs.ReplySetStatus) error {
  56. rpchost, err := getAccessRPCHost(args.DeviceId)
  57. if err != nil {
  58. return err
  59. }
  60. return server.RPCCallByHost(rpchost, "Access.SetStatus", args, reply)
  61. }
  62. func (c *Controller) GetStatus(args rpcs.ArgsGetStatus, reply *rpcs.ReplyGetStatus) error {
  63. rpchost, err := getAccessRPCHost(args.Id)
  64. if err != nil {
  65. return err
  66. }
  67. return server.RPCCallByHost(rpchost, "Access.GetStatus", args, reply)
  68. }
  69. func (c *Controller) Online(args rpcs.ArgsGetStatus, reply *rpcs.ReplyEmptyResult) error {
  70. data := gjson.New(nil)
  71. _ = data.Set("device_id", args.Id)
  72. t := time.Now()
  73. msg := &protocol.Message{
  74. Id: guid.S(),
  75. Ts: &t,
  76. Type: protocol.CONNECT_EVENT,
  77. Data: data.MustToJsonString(),
  78. Callback: nil,
  79. MetaData: map[string]interface{}{
  80. "device_id": args.Id,
  81. "vendor_id": args.VendorId,
  82. },
  83. Originator: "device",
  84. }
  85. tpi := queue.ResolvePartition("RULE_ENGINE",
  86. msg.GetQueueName(),
  87. args.VendorId,
  88. args.Id)
  89. g, err := queue.NewGobQueueMessage(msg)
  90. if err != nil {
  91. return err
  92. }
  93. g.Headers.Put("tenant_id", []byte(args.VendorId))
  94. return c.producer.Send(tpi, g, nil)
  95. }
  96. func (c *Controller) Offline(args rpcs.ArgsGetStatus, reply *rpcs.ReplyEmptyResult) error {
  97. data := gjson.New(nil)
  98. _ = data.Set("device_id", args.Id)
  99. t := time.Now()
  100. msg := &protocol.Message{
  101. Id: guid.S(),
  102. Ts: &t,
  103. Type: protocol.DISCONNECT_EVENT,
  104. Data: data.MustToJsonString(),
  105. Callback: nil,
  106. MetaData: map[string]interface{}{
  107. "device_id": args.Id,
  108. "vendor_id": args.VendorId,
  109. },
  110. Originator: "device",
  111. }
  112. tpi := queue.ResolvePartition("RULE_ENGINE",
  113. msg.GetQueueName(),
  114. args.VendorId,
  115. args.Id)
  116. g, err := queue.NewGobQueueMessage(msg)
  117. if err != nil {
  118. return err
  119. }
  120. g.Headers.Put("tenant_id", []byte(args.VendorId))
  121. return c.producer.Send(tpi, g, nil)
  122. }
  123. func (c *Controller) OnStatus(args rpcs.ArgsOnStatus, reply *rpcs.ReplyOnStatus) error {
  124. t := time.Unix(int64(args.Timestamp/1000), 0)
  125. data, err := c.processStatusToQueue(args)
  126. if err != nil {
  127. return err
  128. }
  129. msg := &protocol.Message{
  130. Id: guid.S(),
  131. Ts: &t,
  132. Type: protocol.POST_ATTRIBUTES_REQUEST,
  133. Data: data,
  134. Callback: nil,
  135. MetaData: map[string]interface{}{
  136. "tenant_id": args.VendorId,
  137. "device_id": args.DeviceId,
  138. "sub_device_id": args.SubDeviceId,
  139. },
  140. Originator: "device",
  141. }
  142. tpi := queue.ResolvePartition("RULE_ENGINE",
  143. msg.GetQueueName(),
  144. args.VendorId,
  145. args.DeviceId)
  146. g, err := queue.NewGobQueueMessage(msg)
  147. if err != nil {
  148. return err
  149. }
  150. g.Headers.Put("tenant_id", []byte(args.VendorId))
  151. return c.producer.Send(tpi, g, nil)
  152. }
  153. func (c *Controller) processStatusToQueue(args rpcs.ArgsOnStatus) (string, error) {
  154. result := gjson.New(nil)
  155. j, err := gjson.DecodeToJson(args.SubData)
  156. if err != nil {
  157. return "", err
  158. }
  159. switch args.Action {
  160. case klink.DevSendAction:
  161. params := j.GetMap("params")
  162. if err = result.Set(j.GetString("cmd"), params); err != nil {
  163. return "", err
  164. }
  165. }
  166. return result.MustToJsonString(), nil
  167. }
  168. func (c *Controller) processEventToQueue(args rpcs.ArgsOnEvent) (string, error) {
  169. result := gjson.New(nil)
  170. j, err := gjson.DecodeToJson(args.SubData)
  171. if err != nil {
  172. return "", nil
  173. }
  174. params := j.GetMap("params")
  175. if err = result.Set(j.GetString("cmd"), params); err != nil {
  176. return "", err
  177. }
  178. return result.MustToJsonString(), nil
  179. }
  180. func (c *Controller) OnEvent(args rpcs.ArgsOnEvent, reply *rpcs.ReplyOnEvent) error {
  181. t := time.Unix(int64(args.TimeStamp/1000), 0)
  182. data, err := c.processEventToQueue(args)
  183. if err != nil {
  184. return err
  185. }
  186. msg := &protocol.Message{
  187. Id: guid.S(),
  188. Ts: &t,
  189. Type: protocol.POST_EVENT_REQUEST,
  190. Data: data,
  191. Callback: nil,
  192. MetaData: map[string]interface{}{
  193. "tenant_id": args.VendorId,
  194. "device_id": args.DeviceId,
  195. "sub_device_id": args.SubDeviceId,
  196. },
  197. Originator: "device",
  198. }
  199. tpi := queue.ResolvePartition("RULE_ENGINE",
  200. msg.GetQueueName(),
  201. args.VendorId,
  202. args.DeviceId)
  203. g, err := queue.NewGobQueueMessage(msg)
  204. if err != nil {
  205. return err
  206. }
  207. g.Headers.Put("tenant_id", []byte(args.VendorId))
  208. return c.producer.Send(tpi, g, nil)
  209. }
  210. func (c *Controller) SendCommand(args rpcs.ArgsSendCommand, reply *rpcs.ReplySendCommand) error {
  211. rpchost, err := getAccessRPCHost(args.DeviceId)
  212. if err != nil {
  213. return err
  214. }
  215. return server.RPCCallByHost(rpchost, "Access.SendCommand", args, reply)
  216. }
  217. func getAccessRPCHost(deviceid string) (string, error) {
  218. args := rpcs.ArgsGetDeviceOnlineStatus{
  219. Id: deviceid,
  220. }
  221. reply := &rpcs.ReplyGetDeviceOnlineStatus{}
  222. err := server.RPCCallByName(nil, rpcs.DeviceManagerName, "DeviceManager.GetDeviceOnlineStatus", args, reply)
  223. if err != nil {
  224. return "", err
  225. }
  226. return reply.AccessRPCHost, nil
  227. }
  228. type ActorSystem struct {
  229. rootActor ruleEngine.Ref
  230. }
  231. // 初始化actor system
  232. func (c *Controller) initActorSystem() (*ActorSystem, error) {
  233. system := ruleEngine.NewDefaultActorSystem(&ruleEngine.DefaultActorSystemConfig{
  234. SchedulerPoolSize: 5,
  235. AppDispatcherPoolSize: 4,
  236. TenantDispatcherPoolSize: 4,
  237. RuleEngineDispatcherPoolSize: 4,
  238. })
  239. _ = system.CreateDispatcher(ruleEngine.APP_DISPATCHER_NAME, ruleEngine.NewPoolDispatcher(0))
  240. _ = system.CreateDispatcher(ruleEngine.TENANT_DISPATCHER_NAME, ruleEngine.NewPoolDispatcher(0))
  241. _ = system.CreateDispatcher(ruleEngine.RULE_DISPATCHER_NAME, ruleEngine.NewPoolDispatcher(0))
  242. // init services
  243. tenantService := &TenantService{}
  244. ruleChainService := &RuleChainService{}
  245. actorContext := ruleEngine.NewSystemContext(system, ruleEngine.SystemContextServiceConfig{
  246. ClusterService: c.cluster,
  247. RuleChainService: ruleChainService,
  248. TenantService: tenantService,
  249. EventService: NewEventService(),
  250. })
  251. appActor, err := system.CreateRootActor(ruleEngine.APP_DISPATCHER_NAME,
  252. actors.NewAppActorCreator(actorContext))
  253. if err != nil {
  254. return nil, err
  255. }
  256. actorContext.AppActor = appActor
  257. server.Log.Debugln("actor system initialized")
  258. time.Sleep(time.Second * 1)
  259. appActor.Tell(&ruleEngine.AppInitMsg{})
  260. c.actorContext = actorContext
  261. return &ActorSystem{rootActor: appActor}, nil
  262. }
  263. // 启动mq consumers
  264. func (c *Controller) launchConsumer() {
  265. msgs, err := c.consumer.Pop(100 * time.Millisecond)
  266. if err != nil {
  267. server.Log.Error(err)
  268. }
  269. for {
  270. select {
  271. case msg := <-msgs:
  272. ruleEngineMsg := &protocol.Message{}
  273. if err := ruleEngineMsg.Decode(msg.GetData()); err != nil {
  274. fmt.Println("解析消息失败")
  275. }
  276. tanantId := msg.GetHeaders().Get("tenant_id")
  277. if c.actorContext != nil {
  278. c.actorContext.Tell(&ruleEngine.QueueToRuleEngineMsg{
  279. TenantId: string(tanantId),
  280. Message: ruleEngineMsg,
  281. })
  282. }
  283. }
  284. }
  285. }