Skip to content

Commit 953d42a

Browse files
authored
Merge pull request #83 from HyperloopUPV-H8/backend/client-broker
[backend] [PR 2] Missing broker topics
2 parents 972a0c8 + aba92cf commit 953d42a

10 files changed

Lines changed: 378 additions & 56 deletions

File tree

backend/cmd/main.go

Lines changed: 84 additions & 49 deletions
Original file line numberDiff line numberDiff line change
@@ -17,14 +17,10 @@ import (
1717

1818
blcuPackage "github.com/HyperloopUPV-H8/h9-backend/internal/blcu"
1919
"github.com/HyperloopUPV-H8/h9-backend/internal/common"
20-
"github.com/HyperloopUPV-H8/h9-backend/internal/connection_transfer"
21-
"github.com/HyperloopUPV-H8/h9-backend/internal/data_transfer"
2220
"github.com/HyperloopUPV-H8/h9-backend/internal/excel"
2321
"github.com/HyperloopUPV-H8/h9-backend/internal/excel/ade"
2422
"github.com/HyperloopUPV-H8/h9-backend/internal/excel/utils"
2523
"github.com/HyperloopUPV-H8/h9-backend/internal/info"
26-
"github.com/HyperloopUPV-H8/h9-backend/internal/message_transfer"
27-
"github.com/HyperloopUPV-H8/h9-backend/internal/order_transfer"
2824
"github.com/HyperloopUPV-H8/h9-backend/internal/pod_data"
2925
"github.com/HyperloopUPV-H8/h9-backend/internal/server"
3026
"github.com/HyperloopUPV-H8/h9-backend/internal/update_factory"
@@ -34,6 +30,9 @@ import (
3430
"github.com/HyperloopUPV-H8/h9-backend/pkg/broker"
3531
connection_topic "github.com/HyperloopUPV-H8/h9-backend/pkg/broker/topics/connection"
3632
data_topic "github.com/HyperloopUPV-H8/h9-backend/pkg/broker/topics/data"
33+
logger_topic "github.com/HyperloopUPV-H8/h9-backend/pkg/broker/topics/logger"
34+
message_topic "github.com/HyperloopUPV-H8/h9-backend/pkg/broker/topics/message"
35+
order_topic "github.com/HyperloopUPV-H8/h9-backend/pkg/broker/topics/order"
3736
"github.com/HyperloopUPV-H8/h9-backend/pkg/transport"
3837
"github.com/HyperloopUPV-H8/h9-backend/pkg/transport/network/sniffer"
3938
"github.com/HyperloopUPV-H8/h9-backend/pkg/transport/network/tcp"
@@ -111,20 +110,11 @@ func main() {
111110
}
112111
config.Vehicle.Network.Interface = dev.Name
113112

114-
connectionTransfer := connection_transfer.New(config.Connections)
115-
116113
vehicleOrders, err := vehicle_models.NewVehicleOrders(podData.Boards, config.Excel.Parse.Global.BLCUAddressKey)
117114
if err != nil {
118115
trace.Fatal().Err(err).Msg("creating vehicleOrders")
119116
}
120117

121-
// <--- data transfer --->
122-
dataTransfer := data_transfer.New(config.DataTransfer)
123-
go dataTransfer.Run()
124-
125-
// <--- message transfer --->
126-
messageTransfer := message_transfer.New(config.Messages)
127-
128118
// <--- update factory --->
129119
updateFactory := update_factory.NewFactory()
130120

@@ -146,24 +136,30 @@ func main() {
146136
idToBoard[packet.Id] = board.Name
147137
}
148138
}
149-
orderTransfer, orderChannel := order_transfer.New(idToBoard)
150139

151140
// <--- blcu --->
152141
var blcu blcuPackage.BLCU
153142
blcuAddr, useBlcu := info.Addresses.Boards["BLCU"]
154143

155144
// <--- broker --->
156145
broker := broker.New()
157-
broker.SetAPI(&brokerAPI{
158-
OnUserPush: func(push abstraction.BrokerPush) {},
159-
})
160146

161147
dataTopic := data_topic.NewUpdateTopic(time.Second / 10)
162148
defer dataTopic.Stop()
163149
connectionTopic := connection_topic.NewUpdateTopic()
150+
orderTopic := order_topic.NewSendTopic()
151+
loggerTopic := logger_topic.NewEnableTopic()
152+
boardIdToBoard := make(map[abstraction.BoardId]string)
153+
for name, id := range info.BoardIds {
154+
boardIdToBoard[abstraction.BoardId(id)] = name
155+
}
156+
messageTopic := message_topic.NewUpdateTopic(boardIdToBoard)
164157

165158
broker.AddTopic(data_topic.UpdateName, dataTopic)
166159
broker.AddTopic(connection_topic.UpdateName, connectionTopic)
160+
broker.AddTopic(order_topic.SendName, orderTopic)
161+
broker.AddTopic(logger_topic.EnableName, loggerTopic)
162+
broker.AddTopic(message_topic.UpdateName, messageTopic)
167163

168164
connections := make(chan *websocket.Client)
169165
upgrader := websocket.NewUpgrader(connections)
@@ -206,7 +202,10 @@ func main() {
206202
}
207203

208204
case *info_packet.Packet:
209-
messageTransfer.SendMessage(p)
205+
err := broker.Push(message_topic.Push(p))
206+
if err != nil {
207+
fmt.Println(err)
208+
}
210209

211210
err = loggerHandler.PushRecord(&messages_logger.Record{
212211
Packet: p,
@@ -220,7 +219,10 @@ func main() {
220219
}
221220

222221
case *protection.Packet:
223-
messageTransfer.SendMessage(p)
222+
err := broker.Push(message_topic.Push(p))
223+
if err != nil {
224+
fmt.Println(err)
225+
}
224226

225227
newPacket := info_packet.NewPacket(p.Id())
226228
newPacket.BoardId = p.BoardId
@@ -256,10 +258,9 @@ func main() {
256258
}
257259

258260
case *order.Add:
259-
orderTransfer.AddStateOrders(*p)
260-
261+
trace.Debug().Msg("adding order")
261262
case *order.Remove:
262-
orderTransfer.RemoveStateOrders(*p)
263+
trace.Debug().Msg("removing order")
263264
}
264265
},
265266

@@ -268,6 +269,65 @@ func main() {
268269
},
269270
})
270271

