Skip to content

Commit 3041d23

Browse files
authored
🛠️ Fix meilisearch task monitoring to unzip responses (#2064)
1 parent 9865875 commit 3041d23

3 files changed

Lines changed: 63 additions & 32 deletions

File tree

lib/sequin/sinks/meilisearch/client.ex

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -33,7 +33,7 @@ defmodule Sequin.Sinks.Meilisearch.Client do
3333
case response_or_exception do
3434
%Req.Response{status: 200, body: encoded_body} ->
3535
# NOTE: Req does not automatically decode on retry functions
36-
case Jason.decode(encoded_body) do
36+
case encoded_body |> :zlib.gunzip() |> Jason.decode() do
3737
{:ok, %{"status" => status}} when status in ["enqueued", "processing"] ->
3838
true
3939

test/sequin/meilisearch_client_test.exs

Lines changed: 42 additions & 25 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,14 @@ defmodule Sequin.Sinks.Meilisearch.ClientTest do
55
alias Sequin.Factory.SinkFactory
66
alias Sequin.Sinks.Meilisearch.Client
77

8+
defp send_gzipped_response(conn, status_code, response_data) do
9+
gzipped_body = response_data |> Jason.encode!() |> :zlib.gzip()
10+
11+
conn
12+
|> Plug.Conn.put_resp_header("content-encoding", "gzip")
13+
|> Plug.Conn.send_resp(status_code, gzipped_body)
14+
end
15+
816
@sink %MeilisearchSink{
917
type: :meilisearch,
1018
endpoint_url: "http://127.0.0.1:7700",
@@ -51,9 +59,8 @@ defmodule Sequin.Sinks.Meilisearch.ClientTest do
5159
assert conn.method == "GET"
5260
assert conn.request_path == "/tasks/1"
5361

54-
Req.Test.json(conn, %{
55-
"status" => "success"
56-
})
62+
response_data = %{"status" => "succeeded"}
63+
send_gzipped_response(conn, 200, response_data)
5764
end)
5865

5966
assert :ok = Client.import_documents(@sink, "test", records)
@@ -77,9 +84,8 @@ defmodule Sequin.Sinks.Meilisearch.ClientTest do
7784
assert conn.method == "GET"
7885
assert conn.request_path == "/tasks/1"
7986

80-
Req.Test.json(conn, %{
81-
"status" => "success"
82-
})
87+
response_data = %{"status" => "succeeded"}
88+
send_gzipped_response(conn, 200, response_data)
8389
end)
8490

8591
ids = Enum.map(records, & &1["id"])
@@ -158,7 +164,8 @@ defmodule Sequin.Sinks.Meilisearch.ClientTest do
158164
assert conn.method == "GET"
159165
assert conn.request_path == "/tasks/123"
160166

161-
Req.Test.json(conn, %{"status" => "succeeded"})
167+
response_data = %{"status" => "succeeded"}
168+
send_gzipped_response(conn, 200, response_data)
162169
end)
163170

164171
assert :ok =
@@ -198,7 +205,8 @@ defmodule Sequin.Sinks.Meilisearch.ClientTest do
198205
assert conn.method == "GET"
199206
assert conn.request_path == "/tasks/456"
200207

201-
Req.Test.json(conn, %{"status" => "succeeded"})
208+
response_data = %{"status" => "succeeded"}
209+
send_gzipped_response(conn, 200, response_data)
202210
end)
203211

204212
assert :ok =
@@ -222,13 +230,15 @@ defmodule Sequin.Sinks.Meilisearch.ClientTest do
222230
assert conn.method == "GET"
223231
assert conn.request_path == "/tasks/789"
224232

225-
Req.Test.json(conn, %{
233+
response_data = %{
226234
"status" => "failed",
227235
"error" => %{
228236
"message" => "Invalid filter expression",
229237
"code" => "invalid_search_filter"
230238
}
231-
})
239+
}
240+
241+
send_gzipped_response(conn, 200, response_data)
232242
end)
233243

234244
assert {:error, error} =
@@ -268,19 +278,22 @@ defmodule Sequin.Sinks.Meilisearch.ClientTest do
268278
:counters.add(call_count, 1, 1)
269279
send(test_pid, {:task_check, count})
270280

271-
cond do
272-
count <= 2 ->
273-
# First two checks show task is still processing
274-
Req.Test.json(conn, %{"status" => "processing", "taskUid" => 123})
281+
response_data =
282+
cond do
283+
count <= 2 ->
284+
# First two checks show task is still processing
285+
%{"status" => "processing", "taskUid" => 123}
275286

