Pārlūkot izejas kodu

增加生命周期管理

lijian 2 gadi atpakaļ
vecāks
revīzija
9e041f408c

+ 5 - 0
pkg/rpcs/controller.go

@@ -22,3 +22,8 @@ type ArgsOnEvent struct {
 	SubData     []byte
 }
 type ReplyOnEvent ReplyEmptyResult
+
+type ArgsRuleChainAct struct {
+	RuleChainId string
+	VendorId    string
+}

+ 51 - 0
services/controller/controller.go

@@ -6,6 +6,7 @@ import (
 	"github.com/gogf/gf/os/grpool"
 	"github.com/gogf/gf/util/guid"
 	"sparrow/pkg/actors"
+	"sparrow/pkg/entities"
 	"sparrow/pkg/klink"
 	"sparrow/pkg/protocol"
 	"sparrow/pkg/queue"
@@ -17,6 +18,7 @@ import (
 	"time"
 )
 
+//Controller 消息控制器
 type Controller struct {
 	producer     queue.QueueProducer
 	timer        *rule.Timer
@@ -27,6 +29,7 @@ type Controller struct {
 	pool         *grpool.Pool
 }
 
+//NewController 新建消息控制器
 func NewController(rabbithost string) (*Controller, error) {
 	admin := msgQueue.NewRabbitMessageQueueAdmin(&msgQueue.RabbitMqSettings{Host: rabbithost}, nil)
 	producer := msgQueue.NewRabbitMqProducer(admin, "default")
@@ -58,6 +61,7 @@ func NewController(rabbithost string) (*Controller, error) {
 	}, nil
 }
 
+// SetStatus 设置设备状态
 func (c *Controller) SetStatus(args rpcs.ArgsSetStatus, reply *rpcs.ReplySetStatus) error {
 	rpchost, err := getAccessRPCHost(args.DeviceId)
 	if err != nil {
@@ -67,6 +71,7 @@ func (c *Controller) SetStatus(args rpcs.ArgsSetStatus, reply *rpcs.ReplySetStat
 	return server.RPCCallByHost(rpchost, "Access.SetStatus", args, reply)
 }
 
+//GetStatus 获取设备状态
 func (c *Controller) GetStatus(args rpcs.ArgsGetStatus, reply *rpcs.ReplyGetStatus) error {
 	rpchost, err := getAccessRPCHost(args.Id)
 	if err != nil {
@@ -76,6 +81,7 @@ func (c *Controller) GetStatus(args rpcs.ArgsGetStatus, reply *rpcs.ReplyGetStat
 	return server.RPCCallByHost(rpchost, "Access.GetStatus", args, reply)
 }
 
+//Online 设备上线
 func (c *Controller) Online(args rpcs.ArgsGetStatus, reply *rpcs.ReplyEmptyResult) error {
 	data := gjson.New(nil)
 	_ = data.Set("device_id", args.Id)
@@ -104,6 +110,7 @@ func (c *Controller) Online(args rpcs.ArgsGetStatus, reply *rpcs.ReplyEmptyResul
 	return c.producer.Send(tpi, g, nil)
 }
 
+//Offline 设备下线
 func (c *Controller) Offline(args rpcs.ArgsGetStatus, reply *rpcs.ReplyEmptyResult) error {
 	if args.Id == "" || args.VendorId == "" {
 		return nil
@@ -135,6 +142,7 @@ func (c *Controller) Offline(args rpcs.ArgsGetStatus, reply *rpcs.ReplyEmptyResu
 	return c.producer.Send(tpi, g, nil)
 }
 
+//OnStatus 状态上报消息处理
 func (c *Controller) OnStatus(args rpcs.ArgsOnStatus, reply *rpcs.ReplyOnStatus) error {
 	t := time.Unix(int64(args.Timestamp/1000), 0)
 	data, err := c.processStatusToQueue(args)
@@ -195,6 +203,7 @@ func (c *Controller) processEventToQueue(args rpcs.ArgsOnEvent) (string, error)
 	return result.MustToJsonString(), nil
 }
 
+//OnEvent 事件消息处理
 func (c *Controller) OnEvent(args rpcs.ArgsOnEvent, reply *rpcs.ReplyOnEvent) error {
 	t := time.Unix(int64(args.TimeStamp/1000), 0)
 	data, err := c.processEventToQueue(args)
@@ -226,6 +235,47 @@ func (c *Controller) OnEvent(args rpcs.ArgsOnEvent, reply *rpcs.ReplyOnEvent) er
 	return c.producer.Send(tpi, g, nil)
 }
 
+// OnCreateRuleChain 规则链生命周期-创建
+func (c *Controller) CreateRuleChain(args rpcs.ArgsRuleChainAct, reply *rpcs.ReplyEmptyResult) error {
+	if c.actorContext != nil {
+		msg := &ruleEngine.ComponentLifecycleMsg{
+			TenantId:  args.VendorId,
+			EntityId:  &entities.RuleChainId{Id: args.RuleChainId},
+			EventType: ruleEngine.CREATED,
+		}
+		c.actorContext.AppActor.TellWithHighPriority(msg)
+	}
+	return nil
+}
+
+// DeleteRuleChain 规则链生命周期-删除
+func (c *Controller) DeleteRuleChain(args rpcs.ArgsRuleChainAct, reply *rpcs.ReplyEmptyResult) error {
+	if c.actorContext != nil {
+		msg := &ruleEngine.ComponentLifecycleMsg{
+			TenantId:  args.VendorId,
+			EntityId:  &entities.RuleChainId{Id: args.RuleChainId},
+			EventType: ruleEngine.DELETED,
+		}
+		c.actorContext.AppActor.TellWithHighPriority(msg)
+	}
+	return nil
+}
+
+// UpdateRuleChain 规则链生命周期-更新
+func (c *Controller) UpdateRuleChain(args rpcs.ArgsRuleChainAct, reply *rpcs.ReplyEmptyResult) error {
+	if c.actorContext != nil {
+		msg := &ruleEngine.ComponentLifecycleMsg{
+			TenantId:  args.VendorId,
+			EntityId:  &entities.RuleChainId{Id: args.RuleChainId},
+			EventType: ruleEngine.UPDATED,
+		}
+		c.actorContext.AppActor.TellWithHighPriority(msg)
+	}
+
+	return nil
+}
+
+//SendCommand 下发设备控制指令
 func (c *Controller) SendCommand(args rpcs.ArgsSendCommand, reply *rpcs.ReplySendCommand) error {
 	rpchost, err := getAccessRPCHost(args.DeviceId)
 	if err != nil {
@@ -247,6 +297,7 @@ func getAccessRPCHost(deviceid string) (string, error) {
 	return reply.AccessRPCHost, nil
 }
 
+//ActorSystem actor system
 type ActorSystem struct {
 	rootActor ruleEngine.Ref
 }

+ 14 - 0
vendor/github.com/jinzhu/gorm/go.mod

@@ -0,0 +1,14 @@
+module github.com/jinzhu/gorm
+
+go 1.12
+
+require (
+	github.com/denisenkom/go-mssqldb v0.0.0-20191124224453-732737034ffd
+	github.com/erikstmartin/go-testdb v0.0.0-20160219214506-8d10e4a1bae5
+	github.com/go-sql-driver/mysql v1.5.0
+	github.com/jinzhu/inflection v1.0.0
+	github.com/jinzhu/now v1.0.1
+	github.com/lib/pq v1.1.1
+	github.com/mattn/go-sqlite3 v1.14.0
+	golang.org/x/crypto v0.0.0-20191205180655-e7c4368fe9dd // indirect
+)

+ 33 - 0
vendor/github.com/jinzhu/gorm/go.sum

@@ -0,0 +1,33 @@
+github.com/PuerkitoBio/goquery v1.5.1/go.mod h1:GsLWisAFVj4WgDibEWF4pvYnkVQBpKBKeU+7zCJoLcc=
+github.com/andybalholm/cascadia v1.1.0/go.mod h1:GsXiBklL0woXo1j/WYWtSYYC4ouU9PqHO0sqidkEA4Y=
+github.com/denisenkom/go-mssqldb v0.0.0-20191124224453-732737034ffd h1:83Wprp6ROGeiHFAP8WJdI2RoxALQYgdllERc3N5N2DM=
+github.com/denisenkom/go-mssqldb v0.0.0-20191124224453-732737034ffd/go.mod h1:xbL0rPBG9cCiLr28tMa8zpbdarY27NDyej4t/EjAShU=
+github.com/erikstmartin/go-testdb v0.0.0-20160219214506-8d10e4a1bae5 h1:Yzb9+7DPaBjB8zlTR87/ElzFsnQfuHnVUVqpZZIcV5Y=
+github.com/erikstmartin/go-testdb v0.0.0-20160219214506-8d10e4a1bae5/go.mod h1:a2zkGnVExMxdzMo3M0Hi/3sEU+cWnZpSni0O6/Yb/P0=
+github.com/go-sql-driver/mysql v1.5.0 h1:ozyZYNQW3x3HtqT1jira07DN2PArx2v7/mN66gGcHOs=
+github.com/go-sql-driver/mysql v1.5.0/go.mod h1:DCzpHaOWr8IXmIStZouvnhqoel9Qv2LBy8hT2VhHyBg=
+github.com/golang-sql/civil v0.0.0-20190719163853-cb61b32ac6fe h1:lXe2qZdvpiX5WZkZR4hgp4KJVfY3nMkvmwbVkpv1rVY=
+github.com/golang-sql/civil v0.0.0-20190719163853-cb61b32ac6fe/go.mod h1:8vg3r2VgvsThLBIFL93Qb5yWzgyZWhEmBwUJWevAkK0=
+github.com/jinzhu/inflection v1.0.0 h1:K317FqzuhWc8YvSVlFMCCUb36O/S9MCKRDI7QkRKD/E=
+github.com/jinzhu/inflection v1.0.0/go.mod h1:h+uFLlag+Qp1Va5pdKtLDYj+kHp5pxUVkryuEj+Srlc=
+github.com/jinzhu/now v1.0.1 h1:HjfetcXq097iXP0uoPCdnM4Efp5/9MsM0/M+XOTeR3M=
+github.com/jinzhu/now v1.0.1/go.mod h1:d3SSVoowX0Lcu0IBviAWJpolVfI5UJVZZ7cO71lE/z8=
+github.com/lib/pq v1.1.1 h1:sJZmqHoEaY7f+NPP8pgLB/WxulyR3fewgCM2qaSlBb4=
+github.com/lib/pq v1.1.1/go.mod h1:5WUZQaWbwv1U+lTReE5YruASi9Al49XbQIvNi/34Woo=
+github.com/mattn/go-sqlite3 v1.14.0 h1:mLyGNKR8+Vv9CAU7PphKa2hkEqxxhn8i32J6FPj1/QA=
+github.com/mattn/go-sqlite3 v1.14.0/go.mod h1:JIl7NbARA7phWnGvh0LKTyg7S9BA+6gx71ShQilpsus=
+github.com/mattn/go-sqlite3 v2.0.1+incompatible h1:xQ15muvnzGBHpIpdrNi1DA5x0+TcBZzsIDwmw9uTHzw=
+github.com/mattn/go-sqlite3 v2.0.1+incompatible/go.mod h1:FPy6KqzDD04eiIsT53CuJW3U88zkxoIYsOqkbpncsNc=
+golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w=
+golang.org/x/crypto v0.0.0-20190325154230-a5d413f7728c h1:Vj5n4GlwjmQteupaxJ9+0FNOmBrHfq7vN4btdGoDZgI=
+golang.org/x/crypto v0.0.0-20190325154230-a5d413f7728c/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w=
+golang.org/x/crypto v0.0.0-20191205180655-e7c4368fe9dd h1:GGJVjV8waZKRHrgwvtH66z9ZGVurTD1MT0n1Bb+q4aM=
+golang.org/x/crypto v0.0.0-20191205180655-e7c4368fe9dd/go.mod h1:LzIPMQfyMNhhGPhUkYOs5KpL4U8rLKemX1yGLhDgUto=
+golang.org/x/net v0.0.0-20180218175443-cbe0f9307d01/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4=
+golang.org/x/net v0.0.0-20190404232315-eb5bcb51f2a3/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg=
+golang.org/x/net v0.0.0-20200202094626-16171245cfb2/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s=
+golang.org/x/net v0.0.0-20200324143707-d3edc9973b7e/go.mod h1:qpuaurCH72eLCgpAm/N6yyVIVM9cpaDIP3A8BGJEC5A=
+golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
+golang.org/x/sys v0.0.0-20190412213103-97732733099d/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
+golang.org/x/sys v0.0.0-20200323222414-85ca7c5b95cd/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
+golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ=