rogerweb

rogerweb

Hi,

Scenario

I have a SQS queue with messages, each message is targeted to a specific user, each user joins his own Phoenix channel via websocket. I’m using Broadway to keep pulling the messages. In it’s handle_message/3 I use Phoenix’s broadcast/3 to push the message to the user:

Endpoint.broadcast("user:#{user_id}", "message", message)

I would like to make sure the user has received the message before I acknowledge it to Broadway/SQS.

Since the broadcast/3 doesn’t return the delivery result, I’m planning to make my JavaScript library to push an “ACK” back to the channel. I would include Broadway’s message_id in both messages.

The problem

How can I acknowledge the message to Broadway from the Phoenix’s channel handle_in/3?

As a side note, I started wondering if Broadway is the right tool for my use case, given the handle_message/3 documentation says:

Basically, any CPU bounded task that runs against a single message should be processed here.

but my task is more an I/O thing.

Any help is much appreciated.

Showing Posts 1 to 2

benwilson512

benwilson512

Author of Craft GraphQL APIs in Elixir with Absinthe

Yeah you probably want to do the broadcast in handle_batch because it would let you push a batch of messages from SQS to the end user, and then do a more efficient ACK back to SQS.

As far as getting the ACK from the client side, I think your only option is to block handle_batch until you get a message back from the client. Getting this ack back might be a bit tricky, perhaps you’ll have to have the broadway producer subscribe to an ack topic before pushing the messages out?

rogerweb

rogerweb OP

Yeah, I have no idea.

I was hoping Broadway could allow us to configure it to do not call handle_batch/4 automatically. Then I would call it myself once the last ACK arrives or times out.

— All posts loaded —

Where Next? Top

Trending in Questions Top

RSP87
I’m working on a project that simulates the bumbl example in the programming phoenix book. It acts almost like an email client. We have a...
New
kszambelanczyk
Hello! Could someone please give me a help/sample code, how to delete a file from s3 using waffle/waffle_ecto from Phoenix app. I creat...
New
RemyXRenard
I’m seeing that a list inside a Kino.DataTable will be interpreted as a charlist, even if the Kino.configure() is set to charlists: :as_l...
New
velrest
So my question is quite simple and i have found no conclusive answer on forum, google or AI. Should we use :erlang.float for Integer to ...
New
samoloth
Hi, I’ve just set up an application with ash_authentication. There is only magic link strategy for now, so there is no confirmation add o...
New
FlyingNoodle
If a change or preparation module uses Ash.Changeset.get_argument/2 or Ash.Query.get_argument/2 (or any of the other get_argument functio...
New
ryanwinchester
apply_graft/2 doesn’t rewrite an add_many sub-workflow’s deps on an add step. Grafted jobs cancel with “upstream job was deleted” Version...
New

Other Trending Topics Top

mudasobwa
I am happy to introduce the very α version of the new programming language compiled to BEAM. Welcome Cure. It has literally three kille...
New
garrison
Hobbes is a low-level distributed database for the Elixir programming language. Hobbes provides a simple, safe, and scalable storage lay...
New
marciok
Hi there! We created Gust: A task orchestrator inspired by Airflow. For those who have never heard about Aiflow, it’s a Python-based wor...
New
jimsynz
Beam Bots (or just BB for short) is a framework for building fault-tolerant robotics applications in Elixir using familiar OTP patterns. ...
New
Dmk
Xamal is a deployment tool for Elixir apps that deploys native releases to bare metal servers over SSH. It’s a port of GitHub - basecamp/...
New
Damirados
Hello everyone. After busy few months I am happy to announce v0.1.0 of Emerge & Solve. They are GUI (Emerge) and State management (S...
New

We're in Beta

About us Mission Statement

Options

Thread Display Mode




Thread Preview

Skip Thread Previews