add beginning of CLI app and some backend fixes.
This commit is contained in:
+44
-3
@@ -6,6 +6,7 @@ import (
|
||||
"git.aiterp.net/lucifer/new-server/models"
|
||||
"github.com/gin-gonic/gin"
|
||||
"log"
|
||||
"strconv"
|
||||
"strings"
|
||||
)
|
||||
|
||||
@@ -14,9 +15,18 @@ func fetchDevices(ctx context.Context, fetchStr string) ([]models.Device, error)
|
||||
return config.DeviceRepository().FetchByReference(ctx, models.RKTag, fetchStr[4:])
|
||||
} else if strings.HasPrefix(fetchStr, "bridge:") {
|
||||
return config.DeviceRepository().FetchByReference(ctx, models.RKBridgeID, fetchStr[7:])
|
||||
} else if fetchStr == "all" {
|
||||
} else if strings.HasPrefix(fetchStr, "id:") {
|
||||
return config.DeviceRepository().FetchByReference(ctx, models.RKDeviceID, fetchStr[7:])
|
||||
} else if strings.HasPrefix(fetchStr, "name:") {
|
||||
return config.DeviceRepository().FetchByReference(ctx, models.RKName, fetchStr[7:])
|
||||
}else if fetchStr == "all" {
|
||||
return config.DeviceRepository().FetchByReference(ctx, models.RKAll, "")
|
||||
} else {
|
||||
_, err := strconv.Atoi(fetchStr)
|
||||
if err != nil {
|
||||
return config.DeviceRepository().FetchByReference(ctx, models.RKName, fetchStr)
|
||||
}
|
||||
|
||||
return config.DeviceRepository().FetchByReference(ctx, models.RKDeviceID, fetchStr)
|
||||
}
|
||||
}
|
||||
@@ -30,9 +40,9 @@ func Devices(r gin.IRoutes) {
|
||||
return fetchDevices(ctxOf(c), c.Param("fetch"))
|
||||
}))
|
||||
|
||||
r.PUT("/batch", handler(func(c *gin.Context) (interface{}, error) {
|
||||
r.PUT("", handler(func(c *gin.Context) (interface{}, error) {
|
||||
var body []struct {
|
||||
Fetch string `json:"fetch"`
|
||||
Fetch string `json:"fetch"`
|
||||
SetState models.NewDeviceState `json:"setState"`
|
||||
}
|
||||
err := parseBody(c, &body)
|
||||
@@ -81,6 +91,34 @@ func Devices(r gin.IRoutes) {
|
||||
return changed, nil
|
||||
}))
|
||||
|
||||
r.PUT("/:fetch", handler(func(c *gin.Context) (interface{}, error) {
|
||||
update := models.DeviceUpdate{}
|
||||
err := parseBody(c, &update)
|
||||
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
|
||||
}
|
||||
|
||||
for i := range devices {
|
||||
devices[i].ApplyUpdate(update)
|
||||
|
||||
err := config.DeviceRepository().Save(context.Background(), &devices[i])
|
||||
if err != nil {
|
||||
log.Println("Failed to save device for state:", err)
|
||||
continue
|
||||
}
|
||||
}
|
||||
|
||||
return devices, nil
|
||||
}))
|
||||
|
||||
r.PUT("/:fetch/state", handler(func(c *gin.Context) (interface{}, error) {
|
||||
state := models.NewDeviceState{}
|
||||
err := parseBody(c, &state)
|
||||
@@ -159,6 +197,9 @@ func Devices(r gin.IRoutes) {
|
||||
index = i
|
||||
}
|
||||
}
|
||||
if index == -1 {
|
||||
continue
|
||||
}
|
||||
|
||||
device.Tags = append(device.Tags[:index], device.Tags[index+1:]...)
|
||||
}
|
||||
|
||||
@@ -0,0 +1,136 @@
|
||||
package client
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"git.aiterp.net/lucifer/new-server/models"
|
||||
"io"
|
||||
"net"
|
||||
"net/http"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
type Client struct {
|
||||
APIRoot string
|
||||
}
|
||||
|
||||
func (client *Client) GetDevices(ctx context.Context, fetchStr string) ([]models.Device, error) {
|
||||
devices := make([]models.Device, 0, 16)
|
||||
err := client.Fetch(ctx, "GET", "/api/devices/"+fetchStr, &devices, nil)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return devices, nil
|
||||
}
|
||||
|
||||
func (client *Client) PutDevice(ctx context.Context, fetchStr string, update models.DeviceUpdate) ([]models.Device, error) {
|
||||
devices := make([]models.Device, 0, 16)
|
||||
err := client.Fetch(ctx, "PUT", "/api/devices/"+fetchStr, &devices, update)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return devices, nil
|
||||
}
|
||||
|
||||
func (client *Client) PutDeviceState(ctx context.Context, fetchStr string, update models.NewDeviceState) ([]models.Device, error) {
|
||||
devices := make([]models.Device, 0, 16)
|
||||
err := client.Fetch(ctx, "PUT", "/api/devices/"+fetchStr+"/state", &devices, update)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return devices, nil
|
||||
}
|
||||
|
||||
func (client *Client) PutDeviceTags(ctx context.Context, fetchStr string, addTags []string, removeTags []string) ([]models.Device, error) {
|
||||
devices := make([]models.Device, 0, 16)
|
||||
err := client.Fetch(ctx, "PUT", "/api/devices/"+fetchStr+"/tags", &devices, map[string][]string{
|
||||
"add": addTags,
|
||||
"remove": removeTags,
|
||||
})
|
||||
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 {
|
||||
return err
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (client *Client) Fetch(ctx context.Context, method string, path string, dst interface{}, body interface{}) error {
|
||||
var reqBody io.ReadWriter
|
||||
if body != nil && method != "GET" {
|
||||
reqBody = bytes.NewBuffer(make([]byte, 0, 512))
|
||||
|
||||
err := json.NewEncoder(reqBody).Encode(body)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
req, err := http.NewRequest(method, client.APIRoot+path, reqBody)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
res, err := httpClient.Do(req.WithContext(ctx))
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer res.Body.Close()
|
||||
|
||||
if !strings.HasPrefix(res.Header.Get("Content-Type"), "application/json") {
|
||||
return fmt.Errorf("%s: %s", path, res.Status)
|
||||
}
|
||||
|
||||
var resJson struct {
|
||||
Code int `json:"code"`
|
||||
Message *string `json:"message"`
|
||||
Data json.RawMessage `json:"data"`
|
||||
}
|
||||
err = json.NewDecoder(res.Body).Decode(&resJson)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if resJson.Code != 200 {
|
||||
msg := ""
|
||||
if resJson.Message != nil {
|
||||
msg = *resJson.Message
|
||||
}
|
||||
|
||||
return fmt.Errorf("%d: %s", resJson.Code, msg)
|
||||
}
|
||||
|
||||
if dst == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
return json.Unmarshal(resJson.Data, dst)
|
||||
}
|
||||
|
||||
var httpClient = &http.Client{
|
||||
Transport: &http.Transport{
|
||||
Proxy: http.ProxyFromEnvironment,
|
||||
DialContext: (&net.Dialer{
|
||||
Timeout: 30 * time.Second,
|
||||
KeepAlive: 30 * time.Second,
|
||||
}).DialContext,
|
||||
MaxIdleConns: 16,
|
||||
MaxIdleConnsPerHost: 16,
|
||||
IdleConnTimeout: time.Minute,
|
||||
},
|
||||
Timeout: time.Minute,
|
||||
}
|
||||
@@ -46,6 +46,8 @@ var loc, _ = time.LoadLocation("Europe/Oslo")
|
||||
var ctx = context.Background()
|
||||
|
||||
func handleEvent(event models.Event) (responses []models.Event) {
|
||||
var deadHandlers []models.EventHandler
|
||||
|
||||
startTime := time.Now()
|
||||
defer func() {
|
||||
duration := time.Since(startTime)
|
||||
@@ -102,6 +104,10 @@ func handleEvent(event models.Event) (responses []models.Event) {
|
||||
continue
|
||||
}
|
||||
|
||||
if handler.OneShot {
|
||||
deadHandlers = append(deadHandlers, handler)
|
||||
}
|
||||
|
||||
if handler.Priority > highestPriority {
|
||||
highestPriority = handler.Priority
|
||||
prioritizedEvent = handler.Actions.FireEvent
|
||||
@@ -159,6 +165,18 @@ func handleEvent(event models.Event) (responses []models.Event) {
|
||||
wg.Done()
|
||||
}(device)
|
||||
}
|
||||
for _, handler := range deadHandlers {
|
||||
wg.Add(1)
|
||||
|
||||
go func(handler models.EventHandler) {
|
||||
err := config.EventHandlerRepository().Delete(context.Background(), &handler)
|
||||
if err != nil {
|
||||
log.Println("Failed to delete spent one-shot event handler:", err)
|
||||
}
|
||||
|
||||
wg.Done()
|
||||
}(handler)
|
||||
}
|
||||
|
||||
wg.Wait()
|
||||
|
||||
|
||||
Reference in New Issue
Block a user