Skip to content

Commit 27ed4f6

Browse files
authored
Merge pull request #153 from HyperloopUPV-H8/control-station/second-version
[Control Station] Second version
2 parents 4e00e0a + 8c9d597 commit 27ed4f6

73 files changed

Lines changed: 1555 additions & 655 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

backend/cmd/main.go

Lines changed: 18 additions & 39 deletions
Original file line numberDiff line numberDiff line change
@@ -63,6 +63,7 @@ var cpuprofile = flag.String("cpuprofile", "", "write cpu profile to file")
6363
var enableSNTP = flag.Bool("sntp", false, "enables a simple SNTP server on port 123")
6464
var networkDevice = flag.Int("dev", -1, "index of the network device to use, overrides device prompt")
6565
var blockprofile = flag.Int("blockprofile", 0, "number of block profiles to include")
66+
var playbackFile = flag.String("playback", "", "")
6667

6768
func main() {
6869
flag.Parse()
@@ -244,21 +245,25 @@ func main() {
244245
if err != nil {
245246
panic("failed to obtain sniffer source: " + err.Error())
246247
}
248+
249+
if *playbackFile != "" {
250+
source, err = pcap.OpenOffline(*playbackFile)
251+
if err != nil {
252+
panic("failed to obtain sniffer source: " + err.Error())
253+
}
254+
}
255+
247256
boardIps := make([]net.IP, 0)
248257
for _, board := range info.Addresses.Boards {
249258
boardIps = append(boardIps, board)
250259
}
251-
err = source.SetBPFFilter(getFilter(boardIps, info.Addresses.Backend, info.Ports.UDP, info.Ports.TcpClient, info.Ports.TcpServer))
260+
filter := getFilter(boardIps, info.Addresses.Backend, info.Ports.UDP)
261+
trace.Warn().Str("filter", filter).Msg("filter")
262+
err = source.SetBPFFilter(filter)
252263
if err != nil {
253264
panic("failed to compile bpf filter")
254265
}
255-
go func() {
256-
sniffer := sniffer.New(source, &layers.LayerTypeEthernet, trace.Logger)
257-
for {
258-
errChan := transp.HandleSniffer(sniffer)
259-
trace.Error().Stack().Err(<-errChan).Msg("sniffer crashed, restarting...")
260-
}
261-
}()
266+
go transp.HandleSniffer(sniffer.New(source, &layers.LayerTypeEthernet, trace.Logger))
262267

263268
// <--- http server --->
264269
podDataHandle, err := h.HandleDataJSON("podData.json", pod_data.GetDataOnlyPodData(podData))
@@ -534,15 +539,11 @@ func getOps(units utils.Units) data.ConversionDescriptor {
534539
return output
535540
}
536541

537-
func getFilter(boardAddrs []net.IP, backendAddr net.IP, udpPort uint16, tcpClientPort uint16, tcpServerPort uint16) string {
542+
func getFilter(boardAddrs []net.IP, backendAddr net.IP, udpPort uint16) string {
538543
ipipFilter := getIPIPfilter()
539-
udpFilter := getUDPFilter(boardAddrs, udpPort)
540-
tcpFilter := getTCPFilter(boardAddrs, tcpServerPort, tcpClientPort)
541-
// noBackend := "not host 192.168.0.9"
542-
543-
// filter := fmt.Sprintf("((%s) or (%s) or (%s)) and (%s)", ipipFilter, udpFilter, tcpFilter, noBackend)
544+
udpFilter := getUDPFilter(boardAddrs, backendAddr, udpPort)
544545

545-
filter := fmt.Sprintf("(%s) or (%s) or (%s)", ipipFilter, udpFilter, tcpFilter)
546+
filter := fmt.Sprintf("(%s) or (%s)", ipipFilter, udpFilter)
546547

547548
trace.Trace().Any("addrs", boardAddrs).Str("filter", filter).Msg("new filter")
548549
return filter
@@ -552,35 +553,13 @@ func getIPIPfilter() string {
552553
return "ip[9] == 4"
553554
}
554555

555-
func getUDPFilter(addrs []net.IP, port uint16) string {
556+
func getUDPFilter(addrs []net.IP, backendAddr net.IP, port uint16) string {
556557
udpPort := fmt.Sprintf("udp port %d", port)
557558
udpAddrs := common.Map(addrs, func(addr net.IP) string {
558559
return fmt.Sprintf("(src host %s)", addr)
559560
})
560561

561562
udpAddrsStr := strings.Join(udpAddrs, " or ")
562563

563-
return fmt.Sprintf("(%s) and (%s)", udpPort, udpAddrsStr)
564-
}
565-
566-
func getTCPFilter(addrs []net.IP, serverPort uint16, clientPort uint16) string {
567-
ports := fmt.Sprintf("tcp port %d or %d", serverPort, clientPort)
568-
notSynFinRst := "tcp[tcpflags] & (tcp-fin | tcp-syn | tcp-rst) == 0"
569-
notJustAck := "tcp[tcpflags] | tcp-ack != 16"
570-
nonZeroPayload := "tcp[tcpflags] & tcp-push != 0"
571-
572-
srcAddresses := common.Map(addrs, func(addr net.IP) string {
573-
return fmt.Sprintf("(src host %s)", addr)
574-
})
575-
576-
srcAddressesStr := strings.Join(srcAddresses, " or ")
577-
578-
dstAddresses := common.Map(addrs, func(addr net.IP) string {
579-
return fmt.Sprintf("(dst host %s)", addr)
580-
})
581-
582-
dstAddressesStr := strings.Join(dstAddresses, " or ")
583-
584-
filter := fmt.Sprintf("(%s) and (%s) and (%s) and (%s) and (%s) and (%s)", ports, notSynFinRst, notJustAck, nonZeroPayload, srcAddressesStr, dstAddressesStr)
585-
return filter
564+
return fmt.Sprintf("(%s) and (%s) and (dst host %s)", udpPort, udpAddrsStr, backendAddr)
586565
}

backend/internal/common/moving_average.go

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,10 @@ func NewMovingAverage[N Numeric](order uint) *MovingAverage[N] {
1515
}
1616

1717
func (avg *MovingAverage[N]) Add(value N) N {
18+
if avg == nil {
19+
return 0
20+
}
21+
1822
if avg.length < avg.Order() {
1923
avg.addElem(value)
2024
} else {
@@ -38,6 +42,10 @@ func (avg *MovingAverage[N]) Order() int {
3842
}
3943

4044
func (avg *MovingAverage[N]) Resize(order uint) N {
45+
if avg == nil {
46+
return 0
47+
}
48+
4149
if order > uint(avg.Order()) {
4250
avg.Grow(order - uint(avg.Order()))
4351
} else if order < uint(avg.Order()) {

backend/pkg/transport/constructor.go

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -11,14 +11,17 @@ import (
1111
)
1212

1313
func NewTransport(baseLogger zerolog.Logger) *Transport {
14-
return &Transport{
14+
transport := &Transport{
1515
connectionsMx: &sync.Mutex{},
1616
connections: make(map[abstraction.TransportTarget]net.Conn),
1717
idToTarget: make(map[abstraction.PacketId]abstraction.TransportTarget),
1818
ipToTarget: make(map[string]abstraction.TransportTarget),
1919

20-
logger: baseLogger,
20+
logger: baseLogger,
21+
errChan: make(chan error, 100),
2122
}
23+
go transport.consumeErrors()
24+
return transport
2225
}
2326

2427
func (transport *Transport) WithDecoder(decoder *presentation.Decoder) *Transport {

backend/pkg/transport/messages.go

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,8 @@ import (
99
const (
1010
// PacketEvent is triggered when a packet is sent or received
1111
PacketEvent abstraction.TransportEvent = "packet"
12+
// ErrorEvent is triggered when an error arises somewhere
13+
ErrorEvent abstraction.TransportEvent = "error"
1214
// FileWriteEvent is used to request a file write through tftp
1315
FileWriteEvent abstraction.TransportEvent = "file-push"
1416
// FileReadEvent is used to request a file read through tftp

backend/pkg/transport/network/sniffer/decoder.go

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -34,14 +34,14 @@ func (dec *decoder) IPv4() layers.IPv4 {
3434
return dec.ipipv4
3535
}
3636

37-
func (dec *decoder) TCP() layers.TCP {
38-
return dec.tcp
39-
}
40-
4137
func (dec *decoder) UDP() layers.UDP {
4238
return dec.udp
4339
}
4440

41+
func (dec *decoder) TCP() layers.TCP {
42+
return dec.tcp
43+
}
44+
4545
func (dec *decoder) Payload() []byte {
4646
return dec.payload
4747
}

backend/pkg/transport/network/sniffer/sniffer.go

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -76,6 +76,7 @@ func (sniffer *Sniffer) ReadNext() (network.Payload, error) {
7676
gotPorts := false
7777
gotIp := false
7878
layerArr := zerolog.Arr()
79+
payloadData := []byte{}
7980
for _, layer := range packetLayers {
8081
if !gotIp && layer == layers.LayerTypeIPv4 {
8182
ip := sniffer.decoder.IPv4()
@@ -97,6 +98,10 @@ func (sniffer *Sniffer) ReadNext() (network.Payload, error) {
9798
}
9899
}
99100

101+
if layer == gopacket.LayerTypePayload {
102+
payloadData = sniffer.decoder.Payload()
103+
}
104+
100105
layerArr = layerArr.Str(layer.String())
101106
}
102107

@@ -123,7 +128,7 @@ func (sniffer *Sniffer) ReadNext() (network.Payload, error) {
123128

124129
return network.Payload{
125130
Socket: socket,
126-
Data: sniffer.decoder.Payload(),
131+
Data: payloadData,
127132
Timestamp: captureInfo.Timestamp,
128133
}, nil
129134
}

backend/pkg/transport/network/tcp/config.go

Lines changed: 7 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -37,16 +37,20 @@ func NewClientConfig(laddr net.Addr) ClientConfig {
3737
type backoffFunction = func(int) time.Duration
3838

3939
const (
40-
defaultBackoffMin time.Duration = 100 * time.Millisecond
41-
defaultBackoffExp float64 = 1.5
42-
defaultBackoffMax time.Duration = 5 * time.Second
40+
defaultBackoffMin time.Duration = 100 * time.Millisecond
41+
defaultBackoffExp float64 = 1.5
42+
defaultBackoffMax time.Duration = 5 * time.Second
43+
defaultBackoffRetriesClamp int = 12
4344
)
4445

4546
// NewExponentialBackoff returns an exponential backoff function with the given paramenters.
4647
//
4748
// It follows this formula: delay = (min * (exp ^ n); delay < max ? delay : max
4849
func NewExponentialBackoff(min time.Duration, exp float64, max time.Duration) backoffFunction {
4950
return func(n int) time.Duration {
51+
if n > defaultBackoffRetriesClamp {
52+
n = defaultBackoffRetriesClamp
53+
}
5054
curr := time.Duration(float64(min) * math.Trunc(math.Pow(exp, float64(n))))
5155
if curr >= max {
5256
return max

backend/pkg/transport/notifications.go

Lines changed: 17 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,9 @@
11
package transport
22

33
import (
4-
"github.com/HyperloopUPV-H8/h9-backend/pkg/abstraction"
54
"time"
5+
6+
"github.com/HyperloopUPV-H8/h9-backend/pkg/abstraction"
67
)
78

89
// PacketNotification notifies of an incoming message
@@ -26,3 +27,18 @@ func NewPacketNotification(packet abstraction.Packet, from string, to string, ti
2627
func (notification PacketNotification) Event() abstraction.TransportEvent {
2728
return PacketEvent
2829
}
30+
31+
// ErrorNotification notifies of an error that arised
32+
type ErrorNotification struct {
33+
Err error
34+
}
35+
36+
func NewErrorNotification(err error) ErrorNotification {
37+
return ErrorNotification{
38+
Err: err,
39+
}
40+
}
41+
42+
func (notification ErrorNotification) Event() abstraction.TransportEvent {
43+
return ErrorEvent
44+
}

backend/pkg/transport/packet/protection/packet.go

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -140,6 +140,19 @@ type Timestamp struct {
140140
Year uint16 `json:"year"`
141141
}
142142

143+
func NowTimestamp() *Timestamp {
144+
t := time.Now()
145+
return &Timestamp{
146+
Counter: 0,
147+
Second: uint8(t.Second()),
148+
Minute: uint8(t.Minute()),
149+
Hour: uint8(t.Hour()),
150+
Day: uint8(t.Day()),
151+
Month: uint8(t.Month()),
152+
Year: uint16(t.Year()),
153+
}
154+
}
155+
143156
func decodeTimestamp(reader io.Reader, endianness binary.ByteOrder) (*Timestamp, error) {
144157
packet := new(Timestamp)
145158
err := binary.Read(reader, endianness, packet)

backend/pkg/transport/presentation/decoder.go

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -72,7 +72,8 @@ func (decoder *Decoder) DecodeNext(reader io.Reader) (abstraction.Packet, error)
7272
dec, ok := decoder.idToDecoder[id]
7373
if !ok {
7474
decoder.logger.Warn().Uint16("id", uint16(id)).Msg("no decoder set")
75-
return nil, ErrUnexpectedId{Id: id}
75+
err := ErrUnexpectedId{Id: id}
76+
return nil, err
7677
}
7778

7879
decoder.logger.Debug().Uint16("id", uint16(id)).Type("decoder", dec).Msg("decoding")

0 commit comments

Comments
 (0)