diff --git a/client/internal/pqkem/convergence.go b/client/internal/pqkem/convergence.go index d7fc65c99..76f21b8af 100644 --- a/client/internal/pqkem/convergence.go +++ b/client/internal/pqkem/convergence.go @@ -108,7 +108,8 @@ func (m *Manager) processOffer(remoteID RemoteID, o *OfferMsg, via string) ([]by } m.mu.Lock() - if ex := m.exchanges[remoteID]; ex != nil && ex.id == o.ExchangeID { + ex := m.exchanges[remoteID] + if ex != nil && ex.id == o.ExchangeID { state, last := ex.state, ex.lastSent m.mu.Unlock() if state == stateReserved { @@ -117,6 +118,17 @@ func (m *Manager) processOffer(remoteID RemoteID, o *OfferMsg, via string) ([]by m.trace("pqkem: duplicate offer, resending cached answer", "peer", remoteID, "exchange", idHex(o.ExchangeID)) return last, nil } + // A data-path chain offer acknowledges the exchange we are awaiting-ack on, and + // ackConverged above deletes that exchange when the ack matches. So if a chain offer + // (non-zero AckID) did NOT match — our current exchange is still here and is a + // different one — it is a rotation a newer signal re-bootstrap has already superseded. + // Drop it, so a late stale offer can't revert us off the new round. A bootstrap + // (zero AckID) is an authoritative fresh start and always proceeds. + if o.AckID != (ExchangeID{}) && ex != nil { + m.mu.Unlock() + m.trace("pqkem: stale chain offer superseded by a newer round, dropping", "peer", remoteID, "exchange", idHex(o.ExchangeID), "acks", idHex(o.AckID)) + return nil, nil + } // Reserve the slot so a concurrent duplicate offer bails. m.exchanges[remoteID] = &exchangeCtl{id: o.ExchangeID, state: stateReserved, gen: m.nextGenLocked(remoteID)} m.mu.Unlock() @@ -139,7 +151,7 @@ func (m *Manager) processOffer(remoteID RemoteID, o *OfferMsg, via string) ([]by } m.mu.Lock() - ex := m.exchanges[remoteID] + ex = m.exchanges[remoteID] if ex == nil || ex.id != o.ExchangeID { m.mu.Unlock() m.trace("pqkem: exchange superseded during respond, dropping answer", "peer", remoteID, "exchange", idHex(o.ExchangeID)) diff --git a/client/internal/pqkem/manager.go b/client/internal/pqkem/manager.go index 76d07a1a7..8912a27c5 100644 --- a/client/internal/pqkem/manager.go +++ b/client/internal/pqkem/manager.go @@ -446,6 +446,10 @@ func (m *Manager) OnDataPathRekeyed(remoteID RemoteID, sinceActivity time.Durati // Hold the lock across the awaitingRekey check and the install so two rekey clocks // can't each start a chained exchange for the same peer. offer, err := m.startExchangeLocked(remoteID, false, ex.id) + var chainID ExchangeID + if nc := m.exchanges[remoteID]; nc != nil { + chainID = nc.id + } m.mu.Unlock() m.trace("pqkem: data-path rekey signal", "peer", remoteID, "chaining", true) @@ -453,6 +457,18 @@ func (m *Manager) OnDataPathRekeyed(remoteID RemoteID, sinceActivity time.Durati m.logger.Error("pqkem: chain offer failed to start", "peer", remoteID, "err", err) return } + // A signal re-bootstrap can supersede this chain exchange in the gap between building + // the offer and sending it. Don't put the stale offer on the wire: the responder would + // otherwise reserve and commit an abandoned exchange, splitting the keys. A newer + // signal round always wins. (A supersede landing after this check still leaks one + // harmless offer; the responder's answer for it is rejected as superseded.) + m.mu.Lock() + stillCurrent := m.exchanges[remoteID] != nil && m.exchanges[remoteID].id == chainID + m.mu.Unlock() + if !stillCurrent { + m.trace("pqkem: chain offer superseded before send, dropping", "peer", remoteID, "exchange", idHex(chainID)) + return + } if err := m.pushDataPath(remoteID, offer); err != nil { m.logger.Warn("pqkem: send chain offer failed", "peer", remoteID, "err", err) return