MCPcopy Create free account
hub / github.com/Hidden-Node/GooseRelayVPN-AndroidClient / drainAll

Method drainAll

internal/exit/exit.go:689–781  ·  view source on GitHub ↗

drainAll returns all currently-buffered TX frames belonging to owner, plus an `urgent` flag signalling that at least one drained session is delivering its first downstream batch (e.g. TLS server hello after SYN). The caller skips the normal coalesce wait when urgent is set so connection setup isn't

(owner [frame.ClientIDLen]byte, byteBudget int)

Source from the content-addressed store, hash-verified

687// first would receive every other client's downstream frames and silently
688// drop them, breaking every TLS stream in flight.
689func (s *Server) drainAll(owner [frame.ClientIDLen]byte, byteBudget int) ([]*frame.Frame, bool) {
690 s.mu.Lock()
691 defer s.mu.Unlock()
692 var out []*frame.Frame
693 var urgent bool
694 if ctrl := s.pendingCtrl[owner]; len(ctrl) > 0 {
695 out = append(out, ctrl...)
696 delete(s.pendingCtrl, owner)
697 urgent = true
698 }
699 if rsts := s.pendingRSTs[owner]; len(rsts) > 0 {
700 out = append(out, rsts...)
701 delete(s.pendingRSTs, owner)
702 urgent = true // RSTs are always urgent — client should know immediately
703 }
704 batchCap := maxDrainFramesPerBatch
705 if len(s.sessions) >= busySessionThreshold {
706 batchCap = maxDrainFramesPerBatchBusy
707 }
708 remaining := batchCap
709 remainingBytes := byteBudget
710
711 // Snapshot and sort active sessions by queue age to ensure fairness.
712 type sessionRef struct {
713 id [frame.SessionIDLen]byte
714 queuedAt time.Time
715 }
716 refs := make([]sessionRef, 0, len(s.txReady))
717 for id := range s.txReady {
718 if sess, ok := s.sessions[id]; ok {
719 if s.sessionOwners[id] != owner {
720 continue
721 }
722 refs = append(refs, sessionRef{id: id, queuedAt: sess.FirstQueuedAt()})
723 } else {
724 delete(s.txReady, id)
725 }
726 }
727 sort.Slice(refs, func(i, j int) bool {
728 return refs[i].queuedAt.Before(refs[j].queuedAt)
729 })
730
731 for _, r := range refs {
732 id := r.id
733 if remaining <= 0 || remainingBytes <= 0 {
734 break
735 }
736 sess, ok := s.sessions[id]
737 if !ok {
738 delete(s.txReady, id)
739 continue
740 }
741 perSessionCap := maxDrainFramesPerSession
742 if remaining < perSessionCap {
743 perSessionCap = remaining
744 }
745 maxPayload := MaxFramePayload
746 if _, isFirst := s.firstReply[id]; isFirst && perSessionCap > 0 {

Calls 3

FirstQueuedAtMethod · 0.80
DrainTxLimitedMethod · 0.80
HasPendingTxMethod · 0.80