Skip to content

Commit a8e8037

Browse files
committed
reactor: retry io_uring socket send on EAGAIN
The io_uring backend can complete a socket send with -EAGAIN and leak it to the caller: when arming a poll fails the request goes to io-wq, which does not retry EAGAIN for SOCK_NONBLOCK sockets. A send can complete with -EAGAIN even without MSG_DONTWAIT; this is expected io_uring behaviour and must be retried (documentation requested in axboe/liburing#1601). Intercept the -EAGAIN completion in complete_with() and re-issue through the poll-based do_sendmsg()/do_send(), like the aio/epoll backends. The handling lives in the shared completion base classes so both io_uring backends recover from it. Observed under load on riscv64 (io_uring is the auto-selected backend): unittest-seastar-socket's test_preemptive_down() aborts with EAGAIN. Signed-off-by: Sun Yuechi <sunyuechi@iscas.ac.cn>
1 parent b2d2a3f commit a8e8037

2 files changed

Lines changed: 58 additions & 17 deletions

File tree

include/seastar/core/reactor.hh

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -146,7 +146,9 @@ class io_intent;
146146

147147
class io_completion : public kernel_completion {
148148
public:
149-
virtual void complete_with(ssize_t res) final override;
149+
// Not final: a backend can intercept the raw result (e.g. io_uring -EAGAIN)
150+
// and reissue instead of completing.
151+
virtual void complete_with(ssize_t res) override;
150152

151153
virtual void complete(size_t res) noexcept = 0;
152154
virtual void set_exception(std::exception_ptr eptr) noexcept = 0;

src/core/reactor_backend.cc

Lines changed: 55 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -1604,11 +1604,14 @@ class reactor_backend_uring_base : public reactor_backend {
16041604

16051605
class sendmsg_completion_base : public sized_promise_completion_base {
16061606
protected:
1607+
reactor_backend_uring_base& _be;
1608+
pollable_fd_state& _fd;
1609+
std::span<iovec> _iovs;
16071610
::msghdr _mh = {};
16081611
const size_t _to_write;
16091612
public:
1610-
sendmsg_completion_base(std::span<iovec> iovs, size_t to_write)
1611-
: _to_write(to_write) {
1613+
sendmsg_completion_base(reactor_backend_uring_base& be, pollable_fd_state& fd, std::span<iovec> iovs, size_t to_write)
1614+
: _be(be), _fd(fd), _iovs(iovs), _to_write(to_write) {
16121615
_mh.msg_iov = iovs.data();
16131616
_mh.msg_iovlen = std::min<size_t>(iovs.size(), IOV_MAX);
16141617
}
@@ -1618,18 +1621,27 @@ class reactor_backend_uring_base : public reactor_backend {
16181621
size_t to_write() const noexcept {
16191622
return _to_write;
16201623
}
1624+
// io_uring can complete a socket send with -EAGAIN even without
1625+
// MSG_DONTWAIT (axboe/liburing#1601); reissue via the poll-based path.
1626+
void complete_with(ssize_t res) override;
16211627
};
16221628

1629+
#if SEASTAR_API_LEVEL < 9
16231630
class send_completion_base : public sized_promise_completion_base {
16241631
protected:
1632+
reactor_backend_uring_base& _be;
1633+
pollable_fd_state& _fd;
1634+
const void* _buffer;
16251635
const size_t _to_write;
16261636
public:
1627-
explicit send_completion_base(size_t to_write)
1628-
: _to_write(to_write) {}
1637+
send_completion_base(reactor_backend_uring_base& be, pollable_fd_state& fd, const void* buffer, size_t to_write)
1638+
: _be(be), _fd(fd), _buffer(buffer), _to_write(to_write) {}
16291639
size_t to_write() const noexcept {
16301640
return _to_write;
16311641
}
1642+
void complete_with(ssize_t res) override;
16321643
};
1644+
#endif
16331645

16341646
bool do_flush_submission_ring() {
16351647
if (_has_pending_submissions) {
@@ -1711,6 +1723,16 @@ class reactor_backend_uring_base : public reactor_backend {
17111723
_r._io_sink.submit(desc.release(), std::move(req));
17121724
return fut;
17131725
}
1726+
1727+
// do_sendmsg()/do_send() are private to reactor; reach them via friendship.
1728+
future<size_t> resend_sendmsg(pollable_fd_state& fd, std::span<iovec> iovs, size_t len) {
1729+
return _r.do_sendmsg(fd, iovs, len);
1730+
}
1731+
#if SEASTAR_API_LEVEL < 9
1732+
future<size_t> resend_send(pollable_fd_state& fd, const void* buffer, size_t len) {
1733+
return _r.do_send(fd, buffer, len);
1734+
}
1735+
#endif
17141736
public:
17151737
explicit reactor_backend_uring_base(reactor& r, ::io_uring uring)
17161738
: reactor_backend(uses_blocking_io::yes, supports_aio_fdatasync::yes)
@@ -1811,6 +1833,27 @@ class reactor_backend_uring_base : public reactor_backend {
18111833
}
18121834
};
18131835

1836+
// Out-of-line: reactor_backend_uring_base must be complete to reach resend_*().
1837+
void reactor_backend_uring_base::sendmsg_completion_base::complete_with(ssize_t res) {
1838+
if (res == -EAGAIN) {
1839+
_be.resend_sendmsg(_fd, _iovs, _to_write).forward_to(std::move(_result));
1840+
delete this;
1841+
return;
1842+
}
1843+
io_completion::complete_with(res);
1844+
}
1845+
1846+
#if SEASTAR_API_LEVEL < 9
1847+
void reactor_backend_uring_base::send_completion_base::complete_with(ssize_t res) {
1848+
if (res == -EAGAIN) {
1849+
_be.resend_send(_fd, _buffer, _to_write).forward_to(std::move(_result));
1850+
delete this;
1851+
return;
1852+
}
1853+
io_completion::complete_with(res);
1854+
}
1855+
#endif
1856+
18141857
class reactor_backend_uring final : public reactor_backend_uring_base {
18151858
public:
18161859
explicit reactor_backend_uring(reactor& r)
@@ -1958,20 +2001,18 @@ class reactor_backend_uring final : public reactor_backend_uring_base {
19582001
return current_exception_as_future<size_t>();
19592002
}
19602003
}
2004+
// EAGAIN reissue lives in the base, shared with the asymmetric backend.
19612005
class sendmsg_completion final : public sendmsg_completion_base {
1962-
pollable_fd_state& _fd;
19632006
public:
1964-
sendmsg_completion(pollable_fd_state& fd, std::span<iovec> iovs, size_t len)
1965-
: sendmsg_completion_base(iovs, len)
1966-
, _fd(fd) {}
2007+
using sendmsg_completion_base::sendmsg_completion_base;
19672008
void complete(size_t bytes) noexcept final {
19682009
if (bytes == to_write()) {
19692010
_fd.speculate_epoll(EPOLLOUT);
19702011
}
19712012
sendmsg_completion_base::complete(bytes);
19722013
}
19732014
};
1974-
auto desc = std::make_unique<sendmsg_completion>(fd, iovs, len);
2015+
auto desc = std::make_unique<sendmsg_completion>(*this, fd, iovs, len);
19752016
auto req = internal::io_request::make_sendmsg(fd.fd.get(), desc->msghdr(), MSG_NOSIGNAL);
19762017
return submit_request(std::move(desc), std::move(req));
19772018
}
@@ -1995,20 +2036,18 @@ class reactor_backend_uring final : public reactor_backend_uring_base {
19952036
return current_exception_as_future<size_t>();
19962037
}
19972038
}
2039+
// As in sendmsg() above: EAGAIN reissue lives in the base.
19982040
class send_completion final : public send_completion_base {
1999-
pollable_fd_state& _fd;
20002041
public:
2001-
send_completion(pollable_fd_state& fd, size_t to_write)
2002-
: send_completion_base(to_write)
2003-
, _fd(fd) {}
2042+
using send_completion_base::send_completion_base;
20042043
void complete(size_t bytes) noexcept final {
20052044
if (bytes == to_write()) {
20062045
_fd.speculate_epoll(EPOLLOUT);
20072046
}
20082047
send_completion_base::complete(bytes);
20092048
}
20102049
};
2011-
auto desc = std::make_unique<send_completion>(fd, len);
2050+
auto desc = std::make_unique<send_completion>(*this, fd, buffer, len);
20122051
auto req = internal::io_request::make_send(fd.fd.get(), buffer, len, MSG_NOSIGNAL);
20132052
return submit_request(std::move(desc), std::move(req));
20142053
}
@@ -2350,14 +2389,14 @@ class reactor_backend_asymmetric_uring final : public reactor_backend_uring_base
23502389
}
23512390

23522391
virtual future<size_t> sendmsg(pollable_fd_state& fd, std::span<iovec> iovs, size_t len) final {
2353-
auto desc = std::make_unique<sendmsg_completion_base>(iovs, len);
2392+
auto desc = std::make_unique<sendmsg_completion_base>(*this, fd, iovs, len);
23542393
auto req = internal::io_request::make_sendmsg(fd.fd.get(), desc->msghdr(), MSG_NOSIGNAL);
23552394
return submit_request(std::move(desc), std::move(req));
23562395
}
23572396

23582397
#if SEASTAR_API_LEVEL < 9
23592398
virtual future<size_t> send(pollable_fd_state& fd, const void* buffer, size_t len) override {
2360-
auto desc = std::make_unique<send_completion_base>(len);
2399+
auto desc = std::make_unique<send_completion_base>(*this, fd, buffer, len);
23612400
auto req = internal::io_request::make_send(fd.fd.get(), buffer, len, MSG_NOSIGNAL);
23622401
return submit_request(std::move(desc), std::move(req));
23632402
}

0 commit comments

Comments
 (0)