Skip to content

Commit d521d9a

Browse files
authored
Merge branch 'develop' into control-station/gauge
2 parents 8dada58 + b83285e commit d521d9a

14 files changed

Lines changed: 414 additions & 59 deletions

File tree

.gitignore

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,11 +1,12 @@
11
build
2-
.vscode
32

43
# GOOGLE API KEY
54
secret.json
65

76
# MacOS Files
87
.DS_Store
98

10-
# JetBrains IDE
11-
.idea/
9+
# Code Editor
10+
.idea/
11+
.vscode/
12+
.prettierrc

backend/cmd/main.go

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -156,11 +156,14 @@ func main() {
156156
boardIdToBoard[abstraction.BoardId(id)] = name
157157
}
158158
messageTopic := message_topic.NewUpdateTopic(boardIdToBoard)
159+
stateOrderTopic := order_topic.NewState(idToBoard, trace.Logger)
159160

160161
broker.AddTopic(data_topic.UpdateName, dataTopic)
161162
broker.AddTopic(connection_topic.UpdateName, connectionTopic)
162163
broker.AddTopic(order_topic.SendName, orderTopic)
164+
broker.AddTopic(order_topic.StateName, stateOrderTopic)
163165
broker.AddTopic(logger_topic.EnableName, loggerTopic)
166+
broker.AddTopic(logger_topic.ResponseName, loggerTopic)
164167
broker.AddTopic(message_topic.UpdateName, messageTopic)
165168

166169
connections := make(chan *websocket.Client)

backend/pkg/boards/events.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,7 @@ type UploadEvent struct {
3030
BoardEvent abstraction.BoardEvent
3131
Board string
3232
Data []byte
33+
Length int
3334
}
3435

