Skip to content

Commit dc81ea1

Browse files
Switch rabbitmq message nacks for rejections (#266)
A way to resolve the rabbitmq infinite loops we have been seeing. This creates a new functions based on python-workflows methods to reject messages rather than nack them. Then replaces (almost) all instances of nack with reject. I'm not certain I've got every requeue choice right but the general policy is: - requeue=False if the input message is wrong (e.g. missing parameter, malformed directory) - requeue=True if the failure is external (e.g. failed subprocess call, missing file) There is a resort back to nacking if not using pika transport, mostly to avoid having to rewrite all the tests
1 parent 701f75c commit dc81ea1

33 files changed

Lines changed: 213 additions & 150 deletions

src/cryoemservices/services/bfactor_setup.py

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -54,7 +54,7 @@ def bfactor_setup(self, rw, header: dict, message: dict):
5454
self.log.info("Received a simple message")
5555
if not isinstance(message, dict):
5656
self.log.error("Rejected invalid simple message")
57-
self._transport.nack(header)
57+
self._reject_message(header, requeue=False)
5858
return
5959

6060
# Create a wrapper-like object that can be passed to functions
@@ -77,7 +77,7 @@ def bfactor_setup(self, rw, header: dict, message: dict):
7777
f"and recipe parameters: {rw.recipe_step.get('parameters', {})} "
7878
f"with exception: {e}"
7979
)
80-
rw.transport.nack(header)
80+
self._reject_message(header, transport=rw.transport, requeue=False)
8181
return
8282

8383
self.log.info(

src/cryoemservices/services/class2d.py

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -33,7 +33,7 @@ def class2d(self, rw, header: dict, message: dict):
3333
self.log.info("Received a simple message")
3434
if not isinstance(message, dict):
3535
self.log.error("Rejected invalid simple message")
36-
self._transport.nack(header)
36+
self._reject_message(header, requeue=False)
3737
return
3838

3939
# Create a wrapper-like object that can be passed to functions
@@ -57,14 +57,14 @@ def class2d(self, rw, header: dict, message: dict):
5757
f"and recipe parameters: {rw.recipe_step.get('parameters', {})} "
5858
f"with exception: {e}"
5959
)
60-
rw.transport.nack(header)
60+
self._reject_message(header, transport=rw.transport, requeue=False)
6161
return
6262

63-
# In this setup we cannot nack messages on failure, so instead check here
63+
# In this setup we cannot reject messages on failure, so instead check here
6464
if message.get("requeue", 0) >= 5:
65-
self.log.warning(f"Nacking requeued file {class2d_params.particles_file}")
66-
rw.transport.nack(header)
67-
return False
65+
self.log.warning(f"Rejecting requeued file {class2d_params.particles_file}")
66+
self._reject_message(header, transport=rw.transport, requeue=False)
67+
return
6868

6969
# Acknowledge the message and disconnect from rabbitmq
7070
self.log.info(

src/cryoemservices/services/class3d.py

Lines changed: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -32,8 +32,8 @@ def class3d(self, rw, header: dict, message: dict):
3232
if not rw:
3333
self.log.info("Received a simple message")
3434
if not isinstance(message, dict):
35-
self.log.error("Rejected invalid simple message")
36-
self._transport.nack(header)
35+
self.log.error("Rejected invaid simple message")
36+
self._reject_message(header, requeue=False)
3737
return
3838

3939
# Create a wrapper-like object that can be passed to functions
@@ -57,14 +57,14 @@ def class3d(self, rw, header: dict, message: dict):
5757
f"and recipe parameters: {rw.recipe_step.get('parameters', {})} "
5858
f"with exception: {e}"
5959
)
60-
rw.transport.nack(header)
60+
self._reject_message(header, transport=rw.transport, requeue=False)
6161
return
6262

63-
# In this setup we cannot nack messages on failure, so instead check here
63+
# In this setup we cannot reject messages on failure, so instead check here
6464
if message.get("requeue", 0) >= 5:
65-
self.log.warning(f"Nacking requeued file {class3d_params.particles_file}")
66-
rw.transport.nack(header)
67-
return False
65+
self.log.warning(f"Rejecting requeued file {class3d_params.particles_file}")
66+
self._reject_message(header, transport=rw.transport, requeue=False)
67+
return
6868

