Improve OTP shutdown behavior for consumers - #103

Open
nsweeting wants to merge 4 commits into
masterfrom
cleanup-shutdown
Open

Improve OTP shutdown behavior for consumers#103
nsweeting wants to merge 4 commits into
masterfrom
cleanup-shutdown

Conversation

@nsweeting

Copy link
Copy Markdown
Owner

Summary

Improves the consumer shutdown sequence to follow proper OTP conventions. Previously, incorrect child spec types, missing consumption cancellation, and lack of graceful drain logic meant that a SIGTERM (e.g. k8s pod termination) would result in abrupt process kills rather than an orderly shutdown.

Problems Fixed

  • Consumer.Supervisor was typed as :worker with a 5s shutdown — the entire consumer subtree had only 5 seconds before being :killed. Now correctly typed as :supervisor with :infinity shutdown, allowing children to drain properly.
  • No consumption cancellation during shutdown — the broker continued delivering new messages while workers were shutting down. Now calls AMQP.Basic.cancel first.
  • Workers stopped sequentially — on an N-core machine with N workers, sequential DynamicSupervisor.stop calls could easily exceed the shutdown budget. Now stops all workers in parallel.
  • Executer had no terminate/2 — when shut down externally, the spawned message-processing process was killed with no opportunity to finish. Now uses Task.async + Task.shutdown/2 to give in-flight work a 5s grace period before escalating, and nacks with requeue: true as a safety net.
  • No explicit shutdown timeouts — all processes used OTP defaults (5s). Now uses a proper timeout chain: Consumer.Supervisor (:infinity) > Consumer.Server (30s) > Workers/Executers (25s).

Corrected Shutdown Sequence

SIGTERM
-> Consumer.Supervisor gets :shutdown (timeout: :infinity)
-> Each Consumer.Server gets :shutdown (timeout: 30s)
-> terminate/2:
1. AMQP.Basic.cancel — stop new message delivery
2. Stop all Workers in parallel (25s timeout)
-> Each Executer.terminate/2:
a. Task.shutdown(task, 5_000) — grace period for in-flight work
b. Safety-net nack with requeue: true if not completed
3. AMQP.Channel.close — clean channel shutdown
-> Producer.Pool, Topology.Server, Connection.Pool shut down after

Changes

  • lib/rabbit/broker/supervisor.ex — Consumer.Supervisor child spec gets type: :supervisor
  • lib/rabbit/consumer/supervisor.ex — Consumer.Server child specs get shutdown: 30_000
  • lib/rabbit/consumer/server.ex — Add cancel_consumer/1, parallel stop_workers/1
  • lib/rabbit/consumer/executer.ex — Refactor to Task.async, add terminate/2 with grace period and safety-net nack, shutdown: 25_000, completed state tracking
  • mix.exs — Version bump to 0.22.0

- Set Consumer.Supervisor child spec type to :supervisor so it gets
shutdown: :infinity instead of being killed after 5 seconds
- Set explicit shutdown: 30_000 on each Consumer.Server child spec
- Cancel AMQP consumption (Basic.cancel) before stopping workers so
no new messages arrive during drain
- Stop workers in parallel instead of sequentially to fit within
the shutdown budget
- Add terminate/2 to Executer that nacks unfinished messages with
requeue: true, preventing unacked messages from accumulating on
quorum queues
- Set Executer child_spec shutdown: 25_000 to give in-flight messages
time to complete before the safety-net nack
- Bump version to 0.22.0
Replace spawn_link with Task.async so terminate/2 can use
Task.shutdown/2 to give in-flight messages a 5s grace period
to complete before escalating to :kill and safety-net nacking.
Previously the spawned process was killed immediately on shutdown
with no chance to finish. Now the sequence is:
1. Task.shutdown(task, 5_000) - sends :shutdown, waits 5s
2. If task doesn't finish, brutally kills it
3. Safety-net nack with requeue: true
@nsweetingnsweeting changed the title fix: improve OTP shutdown behavior for consumersImprove OTP shutdown behavior for consumersMar 25, 2026
Check Task.shutdown/2 return value in terminate/2. If it returns
{:ok, _}, the task finished within the grace period and already
acked/nacked inside its body - skip the redundant nack.
- Executer lifecycle: normal completion, crash error handling, timeout
- Executer terminate/2: shutdown with in-flight task, task-never-started
edge case, skip nack when task completes before shutdown
- Child spec assertions: Executer shutdown/restart, Consumer.Server
shutdown timeout of 30_000
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant

