Powered by AppSignal & Oban Pro

Queue Storage with AzureSDK

livebooks/queue_storage.livemd

Queue Storage with AzureSDK

Mix.install([
  # {:azure_sdk, path: Path.expand("..", __DIR__)},
  {:azure_sdk, "~> 0.4.1"}
])

Overview

Creates a queue against local Azurite (port 10001), enqueues a message, peeks, gets, extends visibility without rewriting the body, and deletes it. A second section shows plain-text (message_encoding: :none) round-trips for queues shared with the v12 Python or .NET SDKs.

Requires Azurite:

docker compose up -d
credential =
  AzureSDK.Identity.SharedKeyCredential.new(
    "devstoreaccount1",
    "Eby8vdM02xNOcqFlqUwJPLlmEtlCDXJ1OUzFT50uSRZ6IFsuFq2UVErCz4I6tq/K1SZFPTOtr/KBHBeksoGMGw=="
  )

client =
  AzureSDK.Storage.Client.new(
    account: "devstoreaccount1",
    credential: credential,
    endpoint: "http://127.0.0.1:10001/devstoreaccount1"
  )

queue = "livebook-jobs"

case AzureSDK.Storage.Queue.exists?(client, queue) do
  true -> :ok
  false -> {:ok, _} = AzureSDK.Storage.Queue.create(client, queue)
  {:error, error} -> raise "Azurite Queue is not reachable: #{error.message}"
end

Put, peek, get, update, delete

{:ok, %{id: _id, content: nil}} =
  AzureSDK.Storage.Queue.Message.put(client, queue, "hello from livebook")

{:ok, [peeked]} = AzureSDK.Storage.Queue.Message.peek(client, queue)
peeked.content
{:ok, [msg]} =
  AzureSDK.Storage.Queue.Message.get(client, queue, visibility_timeout: 30)

# Extend visibility only; omit :content so the body is kept.
{:ok, msg} =
  AzureSDK.Storage.Queue.Message.update(client, queue, msg.id,
    pop_receipt: msg.pop_receipt,
    visibility_timeout: 60
  )

{:ok, :deleted} =
  AzureSDK.Storage.Queue.Message.delete(client, queue, msg.id,
    pop_receipt: msg.pop_receipt
  )

msg.id

Plain-text messages

Producers on the v12 Python or .NET SDKs send plain text by default. Read and write those queues with message_encoding: :none on every call that touches the body.

{:ok, _} =
  AzureSDK.Storage.Queue.Message.put(client, queue, ~s({"job": 1}),
    message_encoding: :none
  )

{:ok, [plain]} =
  AzureSDK.Storage.Queue.Message.get(client, queue, message_encoding: :none)

{:ok, :deleted} =
  AzureSDK.Storage.Queue.Message.delete(client, queue, plain.id,
    pop_receipt: plain.pop_receipt
  )

plain.content

List and properties

{:ok, %{approximate_message_count: _count}} =
  AzureSDK.Storage.Queue.properties(client, queue)

{:ok, queues} =
  AzureSDK.Storage.Queue.list(client, prefix: "livebook-", include_metadata: true)

Enum.map(queues, & &1.name)

Clean up

{:ok, :cleared} = AzureSDK.Storage.Queue.clear_messages(client, queue)
{:ok, :deleted} = AzureSDK.Storage.Queue.delete(client, queue)
:ok