6969
# Acknowledge the message and disconnect from rabbitmq
7070
self.log.info(

src/cryoemservices/services/clem_align_and_merge.py

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -38,7 +38,7 @@ def call_align_and_merge(self, rw, header, message):
3838
self.log.info("Received a simple message")
3939
if not isinstance(message, dict):
4040
self.log.error("Rejected invalid simple message")
41-
self._transport.nack(header)
41+
self._reject_message(header, requeue=False)
4242
return
4343

4444
# Create a wrapper-like object that can be passed to functions
@@ -61,7 +61,7 @@ def call_align_and_merge(self, rw, header, message):
6161
f"and recipe parameters: {rw.recipe_step.get('parameters', {})} "
6262
f"with exception: {e}"
6363
)
64-
rw.transport.nack(header)
64+
self._reject_message(header, transport=rw.transport, requeue=False)
6565
return
6666

6767
# Process files and collect output
@@ -75,20 +75,20 @@ def call_align_and_merge(self, rw, header, message):
7575
align_across=params.align_across,
7676
num_procs=params.num_procs,
7777
)
78-
# Log error and nack message if the command fails to execute
78+
# Log error and reject message if the command fails to execute
7979
except Exception:
8080
self.log.error(
8181
f"Exception encountered while aligning and merging images for {params.series_name!r}: \n",
8282
exc_info=True,
8383
)
84-
rw.transport.nack(header)
84+
self._reject_message(header, transport=rw.transport)
8585
return
8686
if not result:
8787
self.log.error(
8888
"Failed to complete the aligning and merging process for "
8989
f"{params.series_name!r}"
9090
)
91-
rw.transport.nack(header)
91+
self._reject_message(header, transport=rw.transport)
9292
return
9393

9494
# Request for PNG image to be created

src/cryoemservices/services/clem_process_raw_lifs.py

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -38,7 +38,7 @@ def call_process_raw_lifs(self, rw, header, message):
3838
self.log.info("Received a simple message")
3939
if not isinstance(message, dict):
4040
self.log.error("Rejected invalid simple message")
41-
self._transport.nack(header)
41+
self._reject_message(header, requeue=False)
4242
return
4343

4444
# Create a wrapper-like object that can be passed to functions
@@ -61,7 +61,7 @@ def call_process_raw_lifs(self, rw, header, message):
6161
f"and recipe parameters: {rw.recipe_step.get('parameters', {})} "
6262
f"with exception: {e}"
6363
)
64-
rw.transport.nack(header)
64+
self._reject_message(header, transport=rw.transport, requeue=False)
6565
return
6666

6767
# Process files and collect output
@@ -71,19 +71,19 @@ def call_process_raw_lifs(self, rw, header, message):
7171
root_folder=params.root_folder,
7272
number_of_processes=params.num_procs,
7373
)
74-
# Log error and nack message if the command fails to execute
74+
# Log error and reject message if the command fails to execute
7575
except Exception:
7676
self.log.error(
77-
f"Exception encontered while processing LIF file {str(params.lif_file)!r}: \n",
77+
f"Exception encountered while processing LIF file {str(params.lif_file)!r}: \n",
7878
exc_info=True,
7979
)
80-
rw.transport.nack(header)
80+
self._reject_message(header, transport=rw.transport)
8181
return
8282
if not results:
8383
self.log.error(
8484
f"Failed to extract image stacks from {str(params.lif_file)!r}"
8585
)
86-
rw.transport.nack(header)
86+
self._reject_message(header, transport=rw.transport)
8787
return
8888

8989
# Send each subset of output files to Murfey for registration

src/cryoemservices/services/clem_process_raw_tiffs.py

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -38,7 +38,7 @@ def call_process_raw_tiffs(self, rw, header, message):
3838
self.log.info("Received a simple message")
3939
if not isinstance(message, dict):
4040
self.log.error("Rejected invalid simple message")
41-
self._transport.nack(header)
41+
self._reject_message(header, requeue=False)
4242
return
4343

4444
# Create a wrapper-like object that can be passed to functions
@@ -61,7 +61,7 @@ def call_process_raw_tiffs(self, rw, header, message):
6161
f"and recipe parameters: {rw.recipe_step.get('parameters', {})} "
6262
f"with exception: {e}"
6363
)
64-
rw.transport.nack(header)
64+
self._reject_message(header, transport=rw.transport, requeue=False)
6565
return
6666

