Skip to content
Open
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
19 changes: 16 additions & 3 deletions go.mod
Original file line number Diff line number Diff line change
@@ -1,7 +1,20 @@
module github.com/canonical/aproxy

go 1.24.0
go 1.25.0

Comment on lines +3 to 4
toolchain go1.24.1
require (
github.com/prometheus/client_golang v1.24.1
golang.org/x/crypto v0.54.0
)

require golang.org/x/crypto v0.45.0
require (
github.com/beorn7/perks v1.0.1 // indirect
github.com/cespare/xxhash/v2 v2.3.0 // indirect
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect
github.com/prometheus/client_model v0.6.2 // indirect
github.com/prometheus/common v0.70.1 // indirect
github.com/prometheus/procfs v0.21.1 // indirect
go.yaml.in/yaml/v2 v2.4.4 // indirect
golang.org/x/sys v0.47.0 // indirect
google.golang.org/protobuf v1.36.11 // indirect
)
35 changes: 35 additions & 0 deletions go.sum
Original file line number Diff line number Diff line change
@@ -1,2 +1,37 @@
github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM=
github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6rlkpw=
github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs=
github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs=
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 h1:C3w9PqII01/Oq1c1nUAm88MOHcQC9l5mIlSMApZMrHA=
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822/go.mod h1:+n7T8mK8HuQTcFwEeznm/DIxMOiR9yIdICNftLE1DvQ=
github.com/prometheus/client_golang v1.23.2 h1:Je96obch5RDVy3FDMndoUsjAhG5Edi49h0RJWRi/o0o=
github.com/prometheus/client_golang v1.23.2/go.mod h1:Tb1a6LWHB3/SPIzCoaDXI4I8UHKeFTEQ1YCr+0Gyqmg=
github.com/prometheus/client_golang v1.24.1 h1:JnJkREXzWxUdCuPFpIWZiPispT9xVV59uiuyR2bPlnU=
github.com/prometheus/client_golang v1.24.1/go.mod h1:F+oSRECHg4sse5ucfYpYDeIv/hu68Zo0uoHKetWnzcE=
github.com/prometheus/client_model v0.6.2 h1:oBsgwpGs7iVziMvrGhE53c/GrLUsZdHnqNwqPLxwZyk=
github.com/prometheus/client_model v0.6.2/go.mod h1:y3m2F6Gdpfy6Ut/GBsUqTWZqCUvMVzSfMLjcu6wAwpE=
github.com/prometheus/common v0.66.1 h1:h5E0h5/Y8niHc5DlaLlWLArTQI7tMrsfQjHV+d9ZoGs=
github.com/prometheus/common v0.66.1/go.mod h1:gcaUsgf3KfRSwHY4dIMXLPV0K/Wg1oZ8+SbZk/HH/dA=
github.com/prometheus/common v0.70.1 h1:1HvjP4D5oL3t8RsPlwxA9onvvStjtIHYE5XuuwOi/PY=
github.com/prometheus/common v0.70.1/go.mod h1:VdFUQDMZK3VLkurFUVhia6uys/0suUp86TJz5qbJRhc=
github.com/prometheus/procfs v0.16.1 h1:hZ15bTNuirocR6u0JZ6BAHHmwS1p8B4P6MRqxtzMyRg=
github.com/prometheus/procfs v0.16.1/go.mod h1:teAbpZRB1iIAJYREa1LsoWUXykVXA1KlTmWl8x/U+Is=
github.com/prometheus/procfs v0.21.1 h1:GljZCt+zSTS+NZq88cyQ1LjZ+RCHp3uVuabBWA5+OJI=
github.com/prometheus/procfs v0.21.1/go.mod h1:aB55Cww9pdSJVHk0hUf0inxWyyjPogFIjmHKYgMKmtY=
go.yaml.in/yaml/v2 v2.4.2 h1:DzmwEr2rDGHl7lsFgAHxmNz/1NlQ7xLIrlN2h5d1eGI=
go.yaml.in/yaml/v2 v2.4.2/go.mod h1:081UH+NErpNdqlCXm3TtEran0rJZGxAYx9hb/ELlsPU=
go.yaml.in/yaml/v2 v2.4.4 h1:tuyd0P+2Ont/d6e2rl3be67goVK4R6deVxCUX5vyPaQ=
go.yaml.in/yaml/v2 v2.4.4/go.mod h1:gMZqIpDtDqOfM0uNfy0SkpRhvUryYH0Z6wdMYcacYXQ=
golang.org/x/crypto v0.45.0 h1:jMBrvKuj23MTlT0bQEOBcAE0mjg8mK9RXFhRH6nyF3Q=
golang.org/x/crypto v0.45.0/go.mod h1:XTGrrkGJve7CYK7J8PEww4aY7gM3qMCElcJQ8n8JdX4=
golang.org/x/crypto v0.54.0 h1:YLIA59K4fiNzHzjnZt2tUJQjQtUWfWbeHBqKtk3eScw=
golang.org/x/crypto v0.54.0/go.mod h1:KWL8ny2AZdGR2cWmzeHrp2azQPGogOv+HeQaVEXC2dk=
golang.org/x/sys v0.38.0 h1:3yZWxaJjBmCWXqhN1qh02AkOnCQ1poK6oF+a7xWL6Gc=
golang.org/x/sys v0.38.0/go.mod h1:OgkHotnGiDImocRcuBABYBEXf8A9a87e/uXjp9XT3ks=
golang.org/x/sys v0.47.0 h1:o7XGOvZQCADBQQ4Y7VNq2dRWQR7JmOUW8Kxx4ZsNgWs=
golang.org/x/sys v0.47.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
google.golang.org/protobuf v1.36.8 h1:xHScyCOEuuwZEc6UtSOvPbAT4zRh0xcNRYekJwfqyMc=
google.golang.org/protobuf v1.36.8/go.mod h1:fuxRtAxBytpl4zzqUh6/eyUujkJdNiuEkXntxiD/uRU=
google.golang.org/protobuf v1.36.11 h1:fV6ZwhNocDyBLK0dj+fg8ektcVegBBuEolpbTQyBNVE=
google.golang.org/protobuf v1.36.11/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco=
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
179 changes: 179 additions & 0 deletions internal/conn/conn.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,179 @@
// Package conn provides a TCP connection wrapper that adds preread functionality.
// During the preread phase, the read operations are recorded and can be reverted
// using the Rewind function and the same data can be read again from the beginning.
package conn

