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>
255 lines
7.7 KiB
Go
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()
|
|
}
|