if Code.ensure_loaded?(Oban) do defmodule WalletPasses.Sync.Worker do @moduledoc """ Oban worker that bulk-updates wallet passes. Optional — requires `oban` dependency. """ use Oban.Worker, queue: :wallet_passes_sync, max_attempts: 1, unique: [states: :incomplete] require Logger alias WalletPasses.Apple alias WalletPasses.Config alias WalletPasses.Google alias WalletPasses.Schema @impl Oban.Worker def perform(%Oban.Job{args: args}) do serial_numbers = Map.fetch!(args, "serial_numbers") exclude = args |> Map.get("exclude_statuses", []) |> MapSet.new() provider = Config.pass_data_provider() google_passes = Schema.list_all_google_passes() target_serials = MapSet.new(serial_numbers) excluded_serials = google_passes |> Enum.filter(&MapSet.member?(exclude, &1.status)) |> Enum.map(& &1.serial_number) |> MapSet.new() filtered_google = google_passes |> Enum.filter(&MapSet.member?(target_serials, &1.serial_number)) |> Enum.reject(&MapSet.member?(exclude, &1.status)) {google_ok, google_err} = Enum.reduce(filtered_google, {0, 0}, fn gp, {ok, err} -> case provider.build_pass_data(gp.serial_number) do {:ok, %{pass_data: pass_data, google: google_visual}} -> visual = google_visual || %Google.Visual{} case Google.Api.update_object(pass_data, visual, gp.object_id) do {:ok, _} -> {ok + 1, err} {:error, reason} -> Logger.error( "Sync: Google update failed for #{gp.object_id}: #{inspect(reason)}" ) {ok, err + 1} end {:error, reason} -> Logger.error( "Sync: PassDataProvider failed for #{gp.serial_number}: #{inspect(reason)}" ) {ok, err + 1} end end) apple_tokens = serial_numbers |> Enum.reject(&MapSet.member?(excluded_serials, &1)) |> Enum.flat_map(&Schema.list_push_tokens_for_serial/1) |> Enum.uniq() {apple_ok, apple_err} = case Apple.Push.notify_devices(apple_tokens) do {:ok, counts} -> counts {:error, reason} -> Logger.error("Sync: Apple push failed: #{inspect(reason)}") {0, 0} end {:ok, %{google_ok: google_ok, google_err: google_err, apple_ok: apple_ok, apple_err: apple_err}} end end end