import (
"errors"
"io"
"math"
"net"
"reflect"
"sync"
"sync/atomic"
)

const defaultPrereadLimit = 64 * 1024

var ErrPrereadLimitExceeded = errors.New("exceed preread limit")

type ConnMetrics interface {
ObserveWrite(n int, err error)
ObserveRead(n int, err error)
}

type Conn struct {
*net.TCPConn

prereadLimit int
prereadEnd atomic.Bool
prereadCursor int
prereadBuf []byte

metrics ConnMetrics

mu sync.Mutex

// Connection attributes which can be set by the user
// The conn package doesn't use these attributes
Source net.Addr
Destination net.Addr
OriginalDestination net.Addr
Host string
}

func (c *Conn) read(b []byte) (n int, err error) {
n, err = c.TCPConn.Read(b)
if c.metrics != nil {
c.metrics.ObserveRead(n, err)
}
return n, err
}

// Read reads from the connection.
// During preread phase, data read from the underlying TCP connection are recorded
// and can be read again as if it hasn't been read before using the Rewind function.
// A zero-length read during the preread phase just returns without touching the
// underlying TCP connection.
// Conn starts with preread enabled and the preread phase can be ended with the
// EndPreread method.
func (c *Conn) Read(b []byte) (n int, err error) {
c.mu.Lock()
defer c.mu.Unlock()
prereadEnd := c.prereadEnd.Load()
limit := c.prereadLimit - c.prereadCursor
if prereadEnd {
limit = math.MaxInt
}
rb := b[:min(limit, len(b))]
if len(rb) == 0 && len(b) != 0 {
return 0, ErrPrereadLimitExceeded
}
reuse := min(len(rb), max(len(c.prereadBuf)-c.prereadCursor, 0))
copy(rb, c.prereadBuf[c.prereadCursor:c.prereadCursor+reuse])
c.prereadCursor += reuse
if reuse > 0 || (len(b) == 0 && !prereadEnd) {
return reuse, nil
}
n, err = c.read(rb)
if !prereadEnd {
c.prereadBuf = append(c.prereadBuf, rb[:n]...)
c.prereadCursor += n
} else {
if c.prereadCursor >= len(c.prereadBuf) {
c.prereadCursor = 0
c.prereadBuf = nil
}
}
return n + reuse, err
}

// ReadFrom is not supported.
func (c *Conn) ReadFrom(r io.Reader) (n int64, err error) {
panic("not supported")
}

func (c *Conn) write(b []byte) (n int, err error) {
n, err = c.TCPConn.Write(b)
if c.metrics != nil {
c.metrics.ObserveWrite(n, err)
}
return n, err
}

// Write writes to the connection.
// Write is not allowed during the preread phase.
func (c *Conn) Write(b []byte) (n int, err error) {
if !c.prereadEnd.Load() {
panic("write while in preread")
}
return c.write(b)
}

// WriteTo is not supported.
func (c *Conn) WriteTo(w io.Writer) (n int64, err error) {
panic("not supported")
}

// Rewind undoes the reads performed so far, and allows the connection to be
// read from the beginning again.
// Rewind will panic if called after the preread phase.
func (c *Conn) Rewind() {
c.mu.Lock()
defer c.mu.Unlock()
if c.prereadEnd.Load() {
panic("attempt to rewind while preread ended")
}
c.prereadCursor = 0
}

// EndPreread ends the preread phase and rewinds the connection.
// EndPreread can only be called once.
func (c *Conn) EndPreread() {
c.mu.Lock()
defer c.mu.Unlock()
if c.prereadEnd.Load() {
return
}
c.prereadEnd.Store(true)
c.prereadCursor = 0
}

type Option func(*Conn)

// WithPrereadLimit sets a limit on how many bytes can be read during the preread phase.
// A Read operation exceeding the preread limit will cause an ErrPrereadLimitExceeded error.
func WithPrereadLimit(n int) Option {
if n <= 0 {
panic("invalid preread limit")
}
return func(c *Conn) {
c.prereadLimit = n
}
}

// WithMetrics reports socket activity to observability.
func WithMetrics(m ConnMetrics) Option {
if m == nil {
panic("nil metrics")
}
if v := reflect.ValueOf(m); v.Kind() == reflect.Pointer && v.IsNil() {
panic("nil metrics")
}
return func(c *Conn) {
c.metrics = m
}
}

// New wraps the input TCP connection and returns a connection started in the
// preread phase.
func New(c *net.TCPConn, options ...Option) *Conn {
conn := &Conn{
TCPConn: c,
prereadLimit: defaultPrereadLimit,
}
for _, option := range options {
option(conn)
}
return conn
}
Loading
Loading