Skip to content

Commit 6be7a0c

Browse files
committed
Move vehicle logic to vehicle abstraction
1 parent 28302d6 commit 6be7a0c

3 files changed

Lines changed: 185 additions & 206 deletions

File tree

backend/cmd/main.go

Lines changed: 9 additions & 198 deletions
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,6 @@ package main
33
import (
44
"bufio"
55
"encoding/binary"
6-
"errors"
76
"flag"
87
"fmt"
98
"log"
@@ -16,7 +15,6 @@ import (
1615
"strings"
1716
"time"
1817

19-
blcuPackage "github.com/HyperloopUPV-H8/h9-backend/internal/blcu"
2018
"github.com/HyperloopUPV-H8/h9-backend/internal/common"
2119
"github.com/HyperloopUPV-H8/h9-backend/internal/excel"
2220
"github.com/HyperloopUPV-H8/h9-backend/internal/excel/ade"
@@ -26,7 +24,6 @@ import (
2624
"github.com/HyperloopUPV-H8/h9-backend/internal/server"
2725
"github.com/HyperloopUPV-H8/h9-backend/internal/update_factory"
2826
vehicle_models "github.com/HyperloopUPV-H8/h9-backend/internal/vehicle/models"
29-
"github.com/HyperloopUPV-H8/h9-backend/internal/ws_handle"
3027
"github.com/HyperloopUPV-H8/h9-backend/pkg/abstraction"
3128
"github.com/HyperloopUPV-H8/h9-backend/pkg/broker"
3229
connection_topic "github.com/HyperloopUPV-H8/h9-backend/pkg/broker/topics/connection"
@@ -42,8 +39,8 @@ import (
4239
info_packet "github.com/HyperloopUPV-H8/h9-backend/pkg/transport/packet/info"
4340
"github.com/HyperloopUPV-H8/h9-backend/pkg/transport/packet/order"
4441
"github.com/HyperloopUPV-H8/h9-backend/pkg/transport/packet/protection"
45-
"github.com/HyperloopUPV-H8/h9-backend/pkg/transport/packet/state"
4642
"github.com/HyperloopUPV-H8/h9-backend/pkg/transport/presentation"
43+
"github.com/HyperloopUPV-H8/h9-backend/pkg/vehicle"
4744
"github.com/HyperloopUPV-H8/h9-backend/pkg/websocket"
4845
"github.com/fatih/color"
4946
"github.com/google/gopacket/layers"
@@ -138,10 +135,6 @@ func main() {
138135
}
139136
}
140137

141-
// <--- blcu --->
142-
var blcu blcuPackage.BLCU
143-
blcuAddr, useBlcu := info.Addresses.Boards["BLCU"]
144-
145138
// <--- broker --->
146139
broker := broker.New()
147140

@@ -179,156 +172,6 @@ func main() {
179172

180173
transp := transport.NewTransport()
181174

182-
transp.SetAPI(&transportAPI{
183-
OnNotification: func(notification abstraction.TransportNotification) {
184-
packet := notification.(transport.PacketNotification)
185-
186-
switch p := packet.Packet.(type) {
187-
case *data.Packet:
188-
update := updateFactory.NewUpdate(p)
189-
err := broker.Push(data_topic.NewPush(&update))
190-
if err != nil {
191-
fmt.Println(err)
192-
}
193-
194-
err = loggerHandler.PushRecord(&data_logger.Record{
195-
Packet: p,
196-
From: packet.From,
197-
To: packet.To,
198-
Timestamp: packet.Timestamp,
199-
})
200-
201-
if err != nil && !errors.Is(err, logger.ErrLoggerNotRunning{}) {
202-
fmt.Println("Error pushing record to data logger: ", err)
203-
}
204-
205-
case *info_packet.Packet:
206-
err := broker.Push(message_topic.Push(p))
207-
if err != nil {
208-
fmt.Println(err)
209-
}
210-
211-
err = loggerHandler.PushRecord(&messages_logger.Record{
212-
Packet: p,
213-
From: packet.From,
214-
To: packet.To,
215-
Timestamp: packet.Timestamp,
216-
})
217-
218-
if err != nil && !errors.Is(err, logger.ErrLoggerNotRunning{}) {
219-
fmt.Println("Error pushing record to info logger: ", err)
220-
}
221-
222-
case *protection.Packet:
223-
err := broker.Push(message_topic.Push(p))
224-
if err != nil {
225-
fmt.Println(err)
226-
}
227-
228-
newPacket := info_packet.NewPacket(p.Id())
229-
newPacket.BoardId = p.BoardId
230-
newPacket.Timestamp = p.Timestamp
231-
newPacket.Msg = info_packet.InfoData(fmt.Sprint(p))
232-
233-
err = loggerHandler.PushRecord(&messages_logger.Record{
234-
Packet: newPacket,
235-
From: packet.From,
236-
To: packet.To,
237-
Timestamp: packet.Timestamp,
238-
})
239-
240-
if err != nil && !errors.Is(err, logger.ErrLoggerNotRunning{}) {
241-
fmt.Println("Error pushing record to info logger: ", err)
242-
}
243-
244-
case *blcu_packet.Ack:
245-
if useBlcu {
246-
blcu.NotifyAck()
247-
}
248-
249-
case *state.Space:
250-
err = loggerHandler.PushRecord(&state_logger.Record{
251-
Packet: p,
252-
From: packet.From,
253-
To: packet.To,
254-
Timestamp: packet.Timestamp,
255-
})
256-
257-
if err != nil && !errors.Is(err, logger.ErrLoggerNotRunning{}) {
258-
fmt.Println("Error pushing record to state logger: ", err)
259-
}
260-
261-
case *order.Add:
262-
trace.Debug().Msg("adding order")
263-
case *order.Remove:
264-
trace.Debug().Msg("removing order")
265-
}
266-
},
267-
268-
OnConnectionUpdate: func(target abstraction.TransportTarget, isConnected bool) {
269-
connectionTopic.Push(connection_topic.NewConnection(string(target), isConnected))
270-
},
271-
})
272-
273-
// this is here because we need to use transport to send messages
274-
broker.SetAPI(&brokerAPI{
275-
OnUserPush: func(push abstraction.BrokerPush) {
276-
switch push.Topic() {
277-
case order_topic.SendName:
278-
order, ok := push.(*order_topic.Order)
279-
if !ok {
280-
trace.Error().Any("push", push).Msg("error casting push to order")
281-
return
282-
}
283-
284-
packet, err := order.ToPacket()
285-
if err != nil {
286-
trace.Error().Any("order", order).Err(err).Msg("error converting order to packet")
287-
return
288-
}
289-
290-
err = transp.SendMessage(transport.NewPacketMessage(packet))
291-
if err != nil {
292-
trace.Error().Any("order", order).Err(err).Msg("error sending order")
293-
return
294-
}
295-
296-
err = loggerHandler.PushRecord(&order_logger.Record{
297-
Packet: packet,
298-
From: "backend",
299-
To: idToBoard[uint16(packet.Id())],
300-
Timestamp: packet.Timestamp(),
301-
})
302-
303-
if err != nil && !errors.Is(err, logger.ErrLoggerNotRunning{}) {
304-
fmt.Println("Error pushing record to logger: ", err)
305-
}
306-
case logger_topic.EnableName:
307-
status, ok := push.(*logger_topic.Status)
308-
if !ok {
309-
trace.Error().Any("push", push).Msg("error casting push to enable")
310-
fmt.Printf("Push Type: %v\n", push)
311-
return
312-
}
313-
314-
var err error
315-
if status.Enable() {
316-
err = loggerHandler.Start()
317-
} else {
318-
err = loggerHandler.Stop()
319-
}
320-
321-
if err != nil {
322-
status.Fulfill(!status.Enable())
323-
} else {
324-
status.Fulfill(status.Enable())
325-
}
326-
default:
327-
fmt.Printf("unknow topic %s\n", push.Topic())
328-
}
329-
},
330-
})
331-
332175
// Load and set packet decoder and encoder
333176
decoder, encoder := getTransportDecEnc(info, podData)
334177
transp.WithDecoder(decoder).WithEncoder(encoder)
@@ -372,26 +215,15 @@ func main() {
372215
}
373216
go transp.HandleSniffer(sniffer.New(source, &layers.LayerTypeEthernet))
374217

375-
// <--- blcu --->
376-
if useBlcu {
377-
blcu = blcuPackage.NewBLCU(net.TCPAddr{
378-
IP: blcuAddr,
379-
Port: int(info.Ports.TFTP),
380-
}, info.BoardIds, config.BLCU)
381-
382-
blcu.SetSendOrder(func(o *data.Packet) error {
383-
return transp.SendMessage(transport.NewPacketMessage(o))
384-
})
385-
}
386-
387-
// <--- websocket broker --->
388-
websocketBroker := ws_handle.New()
389-
defer websocketBroker.Close()
390-
391-
if useBlcu {
392-
websocketBroker.RegisterHandle(&blcu, config.BLCU.Topics.Upload, config.BLCU.Topics.Download)
393-
}
218+
// <--- vehicle --->
219+
vehicle := vehicle.New()
220+
vehicle.SetBroker(broker)
221+
vehicle.SetLogger(loggerHandler)
222+
vehicle.SetUpdateFactory(updateFactory)
223+
vehicle.SetIdToBoardName(idToBoard)
224+
vehicle.SetTransport(transp)
394225

226+
// <--- http server --->
395227
uploadableBords := common.Filter(common.Keys(info.Addresses.Boards), func(item string) bool {
396228
return item != config.Excel.Parse.Global.BLCUAddressKey
397229
})
@@ -619,19 +451,6 @@ func getOps(units utils.Units) data.ConversionDescriptor {
619451
return output
620452
}
621453

622-
type transportAPI struct {
623-
OnNotification func(abstraction.TransportNotification)
624-
OnConnectionUpdate func(abstraction.TransportTarget, bool)
625-
}
626-
627-
func (api *transportAPI) Notification(notification abstraction.TransportNotification) {
628-
api.OnNotification(notification)
629-
}
630-
631-
func (api *transportAPI) ConnectionUpdate(target abstraction.TransportTarget, isConnected bool) {
632-
api.OnConnectionUpdate(target, isConnected)
633-
}
634-
635454
func getFilter(boardAddrs []net.IP, backendAddr net.IP, udpPort uint16, tcpClientPort uint16, tcpServerPort uint16) string {
636455
ipipFilter := getIPIPfilter()
637456
udpFilter := getUDPFilter(boardAddrs, udpPort)
@@ -682,11 +501,3 @@ func getTCPFilter(addrs []net.IP, serverPort uint16, clientPort uint16) string {
682501
filter := fmt.Sprintf("(%s) and (%s) and (%s) and (%s) and (%s) and (%s)", ports, notSynFinRst, notJustAck, nonZeroPayload, srcAddressesStr, dstAddressesStr)
683502
return filter
684503
}
685-
686-
type brokerAPI struct {
687-
OnUserPush func(abstraction.BrokerPush)
688-
}
689-
690-
func (api *brokerAPI) UserPush(push abstraction.BrokerPush) {
691-
api.OnUserPush(push)
692-
}

backend/pkg/vehicle/constructor.go

Lines changed: 12 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,9 @@
11
package vehicle
22

3-
import "github.com/HyperloopUPV-H8/h9-backend/pkg/abstraction"
3+
import (
4+
"github.com/HyperloopUPV-H8/h9-backend/internal/update_factory"
5+
"github.com/HyperloopUPV-H8/h9-backend/pkg/abstraction"
6+
)
47

58
// New creates a new Vehicle with no modules registered on it
69
func New() Vehicle {
@@ -41,3 +44,11 @@ func (vehicle *Vehicle) SetTransport(transport abstraction.Transport) {
4144
func (vehicle *Vehicle) SetLogger(logger abstraction.Logger) {
4245
vehicle.logger = logger
4346
}
47+
48+
func (vehicle *Vehicle) SetUpdateFactory(updateFactory *update_factory.UpdateFactory) {
49+
vehicle.updateFactory = updateFactory
50+
}
51+
52+
func (vehicle *Vehicle) SetIdToBoardName(idToBoardName map[uint16]string) {
53+
vehicle.idToBoardName = idToBoardName
54+
}

0 commit comments

Comments
 (0)