Files
echolot/server/internal/session/session.go
T
mrambossekandClaude Fable 5 3333788d9e
server-release / image (push) Successful in 15s
server-test / test (push) Successful in 32s
server-release / release (push) Successful in 32s
throughput: paced downstream rate, with the qualifier that makes it honest
A throughput number reports the smallest limit on the path, and the sender's own
ceiling is one of the candidates. If the server was asked for 50 Mbps and 50
Mbps arrived, the network was never the constraint and "50 Mbps" says nothing
about it. So the result always carries limited_by and measures_network, and a
finding is raised only when the path is actually implicated.

Loss is computed against the *sender's* count, not the requested rate: the
server reports what it put on the wire, and the gap is the loss. A receiver
alone cannot tell "the network dropped it" from "the sender never sent it", and
guessing turns a healthy server-side limit into a phantom network fault. The
count is stored per action, not per packet — half a million packets of structs
would turn a measurement into memory exhaustion.

Sending is paced rather than flat out. An unpaced burst measures the server's
NIC and the first queue it meets, then collapses into loss that reads as a
network fault. The schedule is absolute rather than sleep-per-packet, which
would accumulate scheduler error and drift the rate down over a ten-second run.

Throughput gets its own grant budget sized from the request, so every other
action stays bounded at 8 MiB. When the byte cap binds before the clock does,
the *duration* is shortened and reported, rather than the run being truncated
halfway: promising thirty seconds and delivering twenty-one is the same
information with a surprise attached, and it keeps "the clock ended the run" as
the normal case — the only case where the rate is a clean property of the path.

That last behaviour came out of a test that failed honestly: 30 s at 100 Mbps
needs 375 MB against a 256 MB cap.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-08-01 14:05:23 +02:00

255 lines
7.7 KiB
Go

