small fixed, device repo and some device api
This commit is contained in:
+30
-4
@@ -3,11 +3,15 @@ package services
|
||||
import (
|
||||
"context"
|
||||
"git.aiterp.net/lucifer/new-server/app/config"
|
||||
"git.aiterp.net/lucifer/new-server/models"
|
||||
"log"
|
||||
"strconv"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
var cancelMap = make(map[int]context.CancelFunc, 8)
|
||||
var cancelMutex sync.Mutex
|
||||
|
||||
func ConnectToBridges() {
|
||||
go func() {
|
||||
@@ -29,7 +33,10 @@ func runConnectToBridges() error {
|
||||
}
|
||||
|
||||
for _, bridge := range bridges {
|
||||
if cancelMap[bridge.ID] != nil {
|
||||
cancelMutex.Lock()
|
||||
isRunning := cancelMap[bridge.ID] != nil
|
||||
cancelMutex.Unlock()
|
||||
if isRunning {
|
||||
continue
|
||||
}
|
||||
|
||||
@@ -39,11 +46,30 @@ func runConnectToBridges() error {
|
||||
}
|
||||
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
err = driver.Run(ctx, bridge, config.EventChannel)
|
||||
|
||||
cancelMutex.Lock()
|
||||
cancelMap[bridge.ID] = cancel
|
||||
cancelMutex.Unlock()
|
||||
|
||||
log.Printf("Connected to bridge \"%s\" (%d)", bridge.Name, bridge.ID)
|
||||
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)
|
||||
|
||||
cancelMutex.Lock()
|
||||
cancel()
|
||||
cancelMap[bridge.ID] = nil
|
||||
cancelMutex.Unlock()
|
||||
}(bridge, cancel)
|
||||
}
|
||||
|
||||
return nil
|
||||
|
||||
@@ -2,10 +2,12 @@ package services
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"git.aiterp.net/lucifer/new-server/app/config"
|
||||
"git.aiterp.net/lucifer/new-server/models"
|
||||
"log"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
@@ -18,7 +20,7 @@ func StartEventHandler() {
|
||||
}
|
||||
}()
|
||||
|
||||
// Dispatch an HourChanged event at every hour
|
||||
// Generate TimeChanged event
|
||||
go func() {
|
||||
drift := time.Now().Add(time.Minute).Truncate(time.Minute).Sub(time.Now())
|
||||
time.Sleep(drift + time.Millisecond * 5)
|
||||
@@ -48,7 +50,12 @@ func handleEvent(event models.Event) {
|
||||
}
|
||||
|
||||
if !X {
|
||||
log.Println("Unhandled event: " + event.Name)
|
||||
paramStrings := make([]string, 0, 8)
|
||||
for key, value := range event.Payload {
|
||||
paramStrings = append(paramStrings, fmt.Sprintf("%s=%s", key, value))
|
||||
}
|
||||
|
||||
log.Printf("Unhandled event %s(%s)", event.Name, strings.Join(paramStrings, ", "))
|
||||
return
|
||||
}
|
||||
|
||||
@@ -68,7 +75,6 @@ func handleEvent(event models.Event) {
|
||||
if !handler.MatchesEvent(event, devices) {
|
||||
continue
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
+44
-16
@@ -5,6 +5,8 @@ import (
|
||||
"git.aiterp.net/lucifer/new-server/app/config"
|
||||
"git.aiterp.net/lucifer/new-server/models"
|
||||
"log"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
func StartPublisher() {
|
||||
@@ -16,32 +18,58 @@ func StartPublisher() {
|
||||
continue
|
||||
}
|
||||
|
||||
// Emergency solution! Please avoid!
|
||||
// Send devices not belonging to the first channel separately
|
||||
bridgeID := devices[0].BridgeID
|
||||
lists := make(map[int][]models.Device, 4)
|
||||
for _, device := range devices {
|
||||
if device.BridgeID != bridgeID {
|
||||
config.PublishChannel<-[]models.Device{device}
|
||||
}
|
||||
lists[device.BridgeID] = append(lists[device.BridgeID], device)
|
||||
}
|
||||
|
||||
bridge, err := config.BridgeRepository().Find(ctx, devices[0].BridgeID)
|
||||
ctx, cancel := context.WithTimeout(ctx, time.Second * 30)
|
||||
|
||||
bridges, err := config.BridgeRepository().FetchAll(ctx)
|
||||
if err != nil {
|
||||
log.Println("Publishing error (1): " + err.Error())
|
||||
continue
|
||||
}
|
||||
|
||||
driver, err := config.DriverProvider().Provide(bridge.Driver)
|
||||
if err != nil {
|
||||
log.Println("Publishing error (2): " + err.Error())
|
||||
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)
|
||||
}
|
||||
|
||||
err = driver.Publish(ctx, bridge, devices)
|
||||
if err != nil {
|
||||
log.Println("Publishing error (3): " + err.Error())
|
||||
continue
|
||||
}
|
||||
wg.Wait()
|
||||
cancel()
|
||||
}
|
||||
}()
|
||||
}
|
||||
|
||||
@@ -0,0 +1,95 @@
|
||||
package services
|
||||
|
||||
import (
|
||||
"context"
|
||||
"git.aiterp.net/lucifer/new-server/app/config"
|
||||
"git.aiterp.net/lucifer/new-server/models"
|
||||
"log"
|
||||
"strconv"
|
||||
"time"
|
||||
)
|
||||
|
||||
func CheckNewDevices() {
|
||||
go func() {
|
||||
// Wait a bit before the first to let bridges connect.
|
||||
time.Sleep(time.Second * 5)
|
||||
err := checkNewDevices()
|
||||
if err != nil {
|
||||
log.Println("Failed to sync lights:", err)
|
||||
}
|
||||
|
||||
for range time.NewTicker(time.Second * 30).C {
|
||||
err := checkNewDevices()
|
||||
if err != nil {
|
||||
log.Println("Failed to sync lights:", err)
|
||||
}
|
||||
}
|
||||
}()
|
||||
}
|
||||
|
||||
func checkNewDevices() error {
|
||||
ctx, cancel := context.WithTimeout(context.Background(), time.Second*27)
|
||||
defer cancel()
|
||||
|
||||
bridges, err := config.BridgeRepository().FetchAll(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
for _, bridge := range bridges {
|
||||
driver, err := config.DriverProvider().Provide(bridge.Driver)
|
||||
if err != nil {
|
||||
log.Println("Unknown/unsupported driver:", bridge.Driver)
|
||||
continue
|
||||
}
|
||||
|
||||
savedDevices, err := config.DeviceRepository().FetchByReference(ctx, models.RKBridgeID, strconv.Itoa(bridge.ID))
|
||||
if err != nil {
|
||||
log.Println("Failed to list devices from db:", err)
|
||||
continue
|
||||
}
|
||||
|
||||
driverDevices, err := driver.ListDevices(ctx, bridge)
|
||||
if err != nil {
|
||||
log.Println("Failed to list devices from driver:", err)
|
||||
continue
|
||||
}
|
||||
|
||||
foundNewDevices := false
|
||||
SaveLoop:
|
||||
for _, driverDevice := range driverDevices {
|
||||
for _, savedDevice := range savedDevices {
|
||||
if savedDevice.InternalID == driverDevice.InternalID {
|
||||
continue SaveLoop
|
||||
}
|
||||
}
|
||||
|
||||
log.Println("Saving new device", driverDevice.InternalID)
|
||||
|
||||
err := config.DeviceRepository().Save(ctx, &driverDevice)
|
||||
if err != nil {
|
||||
log.Println("Failed to save device:", err)
|
||||
continue
|
||||
}
|
||||
|
||||
foundNewDevices = true
|
||||
}
|
||||
|
||||
// If new devices were found, publish them so that the driver can be set up.
|
||||
if foundNewDevices {
|
||||
savedDevices, err := config.DeviceRepository().FetchByReference(ctx, models.RKBridgeID, strconv.Itoa(bridge.ID))
|
||||
if err != nil {
|
||||
log.Println("Failed to fetch devices from db second time:", err)
|
||||
continue
|
||||
}
|
||||
|
||||
err = driver.Publish(ctx, bridge, savedDevices)
|
||||
if err != nil {
|
||||
log.Println("Failed to list devices from db:", err)
|
||||
continue
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
Reference in New Issue
Block a user