mirror of
https://github.com/soypat/lneto.git
synced 2026-09-03 21:39:04 +00:00
restructure package, rename module/repo
This commit is contained in:
+536
@@ -0,0 +1,536 @@
|
||||
package tcp
|
||||
|
||||
import (
|
||||
"io"
|
||||
"log/slog"
|
||||
"math"
|
||||
"net"
|
||||
|
||||
"github.com/soypat/lneto/internal"
|
||||
)
|
||||
|
||||
// ControlBlock is a partial Transmission Control Block (TCB) implementation as
|
||||
// per RFC 9293 in section 3.3.1. In contrast with the description in RFC9293,
|
||||
// this implementation is limited to receiving only sequential segments.
|
||||
// This means buffer management is left up entirely to the user of the ControlBlock.
|
||||
// Use ControlBlock as the building block that solves Sequence Number calculation
|
||||
// and validation in a full TCP implementation.
|
||||
//
|
||||
// A ControlBlock's internal state is modified by the available "System Calls" as defined in
|
||||
// RFC9293, such as Close, Listen/Open, Send, and Receive.
|
||||
// Sent and received data is represented with the [Segment] struct type.
|
||||
type ControlBlock struct {
|
||||
// # Send Sequence Space
|
||||
//
|
||||
// 'Send' sequence numbers correspond to local data being sent.
|
||||
//
|
||||
// 1 2 3 4
|
||||
// ----------|----------|----------|----------
|
||||
// SND.UNA SND.NXT SND.UNA
|
||||
// +SND.WND
|
||||
// 1. old sequence numbers which have been acknowledged
|
||||
// 2. sequence numbers of unacknowledged data
|
||||
// 3. sequence numbers allowed for new data transmission
|
||||
// 4. future sequence numbers which are not yet allowed
|
||||
snd sendSpace
|
||||
// # Receive Sequence Space
|
||||
//
|
||||
// 'Receive' sequence numbers correspond to remote data being received.
|
||||
//
|
||||
// 1 2 3
|
||||
// ----------|----------|----------
|
||||
// RCV.NXT RCV.NXT
|
||||
// +RCV.WND
|
||||
// 1 - old sequence numbers which have been acknowledged
|
||||
// 2 - sequence numbers allowed for new reception
|
||||
// 3 - future sequence numbers which are not yet allowed
|
||||
rcv recvSpace
|
||||
// When FlagRST is set in pending flags rstPtr will contain the sequence number of the RST segment to make it "believable" (See RFC9293)
|
||||
rstPtr Value
|
||||
// pending is the queue of pending flags to be sent in the next 2 segments.
|
||||
// On a call to Send the queue is advanced and flags set in the segment are unset.
|
||||
// The second position of the queue is used for FIN segments.
|
||||
pending [2]Flags
|
||||
state State
|
||||
challengeAck bool
|
||||
log *slog.Logger
|
||||
}
|
||||
|
||||
// State returns the current state of the TCP connection.
|
||||
func (tcb *ControlBlock) State() State { return tcb.state }
|
||||
|
||||
// RecvNext returns the next sequence number expected to be received from remote.
|
||||
// This implementation will reject segments that are not the next expected sequence.
|
||||
// RecvNext returns 0 before StateSynRcvd.
|
||||
func (tcb *ControlBlock) RecvNext() Value { return tcb.rcv.NXT }
|
||||
|
||||
// RecvWindow returns the receive window size. If connection is closed will return 0.
|
||||
func (tcb *ControlBlock) RecvWindow() Size { return tcb.rcv.WND }
|
||||
|
||||
// ISS returns the initial sequence number of the connection that was defined on a call to Open by user.
|
||||
func (tcb *ControlBlock) ISS() Value { return tcb.snd.ISS }
|
||||
|
||||
// MaxInFlightData returns the maximum size of a segment that can be sent by taking into account
|
||||
// the send window size and the unacked data. Returns 0 before StateSynRcvd.
|
||||
func (tcb *ControlBlock) MaxInFlightData() Size {
|
||||
if !tcb.state.hasIRS() {
|
||||
return 0 // SYN not yet received.
|
||||
}
|
||||
unacked := Sizeof(tcb.snd.UNA, tcb.snd.NXT)
|
||||
return tcb.snd.WND - unacked - 1 // TODO: is this -1 supposed to be here?
|
||||
}
|
||||
|
||||
// SetWindow sets the local receive window size. This represents the maximum amount of data
|
||||
// that is permitted to be in flight.
|
||||
func (tcb *ControlBlock) SetRecvWindow(wnd Size) {
|
||||
tcb.rcv.WND = wnd
|
||||
}
|
||||
|
||||
// SetLogger sets the logger to be used by the ControlBlock.
|
||||
func (tcb *ControlBlock) SetLogger(log *slog.Logger) {
|
||||
tcb.log = log
|
||||
}
|
||||
|
||||
// IncomingIsKeepalive checks if an incoming segment is a keepalive segment.
|
||||
// Segments which are keepalives should not be passed into Recv or Send methods.
|
||||
func (tcb *ControlBlock) IncomingIsKeepalive(incomingSegment Segment) bool {
|
||||
return incomingSegment.SEQ == tcb.rcv.NXT-1 &&
|
||||
incomingSegment.Flags == FlagACK &&
|
||||
incomingSegment.ACK == tcb.snd.NXT && incomingSegment.DATALEN == 0
|
||||
}
|
||||
|
||||
// MakeKeepalive creates a TCP keepalive segment. This segment
|
||||
// should not be passed into Recv or Send methods.
|
||||
func (tcb *ControlBlock) MakeKeepalive() Segment {
|
||||
return Segment{
|
||||
SEQ: tcb.snd.NXT - 1,
|
||||
ACK: tcb.rcv.NXT,
|
||||
Flags: FlagACK,
|
||||
WND: tcb.rcv.WND,
|
||||
DATALEN: 0,
|
||||
}
|
||||
}
|
||||
|
||||
// sendSpace contains Send Sequence Space data. Its sequence numbers correspond to local data.
|
||||
type sendSpace struct {
|
||||
ISS Value // initial send sequence number, defined locally on connection start
|
||||
UNA Value // send unacknowledged. Seqs equal to UNA and above have NOT been acked by remote. Corresponds to local data.
|
||||
NXT Value // send next. This seq and up to UNA+WND-1 are allowed to be sent. Corresponds to local data.
|
||||
WND Size // send window defined by remote. Permitted number of local unacked octets in flight.
|
||||
// WL1 Value // segment sequence number used for last window update
|
||||
// WL2 Value // segment acknowledgment number used for last window update
|
||||
}
|
||||
|
||||
// inFlight returns amount of unacked bytes sent out.
|
||||
func (snd *sendSpace) inFlight() Size {
|
||||
return Sizeof(snd.UNA, snd.NXT)
|
||||
}
|
||||
|
||||
// maxSend returns maximum segment datalength receivable by remote peer.
|
||||
func (snd *sendSpace) maxSend() Size {
|
||||
return snd.WND - snd.inFlight()
|
||||
}
|
||||
|
||||
// recvSpace contains Receive Sequence Space data. Its sequence numbers correspond to remote data.
|
||||
type recvSpace struct {
|
||||
IRS Value // initial receive sequence number, defined by remote in SYN segment received.
|
||||
NXT Value // receive next. seqs before this have been acked. this seq and up to NXT+WND-1 are allowed to be sent. Corresponds to remote data.
|
||||
WND Size // receive window defined by local. Permitted number of remote unacked octets in flight.
|
||||
}
|
||||
|
||||
// Open implements a passive/active opening of a connection.
|
||||
// state must be StateListen or StateSynSent.
|
||||
func (tcb *ControlBlock) Open(iss Value, wnd Size, state State) (err error) {
|
||||
switch {
|
||||
case tcb.state != StateClosed && tcb.state != StateListen:
|
||||
err = errTCBNotClosed
|
||||
case state != StateListen && state != StateSynSent:
|
||||
err = errInvalidState
|
||||
case wnd > math.MaxUint16:
|
||||
err = errWindowTooLarge
|
||||
}
|
||||
if err != nil {
|
||||
tcb.logerr("tcb:open", slog.String("err", err.Error()))
|
||||
return err
|
||||
}
|
||||
tcb.state = state
|
||||
tcb.resetRcv(wnd, 0)
|
||||
tcb.resetSnd(iss, 1)
|
||||
tcb.pending = [2]Flags{}
|
||||
if state == StateSynSent {
|
||||
tcb.pending[0] = FlagSYN
|
||||
}
|
||||
tcb.trace("tcb:open", slog.String("state", tcb.state.String()))
|
||||
return nil
|
||||
}
|
||||
|
||||
// HasPending returns true if there is a pending control segment to send. Calls to Send will advance the pending queue.
|
||||
func (tcb *ControlBlock) HasPending() bool { return tcb.pending[0] != 0 }
|
||||
|
||||
// PendingSegment calculates a suitable next segment to send from a payload length.
|
||||
// It does not modify the ControlBlock state or pending segment queue.
|
||||
func (tcb *ControlBlock) PendingSegment(payloadLen int) (_ Segment, ok bool) {
|
||||
if tcb.challengeAck {
|
||||
tcb.challengeAck = false
|
||||
return Segment{SEQ: tcb.snd.NXT, ACK: tcb.rcv.NXT, Flags: FlagACK, WND: tcb.rcv.WND}, true
|
||||
}
|
||||
pending := tcb.pending[0]
|
||||
established := tcb.state == StateEstablished
|
||||
if !established && tcb.state != StateCloseWait {
|
||||
payloadLen = 0 // Can't send data if not established.
|
||||
}
|
||||
if pending == 0 && payloadLen == 0 {
|
||||
return Segment{}, false // No pending segment.
|
||||
}
|
||||
|
||||
// Limit payload to what send window allows.
|
||||
inFlight := tcb.snd.inFlight()
|
||||
_ = inFlight
|
||||
maxPayload := tcb.snd.maxSend()
|
||||
if payloadLen > int(maxPayload) {
|
||||
if maxPayload == 0 && !tcb.pending[0].HasAny(FlagFIN|FlagRST|FlagSYN) {
|
||||
return Segment{}, false
|
||||
} else if maxPayload > tcb.snd.WND {
|
||||
panic("seqs: bad calculation")
|
||||
}
|
||||
payloadLen = int(maxPayload)
|
||||
}
|
||||
|
||||
if established {
|
||||
pending |= FlagACK // ACK is always set in established state. Not in RFC9293 but somehow expected?
|
||||
} else {
|
||||
payloadLen = 0 // Can't send data if not established.
|
||||
}
|
||||
|
||||
var ack Value
|
||||
if pending.HasAny(FlagACK) {
|
||||
ack = tcb.rcv.NXT
|
||||
}
|
||||
|
||||
var seq Value = tcb.snd.NXT
|
||||
if pending.HasAny(FlagRST) {
|
||||
seq = tcb.rstPtr
|
||||
}
|
||||
|
||||
seg := Segment{
|
||||
SEQ: seq,
|
||||
ACK: ack,
|
||||
WND: tcb.rcv.WND,
|
||||
Flags: pending,
|
||||
DATALEN: Size(payloadLen),
|
||||
}
|
||||
tcb.traceSeg("tcb:pending-out", seg)
|
||||
return seg, true
|
||||
}
|
||||
|
||||
// Recv processes a segment that is being received from the network. It updates the TCB
|
||||
// if there is no error. The ControlBlock can only receive segments that are the next
|
||||
// expected sequence number which means the caller must handle the out-of-order case
|
||||
// and buffering that comes with it.
|
||||
func (tcb *ControlBlock) Recv(seg Segment) (err error) {
|
||||
err = tcb.validateIncomingSegment(seg)
|
||||
if err != nil {
|
||||
tcb.traceRcv("tcb:rcv.reject")
|
||||
tcb.traceSeg("tcb:rcv.reject", seg)
|
||||
tcb.logerr("tcb:rcv.reject", slog.String("err", err.Error()))
|
||||
return err
|
||||
}
|
||||
|
||||
prevNxt := tcb.snd.NXT
|
||||
var pending Flags
|
||||
switch tcb.state {
|
||||
case StateListen:
|
||||
pending, err = tcb.rcvListen(seg)
|
||||
case StateSynSent:
|
||||
pending, err = tcb.rcvSynSent(seg)
|
||||
case StateSynRcvd:
|
||||
pending, err = tcb.rcvSynRcvd(seg)
|
||||
case StateEstablished:
|
||||
pending, err = tcb.rcvEstablished(seg)
|
||||
case StateFinWait1:
|
||||
pending, err = tcb.rcvFinWait1(seg)
|
||||
case StateFinWait2:
|
||||
pending, err = tcb.rcvFinWait2(seg)
|
||||
case StateCloseWait:
|
||||
case StateLastAck:
|
||||
if seg.Flags.HasAny(FlagACK) {
|
||||
tcb.close()
|
||||
}
|
||||
case StateClosing:
|
||||
// Thanks to @knieriem for finding and reporting this bug.
|
||||
if seg.Flags.HasAny(FlagACK) {
|
||||
tcb.state = StateTimeWait
|
||||
}
|
||||
default:
|
||||
panic("unexpected recv state:" + tcb.state.String())
|
||||
}
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
tcb.pending[0] |= pending
|
||||
if prevNxt != 0 && tcb.snd.NXT != prevNxt && tcb.logenabled(slog.LevelDebug) {
|
||||
tcb.debug("tcb:snd.nxt-change", slog.String("state", tcb.state.String()),
|
||||
slog.Uint64("seg.ack", uint64(seg.ACK)), slog.Uint64("snd.nxt", uint64(tcb.snd.NXT)),
|
||||
slog.Uint64("prevnxt", uint64(prevNxt)), slog.Uint64("seg.seq", uint64(seg.SEQ)))
|
||||
}
|
||||
|
||||
// We accept the segment and update TCB state.
|
||||
tcb.snd.WND = seg.WND
|
||||
if seg.Flags.HasAny(FlagACK) {
|
||||
tcb.snd.UNA = seg.ACK
|
||||
}
|
||||
seglen := seg.LEN()
|
||||
tcb.rcv.NXT.UpdateForward(seglen)
|
||||
|
||||
if tcb.logenabled(internal.LevelTrace) {
|
||||
tcb.traceRcv("tcb:rcv")
|
||||
tcb.traceSeg("recv:seg", seg)
|
||||
}
|
||||
return err
|
||||
}
|
||||
|
||||
// Send processes a segment that is being sent to the network. It updates the TCB
|
||||
// if there is no error.
|
||||
func (tcb *ControlBlock) Send(seg Segment) error {
|
||||
err := tcb.validateOutgoingSegment(seg)
|
||||
if err != nil {
|
||||
tcb.traceSnd("tcb:snd.reject")
|
||||
tcb.traceSeg("tcb:snd.reject", seg)
|
||||
tcb.logerr("tcb:snd.reject", slog.String("err", err.Error()))
|
||||
return err
|
||||
}
|
||||
|
||||
hasFIN := seg.Flags.HasAny(FlagFIN)
|
||||
hasACK := seg.Flags.HasAny(FlagACK)
|
||||
var newPending Flags
|
||||
switch tcb.state {
|
||||
case StateSynRcvd:
|
||||
if hasFIN {
|
||||
tcb.state = StateFinWait1 // RFC 9293: 3.10.4 CLOSE call.
|
||||
}
|
||||
case StateClosing:
|
||||
if hasACK {
|
||||
tcb.state = StateTimeWait
|
||||
}
|
||||
case StateEstablished:
|
||||
if hasFIN {
|
||||
tcb.state = StateFinWait1
|
||||
}
|
||||
case StateCloseWait:
|
||||
if hasFIN {
|
||||
tcb.state = StateLastAck
|
||||
} else if hasACK {
|
||||
newPending = finack // Queue finack.
|
||||
}
|
||||
}
|
||||
|
||||
// Advance pending flags queue.
|
||||
tcb.pending[0] &^= seg.Flags
|
||||
if tcb.pending[0] == 0 {
|
||||
// Ensure we don't queue a FINACK if we have already sent a FIN.
|
||||
tcb.pending = [2]Flags{tcb.pending[1] &^ (seg.Flags & (FlagFIN)), 0}
|
||||
}
|
||||
tcb.pending[0] |= newPending
|
||||
|
||||
// The segment is valid, we can update TCB state.
|
||||
seglen := seg.LEN()
|
||||
tcb.snd.NXT.UpdateForward(seglen)
|
||||
tcb.rcv.WND = seg.WND
|
||||
|
||||
if tcb.logenabled(internal.LevelTrace) {
|
||||
tcb.traceSnd("tcb:snd")
|
||||
tcb.traceSeg("tcb:snd", seg)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (tcb *ControlBlock) validateOutgoingSegment(seg Segment) (err error) {
|
||||
hasAck := seg.Flags.HasAny(FlagACK)
|
||||
checkSeq := !seg.Flags.HasAny(FlagRST)
|
||||
seglast := seg.Last()
|
||||
// Extra check for when send Window is zero and no data is being sent.
|
||||
zeroWindowOK := tcb.snd.WND == 0 && seg.DATALEN == 0 && seg.SEQ == tcb.snd.NXT
|
||||
outOfWindow := checkSeq && !InWindow(seg.SEQ, tcb.snd.NXT, tcb.snd.WND) &&
|
||||
!zeroWindowOK
|
||||
switch {
|
||||
case tcb.state == StateClosed:
|
||||
err = io.ErrClosedPipe
|
||||
case seg.WND > math.MaxUint16:
|
||||
err = errWindowTooLarge
|
||||
case hasAck && seg.ACK != tcb.rcv.NXT:
|
||||
err = errAckNotNext
|
||||
|
||||
case outOfWindow:
|
||||
if tcb.snd.WND == 0 {
|
||||
err = errZeroWindow
|
||||
} else {
|
||||
err = errSeqNotInWindow
|
||||
}
|
||||
|
||||
case seg.DATALEN > 0 && (tcb.state == StateFinWait1 || tcb.state == StateFinWait2):
|
||||
err = errConnectionClosing // Case 1: No further SENDs from the user will be accepted by the TCP implementation.
|
||||
|
||||
case checkSeq && tcb.snd.WND == 0 && seg.DATALEN > 0 && seg.SEQ == tcb.snd.NXT:
|
||||
err = errZeroWindow
|
||||
|
||||
case checkSeq && !InWindow(seglast, tcb.snd.NXT, tcb.snd.WND) && !zeroWindowOK:
|
||||
err = errLastNotInWindow
|
||||
}
|
||||
return err
|
||||
}
|
||||
|
||||
func (tcb *ControlBlock) validateIncomingSegment(seg Segment) (err error) {
|
||||
flags := seg.Flags
|
||||
hasAck := flags.HasAll(FlagACK)
|
||||
// Short circuit SEQ checks if SYN present since the incoming segment initialize1s connection.
|
||||
checkSEQ := !flags.HasAny(FlagSYN)
|
||||
established := tcb.state == StateEstablished
|
||||
preestablished := tcb.state.IsPreestablished()
|
||||
acksOld := hasAck && !LessThan(tcb.snd.UNA, seg.ACK)
|
||||
acksUnsentData := hasAck && !LessThanEq(seg.ACK, tcb.snd.NXT)
|
||||
ctlOrDataSegment := established && (seg.DATALEN > 0 || flags.HasAny(FlagFIN|FlagRST))
|
||||
zeroWindowOK := tcb.rcv.WND == 0 && seg.DATALEN == 0 && seg.SEQ == tcb.rcv.NXT
|
||||
// See section 3.4 of RFC 9293 for more on these checks.
|
||||
switch {
|
||||
case seg.WND > math.MaxUint16:
|
||||
err = errWindowOverflow
|
||||
case tcb.state == StateClosed:
|
||||
err = io.ErrClosedPipe
|
||||
|
||||
case checkSEQ && tcb.rcv.WND == 0 && seg.DATALEN > 0 && seg.SEQ == tcb.rcv.NXT:
|
||||
err = errZeroWindow
|
||||
|
||||
case checkSEQ && !InWindow(seg.SEQ, tcb.rcv.NXT, tcb.rcv.WND) && !zeroWindowOK:
|
||||
err = errSeqNotInWindow
|
||||
|
||||
case checkSEQ && !InWindow(seg.Last(), tcb.rcv.NXT, tcb.rcv.WND) && !zeroWindowOK:
|
||||
err = errLastNotInWindow
|
||||
|
||||
case checkSEQ && seg.SEQ != tcb.rcv.NXT:
|
||||
// This part diverts from TCB as described in RFC 9293. We want to support
|
||||
// only sequential segments to keep implementation simple and maintainable. See SHLD-31.
|
||||
err = errRequireSequential
|
||||
}
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if flags.HasAny(FlagRST) {
|
||||
return tcb.handleRST(seg.SEQ)
|
||||
}
|
||||
|
||||
isDebug := tcb.logenabled(slog.LevelDebug)
|
||||
// Drop-segment checks.
|
||||
switch {
|
||||
// Special treatment of duplicate ACKs on established connection and of ACKs of unsent data.
|
||||
// https://www.rfc-editor.org/rfc/rfc9293.html#section-3.10.7.4-2.5.2.2.2.3.2.1
|
||||
case established && acksOld && !ctlOrDataSegment:
|
||||
err = errDropSegment
|
||||
tcb.pending[0] &= FlagFIN // Completely ignore duplicate ACKs but do not erase fin bit.
|
||||
if isDebug {
|
||||
tcb.debug("rcv:ACK-dup", slog.String("state", tcb.state.String()),
|
||||
slog.Uint64("seg.ack", uint64(seg.ACK)), slog.Uint64("snd.una", uint64(tcb.snd.UNA)))
|
||||
}
|
||||
|
||||
case established && acksUnsentData:
|
||||
err = errDropSegment
|
||||
tcb.pending[0] = FlagACK // Send ACK for unsent data.
|
||||
if isDebug {
|
||||
tcb.debug("rcv:ACK-unsent", slog.String("state", tcb.state.String()),
|
||||
slog.Uint64("seg.ack", uint64(seg.ACK)), slog.Uint64("snd.nxt", uint64(tcb.snd.NXT)))
|
||||
}
|
||||
|
||||
case preestablished && (acksOld || acksUnsentData):
|
||||
err = errDropSegment
|
||||
tcb.pending[0] = FlagRST
|
||||
tcb.rstPtr = seg.ACK
|
||||
tcb.resetSnd(tcb.snd.ISS, seg.WND)
|
||||
if isDebug {
|
||||
tcb.debug("rcv:RST-old", slog.String("state", tcb.state.String()), slog.Uint64("ack", uint64(seg.ACK)))
|
||||
}
|
||||
}
|
||||
return err
|
||||
}
|
||||
|
||||
func (tcb *ControlBlock) resetSnd(localISS Value, remoteWND Size) {
|
||||
tcb.snd = sendSpace{
|
||||
ISS: localISS,
|
||||
UNA: localISS,
|
||||
NXT: localISS,
|
||||
WND: remoteWND,
|
||||
// UP, WL1, WL2 defaults to zero values.
|
||||
}
|
||||
}
|
||||
|
||||
func (tcb *ControlBlock) resetRcv(localWND Size, remoteISS Value) {
|
||||
tcb.rcv = recvSpace{
|
||||
IRS: remoteISS,
|
||||
NXT: remoteISS,
|
||||
WND: localWND,
|
||||
}
|
||||
}
|
||||
|
||||
func (tcb *ControlBlock) handleRST(seq Value) error {
|
||||
tcb.debug("rcv:RST", slog.String("state", tcb.state.String()))
|
||||
if seq != tcb.rcv.NXT {
|
||||
// See RFC9293: If the RST bit is set and the sequence number does not exactly match the next expected sequence value, yet is within the current receive window, TCP endpoints MUST send an acknowledgment (challenge ACK).
|
||||
tcb.challengeAck = true
|
||||
tcb.pending[0] |= FlagACK
|
||||
return errDropSegment
|
||||
}
|
||||
if tcb.state.IsPreestablished() {
|
||||
tcb.pending[0] = 0
|
||||
tcb.state = StateListen
|
||||
tcb.resetSnd(tcb.snd.ISS+tcb.rstJump(), tcb.snd.WND)
|
||||
tcb.resetRcv(tcb.rcv.WND, 3_14159_2653^tcb.rcv.IRS)
|
||||
} else {
|
||||
tcb.close() // Enter closed state and return.
|
||||
return net.ErrClosed
|
||||
}
|
||||
return errDropSegment
|
||||
}
|
||||
|
||||
func (tcb *ControlBlock) rstJump() Value {
|
||||
return 100
|
||||
}
|
||||
|
||||
// close sets ControlBlock state to closed and resets all sequence numbers and pending flag.
|
||||
func (tcb *ControlBlock) close() {
|
||||
tcb.state = StateClosed
|
||||
tcb.pending = [2]Flags{}
|
||||
tcb.resetRcv(0, 0)
|
||||
tcb.resetSnd(0, 0)
|
||||
tcb.debug("tcb:close")
|
||||
}
|
||||
|
||||
// Close implements a passive/active closing of a connection. It does not immediately
|
||||
// delete the TCB but initiates the process so that pending outgoing segments initiate
|
||||
// the closing process. After a call to Close users should not send more data.
|
||||
// Close returns an error if the connection is already closed or closing.
|
||||
func (tcb *ControlBlock) Close() (err error) {
|
||||
// See RFC 9293: 3.10.4 CLOSE call.
|
||||
switch tcb.state {
|
||||
case StateClosed:
|
||||
err = errConnNotexist
|
||||
case StateCloseWait:
|
||||
tcb.state = StateLastAck
|
||||
tcb.pending = [2]Flags{FlagFIN, FlagACK}
|
||||
case StateListen, StateSynSent:
|
||||
tcb.close()
|
||||
case StateSynRcvd, StateEstablished:
|
||||
// We suppose user has no more pending data to send, so we flag FIN to be sent.
|
||||
// Users of this API should call Close only when they have no more data to send.
|
||||
tcb.pending[0] = (tcb.pending[0] & FlagACK) | FlagFIN
|
||||
case StateFinWait2, StateTimeWait:
|
||||
err = errConnectionClosing
|
||||
default:
|
||||
err = errInvalidState
|
||||
}
|
||||
if err == nil {
|
||||
tcb.trace("tcb:close", slog.String("state", tcb.state.String()))
|
||||
} else {
|
||||
tcb.logerr("tcb:close", slog.String("err", err.Error()))
|
||||
}
|
||||
return err
|
||||
}
|
||||
@@ -0,0 +1,107 @@
|
||||
package tcp
|
||||
|
||||
func (tcb *ControlBlock) rcvListen(seg Segment) (pending Flags, err error) {
|
||||
switch {
|
||||
case !seg.Flags.HasAll(FlagSYN):
|
||||
err = errExpectedSYN
|
||||
}
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
// Initialize all connection state:
|
||||
tcb.resetSnd(tcb.snd.ISS, seg.WND)
|
||||
tcb.resetRcv(tcb.rcv.WND, seg.SEQ)
|
||||
|
||||
// We must respond with SYN|ACK frame after receiving SYN in listen state (three way handshake).
|
||||
tcb.pending[0] = synack
|
||||
tcb.state = StateSynRcvd
|
||||
return synack, nil
|
||||
}
|
||||
|
||||
func (tcb *ControlBlock) rcvSynSent(seg Segment) (pending Flags, err error) {
|
||||
hasSyn := seg.Flags.HasAny(FlagSYN)
|
||||
hasAck := seg.Flags.HasAny(FlagACK)
|
||||
switch {
|
||||
case !hasSyn:
|
||||
err = errExpectedSYN
|
||||
|
||||
case hasAck && seg.ACK != tcb.snd.UNA+1:
|
||||
err = errBadSegack
|
||||
}
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
|
||||
if hasAck {
|
||||
tcb.state = StateEstablished
|
||||
pending = FlagACK
|
||||
tcb.resetRcv(tcb.rcv.WND, seg.SEQ)
|
||||
} else {
|
||||
// Simultaneous connection sync edge case.
|
||||
pending = synack
|
||||
tcb.state = StateSynRcvd
|
||||
tcb.resetSnd(tcb.snd.ISS, seg.WND)
|
||||
tcb.resetRcv(tcb.rcv.WND, seg.SEQ)
|
||||
}
|
||||
return pending, nil
|
||||
}
|
||||
|
||||
func (tcb *ControlBlock) rcvSynRcvd(seg Segment) (pending Flags, err error) {
|
||||
switch {
|
||||
// case !seg.Flags.HasAll(FlagACK):
|
||||
// err = errors.New("rcvSynRcvd: expected ACK")
|
||||
case seg.ACK != tcb.snd.UNA+1:
|
||||
err = errBadSegack
|
||||
}
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
tcb.state = StateEstablished
|
||||
return 0, nil
|
||||
}
|
||||
|
||||
func (tcb *ControlBlock) rcvEstablished(seg Segment) (pending Flags, err error) {
|
||||
flags := seg.Flags
|
||||
|
||||
dataToAck := seg.DATALEN > 0
|
||||
hasFin := flags.HasAny(FlagFIN)
|
||||
if dataToAck || hasFin {
|
||||
pending = FlagACK
|
||||
if hasFin {
|
||||
// See Figure 5: TCP Connection State Diagram of RFC 9293.
|
||||
tcb.state = StateCloseWait
|
||||
tcb.pending[1] = FlagFIN // Queue FIN for after the CloseWait ACK.
|
||||
}
|
||||
}
|
||||
|
||||
return pending, nil
|
||||
}
|
||||
|
||||
func (tcb *ControlBlock) rcvFinWait1(seg Segment) (pending Flags, err error) {
|
||||
flags := seg.Flags
|
||||
hasFin := flags&FlagFIN != 0
|
||||
hasAck := flags&FlagACK != 0
|
||||
switch {
|
||||
case hasFin && hasAck && seg.ACK == tcb.snd.NXT:
|
||||
// Special case: Server sent a FINACK response to our FIN so we enter TimeWait directly.
|
||||
// We have to check ACK against send NXT to avoid simultaneous close sequence edge case.
|
||||
tcb.state = StateTimeWait
|
||||
case hasFin:
|
||||
tcb.state = StateClosing
|
||||
case hasAck:
|
||||
// TODO(soypat): Check if this branch does NOT need ACK queued. Online flowcharts say not needed.
|
||||
tcb.state = StateFinWait2
|
||||
default:
|
||||
return 0, errFinwaitExpectedACK
|
||||
}
|
||||
pending = FlagACK
|
||||
return pending, nil
|
||||
}
|
||||
|
||||
func (tcb *ControlBlock) rcvFinWait2(seg Segment) (pending Flags, err error) {
|
||||
if !seg.Flags.HasAll(finack) {
|
||||
return pending, errFinwaitExpectedFinack
|
||||
}
|
||||
tcb.state = StateTimeWait
|
||||
return FlagACK, nil
|
||||
}
|
||||
@@ -0,0 +1,59 @@
|
||||
package tcp
|
||||
|
||||
import (
|
||||
"context"
|
||||
"log/slog"
|
||||
|
||||
"github.com/soypat/lneto/internal"
|
||||
)
|
||||
|
||||
func (tcb *ControlBlock) logenabled(lvl slog.Level) bool {
|
||||
return internal.HeapAllocDebugging || (tcb.log != nil && tcb.log.Handler().Enabled(context.Background(), lvl))
|
||||
}
|
||||
|
||||
func (tcb *ControlBlock) logattrs(lvl slog.Level, msg string, attrs ...slog.Attr) {
|
||||
internal.LogAttrs(tcb.log, lvl, msg, attrs...)
|
||||
}
|
||||
|
||||
func (tcb *ControlBlock) debug(msg string, attrs ...slog.Attr) {
|
||||
tcb.logattrs(slog.LevelDebug, msg, attrs...)
|
||||
}
|
||||
|
||||
func (tcb *ControlBlock) trace(msg string, attrs ...slog.Attr) {
|
||||
tcb.logattrs(internal.LevelTrace, msg, attrs...)
|
||||
}
|
||||
|
||||
func (tcb *ControlBlock) logerr(msg string, attrs ...slog.Attr) {
|
||||
tcb.logattrs(slog.LevelError, msg, attrs...)
|
||||
}
|
||||
|
||||
func (tcb *ControlBlock) traceSnd(msg string) {
|
||||
tcb.trace(msg,
|
||||
slog.String("state", tcb.state.String()),
|
||||
slog.Uint64("pend", uint64(tcb.pending[0])),
|
||||
slog.Uint64("snd.nxt", uint64(tcb.snd.NXT)),
|
||||
slog.Uint64("snd.una", uint64(tcb.snd.UNA)),
|
||||
slog.Uint64("snd.wnd", uint64(tcb.snd.WND)),
|
||||
)
|
||||
}
|
||||
|
||||
func (tcb *ControlBlock) traceRcv(msg string) {
|
||||
tcb.trace(msg,
|
||||
slog.String("state", tcb.state.String()),
|
||||
slog.Uint64("rcv.nxt", uint64(tcb.rcv.NXT)),
|
||||
slog.Uint64("rcv.wnd", uint64(tcb.rcv.WND)),
|
||||
slog.Bool("challenge", tcb.challengeAck),
|
||||
)
|
||||
}
|
||||
|
||||
func (tcb *ControlBlock) traceSeg(msg string, seg Segment) {
|
||||
if tcb.logenabled(internal.LevelTrace) {
|
||||
tcb.trace(msg,
|
||||
slog.Uint64("seg.seq", uint64(seg.SEQ)),
|
||||
slog.Uint64("seg.ack", uint64(seg.ACK)),
|
||||
slog.Uint64("seg.wnd", uint64(seg.WND)),
|
||||
slog.String("seg.flags", seg.Flags.String()),
|
||||
slog.Uint64("seg.data", uint64(seg.DATALEN)),
|
||||
)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,391 @@
|
||||
package tcp
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"math/bits"
|
||||
"strconv"
|
||||
"strings"
|
||||
"unsafe"
|
||||
)
|
||||
|
||||
//go:generate stringer -type=State,OptionKind -linecomment -output stringers.go .
|
||||
|
||||
var (
|
||||
// errDropSegment is a flag that signals to drop a segment silently.
|
||||
errDropSegment = errors.New("drop segment")
|
||||
errWindowTooLarge = errors.New("invalid window size > 2**16")
|
||||
|
||||
errTCBNotClosed = errors.New("TCB not closed")
|
||||
errInvalidState = errors.New("invalid state")
|
||||
errConnNotexist = errors.New("connection does not exist")
|
||||
errConnectionClosing = errors.New("connection closing")
|
||||
errExpectedSYN = errors.New("seqs:expected SYN")
|
||||
errBadSegack = errors.New("seqs:bad segack")
|
||||
errFinwaitExpectedACK = errors.New("seqs:finwait1 expected ACK")
|
||||
errFinwaitExpectedFinack = errors.New("seqs:finwait2 expected FINACK")
|
||||
|
||||
errWindowOverflow = newRejectErr("wnd > 2**16")
|
||||
errSeqNotInWindow = newRejectErr("seq not in snd/rcv.wnd")
|
||||
errZeroWindow = newRejectErr("zero window")
|
||||
errLastNotInWindow = newRejectErr("last not in snd/rcv.wnd")
|
||||
errRequireSequential = newRejectErr("seq != rcv.nxt (require sequential segments)")
|
||||
errAckNotNext = newRejectErr("ack != snd.nxt")
|
||||
)
|
||||
|
||||
func newRejectErr(err string) *RejectError { return &RejectError{err: "reject in/out seg: " + err} }
|
||||
|
||||
// RejectError represents an error that arises during admission of a segment into the
|
||||
// Transmission Control Block logic in which the packet cannot be processed by the TCB.
|
||||
type RejectError struct {
|
||||
err string
|
||||
}
|
||||
|
||||
func (e *RejectError) Error() string { return e.err }
|
||||
|
||||
// Segment represents an incoming/outgoing TCP segment in the sequence space.
|
||||
type Segment struct {
|
||||
SEQ Value // sequence number of first octet of segment. If SYN is set it is the initial sequence number (ISN) and the first data octet is ISN+1.
|
||||
ACK Value // acknowledgment number. If ACK is set it is sequence number of first octet the sender of the segment is expecting to receive next.
|
||||
DATALEN Size // The number of octets occupied by the data (payload) not counting SYN and FIN.
|
||||
WND Size // segment window
|
||||
Flags Flags // TCP flags.
|
||||
}
|
||||
|
||||
// LEN returns the length of the segment in octets including SYN and FIN flags.
|
||||
func (seg *Segment) LEN() Size {
|
||||
add := Size(seg.Flags>>0) & 1 // Add FIN bit.
|
||||
add += Size(seg.Flags>>1) & 1 // Add SYN bit.
|
||||
return seg.DATALEN + add
|
||||
}
|
||||
|
||||
// End returns the sequence number of the last octet of the segment.
|
||||
func (seg *Segment) Last() Value {
|
||||
seglen := seg.LEN()
|
||||
if seglen == 0 {
|
||||
return seg.SEQ
|
||||
}
|
||||
return Add(seg.SEQ, seglen) - 1
|
||||
}
|
||||
|
||||
// StringExchange returns a string representation of a segment exchange over
|
||||
// a network in RFC9293 styled visualization. invertDir inverts the arrow directions.
|
||||
// i.e:
|
||||
//
|
||||
// SynSent --> <SEQ=300><ACK=91>[SYN,ACK] --> SynRcvd
|
||||
func StringExchange(seg Segment, A, B State, invertDir bool) string {
|
||||
b := make([]byte, 0, 64)
|
||||
b = appendStringExchange(b, seg, A, B, invertDir)
|
||||
return unsafe.String(unsafe.SliceData(b), len(b))
|
||||
}
|
||||
|
||||
// appendStringExchange appends a RFC9293 styled visualization of exchange to buf.
|
||||
// i.e:
|
||||
//
|
||||
// SynSent --> <SEQ=300><ACK=91>[SYN,ACK] --> SynRcvd
|
||||
func appendStringExchange(buf []byte, seg Segment, A, B State, invertDir bool) []byte {
|
||||
const emptySpaces = " "
|
||||
const fill = len(emptySpaces) - 1
|
||||
appendVal := func(buf []byte, name string, i Value) []byte {
|
||||
buf = append(buf, '<')
|
||||
buf = append(buf, name...)
|
||||
buf = append(buf, '=')
|
||||
buf = strconv.AppendInt(buf, int64(i), 10)
|
||||
buf = append(buf, '>')
|
||||
return buf
|
||||
}
|
||||
startLen := len(buf)
|
||||
dirSep := []byte(" --> ")
|
||||
if invertDir {
|
||||
dirSep = []byte(" <-- ")
|
||||
}
|
||||
astr := A.String()
|
||||
buf = append(buf, astr...)
|
||||
if len(astr) < fill {
|
||||
// Space padding.
|
||||
buf = append(buf, emptySpaces[:fill-len(astr)]...)
|
||||
}
|
||||
buf = append(buf, dirSep...)
|
||||
buf = appendVal(buf, "SEQ", seg.SEQ)
|
||||
buf = appendVal(buf, "ACK", seg.ACK)
|
||||
if seg.DATALEN > 0 {
|
||||
buf = appendVal(buf, "DATA", Value(seg.DATALEN))
|
||||
}
|
||||
buf = append(buf, '[')
|
||||
buf = seg.Flags.AppendFormat(buf)
|
||||
buf = append(buf, ']')
|
||||
if len(buf)-startLen < 48 {
|
||||
// More space padding.
|
||||
buf = append(buf, emptySpaces[:48-len(buf)]...)
|
||||
}
|
||||
buf = append(buf, dirSep...)
|
||||
buf = append(buf, B.String()...)
|
||||
return buf
|
||||
}
|
||||
|
||||
// Flags is a TCP flags bit-masked implementation i.e: SYN, FIN, ACK.
|
||||
type Flags uint16
|
||||
|
||||
const (
|
||||
FlagFIN Flags = 1 << iota // FlagFIN - No more data from sender.
|
||||
FlagSYN // FlagSYN - Synchronize sequence numbers.
|
||||
FlagRST // FlagRST - Reset the connection.
|
||||
FlagPSH // FlagPSH - Push function.
|
||||
FlagACK // FlagACK - Acknowledgment field significant.
|
||||
FlagURG // FlagURG - Urgent pointer field significant.
|
||||
FlagECE // FlagECE - ECN-Echo has a nonce-sum in the SYN/ACK.
|
||||
FlagCWR // FlagCWR - Congestion Window Reduced.
|
||||
FlagNS // FlagNS - Nonce Sum flag (see RFC 3540).
|
||||
)
|
||||
|
||||
const flagMask = 0x01ff
|
||||
|
||||
// The union of SYN|FIN|PSH and ACK flags is commonly found throughout the specification, so we define unexported shorthands.
|
||||
const (
|
||||
synack = FlagSYN | FlagACK
|
||||
finack = FlagFIN | FlagACK
|
||||
pshack = FlagPSH | FlagACK
|
||||
)
|
||||
|
||||
// HasAll checks if mask bits are all set in the receiver flags.
|
||||
func (flags Flags) HasAll(mask Flags) bool { return flags&mask == mask }
|
||||
|
||||
// HasAny checks if one or more mask bits are set in receiver flags.
|
||||
func (flags Flags) HasAny(mask Flags) bool { return flags&mask != 0 }
|
||||
|
||||
// Mask returns the flags with non-flag bits unset.
|
||||
func (flags Flags) Mask() Flags { return flags & flagMask }
|
||||
|
||||
// StringFlags returns human readable flag string. i.e:
|
||||
//
|
||||
// "[SYN,ACK]"
|
||||
//
|
||||
// Flags are printed in order from LSB (FIN) to MSB (NS).
|
||||
// All flags are printed with length of 3, so a NS flag will
|
||||
// end with a space i.e. [ACK,NS ]
|
||||
func (flags Flags) String() string {
|
||||
// Cover most common cases without heap allocating.
|
||||
switch flags {
|
||||
case 0:
|
||||
return "[]"
|
||||
case synack:
|
||||
return "[SYN,ACK]"
|
||||
case finack:
|
||||
return "[FIN,ACK]"
|
||||
case pshack:
|
||||
return "[PSH,ACK]"
|
||||
case FlagACK:
|
||||
return "[ACK]"
|
||||
case FlagSYN:
|
||||
return "[SYN]"
|
||||
case FlagFIN:
|
||||
return "[FIN]"
|
||||
case FlagRST:
|
||||
return "[RST]"
|
||||
}
|
||||
buf := make([]byte, 0, 2+3*bits.OnesCount16(uint16(flags)))
|
||||
buf = append(buf, '[')
|
||||
buf = flags.AppendFormat(buf)
|
||||
buf = append(buf, ']')
|
||||
return string(buf)
|
||||
}
|
||||
|
||||
// AppendFormat appends a human readable flag string to b returning the extended buffer.
|
||||
func (flags Flags) AppendFormat(b []byte) []byte {
|
||||
if flags == 0 {
|
||||
return b
|
||||
}
|
||||
// String Flag const
|
||||
const flaglen = 3
|
||||
const strflags = "FINSYNRSTPSHACKURGECECWRNS "
|
||||
var addcommas bool
|
||||
for flags != 0 { // written by Github Copilot- looks OK.
|
||||
i := bits.TrailingZeros16(uint16(flags))
|
||||
if addcommas {
|
||||
b = append(b, ',')
|
||||
} else {
|
||||
addcommas = true
|
||||
}
|
||||
b = append(b, strflags[i*flaglen:i*flaglen+flaglen]...)
|
||||
flags &= ^(1 << i)
|
||||
}
|
||||
return b
|
||||
}
|
||||
|
||||
// State enumerates states a TCP connection progresses through during its lifetime.
|
||||
type State uint8
|
||||
|
||||
const (
|
||||
// CLOSED - represents no connection state at all. Is not a valid state of the TCP state machine but rather a pseudo-state pre-initialization.
|
||||
StateClosed State = iota // CLOSED
|
||||
// LISTEN - represents waiting for a connection request from any remote TCP and port.
|
||||
StateListen // LISTEN
|
||||
// SYN-RECEIVED - represents waiting for a confirming connection request acknowledgment
|
||||
// after having both received and sent a connection request.
|
||||
StateSynRcvd // SYN-RECEIVED
|
||||
// SYN-SENT - represents waiting for a matching connection request after having sent a connection request.
|
||||
StateSynSent // SYN-SENT
|
||||
// ESTABLISHED - represents an open connection, data received can be delivered
|
||||
// to the user. The normal state for the data transfer phase of the connection.
|
||||
StateEstablished // ESTABLISHED
|
||||
// FIN-WAIT-1 - represents waiting for a connection termination request
|
||||
// from the remote TCP, or an acknowledgment of the connection
|
||||
// termination request previously sent.
|
||||
StateFinWait1 // FIN-WAIT-1
|
||||
// FIN-WAIT-2 - represents waiting for a connection termination request
|
||||
// from the remote TCP.
|
||||
StateFinWait2 // FIN-WAIT-2
|
||||
// CLOSING - represents waiting for a connection termination request
|
||||
// acknowledgment from the remote TCP.
|
||||
StateClosing // CLOSING
|
||||
// TIME-WAIT - represents waiting for enough time to pass to be sure the remote
|
||||
// TCP received the acknowledgment of its connection termination request.
|
||||
StateTimeWait // TIME-WAIT
|
||||
// CLOSE-WAIT - represents waiting for a connection termination request
|
||||
// from the local user.
|
||||
StateCloseWait // CLOSE-WAIT
|
||||
// LAST-ACK - represents waiting for an acknowledgment of the
|
||||
// connection termination request previously sent to the remote TCP
|
||||
// (which includes an acknowledgment of its connection termination request).
|
||||
StateLastAck // LAST-ACK
|
||||
)
|
||||
|
||||
// IsPreestablished returns true if the connection is in a state preceding the established state.
|
||||
// Returns false for Closed pseudo state.
|
||||
func (s State) IsPreestablished() bool {
|
||||
return s == StateSynRcvd || s == StateSynSent || s == StateListen
|
||||
}
|
||||
|
||||
// IsClosing returns true if the connection is in a closing state but not yet terminated (relieved of remote connection state).
|
||||
// Returns false for Closed pseudo state.
|
||||
func (s State) IsClosing() bool {
|
||||
return !(s <= StateEstablished)
|
||||
}
|
||||
|
||||
// IsClosed returns true if the connection closed and can possibly relieved of
|
||||
// all state related to the remote connection. It returns true if Closed or in TimeWait.
|
||||
func (s State) IsClosed() bool {
|
||||
return s == StateClosed || s == StateTimeWait
|
||||
}
|
||||
|
||||
// IsSynchronized returns true if the connection has gone through the Established state.
|
||||
func (s State) IsSynchronized() bool {
|
||||
return s >= StateEstablished
|
||||
}
|
||||
|
||||
// IsDataOpen returns true if the connection allows sending and receiving of data.
|
||||
func (s State) isOpen() bool {
|
||||
return !s.IsClosed()
|
||||
}
|
||||
|
||||
// hasIRS checks if the ControlBlock has received a valid initial sequence number (IRS).
|
||||
func (s State) hasIRS() bool {
|
||||
return s.isOpen() && s != StateSynSent && s != StateListen
|
||||
}
|
||||
|
||||
type OptionKind uint8
|
||||
|
||||
const (
|
||||
OptEnd OptionKind = iota // end of option list
|
||||
OptNop // no-operation
|
||||
OptMaxSegmentSize // maximum segment size
|
||||
OptWindowScale // window scale
|
||||
OptSACKPermitted // SACK permitted
|
||||
OptSACK // SACK
|
||||
OptEcho // echo(obsolete)
|
||||
optEchoReply // echo reply(obsolete)
|
||||
OptTimestamps // timestamps
|
||||
optPOCP // partial order connection permitted(obsolete)
|
||||
optPOSP // partial order service profile(obsolete)
|
||||
optCC // CC(obsolete)
|
||||
optCCnew // CC.new(obsolete)
|
||||
optCCecho // CC.echo(obsolete)
|
||||
optACR // alternate checksum request(obsolete)
|
||||
optACD // alternate checksum data(obsolete)
|
||||
optSkeeter // skeeter
|
||||
optBubba // bubba
|
||||
OptTrailerChecksum // trailer checksum
|
||||
optMD5Signature // MD5 signature(obsolete)
|
||||
OptSCPSCapabilities // SCPS capabilities
|
||||
OptSNA // selective negative acks
|
||||
OptRecordBoundaries // record boundaries
|
||||
OptCorruptionExperienced // corruption experienced
|
||||
OptSNAP // SNAP
|
||||
OptUnassigned // unassigned
|
||||
OptCompressionFilter // compression filter
|
||||
OptQuickStartResponse // quick-start response
|
||||
OptUserTimeout // user timeout or unauthorized use
|
||||
OptAuthetication // Authentication TCP-AO
|
||||
OptMultipath // multipath TCP
|
||||
)
|
||||
|
||||
const (
|
||||
OptFastOpenCookie OptionKind = 34 // fast open cookie
|
||||
OptEncryptionNegotiation OptionKind = 69 // encryption negotiation
|
||||
OptAccurateECN0 OptionKind = 172 // accurate ECN order 0
|
||||
OptAccurateECN1 OptionKind = 174 // accurate ECN order 1
|
||||
)
|
||||
|
||||
// IsObsolete returns true if option considered obsolete by newer TCP specifications.
|
||||
func (kind OptionKind) IsObsolete() bool {
|
||||
if kind.IsDefined() {
|
||||
return strings.HasSuffix(kind.String(), "(obsolete)")
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
// IsDefined returns true if the option is a known unreserved option kind.
|
||||
func (kind OptionKind) IsDefined() bool {
|
||||
return kind <= 30 || kind == 34 || kind == 69 || kind == 172 || kind == 174
|
||||
}
|
||||
|
||||
type OptionParser struct {
|
||||
SkipSizeValidation bool
|
||||
SkipObsolete bool
|
||||
}
|
||||
|
||||
func (op *OptionParser) ForEachOption(opts []byte, fn func(OptionKind, []byte) error) error {
|
||||
off := 0
|
||||
skipSizeValidation := op.SkipSizeValidation
|
||||
skipObsolete := op.SkipObsolete
|
||||
for off < len(opts) && opts[off] != 0 {
|
||||
kind := OptionKind(opts[off])
|
||||
off++
|
||||
if kind == OptNop {
|
||||
continue
|
||||
}
|
||||
if len(opts[off:]) < 2 {
|
||||
return errors.New("short TCP options")
|
||||
}
|
||||
size := int(opts[off])
|
||||
off++
|
||||
if len(opts[off:]) < size {
|
||||
return fmt.Errorf("option %q length %d exceeds buffer size %d", kind.String(), size, len(opts[off:]))
|
||||
}
|
||||
|
||||
if !skipSizeValidation {
|
||||
expectSize := -1
|
||||
switch kind {
|
||||
case OptTimestamps:
|
||||
expectSize = 10
|
||||
case OptMaxSegmentSize, OptUserTimeout:
|
||||
expectSize = 4
|
||||
case OptWindowScale:
|
||||
expectSize = 3
|
||||
case OptSACKPermitted:
|
||||
expectSize = 2
|
||||
}
|
||||
if expectSize != -1 && size != expectSize {
|
||||
return fmt.Errorf("bad TCP option %q size want %d got %d", kind.String(), expectSize, opts[off])
|
||||
}
|
||||
}
|
||||
if skipObsolete && kind.IsObsolete() {
|
||||
err := fn(kind, opts[off:off+size])
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
off += size
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,193 @@
|
||||
package tcp
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"testing"
|
||||
)
|
||||
|
||||
// Here we define internal testing helpers that may be used in any *_test.go file
|
||||
// but are not exported.
|
||||
|
||||
// Exchange represents a single exchange of segments.
|
||||
type Exchange struct {
|
||||
Outgoing *Segment
|
||||
Incoming *Segment
|
||||
WantPending *Segment // Expected pending segment. If nil not checked.
|
||||
WantState State // Expected end state.
|
||||
WantPeerState State // Expected end state of peer. Not necessary when calling HelperExchange but can aid with logging information.
|
||||
}
|
||||
|
||||
func (tcb *ControlBlock) HelperExchange(t *testing.T, exchange []Exchange) {
|
||||
t.Helper()
|
||||
var i int
|
||||
var ex Exchange
|
||||
defer func() {
|
||||
if t.Failed() {
|
||||
t.Errorf("exchange failed:\nwant: %s\ngot: %s",
|
||||
ex.RFC9293String(ex.WantState, ex.WantPeerState),
|
||||
ex.RFC9293String(tcb.state, ex.WantPeerState),
|
||||
)
|
||||
}
|
||||
}()
|
||||
const pfx = "exchange"
|
||||
t.Log(tcb.state, "Exchange start")
|
||||
for i, ex = range exchange {
|
||||
if ex.Outgoing != nil && ex.Incoming != nil {
|
||||
t.Fatalf(pfx+"[%d] cannot send and receive in the same exchange, please split into two exchanges.", i)
|
||||
} else if ex.Outgoing == nil && ex.Incoming == nil {
|
||||
t.Fatalf(pfx+"[%d] must send or receive a segment.", i)
|
||||
}
|
||||
if ex.Outgoing != nil {
|
||||
prevInflight := tcb.snd.inFlight()
|
||||
err := tcb.Send(*ex.Outgoing)
|
||||
gotSent := tcb.snd.inFlight() - prevInflight
|
||||
if err != nil {
|
||||
t.Fatalf(pfx+"[%d] snd: %s\nseg=%+v\nrcv=%+v\nsnd=%+v", i, err, *ex.Outgoing, tcb.rcv, tcb.snd)
|
||||
} else if gotSent != ex.Outgoing.LEN() {
|
||||
t.Fatalf(pfx+"[%d] snd: expected %d data sent, calculated inflight %d", i, ex.Outgoing.LEN(), gotSent)
|
||||
}
|
||||
}
|
||||
if ex.Incoming != nil {
|
||||
err := tcb.Recv(*ex.Incoming)
|
||||
if err != nil {
|
||||
msg := fmt.Sprintf(pfx+"[%d] rcv: %s\nseg=%+v\nrcv=%+v\nsnd=%+v", i, err, *ex.Incoming, tcb.rcv, tcb.snd)
|
||||
if IsDroppedErr(err) {
|
||||
t.Log(msg)
|
||||
} else {
|
||||
t.Fatal(msg)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
t.Log(ex.RFC9293String(tcb.state, ex.WantPeerState))
|
||||
|
||||
state := tcb.State()
|
||||
if state != ex.WantState {
|
||||
t.Errorf(pfx+"[%d] unexpected state:\n got=%s\nwant=%s", i, state, ex.WantState)
|
||||
}
|
||||
pending, ok := tcb.PendingSegment(0)
|
||||
if !ok && ex.WantPending != nil {
|
||||
t.Fatalf(pfx+"[%d] pending:got none, want=%+v", i, *ex.WantPending)
|
||||
} else if ex.WantPending != nil && pending != *ex.WantPending {
|
||||
t.Fatalf(pfx+"[%d] pending:\n got=%+v\nwant=%+v", i, pending, *ex.WantPending)
|
||||
} else if ok && ex.WantPending == nil {
|
||||
t.Fatalf(pfx+"[%d] pending:\n got=%+v\nwant=none", i, pending)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (tcb *ControlBlock) HelperInitState(state State, localISS, localNXT Value, localWindow Size) {
|
||||
tcb.state = state
|
||||
tcb.snd = sendSpace{
|
||||
ISS: localISS,
|
||||
UNA: localISS,
|
||||
NXT: localNXT,
|
||||
WND: 1, // 1 byte window, so we can test the SEQ field.
|
||||
// UP, WL1, WL2 defaults to zero values.
|
||||
}
|
||||
tcb.rcv = recvSpace{
|
||||
WND: localWindow,
|
||||
}
|
||||
}
|
||||
|
||||
func (tcb *ControlBlock) HelperInitRcv(irs, nxt Value, remoteWindow Size) {
|
||||
tcb.rcv.IRS = irs
|
||||
tcb.rcv.NXT = nxt
|
||||
tcb.snd.WND = remoteWindow
|
||||
}
|
||||
|
||||
func (tcb *ControlBlock) RelativeSendSpace() sendSpace {
|
||||
snd := tcb.snd
|
||||
snd.NXT -= snd.ISS
|
||||
snd.UNA -= snd.ISS
|
||||
snd.ISS = 0
|
||||
return snd
|
||||
}
|
||||
|
||||
func (tcb *ControlBlock) RelativeRecvSpace() recvSpace {
|
||||
rcv := tcb.rcv
|
||||
rcv.NXT -= rcv.IRS
|
||||
rcv.IRS = 0
|
||||
return rcv
|
||||
}
|
||||
|
||||
func (tcb *ControlBlock) RelativeRecvSegment(seg Segment) Segment {
|
||||
seg.SEQ -= tcb.rcv.IRS
|
||||
seg.ACK -= tcb.snd.ISS
|
||||
return seg
|
||||
}
|
||||
|
||||
func (tcb *ControlBlock) RelativeSendSegment(seg Segment) Segment {
|
||||
seg.SEQ -= tcb.snd.ISS
|
||||
seg.ACK -= tcb.rcv.IRS
|
||||
return seg
|
||||
}
|
||||
|
||||
func (tcb *ControlBlock) RelativeAutoSegment(seg Segment) Segment {
|
||||
rcv := tcb.RelativeRecvSegment(seg)
|
||||
snd := tcb.RelativeSendSegment(seg)
|
||||
if rcv.SEQ > snd.SEQ {
|
||||
return snd
|
||||
}
|
||||
return rcv
|
||||
}
|
||||
|
||||
func (tcb *ControlBlock) HelperPrintSegment(t *testing.T, isReceive bool, seg Segment) {
|
||||
const fmtmsg = "\nSeg=%+v\nRcvSpace=%s\nSndSpace=%s"
|
||||
rcv := tcb.RelativeRecvSpace()
|
||||
rcvStr := rcv.RelativeGoString()
|
||||
snd := tcb.RelativeSendSpace()
|
||||
sndStr := snd.RelativeGoString()
|
||||
t.Helper()
|
||||
if isReceive {
|
||||
t.Logf("RECV:"+fmtmsg, seg.RelativeGoString(tcb.rcv.IRS, tcb.snd.ISS), rcvStr, sndStr)
|
||||
} else {
|
||||
t.Logf("SEND:"+fmtmsg, seg.RelativeGoString(tcb.snd.ISS, tcb.rcv.IRS), rcvStr, sndStr)
|
||||
}
|
||||
}
|
||||
|
||||
func (rcv recvSpace) RelativeGoString() string {
|
||||
return fmt.Sprintf("{NXT:%d} ", rcv.NXT-rcv.IRS)
|
||||
}
|
||||
|
||||
func (rcv sendSpace) RelativeGoString() string {
|
||||
nxt := rcv.NXT - rcv.ISS
|
||||
una := rcv.UNA - rcv.ISS
|
||||
unaLen := Sizeof(una, nxt)
|
||||
if unaLen != 0 {
|
||||
return fmt.Sprintf("{NXT:%d UNA:%d} (%d unacked)", nxt, una, unaLen)
|
||||
}
|
||||
return fmt.Sprintf("{NXT:%d UNA:%d}", nxt, una)
|
||||
}
|
||||
|
||||
func (seg Segment) RelativeGoString(iseq, iack Value) string {
|
||||
seglen := seg.LEN()
|
||||
if seglen != seg.DATALEN {
|
||||
// If SYN/FIN is set print out the length of the segment.
|
||||
return fmt.Sprintf("{SEQ:%d ACK:%d DATALEN:%d Flags:%s} (LEN:%d)", seg.SEQ-iseq, seg.ACK-iack, seg.DATALEN, seg.Flags, seglen)
|
||||
}
|
||||
return fmt.Sprintf("{SEQ:%d ACK:%d DATALEN:%d Flags:%s} ", seg.SEQ-iseq, seg.ACK-iack, seg.DATALEN, seg.Flags)
|
||||
}
|
||||
|
||||
// https://datatracker.ietf.org/doc/html/rfc9293#section-3.8.6.2.1
|
||||
func (tcb *ControlBlock) UsableWindow() Size {
|
||||
return Sizeof(tcb.snd.NXT, tcb.snd.UNA) + tcb.snd.WND
|
||||
}
|
||||
|
||||
func IsDroppedErr(err error) bool {
|
||||
return err != nil && errors.Is(err, errDropSegment)
|
||||
}
|
||||
|
||||
func (ex *Exchange) RFC9293String(A, B State) string {
|
||||
var seg Segment
|
||||
sentByA := ex.Outgoing != nil
|
||||
if sentByA {
|
||||
seg = *ex.Outgoing
|
||||
} else if ex.Incoming != nil {
|
||||
seg = *ex.Incoming
|
||||
} else {
|
||||
return ""
|
||||
}
|
||||
return StringExchange(seg, A, B, !sentByA)
|
||||
}
|
||||
@@ -0,0 +1,102 @@
|
||||
// Code generated by "stringer -type=State,OptionKind -linecomment -output stringers.go ."; DO NOT EDIT.
|
||||
|
||||
package tcp
|
||||
|
||||
import "strconv"
|
||||
|
||||
func _() {
|
||||
// An "invalid array index" compiler error signifies that the constant values have changed.
|
||||
// Re-run the stringer command to generate them again.
|
||||
var x [1]struct{}
|
||||
_ = x[StateClosed-0]
|
||||
_ = x[StateListen-1]
|
||||
_ = x[StateSynRcvd-2]
|
||||
_ = x[StateSynSent-3]
|
||||
_ = x[StateEstablished-4]
|
||||
_ = x[StateFinWait1-5]
|
||||
_ = x[StateFinWait2-6]
|
||||
_ = x[StateClosing-7]
|
||||
_ = x[StateTimeWait-8]
|
||||
_ = x[StateCloseWait-9]
|
||||
_ = x[StateLastAck-10]
|
||||
}
|
||||
|
||||
const _State_name = "CLOSEDLISTENSYN-RECEIVEDSYN-SENTESTABLISHEDFIN-WAIT-1FIN-WAIT-2CLOSINGTIME-WAITCLOSE-WAITLAST-ACK"
|
||||
|
||||
var _State_index = [...]uint8{0, 6, 12, 24, 32, 43, 53, 63, 70, 79, 89, 97}
|
||||
|
||||
func (i State) String() string {
|
||||
if i >= State(len(_State_index)-1) {
|
||||
return "State(" + strconv.FormatInt(int64(i), 10) + ")"
|
||||
}
|
||||
return _State_name[_State_index[i]:_State_index[i+1]]
|
||||
}
|
||||
func _() {
|
||||
// An "invalid array index" compiler error signifies that the constant values have changed.
|
||||
// Re-run the stringer command to generate them again.
|
||||
var x [1]struct{}
|
||||
_ = x[OptEnd-0]
|
||||
_ = x[OptNop-1]
|
||||
_ = x[OptMaxSegmentSize-2]
|
||||
_ = x[OptWindowScale-3]
|
||||
_ = x[OptSACKPermitted-4]
|
||||
_ = x[OptSACK-5]
|
||||
_ = x[OptEcho-6]
|
||||
_ = x[optEchoReply-7]
|
||||
_ = x[OptTimestamps-8]
|
||||
_ = x[optPOCP-9]
|
||||
_ = x[optPOSP-10]
|
||||
_ = x[optCC-11]
|
||||
_ = x[optCCnew-12]
|
||||
_ = x[optCCecho-13]
|
||||
_ = x[optACR-14]
|
||||
_ = x[optACD-15]
|
||||
_ = x[optSkeeter-16]
|
||||
_ = x[optBubba-17]
|
||||
_ = x[OptTrailerChecksum-18]
|
||||
_ = x[optMD5Signature-19]
|
||||
_ = x[OptSCPSCapabilities-20]
|
||||
_ = x[OptSNA-21]
|
||||
_ = x[OptRecordBoundaries-22]
|
||||
_ = x[OptCorruptionExperienced-23]
|
||||
_ = x[OptSNAP-24]
|
||||
_ = x[OptUnassigned-25]
|
||||
_ = x[OptCompressionFilter-26]
|
||||
_ = x[OptQuickStartResponse-27]
|
||||
_ = x[OptUserTimeout-28]
|
||||
_ = x[OptAuthetication-29]
|
||||
_ = x[OptMultipath-30]
|
||||
_ = x[OptFastOpenCookie-34]
|
||||
_ = x[OptEncryptionNegotiation-69]
|
||||
_ = x[OptAccurateECN0-172]
|
||||
_ = x[OptAccurateECN1-174]
|
||||
}
|
||||
|
||||
const (
|
||||
_OptionKind_name_0 = "end of option listno-operationmaximum segment sizewindow scaleSACK permittedSACKecho(obsolete)echo reply(obsolete)timestampspartial order connection permitted(obsolete)partial order service profile(obsolete)CC(obsolete)CC.new(obsolete)CC.echo(obsolete)alternate checksum request(obsolete)alternate checksum data(obsolete)skeeterbubbatrailer checksumMD5 signature(obsolete)SCPS capabilitiesselective negative acksrecord boundariescorruption experiencedSNAPunassignedcompression filterquick-start responseuser timeout or unauthorized useAuthentication TCP-AOmultipath TCP"
|
||||
_OptionKind_name_1 = "fast open cookie"
|
||||
_OptionKind_name_2 = "encryption negotiation"
|
||||
_OptionKind_name_3 = "accurate ECN order 0"
|
||||
_OptionKind_name_4 = "accurate ECN order 1"
|
||||
)
|
||||
|
||||
var (
|
||||
_OptionKind_index_0 = [...]uint16{0, 18, 30, 50, 62, 76, 80, 94, 114, 124, 168, 207, 219, 235, 252, 288, 321, 328, 333, 349, 372, 389, 412, 429, 451, 455, 465, 483, 503, 535, 556, 569}
|
||||
)
|
||||
|
||||
func (i OptionKind) String() string {
|
||||
switch {
|
||||
case i <= 30:
|
||||
return _OptionKind_name_0[_OptionKind_index_0[i]:_OptionKind_index_0[i+1]]
|
||||
case i == 34:
|
||||
return _OptionKind_name_1
|
||||
case i == 69:
|
||||
return _OptionKind_name_2
|
||||
case i == 172:
|
||||
return _OptionKind_name_3
|
||||
case i == 174:
|
||||
return _OptionKind_name_4
|
||||
default:
|
||||
return "OptionKind(" + strconv.FormatInt(int64(i), 10) + ")"
|
||||
}
|
||||
}
|
||||
+912
@@ -0,0 +1,912 @@
|
||||
package tcp_test
|
||||
|
||||
import (
|
||||
"math/rand"
|
||||
"strconv"
|
||||
"testing"
|
||||
|
||||
"github.com/soypat/lneto"
|
||||
"github.com/soypat/lneto/tcp"
|
||||
)
|
||||
|
||||
const (
|
||||
SYNACK = tcp.FlagSYN | tcp.FlagACK
|
||||
FINACK = tcp.FlagFIN | tcp.FlagACK
|
||||
PSHACK = tcp.FlagPSH | tcp.FlagACK
|
||||
)
|
||||
|
||||
/*
|
||||
Section 3.5 of RFC 9293: Basic 3-way handshake for connection synchronization.
|
||||
TCP Peer A TCP Peer B
|
||||
|
||||
1. CLOSED LISTEN
|
||||
|
||||
2. SYN-SENT --> <SEQ=100><CTL=SYN> --> SYN-RECEIVED
|
||||
|
||||
3. ESTABLISHED <-- <SEQ=300><ACK=101><CTL=SYN,ACK> <-- SYN-RECEIVED
|
||||
|
||||
4. ESTABLISHED --> <SEQ=101><ACK=301><CTL=ACK> --> ESTABLISHED
|
||||
|
||||
5. ESTABLISHED --> <SEQ=101><ACK=301><CTL=ACK><DATA> --> ESTABLISHED
|
||||
*/
|
||||
func TestExchange_rfc9293_figure6(t *testing.T) {
|
||||
const issA, issB, windowA, windowB = 100, 300, 1000, 1000
|
||||
exchangeA := []tcp.Exchange{
|
||||
{ // A sends SYN to B.
|
||||
Outgoing: &tcp.Segment{SEQ: issA, Flags: tcp.FlagSYN, WND: windowA},
|
||||
WantState: tcp.StateSynSent,
|
||||
WantPeerState: tcp.StateSynRcvd,
|
||||
},
|
||||
{ // A receives SYNACK from B thus establishing the connection on A's side.
|
||||
Incoming: &tcp.Segment{SEQ: issB, ACK: issA + 1, Flags: SYNACK, WND: windowB},
|
||||
WantState: tcp.StateEstablished,
|
||||
WantPending: &tcp.Segment{SEQ: issA + 1, ACK: issB + 1, Flags: tcp.FlagACK, WND: windowA},
|
||||
WantPeerState: tcp.StateSynRcvd,
|
||||
},
|
||||
{ // A sends ACK to B, which leaves connection established on their side. Three way handshake complete by now.
|
||||
Outgoing: &tcp.Segment{SEQ: issA + 1, ACK: issB + 1, Flags: tcp.FlagACK, WND: windowA},
|
||||
WantState: tcp.StateEstablished,
|
||||
WantPeerState: tcp.StateEstablished,
|
||||
},
|
||||
}
|
||||
var tcbA tcp.ControlBlock
|
||||
tcbA.HelperInitState(tcp.StateSynSent, issA, issA, windowA)
|
||||
tcbA.HelperExchange(t, exchangeA)
|
||||
segA, ok := tcbA.PendingSegment(0)
|
||||
if ok {
|
||||
t.Error("unexpected Client pending segment after establishment: ", segA)
|
||||
}
|
||||
exchangeB := reverseExchange(exchangeA)
|
||||
|
||||
var tcbB tcp.ControlBlock
|
||||
tcbB.HelperInitState(tcp.StateListen, issB, issB, windowB)
|
||||
tcbB.HelperExchange(t, exchangeB) // TODO remove [:3] after snd.UNA bugfix
|
||||
segB, ok := tcbB.PendingSegment(0)
|
||||
if ok {
|
||||
t.Error("unexpected Listener pending segment after establishment: ", segB)
|
||||
}
|
||||
}
|
||||
|
||||
/*
|
||||
Section 3.5 of RFC 9293: Simultaneous Connection Synchronization (SYN).
|
||||
TCP Peer A TCP Peer B
|
||||
|
||||
1. CLOSED CLOSED
|
||||
|
||||
2. SYN-SENT --> <SEQ=100><CTL=SYN> ...
|
||||
|
||||
3. SYN-RECEIVED <-- <SEQ=300><CTL=SYN> <-- SYN-SENT
|
||||
|
||||
4. ... <SEQ=100><CTL=SYN> --> SYN-RECEIVED
|
||||
|
||||
5. SYN-RECEIVED --> <SEQ=100><ACK=301><CTL=SYN,ACK> ...
|
||||
|
||||
6. ESTABLISHED <-- <SEQ=300><ACK=101><CTL=SYN,ACK> <-- SYN-RECEIVED
|
||||
|
||||
7. ... <SEQ=100><ACK=301><CTL=SYN,ACK> --> ESTABLISHED
|
||||
*/
|
||||
func TestExchange_rfc9293_figure7(t *testing.T) {
|
||||
const issA, issB, windowA, windowB = 100, 300, 1000, 1000
|
||||
exchangeA := []tcp.Exchange{
|
||||
0: { // A sends SYN to B.
|
||||
Outgoing: &tcp.Segment{SEQ: issA, Flags: tcp.FlagSYN, WND: windowA},
|
||||
WantState: tcp.StateSynSent,
|
||||
},
|
||||
1: { // A receives a SYN with no ACK from B.
|
||||
Incoming: &tcp.Segment{SEQ: issB, Flags: tcp.FlagSYN, WND: windowB},
|
||||
WantState: tcp.StateSynRcvd,
|
||||
WantPending: &tcp.Segment{SEQ: issA, ACK: issB + 1, Flags: SYNACK, WND: windowA},
|
||||
},
|
||||
2: { // A sends SYNACK to B.
|
||||
Outgoing: &tcp.Segment{SEQ: issA, ACK: issB + 1, Flags: SYNACK, WND: windowA},
|
||||
WantState: tcp.StateSynRcvd,
|
||||
},
|
||||
3: { // A receives ACK from B.
|
||||
Incoming: &tcp.Segment{SEQ: issB, ACK: issA + 1, Flags: SYNACK, WND: windowA},
|
||||
WantState: tcp.StateEstablished,
|
||||
},
|
||||
}
|
||||
var tcbA tcp.ControlBlock
|
||||
tcbA.HelperInitState(tcp.StateSynSent, issA, issA, windowA)
|
||||
tcbA.HelperExchange(t, exchangeA)
|
||||
}
|
||||
|
||||
/*
|
||||
Recovery from Old Duplicate SYN
|
||||
TCP Peer A TCP Peer B
|
||||
|
||||
1. CLOSED LISTEN
|
||||
|
||||
2. SYN-SENT --> <SEQ=100><CTL=SYN> ...
|
||||
|
||||
3. (duplicate) ... <SEQ=90><CTL=SYN> --> SYN-RECEIVED
|
||||
|
||||
4. SYN-SENT <-- <SEQ=300><ACK=91><CTL=SYN,ACK> <-- SYN-RECEIVED
|
||||
|
||||
5. SYN-SENT --> <SEQ=91><CTL=RST> --> LISTEN
|
||||
|
||||
6. ... <SEQ=100><CTL=SYN> --> SYN-RECEIVED
|
||||
|
||||
7. ESTABLISHED <-- <SEQ=400><ACK=101><CTL=SYN,ACK> <-- SYN-RECEIVED
|
||||
|
||||
8. ESTABLISHED --> <SEQ=101><ACK=401><CTL=ACK> --> ESTABLISHED
|
||||
*/
|
||||
func TestExchange_rfc9293_figure8(t *testing.T) {
|
||||
const issA, issB, windowA, windowB = 100, 300, 1000, 1000
|
||||
const issAold = 90
|
||||
const issBNew = issB + 100
|
||||
exchangeA := []tcp.Exchange{
|
||||
0: { // A sends new SYN to B (which is not received).
|
||||
Outgoing: &tcp.Segment{SEQ: issA, Flags: tcp.FlagSYN, WND: windowA},
|
||||
WantState: tcp.StateSynSent,
|
||||
WantPeerState: tcp.StateSynRcvd,
|
||||
},
|
||||
1: { // Receive SYN from B acking an old "duplicate" SYN.
|
||||
Incoming: &tcp.Segment{SEQ: issB, ACK: issAold + 1, Flags: SYNACK, WND: windowB},
|
||||
WantState: tcp.StateSynSent,
|
||||
WantPending: &tcp.Segment{SEQ: issAold + 1, Flags: tcp.FlagRST, WND: windowA},
|
||||
WantPeerState: tcp.StateSynRcvd,
|
||||
},
|
||||
2: { // A sends RST to B and makes segment believable by using the old SEQ.
|
||||
Outgoing: &tcp.Segment{SEQ: issAold + 1, Flags: tcp.FlagRST, WND: windowA},
|
||||
WantState: tcp.StateSynSent,
|
||||
WantPeerState: tcp.StateListen,
|
||||
},
|
||||
3: { // A sends a duplicate SYN to B.
|
||||
Outgoing: &tcp.Segment{SEQ: issA, Flags: tcp.FlagSYN, WND: windowA},
|
||||
WantState: tcp.StateSynSent,
|
||||
WantPeerState: tcp.StateSynRcvd,
|
||||
},
|
||||
4: { // B SYNACKs new SYN.
|
||||
Incoming: &tcp.Segment{SEQ: issBNew, ACK: issA + 1, Flags: SYNACK, WND: windowB},
|
||||
WantState: tcp.StateEstablished,
|
||||
WantPending: &tcp.Segment{SEQ: issA + 1, ACK: issBNew + 1, Flags: tcp.FlagACK, WND: windowA},
|
||||
WantPeerState: tcp.StateSynRcvd,
|
||||
},
|
||||
5: { // B receives ACK from A.
|
||||
Outgoing: &tcp.Segment{SEQ: issA + 1, ACK: issBNew + 1, Flags: tcp.FlagACK, WND: windowA},
|
||||
WantState: tcp.StateEstablished,
|
||||
WantPeerState: tcp.StateEstablished,
|
||||
},
|
||||
}
|
||||
var tcbA tcp.ControlBlock
|
||||
tcbA.HelperInitState(tcp.StateSynSent, issA, issA, windowA)
|
||||
tcbA.HelperExchange(t, exchangeA)
|
||||
|
||||
exchangeB := []tcp.Exchange{
|
||||
0: { // B receives old SYN from A.
|
||||
Incoming: &tcp.Segment{SEQ: issAold, Flags: tcp.FlagSYN, WND: windowA},
|
||||
WantState: tcp.StateSynRcvd,
|
||||
WantPending: &tcp.Segment{SEQ: issB, ACK: issAold + 1, Flags: SYNACK, WND: windowB},
|
||||
},
|
||||
1: { // B SYNACKs old SYN.
|
||||
Outgoing: &tcp.Segment{SEQ: issB, ACK: issAold + 1, Flags: SYNACK, WND: windowB},
|
||||
WantState: tcp.StateSynRcvd,
|
||||
},
|
||||
2: { // B receives RST from A.
|
||||
Incoming: &tcp.Segment{SEQ: issAold + 1, Flags: tcp.FlagRST, WND: windowA},
|
||||
WantState: tcp.StateListen,
|
||||
},
|
||||
3: { // B receives new SYN from A.
|
||||
Incoming: &tcp.Segment{SEQ: issA, Flags: tcp.FlagSYN, WND: windowA},
|
||||
WantState: tcp.StateSynRcvd,
|
||||
WantPending: &tcp.Segment{SEQ: issBNew, ACK: issA + 1, Flags: SYNACK, WND: windowB},
|
||||
},
|
||||
4: { // B SYNACKs new SYN.
|
||||
Outgoing: &tcp.Segment{SEQ: issBNew, ACK: issA + 1, Flags: SYNACK, WND: windowB},
|
||||
WantState: tcp.StateSynRcvd,
|
||||
},
|
||||
5: { // B receives ACK from A.
|
||||
Incoming: &tcp.Segment{SEQ: issA + 1, ACK: issBNew + 1, Flags: tcp.FlagACK, WND: windowA},
|
||||
WantState: tcp.StateEstablished,
|
||||
},
|
||||
}
|
||||
var tcbB tcp.ControlBlock
|
||||
tcbB.HelperInitState(tcp.StateListen, issB, issB, windowB)
|
||||
tcbB.HelperExchange(t, exchangeB)
|
||||
}
|
||||
|
||||
/*
|
||||
Figure 12: Normal Close Sequence
|
||||
TCP Peer A TCP Peer B
|
||||
1. ESTABLISHED ESTABLISHED
|
||||
|
||||
2. (Close)
|
||||
FIN-WAIT-1 --> <SEQ=100><ACK=300><CTL=FIN,ACK> --> CLOSE-WAIT
|
||||
|
||||
3. FIN-WAIT-2 <-- <SEQ=300><ACK=101><CTL=ACK> <-- CLOSE-WAIT
|
||||
|
||||
4. (Close)
|
||||
TIME-WAIT <-- <SEQ=300><ACK=101><CTL=FIN,ACK> <-- LAST-ACK
|
||||
|
||||
5. TIME-WAIT --> <SEQ=101><ACK=301><CTL=ACK> --> CLOSED
|
||||
|
||||
6. (2 MSL)
|
||||
CLOSED
|
||||
*/
|
||||
func TestExchange_rfc9293_figure12(t *testing.T) {
|
||||
const issA, issB, windowA, windowB = 100, 300, 1000, 1000
|
||||
exchangeA := []tcp.Exchange{
|
||||
0: { // A sends FIN|ACK to B to begin closing connection.
|
||||
Outgoing: &tcp.Segment{SEQ: issA, ACK: issB, Flags: FINACK, WND: windowA},
|
||||
WantState: tcp.StateFinWait1,
|
||||
WantPeerState: tcp.StateCloseWait,
|
||||
},
|
||||
1: { // A receives ACK from B.
|
||||
Incoming: &tcp.Segment{SEQ: issB, ACK: issA + 1, Flags: tcp.FlagACK, WND: windowB},
|
||||
WantState: tcp.StateFinWait2,
|
||||
WantPeerState: tcp.StateCloseWait,
|
||||
// TODO(soypat): WantPending should be nil here? Perhaps fix test by modifying rcvFinWait1 pending result.
|
||||
WantPending: &tcp.Segment{SEQ: issA + 1, ACK: issB, Flags: tcp.FlagACK, WND: windowA},
|
||||
},
|
||||
2: { // A receives FIN|ACK from B.
|
||||
Incoming: &tcp.Segment{SEQ: issB, ACK: issA + 1, Flags: FINACK, WND: windowB},
|
||||
WantState: tcp.StateTimeWait,
|
||||
WantPending: &tcp.Segment{SEQ: issA + 1, ACK: issB + 1, Flags: tcp.FlagACK, WND: windowA},
|
||||
WantPeerState: tcp.StateLastAck,
|
||||
},
|
||||
3: { // A sends ACK to B.
|
||||
Outgoing: &tcp.Segment{SEQ: issA + 1, ACK: issB + 1, Flags: tcp.FlagACK, WND: windowA},
|
||||
WantState: tcp.StateTimeWait, // Technically we should be in TimeWait here.
|
||||
WantPeerState: tcp.StateClosed,
|
||||
},
|
||||
}
|
||||
var tcbA tcp.ControlBlock
|
||||
tcbA.HelperInitState(tcp.StateEstablished, issA, issA, windowA)
|
||||
tcbA.HelperInitRcv(issB, issB, windowB)
|
||||
tcbA.HelperExchange(t, exchangeA)
|
||||
// tcbA.HelperExchange(t, exchangeA[:1])
|
||||
// tcbA.HelperExchange(t, exchangeA[1:2])
|
||||
// tcbA.HelperExchange(t, exchangeA[2:])
|
||||
|
||||
return
|
||||
exchangeB := reverseExchange(exchangeA)
|
||||
exchangeB[1].WantPending = &tcp.Segment{SEQ: issB, ACK: issA + 1, Flags: FINACK, WND: windowB}
|
||||
var tcbB tcp.ControlBlock
|
||||
tcbB.HelperInitState(tcp.StateEstablished, issB, issB, windowB)
|
||||
tcbB.HelperInitRcv(issA, issA, windowA)
|
||||
tcbB.HelperExchange(t, exchangeB)
|
||||
}
|
||||
|
||||
/*
|
||||
Figure 12: Simultaneous Close Sequence
|
||||
TCP Peer A TCP Peer B
|
||||
|
||||
1. ESTABLISHED ESTABLISHED
|
||||
|
||||
2. (Close) (Close)
|
||||
FIN-WAIT-1 --> <SEQ=100><ACK=300><CTL=FIN,ACK> ... FIN-WAIT-1
|
||||
<-- <SEQ=300><ACK=100><CTL=FIN,ACK> <--
|
||||
... <SEQ=100><ACK=300><CTL=FIN,ACK> -->
|
||||
|
||||
3. CLOSING --> <SEQ=101><ACK=301><CTL=ACK> ... CLOSING
|
||||
<-- <SEQ=301><ACK=101><CTL=ACK> <--
|
||||
... <SEQ=101><ACK=301><CTL=ACK> -->
|
||||
|
||||
4. TIME-WAIT TIME-WAIT
|
||||
(2 MSL) (2 MSL)
|
||||
CLOSED CLOSED
|
||||
*/
|
||||
func TestExchange_rfc9293_figure13(t *testing.T) {
|
||||
const issA, issB, windowA, windowB = 100, 300, 1000, 1000
|
||||
exchangeA := []tcp.Exchange{
|
||||
0: { // A sends FIN|ACK to B to begin closing connection.
|
||||
Outgoing: &tcp.Segment{SEQ: issA, ACK: issB, Flags: FINACK, WND: windowA},
|
||||
WantState: tcp.StateFinWait1,
|
||||
},
|
||||
1: { // A receives FIN|ACK from B, who sent packet before receiving A's FINACK.
|
||||
Incoming: &tcp.Segment{SEQ: issB, ACK: issA, Flags: FINACK, WND: windowB},
|
||||
WantState: tcp.StateClosing,
|
||||
WantPending: &tcp.Segment{SEQ: issA + 1, ACK: issB + 1, Flags: tcp.FlagACK, WND: windowA},
|
||||
},
|
||||
2: { // A sends ACK to B.
|
||||
Outgoing: &tcp.Segment{SEQ: issA + 1, ACK: issB + 1, Flags: tcp.FlagACK, WND: windowA},
|
||||
WantState: tcp.StateTimeWait,
|
||||
},
|
||||
}
|
||||
var tcbA tcp.ControlBlock
|
||||
tcbA.HelperInitState(tcp.StateEstablished, issA, issA, windowA)
|
||||
tcbA.HelperInitRcv(issB, issB, windowB)
|
||||
tcbA.HelperExchange(t, exchangeA)
|
||||
|
||||
// No need to test B since exchange is completely symmetric.
|
||||
}
|
||||
|
||||
// Check no duplicate ack is sent during establishment.
|
||||
func TestExchange_noDupAckDuringEstablished(t *testing.T) {
|
||||
var tcbA tcp.ControlBlock
|
||||
const issA, issB, windowA, windowB = 300, 334222749, 256, 64240
|
||||
err := tcbA.Open(issA, issA, tcp.StateSynSent)
|
||||
tcbA.SetRecvWindow(windowA)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
establishA := []tcp.Exchange{
|
||||
0: { // B sends SYN to A.
|
||||
Incoming: &tcp.Segment{SEQ: issB, ACK: 0, WND: windowB, Flags: tcp.FlagSYN},
|
||||
WantPending: &tcp.Segment{SEQ: issA, ACK: issB + 1, WND: windowA, Flags: SYNACK},
|
||||
WantState: tcp.StateSynRcvd,
|
||||
},
|
||||
1: { // Send SYNACK to B.
|
||||
Outgoing: &tcp.Segment{SEQ: issA, ACK: issB + 1, WND: windowA, Flags: SYNACK},
|
||||
WantState: tcp.StateSynRcvd,
|
||||
},
|
||||
2: { // B ACKs SYNACK, thus establishing the connection on both sides.
|
||||
Incoming: &tcp.Segment{SEQ: issB + 1, ACK: issA + 1, WND: windowB, Flags: tcp.FlagACK},
|
||||
WantState: tcp.StateEstablished,
|
||||
},
|
||||
}
|
||||
tcbA.HelperExchange(t, establishA)
|
||||
if tcbA.State() != tcp.StateEstablished {
|
||||
t.Fatal("expected established state")
|
||||
}
|
||||
checkNoPending(t, &tcbA)
|
||||
const datasize = 5
|
||||
dataExA := []tcp.Exchange{
|
||||
0: { // B sends PSH|ACK to A with data.
|
||||
Incoming: &tcp.Segment{SEQ: issB + 1, ACK: issA + 1, WND: windowB, Flags: PSHACK, DATALEN: datasize},
|
||||
WantPending: &tcp.Segment{SEQ: issA + 1, ACK: issB + 1 + datasize, WND: windowA, Flags: tcp.FlagACK},
|
||||
WantState: tcp.StateEstablished,
|
||||
},
|
||||
1: { // A ACKs B's data.
|
||||
Outgoing: &tcp.Segment{SEQ: issA + 1, ACK: issB + 1 + datasize, WND: windowA, Flags: tcp.FlagACK},
|
||||
WantState: tcp.StateEstablished,
|
||||
},
|
||||
2: { // A sends PSH|ACK to B with data, same amount, as if echoing.
|
||||
Outgoing: &tcp.Segment{SEQ: issA + 1, ACK: issB + 1 + datasize, WND: windowA, Flags: PSHACK, DATALEN: datasize},
|
||||
WantState: tcp.StateEstablished,
|
||||
},
|
||||
// 3: { // B ACKs A's data.
|
||||
// Incoming: &tcp.Segment{SEQ: issB + 1 + datasize, ACK: issA + 1 + datasize, WND: windowB, Flags: tcp.FlagACK},
|
||||
// WantPending: nil,
|
||||
// WantState: tcp.StateEstablished,
|
||||
// },
|
||||
}
|
||||
tcbA.HelperExchange(t, dataExA)
|
||||
checkNoPending(t, &tcbA)
|
||||
tcbA.Recv(tcp.Segment{SEQ: issB + 1 + datasize, ACK: issA + 1 + datasize, WND: windowB, Flags: tcp.FlagACK})
|
||||
checkNoPending(t, &tcbA)
|
||||
}
|
||||
|
||||
// This test reenacts a full client-server interaction in the sending and receiving
|
||||
// of the 12 byte message "hello world\n" over TCP.
|
||||
func TestExchange_helloworld(t *testing.T) {
|
||||
// Client Transmission Control Block.
|
||||
var tcbA tcp.ControlBlock
|
||||
const windowA, windowB = 502, 4096
|
||||
const issA, issB = 0x5e722b7d, 0xbe6e4c0f
|
||||
const datalen = 12
|
||||
|
||||
exchangeA := []tcp.Exchange{
|
||||
0: { // A sends SYN to B.
|
||||
Outgoing: &tcp.Segment{SEQ: issA, Flags: tcp.FlagSYN, WND: windowA},
|
||||
WantState: tcp.StateSynSent,
|
||||
WantPeerState: tcp.StateSynRcvd,
|
||||
},
|
||||
1: { // A receives SYNACK from B.
|
||||
Incoming: &tcp.Segment{SEQ: issB, ACK: issA + 1, Flags: SYNACK, WND: windowB},
|
||||
WantState: tcp.StateEstablished,
|
||||
WantPending: &tcp.Segment{SEQ: issA + 1, ACK: issB + 1, Flags: tcp.FlagACK, WND: windowA},
|
||||
WantPeerState: tcp.StateSynRcvd,
|
||||
},
|
||||
2: { // A sends ACK to B thus establishing connection.
|
||||
Outgoing: &tcp.Segment{SEQ: issA + 1, ACK: issB + 1, Flags: tcp.FlagACK, WND: windowA},
|
||||
WantState: tcp.StateEstablished,
|
||||
WantPeerState: tcp.StateEstablished,
|
||||
},
|
||||
3: { // A sends PSH|ACK to B with 12 byte message: "hello world\n"
|
||||
Outgoing: &tcp.Segment{SEQ: issA + 1, ACK: issB + 1, Flags: PSHACK, WND: windowA, DATALEN: datalen},
|
||||
WantState: tcp.StateEstablished,
|
||||
WantPeerState: tcp.StateEstablished,
|
||||
},
|
||||
4: { // A receives ACK from B of last message.
|
||||
Incoming: &tcp.Segment{SEQ: issB + 1, ACK: issA + 1 + datalen, Flags: tcp.FlagACK, WND: windowB},
|
||||
WantState: tcp.StateEstablished,
|
||||
WantPeerState: tcp.StateEstablished,
|
||||
},
|
||||
5: { // A receives PSH|ACK from B with echoed 12 byte message: "hello world\n"
|
||||
Incoming: &tcp.Segment{SEQ: issB + 1, ACK: issA + 1 + datalen, Flags: PSHACK, WND: windowB, DATALEN: datalen},
|
||||
WantState: tcp.StateEstablished,
|
||||
WantPending: &tcp.Segment{SEQ: issA + 1 + datalen, ACK: issB + 1 + datalen, Flags: tcp.FlagACK, WND: windowA},
|
||||
WantPeerState: tcp.StateEstablished,
|
||||
},
|
||||
6: { // A ACKs B's message.
|
||||
Outgoing: &tcp.Segment{SEQ: issA + 1 + datalen, ACK: issB + 1 + datalen, Flags: tcp.FlagACK, WND: windowA},
|
||||
WantState: tcp.StateEstablished,
|
||||
WantPeerState: tcp.StateEstablished,
|
||||
},
|
||||
7: { // A sends PSH|ACK to B with SECOND 12 byte message.
|
||||
Outgoing: &tcp.Segment{SEQ: issA + 1 + datalen, ACK: issB + 1 + datalen, Flags: PSHACK, WND: windowA, DATALEN: datalen},
|
||||
WantState: tcp.StateEstablished,
|
||||
WantPeerState: tcp.StateEstablished,
|
||||
},
|
||||
8: { // A receives PSH|ACK that acks last message and contains echoed of SECOND 12 byte message.
|
||||
Incoming: &tcp.Segment{SEQ: issB + 1 + datalen, ACK: issA + 1 + 2*datalen, Flags: PSHACK, WND: windowB, DATALEN: datalen},
|
||||
WantState: tcp.StateEstablished,
|
||||
WantPending: &tcp.Segment{SEQ: issA + 1 + 2*datalen, ACK: issB + 1 + 2*datalen, Flags: tcp.FlagACK, WND: windowA},
|
||||
WantPeerState: tcp.StateEstablished,
|
||||
},
|
||||
9: { // A ACKs B's SECOND message.
|
||||
Outgoing: &tcp.Segment{SEQ: issA + 1 + 2*datalen, ACK: issB + 1 + 2*datalen, Flags: tcp.FlagACK, WND: windowA},
|
||||
WantState: tcp.StateEstablished,
|
||||
WantPeerState: tcp.StateEstablished,
|
||||
},
|
||||
10: { // A sends FIN|ACK to B to close connection.
|
||||
Outgoing: &tcp.Segment{SEQ: issA + 1 + 2*datalen, ACK: issB + 1 + 2*datalen, Flags: FINACK, WND: windowA},
|
||||
WantState: tcp.StateFinWait1,
|
||||
WantPeerState: tcp.StateCloseWait,
|
||||
},
|
||||
11: { // A receives B's ACK of FIN.
|
||||
Incoming: &tcp.Segment{SEQ: issB + 1 + 2*datalen, ACK: issA + 2 + 2*datalen, Flags: tcp.FlagACK, WND: windowB},
|
||||
WantState: tcp.StateFinWait2,
|
||||
WantPending: &tcp.Segment{SEQ: issA + 2 + 2*datalen, ACK: issB + 1 + 2*datalen, Flags: tcp.FlagACK, WND: windowA},
|
||||
WantPeerState: tcp.StateCloseWait,
|
||||
},
|
||||
}
|
||||
// The client starts in the SYN_SENT state with a random sequence number.
|
||||
gotServerSeg, _ := parseSegment(t, exchangeHelloWorld[0])
|
||||
tcbA.HelperInitState(tcp.StateSynSent, gotServerSeg.SEQ, gotServerSeg.SEQ, windowB)
|
||||
tcbA.HelperExchange(t, exchangeA)
|
||||
|
||||
// TODO(soypat): fix exchange reversal.
|
||||
return
|
||||
exchangeB := reverseExchange(exchangeA)
|
||||
|
||||
exchangeB[7].WantPending = nil // Is an unpredicable action.
|
||||
var tcbB tcp.ControlBlock
|
||||
tcbB.HelperInitState(tcp.StateListen, issB, issB, windowB)
|
||||
tcbB.HelperInitRcv(issA, issA, windowA)
|
||||
tcbB.HelperExchange(t, exchangeB)
|
||||
}
|
||||
|
||||
func TestResetEstablished(t *testing.T) {
|
||||
var tcb tcp.ControlBlock
|
||||
const windowA, windowB = 502, 4096
|
||||
const issA, issB = 0x5e722b7d, 0xbe6e4c0f
|
||||
tcb.HelperInitState(tcp.StateEstablished, issA, issA, windowA)
|
||||
tcb.HelperInitRcv(issB, issB, windowB)
|
||||
|
||||
err := tcb.Recv(tcp.Segment{SEQ: issB, ACK: issA, Flags: tcp.FlagRST, WND: windowB})
|
||||
if err == nil {
|
||||
t.Fatal("expected error")
|
||||
}
|
||||
if tcb.State() != tcp.StateClosed {
|
||||
t.Error("expected closed state; got ", tcb.State().String())
|
||||
}
|
||||
checkNoPending(t, &tcb)
|
||||
}
|
||||
|
||||
func TestFinackClose(t *testing.T) {
|
||||
var tcb tcp.ControlBlock
|
||||
const windowA, windowB = 502, 4096
|
||||
const issA, issB = 100, 200
|
||||
tcb.HelperInitState(tcp.StateEstablished, issA, issA, windowA)
|
||||
tcb.HelperInitRcv(issB, issB, windowB)
|
||||
// Start closing process.
|
||||
err := tcb.Close()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
seg, ok := tcb.PendingSegment(0)
|
||||
if !ok {
|
||||
t.Fatal("expected pending segment")
|
||||
}
|
||||
if !seg.Flags.HasAll(tcp.FlagFIN | tcp.FlagACK) {
|
||||
t.Fatalf("expected FIN|ACK; got %s", seg.Flags.String())
|
||||
}
|
||||
err = tcb.Send(seg)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if tcb.State() != tcp.StateFinWait1 {
|
||||
t.Fatalf("expected FinWait1; got %s", tcb.State().String())
|
||||
}
|
||||
// Special case where we receive FINACK all together, we can streamline and go into TimeWait.
|
||||
err = tcb.Recv(tcp.Segment{
|
||||
SEQ: issB,
|
||||
ACK: issA + 1,
|
||||
WND: windowB,
|
||||
Flags: FINACK,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if tcb.State() != tcp.StateTimeWait {
|
||||
t.Fatalf("expected TimeWait after FINACK; got %s", tcb.State().String())
|
||||
}
|
||||
}
|
||||
|
||||
func TestExchange_helloworld_client(t *testing.T) {
|
||||
return
|
||||
// Client Transmission Control Block.
|
||||
var tcb tcp.ControlBlock
|
||||
// The client starts in the SYN_SENT state with a random sequence number.
|
||||
gotClientSeg, _ := parseSegment(t, exchangeHelloWorld[0])
|
||||
|
||||
// We add the SYN state to the client.
|
||||
tcb.HelperInitState(tcp.StateSynSent, gotClientSeg.SEQ, gotClientSeg.SEQ, gotClientSeg.WND)
|
||||
err := tcb.Send(gotClientSeg)
|
||||
if err != nil {
|
||||
|
||||
t.Fatal(err)
|
||||
}
|
||||
tcb.HelperPrintSegment(t, false, gotClientSeg)
|
||||
|
||||
segString := func(seg tcp.Segment) string {
|
||||
return tcb.RelativeAutoSegment(seg).RelativeGoString(0, 0)
|
||||
}
|
||||
for i, packet := range exchangeHelloWorld {
|
||||
if i == 0 {
|
||||
continue // we already processed first packet.
|
||||
}
|
||||
seg, payload := parseSegment(t, packet)
|
||||
if seg.DATALEN > 0 {
|
||||
t.Logf("seg[%d] <%s> payload: %q", i, tcb.State(), string(payload))
|
||||
} else {
|
||||
t.Logf("seg[%d] <%s>", i, tcb.State())
|
||||
}
|
||||
isClient := packet[0] == 0x28
|
||||
if isClient {
|
||||
isPSH := seg.Flags&tcp.FlagPSH != 0
|
||||
gotClientSeg.Flags |= seg.Flags & (tcp.FlagPSH | tcp.FlagFIN) // Can't predict when client will send FIN.
|
||||
if isPSH {
|
||||
gotClientSeg.DATALEN = seg.DATALEN
|
||||
}
|
||||
|
||||
gotClientSeg.WND = seg.WND // Ignore window field, not a core part of control flow.
|
||||
if gotClientSeg != seg {
|
||||
t.Fatalf("client:\n got=%+v\nwant=%+v", segString(gotClientSeg), segString(seg))
|
||||
}
|
||||
err := tcb.Send(gotClientSeg)
|
||||
if err != nil {
|
||||
t.Fatalf("incoming %s:\nseg[%d]=%s\nrcv=%+v\nsnd=%+v", err, i, segString(gotClientSeg), tcb.RelativeRecvSpace(), tcb.RelativeSendSpace())
|
||||
}
|
||||
tcb.HelperPrintSegment(t, false, gotClientSeg)
|
||||
continue // we only pass server packets to the client.
|
||||
}
|
||||
err = tcb.Recv(seg)
|
||||
if err != nil {
|
||||
t.Fatalf("%s:\nseg[%d]=%s\nrcv=%+v\nsnd=%+v", err, i, segString(seg), tcb.RelativeRecvSpace(), tcb.RelativeSendSpace())
|
||||
}
|
||||
tcb.HelperPrintSegment(t, true, seg)
|
||||
var ok bool
|
||||
gotClientSeg, ok = tcb.PendingSegment(0)
|
||||
if !ok {
|
||||
t.Fatalf("[%d]: got no segment state=%s", i, tcb.State())
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func parseSegment(t *testing.T, b []byte) (tcp.Segment, []byte) {
|
||||
t.Helper()
|
||||
efrm, err := lneto.NewEthFrame(b)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if efrm.EtherTypeOrSize() != lneto.EtherTypeIPv4 {
|
||||
t.Fatalf("not IPv4")
|
||||
}
|
||||
err = efrm.ValidateSize()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
ifrm, err := lneto.NewIPv4Frame(efrm.Payload())
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if ifrm.Protocol() != 6 {
|
||||
t.Fatalf("not TCP")
|
||||
}
|
||||
v, _ := ifrm.VersionAndIHL()
|
||||
if v != 4 {
|
||||
t.Fatal("invalid IP version", v)
|
||||
}
|
||||
err = ifrm.ValidateSize()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
ipl := ifrm.Payload()
|
||||
tfrm, err := lneto.NewTCPFrame(ipl)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
} else if err = tfrm.ValidateSize(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
_ = tfrm.String()
|
||||
payload := tfrm.Payload()
|
||||
return tfrm.Segment(len(payload)), payload
|
||||
}
|
||||
|
||||
func reverseExchange(exchange []tcp.Exchange) []tcp.Exchange {
|
||||
if len(exchange) == 0 {
|
||||
panic("len(exchange) != len(states) or empty exchange: " + strconv.Itoa(len(exchange)))
|
||||
}
|
||||
firstIsIn := exchange[0].Incoming != nil
|
||||
if firstIsIn {
|
||||
panic("please start with an outgoing segment to reverse exchange for best test results")
|
||||
}
|
||||
out := make([]tcp.Exchange, len(exchange))
|
||||
for i := range exchange {
|
||||
isLast := i == len(exchange)-1
|
||||
isOut := exchange[i].Outgoing != nil
|
||||
out[i].WantState, out[i].WantPeerState = exchange[i].WantPeerState, exchange[i].WantState
|
||||
if isOut {
|
||||
out[i].Incoming = exchange[i].Outgoing
|
||||
if !isLast {
|
||||
out[i].WantPending = exchange[i+1].Incoming
|
||||
}
|
||||
} else {
|
||||
out[i].Outgoing = exchange[i].Incoming
|
||||
}
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
func checkNoPending(t *testing.T, tcb *tcp.ControlBlock) bool {
|
||||
t.Helper()
|
||||
// We extensively test the API for inadvertent state modification in a HasPending or PendingSegment call.
|
||||
hasPD := tcb.HasPending()
|
||||
pd, ok := tcb.PendingSegment(0)
|
||||
hasPD2 := tcb.HasPending()
|
||||
if hasPD || ok || hasPD2 {
|
||||
t.Errorf("unexpected pending segment: %+v (%v,%v,%v)", pd, hasPD, ok, hasPD2)
|
||||
return false
|
||||
}
|
||||
if hasPD != ok || hasPD != hasPD2 {
|
||||
t.Fatalf("inconsistent pending segment: (%v,%v,%v)", hasPD, ok, hasPD2)
|
||||
}
|
||||
if !ok && pd != (tcp.Segment{}) {
|
||||
t.Fatalf("inconsistent pending segment: %+v (%v,%v,%v)", pd, hasPD, ok, hasPD2)
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
// Full client-server interaction in the sending of "hello world" over TCP in order.
|
||||
var exchangeHelloWorld = [][]byte{
|
||||
// client SYN1
|
||||
0: []byte("\x28\xcd\xc1\x05\x4d\xbb\xd8\x5e\xd3\x43\x03\xeb\x08\x00\x45\x00\x00\x3c\x71\xac\x40\x00\x40\x06\x44\x9b\xc0\xa8\x01\x93\xc0\xa8\x01\x91\x84\x96\x04\xd2\x5e\x72\x2b\x7d\x00\x00\x00\x00\xa0\x02\xfa\xf0\x27\x6d\x00\x00\x02\x04\x05\xb4\x04\x02\x08\x0a\x07\x8b\x86\x4a\x00\x00\x00\x00\x01\x03\x03\x07"),
|
||||
// server SYNACK
|
||||
1: []byte("\xd8\x5e\xd3\x43\x03\xeb\x28\xcd\xc1\x05\x4d\xbb\x08\x00\x45\x00\x00\x34\x00\x00\x40\x00\x40\x06\xb6\x4f\xc0\xa8\x01\x91\xc0\xa8\x01\x93\x04\xd2\x84\x96\xbe\x6e\x4c\x0f\x5e\x72\x2b\x7e\x80\x12\x10\x00\xc0\xbb\x00\x00\x02\x04\x05\xb4\x03\x03\x00\x04\x02\x00\x00\x00"),
|
||||
// client ACK1
|
||||
2: []byte("\x28\xcd\xc1\x05\x4d\xbb\xd8\x5e\xd3\x43\x03\xeb\x08\x00\x45\x00\x00\x28\x71\xad\x40\x00\x40\x06\x44\xae\xc0\xa8\x01\x93\xc0\xa8\x01\x91\x84\x96\x04\xd2\x5e\x72\x2b\x7e\xbe\x6e\x4c\x10\x50\x10\x01\xf6\x0b\x92\x00\x00"),
|
||||
// client PSHACK0
|
||||
3: []byte("\x28\xcd\xc1\x05\x4d\xbb\xd8\x5e\xd3\x43\x03\xeb\x08\x00\x45\x00\x00\x34\x71\xae\x40\x00\x40\x06\x44\xa1\xc0\xa8\x01\x93\xc0\xa8\x01\x91\x84\x96\x04\xd2\x5e\x72\x2b\x7e\xbe\x6e\x4c\x10\x50\x18\x01\xf6\x79\xa5\x00\x00\x68\x65\x6c\x6c\x6f\x20\x77\x6f\x72\x6c\x64\x0a"),
|
||||
// server ACK1
|
||||
4: []byte("\xd8\x5e\xd3\x43\x03\xeb\x28\xcd\xc1\x05\x4d\xbb\x08\x00\x45\x00\x00\x28\x00\x00\x40\x00\x40\x06\xb6\x5b\xc0\xa8\x01\x91\xc0\xa8\x01\x93\x04\xd2\x84\x96\xbe\x6e\x4c\x10\x5e\x72\x2b\x8a\x50\x10\x0f\xf4\xfd\x87\x00\x00\x00\x00\x00\x00\x00\x00"),
|
||||
// server PSHACK1
|
||||
5: []byte("\xd8\x5e\xd3\x43\x03\xeb\x28\xcd\xc1\x05\x4d\xbb\x08\x00\x45\x00\x00\x34\x00\x00\x40\x00\x40\x06\xb6\x4f\xc0\xa8\x01\x91\xc0\xa8\x01\x93\x04\xd2\x84\x96\xbe\x6e\x4c\x10\x5e\x72\x2b\x8a\x50\x18\x10\x00\x6b\x8f\x00\x00\x68\x65\x6c\x6c\x6f\x20\x77\x6f\x72\x6c\x64\x0a"),
|
||||
// client ACK2
|
||||
6: []byte("\x28\xcd\xc1\x05\x4d\xbb\xd8\x5e\xd3\x43\x03\xeb\x08\x00\x45\x00\x00\x28\x71\xaf\x40\x00\x40\x06\x44\xac\xc0\xa8\x01\x93\xc0\xa8\x01\x91\x84\x96\x04\xd2\x5e\x72\x2b\x8a\xbe\x6e\x4c\x1c\x50\x10\x01\xf6\x0b\x7a\x00\x00"),
|
||||
// client PSHACK1
|
||||
7: []byte("\x28\xcd\xc1\x05\x4d\xbb\xd8\x5e\xd3\x43\x03\xeb\x08\x00\x45\x00\x00\x34\x71\xb0\x40\x00\x40\x06\x44\x9f\xc0\xa8\x01\x93\xc0\xa8\x01\x91\x84\x96\x04\xd2\x5e\x72\x2b\x8a\xbe\x6e\x4c\x1c\x50\x18\x01\xf6\x79\x8d\x00\x00\x68\x65\x6c\x6c\x6f\x20\x77\x6f\x72\x6c\x64\x0a"),
|
||||
// server PSHACK2
|
||||
8: []byte("\xd8\x5e\xd3\x43\x03\xeb\x28\xcd\xc1\x05\x4d\xbb\x08\x00\x45\x00\x00\x34\x00\x00\x40\x00\x40\x06\xb6\x4f\xc0\xa8\x01\x91\xc0\xa8\x01\x93\x04\xd2\x84\x96\xbe\x6e\x4c\x1c\x5e\x72\x2b\x96\x50\x18\x10\x00\x6b\x77\x00\x00\x68\x65\x6c\x6c\x6f\x20\x77\x6f\x72\x6c\x64\x0a"),
|
||||
// client ACK3
|
||||
9: []byte("\x28\xcd\xc1\x05\x4d\xbb\xd8\x5e\xd3\x43\x03\xeb\x08\x00\x45\x00\x00\x28\x71\xb1\x40\x00\x40\x06\x44\xaa\xc0\xa8\x01\x93\xc0\xa8\x01\x91\x84\x96\x04\xd2\x5e\x72\x2b\x96\xbe\x6e\x4c\x28\x50\x10\x01\xf6\x0b\x62\x00\x00"),
|
||||
// client FINACK
|
||||
10: []byte("\x28\xcd\xc1\x05\x4d\xbb\xd8\x5e\xd3\x43\x03\xeb\x08\x00\x45\x00\x00\x28\x71\xb2\x40\x00\x40\x06\x44\xa9\xc0\xa8\x01\x93\xc0\xa8\x01\x91\x84\x96\x04\xd2\x5e\x72\x2b\x96\xbe\x6e\x4c\x28\x50\x11\x01\xf6\x0b\x61\x00\x00"),
|
||||
// server ACK
|
||||
11: []byte("\xd8\x5e\xd3\x43\x03\xeb\x28\xcd\xc1\x05\x4d\xbb\x08\x00\x45\x00\x00\x28\x00\x00\x40\x00\x40\x06\xb6\x5b\xc0\xa8\x01\x91\xc0\xa8\x01\x93\x04\xd2\x84\x96\xbe\x6e\x4c\x28\x5e\x72\x2b\x97\x50\x10\x10\x00\xfd\x56\x00\x00\x00\x00\x00\x00\x00\x00"),
|
||||
}
|
||||
|
||||
func TestUnexpectedStateClosing(t *testing.T) {
|
||||
// TCB is a server which returns an HTTP response and receives a FINACK.
|
||||
var tcb tcp.ControlBlock
|
||||
const httpLen = 1192
|
||||
const issA, issB, windowA, windowB = 1, 127, 2000, 2000
|
||||
tcb.HelperInitState(tcp.StateEstablished, issA, issA, windowA)
|
||||
tcb.HelperInitRcv(issB, issB, windowB)
|
||||
|
||||
ex := []tcp.Exchange{
|
||||
0: { // Server sends HTTP response.
|
||||
Outgoing: &tcp.Segment{SEQ: issA, ACK: issB, Flags: PSHACK, WND: windowA, DATALEN: httpLen},
|
||||
WantState: tcp.StateEstablished,
|
||||
},
|
||||
1: { // Client sends an ACK to server.
|
||||
Incoming: &tcp.Segment{SEQ: issB, ACK: issA + httpLen, Flags: tcp.FlagACK, WND: windowB},
|
||||
WantState: tcp.StateEstablished,
|
||||
},
|
||||
2: { // Client sends FIN|ACK to server.
|
||||
Incoming: &tcp.Segment{SEQ: issB, ACK: issA + httpLen, Flags: FINACK, WND: windowB},
|
||||
WantPending: &tcp.Segment{SEQ: issA + httpLen, ACK: issB + 1, Flags: tcp.FlagACK, WND: windowA},
|
||||
WantState: tcp.StateCloseWait,
|
||||
},
|
||||
3: { // Server sends out FINACK.
|
||||
Outgoing: &tcp.Segment{SEQ: issA + httpLen, ACK: issB + 1, Flags: FINACK, WND: windowA},
|
||||
WantState: tcp.StateLastAck,
|
||||
},
|
||||
4: { // Client sends back ACK.
|
||||
Incoming: &tcp.Segment{SEQ: issB + 1, ACK: issA + httpLen + 1, Flags: tcp.FlagACK, WND: windowB},
|
||||
WantState: tcp.StateClosed,
|
||||
},
|
||||
}
|
||||
tcb.HelperExchange(t, ex[:])
|
||||
}
|
||||
|
||||
// This corresponds to https://github.com/soypat/seqs/issues/19
|
||||
// The bug consisted of a panic condition encountered when using wget client with a seqs based server.
|
||||
// Thanks to @knieriem for finding this and the detailed report they submitted.
|
||||
func TestIssue19(t *testing.T) {
|
||||
var tcb tcp.ControlBlock
|
||||
assertState := func(state tcp.State) {
|
||||
t.Helper()
|
||||
if tcb.State() != state {
|
||||
t.Fatalf("want state %s; got %s", state.String(), tcb.State().String())
|
||||
}
|
||||
}
|
||||
const httpLen = 1192
|
||||
const issA, issB, windowA, windowB = 1, 0, 2000, 2000
|
||||
tcb.HelperInitState(tcp.StateEstablished, issA, issA, windowA)
|
||||
tcb.HelperInitRcv(issB, issB, windowB)
|
||||
|
||||
// Send out HTTP request and close connection.
|
||||
err := tcb.Send(tcp.Segment{SEQ: issA, ACK: issB, Flags: PSHACK, WND: windowA, DATALEN: httpLen})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
err = tcb.Close()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
assertState(tcp.StateEstablished)
|
||||
|
||||
pending, ok := tcb.PendingSegment(0)
|
||||
if !ok {
|
||||
t.Fatal("expected pending segment")
|
||||
} else if pending.Flags != FINACK {
|
||||
t.Fatalf("expected FINACK; got %s", pending.Flags.String())
|
||||
}
|
||||
|
||||
// Receive ACK of HTTP segment.
|
||||
err = tcb.Recv(tcp.Segment{SEQ: issB, ACK: issA + httpLen, Flags: tcp.FlagACK, WND: windowB})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
assertState(tcp.StateEstablished)
|
||||
err = tcb.Close()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
pending, ok = tcb.PendingSegment(0)
|
||||
if !ok {
|
||||
t.Fatal("expected pending segment")
|
||||
} else if pending.Flags != FINACK {
|
||||
t.Fatalf("expected FINACK; got %s", pending.Flags.String())
|
||||
}
|
||||
|
||||
// Send out FINACK.
|
||||
err = tcb.Send(pending)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
assertState(tcp.StateFinWait1)
|
||||
|
||||
// Receive FINACK response from client.
|
||||
err = tcb.Recv(tcp.Segment{SEQ: issB, ACK: issA + httpLen, Flags: FINACK, WND: windowB})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
assertState(tcp.StateClosing)
|
||||
pending, ok = tcb.PendingSegment(0)
|
||||
if !ok {
|
||||
t.Fatal("expected pending segment")
|
||||
} else if pending.Flags != tcp.FlagACK {
|
||||
t.Fatalf("expected ACK; got %s", pending.Flags.String())
|
||||
}
|
||||
|
||||
// Before responding we receive an ACK from client. This is where panic is triggered.
|
||||
err = tcb.Recv(tcp.Segment{SEQ: issB + 1, ACK: issA + httpLen + 1, Flags: tcp.FlagACK, WND: windowB})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
assertState(tcp.StateTimeWait)
|
||||
|
||||
// Check we still need to send an ACK.
|
||||
pending, ok = tcb.PendingSegment(0)
|
||||
if !ok {
|
||||
t.Fatal("expected pending segment")
|
||||
} else if pending.Flags != tcp.FlagACK {
|
||||
t.Fatalf("expected ACK; got %s", pending.Flags.String())
|
||||
}
|
||||
// Prepare response to client.
|
||||
err = tcb.Send(pending)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
|
||||
func FuzzTCBActions(f *testing.F) {
|
||||
const mtu = 2048
|
||||
const (
|
||||
actionRecv = iota
|
||||
actionSend
|
||||
actionClose
|
||||
actionMax
|
||||
)
|
||||
f.Add(
|
||||
0x2313_2313,
|
||||
[]byte{actionSend, actionRecv, actionSend, actionRecv, actionSend, actionRecv},
|
||||
)
|
||||
f.Add(
|
||||
0x2fefe_feefe,
|
||||
[]byte{actionSend, actionRecv, actionSend, actionClose, actionSend, actionRecv},
|
||||
)
|
||||
f.Add(
|
||||
0x2fefe_feefe,
|
||||
[]byte{actionClose, actionRecv, actionSend, actionClose, actionSend, actionRecv},
|
||||
)
|
||||
recvsendSize := func(rng *rand.Rand) int {
|
||||
return rng.Int() % mtu
|
||||
}
|
||||
f.Fuzz(func(t *testing.T, seed int, actions []byte) {
|
||||
if len(actions) == 0 || len(actions) > 100 {
|
||||
t.SkipNow()
|
||||
}
|
||||
rng := rand.New(rand.NewSource(int64(seed)))
|
||||
var clientISS tcp.Value = tcp.Value(rng.Int31())
|
||||
var serverISS tcp.Value = tcp.Value(rng.Int31())
|
||||
|
||||
var client tcp.ControlBlock
|
||||
client.HelperInitState(tcp.StateEstablished, clientISS, clientISS, mtu)
|
||||
client.HelperInitRcv(serverISS, serverISS, mtu)
|
||||
|
||||
var server tcp.ControlBlock
|
||||
server.HelperInitState(tcp.StateEstablished, serverISS, serverISS, mtu)
|
||||
server.HelperInitRcv(clientISS, clientISS, mtu)
|
||||
var closeCalled bool
|
||||
// logger := slog.New(slog.NewTextHandler(os.Stdout, &slog.HandlerOptions{
|
||||
// Level: slog.LevelDebug - 2,
|
||||
// }))
|
||||
// client.SetLogger(logger.WithGroup("client"))
|
||||
// server.SetLogger(logger.WithGroup("server"))
|
||||
// var exchanges []tcp.Exchange
|
||||
// hasPanicked := true
|
||||
// defer func() {
|
||||
// if hasPanicked {
|
||||
// for _, ex := range exchanges {
|
||||
// t.Log(ex.RFC9293String(tcp.StateEstablished, tcp.StateEstablished))
|
||||
// }
|
||||
// }
|
||||
// }()
|
||||
for _, action := range actions {
|
||||
v := recvsendSize(rng)
|
||||
switch action % actionMax {
|
||||
case actionSend:
|
||||
seg, ok := client.PendingSegment(v % mtu)
|
||||
if ok {
|
||||
// exchanges = append(exchanges, tcp.Exchange{Outgoing: &seg})
|
||||
err := client.Send(seg)
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
err = server.Recv(seg)
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
}
|
||||
case actionRecv:
|
||||
seg, ok := server.PendingSegment(v % mtu)
|
||||
if ok {
|
||||
// exchanges = append(exchanges, tcp.Exchange{Incoming: &seg})
|
||||
err := server.Send(seg)
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
err = client.Recv(seg)
|
||||
if err != nil && !closeCalled {
|
||||
panic(err)
|
||||
}
|
||||
}
|
||||
case actionClose:
|
||||
err := client.Close()
|
||||
if err != nil && !closeCalled {
|
||||
panic(err)
|
||||
}
|
||||
closeCalled = true
|
||||
return
|
||||
}
|
||||
}
|
||||
// hasPanicked = false
|
||||
})
|
||||
}
|
||||
+143
@@ -0,0 +1,143 @@
|
||||
package tcp
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"time"
|
||||
|
||||
"github.com/soypat/lneto/internal"
|
||||
)
|
||||
|
||||
func newRingTx(buf []byte, maxQueuedPackets int) *ringTx {
|
||||
if maxQueuedPackets <= 0 || len(buf) < 2 || len(buf) < maxQueuedPackets {
|
||||
panic("invalid argument to NewRingTx")
|
||||
}
|
||||
return &ringTx{
|
||||
rawbuf: buf,
|
||||
packets: make([]ringidx, maxQueuedPackets),
|
||||
}
|
||||
}
|
||||
|
||||
// ringTx is a ring buffer with retransmission queue functionality added.
|
||||
type ringTx struct {
|
||||
// rawbuf contains the ring buffer of ordered bytes. It should be the size of the window.
|
||||
rawbuf []byte
|
||||
// packets contains
|
||||
packets []ringidx
|
||||
// firstPkt is the index of the oldest packet in the packets field.
|
||||
firstPkt int
|
||||
lastPkt int
|
||||
// unsentOff is the offset of start of unsent data into rawbuf.
|
||||
unsentoff int
|
||||
// unsentend is the offset of end of unsent data in rawbuf.
|
||||
unsentend int
|
||||
}
|
||||
|
||||
// ringidx represents packet data inside RingTx
|
||||
type ringidx struct {
|
||||
// off is data start offset of packet data inside buf.
|
||||
off int
|
||||
// end is the ringed data end offset, non-inclusive.
|
||||
end int
|
||||
// seq is the sequence number of the packet.
|
||||
seq Value
|
||||
t time.Time
|
||||
// acked flags if this packet has been acknowledged. Useful for SACK (selective acknowledgement)
|
||||
// acked bool
|
||||
}
|
||||
|
||||
// Buffered returns the amount of unsent bytes.
|
||||
func (tx *ringTx) Buffered() int {
|
||||
r := tx.unsentRing()
|
||||
return r.Buffered()
|
||||
}
|
||||
|
||||
// BufferedSent returns the total amount of bytes sent but not acked.
|
||||
func (tx *ringTx) BufferedSent() int {
|
||||
r := tx.sentRing()
|
||||
return r.Buffered()
|
||||
}
|
||||
|
||||
// Write writes data to the underlying unsent data ring buffer.
|
||||
func (tx *ringTx) Write(b []byte) (int, error) {
|
||||
first := tx.packets[tx.firstPkt]
|
||||
r := tx.unsentRing()
|
||||
if first.off < 0 {
|
||||
// No packets in queue case.
|
||||
return r.Write(b)
|
||||
}
|
||||
return r.WriteLimited(b, first.off)
|
||||
}
|
||||
|
||||
// ReadPacket reads from the unsent data ring buffer and generates a new packet segment.
|
||||
// It fails if the sent packet queue is full.
|
||||
func (tx *ringTx) NewPacketAndRead(b []byte) (int, error) {
|
||||
nxtpkt := (tx.lastPkt + 1) % len(tx.packets)
|
||||
if tx.firstPkt == nxtpkt {
|
||||
return 0, errors.New("packet queue full")
|
||||
}
|
||||
|
||||
r := tx.unsentRing()
|
||||
start := r.Off
|
||||
n, err := r.Read(b)
|
||||
if err != nil {
|
||||
return n, err
|
||||
}
|
||||
last := &tx.packets[tx.lastPkt]
|
||||
rlast := tx.packetRing(tx.lastPkt)
|
||||
tx.packets[nxtpkt].off = start
|
||||
tx.packets[nxtpkt].end = r.Off
|
||||
tx.packets[nxtpkt].seq = last.seq + Value(rlast.Buffered())
|
||||
tx.lastPkt = nxtpkt
|
||||
tx.unsentoff = r.Off
|
||||
return n, nil
|
||||
}
|
||||
|
||||
// IsQueueFull returns true if the sent packet queue is full in which
|
||||
// case a call to ReadPacket is guaranteed to fail.
|
||||
func (tx *ringTx) IsQueueFull() bool {
|
||||
return tx.firstPkt == (tx.lastPkt+1)%len(tx.packets)
|
||||
}
|
||||
|
||||
func (tx *ringTx) packetRing(i int) internal.Ring {
|
||||
pkt := tx.packets[i]
|
||||
if pkt.off < 0 {
|
||||
return internal.Ring{}
|
||||
}
|
||||
return tx.ring(pkt.off, pkt.end)
|
||||
}
|
||||
|
||||
// RecvSegment processes an incoming segment and updates the sent packet queue
|
||||
func (tx *ringTx) RecvACK(ack Value) error {
|
||||
i := tx.firstPkt
|
||||
for {
|
||||
pkt := &tx.packets[i]
|
||||
if ack >= pkt.seq {
|
||||
// Packet was received by remote. Mark it as acked.
|
||||
pkt.off = -1
|
||||
tx.firstPkt++
|
||||
continue
|
||||
}
|
||||
if i == tx.lastPkt {
|
||||
break
|
||||
}
|
||||
i = (i + 1) % len(tx.packets)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (tx *ringTx) unsentRing() internal.Ring {
|
||||
return tx.ring(tx.unsentoff, tx.unsentend)
|
||||
}
|
||||
|
||||
func (tx *ringTx) sentRing() internal.Ring {
|
||||
first := tx.packets[tx.firstPkt]
|
||||
if first.off < 0 {
|
||||
return tx.ring(0, 0)
|
||||
}
|
||||
last := tx.packets[tx.lastPkt]
|
||||
return tx.ring(first.off, last.end)
|
||||
}
|
||||
|
||||
func (tx *ringTx) ring(off, end int) internal.Ring {
|
||||
return internal.Ring{Buf: tx.rawbuf, Off: off, End: end}
|
||||
}
|
||||
@@ -0,0 +1,66 @@
|
||||
/*
|
||||
package ltcp implements TCP control flow.
|
||||
|
||||
# Transmission Control Block
|
||||
|
||||
The Transmission Control Block (TCB) is the core data structure of TCP.
|
||||
It stores core state of the TCP connection such as the send and receive
|
||||
sequence number spaces, the current state of the connection, and the
|
||||
pending control segment flags.
|
||||
|
||||
# Values and Sizes
|
||||
|
||||
All arithmetic dealing with sequence numbers must be performed modulo 2**32
|
||||
which brings with it subtleties to computer modulo arithmetic.
|
||||
*/
|
||||
package tcp
|
||||
|
||||
import "time"
|
||||
|
||||
// Value represents the value of a sequence number.
|
||||
type Value uint32
|
||||
|
||||
// Size represents the size (length) of a sequence number window.
|
||||
type Size uint32
|
||||
|
||||
// LessThan checks if v is before w (modulo 32) i.e., v < w.
|
||||
func LessThan(v, w Value) bool {
|
||||
return int32(v-w) < 0
|
||||
}
|
||||
|
||||
// LessThanEq returns true if v==w or v is before (modulo 32) i.e., v < w.
|
||||
func LessThanEq(v, w Value) bool {
|
||||
return v == w || LessThan(v, w)
|
||||
}
|
||||
|
||||
// InRange checks if v is in the range [a,b) (modulo 32), i.e., a <= v < b.
|
||||
func InRange(v, a, b Value) bool {
|
||||
return v-a < b-a
|
||||
}
|
||||
|
||||
// InWindow checks if v is in the window that starts at 'first' and spans 'size'
|
||||
// sequence numbers (modulo 32).
|
||||
func InWindow(v, first Value, size Size) bool {
|
||||
return InRange(v, first, Add(first, size))
|
||||
}
|
||||
|
||||
// Add calculates the sequence number following the [v, v+s) window.
|
||||
func Add(v Value, s Size) Value {
|
||||
return v + Value(s)
|
||||
}
|
||||
|
||||
// Size calculates the size of the window defined by [v, w).
|
||||
func Sizeof(v, w Value) Size {
|
||||
return Size(w - v)
|
||||
}
|
||||
|
||||
// UpdateForward updates v such that it becomes v + s.
|
||||
func (v *Value) UpdateForward(s Size) {
|
||||
*v += Value(s)
|
||||
}
|
||||
|
||||
// DefaultNewISS returns a new initial send sequence number.
|
||||
// It's implementation is suggested by RFC9293.
|
||||
func DefaultNewISS(t time.Time) Value {
|
||||
return Value(t.UnixMicro() / 4)
|
||||
}
|
||||
Reference in New Issue
Block a user