Skip to content

Commit cf86ede

Browse files
committed
fix: fallback to manual track input when catalog is unavailable
1 parent d10edcc commit cf86ede

1 file changed

Lines changed: 59 additions & 21 deletions

File tree

lib/mix/tasks/moqx.moqtail.demo.ex

Lines changed: 59 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -104,17 +104,32 @@ defmodule Mix.Tasks.Moqx.Moqtail.Demo do
104104
cond do
105105
config.list_tracks_only ->
106106
Mix.shell().info("loading catalog (namespace=#{config.namespace})...")
107-
catalog = load_catalog!(subscriber, config.namespace, config.timeout)
108-
print_available_tracks!(catalog, config.show_raw)
107+
108+
case load_catalog(subscriber, config.namespace, config.timeout) do
109+
{:ok, catalog} ->
110+
print_available_tracks!(catalog, config.show_raw)
111+
112+
{:error, reason} ->
113+
Mix.raise("catalog unavailable: #{reason}")
114+
end
109115

110116
is_binary(config.track_name) ->
111117
run_stream_for_track!(subscriber, config, config.track_name)
112118

113119
true ->
114120
Mix.shell().info("loading catalog (namespace=#{config.namespace})...")
115-
catalog = load_catalog!(subscriber, config.namespace, config.timeout)
116-
track = choose_track!(catalog, nil, config.show_raw)
117-
run_stream_for_track!(subscriber, config, track.name)
121+
122+
track_name =
123+
case load_catalog(subscriber, config.namespace, config.timeout) do
124+
{:ok, catalog} ->
125+
choose_track!(catalog, nil, config.show_raw).name
126+
127+
{:error, reason} ->
128+
Mix.shell().error("catalog unavailable: #{reason}")
129+
prompt_track_name_without_catalog!()
130+
end
131+
132+
run_stream_for_track!(subscriber, config, track_name)
118133
end
119134
after
120135
:ok = MOQX.close(subscriber)
@@ -185,20 +200,20 @@ defmodule Mix.Tasks.Moqx.Moqtail.Demo do
185200
end
186201
end
187202

188-
defp load_catalog!(subscriber, namespace, timeout) do
203+
defp load_catalog(subscriber, namespace, timeout) do
189204
case fetch_catalog(subscriber, namespace, timeout) do
190205
{:ok, catalog} ->
191-
catalog
206+
{:ok, catalog}
192207

193208
{:error, reason} ->
194209
if catalog_fetch_fallback_reason?(reason) do
195210
Mix.shell().info(
196211
"catalog fetch was unavailable (#{reason}); falling back to live subscribe..."
197212
)
198213

199-
subscribe_catalog!(subscriber, namespace, timeout)
214+
subscribe_catalog(subscriber, namespace, timeout)
200215
else
201-
Mix.raise("catalog load failed: #{reason}")
216+
{:error, reason}
202217
end
203218
end
204219
end
@@ -220,7 +235,7 @@ defmodule Mix.Tasks.Moqx.Moqtail.Demo do
220235
end
221236
end
222237

223-
defp subscribe_catalog!(subscriber, namespace, timeout) do
238+
defp subscribe_catalog(subscriber, namespace, timeout) do
224239
deadline = System.monotonic_time(:millisecond) + timeout
225240
subscribe_catalog_loop(subscriber, namespace, deadline)
226241
end
@@ -229,24 +244,27 @@ defmodule Mix.Tasks.Moqx.Moqtail.Demo do
229244
remaining = deadline - System.monotonic_time(:millisecond)
230245

231246
if remaining <= 0 do
232-
Mix.raise("timed out waiting for first catalog object")
247+
{:error, "timed out waiting for first catalog object"}
248+
else
249+
do_subscribe_catalog_loop(subscriber, namespace, deadline, remaining)
233250
end
251+
end
234252

253+
defp do_subscribe_catalog_loop(subscriber, namespace, deadline, remaining) do
235254
{:ok, sub_ref} = MOQX.subscribe(subscriber, namespace, "catalog")
236255
await_subscribed!(sub_ref, namespace, "catalog", remaining)
237256

238257
case await_catalog_payload(sub_ref, min(1_000, remaining)) do
239-
{:ok, payload} ->
240-
case MOQX.Catalog.decode(payload) do
241-
{:ok, catalog} -> catalog
242-
{:error, reason} -> Mix.raise("catalog decode failed: #{reason}")
243-
end
244-
245-
:retry ->
246-
subscribe_catalog_loop(subscriber, namespace, deadline)
258+
{:ok, payload} -> decode_catalog_payload(payload)
259+
:retry -> subscribe_catalog_loop(subscriber, namespace, deadline)
260+
{:error, reason} -> {:error, "catalog subscribe failed: #{reason}"}
261+
end
262+
end
247263

248-
{:error, reason} ->
249-
Mix.raise("catalog subscribe failed: #{reason}")
264+
defp decode_catalog_payload(payload) do
265+
case MOQX.Catalog.decode(payload) do
266+
{:ok, catalog} -> {:ok, catalog}
267+
{:error, reason} -> {:error, "catalog decode failed: #{reason}"}
250268
end
251269
end
252270

@@ -363,6 +381,26 @@ defmodule Mix.Tasks.Moqx.Moqtail.Demo do
363381
end
364382
end
365383

384+
defp prompt_track_name_without_catalog! do
385+
Mix.shell().info("catalog is unavailable; enter a track name manually (or q)")
386+
387+
case IO.gets("Track name: ") do
388+
:eof ->
389+
Mix.raise("stdin closed")
390+
391+
nil ->
392+
Mix.raise("stdin closed")
393+
394+
input ->
395+
case String.trim(input) do
396+
"" -> prompt_track_name_without_catalog!()
397+
"q" -> Mix.raise("aborted")
398+
"quit" -> Mix.raise("aborted")
399+
name -> name
400+
end
401+
end
402+
end
403+
366404
defp await_subscribed!(sub_ref, namespace, track_name, timeout) do
367405
receive do
368406
{:moqx_subscribed, ^sub_ref, ^namespace, ^track_name} -> :ok

0 commit comments

Comments
 (0)