Skip to content
Projects
Groups
Snippets
Help
Loading...
Help
Support
Keyboard shortcuts
?
Submit feedback
Contribute to GitLab
Sign in / Register
Toggle navigation
N
nghttp2
Project overview
Project overview
Details
Activity
Releases
Repository
Repository
Files
Commits
Branches
Tags
Contributors
Graph
Compare
Issues
0
Issues
0
List
Boards
Labels
Milestones
Merge Requests
0
Merge Requests
0
CI / CD
CI / CD
Pipelines
Jobs
Schedules
Analytics
Analytics
CI / CD
Repository
Value Stream
Wiki
Wiki
Snippets
Snippets
Members
Members
Collapse sidebar
Close sidebar
Activity
Graph
Create a new issue
Jobs
Commits
Issue Boards
Open sidebar
Libraries
nghttp2
Commits
01bcc72f
Commit
01bcc72f
authored
Feb 03, 2022
by
Tatsuhiro Tsujikawa
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
nghttpx: Handle EAGAIN/EWOULDBLOCK from sendmsg
parent
7ca255ff
Changes
3
Show whitespace changes
Inline
Side-by-side
Showing
3 changed files
with
199 additions
and
32 deletions
+199
-32
src/shrpx_error.h
src/shrpx_error.h
+1
-0
src/shrpx_http3_upstream.cc
src/shrpx_http3_upstream.cc
+173
-32
src/shrpx_http3_upstream.h
src/shrpx_http3_upstream.h
+25
-0
No files found.
src/shrpx_error.h
View file @
01bcc72f
...
@@ -39,6 +39,7 @@ enum ErrorCode {
...
@@ -39,6 +39,7 @@ enum ErrorCode {
SHRPX_ERR_DCONN_CANCELED
=
-
103
,
SHRPX_ERR_DCONN_CANCELED
=
-
103
,
SHRPX_ERR_RETRY
=
-
104
,
SHRPX_ERR_RETRY
=
-
104
,
SHRPX_ERR_TLS_REQUIRED
=
-
105
,
SHRPX_ERR_TLS_REQUIRED
=
-
105
,
SHRPX_ERR_SEND_BLOCKED
=
-
106
,
};
};
}
// namespace shrpx
}
// namespace shrpx
...
...
src/shrpx_http3_upstream.cc
View file @
01bcc72f
...
@@ -127,7 +127,10 @@ Http3Upstream::Http3Upstream(ClientHandler *handler)
...
@@ -127,7 +127,10 @@ Http3Upstream::Http3Upstream(ClientHandler *handler)
downstream_queue_
{
downstream_queue_size
(
handler
->
get_worker
()),
downstream_queue_
{
downstream_queue_size
(
handler
->
get_worker
()),
!
get_config
()
->
http2_proxy
},
!
get_config
()
->
http2_proxy
},
idle_close_
{
false
},
idle_close_
{
false
},
retry_close_
{
false
}
{
retry_close_
{
false
},
tx_
{
.
data
=
std
::
unique_ptr
<
uint8_t
[]
>
(
new
uint8_t
[
64
_k
]),
}
{
ev_timer_init
(
&
timer_
,
timeoutcb
,
0.
,
0.
);
ev_timer_init
(
&
timer_
,
timeoutcb
,
0.
,
0.
);
timer_
.
data
=
this
;
timer_
.
data
=
this
;
...
@@ -700,6 +703,19 @@ int Http3Upstream::init(const UpstreamAddr *faddr, const Address &remote_addr,
...
@@ -700,6 +703,19 @@ int Http3Upstream::init(const UpstreamAddr *faddr, const Address &remote_addr,
int
Http3Upstream
::
on_read
()
{
return
0
;
}
int
Http3Upstream
::
on_read
()
{
return
0
;
}
int
Http3Upstream
::
on_write
()
{
int
Http3Upstream
::
on_write
()
{
int
rv
;
if
(
tx_
.
send_blocked
)
{
rv
=
send_blocked_packet
();
if
(
rv
!=
0
)
{
return
-
1
;
}
if
(
tx_
.
send_blocked
)
{
return
0
;
}
}
if
(
write_streams
()
!=
0
)
{
if
(
write_streams
()
!=
0
)
{
return
-
1
;
return
-
1
;
}
}
...
@@ -711,14 +727,13 @@ int Http3Upstream::on_write() {
...
@@ -711,14 +727,13 @@ int Http3Upstream::on_write() {
int
Http3Upstream
::
write_streams
()
{
int
Http3Upstream
::
write_streams
()
{
std
::
array
<
nghttp3_vec
,
16
>
vec
;
std
::
array
<
nghttp3_vec
,
16
>
vec
;
std
::
array
<
uint8_t
,
64
_k
>
buf
;
auto
max_udp_payload_size
=
std
::
min
(
auto
max_udp_payload_size
=
std
::
min
(
max_udp_payload_size_
,
ngtcp2_conn_get_path_max_udp_payload_size
(
conn_
));
max_udp_payload_size_
,
ngtcp2_conn_get_path_max_udp_payload_size
(
conn_
));
size_t
max_pktcnt
=
size_t
max_pktcnt
=
std
::
min
(
static_cast
<
size_t
>
(
64
_k
),
ngtcp2_conn_get_send_quantum
(
conn_
))
/
std
::
min
(
static_cast
<
size_t
>
(
64
_k
),
ngtcp2_conn_get_send_quantum
(
conn_
))
/
max_udp_payload_size
;
max_udp_payload_size
;
ngtcp2_pkt_info
pi
,
prev_pi
;
ngtcp2_pkt_info
pi
,
prev_pi
;
uint8_t
*
bufpos
=
buf
.
data
();
uint8_t
*
bufpos
=
tx_
.
data
.
get
();
ngtcp2_path_storage
ps
,
prev_ps
;
ngtcp2_path_storage
ps
,
prev_ps
;
size_t
pktcnt
=
0
;
size_t
pktcnt
=
0
;
int
rv
;
int
rv
;
...
@@ -803,8 +818,6 @@ int Http3Upstream::write_streams() {
...
@@ -803,8 +818,6 @@ int Http3Upstream::write_streams() {
last_error_
=
quic
::
err_transport
(
nwrite
);
last_error_
=
quic
::
err_transport
(
nwrite
);
handler_
->
get_connection
()
->
wlimit
.
stopw
();
return
handle_error
();
return
handle_error
();
}
else
if
(
ndatalen
>=
0
)
{
}
else
if
(
ndatalen
>=
0
)
{
rv
=
nghttp3_conn_add_write_offset
(
httpconn_
,
stream_id
,
ndatalen
);
rv
=
nghttp3_conn_add_write_offset
(
httpconn_
,
stream_id
,
ndatalen
);
...
@@ -817,12 +830,26 @@ int Http3Upstream::write_streams() {
...
@@ -817,12 +830,26 @@ int Http3Upstream::write_streams() {
}
}
if
(
nwrite
==
0
)
{
if
(
nwrite
==
0
)
{
if
(
bufpos
-
buf
.
data
())
{
if
(
bufpos
-
tx_
.
data
.
get
())
{
send_packet
(
static_cast
<
UpstreamAddr
*>
(
prev_ps
.
path
.
user_data
),
auto
faddr
=
static_cast
<
UpstreamAddr
*>
(
prev_ps
.
path
.
user_data
);
prev_ps
.
path
.
remote
.
addr
,
prev_ps
.
path
.
remote
.
addrlen
,
auto
data
=
tx_
.
data
.
get
();
prev_ps
.
path
.
local
.
addr
,
prev_ps
.
path
.
local
.
addrlen
,
auto
datalen
=
bufpos
-
data
;
prev_pi
,
buf
.
data
(),
bufpos
-
buf
.
data
(),
rv
=
send_packet
(
faddr
,
prev_ps
.
path
.
remote
.
addr
,
prev_ps
.
path
.
remote
.
addrlen
,
prev_ps
.
path
.
local
.
addr
,
prev_ps
.
path
.
local
.
addrlen
,
prev_pi
,
data
,
datalen
,
max_udp_payload_size
);
max_udp_payload_size
);
if
(
rv
==
SHRPX_ERR_SEND_BLOCKED
)
{
on_send_blocked
(
faddr
,
prev_ps
.
path
.
remote
,
prev_ps
.
path
.
local
,
prev_pi
,
data
,
datalen
,
max_udp_payload_size
);
ngtcp2_conn_update_pkt_tx_time
(
conn_
,
ts
);
reset_idle_timer
();
signal_write_upstream_addr
(
faddr
);
return
0
;
}
reset_idle_timer
();
reset_idle_timer
();
}
}
...
@@ -842,54 +869,99 @@ int Http3Upstream::write_streams() {
...
@@ -842,54 +869,99 @@ int Http3Upstream::write_streams() {
prev_pi
=
pi
;
prev_pi
=
pi
;
}
else
if
(
!
ngtcp2_path_eq
(
&
prev_ps
.
path
,
&
ps
.
path
)
||
}
else
if
(
!
ngtcp2_path_eq
(
&
prev_ps
.
path
,
&
ps
.
path
)
||
prev_pi
.
ecn
!=
pi
.
ecn
)
{
prev_pi
.
ecn
!=
pi
.
ecn
)
{
send_packet
(
static_cast
<
UpstreamAddr
*>
(
prev_ps
.
path
.
user_data
),
auto
faddr
=
static_cast
<
UpstreamAddr
*>
(
prev_ps
.
path
.
user_data
);
prev_ps
.
path
.
remote
.
addr
,
prev_ps
.
path
.
remote
.
addrlen
,
auto
data
=
tx_
.
data
.
get
();
prev_ps
.
path
.
local
.
addr
,
prev_ps
.
path
.
local
.
addrlen
,
prev_pi
,
auto
datalen
=
bufpos
-
data
-
nwrite
;
buf
.
data
(),
bufpos
-
buf
.
data
()
-
nwrite
,
rv
=
send_packet
(
faddr
,
prev_ps
.
path
.
remote
.
addr
,
prev_ps
.
path
.
remote
.
addrlen
,
prev_ps
.
path
.
local
.
addr
,
prev_ps
.
path
.
local
.
addrlen
,
prev_pi
,
data
,
datalen
,
max_udp_payload_size
);
max_udp_payload_size
);
switch
(
rv
)
{
case
SHRPX_ERR_SEND_BLOCKED
:
on_send_blocked
(
faddr
,
prev_ps
.
path
.
remote
,
prev_ps
.
path
.
local
,
prev_pi
,
data
,
datalen
,
max_udp_payload_size
);
send_packet
(
static_cast
<
UpstreamAddr
*>
(
ps
.
path
.
user_data
),
on_send_blocked
(
static_cast
<
UpstreamAddr
*>
(
ps
.
path
.
user_data
),
ps
.
path
.
remote
.
addr
,
ps
.
path
.
remote
.
addrlen
,
ps
.
path
.
remote
,
ps
.
path
.
local
,
pi
,
bufpos
-
nwrite
,
ps
.
path
.
local
.
addr
,
ps
.
path
.
local
.
addrlen
,
pi
,
nwrite
,
max_udp_payload_size
);
bufpos
-
nwrite
,
nwrite
,
max_udp_payload_size
);
signal_write_upstream_addr
(
faddr
);
break
;
default:
{
auto
faddr
=
static_cast
<
UpstreamAddr
*>
(
ps
.
path
.
user_data
);
auto
data
=
bufpos
-
nwrite
;
rv
=
send_packet
(
faddr
,
ps
.
path
.
remote
.
addr
,
ps
.
path
.
remote
.
addrlen
,
ps
.
path
.
local
.
addr
,
ps
.
path
.
local
.
addrlen
,
pi
,
data
,
nwrite
,
max_udp_payload_size
);
if
(
rv
==
SHRPX_ERR_SEND_BLOCKED
)
{
on_send_blocked
(
faddr
,
ps
.
path
.
remote
,
ps
.
path
.
local
,
pi
,
data
,
nwrite
,
max_udp_payload_size
);
}
signal_write_upstream_addr
(
faddr
);
}
}
ngtcp2_conn_update_pkt_tx_time
(
conn_
,
ts
);
ngtcp2_conn_update_pkt_tx_time
(
conn_
,
ts
);
reset_idle_timer
();
reset_idle_timer
();
handler_
->
signal_write
();
return
0
;
return
0
;
}
}
if
(
++
pktcnt
==
max_pktcnt
||
if
(
++
pktcnt
==
max_pktcnt
||
static_cast
<
size_t
>
(
nwrite
)
<
max_udp_payload_size
)
{
static_cast
<
size_t
>
(
nwrite
)
<
max_udp_payload_size
)
{
send_packet
(
static_cast
<
UpstreamAddr
*>
(
ps
.
path
.
user_data
),
auto
faddr
=
static_cast
<
UpstreamAddr
*>
(
ps
.
path
.
user_data
);
ps
.
path
.
remote
.
addr
,
ps
.
path
.
remote
.
addrlen
,
auto
data
=
tx_
.
data
.
get
();
ps
.
path
.
local
.
addr
,
ps
.
path
.
local
.
addrlen
,
pi
,
buf
.
data
(),
auto
datalen
=
bufpos
-
data
;
bufpos
-
buf
.
data
(),
max_udp_payload_size
);
rv
=
send_packet
(
faddr
,
ps
.
path
.
remote
.
addr
,
ps
.
path
.
remote
.
addrlen
,
ps
.
path
.
local
.
addr
,
ps
.
path
.
local
.
addrlen
,
pi
,
data
,
datalen
,
max_udp_payload_size
);
if
(
rv
==
SHRPX_ERR_SEND_BLOCKED
)
{
on_send_blocked
(
faddr
,
ps
.
path
.
remote
,
ps
.
path
.
local
,
pi
,
data
,
datalen
,
max_udp_payload_size
);
}
ngtcp2_conn_update_pkt_tx_time
(
conn_
,
ts
);
ngtcp2_conn_update_pkt_tx_time
(
conn_
,
ts
);
reset_idle_timer
();
reset_idle_timer
();
handler_
->
signal_write
(
);
signal_write_upstream_addr
(
faddr
);
return
0
;
return
0
;
}
}
#else // !UDP_SEGMENT
#else // !UDP_SEGMENT
send_packet
(
static_cast
<
UpstreamAddr
*>
(
ps
.
path
.
user_data
),
auto
faddr
=
static_cast
<
UpstreamAddr
*>
(
ps
.
path
.
user_data
);
ps
.
path
.
remote
.
addr
,
ps
.
path
.
remote
.
addrlen
,
ps
.
path
.
local
.
addr
,
auto
data
=
tx_
.
data
.
get
();
ps
.
path
.
local
.
addrlen
,
pi
,
buf
.
data
(),
bufpos
-
buf
.
data
(),
0
);
auto
datalen
=
bufpos
-
data
;
rv
=
send_packet
(
faddr
,
ps
.
path
.
remote
.
addr
,
ps
.
path
.
remote
.
addrlen
,
ps
.
path
.
local
.
addr
,
ps
.
path
.
local
.
addrlen
,
pi
,
data
,
datalen
,
0
);
if
(
rv
==
SHRPX_ERR_SEND_BLOCKED
)
{
on_send_blocked
(
faddr
,
ps
.
path
.
remote
,
ps
.
path
.
local
,
pi
,
data
,
datalen
,
0
);
ngtcp2_conn_update_pkt_tx_time
(
conn_
,
ts
);
reset_idle_timer
();
signal_write_upstream_addr
(
faddr
);
return
0
;
}
if
(
++
pktcnt
==
max_pktcnt
)
{
if
(
++
pktcnt
==
max_pktcnt
)
{
ngtcp2_conn_update_pkt_tx_time
(
conn_
,
ts
);
ngtcp2_conn_update_pkt_tx_time
(
conn_
,
ts
);
reset_idle_timer
();
reset_idle_timer
();
handler_
->
signal_write
(
);
signal_write_upstream_addr
(
faddr
);
return
0
;
return
0
;
}
}
bufpos
=
buf
.
data
();
bufpos
=
tx_
.
data
.
get
();
#endif // !UDP_SEGMENT
#endif // !UDP_SEGMENT
}
}
...
@@ -1760,6 +1832,11 @@ int Http3Upstream::send_packet(const UpstreamAddr *faddr,
...
@@ -1760,6 +1832,11 @@ int Http3Upstream::send_packet(const UpstreamAddr *faddr,
case
-
EMSGSIZE
:
case
-
EMSGSIZE
:
max_udp_payload_size_
=
NGTCP2_MAX_UDP_PAYLOAD_SIZE
;
max_udp_payload_size_
=
NGTCP2_MAX_UDP_PAYLOAD_SIZE
;
break
;
break
;
case
-
EAGAIN
:
#if EAGAIN != EWOULDBLOCK
case
-
EWOULDBLOCK
:
#endif // EAGAIN != EWOULDBLOCK
return
SHRPX_ERR_SEND_BLOCKED
;
default:
default:
break
;
break
;
}
}
...
@@ -1767,6 +1844,70 @@ int Http3Upstream::send_packet(const UpstreamAddr *faddr,
...
@@ -1767,6 +1844,70 @@ int Http3Upstream::send_packet(const UpstreamAddr *faddr,
return
-
1
;
return
-
1
;
}
}
void
Http3Upstream
::
on_send_blocked
(
const
UpstreamAddr
*
faddr
,
const
ngtcp2_addr
&
remote_addr
,
const
ngtcp2_addr
&
local_addr
,
const
ngtcp2_pkt_info
&
pi
,
const
uint8_t
*
data
,
size_t
datalen
,
size_t
max_udp_payload_size
)
{
assert
(
tx_
.
num_blocked
||
!
tx_
.
send_blocked
);
assert
(
tx_
.
num_blocked
<
2
);
tx_
.
send_blocked
=
true
;
auto
&
p
=
tx_
.
blocked
[
tx_
.
num_blocked
++
];
memcpy
(
&
p
.
local_addr
.
su
,
local_addr
.
addr
,
local_addr
.
addrlen
);
memcpy
(
&
p
.
remote_addr
.
su
,
remote_addr
.
addr
,
remote_addr
.
addrlen
);
p
.
local_addr
.
len
=
local_addr
.
addrlen
;
p
.
remote_addr
.
len
=
remote_addr
.
addrlen
;
p
.
faddr
=
faddr
;
p
.
pi
=
pi
;
p
.
data
=
data
;
p
.
datalen
=
datalen
;
p
.
max_udp_payload_size
=
max_udp_payload_size
;
}
int
Http3Upstream
::
send_blocked_packet
()
{
int
rv
;
assert
(
tx_
.
send_blocked
);
for
(;
tx_
.
num_blocked_sent
<
tx_
.
num_blocked
;
++
tx_
.
num_blocked_sent
)
{
auto
&
p
=
tx_
.
blocked
[
tx_
.
num_blocked_sent
];
rv
=
send_packet
(
p
.
faddr
,
&
p
.
remote_addr
.
su
.
sa
,
p
.
remote_addr
.
len
,
&
p
.
local_addr
.
su
.
sa
,
p
.
local_addr
.
len
,
p
.
pi
,
p
.
data
,
p
.
datalen
,
p
.
max_udp_payload_size
);
if
(
rv
==
SHRPX_ERR_SEND_BLOCKED
)
{
signal_write_upstream_addr
(
p
.
faddr
);
return
0
;
}
}
tx_
.
send_blocked
=
false
;
tx_
.
num_blocked
=
0
;
tx_
.
num_blocked_sent
=
0
;
return
0
;
}
void
Http3Upstream
::
signal_write_upstream_addr
(
const
UpstreamAddr
*
faddr
)
{
auto
conn
=
handler_
->
get_connection
();
if
(
faddr
->
fd
!=
conn
->
wev
.
fd
)
{
if
(
ev_is_active
(
&
conn
->
wev
))
{
ev_io_stop
(
handler_
->
get_loop
(),
&
conn
->
wev
);
}
ev_io_set
(
&
conn
->
wev
,
faddr
->
fd
,
EV_WRITE
);
}
conn
->
wlimit
.
startw
();
}
int
Http3Upstream
::
handle_error
()
{
int
Http3Upstream
::
handle_error
()
{
if
(
ngtcp2_conn_is_in_closing_period
(
conn_
))
{
if
(
ngtcp2_conn_is_in_closing_period
(
conn_
))
{
return
-
1
;
return
-
1
;
...
...
src/shrpx_http3_upstream.h
View file @
01bcc72f
...
@@ -157,6 +157,14 @@ public:
...
@@ -157,6 +157,14 @@ public:
void
qlog_write
(
const
void
*
data
,
size_t
datalen
,
bool
fin
);
void
qlog_write
(
const
void
*
data
,
size_t
datalen
,
bool
fin
);
int
open_qlog_file
(
const
StringRef
&
dir
,
const
ngtcp2_cid
&
scid
)
const
;
int
open_qlog_file
(
const
StringRef
&
dir
,
const
ngtcp2_cid
&
scid
)
const
;
void
on_send_blocked
(
const
UpstreamAddr
*
faddr
,
const
ngtcp2_addr
&
remote_addr
,
const
ngtcp2_addr
&
local_addr
,
const
ngtcp2_pkt_info
&
pi
,
const
uint8_t
*
data
,
size_t
datalen
,
size_t
max_udp_payload_size
);
int
send_blocked_packet
();
void
signal_write_upstream_addr
(
const
UpstreamAddr
*
faddr
);
private:
private:
ClientHandler
*
handler_
;
ClientHandler
*
handler_
;
ev_timer
timer_
;
ev_timer
timer_
;
...
@@ -174,6 +182,23 @@ private:
...
@@ -174,6 +182,23 @@ private:
bool
idle_close_
;
bool
idle_close_
;
bool
retry_close_
;
bool
retry_close_
;
std
::
vector
<
uint8_t
>
conn_close_
;
std
::
vector
<
uint8_t
>
conn_close_
;
struct
{
bool
send_blocked
;
size_t
num_blocked
;
size_t
num_blocked_sent
;
// blocked field is effective only when send_blocked is true.
struct
{
const
UpstreamAddr
*
faddr
;
Address
local_addr
;
Address
remote_addr
;
ngtcp2_pkt_info
pi
;
const
uint8_t
*
data
;
size_t
datalen
;
size_t
max_udp_payload_size
;
}
blocked
[
2
];
std
::
unique_ptr
<
uint8_t
[]
>
data
;
}
tx_
;
}
;
}
;
}
// namespace shrpx
}
// namespace shrpx
...
...
Write
Preview
Markdown
is supported
0%
Try again
or
attach a new file
Attach a file
Cancel
You are about to add
0
people
to the discussion. Proceed with caution.
Finish editing this message first!
Cancel
Please
register
or
sign in
to comment