summaryrefslogtreecommitdiff
path: root/internal/packet/framer.go
diff options
context:
space:
mode:
authorKyren223 <ulmliad223@gmail.com>2024-10-20 11:36:43 +0300
committerKyren223 <ulmliad223@gmail.com>2024-10-20 11:36:43 +0300
commit1997bc8a150b92783060cc7129e1ee72a761182b (patch)
tree35d4df750feb2f1f36883ce8d7f2da9a72c59a25 /internal/packet/framer.go
parent6e6ca1d7a10a3a4a1decbdb78e2ea11ce24974e3 (diff)
refactor: mid refactor
Diffstat (limited to 'internal/packet/framer.go')
-rw-r--r--internal/packet/framer.go106
1 files changed, 0 insertions, 106 deletions
diff --git a/internal/packet/framer.go b/internal/packet/framer.go
deleted file mode 100644
index 1a76ea6..0000000
--- a/internal/packet/framer.go
+++ /dev/null
@@ -1,106 +0,0 @@
-package packet
-
-import (
- "context"
- "encoding/binary"
- "errors"
- "fmt"
- "io"
-
- "github.com/kyren223/eko/pkg/assert"
- "github.com/kyren223/eko/pkg/util"
-)
-
-const framerPacketCapacity = 10
-
-var (
- PacketUnsupportedVersion error = errors.New("packet error: unsupported version")
- PacketUnsupportedEncoding error = errors.New("packet error: unsupported encoding")
- PacketUnsupportedType error = errors.New("packet error: unsupported type")
-)
-
-type packetFramer struct {
- buffer []byte
- len uint16
- in chan<- Packet
- inErr chan<- error
-}
-
-func RunFramer(ctx context.Context, reader io.Reader) (out <-chan Packet, outErr <-chan error) {
- ch := make(chan Packet, framerPacketCapacity)
- errCh := make(chan error)
-
- framer := packetFramer{
- buffer: make([]byte, PACKET_MAX_SIZE),
- len: 0,
- in: ch,
- inErr: errCh,
- }
-
- go framer.run(ctx, util.NewChannelReader(ctx, reader))
-
- return ch, errCh
-}
-
-func (f *packetFramer) run(ctx context.Context, reader util.ChannelReader) {
- defer close(f.in)
- defer close(f.inErr)
-
-outer:
- for {
- select {
- case data := <-reader.Out:
- dataRead := 0
- for dataRead < len(data) {
- n := copy(f.buffer[f.len:], data[dataRead:])
- assert.Assert(0 <= n && n <= int(PACKET_MAX_SIZE), "n must fit in a u16")
- f.len += uint16(n)
- dataRead += n
- if err := f.parse(); err != nil {
- f.inErr <- err
- break outer
- }
- }
-
- case err := <-reader.Err:
- f.inErr <- err
- break outer
-
- case <-ctx.Done():
- f.inErr <- ctx.Err()
- break outer
- }
- }
-}
-
-func (f *packetFramer) parse() error {
- for f.len > HEADER_SIZE {
- if f.buffer[VERSION_OFFSET] != VERSION {
- return fmt.Errorf("%w version=%v", PacketUnsupportedVersion, f.buffer[VERSION_OFFSET])
- }
-
- encoding := Encoding(f.buffer[ENCODING_OFFSET] >> 6)
- packetType := PacketType(f.buffer[TYPE_OFFSET] & 63)
- if !encoding.IsSupported() {
- return PacketUnsupportedEncoding
- }
- if !packetType.IsSupported() {
- return PacketUnsupportedType
- }
-
- length := binary.BigEndian.Uint16(f.buffer[LENGTH_OFFSET:])
- if f.len-HEADER_SIZE < length {
- // Wait for more data to arrive
- return nil
- }
-
- fullLength := HEADER_SIZE + length
- packetBuffer := make([]byte, fullLength)
- copy(packetBuffer, f.buffer[:fullLength])
-
- f.len = uint16(copy(f.buffer, f.buffer[fullLength:f.len]))
-
- f.in <- Packet{packetBuffer}
- }
- return nil
-}