Skip to content

Commit

Permalink
fix: process loop and locks (#31)
Browse files Browse the repository at this point in the history
* fix: remove full status Update Requests

Signed-off-by: Martin Buchleitner <mbuchleitner@infralovers.com>

* fix: remove full status Update Requests

Signed-off-by: Martin Buchleitner <mbuchleitner@infralovers.com>

* merge update

Signed-off-by: Martin Buchleitner <mbuchleitner@infralovers.com>

* fix: update that loop ctrl

Signed-off-by: Martin Buchleitner <mbuchleitner@infralovers.com>

* fix: atomic increment

Signed-off-by: Martin Buchleitner <mbuchleitner@infralovers.com>

* refactor: split helper code to dedicated files

Signed-off-by: Martin Buchleitner <mbuchleitner@infralovers.com>

---------

Signed-off-by: Martin Buchleitner <mbuchleitner@infralovers.com>
  • Loading branch information
mabunixda authored Feb 9, 2024
1 parent c0ba28f commit 1c59e45
Show file tree
Hide file tree
Showing 6 changed files with 196 additions and 161 deletions.
8 changes: 4 additions & 4 deletions go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -3,16 +3,16 @@ module github.com/mabunixda/wattpilot
go 1.21

require (
github.com/gobwas/ws v1.2.0
github.com/sirupsen/logrus v1.9.0
golang.org/x/crypto v0.17.0
github.com/gobwas/ws v1.3.2
github.com/sirupsen/logrus v1.9.3
golang.org/x/crypto v0.18.0
gopkg.in/yaml.v2 v2.4.0
)

require (
github.com/gobwas/httphead v0.1.0 // indirect
github.com/gobwas/pool v0.2.1 // indirect
golang.org/x/sys v0.15.0 // indirect
golang.org/x/sys v0.16.0 // indirect
)

retract v1.6.3
16 changes: 8 additions & 8 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -5,21 +5,21 @@ github.com/gobwas/httphead v0.1.0 h1:exrUm0f4YX0L7EBwZHuCF4GDp8aJfVeBrlLQrs6NqWU
github.com/gobwas/httphead v0.1.0/go.mod h1:O/RXo79gxV8G+RqlR/otEwx4Q36zl9rqC5u12GKvMCM=
github.com/gobwas/pool v0.2.1 h1:xfeeEhW7pwmX8nuLVlqbzVc7udMDrwetjEv+TZIz1og=
github.com/gobwas/pool v0.2.1/go.mod h1:q8bcK0KcYlCgd9e7WYLm9LpyS+YeLd8JVDW6WezmKEw=
github.com/gobwas/ws v1.2.0 h1:u0p9s3xLYpZCA1z5JgCkMeB34CKCMMQbM+G8Ii7YD0I=
github.com/gobwas/ws v1.2.0/go.mod h1:hRKAFb8wOxFROYNsT1bqfWnhX+b5MFeJM9r2ZSwg/KY=
github.com/gobwas/ws v1.3.2 h1:zlnbNHxumkRvfPWgfXu8RBwyNR1x8wh9cf5PTOCqs9Q=
github.com/gobwas/ws v1.3.2/go.mod h1:hRKAFb8wOxFROYNsT1bqfWnhX+b5MFeJM9r2ZSwg/KY=
github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM=
github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
github.com/sirupsen/logrus v1.9.0 h1:trlNQbNUG3OdDrDil03MCb1H2o9nJ1x4/5LYw7byDE0=
github.com/sirupsen/logrus v1.9.0/go.mod h1:naHLuLoDiP4jHNo9R0sCBMtWGeIprob74mVsIT4qYEQ=
github.com/sirupsen/logrus v1.9.3 h1:dueUQJ1C2q9oE3F7wvmSGAaVtTmUizReu6fjN8uqzbQ=
github.com/sirupsen/logrus v1.9.3/go.mod h1:naHLuLoDiP4jHNo9R0sCBMtWGeIprob74mVsIT4qYEQ=
github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=
github.com/stretchr/testify v1.7.0 h1:nwc3DEeHmmLAfoZucVR881uASk0Mfjw8xYJ99tb5CcY=
github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg=
golang.org/x/crypto v0.17.0 h1:r8bRNjWL3GshPW3gkd+RpvzWrZAwPS49OmTGZ/uhM4k=
golang.org/x/crypto v0.17.0/go.mod h1:gCAAfMLgwOJRpTjQ2zCCt2OcSfYMTeZVSRtQlPC7Nq4=
golang.org/x/crypto v0.18.0 h1:PGVlW0xEltQnzFZ55hkuX5+KLyrMYhHld1YHO4AKcdc=
golang.org/x/crypto v0.18.0/go.mod h1:R0j02AL6hcrfOiy9T4ZYp/rcWeMxM3L6QYxlOuEG1mg=
golang.org/x/sys v0.0.0-20220715151400-c0bba94af5f8/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.15.0 h1:h48lPFYpsTvQJZF4EKyI4aLHaev3CxivZmv7yZig9pc=
golang.org/x/sys v0.15.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA=
golang.org/x/sys v0.16.0 h1:xWw16ngr6ZMtmxDyKyIgsE93KNKz5HKmMa3b8ALHidU=
golang.org/x/sys v0.16.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA=
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405 h1:yhCVgyC4o1eVCa2tZl7eS0r+SDo693bJlVdllGtEeKM=
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
gopkg.in/yaml.v2 v2.4.0 h1:D8xgwECY7CYvx+Y2n4sBz93Jn9JRvxdiyyo8CTfuKaY=
Expand Down
38 changes: 38 additions & 0 deletions helper.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,38 @@
package wattpilot

import (
"crypto/sha256"
"encoding/hex"
"fmt"
"math/rand"
"time"
)

var randomSource = rand.New(rand.NewSource(time.Now().UnixNano()))

func Keys[K comparable, V any](m map[K]V) []K {
keys := make([]K, 0, len(m))
for k := range m {
keys = append(keys, k)
}
return keys
}

func hasKey(data map[string]interface{}, key string) bool {
_, isKnown := data[key]
return isKnown
}

func sha256sum(data string) string {
bs := sha256.Sum256([]byte(data))
return fmt.Sprintf("%x", bs)
}
func randomHexString(n int) string {
b := make([]byte, (n+2)/2) // can be simplified to n/2 if n is always even

if _, err := randomSource.Read(b); err != nil {
panic(err)
}

return hex.EncodeToString(b)[1 : n+1]
}
58 changes: 58 additions & 0 deletions pubsub.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,58 @@
package wattpilot

import "sync"

type Pubsub struct {
mu sync.RWMutex
subs map[string][]chan interface{}
closed bool
}

func NewPubsub() *Pubsub {
ps := &Pubsub{}
ps.subs = make(map[string][]chan interface{})
return ps
}

func (ps *Pubsub) IsEmpty() bool {
ps.mu.Lock()
defer ps.mu.Unlock()

return len(ps.subs) == 0
}

func (ps *Pubsub) Subscribe(topic string) <-chan interface{} {
ps.mu.Lock()
defer ps.mu.Unlock()

ch := make(chan interface{}, 1)
ps.subs[topic] = append(ps.subs[topic], ch)
return ch
}

func (ps *Pubsub) Publish(topic string, msg interface{}) {
ps.mu.RLock()
defer ps.mu.RUnlock()

if ps.closed {
return
}
for _, ch := range ps.subs[topic] {
go func(ch chan interface{}) {
ch <- msg
}(ch)
}
}
func (ps *Pubsub) Close() {
ps.mu.Lock()
defer ps.mu.Unlock()

if !ps.closed {
ps.closed = true
for _, subs := range ps.subs {
for _, ch := range subs {
close(ch)
}
}
}
}
11 changes: 6 additions & 5 deletions shell/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@ var inputs = map[string]InputFunc{
"disconnect": inDisconnect,
"properties": inProperties,
"dump": dumpData,
"level": setLevel,
"log": setLevel,
"update": inUpdateStatus,
}

Expand Down Expand Up @@ -134,6 +134,7 @@ func inConnect(w *api.Wattpilot, data []string) {
err := w.Connect()
if err != nil {
log.Println("Could not connect", err)
return
}
log.Printf("Connected to WattPilot %s, Serial %s", w.GetName(), w.GetSerial())
}
Expand All @@ -142,9 +143,9 @@ func inDisconnect(w *api.Wattpilot, data []string) {
w.Disconnect()
}

func processUpdates(ups <-chan interface{}) {
updates = ups
}
// func processUpdates(ups <-chan interface{}) {
// updates = ups
// }

var interrupt chan os.Signal

Expand All @@ -165,7 +166,7 @@ func main() {
fmt.Println("Could not update loglevel to", level, err)
}
// just a sample to test notification updates
processUpdates(w.GetNotifications("fhz"))
// processUpdates(w.GetNotifications("fhz"))
inConnect(w, nil)

w.StatusInfo()
Expand Down
Loading

0 comments on commit 1c59e45

Please sign in to comment.