Backtests were getting slow enough to matter, so they moved to a couple of old laptops at my house. The problem was that a backtest downloading historical bars uses the same broker account, the same OAuth token, and the same request quota as the live session trading real money on a small cloud droplet. Two machines, one budget.


Where this fits

The last several posts lived inside a single node — account state, portfolio heat, the processes that own them. This one steps outside the trading node for the first time. It's the first of a short run on infrastructure: what I chose, why it was enough, and what it cost.


Two things that must exist exactly once

Splitting backtests onto a worker was easy in principle. The jobs are already Oban jobs; put another node on the same Postgres and it picks them up. What made it a design problem was two pieces of state that are only correct if there's one of them.

The token stores. Each broker has a TokenStore GenServer that holds the access tokens for every credential on that broker, keeps them refreshed, and hands one to whoever asks. TSSession doesn't manage tokens — it asks the store on every call. Run a second store on the worker and you have two processes independently deciding when to refresh the same credential. For a broker that rotates refresh tokens on use, the second refresh invalidates the first holder's token.

The rate limiter. TradeStation's barchart endpoint allows a fixed number of requests per five-minute sliding window, and that window is per broker user, not per machine. Live sessions pull bars from it. Backtests pull bars from it far more aggressively. A rate limiter on each node would let each node spend the whole budget, and the broker would enforce the real limit for me — on the live session, at whatever moment the worker happened to be busy.

  def check_rate(config, :barchart) do
    case RateLimiter.check_sliding_rate(get_key(config) <> ":barchart", 300_000, 490) do
      {:ok, num} -> {:ok, num}
      # ...

Both are cluster-wide singletons.

The usual answer is another server

Rate limit counters in Redis, tokens in Redis or a secrets service, and probably a message broker in between once you're there. That's the conventional shape, and it works.

It also means the token store stops being a process that does something — schedules a refresh, retries on a network error, drops a token the broker has revoked — and becomes a key that two clients race to update. The behavior has to move into both clients, or into a lock, or into a third service whose job is to be the process I already had.

I didn't run a bake-off. Distributed Erlang was already in the VM I was deploying, and every alternative started with adding something: a service to run, a client library, a second place for state to live and a second thing to monitor. For a system with one developer, the dependency you don't add is the one that never pages you.

What I had was a database both nodes connect to, GenServers that already owned this state correctly on one node, and a VM that can address a process on another machine as easily as one on its own.

The other reason is less defensible on a whiteboard and more honest: I wanted to learn how distribution actually behaves. I'd used it in passing, never with something depending on it.


One app, two releases

The worker wasn't a separate codebase. It was the same OTP application built as a second release, with one environment variable choosing which tree started:

myapp: [include_executables_for: [:unix], applications: [runtime_tools: :permanent]],
myapp_worker: [
  include_executables_for: [:unix],
  applications: [myapp: :permanent, runtime_tools: :permanent],
  rel_overlays_path: "rel/worker_overlays"
]
  def start(_type, _args) do
    release_mode = System.get_env("RELEASE_MODE", "web")

    case release_mode do
      "worker" -> start_worker()
      _web -> start_web()
    end
  end

The web tree started everything: the endpoint, the trading supervisors, the three broker token stores, the rate limiter. The worker tree started the database, caches, Oban, PubSub, an HTTP client — and conspicuously not the token stores or the rate limiter. The comment in the worker's child list was the whole design in one line:

      # NOTE: Token stores and RateLimiter are accessed via the web server through this cluster connection

Moving the token stores to :global

Making each token store reachable from the other node was, mechanically, this — applied to all three stores:

# before
GenServer.start(__MODULE__, opts, name: __MODULE__)
GenServer.call(__MODULE__, {:get_token, token.id})

# after
GenServer.start(__MODULE__, opts, name: {:global, __MODULE__})
GenServer.call({:global, __MODULE__}, {:get_token, token.id})

{:global, name} registers the process in :global, which every connected node shares. A call from the worker resolves the name to a pid on the web node and sends the message over the distribution connection. The caller's code didn't change. The TradeStation data source running inside a backtest on the worker asked for a token exactly the way it did on the web node, and got the one the web node's store was keeping fresh.

That's the property I wanted: the singleton kept its behavior. The refresh schedule, the retry logic, the "this credential has been revoked" path — all of it stayed in one process. Nothing about it had to be redesigned for a second client, because there is no second client. There's a second caller.

The rate limiter had to give something up

The rate limiter needed more than a name change. It was built the way you'd build a fast local limiter: an ETS table, written directly from the caller's process. No message, no bottleneck.