3536
func (upload UploadEvent) Topic() abstraction.BrokerTopic {

backend/pkg/broker/topics/logger/enable.go

Lines changed: 44 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -3,24 +3,31 @@ package logger
33
import (
44
"encoding/json"
55
"fmt"
6+
"sync"
67
"sync/atomic"
78

89
"github.com/HyperloopUPV-H8/h9-backend/pkg/abstraction"
910
"github.com/HyperloopUPV-H8/h9-backend/pkg/websocket"
11+
"github.com/google/uuid"
12+
ws "github.com/gorilla/websocket"
1013
)
1114

1215
const EnableName abstraction.BrokerTopic = "logger/enable"
1316
const ResponseName abstraction.BrokerTopic = "logger/response"
1417

1518
type Enable struct {
16-
isRunning *atomic.Bool
17-
pool *websocket.Pool
18-
api abstraction.BrokerAPI
19+
isRunning *atomic.Bool
20+
pool *websocket.Pool
21+
connectionMx *sync.Mutex
22+
subscribers map[websocket.ClientId]struct{}
23+
api abstraction.BrokerAPI
1924
}
2025

2126
func NewEnableTopic() *Enable {
2227
enable := &Enable{
23-
isRunning: &atomic.Bool{},
28+
isRunning: &atomic.Bool{},
29+
connectionMx: new(sync.Mutex),
30+
subscribers: make(map[websocket.ClientId]struct{}),
2431
}
2532
enable.isRunning.Store(false)
2633
return enable
@@ -45,10 +52,23 @@ func (enable *Enable) ClientMessage(id websocket.ClientId, message *websocket.Me
4552
if err != nil {
4653
fmt.Printf("error handling logger: %v\n", err)
4754
}
55+
case ResponseName:
56+
enable.connectionMx.Lock()
57+
defer enable.connectionMx.Unlock()
58+
59+
fmt.Printf("logger/response subscribed %s\n", uuid.UUID(id).String())
60+
enable.subscribers[id] = struct{}{}
61+
default:
62+
enable.connectionMx.Lock()
63+
defer enable.connectionMx.Unlock()
64+
65+
enable.pool.Disconnect(id, ws.CloseUnsupportedData, "unsupported topic")
66+
delete(enable.subscribers, id)
67+
fmt.Printf("logger/response unsubscribed %s\n", uuid.UUID(id).String())
4868
}
4969
}
5070

51-
func (enable *Enable) handleToggle(id websocket.ClientId, message *websocket.Message) error {
71+
func (enable *Enable) handleToggle(_ websocket.ClientId, message *websocket.Message) error {
5272
var request bool
5373
err := json.Unmarshal(message.Payload, &request)
5474
if err != nil {
@@ -71,10 +91,27 @@ func (enable *Enable) broadcastState() error {
7191
return err
7292
}
7393

74-
enable.pool.Broadcast(websocket.Message{
94+
message := websocket.Message{
7595
Topic: ResponseName,
7696
Payload: payload,
77-
})
97+
}
98+
99+
enable.connectionMx.Lock()
100+
defer enable.connectionMx.Unlock()
101+
flaged := make([]websocket.ClientId, 0, len(enable.subscribers))
102+
for id := range enable.subscribers {
103+
err := enable.pool.Write(id, message)
104+
if err != nil {
105+
flaged = append(flaged, id)
106+
}
107+
}
108+
109+
for _, id := range flaged {
110+
enable.pool.Disconnect(id, ws.CloseInternalServerErr, "client disconnected")
111+
delete(enable.subscribers, id)
112+
fmt.Printf("logger/response unsubscribed %s\n", uuid.UUID(id).String())
113+
}
114+
78115
return nil
79116
}
80117

backend/pkg/broker/topics/order/send.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -42,7 +42,7 @@ func (send *Send) ClientMessage(id websocket.ClientId, message *websocket.Messag
4242
}
4343
}
4444

45-
func (send *Send) handleOrder(id websocket.ClientId, message *websocket.Message) error {
45+
func (send *Send) handleOrder(_ websocket.ClientId, message *websocket.Message) error {
4646
var order Order
4747
err := json.Unmarshal(message.Payload, &order)
4848
if err != nil {
Lines changed: 179 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,33 +1,210 @@
11
package order
22

33
import (
4+
"encoding/json"
5+
"fmt"
6+
"sync"
7+
48
"github.com/HyperloopUPV-H8/h9-backend/pkg/abstraction"
9+
"github.com/HyperloopUPV-H8/h9-backend/pkg/broker/topics"
10+
"github.com/HyperloopUPV-H8/h9-backend/pkg/transport/packet/order"
511
"github.com/HyperloopUPV-H8/h9-backend/pkg/websocket"
12+
"github.com/google/uuid"
13+
ws "github.com/gorilla/websocket"
14+
"github.com/rs/zerolog"
615
)
716

817
const StateName abstraction.BrokerTopic = "order/stateOrders"
918

1019
type State struct {
11-
pool *websocket.Pool
12-
api abstraction.BrokerAPI
20+
enabledOrders map[string]map[abstraction.PacketId]struct{}
21+
idToBoard map[uint16]string
22+
connectionMx *sync.Mutex
23+
subscribers map[websocket.ClientId]struct{}
24+
pool *websocket.Pool
25+
api abstraction.BrokerAPI
26+
logger zerolog.Logger
27+
}
28+
29+
func NewState(idToBoard map[uint16]string, baseLogger zerolog.Logger) *State {
30+
enabled := make(map[string]map[abstraction.PacketId]struct{})
31+
for _, board := range idToBoard {
32+
enabled[board] = make(map[abstraction.PacketId]struct{})
33+
}
34+
35+
return &State{
36+
enabledOrders: enabled,
37+
idToBoard: idToBoard,
38+
connectionMx: new(sync.Mutex),
39+
subscribers: make(map[websocket.ClientId]struct{}),
40+
logger: baseLogger,
41+
}
1342
}
1443

1544
func (state *State) Topic() abstraction.BrokerTopic {
1645
return StateName
1746
}
1847

1948
func (state *State) Push(push abstraction.BrokerPush) error {
49+
switch action := push.(type) {
50+
case *StateAdd:
51+
state.logger.Info().Msg("add")
52+
return state.addOrders(action.update)
53+
case *StateRemove:
54+
state.logger.Info().Msg("remove")
55+
return state.removeOrders(action.update)
56+
case *StateClear:
57+
state.logger.Info().Msg("clear board")
58+
return state.clearBoard(action.board)
59+
default:
60+
err := topics.ErrUnexpectedPush{Push: push}
61+
state.logger.Warn().Stack().Err(err).Msg("unexpected topic")
62+
return err
63+
}
64+
}
65+
66+
func (state *State) addOrders(add *order.Add) error {
67+
for _, addedOrder := range add.Orders() {
68+
board, ok := state.idToBoard[uint16(addedOrder)]
69+
if !ok {
70+
state.logger.Warn().Uint16("id", uint16(addedOrder)).Msg("unrecognized topic")
71+
continue
72+
}
73+
74+
state.enabledOrders[board][addedOrder] = struct{}{}
75+
}
76+
77+
return state.updateOrders()
78+
}
79+
80+
func (state *State) removeOrders(remove *order.Remove) error {
81+
for _, removedOrder := range remove.Orders() {
82+
board, ok := state.idToBoard[uint16(removedOrder)]
83+
if !ok {
84+
state.logger.Warn().Uint16("id", uint16(removedOrder)).Msg("unrecognized topic")
85+
continue
86+
}
87+
88+
delete(state.enabledOrders[board], removedOrder)
89+
}
90+
91+
return state.updateOrders()
92+
}
93+
94+
func (state *State) clearBoard(board string) error {
95+
if _, ok := state.enabledOrders[board]; !ok {
96+
state.logger.Warn().Str("board", board).Msg("unknown board")
97+
return nil
98+
}
99+
100+
state.enabledOrders[board] = make(map[abstraction.PacketId]struct{}, len(state.enabledOrders[board]))
101+
102+
return state.updateOrders()
103+
}
104+
105+
func (state *State) updateOrders() error {
106+
orderList := make(map[string][]abstraction.PacketId, len(state.enabledOrders))
107+
for board, enabledOrders := range state.enabledOrders {
108+
orderList[board] = make([]abstraction.PacketId, 0, len(enabledOrders))
109+
for order := range enabledOrders {
110+
orderList[board] = append(orderList[board], order)
111+
}
112+
}
113+
114+
payload, err := json.Marshal(orderList)
115+
if err != nil {
116+
return err
117+
}
118+
119+
message := websocket.Message{
120+
Topic: StateName,
121+
Payload: payload,
122+
}
123+
124+
state.connectionMx.Lock()
125+
defer state.connectionMx.Unlock()
126+
flaged := make([]websocket.ClientId, 0, len(state.subscribers))
127+
for id := range state.subscribers {
128+
err := state.pool.Write(id, message)
129+
if err != nil {
130+
flaged = append(flaged, id)
131+
}
132+
}
133+
134+
for _, id := range flaged {
135+
state.pool.Disconnect(id, ws.CloseInternalServerErr, "client disconnected")
136+
delete(state.subscribers, id)
137+
fmt.Printf("order/stateOrders unsubscribed %s\n", uuid.UUID(id).String())
138+
}
139+
20140
return nil
21141
}
22142

23143
func (state *State) Pull(request abstraction.BrokerRequest) (abstraction.BrokerResponse, error) {
24144
return nil, nil
25145
}
26146

147+
func (state *State) ClientMessage(id websocket.ClientId, message *websocket.Message) {
148+
state.connectionMx.Lock()
149+
defer state.connectionMx.Unlock()
150+
151+
switch message.Topic {
152+
case StateName:
153+
fmt.Printf("order/stateOrders subscribed %s\n", uuid.UUID(id).String())
154+
state.subscribers[id] = struct{}{}
155+
default:
156+
state.pool.Disconnect(id, ws.CloseUnsupportedData, "unsupported topic")
157+
delete(state.subscribers, id)
158+
fmt.Printf("order/stateOrders unsubscribed %s\n", uuid.UUID(id).String())
159+
}
160+
}
161+
27162
func (state *State) SetPool(pool *websocket.Pool) {
28163
state.pool = pool
29164
}
30165

31166
func (state *State) SetAPI(api abstraction.BrokerAPI) {
32167
state.api = api
33168
}
169+
170+
type StateAdd struct {
171+
update *order.Add
172+
}
173+
174+
func NewAdd(diff *order.Add) *StateAdd {
175+
return &StateAdd{
176+
update: diff,
177+
}
178+
}
179+
180+
func (add *StateAdd) Topic() abstraction.BrokerTopic {
181+
return StateName
182+
}
183+
184+
type StateRemove struct {
185+
update *order.Remove
186+
}
187+
188+
func NewRemove(diff *order.Remove) *StateRemove {
189+
return &StateRemove{
190+
update: diff,
191+
}
192+
}
193+
194+
func (remove *StateRemove) Topic() abstraction.BrokerTopic {
195+
return StateName
196+
}
197+
198+
type StateClear struct {
199+
board string
200+
}
201+
202+
func NewStateClear(boardName string) *StateClear {
203+
return &StateClear{
204+
board: boardName,
205+
}
206+
}
207+
208+
func (clear *StateClear) Topic() abstraction.BrokerTopic {
209+
return StateName
210+
}

0 commit comments

Comments
 (0)