Packages
phoenix_kit
1.7.63
1.7.207
1.7.206
1.7.205
1.7.204
1.7.203
1.7.202
1.7.201
1.7.200
1.7.199
1.7.198
1.7.197
1.7.196
1.7.194
1.7.193
1.7.192
1.7.191
1.7.190
1.7.189
1.7.187
1.7.186
1.7.185
1.7.184
1.7.183
1.7.182
1.7.181
1.7.180
1.7.179
1.7.178
1.7.177
1.7.176
1.7.175
1.7.174
1.7.173
1.7.172
1.7.171
1.7.170
1.7.169
1.7.168
1.7.167
1.7.166
1.7.165
1.7.164
1.7.162
1.7.161
1.7.160
1.7.159
1.7.157
1.7.156
1.7.155
1.7.154
1.7.153
1.7.152
1.7.151
1.7.150
1.7.149
1.7.146
1.7.145
1.7.144
1.7.143
1.7.138
1.7.133
1.7.132
1.7.131
1.7.130
1.7.128
1.7.126
1.7.125
1.7.121
1.7.120
1.7.119
1.7.118
1.7.117
1.7.116
1.7.115
1.7.114
1.7.113
1.7.112
1.7.111
1.7.110
1.7.109
1.7.108
1.7.107
1.7.106
1.7.105
1.7.104
1.7.103
1.7.102
1.7.101
1.7.100
1.7.99
1.7.98
1.7.97
1.7.96
1.7.95
1.7.94
1.7.93
1.7.92
1.7.91
1.7.90
1.7.89
1.7.88
1.7.87
1.7.86
1.7.85
1.7.84
1.7.83
1.7.82
1.7.81
1.7.80
1.7.79
1.7.78
1.7.77
1.7.76
1.7.75
1.7.74
1.7.71
1.7.70
1.7.69
1.7.66
1.7.65
1.7.64
1.7.63
1.7.62
1.7.61
1.7.59
1.7.58
1.7.57
1.7.56
1.7.55
1.7.54
1.7.53
1.7.52
1.7.51
1.7.49
1.7.44
1.7.43
1.7.42
1.7.41
1.7.39
1.7.38
1.7.37
1.7.36
1.7.34
1.7.33
1.7.31
1.7.30
1.7.29
1.7.28
1.7.27
1.7.26
1.7.25
1.7.24
1.7.23
1.7.22
1.7.21
1.7.20
1.7.19
1.7.18
1.7.17
1.7.16
1.7.15
1.7.14
1.7.13
1.7.12
1.7.11
1.7.10
1.7.9
1.7.8
1.7.7
1.7.6
1.7.5
1.7.4
1.7.3
1.7.2
1.7.1
1.7.0
1.6.20
1.6.19
1.6.18
1.6.17
1.6.16
1.6.15
1.6.14
1.6.13
1.6.12
1.6.11
1.6.10
1.6.9
1.6.8
1.6.7
1.6.6
1.6.5
1.6.4
1.6.3
1.5.2
1.5.1
1.5.0
1.4.9
1.4.8
1.4.7
1.4.6
1.4.5
1.4.4
1.4.3
1.4.2
1.4.1
1.4.0
1.3.2
1.3.1
1.3.0
1.2.10
1.2.9
1.2.8
1.2.7
1.2.5
1.2.4
1.2.2
1.2.1
1.2.0
1.1.0
1.0.0
A foundation for building Elixir Phoenix apps — SaaS, social networks, ERP systems, marketplaces, and more
Current section
Files
Jump to
Current section
Files
lib/modules/sync/workers/import_worker.ex
defmodule PhoenixKit.Modules.Sync.Workers.ImportWorker do
@moduledoc """
Oban worker for background data import from DB Sync.
This worker handles the actual import of records received from a sender,
processing them in the background so the user doesn't have to wait.
## Job Arguments
- `table` - The table name to import into
- `records` - List of records to import (JSON-serialized)
- `strategy` - Conflict resolution strategy ("skip", "overwrite", "merge")
- `session_code` - The sync session code (for tracking)
- `batch_index` - Optional batch index for large transfers
## Usage
The Receiver LiveView enqueues jobs after receiving data:
ImportWorker.new(%{
table: "users",
records: records,
strategy: "skip",
session_code: "ABC12345"
})
|> Oban.insert()
## Queue Configuration
Add the sync queue to your Oban config:
config :my_app, Oban,
queues: [default: 10, sync: 5]
"""
use Oban.Worker, queue: :sync, max_attempts: 3
alias PhoenixKit.Modules.Sync.DataImporter
alias PhoenixKit.Modules.Sync.SchemaInspector
require Logger
@impl Oban.Worker
def perform(%Oban.Job{args: args}) do
table = Map.fetch!(args, "table")
records = Map.fetch!(args, "records")
strategy = args |> Map.get("strategy", "skip") |> String.to_existing_atom()
session_code = Map.get(args, "session_code", "unknown")
batch_index = Map.get(args, "batch_index", 0)
schema = Map.get(args, "schema")
Logger.info(
"Sync.ImportWorker: Starting import for #{table} " <>
"(batch #{batch_index}, #{length(records)} records, strategy: #{strategy})"
)
# Create table if it doesn't exist and we have a schema
with :ok <- ensure_table_exists(table, schema) do
case DataImporter.import_records(table, records, strategy) do
{:ok, result} ->
Logger.info(
"Sync.ImportWorker: Completed import for #{table} (session: #{session_code}) - " <>
"created: #{result.created}, updated: #{result.updated}, " <>
"skipped: #{result.skipped}, errors: #{length(result.errors)}"
)
# Log any errors for debugging
log_import_errors(result.errors, table)
# Return success even if some records had errors
# (we've logged them and don't want to retry the whole batch)
:ok
{:error, reason} ->
Logger.error(
"Sync.ImportWorker: Failed import for #{table} (session: #{session_code}) - " <>
"#{inspect(reason)}"
)
# Return error to trigger Oban retry
{:error, reason}
end
end
end
defp log_import_errors(errors, table) do
for {record, error} <- errors do
Logger.warning("Sync.ImportWorker: Error importing record in #{table}: #{inspect(error)}")
pk_info = extract_record_pk(record)
Logger.debug("Sync.ImportWorker: Failed record#{pk_info}: #{inspect(record)}")
end
end
defp extract_record_pk(record) do
pk =
Map.get(record, "uuid") || Map.get(record, :uuid) ||
Map.get(record, "id") || Map.get(record, :id)
case pk do
nil -> ""
value -> " (uuid: #{value})"
end
end
defp ensure_table_exists(table, schema) when is_map(schema) do
if SchemaInspector.table_exists?(table) do
:ok
else
Logger.info("Sync.ImportWorker: Creating missing table #{table}")
case SchemaInspector.create_table(table, schema) do
:ok ->
Logger.info("Sync.ImportWorker: Created table #{table}")
:ok
{:error, reason} ->
Logger.error("Sync.ImportWorker: Failed to create table #{table}: #{inspect(reason)}")
{:error, {:table_creation_failed, reason}}
end
end
end
defp ensure_table_exists(_table, nil), do: :ok
defp ensure_table_exists(_table, _), do: :ok
@doc """
Creates a new import job for the specified table and records.
## Parameters
- `table` - Table name to import into
- `records` - List of record maps
- `strategy` - Conflict strategy (atom or string)
- `session_code` - Transfer session code for tracking
- `opts` - Additional options:
- `:batch_index` - Index of the batch for multi-batch imports
- `:schema` - Table schema definition (for auto-creating missing tables)
## Returns
An Oban.Job changeset ready for insertion.
"""
@spec create_job(String.t(), list(map()), atom() | String.t(), String.t(), keyword()) ::
Oban.Job.changeset()
def create_job(table, records, strategy, session_code, opts \\ []) do
strategy_str = if is_atom(strategy), do: Atom.to_string(strategy), else: strategy
args =
%{
"table" => table,
"records" => records,
"strategy" => strategy_str,
"session_code" => session_code,
"batch_index" => Keyword.get(opts, :batch_index, 0)
}
|> maybe_add_schema(Keyword.get(opts, :schema))
new(args)
end
defp maybe_add_schema(args, nil), do: args
defp maybe_add_schema(args, schema), do: Map.put(args, "schema", schema)
@doc """
Enqueues import jobs for multiple tables.
Splits large record sets into batches to avoid memory issues.
## Parameters
- `table_data` - Map of table name to one of:
- `{records, strategy}` - Records and strategy without schema
- `{records, strategy, schema}` - Records, strategy, and schema for auto-creating tables
- `session_code` - Transfer session code
- `batch_size` - Maximum records per job (default: 500)
## Returns
- `{:ok, job_count}` - Number of jobs enqueued
- `{:error, reason}` - If any job failed to enqueue
"""
@spec enqueue_imports(map(), String.t(), pos_integer()) ::
{:ok, non_neg_integer()} | {:error, term()}
def enqueue_imports(table_data, session_code, batch_size \\ 500) do
jobs =
table_data
|> Enum.flat_map(fn {table, table_info} ->
{records, strategy, schema} = normalize_table_info(table_info)
records
|> Enum.chunk_every(batch_size)
|> Enum.with_index()
|> Enum.map(fn {batch, index} ->
opts = [batch_index: index]
opts = if schema, do: Keyword.put(opts, :schema, schema), else: opts
create_job(table, batch, strategy, session_code, opts)
end)
end)
# Insert all jobs
results =
Enum.map(jobs, fn job ->
Oban.insert(job)
end)
# Check for failures
errors =
results
|> Enum.filter(fn
{:error, _} -> true
_ -> false
end)
if Enum.empty?(errors) do
{:ok, length(results)}
else
{:error, {:some_jobs_failed, length(errors)}}
end
end
defp normalize_table_info({records, strategy}), do: {records, strategy, nil}
defp normalize_table_info({records, strategy, schema}), do: {records, strategy, schema}
end