Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114// package millproto carries the mill<->executor session protocol. a new message// vocabulary over the same length-prefixed protobuf framing the spindle already// uses to talk to the microVM guest (see spindle/agentproto). only the framing// pattern is shared. the messages are entirely separatepackage millproto
import ( "encoding/binary" "fmt" "io" "sync"
"buf.build/go/protovalidate" "google.golang.org/protobuf/proto"
millv1 "tangled.org/core/spindle/mill/proto/gen")
const ( // ProtocolVersion is the wire version selected by the current implementation. ProtocolVersion = 4 // ProtocolMinVersion and ProtocolMaxVersion delimit versions this endpoint // can negotiate during rolling upgrades. ProtocolMinVersion = ProtocolVersion ProtocolMaxVersion = ProtocolVersion // CacheSentinelKey is initialized by the mill and read by executors before // they advertise cache support. Its value prevents an executor-local store // from passing the shared-store check by merely containing the key. CacheSentinelKey = "spindle-cache-sentinel" CacheSentinelValue = "spindle-cache-sentinel-v1" // generous vs agentproto's 1 MiB. a ReserveSeat carries the raw pipeline and // workflow JSON, and streamed log lines can be chunky MaxMessageBytes = 8 * 1024 * 1024)
type Message = millv1.Message
var validator protovalidate.Validator
func init() { var err error validator, err = protovalidate.New() if err != nil { panic(fmt.Errorf("failed to initialize protovalidate validator: %w", err)) }}
type Encoder struct { mu sync.Mutex w io.Writer}
func NewEncoder(w io.Writer) *Encoder { return &Encoder{w: w}}
func (e *Encoder) Encode(msg *Message) error { if err := validator.Validate(msg); err != nil { return fmt.Errorf("validate fleet message: %w", err) }
data, err := proto.Marshal(msg) if err != nil { return fmt.Errorf("marshal fleet message: %w", err) } if len(data) > MaxMessageBytes { return fmt.Errorf("fleet message exceeded %d bytes", MaxMessageBytes) }
// single write of header and payload maps to exactly one websocket binary // frame when the writer is a ws stream frame := make([]byte, 4+len(data)) binary.BigEndian.PutUint32(frame[:4], uint32(len(data))) copy(frame[4:], data)
e.mu.Lock() defer e.mu.Unlock() _, err = e.w.Write(frame) return err}
type Decoder struct { r io.Reader}
func NewDecoder(r io.Reader) *Decoder { return &Decoder{r: r}}
func (d *Decoder) Decode() (*Message, error) { msg := &Message{} var header [4]byte if _, err := io.ReadFull(d.r, header[:]); err != nil { return msg, err }
size := binary.BigEndian.Uint32(header[:]) if size > MaxMessageBytes { return msg, fmt.Errorf("fleet message exceeded %d bytes", MaxMessageBytes) }
data := make([]byte, size) if _, err := io.ReadFull(d.r, data); err != nil { return msg, err } if err := proto.Unmarshal(data, msg); err != nil { return msg, fmt.Errorf("parse fleet message: %w", err) } if err := validator.Validate(msg); err != nil { return msg, fmt.Errorf("validate fleet message: %w", err) } return msg, nil}