WIP
This commit is contained in:
@@ -1,67 +1,35 @@
|
||||
package peer
|
||||
|
||||
import (
|
||||
"io"
|
||||
"log"
|
||||
"net/netip"
|
||||
"sync/atomic"
|
||||
)
|
||||
|
||||
type ifReader struct {
|
||||
iface io.Reader
|
||||
writeToUDPAddrPort func([]byte, netip.AddrPort) (int, error)
|
||||
rt *atomic.Pointer[routingTable]
|
||||
buf1 []byte
|
||||
buf2 []byte
|
||||
type IFReader struct {
|
||||
Globals
|
||||
}
|
||||
|
||||
func newIFReader(
|
||||
iface io.Reader,
|
||||
writeToUDPAddrPort func([]byte, netip.AddrPort) (int, error),
|
||||
rt *atomic.Pointer[routingTable],
|
||||
) *ifReader {
|
||||
return &ifReader{iface, writeToUDPAddrPort, rt, newBuf(), newBuf()}
|
||||
func NewIFReader(g Globals) *IFReader {
|
||||
return &IFReader{Globals: g}
|
||||
}
|
||||
|
||||
func (r *ifReader) Run() {
|
||||
packet := newBuf()
|
||||
func (r *IFReader) Run() {
|
||||
packet := make([]byte, bufferSize)
|
||||
for {
|
||||
r.handleNextPacket(packet)
|
||||
}
|
||||
}
|
||||
|
||||
func (r *ifReader) handleNextPacket(packet []byte) {
|
||||
func (r *IFReader) handleNextPacket(packet []byte) {
|
||||
packet = r.readNextPacket(packet)
|
||||
remoteIP, ok := r.parsePacket(packet)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
|
||||
rt := r.rt.Load()
|
||||
peer := rt.Peers[remoteIP]
|
||||
if !peer.Up {
|
||||
r.logf("Peer %d not up.", peer.IP)
|
||||
return
|
||||
}
|
||||
|
||||
enc := peer.EncryptDataPacket(peer.IP, packet, r.buf1)
|
||||
if peer.Direct {
|
||||
r.writeToUDPAddrPort(enc, peer.DirectAddr)
|
||||
return
|
||||
}
|
||||
|
||||
relay, ok := rt.GetRelay()
|
||||
if !ok {
|
||||
r.logf("Relay not available for peer %d.", peer.IP)
|
||||
return
|
||||
}
|
||||
|
||||
enc = relay.EncryptDataPacket(peer.IP, enc, r.buf2)
|
||||
r.writeToUDPAddrPort(enc, relay.DirectAddr)
|
||||
r.RemotePeers[remoteIP].Load().SendDataTo(packet)
|
||||
}
|
||||
|
||||
func (r *ifReader) readNextPacket(buf []byte) []byte {
|
||||
n, err := r.iface.Read(buf[:cap(buf)])
|
||||
func (r *IFReader) readNextPacket(buf []byte) []byte {
|
||||
n, err := r.IFace.Read(buf[:cap(buf)])
|
||||
if err != nil {
|
||||
log.Fatalf("Failed to read from interface: %v", err)
|
||||
}
|
||||
@@ -69,7 +37,9 @@ func (r *ifReader) readNextPacket(buf []byte) []byte {
|
||||
return buf[:n]
|
||||
}
|
||||
|
||||
func (r *ifReader) parsePacket(buf []byte) (byte, bool) {
|
||||
// parsePacket returns the VPN ip for the packet, and a boolean indicating
|
||||
// success.
|
||||
func (r *IFReader) parsePacket(buf []byte) (byte, bool) {
|
||||
n := len(buf)
|
||||
if n == 0 {
|
||||
return 0, false
|
||||
@@ -98,6 +68,6 @@ func (r *ifReader) parsePacket(buf []byte) (byte, bool) {
|
||||
}
|
||||
}
|
||||
|
||||
func (*ifReader) logf(s string, args ...any) {
|
||||
func (*IFReader) logf(s string, args ...any) {
|
||||
log.Printf("[IFReader] "+s, args...)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user