Packages

phoenix_kit

1.7.57
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
phoenix_kit lib modules sync data_exporter.ex
Raw

lib/modules/sync/data_exporter.ex

defmodule PhoenixKit.Modules.Sync.DataExporter do
@moduledoc """
Exports data from database tables for DB Sync module.
Provides functions to fetch records from tables with pagination,
handle data serialization, and stream large datasets.
## Security Considerations
- Uses parameterized queries to prevent SQL injection
- Validates table names against actual database tables
- Respects configured table whitelist/blacklist (future feature)
## Example
iex> DataExporter.get_count("users")
{:ok, 150}
iex> DataExporter.fetch_records("users", offset: 0, limit: 100)
{:ok, [
%{"id" => 1, "email" => "user@example.com", ...},
...
]}
"""
alias PhoenixKit.Modules.Sync.SchemaInspector
alias PhoenixKit.RepoHelper
@default_limit 100
@max_limit 1000
@doc """
Gets the exact count of records in a table.
## Options
- `:schema` - Database schema (default: "public")
"""
@spec get_count(String.t(), keyword()) :: {:ok, non_neg_integer()} | {:error, any()}
def get_count(table_name, opts \\ []) do
schema = Keyword.get(opts, :schema, "public")
# Validate table exists first
if SchemaInspector.table_exists?(table_name, schema: schema) do
# Use exact count for accuracy
# For very large tables, consider using estimated count from pg_stat
query = "SELECT COUNT(*) FROM #{quote_identifier(schema)}.#{quote_identifier(table_name)}"
case RepoHelper.query(query, []) do
{:ok, %{rows: [[count]]}} ->
{:ok, count}
{:error, reason} ->
{:error, reason}
end
else
{:error, :table_not_found}
end
end
@doc """
Fetches records from a table with pagination.
Records are returned as maps with string keys matching column names.
All values are JSON-serializable.
## Options
- `:offset` - Number of records to skip (default: 0)
- `:limit` - Maximum records to return (default: 100, max: 1000)
- `:schema` - Database schema (default: "public")
- `:order_by` - Column(s) to order by (default: primary key)
"""
@spec fetch_records(String.t(), keyword()) :: {:ok, [map()]} | {:error, any()}
def fetch_records(table_name, opts \\ []) do
schema = Keyword.get(opts, :schema, "public")
offset = max(Keyword.get(opts, :offset, 0), 0)
limit = min(Keyword.get(opts, :limit, @default_limit), @max_limit)
# Validate table exists and get schema info
case SchemaInspector.get_schema(table_name, schema: schema) do
{:ok, table_schema} ->
do_fetch_records(table_name, schema, table_schema, offset, limit, opts)
{:error, :not_found} ->
{:error, :table_not_found}
{:error, reason} ->
{:error, reason}
end
end
@doc """
Exports all records from a table as a stream.
Useful for large tables where loading all records into memory
is not practical. Returns a stream that yields batches of records.
## Options
- `:batch_size` - Records per batch (default: 500)
- `:schema` - Database schema (default: "public")
"""
@spec stream_records(String.t(), keyword()) :: {:ok, Enumerable.t()} | {:error, any()}
def stream_records(table_name, opts \\ []) do
schema = Keyword.get(opts, :schema, "public")
batch_size = Keyword.get(opts, :batch_size, 500)
case SchemaInspector.get_schema(table_name, schema: schema) do
{:ok, table_schema} ->
stream =
Stream.resource(
fn -> 0 end,
fn offset ->
case do_fetch_records(table_name, schema, table_schema, offset, batch_size, opts) do
{:ok, []} ->
{:halt, offset}
{:ok, records} ->
{[records], offset + length(records)}
{:error, _reason} ->
{:halt, offset}
end
end,
fn _offset -> :ok end
)
{:ok, stream}
{:error, reason} ->
{:error, reason}
end
end
# ===========================================
# PRIVATE FUNCTIONS
# ===========================================
defp do_fetch_records(table_name, schema, table_schema, offset, limit, opts) do
columns = Enum.map(table_schema.columns, & &1.name)
order_by = Keyword.get(opts, :order_by) || table_schema.primary_key
# Build column list for SELECT
column_list = Enum.map_join(columns, ", ", &quote_identifier/1)
# Build ORDER BY clause
order_clause = build_order_clause(order_by)
query = """
SELECT #{column_list}
FROM #{quote_identifier(schema)}.#{quote_identifier(table_name)}
#{order_clause}
LIMIT $1 OFFSET $2
"""
case RepoHelper.query(query, [limit, offset]) do
{:ok, %{rows: rows}} ->
records =
Enum.map(rows, fn row ->
columns
|> Enum.zip(row)
|> Enum.map(fn {col, val} -> {col, serialize_value(val)} end)
|> Map.new()
end)
{:ok, records}
{:error, reason} ->
{:error, reason}
end
end
defp build_order_clause([]), do: ""
defp build_order_clause(columns) when is_list(columns) do
order_cols = Enum.map_join(columns, ", ", &quote_identifier/1)
"ORDER BY #{order_cols}"
end
defp build_order_clause(column) when is_binary(column) do
"ORDER BY #{quote_identifier(column)}"
end
# Quote identifier to prevent SQL injection
# PostgreSQL uses double quotes for identifiers
defp quote_identifier(name) when is_binary(name) do
# Escape any double quotes in the identifier
escaped = String.replace(name, "\"", "\"\"")
"\"#{escaped}\""
end
# Serialize values to JSON-compatible format
defp serialize_value(nil), do: nil
defp serialize_value(value) when is_binary(value), do: value
defp serialize_value(value) when is_number(value), do: value
defp serialize_value(value) when is_boolean(value), do: value
defp serialize_value(value) when is_list(value), do: Enum.map(value, &serialize_value/1)
defp serialize_value(value) when is_map(value), do: serialize_map(value)
defp serialize_value(%Date{} = date), do: Date.to_iso8601(date)
defp serialize_value(%Time{} = time), do: Time.to_iso8601(time)
defp serialize_value(%DateTime{} = dt), do: DateTime.to_iso8601(dt)
defp serialize_value(%NaiveDateTime{} = ndt), do: NaiveDateTime.to_iso8601(ndt)
defp serialize_value(%Decimal{} = decimal), do: Decimal.to_string(decimal)
# Handle Ecto types
defp serialize_value(%{__struct__: _} = struct) do
if function_exported?(struct.__struct__, :to_string, 1) do
to_string(struct)
else
inspect(struct)
end
end
# Fallback for other types
defp serialize_value(value), do: inspect(value)
defp serialize_map(map) do
Map.new(map, fn {k, v} -> {to_string(k), serialize_value(v)} end)
end
end