Co-authored-by: Gisle Aune <dev@gisle.me> Reviewed-on: #1 Co-authored-by: gisle <gisle@hidden-email> Co-committed-by: gisle <gisle@hidden-email>
This commit was merged in pull request #1.
This commit is contained in:
+95
-12
@@ -3,9 +3,11 @@ package api
|
||||
import (
|
||||
"context"
|
||||
"git.aiterp.net/lucifer/new-server/app/config"
|
||||
"git.aiterp.net/lucifer/new-server/app/services/publisher"
|
||||
"git.aiterp.net/lucifer/new-server/models"
|
||||
"github.com/gin-gonic/gin"
|
||||
"log"
|
||||
"time"
|
||||
)
|
||||
|
||||
func fetchDevices(ctx context.Context, fetchStr string) ([]models.Device, error) {
|
||||
@@ -15,11 +17,21 @@ func fetchDevices(ctx context.Context, fetchStr string) ([]models.Device, error)
|
||||
|
||||
func Devices(r gin.IRoutes) {
|
||||
r.GET("", handler(func(c *gin.Context) (interface{}, error) {
|
||||
return config.DeviceRepository().FetchByReference(ctxOf(c), models.RKAll, "")
|
||||
devices, err := config.DeviceRepository().FetchByReference(ctxOf(c), models.RKAll, "")
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return withSceneState(devices), nil
|
||||
}))
|
||||
|
||||
r.GET("/:fetch", handler(func(c *gin.Context) (interface{}, error) {
|
||||
return fetchDevices(ctxOf(c), c.Param("fetch"))
|
||||
devices, err := fetchDevices(ctxOf(c), c.Param("fetch"))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return withSceneState(devices), nil
|
||||
}))
|
||||
|
||||
r.PUT("", handler(func(c *gin.Context) (interface{}, error) {
|
||||
@@ -70,7 +82,7 @@ func Devices(r gin.IRoutes) {
|
||||
}
|
||||
}()
|
||||
|
||||
return changed, nil
|
||||
return withSceneState(changed), nil
|
||||
}))
|
||||
|
||||
r.PUT("/:fetch", handler(func(c *gin.Context) (interface{}, error) {
|
||||
@@ -98,7 +110,7 @@ func Devices(r gin.IRoutes) {
|
||||
}
|
||||
}
|
||||
|
||||
return devices, nil
|
||||
return withSceneState(devices), nil
|
||||
}))
|
||||
|
||||
r.PUT("/:fetch/state", handler(func(c *gin.Context) (interface{}, error) {
|
||||
@@ -126,16 +138,13 @@ func Devices(r gin.IRoutes) {
|
||||
config.PublishChannel <- devices
|
||||
|
||||
go func() {
|
||||
for _, device := range devices {
|
||||
err := config.DeviceRepository().Save(context.Background(), &device, models.SMState)
|
||||
if err != nil {
|
||||
log.Println("Failed to save device for state:", err)
|
||||
continue
|
||||
}
|
||||
err = config.DeviceRepository().SaveMany(ctxOf(c), models.SMState, devices)
|
||||
if err != nil {
|
||||
log.Println("Failed to save devices states")
|
||||
}
|
||||
}()
|
||||
|
||||
return devices, nil
|
||||
return withSceneState(devices), nil
|
||||
}))
|
||||
|
||||
r.PUT("/:fetch/tags", handler(func(c *gin.Context) (interface{}, error) {
|
||||
@@ -192,6 +201,80 @@ func Devices(r gin.IRoutes) {
|
||||
}
|
||||
}
|
||||
|
||||
return devices, nil
|
||||
return withSceneState(devices), nil
|
||||
}))
|
||||
|
||||
r.PUT("/:fetch/scene", handler(func(c *gin.Context) (interface{}, error) {
|
||||
var body models.DeviceSceneAssignment
|
||||
err := parseBody(c, &body)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
devices, err := fetchDevices(ctxOf(c), c.Param("fetch"))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if len(devices) == 0 {
|
||||
return []models.Device{}, nil
|
||||
}
|
||||
|
||||
_, err = config.SceneRepository().Find(ctxOf(c), body.SceneID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if body.DurationMS < 0 {
|
||||
body.DurationMS = 0
|
||||
}
|
||||
body.StartTime = time.Now()
|
||||
|
||||
pushMode := c.Query("push") == "true"
|
||||
for i := range devices {
|
||||
if pushMode {
|
||||
devices[i].SceneAssignments = append(devices[i].SceneAssignments, body)
|
||||
} else {
|
||||
devices[i].SceneAssignments = []models.DeviceSceneAssignment{body}
|
||||
}
|
||||
}
|
||||
config.PublishChannel <- devices
|
||||
|
||||
err = config.DeviceRepository().SaveMany(ctxOf(c), 0, devices)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return withSceneState(devices), nil
|
||||
}))
|
||||
|
||||
r.DELETE("/:fetch/scene", handler(func(c *gin.Context) (interface{}, error) {
|
||||
devices, err := fetchDevices(ctxOf(c), c.Param("fetch"))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if len(devices) == 0 {
|
||||
return []models.Device{}, nil
|
||||
}
|
||||
|
||||
for i := range devices {
|
||||
devices[i].SceneAssignments = []models.DeviceSceneAssignment{}
|
||||
}
|
||||
config.PublishChannel <- devices
|
||||
|
||||
err = config.DeviceRepository().SaveMany(ctxOf(c), 0, devices)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return withSceneState(devices), nil
|
||||
}))
|
||||
}
|
||||
|
||||
func withSceneState(devices []models.Device) []models.Device {
|
||||
res := make([]models.Device, 0, len(devices))
|
||||
for _, device := range devices {
|
||||
device.SceneState = publisher.Global().SceneState(device.ID)
|
||||
res = append(res, device)
|
||||
}
|
||||
|
||||
return res
|
||||
}
|
||||
|
||||
@@ -0,0 +1,68 @@
|
||||
package api
|
||||
|
||||
import (
|
||||
"git.aiterp.net/lucifer/new-server/app/config"
|
||||
"git.aiterp.net/lucifer/new-server/app/services/publisher"
|
||||
"git.aiterp.net/lucifer/new-server/models"
|
||||
"github.com/gin-gonic/gin"
|
||||
)
|
||||
|
||||
func Scenes(r gin.IRoutes) {
|
||||
r.GET("", handler(func(c *gin.Context) (interface{}, error) {
|
||||
return config.SceneRepository().FetchAll(ctxOf(c))
|
||||
}))
|
||||
|
||||
r.GET("/:id", handler(func(c *gin.Context) (interface{}, error) {
|
||||
return config.SceneRepository().Find(ctxOf(c), intParam(c, "id"))
|
||||
}))
|
||||
|
||||
r.POST("", handler(func(c *gin.Context) (interface{}, error) {
|
||||
var body models.Scene
|
||||
err := parseBody(c, &body)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
err = body.Validate()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
err = config.SceneRepository().Save(ctxOf(c), &body)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
publisher.Global().UpdateScene(body)
|
||||
|
||||
return body, nil
|
||||
}))
|
||||
|
||||
r.PUT("/:id", handler(func(c *gin.Context) (interface{}, error) {
|
||||
var body models.Scene
|
||||
err := parseBody(c, &body)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
scene, err := config.SceneRepository().Find(ctxOf(c), intParam(c, "id"))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
body.ID = scene.ID
|
||||
err = body.Validate()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
err = config.SceneRepository().Save(ctxOf(c), &body)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
publisher.Global().UpdateScene(body)
|
||||
|
||||
return body, nil
|
||||
}))
|
||||
}
|
||||
@@ -15,6 +15,13 @@ var errorMap = map[error]int{
|
||||
models.ErrBadColor: 400,
|
||||
models.ErrInternal: 500,
|
||||
models.ErrUnknownColorFormat: 400,
|
||||
|
||||
models.ErrSceneInvalidInterval: 400,
|
||||
models.ErrSceneNoRoles: 400,
|
||||
models.ErrSceneRoleNoStates: 400,
|
||||
models.ErrSceneRoleUnsupportedOrdering: 422,
|
||||
models.ErrSceneRoleUnknownEffect: 422,
|
||||
models.ErrSceneRoleUnknownPowerMode: 422,
|
||||
}
|
||||
|
||||
type response struct {
|
||||
|
||||
@@ -60,6 +60,31 @@ func (client *Client) PutDeviceTags(ctx context.Context, fetchStr string, addTag
|
||||
return devices, nil
|
||||
}
|
||||
|
||||
func (client *Client) AssignDevice(ctx context.Context, fetchStr string, push bool, assignment models.DeviceSceneAssignment) ([]models.Device, error) {
|
||||
query := ""
|
||||
if push {
|
||||
query = "?push=true"
|
||||
}
|
||||
|
||||
devices := make([]models.Device, 0, 16)
|
||||
err := client.Fetch(ctx, "PUT", "/api/devices/"+fetchStr+"/scene"+query, &devices, assignment)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return devices, nil
|
||||
}
|
||||
|
||||
func (client *Client) ClearDevice(ctx context.Context, fetchStr string) ([]models.Device, error) {
|
||||
devices := make([]models.Device, 0, 16)
|
||||
err := client.Fetch(ctx, "DELETE", "/api/devices/"+fetchStr+"/scene", &devices, nil)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return devices, nil
|
||||
}
|
||||
|
||||
func (client *Client) FireEvent(ctx context.Context, event models.Event) error {
|
||||
err := client.Fetch(ctx, "POST", "/api/events", nil, event)
|
||||
if err != nil {
|
||||
|
||||
@@ -28,7 +28,7 @@ func (client *Client) GetHandler(ctx context.Context, id int) (*models.EventHand
|
||||
|
||||
func (client *Client) PutHandler(ctx context.Context, handler *models.EventHandler) (*models.EventHandler, error) {
|
||||
var response models.EventHandler
|
||||
err := client.Fetch(ctx, "PUT", fmt.Sprintf("/api/event-handlers/%d", handler.ID), &response, handler)
|
||||
err := client.Fetch(ctx, "PUT", fmt.Sprintf("/api/event-handlers/%d?hard=true", handler.ID), &response, handler)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
@@ -0,0 +1,16 @@
|
||||
package client
|
||||
|
||||
import (
|
||||
"context"
|
||||
"git.aiterp.net/lucifer/new-server/models"
|
||||
)
|
||||
|
||||
func (client *Client) GetScenes(ctx context.Context) ([]models.Scene, error) {
|
||||
scenes := make([]models.Scene, 0, 16)
|
||||
err := client.Fetch(ctx, "GET", "/api/scenes", &scenes, nil)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return scenes, nil
|
||||
}
|
||||
+3
-2
@@ -24,8 +24,9 @@ func DBX() *sqlx.DB {
|
||||
MySqlSchema(),
|
||||
))
|
||||
|
||||
dbx.SetMaxIdleConns(50)
|
||||
dbx.SetMaxOpenConns(100)
|
||||
dbx.SetMaxIdleConns(20)
|
||||
dbx.SetMaxOpenConns(40)
|
||||
dbx.SetConnMaxLifetime(0)
|
||||
}
|
||||
|
||||
return dbx
|
||||
|
||||
@@ -20,3 +20,7 @@ func DeviceRepository() models.DeviceRepository {
|
||||
func EventHandlerRepository() models.EventHandlerRepository {
|
||||
return &mysql.EventHandlerRepo{DBX: DBX()}
|
||||
}
|
||||
|
||||
func SceneRepository() models.SceneRepository {
|
||||
return &mysql.SceneRepo{DBX: DBX()}
|
||||
}
|
||||
|
||||
+13
-1
@@ -1,17 +1,28 @@
|
||||
package app
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"git.aiterp.net/lucifer/new-server/app/api"
|
||||
"git.aiterp.net/lucifer/new-server/app/config"
|
||||
"git.aiterp.net/lucifer/new-server/app/services"
|
||||
"git.aiterp.net/lucifer/new-server/app/services/publisher"
|
||||
"github.com/gin-gonic/gin"
|
||||
"log"
|
||||
"time"
|
||||
)
|
||||
|
||||
func StartServer() {
|
||||
setupCtx, cancel := context.WithTimeout(context.Background(), time.Second * 10)
|
||||
defer cancel()
|
||||
|
||||
err := publisher.Initialize(setupCtx)
|
||||
if err != nil {
|
||||
log.Fatalln("Publish init failed:", err)
|
||||
return
|
||||
}
|
||||
|
||||
services.StartEventHandler()
|
||||
services.StartPublisher()
|
||||
services.ConnectToBridges()
|
||||
services.CheckNewDevices()
|
||||
|
||||
@@ -25,6 +36,7 @@ func StartServer() {
|
||||
api.DriverKinds(apiGin.Group("/driver-kinds"))
|
||||
api.Events(apiGin.Group("/events"))
|
||||
api.EventHandlers(apiGin.Group("/event-handlers"))
|
||||
api.Scenes(apiGin.Group("/scenes"))
|
||||
|
||||
log.Fatal(ginny.Run(fmt.Sprintf("0.0.0.0:%d", config.ServerPort())))
|
||||
}
|
||||
|
||||
@@ -5,7 +5,6 @@ import (
|
||||
"git.aiterp.net/lucifer/new-server/app/config"
|
||||
"git.aiterp.net/lucifer/new-server/models"
|
||||
"log"
|
||||
"strconv"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
@@ -53,15 +52,6 @@ func runConnectToBridges() error {
|
||||
log.Printf("Running bridge \"%s\" (%d)", bridge.Name, bridge.ID)
|
||||
|
||||
go func(bridge models.Bridge, cancel func()) {
|
||||
savedDevices, err := config.DeviceRepository().FetchByReference(ctx, models.RKBridgeID, strconv.Itoa(bridge.ID))
|
||||
if err != nil {
|
||||
log.Println("Failed to fetch devices from db for refresh:", err)
|
||||
}
|
||||
err = driver.Publish(ctx, bridge, savedDevices)
|
||||
if err != nil {
|
||||
log.Println("Failed to publish devices from db before run:", err)
|
||||
}
|
||||
|
||||
err = driver.Run(ctx, bridge, config.EventChannel)
|
||||
log.Printf("Bridge \"%s\" (%d) stopped: %s", bridge.Name, bridge.ID, err)
|
||||
|
||||
|
||||
+16
-6
@@ -155,23 +155,33 @@ func handleEvent(event models.Event) (responses []models.Event) {
|
||||
if err != nil {
|
||||
log.Println("Error updating state for device", device.ID, "err:", err)
|
||||
}
|
||||
|
||||
if action.SetScene != nil {
|
||||
action.SetScene.StartTime = time.Now()
|
||||
allDevices[i].SceneAssignments = []models.DeviceSceneAssignment{*action.SetScene}
|
||||
}
|
||||
if action.PushScene != nil {
|
||||
action.PushScene.StartTime = time.Now()
|
||||
allDevices[i].SceneAssignments = append(allDevices[i].SceneAssignments, *action.PushScene)
|
||||
}
|
||||
}
|
||||
|
||||
config.PublishChannel <- allDevices
|
||||
|
||||
wg := sync.WaitGroup{}
|
||||
for _, device := range allDevices {
|
||||
wg.Add(1)
|
||||
|
||||
go func(device models.Device) {
|
||||
err := config.DeviceRepository().Save(context.Background(), &device, models.SMState)
|
||||
if len(allDevices) > 0 {
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
err := config.DeviceRepository().SaveMany(context.Background(), models.SMState, allDevices)
|
||||
if err != nil {
|
||||
log.Println("Failed to save device for state:", err)
|
||||
log.Println("Failed to save devices' state:", err)
|
||||
}
|
||||
|
||||
wg.Done()
|
||||
}(device)
|
||||
}()
|
||||
}
|
||||
|
||||
for _, handler := range deadHandlers {
|
||||
wg.Add(1)
|
||||
|
||||
|
||||
@@ -1,76 +0,0 @@
|
||||
package services
|
||||
|
||||
import (
|
||||
"context"
|
||||
"git.aiterp.net/lucifer/new-server/app/config"
|
||||
"git.aiterp.net/lucifer/new-server/models"
|
||||
"log"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
func StartPublisher() {
|
||||
ctx := context.Background()
|
||||
|
||||
go func() {
|
||||
for devices := range config.PublishChannel {
|
||||
if len(devices) == 0 {
|
||||
continue
|
||||
}
|
||||
|
||||
lists := make(map[int][]models.Device, 4)
|
||||
for _, device := range devices {
|
||||
lists[device.BridgeID] = append(lists[device.BridgeID], device)
|
||||
}
|
||||
|
||||
ctx, cancel := context.WithTimeout(ctx, time.Second * 30)
|
||||
|
||||
bridges, err := config.BridgeRepository().FetchAll(ctx)
|
||||
if err != nil {
|
||||
log.Println("Publishing error (1): " + err.Error())
|
||||
cancel()
|
||||
continue
|
||||
}
|
||||
|
||||
wg := sync.WaitGroup{}
|
||||
for _, devices := range lists {
|
||||
wg.Add(1)
|
||||
|
||||
go func(devices []models.Device) {
|
||||
defer wg.Done()
|
||||
|
||||
var bridge models.Bridge
|
||||
for _, bridge2 := range bridges {
|
||||
if bridge2.ID == devices[0].BridgeID {
|
||||
bridge = bridge2
|
||||
}
|
||||
}
|
||||
if bridge.ID == 0 {
|
||||
log.Println("Unknown bridge")
|
||||
}
|
||||
|
||||
bridge, err := config.BridgeRepository().Find(ctx, devices[0].BridgeID)
|
||||
if err != nil {
|
||||
log.Println("Publishing error (1): " + err.Error())
|
||||
return
|
||||
}
|
||||
|
||||
driver, err := config.DriverProvider().Provide(bridge.Driver)
|
||||
if err != nil {
|
||||
log.Println("Publishing error (2): " + err.Error())
|
||||
return
|
||||
}
|
||||
|
||||
err = driver.Publish(ctx, bridge, devices)
|
||||
if err != nil {
|
||||
log.Println("Publishing error (3): " + err.Error())
|
||||
return
|
||||
}
|
||||
}(devices)
|
||||
}
|
||||
|
||||
wg.Wait()
|
||||
cancel()
|
||||
}
|
||||
}()
|
||||
}
|
||||
@@ -0,0 +1,315 @@
|
||||
package publisher
|
||||
|
||||
import (
|
||||
"context"
|
||||
"git.aiterp.net/lucifer/new-server/app/config"
|
||||
"git.aiterp.net/lucifer/new-server/models"
|
||||
"log"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
type Publisher struct {
|
||||
mu sync.Mutex
|
||||
sceneData map[int]*models.Scene
|
||||
scenes []*Scene
|
||||
sceneAssignment map[int]*Scene
|
||||
started map[int]bool
|
||||
pending map[int][]models.Device
|
||||
waiting map[int]chan struct{}
|
||||
}
|
||||
|
||||
func (p *Publisher) SceneState(deviceID int) *models.DeviceState {
|
||||
p.mu.Lock()
|
||||
defer p.mu.Unlock()
|
||||
|
||||
if s := p.sceneAssignment[deviceID]; s != nil {
|
||||
return s.LastState(deviceID)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (p *Publisher) UpdateScene(data models.Scene) {
|
||||
p.mu.Lock()
|
||||
p.sceneData[data.ID] = &data
|
||||
|
||||
for _, scene := range p.scenes {
|
||||
if scene.data.ID == data.ID {
|
||||
scene.UpdateScene(data)
|
||||
}
|
||||
}
|
||||
p.mu.Unlock()
|
||||
}
|
||||
|
||||
func (p *Publisher) ReloadScenes(ctx context.Context) error {
|
||||
scenes, err := config.SceneRepository().FetchAll(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
p.mu.Lock()
|
||||
for i, scene := range scenes {
|
||||
p.sceneData[scene.ID] = &scenes[i]
|
||||
}
|
||||
|
||||
for _, scene := range p.scenes {
|
||||
scene.UpdateScene(*p.sceneData[scene.data.ID])
|
||||
}
|
||||
p.mu.Unlock()
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (p *Publisher) ReloadDevices(ctx context.Context) error {
|
||||
devices, err := config.DeviceRepository().FetchByReference(ctx, models.RKAll, "")
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
p.Publish(devices...)
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (p *Publisher) Publish(devices ...models.Device) {
|
||||
if len(devices) == 0 {
|
||||
return
|
||||
}
|
||||
|
||||
p.mu.Lock()
|
||||
defer p.mu.Unlock()
|
||||
|
||||
for _, device := range devices {
|
||||
if !p.started[device.BridgeID] {
|
||||
p.started[device.BridgeID] = true
|
||||
go p.runBridge(device.BridgeID)
|
||||
}
|
||||
|
||||
p.reassignDevice(device)
|
||||
|
||||
if p.sceneAssignment[device.ID] != nil {
|
||||
p.sceneAssignment[device.ID].UpsertDevice(device)
|
||||
} else {
|
||||
p.pending[device.BridgeID] = append(p.pending[device.BridgeID], device)
|
||||
if p.waiting[device.BridgeID] != nil {
|
||||
close(p.waiting[device.BridgeID])
|
||||
p.waiting[device.BridgeID] = nil
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (p *Publisher) PublishChannel(ch <-chan []models.Device) {
|
||||
for list := range ch {
|
||||
p.Publish(list...)
|
||||
}
|
||||
}
|
||||
|
||||
func (p *Publisher) Run() {
|
||||
ticker := time.NewTicker(time.Millisecond * 100)
|
||||
deleteList := make([]int, 0, 8)
|
||||
updatedList := make([]models.Device, 0, 16)
|
||||
|
||||
for range ticker.C {
|
||||
deleteList = deleteList[:0]
|
||||
updatedList = updatedList[:0]
|
||||
|
||||
p.mu.Lock()
|
||||
for i, scene := range p.scenes {
|
||||
if (!scene.endTime.IsZero() && time.Now().After(scene.endTime)) || scene.Empty() {
|
||||
deleteList = append(deleteList, i-len(deleteList))
|
||||
updatedList = append(updatedList, scene.AllDevices()...)
|
||||
|
||||
log.Printf("Removing scene instance for %s (%d)", scene.data.Name, scene.data.ID)
|
||||
|
||||
for _, device := range scene.AllDevices() {
|
||||
p.sceneAssignment[device.ID] = nil
|
||||
}
|
||||
continue
|
||||
}
|
||||
|
||||
if scene.Due() {
|
||||
updatedList = append(updatedList, scene.Run()...)
|
||||
updatedList = append(updatedList, scene.UnaffectedDevices()...)
|
||||
}
|
||||
}
|
||||
|
||||
for _, i := range deleteList {
|
||||
p.scenes = append(p.scenes[:i], p.scenes[i+1:]...)
|
||||
}
|
||||
|
||||
for _, device := range updatedList {
|
||||
if !p.started[device.BridgeID] {
|
||||
p.started[device.BridgeID] = true
|
||||
go p.runBridge(device.BridgeID)
|
||||
}
|
||||
|
||||
p.pending[device.BridgeID] = append(p.pending[device.BridgeID], device)
|
||||
if p.waiting[device.BridgeID] != nil {
|
||||
close(p.waiting[device.BridgeID])
|
||||
p.waiting[device.BridgeID] = nil
|
||||
}
|
||||
}
|
||||
p.mu.Unlock()
|
||||
}
|
||||
}
|
||||
|
||||
// reassignDevice re-evaluates the device's scene assignment config. It will return whether the scene changed, which
|
||||
// should trigger an update.
|
||||
func (p *Publisher) reassignDevice(device models.Device) bool {
|
||||
var selectedAssignment *models.DeviceSceneAssignment
|
||||
for _, assignment := range device.SceneAssignments {
|
||||
duration := time.Duration(assignment.DurationMS) * time.Millisecond
|
||||
if duration <= 0 || time.Now().Before(assignment.StartTime.Add(duration)) {
|
||||
selectedAssignment = &assignment
|
||||
}
|
||||
}
|
||||
|
||||
if selectedAssignment == nil {
|
||||
if p.sceneAssignment[device.ID] != nil {
|
||||
p.sceneAssignment[device.ID].RemoveDevice(device)
|
||||
delete(p.sceneAssignment, device.ID)
|
||||
|
||||
// Scene changed
|
||||
return true
|
||||
}
|
||||
|
||||
// Stop here, no scene should be assigned.
|
||||
return false
|
||||
} else {
|
||||
if p.sceneData[selectedAssignment.SceneID] == nil {
|
||||
// Freeze until scene becomes available.
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
if p.sceneAssignment[device.ID] != nil {
|
||||
scene := p.sceneAssignment[device.ID]
|
||||
if scene.data.ID == selectedAssignment.SceneID && scene.group == selectedAssignment.Group {
|
||||
// Current assignment is good.
|
||||
return false
|
||||
}
|
||||
|
||||
p.sceneAssignment[device.ID].RemoveDevice(device)
|
||||
delete(p.sceneAssignment, device.ID)
|
||||
}
|
||||
|
||||
for _, scene := range p.scenes {
|
||||
if scene.data.ID == selectedAssignment.SceneID && scene.group == selectedAssignment.Group {
|
||||
p.sceneAssignment[device.ID] = scene
|
||||
return true
|
||||
}
|
||||
}
|
||||
|
||||
newScene := &Scene{
|
||||
data: p.sceneData[selectedAssignment.SceneID],
|
||||
group: selectedAssignment.Group,
|
||||
startTime: selectedAssignment.StartTime,
|
||||
endTime: selectedAssignment.StartTime.Add(time.Duration(selectedAssignment.DurationMS) * time.Millisecond),
|
||||
roleMap: make(map[int]int, 16),
|
||||
roleList: make(map[int][]models.Device, 16),
|
||||
lastStates: make(map[int]models.DeviceState, 16),
|
||||
due: true,
|
||||
}
|
||||
p.sceneAssignment[device.ID] = newScene
|
||||
p.scenes = append(p.scenes, newScene)
|
||||
|
||||
if selectedAssignment.DurationMS <= 0 {
|
||||
newScene.endTime = time.Time{}
|
||||
}
|
||||
|
||||
return true
|
||||
}
|
||||
|
||||
func (p *Publisher) runBridge(id int) {
|
||||
defer func() {
|
||||
p.mu.Lock()
|
||||
p.started[id] = false
|
||||
p.mu.Unlock()
|
||||
}()
|
||||
|
||||
bridge, err := config.BridgeRepository().Find(context.Background(), id)
|
||||
if err != nil {
|
||||
log.Println("Failed to get bridge data:", err)
|
||||
return
|
||||
}
|
||||
|
||||
driver, err := config.DriverProvider().Provide(bridge.Driver)
|
||||
if err != nil {
|
||||
log.Println("Failed to get bridge driver:", err)
|
||||
log.Println("Maybe Lucifer needs to be updated.")
|
||||
return
|
||||
}
|
||||
|
||||
devices := make(map[int]models.Device)
|
||||
|
||||
for {
|
||||
p.mu.Lock()
|
||||
if len(p.pending[id]) == 0 {
|
||||
if p.waiting[id] == nil {
|
||||
p.waiting[id] = make(chan struct{})
|
||||
}
|
||||
waitCh := p.waiting[id]
|
||||
p.mu.Unlock()
|
||||
<-waitCh
|
||||
p.mu.Lock()
|
||||
}
|
||||
|
||||
updates := p.pending[id]
|
||||
p.pending[id] = p.pending[id][:0:0]
|
||||
p.mu.Unlock()
|
||||
|
||||
// Only allow the latest update per device (this avoids slow bridges causing a backlog of cations).
|
||||
for key := range devices {
|
||||
delete(devices, key)
|
||||
}
|
||||
for _, update := range updates {
|
||||
devices[update.ID] = update
|
||||
}
|
||||
updates = updates[:0]
|
||||
for _, value := range devices {
|
||||
updates = append(updates, value)
|
||||
}
|
||||
|
||||
err := driver.Publish(context.Background(), bridge, updates)
|
||||
if err != nil {
|
||||
log.Println("Failed to publish to driver:", err)
|
||||
|
||||
p.mu.Lock()
|
||||
p.pending[id] = append(updates, p.pending[id]...)
|
||||
p.mu.Unlock()
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
var publisher = Publisher{
|
||||
sceneData: make(map[int]*models.Scene),
|
||||
scenes: make([]*Scene, 0, 16),
|
||||
sceneAssignment: make(map[int]*Scene, 16),
|
||||
started: make(map[int]bool),
|
||||
pending: make(map[int][]models.Device),
|
||||
waiting: make(map[int]chan struct{}),
|
||||
}
|
||||
|
||||
func Global() *Publisher {
|
||||
return &publisher
|
||||
}
|
||||
|
||||
func Initialize(ctx context.Context) error {
|
||||
err := publisher.ReloadScenes(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
err = publisher.ReloadDevices(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
go publisher.Run()
|
||||
time.Sleep(time.Millisecond * 50)
|
||||
go publisher.PublishChannel(config.PublishChannel)
|
||||
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,186 @@
|
||||
package publisher
|
||||
|
||||
import (
|
||||
"git.aiterp.net/lucifer/new-server/models"
|
||||
"time"
|
||||
)
|
||||
|
||||
type Scene struct {
|
||||
data *models.Scene
|
||||
group string
|
||||
startTime time.Time
|
||||
endTime time.Time
|
||||
roleMap map[int]int
|
||||
roleList map[int][]models.Device
|
||||
lastStates map[int]models.DeviceState
|
||||
|
||||
due bool
|
||||
lastInterval int64
|
||||
}
|
||||
|
||||
// UpdateScene updates the scene data and re-seats all devices.
|
||||
func (s *Scene) UpdateScene(data models.Scene) {
|
||||
devices := make([]models.Device, 0, 16)
|
||||
|
||||
// Collect all devices into the undefined role (-1)
|
||||
for _, list := range s.roleList {
|
||||
for _, device := range list {
|
||||
devices = append(devices, device)
|
||||
s.roleMap[device.ID] = -1
|
||||
}
|
||||
}
|
||||
s.roleList = map[int][]models.Device{-1: devices}
|
||||
|
||||
// Update data and reset devices.
|
||||
s.data = &data
|
||||
for _, device := range append(devices[:0:0], devices...) {
|
||||
s.UpsertDevice(device)
|
||||
}
|
||||
}
|
||||
|
||||
// UpsertDevice moves the device if necessary and updates its state.
|
||||
func (s *Scene) UpsertDevice(device models.Device) {
|
||||
if s.data == nil {
|
||||
s.roleMap[device.ID] = -1
|
||||
s.roleList[-1] = append(s.roleList[-1], device)
|
||||
return
|
||||
}
|
||||
|
||||
oldIndex, hasOldIndex := s.roleMap[device.ID]
|
||||
newIndex := s.data.RoleIndex(&device)
|
||||
|
||||
s.roleMap[device.ID] = newIndex
|
||||
|
||||
if hasOldIndex {
|
||||
for i, device2 := range s.roleList[oldIndex] {
|
||||
if device2.ID == device.ID {
|
||||
s.roleList[oldIndex] = append(s.roleList[oldIndex][:i], s.roleList[oldIndex][i+1:]...)
|
||||
break
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
s.due = true
|
||||
|
||||
s.roleList[newIndex] = append(s.roleList[newIndex], device)
|
||||
if newIndex != -1 {
|
||||
s.data.Roles[newIndex].ApplyOrder(s.roleList[newIndex])
|
||||
}
|
||||
}
|
||||
|
||||
// RemoveDevice finds and remove a device. It's a noop if the device does not exist in this scene.
|
||||
func (s *Scene) RemoveDevice(device models.Device) {
|
||||
roleIndex, hasRoleIndex := s.roleMap[device.ID]
|
||||
if !hasRoleIndex {
|
||||
return
|
||||
}
|
||||
|
||||
for i, device2 := range s.roleList[roleIndex] {
|
||||
if device2.ID == device.ID {
|
||||
s.roleList[roleIndex] = append(s.roleList[roleIndex][:i], s.roleList[roleIndex][i+1:]...)
|
||||
break
|
||||
}
|
||||
}
|
||||
|
||||
s.due = true
|
||||
|
||||
delete(s.roleMap, device.ID)
|
||||
delete(s.lastStates, device.ID)
|
||||
}
|
||||
|
||||
func (s *Scene) Empty() bool {
|
||||
for _, list := range s.roleList {
|
||||
if len(list) > 0 {
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
return true
|
||||
}
|
||||
|
||||
func (s *Scene) Due() bool {
|
||||
if s.due {
|
||||
return true
|
||||
}
|
||||
|
||||
if s.data.IntervalMS > 0 {
|
||||
interval := time.Duration(s.data.IntervalMS) * time.Millisecond
|
||||
return int64(time.Since(s.startTime)/interval) != s.lastInterval
|
||||
}
|
||||
|
||||
return false
|
||||
}
|
||||
|
||||
func (s *Scene) UnaffectedDevices() []models.Device {
|
||||
return append(s.roleList[-1][:0:0], s.roleList[-1]...)
|
||||
}
|
||||
|
||||
func (s *Scene) AllDevices() []models.Device {
|
||||
res := make([]models.Device, 0, 16)
|
||||
for _, list := range s.roleList {
|
||||
res = append(res, list...)
|
||||
}
|
||||
|
||||
return res
|
||||
}
|
||||
|
||||
// Run runs the scene
|
||||
func (s *Scene) Run() []models.Device {
|
||||
if s.data == nil {
|
||||
return []models.Device{}
|
||||
}
|
||||
|
||||
intervalNumber := int64(0)
|
||||
intervalMax := int64(1)
|
||||
if s.data.IntervalMS > 0 {
|
||||
interval := time.Duration(s.data.IntervalMS) * time.Millisecond
|
||||
intervalNumber = int64(time.Since(s.startTime) / interval)
|
||||
|
||||
if !s.endTime.IsZero() {
|
||||
intervalMax = int64(s.endTime.Sub(s.startTime) / interval)
|
||||
} else {
|
||||
intervalMax = intervalNumber + 1
|
||||
}
|
||||
}
|
||||
|
||||
updatedDevices := make([]models.Device, 0, 16)
|
||||
for i, list := range s.roleList {
|
||||
if i == -1 {
|
||||
continue
|
||||
}
|
||||
|
||||
role := s.data.Roles[i]
|
||||
|
||||
for j, device := range list {
|
||||
newState := role.ApplyEffect(&device, models.SceneRunContext{
|
||||
Index: j,
|
||||
Length: len(list),
|
||||
IntervalNumber: intervalNumber,
|
||||
IntervalMax: intervalMax,
|
||||
})
|
||||
|
||||
err := device.SetState(newState)
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
|
||||
s.lastStates[device.ID] = device.State
|
||||
|
||||
updatedDevices = append(updatedDevices, device)
|
||||
}
|
||||
}
|
||||
|
||||
s.due = false
|
||||
s.lastInterval = intervalNumber
|
||||
|
||||
return updatedDevices
|
||||
}
|
||||
|
||||
func (s *Scene) LastState(id int) *models.DeviceState {
|
||||
lastState, ok := s.lastStates[id]
|
||||
if !ok {
|
||||
return nil
|
||||
}
|
||||
|
||||
return &lastState
|
||||
}
|
||||
Reference in New Issue
Block a user