Skip to content
This repository was archived by the owner on Feb 22, 2022. It is now read-only.
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 8 additions & 4 deletions funcx_forwarder/taskqueue.py
Original file line number Diff line number Diff line change
Expand Up @@ -63,21 +63,25 @@ def __init__(self,
self.zmq_socket.set(zmq.ROUTER_MANDATORY, 1)
self.zmq_socket.set(zmq.ROUTER_HANDOVER, 1)

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

We need this because we want new zmq socket ids (we use endpoint id) to replace the old sockets with the same id. Not doing so means an old socket with a given id can block a new socket from communicating that is given that same id on a new registration.

@knagaitsev knagaitsev Aug 5, 2021 •

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I believe this was not an issue with the tasks_q from the forwarder because since it was outgoing the forwarder could realize the client was not reachable, so when it sends heartbeats over this channel it simply thinks the same client socket has reconnected. This is why heartbeats were still reaching newly connected sockets, since they were going over this channel.

However, with the results_q, the forwarder has no way of realizing that the socket is unreachable since it is not sending messages over this channel, only receiving messages. So it still thinks the old socket is connected when a new socket tries to connect.

self.setup_server_auth()
self.zmq_socket.bind("tcp://*:{}".format(port))
elif self.mode == 'client':
self.zmq_socket = self.context.socket(zmq.DEALER)
self.setup_client_auth()
self.zmq_socket.setsockopt(zmq.IDENTITY, identity.encode('utf-8'))
self.zmq_socket.connect("tcp://{}:{}".format(address, port))
else:
raise ValueError("TaskQueue must be initialized with mode set to 'server' or 'client'")

if set_hwm:
self.zmq_socket.set_hwm(0)
if RCVTIMEO is not None:
self.zmq_socket.RCVTIMEO = RCVTIMEO
self.zmq_socket.setsockopt(zmq.RCVTIMEO, RCVTIMEO)
if SNDTIMEO is not None:
self.zmq_socket.SNDTIMEO = SNDTIMEO
self.zmq_socket.setsockopt(zmq.SNDTIMEO, SNDTIMEO)

# all zmq setsockopt calls must be done before bind/connect is called
if self.mode == 'server':
self.zmq_socket.bind("tcp://*:{}".format(port))
elif self.mode == 'client':
self.zmq_socket.connect("tcp://{}:{}".format(address, port))

self.poller = zmq.Poller()
self.poller.register(self.zmq_socket, zmq.POLLOUT)
Expand Down