Packages
exq
0.24.0
0.24.0
0.23.0
0.22.0
0.21.0
0.20.0
0.19.0
0.18.0
0.17.0
0.16.2
0.16.1
0.16.0
0.15.0
0.14.0
0.13.5
0.13.4
0.13.3
0.13.2
0.13.1
0.13.0
0.12.2
0.12.1
0.12.0
0.11.0
0.10.1
0.10.0
0.9.1
0.9.0
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.3
0.7.2
0.7.1
0.7.0
0.6.5
0.6.4
0.6.3
0.6.2
0.6.1
0.6.0
0.5.0
0.4.4
0.4.3
0.4.2
0.4.1
0.4.0
0.3.0
0.2.3
0.2.2
0.2.1
0.2.0
0.1.4
0.1.3
0.1.2
0.1.1
0.1.0
0.0.2
Exq is a job processing library compatible with Resque / Sidekiq for the Elixir language.
Current section
Files
Jump to
Current section
Files
lib/exq/serializers/json_serializer.ex
defmodule Exq.Serializers.JsonSerializer do
@behaviour Exq.Serializers.Behaviour
alias Exq.Support.Config
defp json_library do
Config.get(:json_library)
end
def decode(json) do
json_library().decode(json)
end
def encode(e) do
json_library().encode(e)
end
def decode!(json) do
json_library().decode!(json)
end
def encode!(e) do
json_library().encode!(e)
end
def decode_job(serialized) do
deserialized = decode!(serialized)
%Exq.Support.Job{
args: Map.get(deserialized, "args"),
class: Map.get(deserialized, "class"),
enqueued_at: Map.get(deserialized, "enqueued_at"),
error_message: Map.get(deserialized, "error_message"),
error_class: Map.get(deserialized, "error_class"),
failed_at: Map.get(deserialized, "failed_at"),
retried_at: Map.get(deserialized, "retried_at"),
finished_at: Map.get(deserialized, "finished_at"),
jid: Map.get(deserialized, "jid"),
processor: Map.get(deserialized, "processor"),
queue: Map.get(deserialized, "queue"),
retry: Map.get(deserialized, "retry"),
retry_count: Map.get(deserialized, "retry_count"),
unique_for: Map.get(deserialized, "unique_for"),
unique_until: Map.get(deserialized, "unique_until"),
unique_token: Map.get(deserialized, "unique_token"),
unlocks_at: Map.get(deserialized, "unlocks_at")
}
end
def encode_job(job) do
deserialized = %{
args: job.args,
class: job.class,
enqueued_at: job.enqueued_at,
error_message: job.error_message,
error_class: job.error_class,
failed_at: job.failed_at,
retried_at: job.retried_at,
finished_at: job.finished_at,
jid: job.jid,
processor: job.processor,
queue: job.queue,
retry: job.retry,
retry_count: job.retry_count
}
deserialized =
if job.unique_for do
Map.merge(deserialized, %{
unique_for: job.unique_for,
unique_until: job.unique_until,
unique_token: job.unique_token,
unlocks_at: job.unlocks_at
})
else
deserialized
end
encode!(deserialized)
end
def decode_process(serialized) do
deserialized = decode!(serialized)
%Exq.Support.Process{
pid: Map.get(deserialized, "pid"),
host: Map.get(deserialized, "host"),
payload:
Map.get(deserialized, "payload")
|> Exq.Support.Job.decode(),
run_at: Map.get(deserialized, "run_at"),
queue: Map.get(deserialized, "queue")
}
end
def encode_process(process) do
deserialized =
Enum.into(
[
pid: process.pid,
host: process.host,
payload: process.payload,
run_at: process.run_at,
queue: process.queue
],
Map.new()
)
encode!(deserialized)
end
def encode_node(node) do
encode!(Map.from_struct(node))
end
def decode_node(serialized) do
deserialized = decode!(serialized)
%Exq.Support.Node{
hostname: Map.get(deserialized, "hostname"),
identity: Map.get(deserialized, "identity"),
started_at: Map.get(deserialized, "started_at"),
pid: Map.get(deserialized, "pid"),
queues: Map.get(deserialized, "queues"),
labels: Map.get(deserialized, "labels"),
tag: Map.get(deserialized, "tag"),
busy: Map.get(deserialized, "busy"),
quiet: Map.get(deserialized, "quiet"),
concurrency: Map.get(deserialized, "concurrency")
}
end
end