6767
# Reconstruct series name using reference file from list
@@ -74,7 +74,7 @@ def call_process_raw_tiffs(self, rw, header, message):
7474
f"Subpath {params.root_folder!r} was not found in file path "
7575
f"{str(ref_file.parent / ref_file.stem.split('--')[0])!r}"
7676
)
77-
rw.transport.nack(header)
77+
self._reject_message(header, transport=rw.transport, requeue=False)
7878
return
7979
series_name = "--".join(
8080
[p.replace(" ", "_") if " " in p else p for p in path_parts][
@@ -90,19 +90,19 @@ def call_process_raw_tiffs(self, rw, header, message):
9090
metadata_file=params.metadata,
9191
number_of_processes=params.num_procs,
9292
)
93-
# Log error and nack message if the command fails to execute
93+
# Log error and reject message if the command fails to execute
9494
except Exception:
9595
self.log.error(
9696
f"Exception encountered while processing TIFF files for series {series_name}: \n",
9797
exc_info=True,
9898
)
99-
rw.transport.nack(header)
99+
self._reject_message(header, transport=rw.transport)
100100
return
101101
if not result:
102102
self.log.error(
103103
f"No processing results were returned for TIFF series {series_name!r}"
104104
)
105-
rw.transport.nack(header)
105+
self._reject_message(header, transport=rw.transport)
106106
return
107107

108108
# Request for PNG images to be created

src/cryoemservices/services/cluster_submission.py

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -46,7 +46,7 @@ def run_submit_job(self, rw, header, message):
4646
Path(wrapper).parent.mkdir(parents=True, exist_ok=True)
4747
except OSError:
4848
self.log.error(f"Cannot make directory for {wrapper}")
49-
self._transport.nack(header)
49+
self._reject_message(header, transport=rw.transport)
5050
return
5151
self.log.info(f"Storing serialized recipe wrapper in {wrapper}")
5252
cluster_params.commands = cluster_params.commands.replace(
@@ -72,7 +72,7 @@ def run_submit_job(self, rw, header, message):
7272
self.log.error(
7373
"No absolute working directory specified. Will not run cluster job"
7474
)
75-
self._transport.nack(header)
75+
self._reject_message(header, transport=rw.transport)
7676
return
7777
working_directory = Path(parameters["workingdir"])
7878
try:
@@ -81,7 +81,7 @@ def run_submit_job(self, rw, header, message):
8181
self.log.error(
8282
"Could not create working directory: %s", str(e), exc_info=True
8383
)
84-
self._transport.nack(header)
84+
self._reject_message(header, transport=rw.transport)
8585
return
8686

8787
if parameters.get("standard_output"):
@@ -105,7 +105,7 @@ def run_submit_job(self, rw, header, message):
105105
)
106106
if not jobnumber:
107107
self.log.error("Job was not submitted")
108-
self._transport.nack(header)
108+
self._reject_message(header, transport=rw.transport)
109109
return
110110

111111
# Conditionally acknowledge receipt of the message

src/cryoemservices/services/common_service.py

Lines changed: 31 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,8 +2,10 @@
22

33
import logging
44
import queue
5+
from functools import partial
56

67
from workflows.transport.common_transport import CommonTransport
8+
from workflows.transport.pika_transport import PikaTransport
79

810

911
class CommonService:
@@ -66,3 +68,32 @@ def start(self):
6668
self._transport.disconnect()
6769
except Exception as e:
6870
self.log.error(f"Could not disconnect transport: {e}", exc_info=True)
71+
72+
def _reject_message(
73+
self,
74+
header: dict,
75+
transport: CommonTransport | None = None,
76+
requeue: bool = True,
77+
):
78+
"""Reject failed messages back to rabbitmq"""
79+
message_id = header.get("message-id")
80+
subscription_id = header.get("subscription")
81+
if transport is None:
82+
transport = self._transport
83+
if (
84+
isinstance(transport, PikaTransport)
85+
and message_id is not None
86+
and subscription_id is not None
87+
):
88+
pika_thread = transport._pika_thread
89+
channel = pika_thread._pika_channels[subscription_id]
90+
pika_thread._connection.add_callback_threadsafe(
91+
partial(channel.basic_reject, delivery_tag=message_id, requeue=requeue)
92+
)
93+
else:
94+
# Resort back to nacking if this isn't pika or the header is invalid
95+
# Mostly just for tests compatibility
96+
self.log.warning(
97+
f"Message {message_id} in {subscription_id} is not valid for rabbitmq"
98+
)
99+
transport.nack(header, requeue=requeue)

src/cryoemservices/services/cryolo.py

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -112,7 +112,7 @@ def cryolo(self, rw, header: dict, message: dict):
112112
self.log.info("Received a simple message")
113113
if not isinstance(message, dict):
114114
self.log.error("Rejected invalid simple message")
115-
self._transport.nack(header)
115+
self._reject_message(header, requeue=False)
116116
return
117117

118118
# Create a wrapper-like object that can be passed to functions
@@ -138,7 +138,7 @@ def cryolo(self, rw, header: dict, message: dict):
138138
f"and recipe parameters: {rw.recipe_step.get('parameters', {})} "
139139
f"with exception: {e}"
140140
)
141-
rw.transport.nack(header)
141+
self._reject_message(header, transport=rw.transport, requeue=False)
142142
return
143143

