Packages
snakepit
0.8.9
0.13.0
0.12.0
0.11.1
0.11.0
0.10.1
0.10.0
0.9.1
0.9.0
0.8.9
0.8.8
0.8.7
0.8.6
0.8.5
0.8.4
0.8.3
0.8.2
0.8.1
0.8.0
0.7.7
0.7.6
0.7.5
0.7.4
0.7.3
0.7.2
0.7.1
0.7.0
0.6.11
0.6.10
0.6.9
0.6.8
0.6.7
0.6.6
0.6.5
0.6.4
0.6.3
0.6.2
0.6.1
0.6.0
0.5.1
0.5.0
0.4.3
0.4.2
0.4.1
0.4.0
0.3.3
0.3.2
0.3.1
0.3.0
0.2.1
0.2.0
0.1.2
0.1.1
0.1.0
High-performance pooler and session manager for external language integrations. Supports Python, Node.js, Ruby, and more with gRPC streaming, session management, and production-ready process cleanup.
Current section
Files
Jump to
Current section
Files
lib/snakepit/worker/taint_registry.ex
defmodule Snakepit.Worker.TaintRegistry do
@moduledoc """
Tracks tainted workers and devices after crash classification.
"""
@table :snakepit_worker_taints
def taint_worker(worker_id, opts) when is_binary(worker_id) do
ensure_table()
now_ms = System.monotonic_time(:millisecond)
duration = Keyword.get(opts, :duration_ms, 60_000)
record = %{
tainted_until: now_ms + duration,
reason: Keyword.get(opts, :reason),
exit_code: Keyword.get(opts, :exit_code),
device: Keyword.get(opts, :device),
crashed_at: now_ms,
restart_notified: false
}
:ets.insert(@table, {worker_id, record})
:ok
end
def worker_tainted?(worker_id) when is_binary(worker_id) do
ensure_table()
case :ets.lookup(@table, worker_id) do
[{^worker_id, record}] ->
if expired?(record) do
:ets.delete(@table, worker_id)
false
else
true
end
_ ->
false
end
end
def worker_info(worker_id) when is_binary(worker_id) do
ensure_table()
case :ets.lookup(@table, worker_id) do
[{^worker_id, record}] ->
if expired?(record) do
:ets.delete(@table, worker_id)
:error
else
{:ok, record}
end
_ ->
:error
end
end
def consume_restart(worker_id) when is_binary(worker_id) do
ensure_table()
case :ets.lookup(@table, worker_id) do
[{^worker_id, record}] ->
cond do
expired?(record) ->
:ets.delete(@table, worker_id)
:error
record.restart_notified ->
:error
true ->
updated = Map.put(record, :restart_notified, true)
:ets.insert(@table, {worker_id, updated})
{:ok, record}
end
_ ->
:error
end
end
def clear_worker(worker_id) when is_binary(worker_id) do
ensure_table()
:ets.delete(@table, worker_id)
:ok
end
defp expired?(record) do
now_ms = System.monotonic_time(:millisecond)
now_ms > Map.get(record, :tainted_until, 0)
end
defp ensure_table do
case :ets.whereis(@table) do
:undefined ->
try do
:ets.new(@table, [:named_table, :set, :public, {:read_concurrency, true}])
rescue
ArgumentError ->
:ok
end
@table
_ ->
@table
end
end
end