Repository navigation
Create Subscription Option for NATS #78
Description
Activity
- added a parent issue
on Aug 19, 2026 On reviewing, this looks an awful lot like our current polling strategy in a different hat. Leaving this open for right this moment, but tempted to just close.
If I can't think of a reason to implement this by the time I've gotten through this issue's more obvious siblings, then I'm going to mark this issue as Invalid and close it without any code changes.
Found a reason to keep this active. Short answer: I was getting confused between the distinction between Pull and Poll. Everything this template does is Pull, but only some non-NATS libraries have subscribe modes that still Poll.
See recent comments in parent #72 for more context.
Fed the prompt through Cursor. With the MQs as an example, it did decent with the already-established classes. The
NatsMessageSourceneeded a bit of TLC. Subscribe mode didn't quite hit the mark. See below for more context.Tests so far:
- Container builds with tests enabled
- Polling mode continues to work.
- Polling mode recovers from a connection blip (read: killing the NATS container and bringing it back up) without any issue.
- Because of some structuring issues (see below), the new subscribe mode is not able to handle a message at the moment.
On first blush, the main problem seems to be that the NATS library doesn't handle message consumption in an under-the-hood manner like RabbitMQ/ActiveMQ. Cursor's take on the issue is that our application code is expected to run
ConsumeAsyncourselves on a loop. This is consistent with my exploration on the topic in the whole Pull vs Poll dilemma, but it just never fully clicked until now. This wouldn't be a problem... for re-subscribes, but the secret hidden trick is that in order to re-subscribe you need to subscribe first (speaking of, I haven't explored where the exception handling figures into this) . As a result,StartSubscriberAsyncnever exits, and the JobSourceSubscriptionManager never pivots to its second role of reading fromIJobSubscriberIntakeQueueand therefore nothing ever gets loaded into intake. As I type this, the message keeps getting re-received (which is a wee bit of a separate red flag)So then, I have two options:
- Rework the subscription manager setup so that subscribing and reading from the intake queue are two separate channels.
- Rework the NATS subscribe job source to be non-blocking.
I like door number 2 better. Door number 1 would imply a new value of sub-handler return, and it'd be a whole thing. Door number 2 has a smaller blast radius of just NATS.
Items for inspection:
- Follow Door number 2 and change
StartSubscriberAsyncto be non-blocking. - Confirm where the exception handling methods fit into things. The exception-handling method itself looks sane enough, but as of this moment I haven't looked at where it's used.
- Wait a second, are these NATS jobs timing out and being re-fetched after 30-odd seconds, or is that just a side effect of my current grief? It feels like it's a timeout, because being stuck in-memory the intake queue should be functionally identical to a job sitting in the repository for whatever reason.
Off-topic side-bar: I should add some more comments to
JobSubscriberIntakeQueue. I was double-checking some stuff in there, and I left zero guidance to explain my reasoning for the cancellation strategy.Created #108 for the side-bar stuff.
Oof. I take back everything I said about exception-handling looking sane, we're 100% in bonkers hallucination territory here when it comes to NATS' subscribe mode.
- On a new connection, set handlers. This is the sanest-sounding part, but as we'll see it doesn't make sense.
- First yellow flag,
OnReconnectFailedAsyncis just a wrapper aroundOnConnectionDisconnectedAsync. I feel like everything could have just routed toOnConnectionDisconnectedAsyncto begin with. OnConnectionDisconnectedAsyncmimics the exception judgment pattern. However, since theNatsEventArgsdoesn't include an exception, we make our own and it's alwaysNatsConnectionFailedException. Sure, that's implied from this specifically being attached to connection-related error handlers versus something like ActiveMQ. However, it makes some of the exception handling unnecessary because every exception is reconnect-flavoured.- If
OnConnectionDisconnectedAsyncdetects a reason to reconnect (and it always does), then it feeds everything into a new thread that runsSubscribeWithRetryLoopAsync. Except if the connection was disconnected, then the originalStartSubscriberAsynccall is still going to be doing it's ownSubscribeWithRetryLoopAsynccall, which is dutifully making sure that there's only one invocation ofSubscribeWithRetryLoopAsyncrunning. The logic that might lead to a new thread splatters and does nothing because of the thread safety guardrails. - I don't think that
WaitThenStopSubscriberAsyncis even doing anything functional.
So yeah, fun times!
Next steps (it's just past midnight, so in the name of a sane sleep schedule these are next-session steps):
StartSubscriberAsyncneeds to start a separate worker thread that will do the retry loop. Even with an AI-shaped hammer, not looking forward to double-checking the unit tests on this.- See if there's some way to salvage something out of
OnConnectionDisconnectedAsyncas an exception handler for the loop. - The thread safety in
SubscribeWithRetryLoopAsyncis completely unnecessary if we know that there will only be one instance ever. - I'm getting the feeling that we need to wire in an
IExecutionEndArbitercallback. Something will need to cancel theConsumeAsynccall running within our thread worker.- On a related side-note: Now that I think about it, the RabbitMQ/ActiveMQ handlers could cut an entire thread worker out of the equation by also using an
IExecutionEndArbitercallback. Will make other cleanup spin-off tasks for this.
- On a related side-note: Now that I think about it, the RabbitMQ/ActiveMQ handlers could cut an entire thread worker out of the equation by also using an
This doesn't even touch on the timeout thing. Leaving that problem until subscribe mode is at least a bit functional.
One other action item: The current implementation doesn't respect the fetch count. It just consumes endlessly.
Since the subscription queue is not an immediate filing (it's kind of right in the name), the job source will need to maintain its own count of fetched jobs. I think a count will be sufficient, but I'm not completely ruling out having a set of message IDs.
In-progress notes:
- Nothing about the on-disconnect handler method ended up being able to be salvaged.
- I take back what I said about the fetch count thing, was driven by a misc thought in between sessions. Given the quality of the rest of subscribe mode, definitely needs to be double-checked of course.
Progress report:
- Basic message handling is fixed.
ConsumeAsyncis 100% respecting the fetch count limit- Confirmed the timeout issue, will investigate further.
Well, it turns out that NATS has a heartbeat requirement after all. Fixing that.
Currently running a confirmation test on heartbeats:
- Fetch count: 10 (local testing default)
- Thread count: 2 (local testing default)
- Queued 12 jobs for 120 seconds apiece.
Test successful with subscribe mode, repeating test in polling mode.
Confirming that the polling mode test went well too.
In my next session on this, the plan is to confirm subscribe-mode's recovery ability.
So far, it's not recovering after a NATS container restart. My off the cuff theory is that the loop is gracefully stopping after the connection explodes.
Oops, I take that back. Somewhat. For starters, it just took a darned long time to trigger an exception:
worker-1 | RedShirt.Example.JobWorker.Core.Exceptions.WorkerJobSourceException: Consumer appears to be deleted after 10 consecutive 503 errors worker-1 | ---> NATS.Client.JetStream.NatsJSException: Consumer appears to be deleted after 10 consecutive 503 errorsYou can adjust the threshold for the number of consecutive 503s, so moved down to 2 from default of 10. A downed container and brought back up takes ~60s. I'm not 100% jazzed about this recovery time compared with other job sources, but it still recovers at least. Also, my particular method of testing an interruption with a container death is perhaps not the best example of a real-world connection interruption. For this branch, I think that this is acceptable.
Picking this up again, recapping:
- Subscribe-mode works.
- Subscribe-mode is more responsive than polling mode.
- Respects fetch count.
- Correct heartbeat handling.
- Respects halt-on-failure.
- Takes a moment to recover from a major issue that I'm not entirely happy with, but they key point is that it is able to recover.
Overall, I don't see a reason not to merge in what I've got at present.
Running one last solution cleanup sweep to quadruple-check for code standards.
Below is Cursor's opinion on the viability of subscribing using NATS.
NATS JetStream. Today it uses
FetchNoWaitAsync/NextAsync. NATS.Net’sConsumeAsyncis the documented long-running worker API: overlapping pull requests, per-message ack,MaxAckPendingas prefetch. Protocol-wise it is still pull, but at theIJobSourcelayer it is the same structure (start a consumer, push into the intake queue, stop polling).Testing
Similar to RabbitMQ, testing criteria is:
developalongside the polling implementation, it should be immediately responsive to new messagesJOBS__HALT_ON_FAILUREenvironment variable being set totruefor major problems. Break the exception arbiter in a local test if you have to in order to confirm this.