This commit is contained in:
matst80
2024-11-08 21:58:28 +01:00
parent dfcdf0939f
commit 65a969443a
13 changed files with 437 additions and 173 deletions

View File

@@ -1,109 +1,140 @@
package main
import (
"bufio"
"bytes"
"encoding/binary"
"fmt"
"io"
"time"
"git.tornberg.me/go-cart-actor/git.tornberg.me/go-cart-actor/messages"
messages "git.tornberg.me/go-cart-actor/proto"
"google.golang.org/protobuf/proto"
)
type StorableMessage interface {
GetBytes() ([]byte, error)
FromReader(io.Reader, *Message) error
Write(w io.Writer) error
}
type Message struct {
Type uint64
Type uint16
TimeStamp *int64
Content interface{}
}
type MessageWriter struct {
writer io.Writer
io.Writer
}
func NewMessageWriter(b *bytes.Buffer) *MessageWriter {
return &MessageWriter{writer: bufio.NewWriter(b)}
type StorableMessageHeader struct {
Version uint16
Type uint16
TimeStamp int64
DataLength uint64
}
func (w *MessageWriter) WriteUint64(value uint64) error {
bytes := make([]byte, 8)
binary.LittleEndian.PutUint64(bytes, value)
_, err := w.writer.Write(bytes)
return err
}
func (w *MessageWriter) WriteInt64(value int64) error {
return w.WriteUint64(uint64(value))
}
func (w *MessageWriter) WriteMessage(m *Message) error {
if err := w.WriteUint64(m.Type); err != nil {
return err
}
if err := w.WriteInt64(*m.TimeStamp); err != nil {
return err
}
var messageBytes []byte
var err error
if m.Type == AddRequestType {
messageBytes, err = proto.Marshal(m.Content.(*messages.AddRequest))
} else if m.Type == AddItemType {
messageBytes, err = proto.Marshal(m.Content.(*messages.AddItem))
} else {
return fmt.Errorf("unknown message type")
func GetData(fn func(w io.Writer) error) ([]byte, error) {
var buf bytes.Buffer
err := fn(&buf)
if err != nil {
return nil, err
}
b := buf.Bytes()
return b, nil
}
// func (w *MessageWriter) WriteUint64(value uint64) (int, error) {
// bytes := make([]byte, 8)
// binary.LittleEndian.PutUint64(bytes, value)
// return w.Write(bytes)
// }
// func (w *MessageWriter) WriteInt64(value int64) (int, error) {
// return w.WriteUint64(uint64(value))
// }
// func (w *MessageWriter) WriteMessage(m *Message) (int, error) {
// var i, l int
// var err error
// i, err = w.WriteUint64(m.Type)
// l += i
// i, err = w.WriteInt64(*m.TimeStamp)
// l += i
// var messageBytes []byte
// var err error
// if m.Type == AddRequestType {
// messageBytes, err = proto.Marshal(m.Content.(*messages.AddRequest))
// } else if m.Type == AddItemType {
// messageBytes, err = proto.Marshal(m.Content.(*messages.AddItem))
// } else {
// return fmt.Errorf("unknown message type")
// }
// if err != nil {
// return err
// }
// if err := w.WriteUint64(uint64(len(messageBytes))); err != nil {
// return err
// }
// _, err = w.Write(messageBytes)
// return err
// }
func (m Message) Write(w io.Writer) error {
data, err := GetData(func(wr io.Writer) error {
if m.Type == AddRequestType {
messageBytes, err := proto.Marshal(m.Content.(*messages.AddRequest))
if err != nil {
return err
}
wr.Write(messageBytes)
} else if m.Type == AddItemType {
messageBytes, err := proto.Marshal(m.Content.(*messages.AddItem))
if err != nil {
return err
}
wr.Write(messageBytes)
}
return nil
})
if err != nil {
return err
}
if err := w.WriteUint64(uint64(len(messageBytes))); err != nil {
return err
ts := time.Now().Unix()
if m.TimeStamp != nil {
ts = *m.TimeStamp
}
_, err = w.writer.Write(messageBytes)
err = binary.Write(w, binary.LittleEndian, StorableMessageHeader{
Version: 1,
Type: m.Type,
TimeStamp: ts,
DataLength: uint64(len(data)),
})
w.Write(data)
return err
}
func (m Message) GetBytes() ([]byte, error) {
var b bytes.Buffer
mw := NewMessageWriter(&b)
err := mw.WriteMessage(&m)
return b.Bytes(), err
}
func (i Message) FromReader(reader io.Reader, m *Message) error {
bytes := make([]byte, 8)
if _, err := reader.Read(bytes); err != nil {
func MessageFromReader(reader io.Reader, m *Message) error {
header := StorableMessageHeader{}
err := binary.Read(reader, binary.LittleEndian, &header)
if err != nil {
return err
}
m.Type = binary.LittleEndian.Uint64(bytes)
if _, err := reader.Read(bytes); err != nil {
messageBytes := make([]byte, header.DataLength)
_, err = reader.Read(messageBytes)
if err != nil {
return err
}
timestamp := int64(binary.LittleEndian.Uint64(bytes))
m.TimeStamp = &timestamp
if _, err := reader.Read(bytes); err != nil {
return err
}
messageBytes := make([]byte, binary.LittleEndian.Uint64(bytes))
if _, err := reader.Read(messageBytes); err != nil {
return err
}
var err error
if m.Type == AddRequestType {
switch header.Type {
case AddRequestType:
msg := &messages.AddRequest{}
err = proto.Unmarshal(messageBytes, msg)
m.Content = msg
} else if m.Type == AddItemType {
case AddItemType:
msg := &messages.AddItem{}
err = proto.Unmarshal(messageBytes, msg)
m.Content = msg
} else {
default:
return fmt.Errorf("unknown message type")
}
if err != nil {