core-async 2023-06-19

Hey guys, I’m looking for some opinions on implementing a telegram handler with core.async as I’m quite new to it My app is a openAI telegram client. Telegram sends messages through a webhook. Telegram sends you messages from users in a queue. It won’t send you a new message until you replied with ok for the previous message it sent. Previously I handled the messages and replied directly in the webhook handler but this is a bad experience as messages from new users need to wait until the last message from a user has been handled. I’m using reitit as a router for the webhook. Here’s my initial idea with core.async

(def messages-chan (chan 1024))

;; Threads that handle telegram messages
(dotimes [_ 8]
  (go (loop []
        (when-some [message (<! messages-chan)]
          (handle-telegram-message message))
        (recur))))


(defn telegram-webhook-handler
  [{:keys [db] :as config}]
  (fn [req]
    (log/info "Telegram Webhook Request Body: " (-> req :parameters :body))
    (let [message (-> req :parameters :body :message)]
      (put! messages-chan message)
      {:status 200 :body "ok"})))
Some questions: 1. Will the logical threads block correctly on the (when-some [message (<! messages-chan)] ? 2. I’m doing api calls inside the handle-telegram-message part. from what I read it’s better to use JVM threads with thread for long await times or long running functions. Should I switch to that? 3. I’m using integrant to orchestrate my app starting but it’s not clear to me how I can orchestrate the opening and closing of the handler threads through integrant. The opening part I see, but not exactly the closing. Is closing the channels returned by the threads enough to close those threads? 4. I’m foreseeing that if the openAI service is down it’s very possible that my buffer will fill up and I don’t yet see how I can handle this case. What are some options? Thank you!

1. When using logical threads, it’s called “parking” instead of “blocking”, in your case they will park correctly. Although when putting the message on messages-chan you’re using put! which doesn’t block or park, meaning it could fill up the buffer and then it could led to an exception about >1024 waiting puts. Prefer using >!! which will block the request thread until messages are accepted by the channel. Consider non-blocking buffers if you don’t want the request thread to block (sliding-buffer, dropping-buffer ) 2. Yes, as the core.async go loops use a thread executor, if you block those threads on IO, all core.async go blocks could be blocked. Make (handle-telegram-message message) to do it’s operation in a new core.async thread and then in the go block use (<! (handle-telegram-message message)) 3. First, you could apply a small change to the go-loop that puts the recur inside the when-some . Then, whenever the channel is closed, the loop will end and the go block will end to, closing the chan returned by go. From here, you could close messages-chan and the go blocks will finish when channel is consumed. If you want to wait until all go blocks are finished, you could store the channels returned by the go-blocks somewhere and wait for all of them to close after closing messages-chan . 4. This one is a bit more complex and beyond uniquely core.async, and it depends on what do you want when that happens. You could drop new messages when the queue is full, maybe returning to the user that the queue is full. You could achieve this by a combination of dropping-buffer and offer! . E.g (if (offer! messages-chan msg) {:status 200 :body "ok"} {:body "queue is full" :status 4XX}) . If you want an unbundled queue that waits for openAI to be online, and reprocesses everything, you probably need to go beyond core.async and use some real database to store pending messages