Skip to content
Projects
Groups
Snippets
Help
Loading...
Sign in / Register
Toggle navigation
C
capnproto
Project
Project
Details
Activity
Cycle Analytics
Repository
Repository
Files
Commits
Branches
Tags
Contributors
Graph
Compare
Charts
Issues
0
Issues
0
List
Board
Labels
Milestones
Merge Requests
0
Merge Requests
0
CI / CD
CI / CD
Pipelines
Jobs
Schedules
Charts
Packages
Packages
Wiki
Wiki
Snippets
Snippets
Members
Members
Collapse sidebar
Close sidebar
Activity
Graph
Charts
Create a new issue
Jobs
Commits
Issue Boards
Open sidebar
submodule
capnproto
Commits
559044bf
Commit
559044bf
authored
Dec 01, 2016
by
Kenton Varda
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
Make wrapConnectingSocketFd() have same signature on Unix as Windows.
parent
f044b0a6
Show whitespace changes
Inline
Side-by-side
Showing
2 changed files
with
36 additions
and
55 deletions
+36
-55
async-io.c++
c++/src/kj/async-io.c++
+34
-45
async-io.h
c++/src/kj/async-io.h
+2
-10
No files found.
c++/src/kj/async-io.c++
View file @
559044bf
...
@@ -416,27 +416,6 @@ public:
...
@@ -416,27 +416,6 @@ public:
KJ_SYSCALL
(
::
bind
(
sockfd
,
&
addr
.
generic
,
addrlen
),
toString
());
KJ_SYSCALL
(
::
bind
(
sockfd
,
&
addr
.
generic
,
addrlen
),
toString
());
}
}
void
connect
(
int
sockfd
)
const
{
// Unfortunately connect() doesn't fit the mold of KJ_NONBLOCKING_SYSCALL, since it indicates
// non-blocking using EINPROGRESS.
for
(;;)
{
if
(
::
connect
(
sockfd
,
&
addr
.
generic
,
addrlen
)
<
0
)
{
int
error
=
errno
;
if
(
error
==
EINPROGRESS
)
{
return
;
}
else
if
(
error
!=
EINTR
)
{
KJ_FAIL_SYSCALL
(
"connect()"
,
error
,
toString
())
{
// Recover by returning, since reads/writes will simply fail.
return
;
}
}
}
else
{
// no error
return
;
}
}
}
uint
getPort
()
const
{
uint
getPort
()
const
{
switch
(
addr
.
generic
.
sa_family
)
{
switch
(
addr
.
generic
.
sa_family
)
{
case
AF_INET
:
return
ntohs
(
addr
.
inet4
.
sin_port
);
case
AF_INET
:
return
ntohs
(
addr
.
inet4
.
sin_port
);
...
@@ -901,7 +880,26 @@ public:
...
@@ -901,7 +880,26 @@ public:
Own
<
AsyncIoStream
>
wrapSocketFd
(
int
fd
,
uint
flags
=
0
)
override
{
Own
<
AsyncIoStream
>
wrapSocketFd
(
int
fd
,
uint
flags
=
0
)
override
{
return
heap
<
AsyncStreamFd
>
(
eventPort
,
fd
,
flags
);
return
heap
<
AsyncStreamFd
>
(
eventPort
,
fd
,
flags
);
}
}
Promise
<
Own
<
AsyncIoStream
>>
wrapConnectingSocketFd
(
int
fd
,
uint
flags
=
0
)
override
{
Promise
<
Own
<
AsyncIoStream
>>
wrapConnectingSocketFd
(
int
fd
,
const
struct
sockaddr
*
addr
,
uint
addrlen
,
uint
flags
=
0
)
override
{
// Unfortunately connect() doesn't fit the mold of KJ_NONBLOCKING_SYSCALL, since it indicates
// non-blocking using EINPROGRESS.
for
(;;)
{
if
(
::
connect
(
fd
,
addr
,
addrlen
)
<
0
)
{
int
error
=
errno
;
if
(
error
==
EINPROGRESS
)
{
// Fine.
break
;
}
else
if
(
error
!=
EINTR
)
{
KJ_FAIL_SYSCALL
(
"connect()"
,
error
)
{
break
;
}
return
Own
<
AsyncIoStream
>
();
}
}
else
{
// no error
break
;
}
}
auto
result
=
heap
<
AsyncStreamFd
>
(
eventPort
,
fd
,
flags
);
auto
result
=
heap
<
AsyncStreamFd
>
(
eventPort
,
fd
,
flags
);
auto
connected
=
result
->
waitConnected
();
auto
connected
=
result
->
waitConnected
();
...
@@ -940,7 +938,9 @@ public:
...
@@ -940,7 +938,9 @@ public:
:
lowLevel
(
lowLevel
),
addrs
(
kj
::
mv
(
addrs
))
{}
:
lowLevel
(
lowLevel
),
addrs
(
kj
::
mv
(
addrs
))
{}
Promise
<
Own
<
AsyncIoStream
>>
connect
()
override
{
Promise
<
Own
<
AsyncIoStream
>>
connect
()
override
{
return
connectImpl
(
0
);
auto
addrsCopy
=
heapArray
(
addrs
.
asPtr
());
auto
promise
=
connectImpl
(
lowLevel
,
addrsCopy
);
return
promise
.
attach
(
kj
::
mv
(
addrsCopy
));
}
}
Own
<
ConnectionReceiver
>
listen
()
override
{
Own
<
ConnectionReceiver
>
listen
()
override
{
...
@@ -1010,34 +1010,23 @@ private:
...
@@ -1010,34 +1010,23 @@ private:
Array
<
SocketAddress
>
addrs
;
Array
<
SocketAddress
>
addrs
;
uint
counter
=
0
;
uint
counter
=
0
;
Promise
<
Own
<
AsyncIoStream
>>
connectImpl
(
uint
index
)
{
static
Promise
<
Own
<
AsyncIoStream
>>
connectImpl
(
KJ_ASSERT
(
index
<
addrs
.
size
());
LowLevelAsyncIoProvider
&
lowLevel
,
ArrayPtr
<
SocketAddress
>
addrs
)
{
KJ_ASSERT
(
addrs
.
size
()
>
0
);
int
fd
=
addrs
[
index
].
socket
(
SOCK_STREAM
);
KJ_IF_MAYBE
(
exception
,
kj
::
runCatchingExceptions
([
&
]()
{
int
fd
=
addrs
[
0
].
socket
(
SOCK_STREAM
);
addrs
[
index
].
connect
(
fd
);
}))
{
// Connect failed.
close
(
fd
);
if
(
index
+
1
<
addrs
.
size
())
{
// Try the next address instead.
return
connectImpl
(
index
+
1
);
}
else
{
// No more addresses to try, so propagate the exception.
return
kj
::
mv
(
*
exception
);
}
}
return
lowLevel
.
wrapConnectingSocketFd
(
fd
,
NEW_FD_FLAGS
).
then
(
return
kj
::
evalNow
([
&
]()
{
[](
Own
<
AsyncIoStream
>&&
stream
)
->
Promise
<
Own
<
AsyncIoStream
>>
{
return
lowLevel
.
wrapConnectingSocketFd
(
fd
,
addrs
[
0
].
getRaw
(),
addrs
[
0
].
getRawSize
(),
NEW_FD_FLAGS
);
}).
then
([](
Own
<
AsyncIoStream
>&&
stream
)
->
Promise
<
Own
<
AsyncIoStream
>>
{
// Success, pass along.
// Success, pass along.
return
kj
::
mv
(
stream
);
return
kj
::
mv
(
stream
);
},
[
this
,
index
](
Exception
&&
exception
)
->
Promise
<
Own
<
AsyncIoStream
>>
{
},
[
&
lowLevel
,
addrs
](
Exception
&&
exception
)
mutable
->
Promise
<
Own
<
AsyncIoStream
>>
{
// Connect failed.
// Connect failed.
if
(
index
+
1
<
addrs
.
size
()
)
{
if
(
addrs
.
size
()
>
1
)
{
// Try the next address instead.
// Try the next address instead.
return
connectImpl
(
index
+
1
);
return
connectImpl
(
lowLevel
,
addrs
.
slice
(
1
,
addrs
.
size
())
);
}
else
{
}
else
{
// No more addresses to try, so propagate the exception.
// No more addresses to try, so propagate the exception.
return
kj
::
mv
(
exception
);
return
kj
::
mv
(
exception
);
...
...
c++/src/kj/async-io.h
View file @
559044bf
...
@@ -425,20 +425,12 @@ public:
...
@@ -425,20 +425,12 @@ public:
//
//
// `flags` is a bitwise-OR of the values of the `Flags` enum.
// `flags` is a bitwise-OR of the values of the `Flags` enum.
#if _WIN32
virtual
Promise
<
Own
<
AsyncIoStream
>>
wrapConnectingSocketFd
(
virtual
Promise
<
Own
<
AsyncIoStream
>>
wrapConnectingSocketFd
(
Fd
fd
,
const
struct
sockaddr
*
addr
,
uint
addrlen
,
uint
flags
=
0
)
=
0
;
Fd
fd
,
const
struct
sockaddr
*
addr
,
uint
addrlen
,
uint
flags
=
0
)
=
0
;
#else
// Create an AsyncIoStream wrapping a socket and initiate a connection to the given address.
virtual
Promise
<
Own
<
AsyncIoStream
>>
wrapConnectingSocketFd
(
Fd
fd
,
uint
flags
=
0
)
=
0
;
// The returned promise does not resolve until connection has completed.
#endif
// Create an AsyncIoStream wrapping a socket that is in the process of connecting. The returned
// promise should not resolve until connection has completed -- traditionally indicated by the
// descriptor becoming writable.
//
//
// `flags` is a bitwise-OR of the values of the `Flags` enum.
// `flags` is a bitwise-OR of the values of the `Flags` enum.
//
// On Windows, the callee initiates connect rather than the caller.
// TODO(now): Maybe on all systems?
virtual
Own
<
ConnectionReceiver
>
wrapListenSocketFd
(
Fd
fd
,
uint
flags
=
0
)
=
0
;
virtual
Own
<
ConnectionReceiver
>
wrapListenSocketFd
(
Fd
fd
,
uint
flags
=
0
)
=
0
;
// Create an AsyncIoStream wrapping a listen socket file descriptor. This socket should already
// Create an AsyncIoStream wrapping a listen socket file descriptor. This socket should already
...
...
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