fireproofsocks

fireproofsocks

I just started looking over a Broadway repository and I’m wondering how to take advantage of Broadway.DummyProducer and/or test_messages. I can’t seem to make those run…

I tried to adapt the examples from the docs:

defmodule MyAppTest do
  use ExUnit.Case

  test "example test" do
    ref = Broadway.test_messages(MyApp, [1, 2, 3])
    assert_receive {:ack, ^ref, [_, _, _] = _successful, _failed}, 1000
  end
end

but that always generates an error:

     test/my_app_test.exs:5
     ** (exit) exited in: GenServer.call(MyApp, :producer_names, 5000)
         ** (EXIT) no process: the process is not alive or there's no process currently associated with the given name, possibly because its application isn't started
     code: ref = Broadway.test_messages(BroadwaySQS.Producer, [1, 2, 3])

I’ve tried using various other modules as input to the Broadway.test_messages function, but can’t seem to get past the error. I’m sure I’m missing something silly, but what is a valid argument to that function?

Showing Posts 1 to 4

NobbZ

NobbZ

Have you started your brodway? And how do you do so?

fireproofsocks

fireproofsocks OP

Ah… digging around a bit it’s making a bit more sense.

This is my application.ex:

defmodule MyApp.Application do
  # See https://hexdocs.pm/elixir/Application.html
  # for more information on OTP Applications
  @moduledoc false

  use Application

  def start(_type, _args) do
    children = []

    # See https://hexdocs.pm/elixir/Supervisor.html
    # for other strategies and supported options
    opts = [strategy: :one_for_one, name: MyApp.Supervisor]
    Supervisor.start_link(children, opts)
  end
end

In this case, the various Broadway processes are NOT started with the rest of the app. Instead, they get started in dedicated .exs scripts using something like

alias MyApp.Something

require Logger

# Parse arguments to get a duration value, but for simplicity:
duration = 30_000
{:ok, pid} = Something.start_link([])
Process.sleep(duration)

GenServer.stop(pid)
|> Logger.info()

If I add the Something.start_link([]) to my tests, then the test methods seem to work… at least it’s getting new error messages…

fireproofsocks

fireproofsocks OP

I’ve had to dig much deeper on this, and I’ve found a couple things… first is that the documentation appears to be wrong (but someone should check me on that before I submit a PR).

One thing to clarify is that the Broadway.test_messages/2 is pretty simplistic: the list of arguments that you provide it as the 2nd argument gets mapped onto %Broadway.Message{} structs as the data field. In our case, that was not sufficient to test our pipeline because our message handlers rely on pattern matching of the message metadata field as well.

A useful modification to my start_link/1 function allows me to use the Broadway.DummyProducer in my tests:

defmodule Something do
  use Broadway
  alias Broadway.Message
  
  def start_link(opts) do
    Broadway.start_link(__MODULE__,
      name: __MODULE__,
      producer: [
        module: {
          Keyword.get(opts, :producer, BroadwaySQS.Producer),
          queue_url: Application.get_env(:my_app, :sqs_queue_url),
          message_attribute_names: ["type", "foo"],
          config: Application.fetch_env!(:my_app, :sqs_config)
        }
      ],
      # ... batchers, etc...
    )
  end

The area in particular that does not work as advertised is the assert_receive. Instead of returning a list of all successful messages in one go, the batch_mode: :flush option seems to trigger each message to eject as soon as it’s done processing. So instead of a single assert_receive checking for a list with 3 elements, I had to do 3 assert_receive statements, each checking for a list with 1 element. Here’s how it looked (with my custom modifications to mimic what was in test_messages/2:

# I set up a fixture that supplied fully formed `%Broadway.Message{}` structs
test "pipeline is drained of test messages", %{msg1: msg1, msg2: msg2, msg3: msg3} do
      {:ok, _pid} = Something.start_link(
        producer: Broadway.DummyProducer,
      # ...
      )
      # We cannot make use of `Broadway.test_messages/2` because the list of values you provide it as the 2nd argument
      # maps to each message's `data` field, and that's all it does in the way of creating test messages.
      # Our messages are more complex: our messages must have metadata, so we come up with our own modification of the
      # `Broadway.test_messages/2` function:
      batch_mode = :flush
      ref = make_ref()
      ack = {Broadway.CallerAcknowledger, {self(), ref}, :ok}

      msg1 = %Message{msg1 | acknowledger: ack, batch_mode: batch_mode}
      msg2 = %Message{msg2 | acknowledger: ack, batch_mode: batch_mode}
      msg3 = %Message{msg3 | acknowledger: ack, batch_mode: batch_mode}

      messages = [msg1, msg2, msg3]

      :ok = Broadway.push_messages(Something, messages)

        # Assert that the messages have been consumed
        assert_receive {:ack, ^ref, [_] = _successful, _failed}, 100
        assert_receive {:ack, ^ref, [_] = _successful, _failed}, 100
        assert_receive {:ack, ^ref, [_] = _successful, _failed}, 100
    end

Note that a 4th assert_receive statement would fail because the mailbox would be empty.

I tried doing this using batch_mode: :bulk but that also seems to have returned 1 message per receive.

— All posts loaded —

Where Next? Top

Trending in Questions Top

stjefim
Hello! Suppose you are building workflow (order / task / payment) processing system with the following requirements: Each workflow con...
New
Blokh
Hey guys, I’ve got a huge CSV ( around 10 GB ) that needs to be processed hourly Do you guys have any suggestions what is the best prac...
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
Onor.io
I have what I’ve heard referred to as a “lookup table” in my database. This is a way of assigning codes to common values. One common lo...
New
jaybe78
Hello, I’m developing a online persistent chat system (what’s app) like using elixir/dynamodb/aws for a mobile app(flutter). The diffic...
New
Trolleger
What approach to take when sending live updates to “random” users Hi! I have a question, I have a little chat app, and when I create a DM...
New
widianto
I think I’ve found a small improvement I could contribute to <%= web_namespace %>.CoreComponents (installer/templates/phx_web/compo...
New

Other Trending Topics Top

garrison
Hobbes is a low-level distributed database for the Elixir programming language. Hobbes provides a simple, safe, and scalable storage lay...
New
mcass19
ExRatatui lets you cook up rich terminal UIs in Elixir, powered by Rust’s ratatui via Rustler NIFs. Build interactive terminal applicatio...
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
netoum
Corex is an accessible, unstyled UI component library for Phoenix that integrates Zag.js state machines using Vanilla JavaScript and Live...
New
wintermeyer
There are three potential reasons for members of this forum to have a look at https://vutuv.de You are tired or annoyed of LinkedIn. Yo...
New
webofbits
Aludel - LLM Evaluation Workbench Aludel is an embeddable Phoenix LiveView dashboard for evaluating and comparing LLM prompts across mult...
New

We're in Beta

About us Mission Statement

Options

Thread Display Mode




Thread Preview

Skip Thread Previews