diff --git a/lib/towerops/job_monitoring/events.ex b/lib/towerops/job_monitoring/events.ex new file mode 100644 index 00000000..b581a1a5 --- /dev/null +++ b/lib/towerops/job_monitoring/events.ex @@ -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 diff --git a/test/towerops/job_monitoring/events_test.exs b/test/towerops/job_monitoring/events_test.exs new file mode 100644 index 00000000..b8e82e78 --- /dev/null +++ b/test/towerops/job_monitoring/events_test.exs @@ -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