Skip to content

Create Subscription Option for NATS #78

Description

@adeutscher

Below is Cursor's opinion on the viability of subscribing using NATS.

NATS JetStream. Today it uses FetchNoWaitAsync / NextAsync. NATS.Net’s ConsumeAsync is the documented long-running worker API: overlapping pull requests, per-message ack, MaxAckPending as prefetch. Protocol-wise it is still pull, but at the IJobSource layer it is the same structure (start a consumer, push into the intake queue, stop polling).

Testing

Similar to RabbitMQ, testing criteria is:

  • The implementation must be able to handle messages (spoiler alert)
  • In order for this effort to be worth merging into develop alongside the polling implementation, it should be immediately responsive to new messages
  • The implementation must be able to handle messages after a connection is lost/recovered. Restarting the message source should not be a process-breaker for the JobWorker.
  • The implementation should respect the JOBS__HALT_ON_FAILURE environment variable being set to true for major problems. Break the exception arbiter in a local test if you have to in order to confirm this.

Activity

  1. self-assigned this
    on Aug 19, 2026
  2. adeutscher commented on Aug 19, 2026

    @adeutscher
    OwnerAuthor

    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.

  3. adeutscher commented on Aug 19, 2026

    @adeutscher
    OwnerAuthor

    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.

  4. adeutscher commented on Aug 28, 2026

    @adeutscher
    OwnerAuthor

    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.

  5. adeutscher commented on Aug 30, 2026

    @adeutscher
    OwnerAuthor

    Fed the prompt through Cursor. With the MQs as an example, it did decent with the already-established classes. The NatsMessageSource needed 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 ConsumeAsync ourselves 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, StartSubscriberAsync never exits, and the JobSourceSubscriptionManager never pivots to its second role of reading from IJobSubscriberIntakeQueue and 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:

    1. Rework the subscription manager setup so that subscribing and reading from the intake queue are two separate channels.
    2. 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 StartSubscriberAsync to 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.

  6. adeutscher commented on Aug 30, 2026

    @adeutscher
    OwnerAuthor

    Created #108 for the side-bar stuff.

  7. adeutscher commented on Aug 30, 2026

    @adeutscher
    OwnerAuthor

    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, OnReconnectFailedAsync is just a wrapper around OnConnectionDisconnectedAsync. I feel like everything could have just routed to OnConnectionDisconnectedAsync to begin with.
    • OnConnectionDisconnectedAsync mimics the exception judgment pattern. However, since the NatsEventArgs doesn't include an exception, we make our own and it's always NatsConnectionFailedException. 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 OnConnectionDisconnectedAsync detects a reason to reconnect (and it always does), then it feeds everything into a new thread that runs SubscribeWithRetryLoopAsync. Except if the connection was disconnected, then the original StartSubscriberAsync call is still going to be doing it's own SubscribeWithRetryLoopAsync call, which is dutifully making sure that there's only one invocation of SubscribeWithRetryLoopAsync running. The logic that might lead to a new thread splatters and does nothing because of the thread safety guardrails.
    • I don't think that WaitThenStopSubscriberAsync is 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):

    • StartSubscriberAsync needs 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 OnConnectionDisconnectedAsync as an exception handler for the loop.
    • The thread safety in SubscribeWithRetryLoopAsync is 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 IExecutionEndArbiter callback. Something will need to cancel the ConsumeAsync call 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 IExecutionEndArbiter callback. Will make other cleanup spin-off tasks for this.

    This doesn't even touch on the timeout thing. Leaving that problem until subscribe mode is at least a bit functional.

  8. adeutscher commented on Aug 30, 2026

    @adeutscher
    OwnerAuthor

    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.

  9. adeutscher commented on Aug 30, 2026

    @adeutscher
    OwnerAuthor

    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.
  10. adeutscher commented on Aug 30, 2026

    @adeutscher
    OwnerAuthor

    Progress report:

    • Basic message handling is fixed.
    • ConsumeAsync is 100% respecting the fetch count limit
    • Confirmed the timeout issue, will investigate further.
  11. adeutscher commented on Aug 30, 2026

    @adeutscher
    OwnerAuthor

    Well, it turns out that NATS has a heartbeat requirement after all. Fixing that.

  12. adeutscher commented on Aug 30, 2026

    @adeutscher
    OwnerAuthor

    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.
  13. adeutscher commented on Aug 30, 2026

    @adeutscher
    OwnerAuthor

    Test successful with subscribe mode, repeating test in polling mode.

  14. adeutscher commented on Aug 30, 2026

    @adeutscher
    OwnerAuthor

    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.

  15. adeutscher commented on Aug 31, 2026

    @adeutscher
    OwnerAuthor

    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.

  16. adeutscher commented on Aug 31, 2026

    @adeutscher
    OwnerAuthor

    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 errors
    
  17. adeutscher commented on Aug 31, 2026

    @adeutscher
    OwnerAuthor

    You 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.

  18. linked a pull request that will close this issueNATS Subscribe-Mode #111on Aug 31, 2026
  19. adeutscher commented on Aug 31, 2026

    @adeutscher
    OwnerAuthor

    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.

  20. adeutscher commented on Aug 31, 2026

    @adeutscher
    OwnerAuthor

    Running one last solution cleanup sweep to quadruple-check for code standards.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

enhancementNew feature or request

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions