feat: add PubSub event broadcasting for job lifecycle

Co-Authored-By: Claude Sonnet 4.5 <noreply@anthropic.com>
This commit is contained in:
Graham McIntire 2026-02-06 16:36:34 -06:00 committed by Graham McIntire
parent 9b0b000263
commit 9f30d366b1
2 changed files with 65 additions and 0 deletions

View file

@ -0,0 +1,37 @@
defmodule Towerops.JobMonitoring.Events do
@moduledoc """
PubSub event broadcasting for job lifecycle events.
"""
@topic "job:lifecycle"
@doc """
Broadcasts a job lifecycle event to PubSub.
## Events
- `:started` - Job began executing
- `:completed` - Job finished successfully
- `:failed` - Job failed with error
## Metadata
Optional map with additional context:
- `duration`: Execution time in seconds
- `error`: Error message/reason
- `items_collected`: Number of items discovered/polled
"""
@spec broadcast_job_event(Oban.Job.t(), atom(), map()) :: :ok | {:error, term()}
def broadcast_job_event(%Oban.Job{} = job, event, metadata \\ %{}) do
Phoenix.PubSub.broadcast(
Towerops.PubSub,
@topic,
%{
job_id: job.id,
worker: job.worker,
device_id: get_in(job.args, ["device_id"]),
event: event,
metadata: metadata,
timestamp: DateTime.utc_now()
}
)
end
end

View file

@ -0,0 +1,28 @@
defmodule Towerops.JobMonitoring.EventsTest do
use Towerops.DataCase, async: true
alias Towerops.JobMonitoring.Events
describe "broadcast_job_event/3" do
test "broadcasts job lifecycle event to PubSub" do
Phoenix.PubSub.subscribe(Towerops.PubSub, "job:lifecycle")
job = %Oban.Job{
id: 123,
worker: "Towerops.Workers.DevicePollerWorker",
args: %{"device_id" => 456}
}
Events.broadcast_job_event(job, :started, %{})
assert_receive %{
job_id: 123,
worker: "Towerops.Workers.DevicePollerWorker",
device_id: 456,
event: :started,
metadata: %{},
timestamp: %DateTime{}
}
end
end
end