1defmodule RateLimitedConsumer do
2use GenServer
3
4@max_messages_per_second 5
5@sleep_interval_ms 1000 # 1 second
6
7# Client API
8def start_link(queue_name) do
9GenServer.start_link(__MODULE__, queue_name, name: __MODULE__)
10end
11
12def stop() do
13GenServer.cast(__MODULE__, :stop)
14end
15
16# Server Callbacks
17@impl true
18def init(queue_name) do
19IO.puts("RateLimitedConsumer started. Processing up to #{@max_messages_per_second} messages per second.")
20# In a real scenario, this would connect to a message queue.
21# For simulation, we'll use a GenServer as a mock queue.
22{:ok, %{queue_name: queue_name, message_count_in_interval: 0, last_interval_start: :os.system_time(:milli_seconds)}}
23end
24
25@impl true
26def handle_cast(:stop, state) do
27IO.puts("RateLimitedConsumer stopping.")
28{:stop, :normal, state}
29end
30
31@impl true
32def handle_info(:process_next, state) do
33current_time = :os.system_time(:milli_seconds)
34time_elapsed = current_time - state.last_interval_start
35
36# Reset message count if a new second has started
37if time_elapsed >= @sleep_interval_ms do
38state = %{state | message_count_in_interval: 0, last_interval_start: current_time}
39end
40
41# Check if we can process another message
42if state.message_count_in_interval < @max_messages_per_second do
43# Simulate fetching a message from the queue
44case MockQueue.fetch_message(state.queue_name) do
45{:ok, message} ->
46IO.puts("Processing message: #{inspect(message)}")
47# Simulate message processing time
48Process.sleep(50) # Short sleep for processing
49
50new_state = %{state | message_count_in_interval: state.message_count_in_interval + 1}
51# Schedule the next message processing attempt
52send(self(), :process_next)
53{:noreply, new_state}
54:empty ->
55# Queue is empty, wait a bit before checking again
56Process.sleep(100)
57send(self(), :process_next)
58{:noreply, state}
59end
60else
61# Rate limit reached, wait until the next interval
62time_to_wait = @sleep_interval_ms - time_elapsed
63Process.sleep(max(0, time_to_wait))
64send(self(), :process_next) # Try again after waiting
65{:noreply, state}
66end
67end
68
69# Handle unexpected messages
70def handle_info(message, state) do
71IO.puts("Received unexpected message: #{inspect(message)}")
72{:noreply, state}
73end
74end
75
76defmodule MockQueue do
77use GenServer
78
79@max_messages 10
80
81def start_link(name) do
82GenServer.start_link(__MODULE__, name, name: name)
83end
84
85def add_message(queue_name, message) do
86GenServer.cast(queue_name, {:add, message})
87end
88
89def fetch_message(queue_name) do
90GenServer.call(queue_name, :fetch)
91end
92
93@impl true
94def init(name) do
95IO.puts("MockQueue '#{name}' started.")
96{:ok, %{name: name, messages: :queue.new(), size: 0}}
97end
98
99@impl true
100def handle_cast({:add, message}, state) do
101if state.size < @max_messages do
102new_queue = :queue.in(message, state.messages)
103new_state = %{state | messages: new_queue, size: state.size + 1}
104IO.puts("Message added to queue '#{state.name}'. Current size: #{new_state.size}")
105{:noreply, new_state}
106else
107IO.puts("Queue '#{state.name}' is full. Message not added.")
108{:noreply, state}
109end
110end
111
112@impl true
113def handle_call(:fetch, _from, state) do
114case :queue.out(state.messages) do
115{:empty, _} ->
116{:reply, :empty, state}
117{{:value, message}, new_queue} ->
118new_state = %{state | messages: new_queue, size: state.size - 1}
119{:reply, {:ok, message}, new_state}
120end
121end
122end