:p
atchew
Login
Under memory pressure, a pruning of the MPTCP-level OoO queue might be required as last resort, to avoid too long recoveries, or even stalls. Geliang and Gang managed to reproduce this behaviour, and Paolo improved the situation thanks to the following patches: - Patches 1-2: improve the MPTCP-level retransmission schema to make recoveries from memory pressure/after MPTCP-level drop significantly faster. - Patches 3-4: make the admission check way stricter for incoming packets exceeding the memory limits, with some exceptions for fallback sockets. - Patch 5: implement OoO queue pruning for MPTCP. Signed-off-by: Matthieu Baerts (NGI0) <matttbe@kernel.org> --- Changes in v2: - Address comments from Clashiko. - Patch 3: typo, comment, bump TCP MIB counter. - Patch 5: uniform MIB counter name. - Rebased. - Link to v1: https://patch.msgid.link/20260724-net-next-mptcp-oooq-pruning-v1-0-5dd4dec63a54@kernel.org --- Paolo Abeni (5): mptcp: move the retrans loop to a separate helper mptcp: let the retrans scheduler do its job mptcp: explicitly drop over memory limits mptcp: enforce hard limit on backlog flushing mptcp: implemented OoO queue pruning net/mptcp/mib.c | 3 + net/mptcp/mib.h | 3 + net/mptcp/options.c | 32 ++++++- net/mptcp/protocol.c | 251 +++++++++++++++++++++++++++++++++++++-------------- 4 files changed, 219 insertions(+), 70 deletions(-) --- base-commit: 2fbade66245059c78daeaccfce13ecf499fffb51 change-id: 20260724-net-next-mptcp-oooq-pruning-48566d10dbd0 Best regards, -- Matthieu Baerts (NGI0) <matttbe@kernel.org>
From: Paolo Abeni <pabeni@redhat.com> This is a cleanup in order to make the next patch simpler. No functional change intended. Tested-by: Gang Yan <yangang@kylinos.cn> Tested-by: Geliang Tang <geliang@kernel.org> Acked-by: Geliang Tang <geliang@kernel.org> Signed-off-by: Paolo Abeni <pabeni@redhat.com> Signed-off-by: Matthieu Baerts (NGI0) <matttbe@kernel.org> --- net/mptcp/protocol.c | 74 ++++++++++++++++++++++++++++++---------------------- 1 file changed, 43 insertions(+), 31 deletions(-) diff --git a/net/mptcp/protocol.c b/net/mptcp/protocol.c index XXXXXXX..XXXXXXX 100644 --- a/net/mptcp/protocol.c +++ b/net/mptcp/protocol.c @@ -XXX,XX +XXX,XX @@ static void mptcp_check_fastclose(struct mptcp_sock *msk) sk_error_report(sk); } -static void __mptcp_retrans(struct sock *sk) +/* Retransmit the specified data fragment on all the selected subflows. */ +static int __mptcp_push_retrans(struct sock *sk, struct mptcp_data_frag *dfrag) { struct mptcp_sendmsg_info info = { .data_lock_held = true, }; struct mptcp_sock *msk = mptcp_sk(sk); struct mptcp_subflow_context *subflow; - struct mptcp_data_frag *dfrag; struct sock *ssk; - int ret, err; - u16 len = 0; - - mptcp_clean_una_wakeup(sk); - - /* first check ssk: need to kick "stale" logic */ - err = mptcp_sched_get_retrans(msk); - dfrag = mptcp_rtx_head(sk); - if (!dfrag) { - if (mptcp_data_fin_enabled(msk)) { - struct inet_connection_sock *icsk = inet_csk(sk); - - WRITE_ONCE(icsk->icsk_retransmits, - icsk->icsk_retransmits + 1); - mptcp_set_datafin_timeout(sk); - mptcp_send_ack(msk); - - goto reset_timer; - } - - if (!mptcp_send_head(sk)) - goto clear_scheduled; - - goto reset_timer; - } - - if (err) - goto reset_timer; + int ret, len = 0; mptcp_for_each_subflow(msk, subflow) { if (READ_ONCE(subflow->scheduled)) { @@ -XXX,XX +XXX,XX @@ static void __mptcp_retrans(struct sock *sk) !msk->allow_subflows) { spin_unlock_bh(&msk->fallback_lock); release_sock(ssk); - goto clear_scheduled; + return -1; } while (info.sent < info.limit) { @@ -XXX,XX +XXX,XX @@ static void __mptcp_retrans(struct sock *sk) release_sock(ssk); } } + return len; +} + +static void __mptcp_retrans(struct sock *sk) +{ + struct mptcp_sock *msk = mptcp_sk(sk); + struct mptcp_subflow_context *subflow; + struct mptcp_data_frag *dfrag; + int err, len; + + mptcp_clean_una_wakeup(sk); + + /* first check ssk: need to kick "stale" logic */ + err = mptcp_sched_get_retrans(msk); + dfrag = mptcp_rtx_head(sk); + if (!dfrag) { + if (mptcp_data_fin_enabled(msk)) { + struct inet_connection_sock *icsk = inet_csk(sk); + + WRITE_ONCE(icsk->icsk_retransmits, + icsk->icsk_retransmits + 1); + mptcp_set_datafin_timeout(sk); + mptcp_send_ack(msk); + + goto reset_timer; + } + + if (!mptcp_send_head(sk)) + goto clear_scheduled; + + goto reset_timer; + } + + if (err) + goto reset_timer; + + len = __mptcp_push_retrans(sk, dfrag); + if (len < 0) + goto clear_scheduled; msk->bytes_retrans += len; dfrag->already_sent = max(dfrag->already_sent, len); -- 2.53.0
From: Paolo Abeni <pabeni@redhat.com> Currently the MPTCP core enforces that when MPTCP-level retrans timer fires, at most a single dfrag is retransmitted. In some corner-cases, it may be necessary to retransmit multiple dfrags, and the MPTCP socket will need to wait multiple retrans timeout to accomplish that. Remove the mentioned constraint, allowing to transmit multiple dfrags per retrans period, as long as the scheduler keeps selecting subflows for retransmissions and pending data is available in the rtx queue. The default scheduler will transmit a dfrag per available subflow. Tested-by: Gang Yan <yangang@kylinos.cn> Tested-by: Geliang Tang <geliang@kernel.org> Acked-by: Geliang Tang <geliang@kernel.org> Signed-off-by: Paolo Abeni <pabeni@redhat.com> Signed-off-by: Matthieu Baerts (NGI0) <matttbe@kernel.org> --- net/mptcp/protocol.c | 119 +++++++++++++++++++++++++++++++++++++-------------- 1 file changed, 87 insertions(+), 32 deletions(-) diff --git a/net/mptcp/protocol.c b/net/mptcp/protocol.c index XXXXXXX..XXXXXXX 100644 --- a/net/mptcp/protocol.c +++ b/net/mptcp/protocol.c @@ -XXX,XX +XXX,XX @@ static void __mptcp_clean_una_wakeup(struct sock *sk) mptcp_write_space(sk); } -static void mptcp_clean_una_wakeup(struct sock *sk) -{ - mptcp_data_lock(sk); - __mptcp_clean_una_wakeup(sk); - mptcp_data_unlock(sk); -} - static void mptcp_enter_memory_pressure(struct sock *sk) { struct mptcp_subflow_context *subflow; @@ -XXX,XX +XXX,XX @@ static void mptcp_check_fastclose(struct mptcp_sock *msk) sk_error_report(sk); } -/* Retransmit the specified data fragment on all the selected subflows. */ -static int __mptcp_push_retrans(struct sock *sk, struct mptcp_data_frag *dfrag) +/* + * Retransmit the specified data fragment on all the selected subflows, + * starting from the specified sequence + */ +static int __mptcp_push_retrans(struct sock *sk, struct mptcp_data_frag *dfrag, + u64 sent_seq) { struct mptcp_sendmsg_info info = { .data_lock_held = true, }; struct mptcp_sock *msk = mptcp_sk(sk); @@ -XXX,XX +XXX,XX @@ static int __mptcp_push_retrans(struct sock *sk, struct mptcp_data_frag *dfrag) mptcp_for_each_subflow(msk, subflow) { if (READ_ONCE(subflow->scheduled)) { + u16 offset = sent_seq - dfrag->data_seq; u16 copied = 0; mptcp_subflow_set_scheduled(subflow, false); @@ -XXX,XX +XXX,XX @@ static int __mptcp_push_retrans(struct sock *sk, struct mptcp_data_frag *dfrag) lock_sock(ssk); /* limit retransmission to the bytes already sent on some subflows */ - info.sent = 0; + info.sent = offset; info.limit = READ_ONCE(msk->csum_enabled) ? dfrag->data_len : dfrag->already_sent; @@ -XXX,XX +XXX,XX @@ static void __mptcp_retrans(struct sock *sk) struct mptcp_sock *msk = mptcp_sk(sk); struct mptcp_subflow_context *subflow; struct mptcp_data_frag *dfrag; + bool need_retrans; + u64 retrans_seq; int err, len; - mptcp_clean_una_wakeup(sk); - - /* first check ssk: need to kick "stale" logic */ - err = mptcp_sched_get_retrans(msk); + mptcp_data_lock(sk); + __mptcp_clean_una_wakeup(sk); + retrans_seq = msk->snd_una; dfrag = mptcp_rtx_head(sk); - if (!dfrag) { + need_retrans = !!dfrag; + mptcp_data_unlock(sk); + if (!dfrag) + goto check_data_fin; + + for (;;) { + bool already_retrans; + u64 sent_seq; + + /* The default scheduler will kick "stale" logic, that in + * turn can process incoming acks and clean the RTX queue; + * ensure that the current dfrag will still be around + * afterwards. + */ + get_page(dfrag->page); + err = mptcp_sched_get_retrans(msk); + if (err) { + put_page(dfrag->page); + break; + } + + /* Incoming acks can have moved retrans sequence after + * the current dfrag, if so try to start again from RTX head. + */ + mptcp_data_lock(sk); + already_retrans = !before64(msk->snd_una, dfrag->data_seq + + dfrag->already_sent); + put_page(dfrag->page); + if (already_retrans) { + __mptcp_clean_una_wakeup(sk); + retrans_seq = msk->snd_una; + dfrag = mptcp_rtx_head(sk); + need_retrans = !!dfrag; + } else if (after64(msk->snd_una, retrans_seq)) { + retrans_seq = msk->snd_una; + } + mptcp_data_unlock(sk); + + /* `already_sent` can be 0 for `dfrag` belonging to the RTX + * queue due to __mptcp_retransmit_pending_data(). + */ + if (!dfrag || !dfrag->already_sent) + break; + + /* Can fail only in case of fallback. */ + len = __mptcp_push_retrans(sk, dfrag, retrans_seq); + if (len < 0) + goto clear_scheduled; + + retrans_seq += len; + msk->bytes_retrans += len; + dfrag->already_sent = max_t(u16, dfrag->already_sent, + retrans_seq - dfrag->data_seq); + + /* With csum enabled retransmission can send new data. */ + sent_seq = dfrag->already_sent + dfrag->data_seq; + if (after64(sent_seq, msk->snd_nxt)) + WRITE_ONCE(msk->snd_nxt, sent_seq); + + /* Attempt the next fragment only if the current one is + * completely retransmitted. + */ + if (before64(retrans_seq, dfrag->data_seq + dfrag->data_len)) + break; + + dfrag = list_is_last(&dfrag->list, &msk->rtx_queue) ? + NULL : list_next_entry(dfrag, list); + if (!dfrag) + break; + } + + /* Attempt data-fin retransmission only when the RTX queue is empty. */ + if (!need_retrans) { +check_data_fin: if (mptcp_data_fin_enabled(msk)) { struct inet_connection_sock *icsk = inet_csk(sk); @@ -XXX,XX +XXX,XX @@ static void __mptcp_retrans(struct sock *sk) icsk->icsk_retransmits + 1); mptcp_set_datafin_timeout(sk); mptcp_send_ack(msk); - goto reset_timer; } if (!mptcp_send_head(sk)) goto clear_scheduled; - - goto reset_timer; } - if (err) - goto reset_timer; - - len = __mptcp_push_retrans(sk, dfrag); - if (len < 0) - goto clear_scheduled; - - msk->bytes_retrans += len; - dfrag->already_sent = max(dfrag->already_sent, len); - - /* With csum enabled retransmission can send new data. */ - if (after64(dfrag->already_sent + dfrag->data_seq, msk->snd_nxt)) - WRITE_ONCE(msk->snd_nxt, dfrag->already_sent + dfrag->data_seq); - reset_timer: mptcp_check_and_set_pending(sk); -- 2.53.0
From: Paolo Abeni <pabeni@redhat.com> Currently the enforcement of the rcvbuf constraint is implemented when moving the skbs into the msk receive or OoO queue, keeping the incoming skbs in the subflow queue when over limits. Under significant memory pressure the above can cause permanent data transfer stalls, as the skb needed to make forward progress can be stuck in a subflow queue. Over memory limits, drop the incoming skb, relying on MPTCP-level retransmissions. Note that fallback socket must perform the limit before the skb reaches the subflow-level queue, as dropping an in-sequence already acked skb would break the stream. This is not a complete fix for the stall issue, as the drop strategy needs refinements that will come in the next patches. Signed-off-by: Paolo Abeni <pabeni@redhat.com> Reviewed-by: Matthieu Baerts (NGI0) <matttbe@kernel.org> [ Fix typo, comment, and bump LINUX_MIB_TCPRCVQDROP ] Signed-off-by: Matthieu Baerts (NGI0) <matttbe@kernel.org> --- v2: - mib: typo: "constrains" -> "constraints". - mptcp_over_limit: more than 0-win: retrans, dup or old acks. - mptcp_over_limit: bump LINUX_MIB_TCPRCVQDROP. - Note: Sashiko might point to a possible forward-allocated memory leak: this is a temp leak, and releasing additionally allocated fwd memory in the error path will be fix in a patch for -net. --- net/mptcp/mib.c | 2 ++ net/mptcp/mib.h | 2 ++ net/mptcp/options.c | 32 +++++++++++++++++++++++++++++--- net/mptcp/protocol.c | 31 +++++++++++++++++++++++-------- 4 files changed, 56 insertions(+), 11 deletions(-) diff --git a/net/mptcp/mib.c b/net/mptcp/mib.c index XXXXXXX..XXXXXXX 100644 --- a/net/mptcp/mib.c +++ b/net/mptcp/mib.c @@ -XXX,XX +XXX,XX @@ static const struct snmp_mib mptcp_snmp_list[] = { SNMP_MIB_ITEM("SimultConnectFallback", MPTCP_MIB_SIMULTCONNFALLBACK), SNMP_MIB_ITEM("FallbackFailed", MPTCP_MIB_FALLBACKFAILED), SNMP_MIB_ITEM("WinProbe", MPTCP_MIB_WINPROBE), + SNMP_MIB_ITEM("BacklogDrop", MPTCP_MIB_BACKLOGDROP), + SNMP_MIB_ITEM("RcvPruned", MPTCP_MIB_RCVPRUNED), }; /* mptcp_mib_alloc - allocate percpu mib counters diff --git a/net/mptcp/mib.h b/net/mptcp/mib.h index XXXXXXX..XXXXXXX 100644 --- a/net/mptcp/mib.h +++ b/net/mptcp/mib.h @@ -XXX,XX +XXX,XX @@ enum linux_mptcp_mib_field { MPTCP_MIB_SIMULTCONNFALLBACK, /* Simultaneous connect */ MPTCP_MIB_FALLBACKFAILED, /* Can't fallback due to msk status */ MPTCP_MIB_WINPROBE, /* MPTCP-level zero window probe */ + MPTCP_MIB_BACKLOGDROP, /* Backlog over memory limit */ + MPTCP_MIB_RCVPRUNED, /* Dropped due to memory constraints */ __MPTCP_MIB_MAX }; diff --git a/net/mptcp/options.c b/net/mptcp/options.c index XXXXXXX..XXXXXXX 100644 --- a/net/mptcp/options.c +++ b/net/mptcp/options.c @@ -XXX,XX +XXX,XX @@ static bool add_addr_hmac_valid(struct mptcp_sock *msk, return hmac == mp_opt->ahmac; } -/* Return false in case of error (or subflow has been reset), - * else return true. +static bool mptcp_over_limit(struct sock *sk, struct sock *ssk, + const struct sk_buff *skb) +{ + struct mptcp_sock *msk = mptcp_sk(sk); + u64 mem = sk_rmem_alloc_get(sk); + + mem += READ_ONCE(msk->backlog_len); + if (likely(mem <= READ_ONCE(sk->sk_rcvbuf))) + return false; + + /* Avoid silently dropping pure acks, fin or already-acked segments. */ + if (TCP_SKB_CB(skb)->seq == TCP_SKB_CB(skb)->end_seq || + TCP_SKB_CB(skb)->tcp_flags & TCPHDR_FIN || + !after(TCP_SKB_CB(skb)->end_seq, tcp_sk(ssk)->rcv_nxt)) + return false; + + /* Dropped due to memory constraints, schedule an ack. */ + inet_csk(ssk)->icsk_ack.pending |= ICSK_ACK_NOMEM | ICSK_ACK_NOW; + inet_csk_schedule_ack(ssk); + + /* Plain TCP (fallback) and skb is dropped before the TCP recv queue. */ + NET_INC_STATS(sock_net(sk), LINUX_MIB_TCPRCVQDROP); + + return true; +} + +/* Return false when the caller must drop the packet, i.e. in case of error, + * subflow has been reset, or over memory limits. */ bool mptcp_incoming_options(struct sock *sk, struct sk_buff *skb) { @@ -XXX,XX +XXX,XX @@ bool mptcp_incoming_options(struct sock *sk, struct sk_buff *skb) __mptcp_data_acked(subflow->conn); mptcp_data_unlock(subflow->conn); - return true; + return !mptcp_over_limit(subflow->conn, sk, skb); } mptcp_get_options(skb, &mp_opt); diff --git a/net/mptcp/protocol.c b/net/mptcp/protocol.c index XXXXXXX..XXXXXXX 100644 --- a/net/mptcp/protocol.c +++ b/net/mptcp/protocol.c @@ -XXX,XX +XXX,XX @@ static bool __mptcp_move_skb(struct sock *sk, struct sk_buff *skb) mptcp_borrow_fwdmem(sk, skb); + /* Can't drop packets for fallback socket this late, or the stream + * will break. + */ + if (unlikely(sk_rmem_alloc_get(sk) > READ_ONCE(sk->sk_rcvbuf)) && + !__mptcp_check_fallback(msk)) { + MPTCP_INC_STATS(sock_net(sk), MPTCP_MIB_RCVPRUNED); + mptcp_drop(sk, skb); + return false; + } + if (MPTCP_SKB_CB(skb)->map_seq == msk->ack_seq) { /* in sequence */ msk->bytes_received += copy_len; @@ -XXX,XX +XXX,XX @@ static void __mptcp_add_backlog(struct sock *sk, struct sk_buff *tail = NULL; struct sock *ssk = skb->sk; bool fragstolen; + u64 limit; int delta; if (unlikely(sk->sk_state == TCP_CLOSE)) { @@ -XXX,XX +XXX,XX @@ static void __mptcp_add_backlog(struct sock *sk, return; } + /* Similar additional allowance as plain TCP. */ + limit = READ_ONCE(sk->sk_rcvbuf); + limit += (limit >> 1) + 64 * 1024; + limit = min_t(u64, limit, UINT_MAX); + if (msk->backlog_len > limit && !__mptcp_check_fallback(msk)) { + __MPTCP_INC_STATS(sock_net(sk), MPTCP_MIB_BACKLOGDROP); + kfree_skb_reason(skb, SKB_DROP_REASON_SOCKET_BACKLOG); + return; + } + /* Try to coalesce with the last skb in our backlog */ if (!list_empty(&msk->backlog_list)) tail = list_last_entry(&msk->backlog_list, struct sk_buff, list); @@ -XXX,XX +XXX,XX @@ static bool __mptcp_move_skbs_from_subflow(struct mptcp_sock *msk, mptcp_init_skb(ssk, skb, offset, len); - if (own_msk && sk_rmem_alloc_get(sk) < sk->sk_rcvbuf) { + if (own_msk) { mptcp_subflow_lend_fwdmem(subflow, skb); ret |= __mptcp_move_skb(sk, skb); } else { @@ -XXX,XX +XXX,XX @@ static bool __mptcp_move_skbs(struct sock *sk, struct list_head *skbs, u32 *delt *delta = 0; while (1) { - /* If the msk recvbuf is full stop, don't drop */ - if (sk_rmem_alloc_get(sk) > sk->sk_rcvbuf) - break; - prefetch(skb->next); list_del(&skb->list); *delta += skb->truesize; @@ -XXX,XX +XXX,XX @@ static bool mptcp_can_spool_backlog(struct sock *sk, struct list_head *skbs) DEBUG_NET_WARN_ON_ONCE(msk->backlog_unaccounted && sk->sk_socket && mem_cgroup_from_sk(sk)); - /* Don't spool the backlog if the rcvbuf is full. */ - if (list_empty(&msk->backlog_list) || - sk_rmem_alloc_get(sk) > sk->sk_rcvbuf) + if (list_empty(&msk->backlog_list)) return false; INIT_LIST_HEAD(skbs); -- 2.53.0
From: Paolo Abeni <pabeni@redhat.com> Currently a wild producer could keep the backlog flushing operation spinning for an unbound time. Since the previous patch, the amount of data present in the backlog is hard-limited. Move the backlog len update at the end of the flush loop to prevent it spinning forever. Also, no need to splice back the remaining skbs list into the backlog, as such list is always empty after each backlog processing loop. Signed-off-by: Paolo Abeni <pabeni@redhat.com> Reviewed-by: Matthieu Baerts (NGI0) <matttbe@kernel.org> Signed-off-by: Matthieu Baerts (NGI0) <matttbe@kernel.org> --- net/mptcp/protocol.c | 21 ++++++--------------- 1 file changed, 6 insertions(+), 15 deletions(-) diff --git a/net/mptcp/protocol.c b/net/mptcp/protocol.c index XXXXXXX..XXXXXXX 100644 --- a/net/mptcp/protocol.c +++ b/net/mptcp/protocol.c @@ -XXX,XX +XXX,XX @@ static bool __mptcp_move_skbs(struct sock *sk, struct list_head *skbs, u32 *delt struct mptcp_sock *msk = mptcp_sk(sk); bool moved = false; - *delta = 0; while (1) { prefetch(skb->next); list_del(&skb->list); @@ -XXX,XX +XXX,XX @@ static bool mptcp_can_spool_backlog(struct sock *sk, struct list_head *skbs) return true; } -static void mptcp_backlog_spooled(struct sock *sk, u32 moved, - struct list_head *skbs) -{ - struct mptcp_sock *msk = mptcp_sk(sk); - - WRITE_ONCE(msk->backlog_len, msk->backlog_len - moved); - list_splice(skbs, &msk->backlog_list); -} - static bool mptcp_move_skbs(struct sock *sk) { + struct mptcp_sock *msk = mptcp_sk(sk); struct list_head skbs; bool enqueued = false; - u32 moved; + u32 moved = 0; mptcp_data_lock(sk); while (mptcp_can_spool_backlog(sk, &skbs)) { @@ -XXX,XX +XXX,XX @@ static bool mptcp_move_skbs(struct sock *sk) enqueued |= __mptcp_move_skbs(sk, &skbs, &moved); mptcp_data_lock(sk); - mptcp_backlog_spooled(sk, moved, &skbs); } + WRITE_ONCE(msk->backlog_len, msk->backlog_len - moved); mptcp_data_unlock(sk); if (enqueued && mptcp_epollin_ready(sk)) @@ -XXX,XX +XXX,XX @@ static void mptcp_release_cb(struct sock *sk) __must_hold(&sk->sk_lock.slock) { struct mptcp_sock *msk = mptcp_sk(sk); + u32 moved = 0; for (;;) { unsigned long flags = (msk->cb_flags & MPTCP_FLAGS_PROCESS_CTX_NEED); struct list_head join_list, skbs; bool spool_bl; - u32 moved; spool_bl = mptcp_can_spool_backlog(sk, &skbs); if (!flags && !spool_bl) @@ -XXX,XX +XXX,XX @@ static void mptcp_release_cb(struct sock *sk) cond_resched(); spin_lock_bh(&sk->sk_lock.slock); - if (spool_bl) - mptcp_backlog_spooled(sk, moved, &skbs); } + if (moved) + WRITE_ONCE(msk->backlog_len, msk->backlog_len - moved); if (__test_and_clear_bit(MPTCP_CLEAN_UNA, &msk->cb_flags)) __mptcp_clean_una_wakeup(sk); -- 2.53.0
From: Paolo Abeni <pabeni@redhat.com> When moving incoming skbs in the msk receive queue and the latter is above limits, prune it as needed quite alike what TCP is doing at the subflow level. The main difference relies in the stop condition: since MPTCP does not perform collapsing, it's better off dropping the bare minimum to fit the (newer) incoming packet. Signed-off-by: Paolo Abeni <pabeni@redhat.com> Tested-by: Gang Yan <yangang@kylinos.cn> Reviewed-by: Matthieu Baerts (NGI0) <matttbe@kernel.org> [ Uniform OFO MIB counters ] Signed-off-by: Matthieu Baerts (NGI0) <matttbe@kernel.org> --- v2: - Uniform the new counter with the other OFO ones. --- net/mptcp/mib.c | 1 + net/mptcp/mib.h | 1 + net/mptcp/protocol.c | 46 +++++++++++++++++++++++++++++++++++++++++++++- 3 files changed, 47 insertions(+), 1 deletion(-) diff --git a/net/mptcp/mib.c b/net/mptcp/mib.c index XXXXXXX..XXXXXXX 100644 --- a/net/mptcp/mib.c +++ b/net/mptcp/mib.c @@ -XXX,XX +XXX,XX @@ static const struct snmp_mib mptcp_snmp_list[] = { SNMP_MIB_ITEM("WinProbe", MPTCP_MIB_WINPROBE), SNMP_MIB_ITEM("BacklogDrop", MPTCP_MIB_BACKLOGDROP), SNMP_MIB_ITEM("RcvPruned", MPTCP_MIB_RCVPRUNED), + SNMP_MIB_ITEM("OFOPruned", MPTCP_MIB_OFOPRUNED), }; /* mptcp_mib_alloc - allocate percpu mib counters diff --git a/net/mptcp/mib.h b/net/mptcp/mib.h index XXXXXXX..XXXXXXX 100644 --- a/net/mptcp/mib.h +++ b/net/mptcp/mib.h @@ -XXX,XX +XXX,XX @@ enum linux_mptcp_mib_field { MPTCP_MIB_WINPROBE, /* MPTCP-level zero window probe */ MPTCP_MIB_BACKLOGDROP, /* Backlog over memory limit */ MPTCP_MIB_RCVPRUNED, /* Dropped due to memory constraints */ + MPTCP_MIB_OFOPRUNED, /* MPTCP-level OoO queue pruned */ __MPTCP_MIB_MAX }; diff --git a/net/mptcp/protocol.c b/net/mptcp/protocol.c index XXXXXXX..XXXXXXX 100644 --- a/net/mptcp/protocol.c +++ b/net/mptcp/protocol.c @@ -XXX,XX +XXX,XX @@ static void mptcp_init_skb(struct sock *ssk, struct sk_buff *skb, int offset, skb_dst_drop(skb); } +/* "Inspired" from the TCP version; main difference: stop as soon as the MPTCP + * socket is under memory limit. + */ +static bool mptcp_prune_ofo_queue(struct sock *sk, u64 seq) +{ + struct mptcp_sock *msk = mptcp_sk(sk); + struct rb_node *node, *prev; + bool pruned = false; + u64 mem; + + if (RB_EMPTY_ROOT(&msk->out_of_order_queue)) + goto out; + + node = &msk->ooo_last_skb->rbnode; + + do { + struct sk_buff *skb = rb_to_skb(node); + + /* Stop pruning if the incoming skb would land in OoO tail. */ + if (after64(seq, MPTCP_SKB_CB(skb)->map_seq)) + break; + + pruned = true; + prev = rb_prev(node); + rb_erase(node, &msk->out_of_order_queue); + mptcp_drop(sk, skb); + msk->ooo_last_skb = rb_to_skb(prev); + + mem = (unsigned int)sk_rmem_alloc_get(sk); + if (mem <= sk->sk_rcvbuf) + break; + + node = prev; + } while (node); + + if (pruned) + MPTCP_INC_STATS(sock_net(sk), MPTCP_MIB_OFOPRUNED); + +out: + mem = (unsigned int)sk_rmem_alloc_get(sk); + return mem <= sk->sk_rcvbuf; +} + static bool __mptcp_move_skb(struct sock *sk, struct sk_buff *skb) { u64 copy_len = MPTCP_SKB_CB(skb)->end_seq - MPTCP_SKB_CB(skb)->map_seq; @@ -XXX,XX +XXX,XX @@ static bool __mptcp_move_skb(struct sock *sk, struct sk_buff *skb) * will break. */ if (unlikely(sk_rmem_alloc_get(sk) > READ_ONCE(sk->sk_rcvbuf)) && - !__mptcp_check_fallback(msk)) { + !__mptcp_check_fallback(msk) && + !mptcp_prune_ofo_queue(sk, MPTCP_SKB_CB(skb)->map_seq)) { MPTCP_INC_STATS(sock_net(sk), MPTCP_MIB_RCVPRUNED); mptcp_drop(sk, skb); return false; -- 2.53.0
Under memory pressure, a pruning of the MPTCP-level OoO queue might be required as last resort, to avoid too long recoveries, or even stalls. Geliang and Gang managed to reproduce this behaviour, and Paolo improved the situation thanks to the following patches: - Patches 1-3: improve the MPTCP-level retransmission schema to make recoveries from memory pressure/after MPTCP-level drop significantly faster. - Patches 4-5: make the admission check way stricter for incoming packets exceeding the memory limits, with some exceptions for fallback sockets. - Patches 6-7: implement OoO queue pruning for MPTCP. Signed-off-by: Matthieu Baerts (NGI0) <matttbe@kernel.org> --- Changes in v3: - Address comments from Clashiko. - Patches 2, 6: new - Patch 3: many cleanups, needed after patch 2. - Patch 4: check backlog and rcvbuf separately + do not drop rst - Patch 7: prune only for new data + reorganize code to follow TCP - Link to v2: https://patch.msgid.link/20260731-net-next-mptcp-oooq-pruning-v2-0-24838164fa21@kernel.org Changes in v2: - Address comments from Clashiko. - Patch 3: typo, comment, bump TCP MIB counter. - Patch 5: uniform MIB counter name. - Rebased. - Link to v1: https://patch.msgid.link/20260724-net-next-mptcp-oooq-pruning-v1-0-5dd4dec63a54@kernel.org --- Paolo Abeni (7): mptcp: move the retrans loop to a separate helper mptcp: move the stale logic out of retrans scheduler mptcp: let the retrans scheduler do its job mptcp: explicitly drop over memory limits mptcp: enforce hard limit on backlog flushing mptcp: avoid code duplication in __mptcp_move_skb() mptcp: implemented OoO queue pruning net/mptcp/mib.c | 3 + net/mptcp/mib.h | 3 + net/mptcp/options.c | 32 +++++- net/mptcp/pm.c | 41 +++++--- net/mptcp/protocol.c | 290 ++++++++++++++++++++++++++++++++++++--------------- net/mptcp/protocol.h | 11 +- 6 files changed, 273 insertions(+), 107 deletions(-) --- base-commit: 4fa4977a0d900f936bcae5cd2c510be5554e8dd6 change-id: 20260724-net-next-mptcp-oooq-pruning-48566d10dbd0 Best regards, -- Matthieu Baerts (NGI0) <matttbe@kernel.org>
From: Paolo Abeni <pabeni@redhat.com> This is a cleanup in order to make the next patch simpler. No functional change intended. Tested-by: Gang Yan <yangang@kylinos.cn> Tested-by: Geliang Tang <geliang@kernel.org> Acked-by: Geliang Tang <geliang@kernel.org> Signed-off-by: Paolo Abeni <pabeni@redhat.com> Signed-off-by: Matthieu Baerts (NGI0) <matttbe@kernel.org> --- net/mptcp/protocol.c | 74 ++++++++++++++++++++++++++++++---------------------- 1 file changed, 43 insertions(+), 31 deletions(-) diff --git a/net/mptcp/protocol.c b/net/mptcp/protocol.c index XXXXXXX..XXXXXXX 100644 --- a/net/mptcp/protocol.c +++ b/net/mptcp/protocol.c @@ -XXX,XX +XXX,XX @@ static void mptcp_check_fastclose(struct mptcp_sock *msk) sk_error_report(sk); } -static void __mptcp_retrans(struct sock *sk) +/* Retransmit the specified data fragment on all the selected subflows. */ +static int __mptcp_push_retrans(struct sock *sk, struct mptcp_data_frag *dfrag) { struct mptcp_sendmsg_info info = { .data_lock_held = true, }; struct mptcp_sock *msk = mptcp_sk(sk); struct mptcp_subflow_context *subflow; - struct mptcp_data_frag *dfrag; struct sock *ssk; - int ret, err; - u16 len = 0; - - mptcp_clean_una_wakeup(sk); - - /* first check ssk: need to kick "stale" logic */ - err = mptcp_sched_get_retrans(msk); - dfrag = mptcp_rtx_head(sk); - if (!dfrag) { - if (mptcp_data_fin_enabled(msk)) { - struct inet_connection_sock *icsk = inet_csk(sk); - - WRITE_ONCE(icsk->icsk_retransmits, - icsk->icsk_retransmits + 1); - mptcp_set_datafin_timeout(sk); - mptcp_send_ack(msk); - - goto reset_timer; - } - - if (!mptcp_send_head(sk)) - goto clear_scheduled; - - goto reset_timer; - } - - if (err) - goto reset_timer; + int ret, len = 0; mptcp_for_each_subflow(msk, subflow) { if (READ_ONCE(subflow->scheduled)) { @@ -XXX,XX +XXX,XX @@ static void __mptcp_retrans(struct sock *sk) !msk->allow_subflows) { spin_unlock_bh(&msk->fallback_lock); release_sock(ssk); - goto clear_scheduled; + return -1; } while (info.sent < info.limit) { @@ -XXX,XX +XXX,XX @@ static void __mptcp_retrans(struct sock *sk) release_sock(ssk); } } + return len; +} + +static void __mptcp_retrans(struct sock *sk) +{ + struct mptcp_sock *msk = mptcp_sk(sk); + struct mptcp_subflow_context *subflow; + struct mptcp_data_frag *dfrag; + int err, len; + + mptcp_clean_una_wakeup(sk); + + /* first check ssk: need to kick "stale" logic */ + err = mptcp_sched_get_retrans(msk); + dfrag = mptcp_rtx_head(sk); + if (!dfrag) { + if (mptcp_data_fin_enabled(msk)) { + struct inet_connection_sock *icsk = inet_csk(sk); + + WRITE_ONCE(icsk->icsk_retransmits, + icsk->icsk_retransmits + 1); + mptcp_set_datafin_timeout(sk); + mptcp_send_ack(msk); + + goto reset_timer; + } + + if (!mptcp_send_head(sk)) + goto clear_scheduled; + + goto reset_timer; + } + + if (err) + goto reset_timer; + + len = __mptcp_push_retrans(sk, dfrag); + if (len < 0) + goto clear_scheduled; msk->bytes_retrans += len; dfrag->already_sent = max(dfrag->already_sent, len); -- 2.53.0
From: Paolo Abeni <pabeni@redhat.com> This allow separating the stale logic invocation and the retrans scheduler, and will simplify the next patch. It's also a cleaner design as the retrans scheduler has currently too many side effects. As a possible downside, the retrans work will now traverse the subflows list additional times; that does not matter much, as this is slowpath. While at it, pick more accurate names for the involved helpers and explicitly note that the per subflow stale data is under msk socket lock protection. The scheduler and the stale logic may observe different subflow statues, as no subflow lock is acquired. This is intentional and not harmful, worst case leading to slower retransmissions. Signed-off-by: Paolo Abeni <pabeni@redhat.com> Reviewed-by: Matthieu Baerts (NGI0) <matttbe@kernel.org> Signed-off-by: Matthieu Baerts (NGI0) <matttbe@kernel.org> --- v3: new --- net/mptcp/pm.c | 41 +++++++++++++++++++++++++++-------------- net/mptcp/protocol.c | 4 ++-- net/mptcp/protocol.h | 11 +++++++---- 3 files changed, 36 insertions(+), 20 deletions(-) diff --git a/net/mptcp/pm.c b/net/mptcp/pm.c index XXXXXXX..XXXXXXX 100644 --- a/net/mptcp/pm.c +++ b/net/mptcp/pm.c @@ -XXX,XX +XXX,XX @@ bool mptcp_pm_is_backup(struct mptcp_sock *msk, struct sock_common *skc) return mptcp_pm_nl_is_backup(msk, &skc_local); } -static void mptcp_pm_subflows_chk_stale(const struct mptcp_sock *msk, struct sock *ssk) +static void +mptcp_pm_subflow_chk_stale(const struct mptcp_sock *msk, struct sock *ssk) { struct mptcp_subflow_context *iter, *subflow = mptcp_subflow_ctx(ssk); struct sock *sk = (struct sock *)msk; @@ -XXX,XX +XXX,XX @@ static void mptcp_pm_subflows_chk_stale(const struct mptcp_sock *msk, struct soc } } -void mptcp_pm_subflow_chk_stale(const struct mptcp_sock *msk, struct sock *ssk) +void mptcp_pm_chk_stale(const struct mptcp_sock *msk) { - struct mptcp_subflow_context *subflow = mptcp_subflow_ctx(ssk); - u32 rcv_tstamp = READ_ONCE(tcp_sk(ssk)->rcv_tstamp); + struct mptcp_subflow_context *subflow; - /* keep track of rtx periods with no progress */ - if (!subflow->stale_count) { - subflow->stale_rcv_tstamp = rcv_tstamp; - subflow->stale_count++; - } else if (subflow->stale_rcv_tstamp == rcv_tstamp) { - if (subflow->stale_count < U8_MAX) + mptcp_for_each_subflow(msk, subflow) { + struct sock *ssk = mptcp_subflow_tcp_sock(subflow); + u32 rcv_tstamp; + + if (!__mptcp_subflow_active(subflow)) + continue; + + /* No data outstanding at TCP level? not stale */ + if (tcp_rtx_and_write_queues_empty(ssk)) + continue; + + /* keep track of rtx periods with no progress */ + rcv_tstamp = READ_ONCE(tcp_sk(ssk)->rcv_tstamp); + if (!subflow->stale_count) { + subflow->stale_rcv_tstamp = rcv_tstamp; subflow->stale_count++; - mptcp_pm_subflows_chk_stale(msk, ssk); - } else { - subflow->stale_count = 0; - mptcp_subflow_set_active(subflow); + } else if (subflow->stale_rcv_tstamp == rcv_tstamp) { + if (subflow->stale_count < U8_MAX) + subflow->stale_count++; + mptcp_pm_subflow_chk_stale(msk, ssk); + } else { + subflow->stale_count = 0; + mptcp_subflow_set_active(subflow); + } } } diff --git a/net/mptcp/protocol.c b/net/mptcp/protocol.c index XXXXXXX..XXXXXXX 100644 --- a/net/mptcp/protocol.c +++ b/net/mptcp/protocol.c @@ -XXX,XX +XXX,XX @@ struct sock *mptcp_subflow_get_retrans(struct mptcp_sock *msk) /* still data outstanding at TCP level? skip this */ if (!tcp_rtx_and_write_queues_empty(ssk)) { - mptcp_pm_subflow_chk_stale(msk, ssk); min_stale_count = min_t(int, min_stale_count, subflow->stale_count); continue; } @@ -XXX,XX +XXX,XX @@ static void __mptcp_retrans(struct sock *sk) struct mptcp_data_frag *dfrag; int err, len; + mptcp_pm_chk_stale(msk); + mptcp_clean_una_wakeup(sk); - /* first check ssk: need to kick "stale" logic */ err = mptcp_sched_get_retrans(msk); dfrag = mptcp_rtx_head(sk); if (!dfrag) { diff --git a/net/mptcp/protocol.h b/net/mptcp/protocol.h index XXXXXXX..XXXXXXX 100644 --- a/net/mptcp/protocol.h +++ b/net/mptcp/protocol.h @@ -XXX,XX +XXX,XX @@ struct mptcp_subflow_context { remote_key_valid : 1, /* received the peer key from */ disposable : 1, /* ctx can be free at ulp release time */ closing : 1, /* must not pass rx data to msk anymore */ - stale : 1, /* unable to snd/rcv data, do not use for xmit */ valid_csum_seen : 1, /* at least one csum validated */ is_mptfo : 1, /* subflow is doing TFO */ close_event_done : 1, /* has done the post-closed part */ mpc_drop : 1, /* the MPC option has been dropped in a rtx */ - __unused : 8; + __unused : 9; bool data_avail; bool scheduled; bool pm_listener; /* a listener managed by the kernel PM? */ @@ -XXX,XX +XXX,XX @@ struct mptcp_subflow_context { u8 reset_seen:1; u8 reset_transient:1; u8 reset_reason:4; - u8 stale_count; + u8 stale_count; /* Protected by the msk socket lock */ + u8 stale; /* Protected by the msk socket lock, + * if set the subflow is unable to snd/rcv + * data, the schedule should skip it + */ u32 subflow_id; @@ -XXX,XX +XXX,XX @@ int mptcp_pm_parse_entry(struct nlattr *attr, struct genl_info *info, bool mptcp_pm_addr_families_match(const struct sock *sk, const struct mptcp_addr_info *loc, const struct mptcp_addr_info *rem); -void mptcp_pm_subflow_chk_stale(const struct mptcp_sock *msk, struct sock *ssk); +void mptcp_pm_chk_stale(const struct mptcp_sock *msk); void mptcp_pm_new_connection(struct mptcp_sock *msk, const struct sock *ssk, int server_side); void mptcp_pm_fully_established(struct mptcp_sock *msk, const struct sock *ssk); bool mptcp_pm_allow_new_subflow(struct mptcp_sock *msk); -- 2.53.0
From: Paolo Abeni <pabeni@redhat.com> Currently the MPTCP core enforces that when MPTCP-level retrans timer fires, at most a single dfrag is retransmitted. In some corner-cases, it may be necessary to retransmit multiple dfrags, and the MPTCP socket will need to wait multiple retrans timeout to accomplish that. Remove the mentioned constraint, allowing to transmit multiple dfrags per retrans period, as long as the scheduler keeps selecting subflows for retransmissions and pending data is available in the rtx queue. The default scheduler will transmit a dfrag per available subflow. Tested-by: Gang Yan <yangang@kylinos.cn> Tested-by: Geliang Tang <geliang@kernel.org> Acked-by: Geliang Tang <geliang@kernel.org> Signed-off-by: Paolo Abeni <pabeni@redhat.com> Signed-off-by: Matthieu Baerts (NGI0) <matttbe@kernel.org> --- v3: - Many cleanups, since the stale logic is invoked only once outside the main loop. --- net/mptcp/protocol.c | 106 ++++++++++++++++++++++++++++++++++++--------------- 1 file changed, 75 insertions(+), 31 deletions(-) diff --git a/net/mptcp/protocol.c b/net/mptcp/protocol.c index XXXXXXX..XXXXXXX 100644 --- a/net/mptcp/protocol.c +++ b/net/mptcp/protocol.c @@ -XXX,XX +XXX,XX @@ static void __mptcp_clean_una_wakeup(struct sock *sk) mptcp_write_space(sk); } -static void mptcp_clean_una_wakeup(struct sock *sk) -{ - mptcp_data_lock(sk); - __mptcp_clean_una_wakeup(sk); - mptcp_data_unlock(sk); -} - static void mptcp_enter_memory_pressure(struct sock *sk) { struct mptcp_subflow_context *subflow; @@ -XXX,XX +XXX,XX @@ static void mptcp_check_fastclose(struct mptcp_sock *msk) sk_error_report(sk); } -/* Retransmit the specified data fragment on all the selected subflows. */ -static int __mptcp_push_retrans(struct sock *sk, struct mptcp_data_frag *dfrag) +/* + * Retransmit the specified data fragment on all the selected subflows, + * starting from the specified sequence + */ +static int __mptcp_push_retrans(struct sock *sk, struct mptcp_data_frag *dfrag, + u64 sent_seq) { struct mptcp_sendmsg_info info = { .data_lock_held = true, }; struct mptcp_sock *msk = mptcp_sk(sk); @@ -XXX,XX +XXX,XX @@ static int __mptcp_push_retrans(struct sock *sk, struct mptcp_data_frag *dfrag) mptcp_for_each_subflow(msk, subflow) { if (READ_ONCE(subflow->scheduled)) { + u16 offset = sent_seq - dfrag->data_seq; u16 copied = 0; mptcp_subflow_set_scheduled(subflow, false); @@ -XXX,XX +XXX,XX @@ static int __mptcp_push_retrans(struct sock *sk, struct mptcp_data_frag *dfrag) lock_sock(ssk); /* limit retransmission to the bytes already sent on some subflows */ - info.sent = 0; + info.sent = offset; info.limit = READ_ONCE(msk->csum_enabled) ? dfrag->data_len : dfrag->already_sent; @@ -XXX,XX +XXX,XX @@ static void __mptcp_retrans(struct sock *sk) struct mptcp_sock *msk = mptcp_sk(sk); struct mptcp_subflow_context *subflow; struct mptcp_data_frag *dfrag; + u64 retrans_seq, sent_seq; + bool need_retrans; int err, len; mptcp_pm_chk_stale(msk); - mptcp_clean_una_wakeup(sk); - - err = mptcp_sched_get_retrans(msk); + /* Get an updated and consistent rtx queue status. */ + mptcp_data_lock(sk); + __mptcp_clean_una_wakeup(sk); + retrans_seq = msk->snd_una; dfrag = mptcp_rtx_head(sk); - if (!dfrag) { + need_retrans = !!dfrag; + mptcp_data_unlock(sk); + + for (;;) { + bool already_acked; + + err = mptcp_sched_get_retrans(msk); + if (err) + break; + + /* `already_sent` can be 0 for `dfrag` belonging to the RTX + * queue due to __mptcp_retransmit_pending_data(). + */ + if (!dfrag || !dfrag->already_sent) + break; + + /* Can fail only in case of fallback. */ + len = __mptcp_push_retrans(sk, dfrag, retrans_seq); + if (len < 0) + goto clear_scheduled; + + retrans_seq += len; + msk->bytes_retrans += len; + dfrag->already_sent = max_t(u16, dfrag->already_sent, + retrans_seq - dfrag->data_seq); + + /* With csum enabled, retransmission can send new data. */ + sent_seq = dfrag->already_sent + dfrag->data_seq; + if (after64(sent_seq, msk->snd_nxt)) + WRITE_ONCE(msk->snd_nxt, sent_seq); + + /* Attempt the next fragment only if the current one is + * completely retransmitted. + */ + if (before64(retrans_seq, dfrag->data_seq + dfrag->data_len)) + break; + + dfrag = list_is_last(&dfrag->list, &msk->rtx_queue) ? + NULL : list_next_entry(dfrag, list); + if (!dfrag) + break; + + /* Incoming acks can move snd_una after the current dfrag + * across loop iterations, if so start again from RTX head. + */ + mptcp_data_lock(sk); + already_acked = !before64(msk->snd_una, dfrag->data_seq + + dfrag->already_sent); + if (already_acked) { + __mptcp_clean_una_wakeup(sk); + retrans_seq = msk->snd_una; + dfrag = mptcp_rtx_head(sk); + need_retrans = !!dfrag; + } else if (after64(msk->snd_una, retrans_seq)) { + retrans_seq = msk->snd_una; + } + mptcp_data_unlock(sk); + } + + /* Attempt data-fin retransmission only when the RTX queue is empty. */ + if (!need_retrans) { if (mptcp_data_fin_enabled(msk)) { struct inet_connection_sock *icsk = inet_csk(sk); @@ -XXX,XX +XXX,XX @@ static void __mptcp_retrans(struct sock *sk) icsk->icsk_retransmits + 1); mptcp_set_datafin_timeout(sk); mptcp_send_ack(msk); - goto reset_timer; } if (!mptcp_send_head(sk)) goto clear_scheduled; - - goto reset_timer; } - if (err) - goto reset_timer; - - len = __mptcp_push_retrans(sk, dfrag); - if (len < 0) - goto clear_scheduled; - - msk->bytes_retrans += len; - dfrag->already_sent = max(dfrag->already_sent, len); - - /* With csum enabled retransmission can send new data. */ - if (after64(dfrag->already_sent + dfrag->data_seq, msk->snd_nxt)) - WRITE_ONCE(msk->snd_nxt, dfrag->already_sent + dfrag->data_seq); - reset_timer: mptcp_check_and_set_pending(sk); -- 2.53.0
From: Paolo Abeni <pabeni@redhat.com> Currently the enforcement of the rcvbuf constraint is implemented when moving the skbs into the msk receive or OoO queue, keeping the incoming skbs in the subflow queue when over limits. Under significant memory pressure the above can cause permanent data transfer stalls, as the skb needed to make forward progress can be stuck in a subflow queue. Over memory limits, drop the incoming skb, relying on MPTCP-level retransmissions. Note that fallback socket must perform the limit before the skb reaches the subflow-level queue, as dropping an in-sequence already acked skb would break the stream. This is not a complete fix for the stall issue, as the drop strategy needs refinements that will come in the next patches. Signed-off-by: Paolo Abeni <pabeni@redhat.com> Reviewed-by: Matthieu Baerts (NGI0) <matttbe@kernel.org> Signed-off-by: Matthieu Baerts (NGI0) <matttbe@kernel.org> --- v2: - mib: typo: "constrains" -> "constraints". - mptcp_over_limit: more than 0-win: retrans, dup or old acks. - mptcp_over_limit: bump LINUX_MIB_TCPRCVQDROP. - Note: Sashiko might point to a possible forward-allocated memory leak: this is a temp leak, and releasing additionally allocated fwd memory in the error path will be fix in a patch for -net. v3: - refine mptcp_over_limit() to check separately backlog and rcvbuf - do not drop rst (with data) --- net/mptcp/mib.c | 2 ++ net/mptcp/mib.h | 2 ++ net/mptcp/options.c | 32 +++++++++++++++++++++++++++++--- net/mptcp/protocol.c | 31 +++++++++++++++++++++++-------- 4 files changed, 56 insertions(+), 11 deletions(-) diff --git a/net/mptcp/mib.c b/net/mptcp/mib.c index XXXXXXX..XXXXXXX 100644 --- a/net/mptcp/mib.c +++ b/net/mptcp/mib.c @@ -XXX,XX +XXX,XX @@ static const struct snmp_mib mptcp_snmp_list[] = { SNMP_MIB_ITEM("SimultConnectFallback", MPTCP_MIB_SIMULTCONNFALLBACK), SNMP_MIB_ITEM("FallbackFailed", MPTCP_MIB_FALLBACKFAILED), SNMP_MIB_ITEM("WinProbe", MPTCP_MIB_WINPROBE), + SNMP_MIB_ITEM("BacklogDrop", MPTCP_MIB_BACKLOGDROP), + SNMP_MIB_ITEM("RcvPruned", MPTCP_MIB_RCVPRUNED), }; /* mptcp_mib_alloc - allocate percpu mib counters diff --git a/net/mptcp/mib.h b/net/mptcp/mib.h index XXXXXXX..XXXXXXX 100644 --- a/net/mptcp/mib.h +++ b/net/mptcp/mib.h @@ -XXX,XX +XXX,XX @@ enum linux_mptcp_mib_field { MPTCP_MIB_SIMULTCONNFALLBACK, /* Simultaneous connect */ MPTCP_MIB_FALLBACKFAILED, /* Can't fallback due to msk status */ MPTCP_MIB_WINPROBE, /* MPTCP-level zero window probe */ + MPTCP_MIB_BACKLOGDROP, /* Backlog over memory limit */ + MPTCP_MIB_RCVPRUNED, /* Dropped due to memory constraints */ __MPTCP_MIB_MAX }; diff --git a/net/mptcp/options.c b/net/mptcp/options.c index XXXXXXX..XXXXXXX 100644 --- a/net/mptcp/options.c +++ b/net/mptcp/options.c @@ -XXX,XX +XXX,XX @@ static bool add_addr_hmac_valid(struct mptcp_sock *msk, return hmac == mp_opt->ahmac; } -/* Return false in case of error (or subflow has been reset), - * else return true. +static bool mptcp_over_limit(struct sock *sk, struct sock *ssk, + const struct sk_buff *skb) +{ + struct mptcp_sock *msk = mptcp_sk(sk); + u32 rcvbuf = READ_ONCE(sk->sk_rcvbuf); + + if (likely((u32)sk_rmem_alloc_get(sk) <= rcvbuf && + READ_ONCE(msk->backlog_len) <= rcvbuf)) + return false; + + /* Avoid silently dropping pure acks, fin, rst or already-acked segm. */ + if (TCP_SKB_CB(skb)->seq == TCP_SKB_CB(skb)->end_seq || + TCP_SKB_CB(skb)->tcp_flags & (TCPHDR_FIN | TCPHDR_RST) || + !after(TCP_SKB_CB(skb)->end_seq, tcp_sk(ssk)->rcv_nxt)) + return false; + + /* Dropped due to memory constraints, schedule an ack. */ + inet_csk(ssk)->icsk_ack.pending |= ICSK_ACK_NOMEM | ICSK_ACK_NOW; + inet_csk_schedule_ack(ssk); + + /* Plain TCP (fallback) and skb is dropped before the TCP recv queue. */ + NET_INC_STATS(sock_net(sk), LINUX_MIB_TCPRCVQDROP); + + return true; +} + +/* Return false when the caller must drop the packet, i.e. in case of error, + * subflow has been reset, or over memory limits. */ bool mptcp_incoming_options(struct sock *sk, struct sk_buff *skb) { @@ -XXX,XX +XXX,XX @@ bool mptcp_incoming_options(struct sock *sk, struct sk_buff *skb) __mptcp_data_acked(subflow->conn); mptcp_data_unlock(subflow->conn); - return true; + return !mptcp_over_limit(subflow->conn, sk, skb); } mptcp_get_options(skb, &mp_opt); diff --git a/net/mptcp/protocol.c b/net/mptcp/protocol.c index XXXXXXX..XXXXXXX 100644 --- a/net/mptcp/protocol.c +++ b/net/mptcp/protocol.c @@ -XXX,XX +XXX,XX @@ static bool __mptcp_move_skb(struct sock *sk, struct sk_buff *skb) mptcp_borrow_fwdmem(sk, skb); + /* Can't drop packets for fallback socket this late, or the stream + * will break. + */ + if (unlikely(sk_rmem_alloc_get(sk) > READ_ONCE(sk->sk_rcvbuf)) && + !__mptcp_check_fallback(msk)) { + MPTCP_INC_STATS(sock_net(sk), MPTCP_MIB_RCVPRUNED); + mptcp_drop(sk, skb); + return false; + } + if (MPTCP_SKB_CB(skb)->map_seq == msk->ack_seq) { /* in sequence */ msk->bytes_received += copy_len; @@ -XXX,XX +XXX,XX @@ static void __mptcp_add_backlog(struct sock *sk, struct sk_buff *tail = NULL; struct sock *ssk = skb->sk; bool fragstolen; + u64 limit; int delta; if (unlikely(sk->sk_state == TCP_CLOSE)) { @@ -XXX,XX +XXX,XX @@ static void __mptcp_add_backlog(struct sock *sk, return; } + /* Similar additional allowance as plain TCP. */ + limit = READ_ONCE(sk->sk_rcvbuf); + limit += (limit >> 1) + 64 * 1024; + limit = min_t(u64, limit, UINT_MAX); + if (msk->backlog_len > limit && !__mptcp_check_fallback(msk)) { + __MPTCP_INC_STATS(sock_net(sk), MPTCP_MIB_BACKLOGDROP); + kfree_skb_reason(skb, SKB_DROP_REASON_SOCKET_BACKLOG); + return; + } + /* Try to coalesce with the last skb in our backlog */ if (!list_empty(&msk->backlog_list)) tail = list_last_entry(&msk->backlog_list, struct sk_buff, list); @@ -XXX,XX +XXX,XX @@ static bool __mptcp_move_skbs_from_subflow(struct mptcp_sock *msk, mptcp_init_skb(ssk, skb, offset, len); - if (own_msk && sk_rmem_alloc_get(sk) < sk->sk_rcvbuf) { + if (own_msk) { mptcp_subflow_lend_fwdmem(subflow, skb); ret |= __mptcp_move_skb(sk, skb); } else { @@ -XXX,XX +XXX,XX @@ static bool __mptcp_move_skbs(struct sock *sk, struct list_head *skbs, u32 *delt *delta = 0; while (1) { - /* If the msk recvbuf is full stop, don't drop */ - if (sk_rmem_alloc_get(sk) > sk->sk_rcvbuf) - break; - prefetch(skb->next); list_del(&skb->list); *delta += skb->truesize; @@ -XXX,XX +XXX,XX @@ static bool mptcp_can_spool_backlog(struct sock *sk, struct list_head *skbs) DEBUG_NET_WARN_ON_ONCE(msk->backlog_unaccounted && sk->sk_socket && mem_cgroup_from_sk(sk)); - /* Don't spool the backlog if the rcvbuf is full. */ - if (list_empty(&msk->backlog_list) || - sk_rmem_alloc_get(sk) > sk->sk_rcvbuf) + if (list_empty(&msk->backlog_list)) return false; INIT_LIST_HEAD(skbs); -- 2.53.0
From: Paolo Abeni <pabeni@redhat.com> Currently a wild producer could keep the backlog flushing operation spinning for an unbound time. Since the previous patch, the amount of data present in the backlog is hard-limited. Move the backlog len update at the end of the flush loop to prevent it spinning forever. Also, no need to splice back the remaining skbs list into the backlog, as such list is always empty after each backlog processing loop. Signed-off-by: Paolo Abeni <pabeni@redhat.com> Reviewed-by: Matthieu Baerts (NGI0) <matttbe@kernel.org> Signed-off-by: Matthieu Baerts (NGI0) <matttbe@kernel.org> --- net/mptcp/protocol.c | 21 ++++++--------------- 1 file changed, 6 insertions(+), 15 deletions(-) diff --git a/net/mptcp/protocol.c b/net/mptcp/protocol.c index XXXXXXX..XXXXXXX 100644 --- a/net/mptcp/protocol.c +++ b/net/mptcp/protocol.c @@ -XXX,XX +XXX,XX @@ static bool __mptcp_move_skbs(struct sock *sk, struct list_head *skbs, u32 *delt struct mptcp_sock *msk = mptcp_sk(sk); bool moved = false; - *delta = 0; while (1) { prefetch(skb->next); list_del(&skb->list); @@ -XXX,XX +XXX,XX @@ static bool mptcp_can_spool_backlog(struct sock *sk, struct list_head *skbs) return true; } -static void mptcp_backlog_spooled(struct sock *sk, u32 moved, - struct list_head *skbs) -{ - struct mptcp_sock *msk = mptcp_sk(sk); - - WRITE_ONCE(msk->backlog_len, msk->backlog_len - moved); - list_splice(skbs, &msk->backlog_list); -} - static bool mptcp_move_skbs(struct sock *sk) { + struct mptcp_sock *msk = mptcp_sk(sk); struct list_head skbs; bool enqueued = false; - u32 moved; + u32 moved = 0; mptcp_data_lock(sk); while (mptcp_can_spool_backlog(sk, &skbs)) { @@ -XXX,XX +XXX,XX @@ static bool mptcp_move_skbs(struct sock *sk) enqueued |= __mptcp_move_skbs(sk, &skbs, &moved); mptcp_data_lock(sk); - mptcp_backlog_spooled(sk, moved, &skbs); } + WRITE_ONCE(msk->backlog_len, msk->backlog_len - moved); mptcp_data_unlock(sk); if (enqueued && mptcp_epollin_ready(sk)) @@ -XXX,XX +XXX,XX @@ static void mptcp_release_cb(struct sock *sk) __must_hold(&sk->sk_lock.slock) { struct mptcp_sock *msk = mptcp_sk(sk); + u32 moved = 0; for (;;) { unsigned long flags = (msk->cb_flags & MPTCP_FLAGS_PROCESS_CTX_NEED); struct list_head join_list, skbs; bool spool_bl; - u32 moved; spool_bl = mptcp_can_spool_backlog(sk, &skbs); if (!flags && !spool_bl) @@ -XXX,XX +XXX,XX @@ static void mptcp_release_cb(struct sock *sk) cond_resched(); spin_lock_bh(&sk->sk_lock.slock); - if (spool_bl) - mptcp_backlog_spooled(sk, moved, &skbs); } + if (moved) + WRITE_ONCE(msk->backlog_len, msk->backlog_len - moved); if (__test_and_clear_bit(MPTCP_CLEAN_UNA, &msk->cb_flags)) __mptcp_clean_una_wakeup(sk); -- 2.53.0
From: Paolo Abeni <pabeni@redhat.com> Alike TCP, MPTCP handles in-sequence packets and partially overlapping ones in a very similar way: we can use the same path to handle both, avoiding some code duplication. This will also make the next patch simpler. Signed-off-by: Paolo Abeni <pabeni@redhat.com> Reviewed-by: Matthieu Baerts (NGI0) <matttbe@kernel.org> Signed-off-by: Matthieu Baerts (NGI0) <matttbe@kernel.org> --- v3: new --- net/mptcp/protocol.c | 31 +++++++++++++------------------ 1 file changed, 13 insertions(+), 18 deletions(-) diff --git a/net/mptcp/protocol.c b/net/mptcp/protocol.c index XXXXXXX..XXXXXXX 100644 --- a/net/mptcp/protocol.c +++ b/net/mptcp/protocol.c @@ -XXX,XX +XXX,XX @@ static bool __mptcp_move_skb(struct sock *sk, struct sk_buff *skb) if (MPTCP_SKB_CB(skb)->map_seq == msk->ack_seq) { /* in sequence */ +insert: msk->bytes_received += copy_len; WRITE_ONCE(msk->ack_seq, msk->ack_seq + copy_len); tail = skb_peek_tail(&sk->sk_receive_queue); @@ -XXX,XX +XXX,XX @@ static bool __mptcp_move_skb(struct sock *sk, struct sk_buff *skb) return false; } - /* Completely old data? */ - if (!after64(MPTCP_SKB_CB(skb)->end_seq, msk->ack_seq)) { - MPTCP_INC_STATS(sock_net(sk), MPTCP_MIB_DUPDATA); - mptcp_drop(sk, skb); - return false; + /* Partial packet */ + if (after64(MPTCP_SKB_CB(skb)->end_seq, msk->ack_seq)) { + copy_len = MPTCP_SKB_CB(skb)->end_seq - msk->ack_seq; + MPTCP_SKB_CB(skb)->offset += msk->ack_seq - + MPTCP_SKB_CB(skb)->map_seq; + MPTCP_SKB_CB(skb)->map_seq += msk->ack_seq - + MPTCP_SKB_CB(skb)->map_seq; + goto insert; } - /* Partial packet: map_seq < ack_seq < end_seq. - * Skip the already-acked bytes and enqueue the new data. - */ - copy_len = MPTCP_SKB_CB(skb)->end_seq - msk->ack_seq; - MPTCP_SKB_CB(skb)->offset += msk->ack_seq - MPTCP_SKB_CB(skb)->map_seq; - MPTCP_SKB_CB(skb)->map_seq += msk->ack_seq - - MPTCP_SKB_CB(skb)->map_seq; - msk->bytes_received += copy_len; - WRITE_ONCE(msk->ack_seq, msk->ack_seq + copy_len); - - skb_set_owner_r(skb, sk); - __skb_queue_tail(&sk->sk_receive_queue, skb); - return true; + /* Completely old data */ + MPTCP_INC_STATS(sock_net(sk), MPTCP_MIB_DUPDATA); + mptcp_drop(sk, skb); + return false; } static void mptcp_stop_rtx_timer(struct sock *sk) -- 2.53.0
From: Paolo Abeni <pabeni@redhat.com> When moving incoming skbs in the msk receive queue and the latter is above limits, prune it as needed quite alike what TCP is doing at the subflow level. The main difference relies in the stop condition: since MPTCP does not perform collapsing, it's better off dropping the bare minimum to fit the (newer) incoming packet. Signed-off-by: Paolo Abeni <pabeni@redhat.com> Tested-by: Gang Yan <yangang@kylinos.cn> Reviewed-by: Matthieu Baerts (NGI0) <matttbe@kernel.org> Signed-off-by: Matthieu Baerts (NGI0) <matttbe@kernel.org> --- v2: - Uniform the new counter with the other OFO ones. v3: - prune only for new data - reorganize the code to follow more closely TCP --- net/mptcp/mib.c | 1 + net/mptcp/mib.h | 1 + net/mptcp/protocol.c | 81 +++++++++++++++++++++++++++++++++++++++++++++------- 3 files changed, 73 insertions(+), 10 deletions(-) diff --git a/net/mptcp/mib.c b/net/mptcp/mib.c index XXXXXXX..XXXXXXX 100644 --- a/net/mptcp/mib.c +++ b/net/mptcp/mib.c @@ -XXX,XX +XXX,XX @@ static const struct snmp_mib mptcp_snmp_list[] = { SNMP_MIB_ITEM("WinProbe", MPTCP_MIB_WINPROBE), SNMP_MIB_ITEM("BacklogDrop", MPTCP_MIB_BACKLOGDROP), SNMP_MIB_ITEM("RcvPruned", MPTCP_MIB_RCVPRUNED), + SNMP_MIB_ITEM("OFOPruned", MPTCP_MIB_OFOPRUNED), }; /* mptcp_mib_alloc - allocate percpu mib counters diff --git a/net/mptcp/mib.h b/net/mptcp/mib.h index XXXXXXX..XXXXXXX 100644 --- a/net/mptcp/mib.h +++ b/net/mptcp/mib.h @@ -XXX,XX +XXX,XX @@ enum linux_mptcp_mib_field { MPTCP_MIB_WINPROBE, /* MPTCP-level zero window probe */ MPTCP_MIB_BACKLOGDROP, /* Backlog over memory limit */ MPTCP_MIB_RCVPRUNED, /* Dropped due to memory constraints */ + MPTCP_MIB_OFOPRUNED, /* MPTCP-level OoO queue pruned */ __MPTCP_MIB_MAX }; diff --git a/net/mptcp/protocol.c b/net/mptcp/protocol.c index XXXXXXX..XXXXXXX 100644 --- a/net/mptcp/protocol.c +++ b/net/mptcp/protocol.c @@ -XXX,XX +XXX,XX @@ static bool mptcp_rcvbuf_grow(struct sock *sk, u32 newval) return false; } +/* "Inspired" from the TCP version; main difference: stop as soon as the MPTCP + * socket is under memory limit. + */ +static void mptcp_prune_ofo_queue(struct sock *sk, + const struct sk_buff *in_skb) +{ + struct mptcp_sock *msk = mptcp_sk(sk); + struct rb_node *node, *prev; + bool pruned = false; + u64 mem; + + if (RB_EMPTY_ROOT(&msk->out_of_order_queue)) + return; + + node = &msk->ooo_last_skb->rbnode; + + do { + struct sk_buff *skb = rb_to_skb(node); + + /* Stop pruning if the incoming skb would land in OoO tail. */ + if (after64(MPTCP_SKB_CB(in_skb)->map_seq, + MPTCP_SKB_CB(skb)->map_seq)) + break; + + pruned = true; + prev = rb_prev(node); + rb_erase(node, &msk->out_of_order_queue); + mptcp_drop(sk, skb); + msk->ooo_last_skb = rb_to_skb(prev); + + mem = (unsigned int)sk_rmem_alloc_get(sk); + if (mem <= sk->sk_rcvbuf) + break; + + node = prev; + } while (node); + + if (pruned) + MPTCP_INC_STATS(sock_net(sk), MPTCP_MIB_OFOPRUNED); +} + +/* The stack can't drop packets for fallback socket at the msk level, or the + * stream will break. + */ +static bool mptcp_can_ingest(const struct sock *sk) +{ + return unlikely(sk_rmem_alloc_get(sk) <= READ_ONCE(sk->sk_rcvbuf)) || + __mptcp_check_fallback(mptcp_sk(sk)); +} + +static bool mptcp_try_rmem_schedule(struct sock *sk, const struct sk_buff *skb) +{ + if (!mptcp_can_ingest(sk)) { + mptcp_prune_ofo_queue(sk, skb); + return mptcp_can_ingest(sk); + } + return true; +} + /* "inspired" by tcp_data_queue_ofo(), main differences: * - use mptcp seqs * - don't cope with sacks @@ -XXX,XX +XXX,XX @@ static void mptcp_data_queue_ofo(struct mptcp_sock *msk, struct sk_buff *skb) u64 seq, end_seq, max_seq; struct sk_buff *skb1; + if (!mptcp_try_rmem_schedule(sk, skb)) { + MPTCP_INC_STATS(sock_net(sk), MPTCP_MIB_RCVPRUNED); + mptcp_drop(sk, skb); + return; + } + seq = MPTCP_SKB_CB(skb)->map_seq; end_seq = MPTCP_SKB_CB(skb)->end_seq; max_seq = atomic64_read(&msk->rcv_wnd_sent); @@ -XXX,XX +XXX,XX @@ static bool __mptcp_move_skb(struct sock *sk, struct sk_buff *skb) mptcp_borrow_fwdmem(sk, skb); - /* Can't drop packets for fallback socket this late, or the stream - * will break. - */ - if (unlikely(sk_rmem_alloc_get(sk) > READ_ONCE(sk->sk_rcvbuf)) && - !__mptcp_check_fallback(msk)) { - MPTCP_INC_STATS(sock_net(sk), MPTCP_MIB_RCVPRUNED); - mptcp_drop(sk, skb); - return false; - } - if (MPTCP_SKB_CB(skb)->map_seq == msk->ack_seq) { /* in sequence */ insert: + if (!mptcp_try_rmem_schedule(sk, skb)) { + MPTCP_INC_STATS(sock_net(sk), MPTCP_MIB_RCVPRUNED); + mptcp_drop(sk, skb); + return false; + } + msk->bytes_received += copy_len; WRITE_ONCE(msk->ack_seq, msk->ack_seq + copy_len); tail = skb_peek_tail(&sk->sk_receive_queue); -- 2.53.0