@nsweeting
, 'i'); if (__m === '*' || __re.test(location.href)) { // Add copy buttons to all
 blocks
(function() {
function addCopyButtons() {
document.querySelectorAll('pre code').forEach(function(codeBlock) {
if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;
codeBlock.parentElement.setAttribute('data-copy-added', 'true');
var btn = document.createElement('button');
btn.textContent = 'Copy';
btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';
btn.onmouseover = function() { this.style.opacity = '1'; };
btn.onmouseout = function() { this.style.opacity = '0.7'; };
btn.onclick = function() {
navigator.clipboard.writeText(codeBlock.textContent).then(function() {
btn.textContent = 'Copied!';
setTimeout(function() { btn.textContent = 'Copy'; }, 1500);
});
};
codeBlock.parentElement.style.position = 'relative';
codeBlock.parentElement.appendChild(btn);
});
}
addCopyButtons();
// Re-run on dynamic content
var observer = new MutationObserver(addCopyButtons);
observer.observe(document.body, { childList: true, subtree: true });
})();
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Skip to content

Improve OTP shutdown behavior for consumers - #103

Open
nsweeting wants to merge 4 commits into
masterfrom
cleanup-shutdown
Open

Improve OTP shutdown behavior for consumers#103
nsweeting wants to merge 4 commits into
masterfrom
cleanup-shutdown

Conversation

@nsweeting

Copy link
Copy Markdown
Owner

Summary

Improves the consumer shutdown sequence to follow proper OTP conventions. Previously, incorrect child spec types, missing consumption cancellation, and lack of graceful drain logic meant that a SIGTERM (e.g. k8s pod termination) would result in abrupt process kills rather than an orderly shutdown.

Problems Fixed

  • Consumer.Supervisor was typed as :worker with a 5s shutdown — the entire consumer subtree had only 5 seconds before being :killed. Now correctly typed as :supervisor with :infinity shutdown, allowing children to drain properly.
  • No consumption cancellation during shutdown — the broker continued delivering new messages while workers were shutting down. Now calls AMQP.Basic.cancel first.
  • Workers stopped sequentially — on an N-core machine with N workers, sequential DynamicSupervisor.stop calls could easily exceed the shutdown budget. Now stops all workers in parallel.
  • Executer had no terminate/2 — when shut down externally, the spawned message-processing process was killed with no opportunity to finish. Now uses Task.async + Task.shutdown/2 to give in-flight work a 5s grace period before escalating, and nacks with requeue: true as a safety net.
  • No explicit shutdown timeouts — all processes used OTP defaults (5s). Now uses a proper timeout chain: Consumer.Supervisor (:infinity) > Consumer.Server (30s) > Workers/Executers (25s).

Corrected Shutdown Sequence

SIGTERM
-> Consumer.Supervisor gets :shutdown (timeout: :infinity)
-> Each Consumer.Server gets :shutdown (timeout: 30s)
-> terminate/2:
1. AMQP.Basic.cancel — stop new message delivery
2. Stop all Workers in parallel (25s timeout)
-> Each Executer.terminate/2:
a. Task.shutdown(task, 5_000) — grace period for in-flight work
b. Safety-net nack with requeue: true if not completed
3. AMQP.Channel.close — clean channel shutdown
-> Producer.Pool, Topology.Server, Connection.Pool shut down after

Changes

  • lib/rabbit/broker/supervisor.ex — Consumer.Supervisor child spec gets type: :supervisor
  • lib/rabbit/consumer/supervisor.ex — Consumer.Server child specs get shutdown: 30_000
  • lib/rabbit/consumer/server.ex — Add cancel_consumer/1, parallel stop_workers/1
  • lib/rabbit/consumer/executer.ex — Refactor to Task.async, add terminate/2 with grace period and safety-net nack, shutdown: 25_000, completed state tracking
  • mix.exs — Version bump to 0.22.0

- Set Consumer.Supervisor child spec type to :supervisor so it gets
shutdown: :infinity instead of being killed after 5 seconds
- Set explicit shutdown: 30_000 on each Consumer.Server child spec
- Cancel AMQP consumption (Basic.cancel) before stopping workers so
no new messages arrive during drain
- Stop workers in parallel instead of sequentially to fit within
the shutdown budget
- Add terminate/2 to Executer that nacks unfinished messages with
requeue: true, preventing unacked messages from accumulating on
quorum queues
- Set Executer child_spec shutdown: 25_000 to give in-flight messages
time to complete before the safety-net nack
- Bump version to 0.22.0
Replace spawn_link with Task.async so terminate/2 can use
Task.shutdown/2 to give in-flight messages a 5s grace period
to complete before escalating to :kill and safety-net nacking.
Previously the spawned process was killed immediately on shutdown
with no chance to finish. Now the sequence is:
1. Task.shutdown(task, 5_000) - sends :shutdown, waits 5s
2. If task doesn't finish, brutally kills it
3. Safety-net nack with requeue: true
@nsweetingnsweeting changed the title fix: improve OTP shutdown behavior for consumersImprove OTP shutdown behavior for consumersMar 25, 2026
Check Task.shutdown/2 return value in terminate/2. If it returns
{:ok, _}, the task finished within the grace period and already
acked/nacked inside its body - skip the redundant nack.
- Executer lifecycle: normal completion, crash error handling, timeout
- Executer terminate/2: shutdown with in-flight task, task-never-started
edge case, skip nack when task completes before shutdown
- Child spec assertions: Executer shutdown/restart, Consumer.Server
shutdown timeout of 30_000
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant

@nsweeting
, 'i'); if (__m === '*' || __re.test(location.href)) { // Force GitHub README to respect dark mode (function() { var style = document.createElement('style'); style.textContent = ' .markdown-body { color-scheme: dark light; } .markdown-body pre { background: #161b22 !important; } .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; } .markdown-body table th, .markdown-body table td { border-color: #30363d !important; } .markdown-body img { background: #0d1117; } .markdown-body blockquote { border-left-color: #8b949e; } .markdown-body hr { border-color: #30363d; } '; document.head.appendChild(style); })(); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Improve OTP shutdown behavior for consumers - #103

Open
nsweeting wants to merge 4 commits into
masterfrom
cleanup-shutdown
Open

Improve OTP shutdown behavior for consumers#103
nsweeting wants to merge 4 commits into
masterfrom
cleanup-shutdown

Conversation

@nsweeting

Copy link
Copy Markdown
Owner

Summary

Improves the consumer shutdown sequence to follow proper OTP conventions. Previously, incorrect child spec types, missing consumption cancellation, and lack of graceful drain logic meant that a SIGTERM (e.g. k8s pod termination) would result in abrupt process kills rather than an orderly shutdown.

Problems Fixed

  • Consumer.Supervisor was typed as :worker with a 5s shutdown — the entire consumer subtree had only 5 seconds before being :killed. Now correctly typed as :supervisor with :infinity shutdown, allowing children to drain properly.
  • No consumption cancellation during shutdown — the broker continued delivering new messages while workers were shutting down. Now calls AMQP.Basic.cancel first.
  • Workers stopped sequentially — on an N-core machine with N workers, sequential DynamicSupervisor.stop calls could easily exceed the shutdown budget. Now stops all workers in parallel.
  • Executer had no terminate/2 — when shut down externally, the spawned message-processing process was killed with no opportunity to finish. Now uses Task.async + Task.shutdown/2 to give in-flight work a 5s grace period before escalating, and nacks with requeue: true as a safety net.
  • No explicit shutdown timeouts — all processes used OTP defaults (5s). Now uses a proper timeout chain: Consumer.Supervisor (:infinity) > Consumer.Server (30s) > Workers/Executers (25s).

Corrected Shutdown Sequence

SIGTERM
-> Consumer.Supervisor gets :shutdown (timeout: :infinity)
-> Each Consumer.Server gets :shutdown (timeout: 30s)
-> terminate/2:
1. AMQP.Basic.cancel — stop new message delivery
2. Stop all Workers in parallel (25s timeout)
-> Each Executer.terminate/2:
a. Task.shutdown(task, 5_000) — grace period for in-flight work
b. Safety-net nack with requeue: true if not completed
3. AMQP.Channel.close — clean channel shutdown
-> Producer.Pool, Topology.Server, Connection.Pool shut down after

Changes

  • lib/rabbit/broker/supervisor.ex — Consumer.Supervisor child spec gets type: :supervisor
  • lib/rabbit/consumer/supervisor.ex — Consumer.Server child specs get shutdown: 30_000
  • lib/rabbit/consumer/server.ex — Add cancel_consumer/1, parallel stop_workers/1
  • lib/rabbit/consumer/executer.ex — Refactor to Task.async, add terminate/2 with grace period and safety-net nack, shutdown: 25_000, completed state tracking
  • mix.exs — Version bump to 0.22.0

- Set Consumer.Supervisor child spec type to :supervisor so it gets
shutdown: :infinity instead of being killed after 5 seconds
- Set explicit shutdown: 30_000 on each Consumer.Server child spec
- Cancel AMQP consumption (Basic.cancel) before stopping workers so
no new messages arrive during drain
- Stop workers in parallel instead of sequentially to fit within
the shutdown budget
- Add terminate/2 to Executer that nacks unfinished messages with
requeue: true, preventing unacked messages from accumulating on
quorum queues
- Set Executer child_spec shutdown: 25_000 to give in-flight messages
time to complete before the safety-net nack
- Bump version to 0.22.0
Replace spawn_link with Task.async so terminate/2 can use
Task.shutdown/2 to give in-flight messages a 5s grace period
to complete before escalating to :kill and safety-net nacking.
Previously the spawned process was killed immediately on shutdown
with no chance to finish. Now the sequence is:
1. Task.shutdown(task, 5_000) - sends :shutdown, waits 5s
2. If task doesn't finish, brutally kills it
3. Safety-net nack with requeue: true
@nsweetingnsweeting changed the title fix: improve OTP shutdown behavior for consumersImprove OTP shutdown behavior for consumersMar 25, 2026
Check Task.shutdown/2 return value in terminate/2. If it returns
{:ok, _}, the task finished within the grace period and already
acked/nacked inside its body - skip the redundant nack.
- Executer lifecycle: normal completion, crash error handling, timeout
- Executer terminate/2: shutdown with in-flight task, task-never-started
edge case, skip nack when task completes before shutdown
- Child spec assertions: Executer shutdown/restart, Consumer.Server
shutdown timeout of 30_000
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant

@nsweeting
, 'i'); if (__m === '*' || __re.test(location.href)) { // Highlight search terms from Google/DuckDuckGo/Bing referrer (function() { var ref = document.referrer; var terms = []; if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) { var url = new URL(ref); var q = url.searchParams.get('q') || url.searchParams.get('p'); if (q) { terms = q.split(/\s+/).filter(function(t) { return t.length > 2; }); } } if (terms.length === 0) return; var style = document.createElement('style'); style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }'; document.head.appendChild(style); function highlight(node) { if (node.nodeType === 3) { // text node var text = node.textContent; var found = false; terms.forEach(function(term) { var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\]\\]/g, '\\') + ')', 'gi'); if (regex.test(text)) { found = true; var frag = document.createDocumentFragment(); var parts = text.split(regex); parts.forEach(function(part, i) { if (i % 2 === 0) { frag.appendChild(document.createTextNode(part)); } else { var span = document.createElement('span'); span.className = 'userscript-highlight'; span.textContent = part; frag.appendChild(span); } }); node.parentNode.replaceChild(frag, node); } }); } else if (node.nodeType === 1 && node.childNodes) { // element var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT']; if (!skipTags.includes(node.tagName)) { Array.from(node.childNodes).forEach(highlight); } } } highlight(document.body); // Re-highlight on dynamic content var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1 || node.nodeType === 3) highlight(node); }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Improve OTP shutdown behavior for consumers - #103

Open
nsweeting wants to merge 4 commits into
masterfrom
cleanup-shutdown
Open

Improve OTP shutdown behavior for consumers#103
nsweeting wants to merge 4 commits into
masterfrom
cleanup-shutdown

Conversation

@nsweeting

Copy link
Copy Markdown
Owner

Summary

Improves the consumer shutdown sequence to follow proper OTP conventions. Previously, incorrect child spec types, missing consumption cancellation, and lack of graceful drain logic meant that a SIGTERM (e.g. k8s pod termination) would result in abrupt process kills rather than an orderly shutdown.

Problems Fixed

  • Consumer.Supervisor was typed as :worker with a 5s shutdown — the entire consumer subtree had only 5 seconds before being :killed. Now correctly typed as :supervisor with :infinity shutdown, allowing children to drain properly.
  • No consumption cancellation during shutdown — the broker continued delivering new messages while workers were shutting down. Now calls AMQP.Basic.cancel first.
  • Workers stopped sequentially — on an N-core machine with N workers, sequential DynamicSupervisor.stop calls could easily exceed the shutdown budget. Now stops all workers in parallel.
  • Executer had no terminate/2 — when shut down externally, the spawned message-processing process was killed with no opportunity to finish. Now uses Task.async + Task.shutdown/2 to give in-flight work a 5s grace period before escalating, and nacks with requeue: true as a safety net.
  • No explicit shutdown timeouts — all processes used OTP defaults (5s). Now uses a proper timeout chain: Consumer.Supervisor (:infinity) > Consumer.Server (30s) > Workers/Executers (25s).

Corrected Shutdown Sequence

SIGTERM
-> Consumer.Supervisor gets :shutdown (timeout: :infinity)
-> Each Consumer.Server gets :shutdown (timeout: 30s)
-> terminate/2:
1. AMQP.Basic.cancel — stop new message delivery
2. Stop all Workers in parallel (25s timeout)
-> Each Executer.terminate/2:
a. Task.shutdown(task, 5_000) — grace period for in-flight work
b. Safety-net nack with requeue: true if not completed
3. AMQP.Channel.close — clean channel shutdown
-> Producer.Pool, Topology.Server, Connection.Pool shut down after

Changes

  • lib/rabbit/broker/supervisor.ex — Consumer.Supervisor child spec gets type: :supervisor
  • lib/rabbit/consumer/supervisor.ex — Consumer.Server child specs get shutdown: 30_000
  • lib/rabbit/consumer/server.ex — Add cancel_consumer/1, parallel stop_workers/1
  • lib/rabbit/consumer/executer.ex — Refactor to Task.async, add terminate/2 with grace period and safety-net nack, shutdown: 25_000, completed state tracking
  • mix.exs — Version bump to 0.22.0

- Set Consumer.Supervisor child spec type to :supervisor so it gets
shutdown: :infinity instead of being killed after 5 seconds
- Set explicit shutdown: 30_000 on each Consumer.Server child spec
- Cancel AMQP consumption (Basic.cancel) before stopping workers so
no new messages arrive during drain
- Stop workers in parallel instead of sequentially to fit within
the shutdown budget
- Add terminate/2 to Executer that nacks unfinished messages with
requeue: true, preventing unacked messages from accumulating on
quorum queues
- Set Executer child_spec shutdown: 25_000 to give in-flight messages
time to complete before the safety-net nack
- Bump version to 0.22.0
Replace spawn_link with Task.async so terminate/2 can use
Task.shutdown/2 to give in-flight messages a 5s grace period
to complete before escalating to :kill and safety-net nacking.
Previously the spawned process was killed immediately on shutdown
with no chance to finish. Now the sequence is:
1. Task.shutdown(task, 5_000) - sends :shutdown, waits 5s
2. If task doesn't finish, brutally kills it
3. Safety-net nack with requeue: true
@nsweetingnsweeting changed the title fix: improve OTP shutdown behavior for consumersImprove OTP shutdown behavior for consumersMar 25, 2026
Check Task.shutdown/2 return value in terminate/2. If it returns
{:ok, _}, the task finished within the grace period and already
acked/nacked inside its body - skip the redundant nack.
- Executer lifecycle: normal completion, crash error handling, timeout
- Executer terminate/2: shutdown with in-flight task, task-never-started
edge case, skip nack when task completes before shutdown
- Child spec assertions: Executer shutdown/restart, Consumer.Server
shutdown timeout of 30_000
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant

@nsweeting
, 'i'); if (__m === '*' || __re.test(location.href)) { // Strip utm_, fbclid, gclid, etc. from all links on page (function() { var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content', 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid', 'ref', 'ref_src', 'source', 'medium', 'campaign']; function cleanUrl(url) { try { var u = new URL(url, window.location.origin); var changed = false; trackingParams.forEach(function(p) { if (u.searchParams.has(p)) { u.searchParams.delete(p); changed = true; } }); return changed ? u.toString() : url; } catch (e) { return url; } } function cleanLinks() { document.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } cleanLinks(); var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1) { if (node.tagName === 'A') cleanLinks(); node.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content

Improve OTP shutdown behavior for consumers - #103

Open
nsweeting wants to merge 4 commits into
masterfrom
cleanup-shutdown
Open

Improve OTP shutdown behavior for consumers#103
nsweeting wants to merge 4 commits into
masterfrom
cleanup-shutdown

Conversation

@nsweeting

Copy link
Copy Markdown
Owner

Summary

Improves the consumer shutdown sequence to follow proper OTP conventions. Previously, incorrect child spec types, missing consumption cancellation, and lack of graceful drain logic meant that a SIGTERM (e.g. k8s pod termination) would result in abrupt process kills rather than an orderly shutdown.

Problems Fixed

  • Consumer.Supervisor was typed as :worker with a 5s shutdown — the entire consumer subtree had only 5 seconds before being :killed. Now correctly typed as :supervisor with :infinity shutdown, allowing children to drain properly.
  • No consumption cancellation during shutdown — the broker continued delivering new messages while workers were shutting down. Now calls AMQP.Basic.cancel first.
  • Workers stopped sequentially — on an N-core machine with N workers, sequential DynamicSupervisor.stop calls could easily exceed the shutdown budget. Now stops all workers in parallel.
  • Executer had no terminate/2 — when shut down externally, the spawned message-processing process was killed with no opportunity to finish. Now uses Task.async + Task.shutdown/2 to give in-flight work a 5s grace period before escalating, and nacks with requeue: true as a safety net.
  • No explicit shutdown timeouts — all processes used OTP defaults (5s). Now uses a proper timeout chain: Consumer.Supervisor (:infinity) > Consumer.Server (30s) > Workers/Executers (25s).

Corrected Shutdown Sequence

SIGTERM
-> Consumer.Supervisor gets :shutdown (timeout: :infinity)
-> Each Consumer.Server gets :shutdown (timeout: 30s)
-> terminate/2:
1. AMQP.Basic.cancel — stop new message delivery
2. Stop all Workers in parallel (25s timeout)
-> Each Executer.terminate/2:
a. Task.shutdown(task, 5_000) — grace period for in-flight work
b. Safety-net nack with requeue: true if not completed
3. AMQP.Channel.close — clean channel shutdown
-> Producer.Pool, Topology.Server, Connection.Pool shut down after

Changes

  • lib/rabbit/broker/supervisor.ex — Consumer.Supervisor child spec gets type: :supervisor
  • lib/rabbit/consumer/supervisor.ex — Consumer.Server child specs get shutdown: 30_000
  • lib/rabbit/consumer/server.ex — Add cancel_consumer/1, parallel stop_workers/1
  • lib/rabbit/consumer/executer.ex — Refactor to Task.async, add terminate/2 with grace period and safety-net nack, shutdown: 25_000, completed state tracking
  • mix.exs — Version bump to 0.22.0

- Set Consumer.Supervisor child spec type to :supervisor so it gets
shutdown: :infinity instead of being killed after 5 seconds
- Set explicit shutdown: 30_000 on each Consumer.Server child spec
- Cancel AMQP consumption (Basic.cancel) before stopping workers so
no new messages arrive during drain
- Stop workers in parallel instead of sequentially to fit within
the shutdown budget
- Add terminate/2 to Executer that nacks unfinished messages with
requeue: true, preventing unacked messages from accumulating on
quorum queues
- Set Executer child_spec shutdown: 25_000 to give in-flight messages
time to complete before the safety-net nack
- Bump version to 0.22.0
Replace spawn_link with Task.async so terminate/2 can use
Task.shutdown/2 to give in-flight messages a 5s grace period
to complete before escalating to :kill and safety-net nacking.
Previously the spawned process was killed immediately on shutdown
with no chance to finish. Now the sequence is:
1. Task.shutdown(task, 5_000) - sends :shutdown, waits 5s
2. If task doesn't finish, brutally kills it
3. Safety-net nack with requeue: true
@nsweetingnsweeting changed the title fix: improve OTP shutdown behavior for consumersImprove OTP shutdown behavior for consumersMar 25, 2026
Check Task.shutdown/2 return value in terminate/2. If it returns
{:ok, _}, the task finished within the grace period and already
acked/nacked inside its body - skip the redundant nack.
- Executer lifecycle: normal completion, crash error handling, timeout
- Executer terminate/2: shutdown with in-flight task, task-never-started
edge case, skip nack when task completes before shutdown
- Child spec assertions: Executer shutdown/restart, Consumer.Server
shutdown timeout of 30_000
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant

@nsweeting
, 'i'); if (__m === '*' || __re.test(location.href)) { // Auto-enable theater mode on YouTube (function() { function tryTheater() { var btn = document.querySelector('button[aria-label="Theater mode"], ytd-player #player button[title="Theater mode"]'); if (btn && !btn.classList.contains('activated')) { btn.click(); } } // Try immediately tryTheater(); // Try after navigation (SPA) var lastUrl = location.href; setInterval(function() { if (location.href !== lastUrl) { lastUrl = location.href; setTimeout(tryTheater, 500); } }, 1000); // Also try on player load var observer = new MutationObserver(tryTheater); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Improve OTP shutdown behavior for consumers - #103

Open
nsweeting wants to merge 4 commits into
masterfrom
cleanup-shutdown
Open

Improve OTP shutdown behavior for consumers#103
nsweeting wants to merge 4 commits into
masterfrom
cleanup-shutdown

Conversation

@nsweeting

Copy link
Copy Markdown
Owner

Summary

Improves the consumer shutdown sequence to follow proper OTP conventions. Previously, incorrect child spec types, missing consumption cancellation, and lack of graceful drain logic meant that a SIGTERM (e.g. k8s pod termination) would result in abrupt process kills rather than an orderly shutdown.

Problems Fixed

  • Consumer.Supervisor was typed as :worker with a 5s shutdown — the entire consumer subtree had only 5 seconds before being :killed. Now correctly typed as :supervisor with :infinity shutdown, allowing children to drain properly.
  • No consumption cancellation during shutdown — the broker continued delivering new messages while workers were shutting down. Now calls AMQP.Basic.cancel first.
  • Workers stopped sequentially — on an N-core machine with N workers, sequential DynamicSupervisor.stop calls could easily exceed the shutdown budget. Now stops all workers in parallel.
  • Executer had no terminate/2 — when shut down externally, the spawned message-processing process was killed with no opportunity to finish. Now uses Task.async + Task.shutdown/2 to give in-flight work a 5s grace period before escalating, and nacks with requeue: true as a safety net.
  • No explicit shutdown timeouts — all processes used OTP defaults (5s). Now uses a proper timeout chain: Consumer.Supervisor (:infinity) > Consumer.Server (30s) > Workers/Executers (25s).

Corrected Shutdown Sequence

SIGTERM
-> Consumer.Supervisor gets :shutdown (timeout: :infinity)
-> Each Consumer.Server gets :shutdown (timeout: 30s)
-> terminate/2:
1. AMQP.Basic.cancel — stop new message delivery
2. Stop all Workers in parallel (25s timeout)
-> Each Executer.terminate/2:
a. Task.shutdown(task, 5_000) — grace period for in-flight work
b. Safety-net nack with requeue: true if not completed
3. AMQP.Channel.close — clean channel shutdown
-> Producer.Pool, Topology.Server, Connection.Pool shut down after

Changes

  • lib/rabbit/broker/supervisor.ex — Consumer.Supervisor child spec gets type: :supervisor
  • lib/rabbit/consumer/supervisor.ex — Consumer.Server child specs get shutdown: 30_000
  • lib/rabbit/consumer/server.ex — Add cancel_consumer/1, parallel stop_workers/1
  • lib/rabbit/consumer/executer.ex — Refactor to Task.async, add terminate/2 with grace period and safety-net nack, shutdown: 25_000, completed state tracking
  • mix.exs — Version bump to 0.22.0

- Set Consumer.Supervisor child spec type to :supervisor so it gets
shutdown: :infinity instead of being killed after 5 seconds
- Set explicit shutdown: 30_000 on each Consumer.Server child spec
- Cancel AMQP consumption (Basic.cancel) before stopping workers so
no new messages arrive during drain
- Stop workers in parallel instead of sequentially to fit within
the shutdown budget
- Add terminate/2 to Executer that nacks unfinished messages with
requeue: true, preventing unacked messages from accumulating on
quorum queues
- Set Executer child_spec shutdown: 25_000 to give in-flight messages
time to complete before the safety-net nack
- Bump version to 0.22.0
Replace spawn_link with Task.async so terminate/2 can use
Task.shutdown/2 to give in-flight messages a 5s grace period
to complete before escalating to :kill and safety-net nacking.
Previously the spawned process was killed immediately on shutdown
with no chance to finish. Now the sequence is:
1. Task.shutdown(task, 5_000) - sends :shutdown, waits 5s
2. If task doesn't finish, brutally kills it
3. Safety-net nack with requeue: true
@nsweetingnsweeting changed the title fix: improve OTP shutdown behavior for consumersImprove OTP shutdown behavior for consumersMar 25, 2026
Check Task.shutdown/2 return value in terminate/2. If it returns
{:ok, _}, the task finished within the grace period and already
acked/nacked inside its body - skip the redundant nack.
- Executer lifecycle: normal completion, crash error handling, timeout
- Executer terminate/2: shutdown with in-flight task, task-never-started
edge case, skip nack when task completes before shutdown
- Child spec assertions: Executer shutdown/restart, Consumer.Server
shutdown timeout of 30_000
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant

@nsweeting
, 'i'); if (__m === '*' || __re.test(location.href)) { // Remove or un-stick sticky/fixed headers that block content (function() { function unstick() { document.querySelectorAll('header, nav, [role="banner"], .header, .navbar, .sticky, .fixed-top, [style*="position: fixed"], [style*="position:sticky"]').forEach(function(el) { if (el.style.position === 'fixed' || el.style.position === 'sticky' || getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') { el.style.position = 'static'; el.style.top = 'auto'; el.style.zIndex = 'auto'; } }); } unstick(); var observer = new MutationObserver(unstick); observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] }); })(); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Improve OTP shutdown behavior for consumers - #103

Open
nsweeting wants to merge 4 commits into
masterfrom
cleanup-shutdown
Open

Improve OTP shutdown behavior for consumers#103
nsweeting wants to merge 4 commits into
masterfrom
cleanup-shutdown

Conversation

@nsweeting

Copy link
Copy Markdown
Owner

Summary

Improves the consumer shutdown sequence to follow proper OTP conventions. Previously, incorrect child spec types, missing consumption cancellation, and lack of graceful drain logic meant that a SIGTERM (e.g. k8s pod termination) would result in abrupt process kills rather than an orderly shutdown.

Problems Fixed

  • Consumer.Supervisor was typed as :worker with a 5s shutdown — the entire consumer subtree had only 5 seconds before being :killed. Now correctly typed as :supervisor with :infinity shutdown, allowing children to drain properly.
  • No consumption cancellation during shutdown — the broker continued delivering new messages while workers were shutting down. Now calls AMQP.Basic.cancel first.
  • Workers stopped sequentially — on an N-core machine with N workers, sequential DynamicSupervisor.stop calls could easily exceed the shutdown budget. Now stops all workers in parallel.
  • Executer had no terminate/2 — when shut down externally, the spawned message-processing process was killed with no opportunity to finish. Now uses Task.async + Task.shutdown/2 to give in-flight work a 5s grace period before escalating, and nacks with requeue: true as a safety net.
  • No explicit shutdown timeouts — all processes used OTP defaults (5s). Now uses a proper timeout chain: Consumer.Supervisor (:infinity) > Consumer.Server (30s) > Workers/Executers (25s).

Corrected Shutdown Sequence

SIGTERM
-> Consumer.Supervisor gets :shutdown (timeout: :infinity)
-> Each Consumer.Server gets :shutdown (timeout: 30s)
-> terminate/2:
1. AMQP.Basic.cancel — stop new message delivery
2. Stop all Workers in parallel (25s timeout)
-> Each Executer.terminate/2:
a. Task.shutdown(task, 5_000) — grace period for in-flight work
b. Safety-net nack with requeue: true if not completed
3. AMQP.Channel.close — clean channel shutdown
-> Producer.Pool, Topology.Server, Connection.Pool shut down after

Changes

  • lib/rabbit/broker/supervisor.ex — Consumer.Supervisor child spec gets type: :supervisor
  • lib/rabbit/consumer/supervisor.ex — Consumer.Server child specs get shutdown: 30_000
  • lib/rabbit/consumer/server.ex — Add cancel_consumer/1, parallel stop_workers/1
  • lib/rabbit/consumer/executer.ex — Refactor to Task.async, add terminate/2 with grace period and safety-net nack, shutdown: 25_000, completed state tracking
  • mix.exs — Version bump to 0.22.0

- Set Consumer.Supervisor child spec type to :supervisor so it gets
shutdown: :infinity instead of being killed after 5 seconds
- Set explicit shutdown: 30_000 on each Consumer.Server child spec
- Cancel AMQP consumption (Basic.cancel) before stopping workers so
no new messages arrive during drain
- Stop workers in parallel instead of sequentially to fit within
the shutdown budget
- Add terminate/2 to Executer that nacks unfinished messages with
requeue: true, preventing unacked messages from accumulating on
quorum queues
- Set Executer child_spec shutdown: 25_000 to give in-flight messages
time to complete before the safety-net nack
- Bump version to 0.22.0
Replace spawn_link with Task.async so terminate/2 can use
Task.shutdown/2 to give in-flight messages a 5s grace period
to complete before escalating to :kill and safety-net nacking.
Previously the spawned process was killed immediately on shutdown
with no chance to finish. Now the sequence is:
1. Task.shutdown(task, 5_000) - sends :shutdown, waits 5s
2. If task doesn't finish, brutally kills it
3. Safety-net nack with requeue: true
@nsweetingnsweeting changed the title fix: improve OTP shutdown behavior for consumersImprove OTP shutdown behavior for consumersMar 25, 2026
Check Task.shutdown/2 return value in terminate/2. If it returns
{:ok, _}, the task finished within the grace period and already
acked/nacked inside its body - skip the redundant nack.
- Executer lifecycle: normal completion, crash error handling, timeout
- Executer terminate/2: shutdown with in-flight task, task-never-started
edge case, skip nack when task completes before shutdown
- Child spec assertions: Executer shutdown/restart, Consumer.Server
shutdown timeout of 30_000
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant

@nsweeting
, 'i'); if (__m === '*' || __re.test(location.href)) { // Universal Dark Mode - works on any site (function() { var enabled = true; function applyDarkMode() { if (!enabled) return; // Create style element if it doesn't exist var style = document.getElementById('universal-dark-mode-style'); if (!style) { style = document.createElement('style'); style.id = 'universal-dark-mode-style'; document.head.appendChild(style); } // Dark mode CSS - inverts colors but preserves images/video style.textContent = ' /* Invert everything except media */ html { filter: invert(1) hue-rotate(180deg) !important; background: #1a1a2e !important; } /* Restore images, videos, iframes, canvas */ img, video, iframe, canvas, svg, picture, [style*="background-image"] { filter: invert(1) hue-rotate(180deg) !important; } /* Preserve specific elements that should not be inverted */ .no-dark-mode, .no-dark-mode *, [data-theme="light"], [data-theme="light"], .ace_editor, .ace_editor *, .CodeMirror, .CodeMirror *, .monaco-editor, .monaco-editor *, .markdown-body pre, .markdown-body pre *, .highlight, .highlight *, pre code, pre code * { filter: none !important; } /* Fix common UI elements */ .modal, .popup, .dropdown-menu, .tooltip, .popover { filter: invert(1) hue-rotate(180deg) !important; background: #2d2d44 !important; border-color: #444 !important; } /* Scrollbars */ ::-webkit-scrollbar { background: #1a1a2e !important; } ::-webkit-scrollbar-thumb { background: #444 !important; } ::-webkit-scrollbar-thumb:hover { background: #555 !important; } /* Selection */ ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; } ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; } '; } function removeDarkMode() { var style = document.getElementById('universal-dark-mode-style'); if (style) style.remove(); } // Toggle with Alt+Shift+D document.addEventListener('keydown', function(e) { if (e.altKey && e.shiftKey && e.key === 'D') { e.preventDefault(); enabled = !enabled; if (enabled) { applyDarkMode(); console.log('[Universal Dark Mode] Enabled'); } else { removeDarkMode(); console.log('[Universal Dark Mode] Disabled'); } } }); // Apply on load applyDarkMode(); // Re-apply on dynamic content var observer = new MutationObserver(function(mutations) { if (enabled && !document.getElementById('universal-dark-mode-style')) { applyDarkMode(); } }); observer.observe(document.head, { childList: true }); console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle'); })(); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content

Improve OTP shutdown behavior for consumers - #103

Open
nsweeting wants to merge 4 commits into
masterfrom
cleanup-shutdown
Open

Improve OTP shutdown behavior for consumers#103
nsweeting wants to merge 4 commits into
masterfrom
cleanup-shutdown

Conversation

@nsweeting

Copy link
Copy Markdown
Owner

Summary

Improves the consumer shutdown sequence to follow proper OTP conventions. Previously, incorrect child spec types, missing consumption cancellation, and lack of graceful drain logic meant that a SIGTERM (e.g. k8s pod termination) would result in abrupt process kills rather than an orderly shutdown.

Problems Fixed

  • Consumer.Supervisor was typed as :worker with a 5s shutdown — the entire consumer subtree had only 5 seconds before being :killed. Now correctly typed as :supervisor with :infinity shutdown, allowing children to drain properly.
  • No consumption cancellation during shutdown — the broker continued delivering new messages while workers were shutting down. Now calls AMQP.Basic.cancel first.
  • Workers stopped sequentially — on an N-core machine with N workers, sequential DynamicSupervisor.stop calls could easily exceed the shutdown budget. Now stops all workers in parallel.
  • Executer had no terminate/2 — when shut down externally, the spawned message-processing process was killed with no opportunity to finish. Now uses Task.async + Task.shutdown/2 to give in-flight work a 5s grace period before escalating, and nacks with requeue: true as a safety net.
  • No explicit shutdown timeouts — all processes used OTP defaults (5s). Now uses a proper timeout chain: Consumer.Supervisor (:infinity) > Consumer.Server (30s) > Workers/Executers (25s).

Corrected Shutdown Sequence

SIGTERM
-> Consumer.Supervisor gets :shutdown (timeout: :infinity)
-> Each Consumer.Server gets :shutdown (timeout: 30s)
-> terminate/2:
1. AMQP.Basic.cancel — stop new message delivery
2. Stop all Workers in parallel (25s timeout)
-> Each Executer.terminate/2:
a. Task.shutdown(task, 5_000) — grace period for in-flight work
b. Safety-net nack with requeue: true if not completed
3. AMQP.Channel.close — clean channel shutdown
-> Producer.Pool, Topology.Server, Connection.Pool shut down after

Changes

  • lib/rabbit/broker/supervisor.ex — Consumer.Supervisor child spec gets type: :supervisor
  • lib/rabbit/consumer/supervisor.ex — Consumer.Server child specs get shutdown: 30_000
  • lib/rabbit/consumer/server.ex — Add cancel_consumer/1, parallel stop_workers/1
  • lib/rabbit/consumer/executer.ex — Refactor to Task.async, add terminate/2 with grace period and safety-net nack, shutdown: 25_000, completed state tracking
  • mix.exs — Version bump to 0.22.0

- Set Consumer.Supervisor child spec type to :supervisor so it gets
shutdown: :infinity instead of being killed after 5 seconds
- Set explicit shutdown: 30_000 on each Consumer.Server child spec
- Cancel AMQP consumption (Basic.cancel) before stopping workers so
no new messages arrive during drain
- Stop workers in parallel instead of sequentially to fit within
the shutdown budget
- Add terminate/2 to Executer that nacks unfinished messages with
requeue: true, preventing unacked messages from accumulating on
quorum queues
- Set Executer child_spec shutdown: 25_000 to give in-flight messages
time to complete before the safety-net nack
- Bump version to 0.22.0
Replace spawn_link with Task.async so terminate/2 can use
Task.shutdown/2 to give in-flight messages a 5s grace period
to complete before escalating to :kill and safety-net nacking.
Previously the spawned process was killed immediately on shutdown
with no chance to finish. Now the sequence is:
1. Task.shutdown(task, 5_000) - sends :shutdown, waits 5s
2. If task doesn't finish, brutally kills it
3. Safety-net nack with requeue: true
@nsweetingnsweeting changed the title fix: improve OTP shutdown behavior for consumersImprove OTP shutdown behavior for consumersMar 25, 2026
Check Task.shutdown/2 return value in terminate/2. If it returns
{:ok, _}, the task finished within the grace period and already
acked/nacked inside its body - skip the redundant nack.
- Executer lifecycle: normal completion, crash error handling, timeout
- Executer terminate/2: shutdown with in-flight task, task-never-started
edge case, skip nack when task completes before shutdown
- Child spec assertions: Executer shutdown/restart, Consumer.Server
shutdown timeout of 30_000
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant

@nsweeting