ETS tables don't cross nodes. So every check had to go through the GenServer:

  @moduledoc """
  GenServer providing token bucket and sliding window rate limiting using ETS.

  Globally registered for cluster-wide rate limiting across nodes.
  """

  def start_link(args \\ []) do
    GenServer.start_link(__MODULE__, args, name: {:global, __MODULE__})
  end

  # Public API - all calls go through GenServer for cluster-wide access
  def check_rate(id, scale, limit) do
    GenServer.call({:global, __MODULE__}, {:check_rate, id, scale, limit})
  end

  def check_sliding_rate(id, scale, limit) do
    GenServer.call({:global, __MODULE__}, {:check_sliding_rate, id, scale, limit})
  end

That was a deliberate regression in local performance: every rate check became a serialized message, and from the worker a network round trip. It was the right trade here because of what's being limited. The budget is a few hundred requests per five minutes. A check that costs a millisecond more is irrelevant next to a broker HTTP call that costs a hundred. A rate limiter that's fast and wrong — two nodes each confidently spending the whole quota — is the expensive version.


Postgres is the queue

The other half of "sync" is work distribution, and that needed nothing new. Oban already used Postgres as its queue and its notifier. The worker just ran a different set of queues against the same tables:

  if release_mode == "worker" do
    # Worker release: only process back_tests queue with higher concurrency
    config :myapp, Oban,
      queues: [back_tests: 20],
      plugins: [
        # Only Lifeline for rescuing stuck jobs, let web server handle pruning/reindexing
        {Oban.Plugins.Lifeline, rescue_after: :timer.minutes(30)}
      ]
  end

The web node kept a single slot on back_tests and ran everything else. The worker took twenty — and, it turns out, ten slots on default as well. Elixir's config merges keyword lists rather than replacing them, so the worker's queues: [back_tests: 20] merged with the base config's [default: 10, back_tests: 1]. The comment says “only back_tests.” The running node disagreed. Every job at the time went to back_tests, so it never mattered, but runtime config is a merge, not an override.

Housekeeping plugins ran on the web node only, so the two nodes weren't both pruning and reindexing the same table.

So the full inventory of "infrastructure that keeps two servers in sync" was: Postgres, which was already there; Erlang distribution, which ships with the VM; and a VPN between a droplet in a data center and my house.

What doesn't cross the cluster

Very little crossed the cluster on purpose. Phoenix PubSub's default adapter works across nodes, and both trees started it, but every broadcast that mattered originated and was consumed on the web node — broker sessions, trader events, the dashboard. (The adapter still delivered those broadcasts to the worker, which had no subscribers and dropped them. Traffic, not coordination.) The candle caches were per node. Nothing streamed from worker to web; results landed in Postgres and the web node read them.

The worker made two kinds of call: give me a token and may I make a request. Keeping it that narrow is most of why it worked without a broker. There was no event stream to make durable, no fan-out, no ordering guarantees to reason about. Just synchronous calls to two processes that happened to live elsewhere.


What this cost

Two costs, and a premise I didn't check.

The first cost was visible: the web node was now a hard dependency of the worker. Every backtest that touched the broker called through to it. That's what "singleton" means, and I chose it.

The second was quieter: what GenServer.call({:global, name}, ...) does when the name isn't registered — because the worker isn't connected to the node that registered it. It doesn't wait. It doesn't return an error tuple. It exits the caller. Every job that needed a token would crash, Oban would retry it, and the retry would crash the same way. Nothing in the design above stops that from happening as fast as Oban can pick up jobs. I found out how fast later, and that's its own post.

And the premise. The single-refresher argument is airtight for a broker that rotates refresh tokens. It turned out TradeStation doesn't — a refresh token survives being used, so two independent refreshers on that credential would have been fine. I'd built the strict version of the constraint for every broker because one of them needed it. That ended up mattering when token reads had to stop depending on the cluster.

Why it was still the right call

None of that changes the original decision. A cluster of two, with a narrow, synchronous interface between them, is exactly what Distributed Erlang is good at. The token stores kept their behavior instead of being flattened into keys. The rate limiter became correct at the cost of speed it didn't need. The only new moving part was a network link I'd have needed for any shared service anyway — Redis on the far side of a VPN fails in the same places.

What I underestimated wasn't the approach. It was how much of a distributed system's behavior lives in the failure path, and how little of it I'd written.


Getting two BEAM nodes on different networks to find each other at all — the connector process, node naming, the VPN between them — is the next post.