272+
// this is here because we need to use transport to send messages
273+
broker.SetAPI(&brokerAPI{
274+
OnUserPush: func(push abstraction.BrokerPush) {
275+
switch push.Topic() {
276+
case order_topic.SendName:
277+
order, ok := push.(*order_topic.Order)
278+
if !ok {
279+
trace.Error().Any("push", push).Msg("error casting push to order")
280+
return
281+
}
282+
283+
packet, err := order.ToPacket()
284+
if err != nil {
285+
trace.Error().Any("order", order).Err(err).Msg("error converting order to packet")
286+
return
287+
}
288+
289+
err = transp.SendMessage(transport.NewPacketMessage(packet))
290+
if err != nil {
291+
trace.Error().Any("order", order).Err(err).Msg("error sending order")
292+
return
293+
}
294+
295+
err = loggerHandler.PushRecord(&order_logger.Record{
296+
Packet: packet,
297+
From: "backend",
298+
To: idToBoard[uint16(packet.Id())],
299+
Timestamp: packet.Timestamp(),
300+
})
301+
302+
if err != nil {
303+
fmt.Println("Error pushing record to logger: ", err)
304+
}
305+
case logger_topic.EnableName:
306+
status, ok := push.(*logger_topic.Status)
307+
if !ok {
308+
trace.Error().Any("push", push).Msg("error casting push to enable")
309+
fmt.Printf("Push Type: %v\n", push)
310+
return
311+
}
312+
313+
var err error
314+
if status.Enable() {
315+
err = loggerHandler.Start()
316+
} else {
317+
err = loggerHandler.Stop()
318+
}
319+
320+
if err != nil {
321+
status.Fulfill(!status.Enable())
322+
} else {
323+
status.Fulfill(status.Enable())
324+
}
325+
default:
326+
fmt.Printf("unknow topic %s\n", push.Topic())
327+
}
328+
},
329+
})
330+
271331
// Load and set packet decoder and encoder
272332
decoder, encoder := getTransportDecEnc(info, podData)
273333
transp.WithDecoder(decoder).WithEncoder(encoder)
@@ -311,27 +371,6 @@ func main() {
311371
}
312372
go transp.HandleSniffer(sniffer.New(source, &layers.LayerTypeEthernet))
313373

