fix(coverage): structured stage logging + per-row hang detection

CoverageWorker was silent between progress updates, so a stuck
gdalwarp / runaway pixel loop showed up only as 'stuck at 15%' in
the UI with nothing in k8s logs. Now every stage logs with the
coverage id (terrain fetch, building load, per-tier LOS+NLOS, row
band progress every 10%, finalize) and the row stream uses a
10-minute timeout with on_timeout: :kill_task so a hang surfaces
as an explicit row-failed log line + NaN fill instead of a wedged
worker. Unexpected exceptions now log with stacktrace and fail the
coverage cleanly.
This commit is contained in:
Graham McIntire 2026-05-06 17:43:16 -05:00
parent cb04cdd1d9
commit 16e4997368

View file

@ -75,24 +75,51 @@ defmodule Towerops.Workers.CoverageWorker do
end
defp run(coverage) do
started = System.monotonic_time(:millisecond)
log(coverage, "starting compute (cell=#{coverage.cell_size_m} m, radius=#{coverage.radius_m} m)")
update_progress(coverage, "computing", 1)
with {:ok, antenna} <- resolve_antenna(coverage),
{:ok, {lat, lon}} <- resolve_location(coverage),
{:ok, bbox} <- compute_bbox(coverage, lat, lon),
_ = update_progress(coverage, "computing", 5),
{:ok, grid} <- fetch_terrain(bbox, coverage.cell_size_m, lat),
clutter = Buildings.for_bbox(bbox),
canopy = fetch_canopy(bbox, coverage.cell_size_m, lat),
_ = update_progress(coverage, "computing", 15),
{:ok, tiers} <-
compute_all_tiers(coverage, antenna, lat, lon, bbox, grid, clutter, canopy) do
finalize(coverage, tiers)
else
{:error, reason} -> fail(coverage, reason)
try do
with {:ok, antenna} <- resolve_antenna(coverage),
_ = log(coverage, "antenna resolved: #{antenna.slug}"),
{:ok, {lat, lon}} <- resolve_location(coverage),
{:ok, bbox} <- compute_bbox(coverage, lat, lon),
_ = update_progress(coverage, "computing", 5),
_ = log(coverage, "fetching terrain grid for bbox #{inspect(bbox)}"),
{:ok, grid} <- fetch_terrain(bbox, coverage.cell_size_m, lat),
_ =
log(
coverage,
"terrain grid: #{grid.ncols}x#{grid.nrows} cells (#{grid.ncols * grid.nrows} samples)"
),
clutter = Buildings.for_bbox(bbox),
_ = log(coverage, "loaded #{length(clutter)} building footprints for clutter"),
canopy = fetch_canopy(bbox, coverage.cell_size_m, lat),
_ = update_progress(coverage, "computing", 15),
_ = log(coverage, "starting per-tier compute (#{length(@sm_height_tiers_m)} tiers x LOS+NLOS)"),
{:ok, tiers} <-
compute_all_tiers(coverage, antenna, lat, lon, bbox, grid, clutter, canopy) do
elapsed = System.monotonic_time(:millisecond) - started
log(coverage, "compute finished in #{elapsed} ms; finalizing")
finalize(coverage, tiers)
else
{:error, reason} -> fail(coverage, reason)
end
rescue
e ->
Logger.error(
"CoverageWorker[#{coverage.id}]: unexpected exception: #{Exception.format(:error, e, __STACKTRACE__)}"
)
fail(coverage, {:terrain_unavailable, {:exception, Exception.message(e)}})
end
end
defp log(coverage, message) do
Logger.info("CoverageWorker[#{coverage.id}]: #{message}")
:ok
end
# Best-effort canopy fetch — non-fatal if disabled / out-of-area.
# Only NLOS sampling consumes it, so a nil here just means
# "compute as if no trees are present" for that mode.
@ -160,10 +187,25 @@ defmodule Towerops.Workers.CoverageWorker do
los_ctx = Map.merge(env, %{clutter: [], canopy: nil})
nlos_ctx = Map.merge(env, %{clutter: env.clutter, canopy: env.canopy})
log(coverage, "tier h=#{Float.round(height_m, 2)} m: computing LOS pixels")
t_los = System.monotonic_time(:millisecond)
with {:ok, los_out} <- compute_pixels(cov, los_ctx, env.lat, env.lon),
_ =
log(
coverage,
"tier h=#{Float.round(height_m, 2)} m: LOS pixels done in #{System.monotonic_time(:millisecond) - t_los} ms; writing raster"
),
{:ok, los_paths} <-
Raster.write(coverage, los_out.pixels, los_out.dims, tier_suffix(height_m, :los)),
_ = log(coverage, "tier h=#{Float.round(height_m, 2)} m: computing NLOS pixels"),
t_nlos = System.monotonic_time(:millisecond),
{:ok, nlos_out} <- compute_pixels(cov, nlos_ctx, env.lat, env.lon),
_ =
log(
coverage,
"tier h=#{Float.round(height_m, 2)} m: NLOS pixels done in #{System.monotonic_time(:millisecond) - t_nlos} ms; writing raster"
),
{:ok, nlos_paths} <-
Raster.write(coverage, nlos_out.pixels, nlos_out.dims, tier_suffix(height_m, :nlos)) do
{:ok, %{los: los_paths, nlos: nlos_paths}}
@ -232,6 +274,9 @@ defmodule Towerops.Workers.CoverageWorker do
pixel_ctx = Map.put(ctx, :eirp, eirp)
tx = {{lat, lon}, tx_z, tx_height}
log_every = max(1, div(nrows, 10))
started = System.monotonic_time(:millisecond)
rows =
0..(nrows - 1)
|> Task.async_stream(
@ -245,9 +290,28 @@ defmodule Towerops.Workers.CoverageWorker do
end,
ordered: true,
max_concurrency: System.schedulers_online(),
timeout: :infinity
# 10-minute ceiling per row band — surfaces hangs instead of pinning a
# worker forever. A real row should complete in ms.
timeout: to_timeout(minute: 10),
on_timeout: :kill_task
)
|> Enum.map(fn {:ok, row} -> row end)
|> Stream.with_index()
|> Stream.map(fn
{{:ok, row}, idx} ->
if rem(idx + 1, log_every) == 0 do
elapsed = System.monotonic_time(:millisecond) - started
Logger.info("CoverageWorker[#{coverage.id}]: pixel rows #{idx + 1}/#{nrows} (#{elapsed} ms elapsed)")
end
row
{{:exit, reason}, idx} ->
Logger.error("CoverageWorker[#{coverage.id}]: row #{idx} failed: #{inspect(reason)}; filling with NaN")
List.duplicate(:nan, ncols)
end)
|> Enum.to_list()
pixels = Enum.flat_map(rows, & &1)
@ -393,6 +457,7 @@ defmodule Towerops.Workers.CoverageWorker do
})
Coverages.broadcast(updated, {:coverage_status, :ready, 100})
log(coverage, "marked ready (raster_path=#{tif})")
_ = Notifier.deliver_compute_complete(updated, :ready)
:ok
end