Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 5 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -260,11 +260,14 @@ make check # everything CI runs

## Status and roadmap

The local and SSH transports are complete and covered end to end. Planned:
The local and SSH transports are complete and covered end to end. SSH
planning-phase manifests are gzip-compressed when both sides advertise
the `manifest-gzip` feature (protocol version stays 1, so older helpers
still work). Planned:

- QUIC transport for high-latency links
- Direct Hugging Face Hub and Ollama registry sources
- Optional compression for the metadata files, where it actually pays
- Optional compression for metadata *files* in transit, where it actually pays

## License

Expand Down
33 changes: 25 additions & 8 deletions internal/protocol/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -14,12 +14,13 @@ import (

// Client drives a remote helper over a pair of pipes.
type Client struct {
r *bufio.Reader
w *bufio.Writer
closer io.Closer
tool string
warn func(format string, args ...any)
root string
r *bufio.Reader
w *bufio.Writer
closer io.Closer
tool string
warn func(format string, args ...any)
root string
manifestGzip bool
}

// ClientOptions configures a Client.
Expand All @@ -40,7 +41,12 @@ func NewClient(r io.Reader, w io.Writer, closer io.Closer, opt ClientOptions) (*
warn: opt.Warn,
root: opt.Root,
}
if err := c.send(MsgHello, Hello{Version: Version, Tool: opt.Tool, Root: opt.Root}); err != nil {
if err := c.send(MsgHello, Hello{
Version: Version,
Tool: opt.Tool,
Root: opt.Root,
Features: []string{FeatureManifestGzip},
}); err != nil {
return nil, err
}
payload, err := c.expect(MsgHelloAck)
Expand All @@ -55,6 +61,7 @@ func NewClient(r io.Reader, w io.Writer, closer io.Closer, opt ClientOptions) (*
return nil, fmt.Errorf("protocol: remote speaks version %d, this build speaks %d", ack.Version, Version)
}
c.root = ack.Root
c.manifestGzip = hasFeature(ack.Features, FeatureManifestGzip)
return c, nil
}

Expand Down Expand Up @@ -113,7 +120,17 @@ func (c *Client) Plan(ctx context.Context, m *manifest.Manifest, req PlanRequest
if err := manifest.EncodeBinary(&buf, m); err != nil {
return nil, err
}
if err := WriteFrame(c.w, MsgManifest, buf.Bytes()); err != nil {
payload := buf.Bytes()
typ := MsgManifest
if c.manifestGzip {
gz, err := gzipBytes(payload)
if err != nil {
return nil, err
}
payload = gz
typ = MsgManifestGzip
}
if err := WriteFrame(c.w, typ, payload); err != nil {
return nil, err
}
if err := c.w.Flush(); err != nil {
Expand Down
36 changes: 36 additions & 0 deletions internal/protocol/gzip.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,36 @@
package protocol

import (
"bytes"
"compress/gzip"
"fmt"
"io"
)

func gzipBytes(p []byte) ([]byte, error) {
var buf bytes.Buffer
w, err := gzip.NewWriterLevel(&buf, gzip.BestSpeed)
if err != nil {
return nil, err
}
if _, err := w.Write(p); err != nil {
return nil, err
}
if err := w.Close(); err != nil {
return nil, err
}
return buf.Bytes(), nil
}

func gunzipBytes(p []byte) ([]byte, error) {
r, err := gzip.NewReader(bytes.NewReader(p))
if err != nil {
return nil, fmt.Errorf("protocol: gzip manifest: %w", err)
}
defer r.Close()
out, err := io.ReadAll(r)
if err != nil {
return nil, fmt.Errorf("protocol: gzip manifest: %w", err)
}
return out, nil
}
53 changes: 36 additions & 17 deletions internal/protocol/protocol.go
Original file line number Diff line number Diff line change
Expand Up @@ -29,19 +29,25 @@ type MsgType uint8
const (
MsgHello MsgType = iota + 1
MsgHelloAck
MsgPlanReq // JSON PlanRequest
MsgManifest // binary manifest
MsgPlanResp // JSON receiver.Plan
MsgFileBegin // JSON FileBegin
MsgChunk // 32-byte digest followed by the payload
MsgFileEnd // empty
MsgFileResult // JSON receiver.FileResult
MsgFinish // empty
MsgSummary // JSON receiver.Summary
MsgError // JSON Error
MsgWarn // JSON Warn
MsgPlanReq // JSON PlanRequest
MsgManifest // binary manifest
MsgPlanResp // JSON receiver.Plan
MsgFileBegin // JSON FileBegin
MsgChunk // 32-byte digest followed by the payload
MsgFileEnd // empty
MsgFileResult // JSON receiver.FileResult
MsgFinish // empty
MsgSummary // JSON receiver.Summary
MsgError // JSON Error
MsgWarn // JSON Warn
MsgManifestGzip // gzip-compressed binary manifest (negotiated)
)

// FeatureManifestGzip is advertised in Hello/HelloAck when both ends can
// compress the planning-phase manifest. Version stays 1 so older helpers
// still work; they simply omit the feature and receive MsgManifest.
const FeatureManifestGzip = "manifest-gzip"

func (t MsgType) String() string {
switch t {
case MsgHello:
Expand All @@ -52,6 +58,8 @@ func (t MsgType) String() string {
return "plan-req"
case MsgManifest:
return "manifest"
case MsgManifestGzip:
return "manifest-gzip"
case MsgPlanResp:
return "plan-resp"
case MsgFileBegin:
Expand All @@ -77,16 +85,27 @@ func (t MsgType) String() string {

// Hello is the first frame the client sends.
type Hello struct {
Version int `json:"version"`
Tool string `json:"tool"`
Root string `json:"root"`
Version int `json:"version"`
Tool string `json:"tool"`
Root string `json:"root"`
Features []string `json:"features,omitempty"`
}

// HelloAck is the helper's reply.
type HelloAck struct {
Version int `json:"version"`
Tool string `json:"tool"`
Root string `json:"root"`
Version int `json:"version"`
Tool string `json:"tool"`
Root string `json:"root"`
Features []string `json:"features,omitempty"`
}

func hasFeature(features []string, name string) bool {
for _, f := range features {
if f == name {
return true
}
}
return false
}

// PlanRequest carries the wire-safe subset of the receiver options.
Expand Down
55 changes: 54 additions & 1 deletion internal/protocol/protocol_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -149,7 +149,7 @@ func TestErrorFrame(t *testing.T) {

func TestMsgTypeString(t *testing.T) {
for _, m := range []MsgType{MsgHello, MsgHelloAck, MsgPlanReq, MsgManifest, MsgPlanResp,
MsgFileBegin, MsgChunk, MsgFileEnd, MsgFileResult, MsgFinish, MsgSummary, MsgError, MsgWarn} {
MsgFileBegin, MsgChunk, MsgFileEnd, MsgFileResult, MsgFinish, MsgSummary, MsgError, MsgWarn, MsgManifestGzip} {
if m.String() == "" {
t.Errorf("MsgType(%d) has no name", m)
}
Expand Down Expand Up @@ -513,6 +513,59 @@ func TestCancelledContextStopsClient(t *testing.T) {
}
}

func TestGzipManifestRoundTrip(t *testing.T) {
src := sourceDir(t)
m, err := scan.Build(context.Background(), scan.Options{Root: src, Tool: "test"})
if err != nil {
t.Fatal(err)
}
var raw bytes.Buffer
if err := manifest.EncodeBinary(&raw, m); err != nil {
t.Fatal(err)
}
gz, err := gzipBytes(raw.Bytes())
if err != nil {
t.Fatal(err)
}
if len(gz) == 0 || bytes.Equal(gz, raw.Bytes()) {
t.Fatal("gzip did not change the payload")
}
got, err := gunzipBytes(gz)
if err != nil {
t.Fatal(err)
}
if !bytes.Equal(got, raw.Bytes()) {
t.Fatal("gunzip did not restore the manifest bytes")
}
}

func TestManifestGzipOverPipes(t *testing.T) {
src := sourceDir(t)
dst := filepath.Join(t.TempDir(), "out")
m, err := scan.Build(context.Background(), scan.Options{Root: src, Tool: "test"})
if err != nil {
t.Fatal(err)
}
client, done := pipePair(t, ServerOptions{Root: dst, Tool: "modelmove/test"})
if !client.manifestGzip {
done()
t.Fatal("handshake did not negotiate manifest-gzip")
}
runTransfer(t, client, m, src, defaultRequest())
done()
want, err := os.ReadFile(filepath.Join(src, "config.json"))
if err != nil {
t.Fatal(err)
}
got, err := os.ReadFile(filepath.Join(dst, "config.json"))
if err != nil {
t.Fatal(err)
}
if !bytes.Equal(want, got) {
t.Fatal("gzip plan transfer did not write the source file")
}
}

func TestResumeOverPipes(t *testing.T) {
src := t.TempDir()
payload := randomBytes(400<<10, 6)
Expand Down
21 changes: 17 additions & 4 deletions internal/protocol/server.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package protocol

import (
"bufio"
"bytes"
"context"
"errors"
"fmt"
Expand Down Expand Up @@ -73,7 +74,12 @@ func (s *server) run(ctx context.Context) error {
if hello.Version != Version {
return fmt.Errorf("protocol: client speaks version %d, this helper speaks %d", hello.Version, Version)
}
if err := s.reply(MsgHelloAck, HelloAck{Version: Version, Tool: s.opt.Tool, Root: root}); err != nil {
if err := s.reply(MsgHelloAck, HelloAck{
Version: Version,
Tool: s.opt.Tool,
Root: root,
Features: []string{FeatureManifestGzip},
}); err != nil {
return err
}

Expand Down Expand Up @@ -141,14 +147,21 @@ func (s *server) handlePlan(ctx context.Context, payload []byte) error {
return fmt.Errorf("protocol: this helper was not started with --allow-delete")
}

t, n, err := ReadHeader(s.r)
t, raw, err := ReadFrame(s.r)
if err != nil {
return err
}
if t != MsgManifest {
switch t {
case MsgManifestGzip:
raw, err = gunzipBytes(raw)
if err != nil {
return err
}
case MsgManifest:
default:
return fmt.Errorf("protocol: expected manifest, got %s", t)
}
m, err := manifest.DecodeBinary(io.LimitReader(s.r, int64(n)))
m, err := manifest.DecodeBinary(bytes.NewReader(raw))
if err != nil {
return err
}
Expand Down