Packages
ecto_tablestore
0.8.2
0.15.1
0.15.0
0.14.0
0.13.3
0.13.2
0.13.1
0.13.0
0.12.2
0.12.1
0.12.0
0.11.2
0.11.1
0.11.0
0.10.1
0.10.0
0.9.0
0.8.3
0.8.2
0.8.1
0.8.0
0.7.0
0.6.1
0.6.0
0.5.11
0.5.10
0.5.9
0.5.8
0.5.7
0.5.6
0.5.5
0.5.4
0.5.3
0.5.2
0.5.1
0.5.0
0.4.2
0.4.1
0.4.0
0.3.8
0.3.7
0.3.6
0.3.5
0.3.4
0.3.3
0.3.2
0.3.1
0.3.0
0.2.0
0.1.0
Alibaba Tablestore adapter for Ecto
Current section
Files
Jump to
Current section
Files
lib/ecto_tablestore/migration/runner.ex
defmodule EctoTablestore.Migration.Runner do
@moduledoc false
use Agent, restart: :temporary
require Logger
def run(repo, version, module, operation, opts) do
level = Keyword.get(opts, :log, :info)
args = {self(), repo, module, %{level: level}}
{:ok, runner} =
DynamicSupervisor.start_child(EctoTablestore.MigratorSupervisor, {__MODULE__, args})
Process.put(:ecto_tablestore_runner, runner)
log(level, "== Running #{version} #{inspect(module)}.#{operation}/0")
{time, _} = :timer.tc(fn -> perform_operation(module, operation) end)
time = System.convert_time_unit(time, :microsecond, :second)
log(level, "== Migrated #{version} in #{time}s")
Agent.stop(runner)
end
def start_link({parent, repo, module, log}) do
Agent.start_link(fn ->
Process.link(parent)
%{
repo: repo,
migration: module,
commands: [],
log: log,
config: repo.config()
}
end)
end
def repo do
Agent.get(runner(), & &1.repo)
end
def list_table_names(instance) do
case Process.get(:ecto_tablestore_table_names) do
nil ->
{:ok, %{table_names: table_names}} = ExAliyunOts.list_table(instance)
Process.put(:ecto_tablestore_table_names, table_names)
table_names
table_names when is_list(table_names) ->
table_names
end
end
def repo_config(key, default) do
Agent.get(runner(), &Keyword.get(&1.config, key, default))
end
defp perform_operation(module, operation) do
apply(module, operation, [])
flush()
end
defp log(false, _msg), do: :ok
defp log(level, msg), do: Logger.log(level, msg)
defp runner do
case Process.get(:ecto_tablestore_runner) do
nil -> raise "could not find migration runner process for #{inspect(self())}"
runner -> runner
end
end
def push_command(fun) when is_function(fun, 1) do
Agent.update(runner(), &%{&1 | commands: [fun | &1.commands]})
end
defp flush do
%{commands: commands, repo: repo} = Agent.get(runner(), & &1)
commands
|> Enum.reverse()
|> Enum.each(& &1.(repo))
end
end