Skip to content

Commit cd59915

Browse files
committed
Add message topic
1 parent 1565cda commit cd59915

2 files changed

Lines changed: 138 additions & 17 deletions

File tree

backend/cmd/main.go

Lines changed: 17 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -17,12 +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/data_transfer"
2120
"github.com/HyperloopUPV-H8/h9-backend/internal/excel"
2221
"github.com/HyperloopUPV-H8/h9-backend/internal/excel/ade"
2322
"github.com/HyperloopUPV-H8/h9-backend/internal/excel/utils"
2423
"github.com/HyperloopUPV-H8/h9-backend/internal/info"
25-
"github.com/HyperloopUPV-H8/h9-backend/internal/message_transfer"
2624
"github.com/HyperloopUPV-H8/h9-backend/internal/pod_data"
2725
"github.com/HyperloopUPV-H8/h9-backend/internal/server"
2826
"github.com/HyperloopUPV-H8/h9-backend/internal/update_factory"
@@ -33,6 +31,7 @@ import (
3331
connection_topic "github.com/HyperloopUPV-H8/h9-backend/pkg/broker/topics/connection"
3432
data_topic "github.com/HyperloopUPV-H8/h9-backend/pkg/broker/topics/data"
3533
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"
3635
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"
@@ -116,13 +115,6 @@ func main() {
116115
trace.Fatal().Err(err).Msg("creating vehicleOrders")
117116
}
118117

119-
// <--- data transfer --->
120-
dataTransfer := data_transfer.New(config.DataTransfer)
121-
go dataTransfer.Run()
122-
123-
// <--- message transfer --->
124-
messageTransfer := message_transfer.New(config.Messages)
125-
126118
// <--- update factory --->
127119
updateFactory := update_factory.NewFactory()
128120

@@ -157,11 +149,17 @@ func main() {
157149
connectionTopic := connection_topic.NewUpdateTopic()
158150
orderTopic := order_topic.NewSendTopic()
159151
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)
160157

161158
broker.AddTopic(data_topic.UpdateName, dataTopic)
162159
broker.AddTopic(connection_topic.UpdateName, connectionTopic)
163160
broker.AddTopic(order_topic.SendName, orderTopic)
164161
broker.AddTopic(logger_topic.EnableName, loggerTopic)
162+
broker.AddTopic(message_topic.UpdateName, messageTopic)
165163

166164
connections := make(chan *websocket.Client)
167165
upgrader := websocket.NewUpgrader(connections)
@@ -201,7 +199,10 @@ func main() {
201199
}
202200

203201
case *info_packet.Packet:
204-
messageTransfer.SendMessage(p)
202+
err := broker.Push(message_topic.Push(p))
203+
if err != nil {
204+
fmt.Println(err)
205+
}
205206

206207
err = loggerHandler.PushRecord(&messages_logger.Record{
207208
Packet: p,
@@ -212,7 +213,10 @@ func main() {
212213
}
213214

214215
case *protection.Packet:
215-
messageTransfer.SendMessage(p)
216+
err := broker.Push(message_topic.Push(p))
217+
if err != nil {
218+
fmt.Println(err)
219+
}
216220

217221
packet := info_packet.NewPacket(p.Id())
218222
packet.BoardId = p.BoardId
@@ -372,9 +376,6 @@ func main() {
372376
websocketBroker.RegisterHandle(&blcu, config.BLCU.Topics.Upload, config.BLCU.Topics.Download)
373377
}
374378

375-
websocketBroker.RegisterHandle(loggerHandler, config.LoggerHandler.Topics.Enable)
376-
websocketBroker.RegisterHandle(&messageTransfer, "message/update")
377-
378379
uploadableBords := common.Filter(common.Keys(info.Addresses.Boards), func(item string) bool {
379380
return item != config.Excel.Parse.Global.BLCUAddressKey
380381
})
@@ -585,6 +586,8 @@ func getTransportDecEnc(info info.Info, podData pod_data.PodData) (*presentation
585586
protectionDecoder := protection.NewDecoder()
586587
protectionDecoder.SetSeverity(abstraction.PacketId(info.MessageIds.Warning), protection.SeverityWarning)
587588
protectionDecoder.SetSeverity(abstraction.PacketId(info.MessageIds.Fault), protection.SeverityFault)
589+
decoder.SetPacketDecoder(abstraction.PacketId(info.MessageIds.Warning), protectionDecoder)
590+
decoder.SetPacketDecoder(abstraction.PacketId(info.MessageIds.Fault), protectionDecoder)
588591

589592
return decoder, encoder
590593
}
Lines changed: 121 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,33 +1,151 @@
11
package data
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"
11+
"github.com/HyperloopUPV-H8/h9-backend/pkg/transport/packet/info"
12+
"github.com/HyperloopUPV-H8/h9-backend/pkg/transport/packet/protection"
513
"github.com/HyperloopUPV-H8/h9-backend/pkg/websocket"
14+
"github.com/google/uuid"
15+
ws "github.com/gorilla/websocket"
616
)
717

818
const UpdateName abstraction.BrokerTopic = "message/update"
19+
const SubscribeName abstraction.BrokerTopic = "message/update"
920

1021
type Update struct {
11-
pool *websocket.Pool
12-
api abstraction.BrokerAPI
22+
subscribersMx *sync.Mutex
23+
subscribers map[websocket.ClientId]struct{}
24+
idToBoard map[abstraction.BoardId]string
25+
pool *websocket.Pool
26+
api abstraction.BrokerAPI
27+
}
28+
29+
func NewUpdateTopic(idToBoard map[abstraction.BoardId]string) *Update {
30+
return &Update{
31+
subscribersMx: &sync.Mutex{},
32+
subscribers: make(map[websocket.ClientId]struct{}),
33+
idToBoard: idToBoard,
34+
}
1335
}
1436

1537
func (update *Update) Topic() abstraction.BrokerTopic {
1638
return UpdateName
1739
}
1840

19-
func (update *Update) Push(push abstraction.BrokerPush) error {
41+
func (update *Update) Push(p abstraction.BrokerPush) error {
42+
push, ok := p.(*push)
43+
if !ok {
44+
return topics.ErrUnexpectedPush{Push: p}
45+
}
46+
47+
raw, err := json.Marshal(push.Data(update.idToBoard))
48+
if err != nil {
49+
return err
50+
}
51+
52+
message := websocket.Message{
53+
Topic: UpdateName,
54+
Payload: raw,
55+
}
56+
57+
update.subscribersMx.Lock()
58+
defer update.subscribersMx.Unlock()
59+
60+
flagged := make([]websocket.ClientId, 0, len(update.subscribers))
61+
for client := range update.subscribers {
62+
err := update.pool.Write(client, message)
63+
if err != nil {
64+
flagged = append(flagged, client)
65+
}
66+
}
67+
68+
for _, id := range flagged {
69+
update.pool.Disconnect(id, ws.CloseUnsupportedData, "unsupported topic")
70+
delete(update.subscribers, id)
71+
fmt.Printf("message/update unsubscribed %s\n", uuid.UUID(id).String())
72+
}
73+
2074
return nil
2175
}
2276

2377
func (update *Update) Pull(request abstraction.BrokerRequest) (abstraction.BrokerResponse, error) {
2478
return nil, nil
2579
}
2680

81+
func (update *Update) ClientMessage(id websocket.ClientId, message *websocket.Message) {
82+
update.subscribersMx.Lock()
83+
defer update.subscribersMx.Unlock()
84+
85+
switch message.Topic {
86+
case SubscribeName:
87+
fmt.Printf("message/update subscribed %s\n", uuid.UUID(id).String())
88+
update.subscribers[id] = struct{}{}
89+
default:
90+
update.pool.Disconnect(id, ws.CloseUnsupportedData, "unsupported topic")
91+
delete(update.subscribers, id)
92+
fmt.Printf("message/update unsubscribed %s\n", uuid.UUID(id).String())
93+
}
94+
}
95+
2796
func (update *Update) SetPool(pool *websocket.Pool) {
2897
update.pool = pool
2998
}
3099

31100
func (update *Update) SetAPI(api abstraction.BrokerAPI) {
32101
update.api = api
33102
}
103+
104+
type push struct {
105+
data any
106+
}
107+
108+
func Push(data any) *push {
109+
return &push{data: data}
110+
}
111+
112+
func (push *push) Topic() abstraction.BrokerTopic {
113+
return UpdateName
114+
}
115+
116+
func (push *push) Data(idToBoard map[abstraction.BoardId]string) wrapper {
117+
switch data := push.data.(type) {
118+
case *protection.Packet:
119+
return wrapper{
120+
Kind: string(data.Severity()),
121+
Payload: struct {
122+
Kind string `json:"kind"`
123+
Data any `json:"data"`
124+
}{
125+
Kind: string(data.Protection.Type),
126+
Data: data.Protection.Data,
127+
},
128+
Board: string(idToBoard[data.BoardId]),
129+
Name: string(data.Protection.Name),
130+
Timestamp: data.Timestamp,
131+
}
132+
case *info.Packet:
133+
return wrapper{
134+
Kind: "info",
135+
Payload: string(data.Msg),
136+
Board: string(idToBoard[data.BoardId]),
137+
Name: "info",
138+
Timestamp: data.Timestamp,
139+
}
140+
}
141+
142+
return wrapper{}
143+
}
144+
145+
type wrapper struct {
146+
Kind string `json:"kind"`
147+
Payload any `json:"payload"`
148+
Board string `json:"board"`
149+
Name string `json:"name"`
150+
Timestamp packet.Timestamp `json:"timestamp"`
151+
}

0 commit comments

Comments
 (0)