276-
count == 3 ->
277-
# Third check shows task succeeded
278-
Req.Test.json(conn, %{"status" => "succeeded", "taskUid" => 123})
287+
count == 3 ->
288+
# Third check shows task succeeded
289+
%{"status" => "succeeded", "taskUid" => 123}
290+
291+
true ->
292+
# Should not get here
293+
%{"status" => "succeeded", "taskUid" => 123}
294+
end
279295

280-
true ->
281-
# Should not get here
282-
Req.Test.json(conn, %{"status" => "succeeded", "taskUid" => 123})
283-
end
296+
send_gzipped_response(conn, 200, response_data)
284297
end)
285298

286299
# Should succeed after multiple task status checks
@@ -315,7 +328,9 @@ defmodule Sequin.Sinks.Meilisearch.ClientTest do
315328
send(test_pid, {:task_check, count})
316329

317330
# Always return processing status
318-
Req.Test.json(conn, %{"status" => "processing", "taskUid" => 456})
331+
response_data = %{"status" => "processing", "taskUid" => 456}
332+
333+
send_gzipped_response(conn, 200, response_data)
319334
end)
320335

321336
# Should fail after exhausting retries
@@ -326,7 +341,7 @@ defmodule Sequin.Sinks.Meilisearch.ClientTest do
326341

327342
# Verify we made all 6 attempts (1 initial + 5 retries)
328343
for i <- 0..5 do
329-
assert_receive {:task_check, ^i}, 15_000
344+
assert_receive {:task_check, ^i}, 3_000
330345
end
331346

332347
refute_receive {:task_check, 6}, 100
@@ -348,14 +363,16 @@ defmodule Sequin.Sinks.Meilisearch.ClientTest do
348363
assert conn.method == "GET"
349364
assert conn.request_path == "/tasks/789"
350365

351-
Req.Test.json(conn, %{
366+
response_data = %{
352367
"status" => "failed",
353368
"taskUid" => 789,
354369
"error" => %{
355370
"code" => "invalid_document_id",
356371
"message" => "Document ID is invalid"
357372
}
358-
})
373+
}
374+
375+
send_gzipped_response(conn, 200, response_data)
359376
end)
360377

361378
# Should return error with details

test/sequin/meilisearch_pipeline_test.exs

Lines changed: 20 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,14 @@ defmodule Sequin.Runtime.MeilisearchPipelineTest do
99
alias Sequin.Runtime.SinkPipeline
1010
alias Sequin.Sinks.Meilisearch.Client
1111

12+
defp send_gzipped_response(conn, status_code, response_data) do
13+
gzipped_body = response_data |> Jason.encode!() |> :zlib.gzip()
14+
15+
conn
16+
|> Plug.Conn.put_resp_header("content-encoding", "gzip")
17+
|> Plug.Conn.send_resp(status_code, gzipped_body)
18+
end
19+
1220
describe "meilisearch pipeline" do
1321
setup do
1422
account = AccountsFactory.insert_account!()
@@ -123,7 +131,9 @@ defmodule Sequin.Runtime.MeilisearchPipelineTest do
123131

124132
"GET" ->
125133
assert conn.request_path =~ "/tasks/"
126-
Req.Test.json(conn, %{"status" => "succeeded"})
134+
135+
response_data = %{"status" => "succeeded"}
136+
send_gzipped_response(conn, 200, response_data)
127137
end
128138
end)
129139

@@ -176,13 +186,15 @@ defmodule Sequin.Runtime.MeilisearchPipelineTest do
176186
assert conn.method == "GET"
177187
assert conn.request_path == "/tasks/#{task_uid}"
178188

179-
Req.Test.json(conn, %{
189+
response_data = %{
180190
"status" => "failed",
181191
"error" => %{
182192
"message" => "Invalid function expression",
183193
"code" => "invalid_document_function"
184194
}
185-
})
195+
}
196+
197+
send_gzipped_response(conn, 200, response_data)
186198
end)
187199

188200
start_pipeline!(consumer)
@@ -274,9 +286,11 @@ defmodule Sequin.Runtime.MeilisearchPipelineTest do
274286
assert conn.method == "GET"
275287
assert conn.request_path == "/tasks/#{task_uid}"
276288

277-
Req.Test.json(conn, %{
278-
"status" => "success"
279-
})
289+
response_data = %{
290+
"status" => "succeeded"
291+
}
292+
293+
send_gzipped_response(conn, 200, response_data)
280294
end)
281295

282296
start_pipeline!(consumer)

0 commit comments

Comments
 (0)