144144
# Check if this file has been run before
@@ -160,7 +160,7 @@ def cryolo(self, rw, header: dict, message: dict):
160160
job_number = int(job_num_search[0][4:7])
161161
else:
162162
self.log.warning(f"Invalid job directory in {cryolo_params.output_path}")
163-
rw.transport.nack(header)
163+
self._reject_message(header, transport=rw.transport, requeue=False)
164164
return
165165

166166
# Check job alias
@@ -169,7 +169,7 @@ def cryolo(self, rw, header: dict, message: dict):
169169
job_alias.symlink_to(job_dir)
170170
elif not (job_alias.is_symlink() and job_alias.resolve() == job_dir.resolve()):
171171
self.log.error(f"Symlink {job_alias} already exists")
172-
rw.transport.nack(header)
172+
self._reject_message(header, transport=rw.transport)
173173
return
174174

175175
Path(cryolo_params.output_path).unlink(missing_ok=True)
@@ -302,7 +302,7 @@ def cryolo(self, rw, header: dict, message: dict):
302302
f"crYOLO failed with exitcode {result.returncode}:\n"
303303
+ result.stderr.decode("utf8", "replace")
304304
)
305-
rw.transport.nack(header)
305+
self._reject_message(header, transport=rw.transport)
306306
return
307307

308308
# If this is tomo then make an image and stop here

src/cryoemservices/services/ctffind.py

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -117,7 +117,7 @@ def ctf_find(self, rw, header: dict, message: dict):
117117
self.log.info("Received a simple message")
118118
if not isinstance(message, dict):
119119
self.log.error("Rejected invalid simple message")
120-
self._transport.nack(header)
120+
self._reject_message(header, requeue=False)
121121
return
122122

123123
# Create a wrapper-like object that can be passed to functions
@@ -138,12 +138,12 @@ def ctf_find(self, rw, header: dict, message: dict):
138138
f"and recipe parameters: {rw.recipe_step.get('parameters', {})} "
139139
f"with exception: {e}"
140140
)
141-
rw.transport.nack(header)
141+
self._reject_message(header, transport=rw.transport, requeue=False)
142142
return
143143

144144
if ctf_params.ctffind_version not in [4, 5]:
145145
self.log.error(f"Cannot use CTFFind version {ctf_params.ctffind_version}")
146-
rw.transport.nack(header)
146+
self._reject_message(header, transport=rw.transport, requeue=False)
147147
return
148148
self.log.info(f"Using CTFFind version {ctf_params.ctffind_version}")
149149

@@ -168,7 +168,7 @@ def ctf_find(self, rw, header: dict, message: dict):
168168
self.log.warning(
169169
f"Could not determine job number in {ctf_params.output_image}"
170170
)
171-
rw.transport.nack(header)
171+
self._reject_message(header, transport=rw.transport, requeue=False)
172172
return
173173
job_alias = Path(
174174
re.sub(
@@ -185,7 +185,7 @@ def ctf_find(self, rw, header: dict, message: dict):
185185
== (job_alias.parent / f"job{ctf_job_number:03}").resolve()
186186
):
187187
self.log.error(f"Symlink {job_alias} already exists")
188-
rw.transport.nack(header)
188+
self._reject_message(header, transport=rw.transport)
189189
return
190190

191191
parameters_list = [
@@ -271,7 +271,7 @@ def ctf_find(self, rw, header: dict, message: dict):
271271
f"CTFFind failed with exitcode {result.returncode}:\n"
272272
+ result.stderr.decode("utf8", "replace")
273273
)
274-
rw.transport.nack(header)
274+
self._reject_message(header, transport=rw.transport)
275275
return
276276

277277
# Write stdout to logfile

0 commit comments

Comments
 (0)