// SPDX-FileCopyrightText: 2026 Echolot contributors
// SPDX-License-Identifier: GPL-3.0-or-later
// Package session implements spec §2.4 sessions and the §2.4 key schedule:
// HKDF-SHA256(ikm=device_credential, salt=key_salt, info="echolot-v1/"+session_id).
package session
import (
"crypto/hkdf"
"crypto/rand"
"crypto/sha256"
"encoding/hex"
"net/netip"
"sync"
"time"
)
type Session struct {
ID string // opaque hex; first 8 bytes (16 hex chars) are the wire prefix
Key [32]byte // derived, never crosses the wire
Epoch time.Time // server wall clock at creation (spec: RFC3339 in the response)
Expires time.Time
Device string // device ID
// Observed source of the session-creating request — the ONLY address
// reflected/generated traffic may target (spec §2.5).
ControlSource netip.Addr
// Last data-plane source seen with a valid HMAC (NAT rebinding evidence).
mu sync.Mutex
dataSource netip.AddrPort
dataLocal netip.AddrPort
// Replay window (spec §3.1: 1024-wide seq window). Highest seq seen plus
// a bitmask of the 1024 preceding.
maxSeq uint32
window [16]uint64
// Active asymmetric grants (spec §3.4/§5) — the only licence to send more than we receive.
grants []*Grant
// Observations (spec §6): per-packet UDP view + connect-back results.
packetsSeen uint64
udpObs []UDPObservation // ring, newest last, cap obsCap
connectBack []ConnectBackResult
throughput []ThroughputReport
}
const obsCap = 4096
// UDPObservation is the server's witnessed view of one data-plane packet.
type UDPObservation struct {
Seq uint32 `json:"seq"`
TRxNs int64 `json:"t_rx_ns"`
TTxNs int64 `json:"t_tx_ns"`
Src string `json:"src"`
Size int `json:"size"`
Type uint8 `json:"type"`
}
// ThroughputReport is the server's own account of a sustained send: what it managed to put on
// the wire, and what stopped it. The client needs this to interpret its own count — the gap
// between the two IS the loss, and without the sender's number a receiver can only guess.
type ThroughputReport struct {
ActionID string `json:"action_id"`
Packets int `json:"packets"`
Bytes int64 `json:"bytes"`
DurationMs int64 `json:"duration_ms"`
Kbps int `json:"kbps"`
LimitedBy string `json:"limited_by"`
}
// ConnectBackResult records one connect-back action outcome.
type ConnectBackResult struct {
ActionID string `json:"action_id"`
Result string `json:"result"` // connected | refused | timeout
RttMs float64 `json:"rtt_ms"`
}
// RecordUDP appends a packet observation (ring-capped).
func (s *Session) RecordUDP(o UDPObservation) {
s.mu.Lock()
defer s.mu.Unlock()
s.packetsSeen++
if len(s.udpObs) >= obsCap {
s.udpObs = s.udpObs[1:]
}
s.udpObs = append(s.udpObs, o)
}
// RecordConnectBack appends a connect-back outcome.
func (s *Session) RecordConnectBack(r ConnectBackResult) {
s.mu.Lock()
defer s.mu.Unlock()
s.connectBack = append(s.connectBack, r)
}
// Observations returns a copy of everything witnessed so far.
func (s *Session) Observations() (packetsSeen uint64, udp []UDPObservation, cb []ConnectBackResult) {
s.mu.Lock()
defer s.mu.Unlock()
return s.packetsSeen, append([]UDPObservation(nil), s.udpObs...),
append([]ConnectBackResult(nil), s.connectBack...)
}
// RecordThroughput stores the server's account of one sustained send.
//
// Kept as a per-action summary rather than per-packet records: a ten-second run at 50 Mbps is
// half a million packets, and holding one struct each would turn a measurement into a memory
// exhaustion. The client has the per-packet view; the server only needs to say how many it sent.
func (s *Session) RecordThroughput(actionID string, packets int, bytes, durationMs int64, kbps int, limitedBy string) {
s.mu.Lock()
defer s.mu.Unlock()
s.throughput = append(s.throughput, ThroughputReport{
ActionID: actionID, Packets: packets, Bytes: bytes,
DurationMs: durationMs, Kbps: kbps, LimitedBy: limitedBy,
})
}
// ThroughputReports returns the server's account of every sustained send in this session.
func (s *Session) ThroughputReports() []ThroughputReport {
s.mu.Lock()
defer s.mu.Unlock()
return append([]ThroughputReport(nil), s.throughput...)
}
// DataSource returns the last verified data-plane source (invalid when the
// session has not sent data-plane traffic yet).
func (s *Session) DataSource() netip.AddrPort {
s.mu.Lock()
defer s.mu.Unlock()
return s.dataSource
}
// KeySalt returns nothing — the salt is not retained after derivation; it is
// generated in New and returned once for the response body.
type Manager struct {
mu sync.Mutex
byPrefix map[string]*Session // key: first 16 hex chars of ID
ttl time.Duration
}
func NewManager(ttl time.Duration) *Manager {
return &Manager{byPrefix: map[string]*Session{}, ttl: ttl}
}
// New creates a session for a device credential per the spec key schedule.
// Returns the session and the one-time key_salt for the response.
func (m *Manager) New(deviceID, credential string, controlSource netip.Addr) (*Session, []byte, error) {
idBytes := make([]byte, 16)
if _, err := rand.Read(idBytes); err != nil {
return nil, nil, err
}
salt := make([]byte, 16)
if _, err := rand.Read(salt); err != nil {
return nil, nil, err
}
id := hex.EncodeToString(idBytes)
key, err := hkdf.Key(sha256.New, []byte(credential), salt, "echolot-v1/"+id, 32)
if err != nil {
return nil, nil, err
}
s := &Session{
ID: id,
Epoch: time.Now().UTC(),
Expires: time.Now().Add(m.ttl),
Device: deviceID,
ControlSource: controlSource,
}
copy(s.Key[:], key)
m.mu.Lock()
m.byPrefix[id[:16]] = s
m.mu.Unlock()
return s, salt, nil
}
// ByWirePrefix resolves the 8-byte on-the-wire prefix (as raw bytes).
func (m *Manager) ByWirePrefix(prefix [8]byte) *Session {
m.mu.Lock()
defer m.mu.Unlock()
s := m.byPrefix[hex.EncodeToString(prefix[:])]
if s == nil || time.Now().After(s.Expires) {
return nil
}
return s
}
// ByID resolves a full session id (sessions are keyed by their wire prefix).
func (m *Manager) ByID(id string) *Session {
if len(id) < 16 {
return nil
}
m.mu.Lock()
defer m.mu.Unlock()
s := m.byPrefix[id[:16]]
if s == nil || s.ID != id || time.Now().After(s.Expires) {
return nil
}
return s
}
func (m *Manager) Delete(id string) {
m.mu.Lock()
defer m.mu.Unlock()
delete(m.byPrefix, id[:16])
}
// CheckSeq enforces the 1024-wide anti-replay window. Returns false for
// replays and for packets older than the window.
func (s *Session) CheckSeq(seq uint32) bool {
s.mu.Lock()
defer s.mu.Unlock()
switch {
case seq > s.maxSeq:
shift := seq - s.maxSeq
for i := uint32(0); i < shift && i < 1024; i++ {
idx := (s.maxSeq + 1 + i) % 1024
s.window[idx/64] &^= 1 << (idx % 64)
}
s.maxSeq = seq
case s.maxSeq-seq >= 1024:
return false
}
idx := seq % 1024
if s.window[idx/64]&(1<<(idx%64)) != 0 {
return false
}
s.window[idx/64] |= 1 << (idx % 64)
return true
}
// DataLocal returns the server-side address that received this session's data-plane traffic.
//
// This matters more than it looks: a server bound to several addresses must send granted traffic
// back from the one the client has been talking to. Any stateful firewall or NAT in between has
// a mapping keyed on that exact pair, and a reply from a sibling address is dropped — which the
// client would then measure as downstream loss. See connFor.
func (s *Session) DataLocal() netip.AddrPort {
s.mu.Lock()
defer s.mu.Unlock()
return s.dataLocal
}
// NoteDataLocal records which of our own bound addresses saw this session's traffic.
func (s *Session) NoteDataLocal(ap netip.AddrPort) {
s.mu.Lock()
defer s.mu.Unlock()
s.dataLocal = ap
}
// NoteDataSource records the latest verified data-plane source.
func (s *Session) NoteDataSource(ap netip.AddrPort) {
s.mu.Lock()
s.dataSource = ap
s.mu.Unlock()
}