314-
// <--- order transfer --->
315-
go func() {
316-
for order := range orderChannel {
317-
err := transp.SendMessage(transport.NewPacketMessage(&order))
318-
if err != nil {
319-
trace.Error().Any("order", order).Err(err).Msg("error sending order")
320-
}
321-
322-
err = loggerHandler.PushRecord(&order_logger.Record{
323-
Packet: &order,
324-
From: "backend",
325-
To: idToBoard[uint16(order.Id())],
326-
Timestamp: order.Timestamp(),
327-
})
328-
329-
if err != nil {
330-
fmt.Println("Error pushing record to logger: ", err)
331-
}
332-
}
333-
}()
334-
335374
// <--- blcu --->
336375
if useBlcu {
337376
blcu = blcuPackage.NewBLCU(net.TCPAddr{
@@ -352,12 +391,6 @@ func main() {
352391
websocketBroker.RegisterHandle(&blcu, config.BLCU.Topics.Upload, config.BLCU.Topics.Download)
353392
}
354393

355-
websocketBroker.RegisterHandle(&connectionTransfer, config.Connections.UpdateTopic, "connection/update")
356-
websocketBroker.RegisterHandle(&dataTransfer, "podData/update")
357-
websocketBroker.RegisterHandle(loggerHandler, config.LoggerHandler.Topics.Enable)
358-
websocketBroker.RegisterHandle(&messageTransfer, "message/update")
359-
websocketBroker.RegisterHandle(&orderTransfer, config.Orders.SendTopic, "order/stateOrders")
360-
361394
uploadableBords := common.Filter(common.Keys(info.Addresses.Boards), func(item string) bool {
362395
return item != config.Excel.Parse.Global.BLCUAddressKey
363396
})
@@ -568,6 +601,8 @@ func getTransportDecEnc(info info.Info, podData pod_data.PodData) (*presentation
568601
protectionDecoder := protection.NewDecoder()
569602
protectionDecoder.SetSeverity(abstraction.PacketId(info.MessageIds.Warning), protection.SeverityWarning)
570603
protectionDecoder.SetSeverity(abstraction.PacketId(info.MessageIds.Fault), protection.SeverityFault)
604+
decoder.SetPacketDecoder(abstraction.PacketId(info.MessageIds.Warning), protectionDecoder)
605+
decoder.SetPacketDecoder(abstraction.PacketId(info.MessageIds.Fault), protectionDecoder)
571606

572607
return decoder, encoder
573608
}

backend/internal/vehicle/models/order_data.go

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -121,6 +121,7 @@ func getField(m pod_data.Measurement) (any, error) {
121121
Kind: BooleanKind,
122122
Name: typedMeas.Name,
123123
},
124+
VarType: "bool",
124125
}, nil
125126
case pod_data.EnumMeasurement:
126127
return EnumDescription{
@@ -130,6 +131,7 @@ func getField(m pod_data.Measurement) (any, error) {
130131
Name: typedMeas.Name,
131132
},
132133
Options: typedMeas.Options,
134+
VarType: "enum",
133135
}, nil
134136
default:
135137
return struct{}{}, errors.New("unrecognized measurement type")
@@ -157,9 +159,11 @@ type NumericDescription struct {
157159

158160
type BooleanDescription struct {
159161
fieldDescription
162+
VarType string `json:"type"`
160163
}
161164

162165
type EnumDescription struct {
163166
fieldDescription
164167
Options []string `json:"options"`
168+
VarType string `json:"type"`
165169
}

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

Lines changed: 80 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,15 +1,29 @@
11
package data
22

33
import (
4+
"encoding/json"
5+
"fmt"
6+
"sync/atomic"
7+
48
"github.com/HyperloopUPV-H8/h9-backend/pkg/abstraction"
59
"github.com/HyperloopUPV-H8/h9-backend/pkg/websocket"
610
)
711

812
const EnableName abstraction.BrokerTopic = "logger/enable"
13+
const ResponseName abstraction.BrokerTopic = "logger/response"
914

1015
type Enable struct {
11-
pool *websocket.Pool
12-
api abstraction.BrokerAPI
16+
isRunning *atomic.Bool
17+
pool *websocket.Pool
18+
api abstraction.BrokerAPI
19+
}
20+
21+
func NewEnableTopic() *Enable {
22+
enable := &Enable{
23+
isRunning: &atomic.Bool{},
24+
}
25+
enable.isRunning.Store(false)
26+
return enable
1327
}
1428

1529
func (enable *Enable) Topic() abstraction.BrokerTopic {
@@ -24,10 +38,74 @@ func (enable *Enable) Pull(request abstraction.BrokerRequest) (abstraction.Broke
2438
return nil, nil
2539
}
2640

41+
func (enable *Enable) ClientMessage(id websocket.ClientId, message *websocket.Message) {
42+
switch message.Topic {
43+
case EnableName:
44+
err := enable.handleToggle(id, message)
45+
if err != nil {
46+
fmt.Printf("error handling logger: %v\n", err)
47+
}
48+
}
49+
}
50+
51+
func (enable *Enable) handleToggle(id websocket.ClientId, message *websocket.Message) error {
52+
var request bool
53+
err := json.Unmarshal(message.Payload, &request)
54+
if err != nil {
55+
return err
56+
}
57+
58+
status := newStatus(request)
59+
go enable.api.UserPush(status)
60+
61+
go func() {
62+
enable.isRunning.Store(<-status.response)
63+
enable.broadcastState()
64+
}()
65+
return nil
66+
}
67+
68+
func (enable *Enable) broadcastState() error {
69+
payload, err := json.Marshal(enable.isRunning.Load())
70+
if err != nil {
71+
return err
72+
}
73+
74+
enable.pool.Broadcast(websocket.Message{
75+
Topic: ResponseName,
76+
Payload: payload,
77+
})
78+
return nil
79+
}
80+
2781
func (enable *Enable) SetPool(pool *websocket.Pool) {
2882
enable.pool = pool
2983
}
3084

3185
func (enable *Enable) SetAPI(api abstraction.BrokerAPI) {
3286
enable.api = api
3387
}
88+
89+
type Status struct {
90+
request bool
91+
response chan bool
92+
}
93+
94+
func newStatus(request bool) *Status {
95+
return &Status{
96+
request: request,
97+
response: make(chan bool),
98+
}
99+
}
100+
101+
func (status *Status) Topic() abstraction.BrokerTopic {
102+
return EnableName
103+
}
104+
105+
func (status *Status) Fulfill(response bool) {
106+
status.response <- response
107+
}
108+
109+
func (status *Status) Enable() bool {
110+
return status.request
111+
}

0 commit comments

Comments
 (0)