Skip to the content.

Running Node.js / Python apps behind Conduit as a worker pool

Status: recipe, no dedicated Conduit feature required. The companion supervisor library referenced below (working name only, not final) does not exist yet as a published package — this document describes the pattern and the code it would wrap. See issue #290 (Node.js) and issue #291 (Python) for the background discussion and open questions.

The idea in one sentence

Conduit does what a reverse proxy does best — TLS, routing, load balancing, rate limiting, health checking, retries — and hands the actual request handling off to a pool of Node.js or Python worker processes running your application code. This is not a new Conduit feature: it’s the same “nginx + Node.js” / “nginx + Gunicorn” pattern that’s been standard practice for over a decade, made slightly more convenient with a small supervisor script that uses Conduit’s existing dynamic-upstream Admin API.

What this recipe is not: a CGI/Azure-Functions-style invoke-on-demand model, where a request causes a fresh process to start (or a scaled-to-zero one to wake) and nothing else runs in between. The workers here are a fixed pool of long-lived processes — always warm, sized for steady CPU-bound throughput, not for scale-to-zero or per-request cold starts. If what you want is the latter, that’s a materially different, not-yet-designed problem; see issue #290’s discussion for the distinction.

What Conduit already does for you, today, with zero new code

None of that requires anything beyond a normal proxy: config pointing at http://127.0.0.1:PORT. What’s missing is pool management: who starts N worker processes, watches their health, restarts crashed ones, and tells Conduit when a worker comes up or goes away. That’s what this recipe adds.

The mechanism: Conduit’s dynamic upstream Admin API

conduit upstreams add/remove/weight (the CLI subcommands) are thin clients over three real HTTP endpoints on the Admin API (bound to global.admin.bind, e.g. 127.0.0.1:2019):

POST /upstreams/add     {"route": "/api", "target": "http://127.0.0.1:4001", "weight": 1, "site": "*:8080"}
POST /upstreams/remove  {"route": "/api", "target": "http://127.0.0.1:4001", "site": "*:8080"}
POST /upstreams/weight  {"route": "/api", "target": "http://127.0.0.1:4001", "weight": 3, "site": "*:8080"}

Because these are plain HTTP endpoints, any process in any language can call them directly — not just the conduit binary’s own CLI. That’s the whole trick: a small supervisor script in your worker pool’s own language calls /upstreams/add when a worker becomes ready and /upstreams/remove when it exits, and Conduit’s existing load balancer does the rest.

What actually runs

Three kinds of process, side by side:

  1. conduit (Rust) — listens on the public port, routes, load-balances, handles TLS/rate-limiting/health-checks.
  2. A supervisor process (your code, in Node.js or Python) — starts N worker processes, watches them, restarts crashed ones, registers/ deregisters them with Conduit via the Admin API above.
  3. N worker processes (your code) — the actual application, one instance per process, each on its own loopback port.

Request flow: client → conduit:8080 (TLS, rate limit, route match, pick a healthy worker via the configured strategy) → one of the live worker processes (your application logic) → response flows back through Conduit. Conduit already owns the balancing/health/circuit-breaker decisions; the supervisor’s only job is keeping the worker list accurate.

Node.js example

// pool.js — worker-pool supervisor
const { fork } = require("child_process");
const http = require("http");

const NUM_WORKERS = 4;
const BASE_PORT = 4001;
const ADMIN_URL = "http://127.0.0.1:2019";
const ADMIN_TOKEN = process.env.CONDUIT_ADMIN_TOKEN; // unset if global.admin.token isn't configured
const ROUTE = "/api"; // must match a path prefix in conduit.yaml
// No `site` field below — this example is single-site, so the registration
// applies to whichever site serves ROUTE. Add `site: "host:port"` for a
// multi-site deployment.

function callAdmin(path, body) {
  return new Promise((resolve, reject) => {
    const data = JSON.stringify(body);
    const headers = {
      "Content-Type": "application/json",
      "Content-Length": Buffer.byteLength(data),
    };
    if (ADMIN_TOKEN) headers["Authorization"] = `Bearer ${ADMIN_TOKEN}`;
    const req = http.request(
      ADMIN_URL + path,
      { method: "POST", headers },
      (res) => {
        let chunks = "";
        res.on("data", (c) => (chunks += c));
        res.on("end", () => {
          if (res.statusCode < 200 || res.statusCode >= 300) {
            // A 401 (missing/wrong admin token) returns an empty body --
            // reject before JSON.parse would throw on it.
            reject(
              new Error(
                `Admin API ${path} returned HTTP ${res.statusCode}: ${chunks}`,
              ),
            );
            return;
          }
          resolve(chunks ? JSON.parse(chunks) : {});
        });
      },
    );
    req.on("error", reject);
    req.write(data);
    req.end();
  });
}

// Registers `target` and retries a few times on failure (Admin API
// momentarily unreachable, etc.) rather than leaving a live worker
// silently unregistered. The periodic re-register timer below is the
// longer-term backstop — this is just for the very first attempt.
async function registerWithRetry(target, attempts = 5, delayMs = 1000) {
  for (let i = 1; i <= attempts; i++) {
    try {
      await callAdmin("/upstreams/add", { route: ROUTE, target, weight: 1 });
      return true;
    } catch (err) {
      console.error(
        `worker ${target} registration attempt ${i}/${attempts} failed: ${err.message}`,
      );
      if (i < attempts) await new Promise((r) => setTimeout(r, delayMs));
    }
  }
  return false;
}

function spawnWorker(port) {
  // Workers don't need Admin API access -- strip the token from their env
  // rather than let a compromised application process reuse it to call
  // /upstreams/add, /reload, or anything else on the Admin API.
  const workerEnv = { ...process.env, PORT: port };
  delete workerEnv.CONDUIT_ADMIN_TOKEN;
  const worker = fork("./worker.js", [], { env: workerEnv });
  const target = `http://127.0.0.1:${port}`;
  let reregisterTimer = null;
  let alive = true;
  // Tracks this generation's 'ready' handling (registration attempt, plus
  // any compensating cleanup below) so the exit handler can wait for it to
  // fully settle before respawning -- see the comment on `setTimeout`
  // below for why that ordering guarantee matters.
  let readySettled = Promise.resolve();

  async function handleReady() {
    const ok = await registerWithRetry(target);
    if (!alive) {
      // Exited while registration was in flight. If it landed anyway
      // (after the 'exit' handler's own, necessarily premature
      // /upstreams/remove already ran), undo it -- otherwise a dead
      // target stays registered until the next periodic re-register tick.
      // Safe to await here (not fire-and-forget): the exit handler's
      // setTimeout below waits for this whole function to settle before
      // letting a respawned worker register the same `target`, so this
      // cleanup is always the *last* Admin API call to land for this
      // generation, never racing a replacement's own registration.
      if (ok) {
        await callAdmin("/upstreams/remove", { route: ROUTE, target }).catch(
          () => {},
        );
      }
      return;
    }
    if (!ok) {
      console.error(
        `worker ${port}: giving up on registration, killing and respawning`,
      );
      worker.kill();
      return;
    }
    console.log(`worker ${port} registered with Conduit`);
    // Re-register on a timer so a `conduit reload` (which clears every
    // dynamic registration, even for an unrelated config change) doesn't
    // silently drop this worker until it next crashes and respawns.
    // Idempotent — see the Admin API section above.
    reregisterTimer = setInterval(() => {
      callAdmin("/upstreams/add", { route: ROUTE, target, weight: 1 }).catch(
        (err) => {
          console.error(
            `worker ${port} periodic re-registration failed: ${err.message}`,
          );
        },
      );
    }, 30_000);
  }

  worker.on("message", (msg) => {
    if (msg === "ready") readySettled = handleReady();
  });

  worker.on("exit", async (code) => {
    alive = false;
    if (reregisterTimer) clearInterval(reregisterTimer);
    await callAdmin("/upstreams/remove", { route: ROUTE, target }).catch(
      () => {},
    );
    console.log(`worker ${port} exited (code ${code}), respawning...`);
    setTimeout(async () => {
      // Wait for this generation's own 'ready' handling -- including any
      // compensating cleanup it might still be running -- to fully settle
      // before letting the next generation register the same `target`.
      // Without this, a late compensating remove above could still arrive
      // at the Admin API *after* the replacement worker's own registration
      // and take down a healthy worker instead of a dead one. Bounded by
      // registerWithRetry's own attempts*delayMs (a few seconds, worst
      // case) -- not a real respawn-latency concern in practice.
      await readySettled;
      spawnWorker(port);
    }, 500);
  });
}

for (let i = 0; i < NUM_WORKERS; i++) spawnWorker(BASE_PORT + i);
// worker.js — your application, one instance per worker process
const http = require("http");
const port = process.env.PORT;

http
  .createServer((req, res) => {
    // req.headers['x-request-id'] is already set by Conduit's XRequestIdGuard —
    // propagate it into your own logs for cross-system correlation.
    res.writeHead(200, { "Content-Type": "application/json" });
    res.end(
      JSON.stringify({
        handledBy: `worker-${port}`,
        requestId: req.headers["x-request-id"],
      }),
    );
  })
  .listen(port, "127.0.0.1", () => process.send("ready"));
# conduit.yaml — the `global:`/`sites:` shape is required here: `global.admin`
# is only recognized on the full config shape, not the flat single-site
# shorthand (`{ "port": 8080, "proxy": {...} }`) also shown elsewhere in this
# repo's docs — the Admin API silently doesn't start under the flat shorthand.
global:
  admin:
    bind: "127.0.0.1:2019"
sites:
  - port: 8080
    proxy:
      "/api":
        strategy:
          round-robin # least-conn/other strategies work too; round-robin
          # makes distribution obvious when trying this out —
          # near-instant responses can make least-conn's tie-
          # breaking consistently favor one worker
        targets:
          - http://127.0.0.1:4001 # worker 0 — a static seed target (Conduit
            # rejects an empty targets list); the rest
            # are added dynamically by pool.js

Python example

The same pattern, using multiprocessing and urllib/requests for the Admin API calls instead of child_process/http. A production setup would more likely put Gunicorn/uvicorn workers behind this instead of hand-rolling multiprocessing — the supervisor’s job (register/deregister via the Admin API) stays the same regardless of what actually manages the worker processes underneath it.

# pool.py — worker-pool supervisor
import json
import multiprocessing
import os
import time
import urllib.request

NUM_WORKERS = 4
BASE_PORT = 5001
ADMIN_URL = "http://127.0.0.1:2019"
ROUTE = "/api"
# No `site` field below — this example is single-site, so the registration
# applies to whichever site serves ROUTE. Add site="host:port" for a
# multi-site deployment.


def call_admin(path: str, body: dict) -> dict:
    # Read fresh on every call rather than caching into a module-level
    # global -- a global would still be reachable from worker code after
    # run_worker()'s os.environ.pop() below: "fork" workers inherit it as
    # already-bound memory, and "spawn" workers re-run this module's
    # top-level code (rebinding it from the still-intact parent env)
    # before run_worker() ever gets a chance to strip anything. Reading
    # os.environ directly here means the pop actually takes effect for
    # any code path -- including a supervisor bug that calls call_admin()
    # from inside a worker -- not just the supervisor's own normal use.
    admin_token = os.environ.get("CONDUIT_ADMIN_TOKEN")  # unset if global.admin.token isn't configured
    data = json.dumps(body).encode()
    headers = {"Content-Type": "application/json"}
    if admin_token:
        headers["Authorization"] = f"Bearer {admin_token}"
    req = urllib.request.Request(ADMIN_URL + path, data=data, headers=headers, method="POST")
    with urllib.request.urlopen(req) as resp:
        return json.loads(resp.read())


def run_worker(port: int, ready: multiprocessing.synchronize.Event) -> None:
    # Workers don't need Admin API access -- strip the token from this
    # process's environment before your application code (or anything it
    # imports) can read it and reuse it against the Admin API. Needed on
    # both the "fork" start method (child inherits the parent's full
    # os.environ) and "spawn" (the new interpreter still inherits the OS
    # environment by default). Note this doesn't scrub /proc/<pid>/environ
    # on Linux -- that's a kernel-captured exec-time snapshot, not
    # something a running process can rewrite -- so this stops the token
    # from being read by application code via os.environ, not from a
    # co-resident process with same-uid/root /proc access to this worker.
    os.environ.pop("CONDUIT_ADMIN_TOKEN", None)
    from worker import serve  # your application's entry point

    serve(port, ready)


def register_with_retry(target: str, attempts: int = 5, delay_secs: float = 1.0) -> bool:
    """Register `target` and retry a few times on failure (Admin API
    momentarily unreachable, etc.) rather than leaving a live worker
    silently unregistered. The periodic re-register below is the
    longer-term backstop -- this is just for the very first attempt."""
    for i in range(1, attempts + 1):
        try:
            call_admin("/upstreams/add", {"route": ROUTE, "target": target, "weight": 1})
            return True
        except Exception as e:  # noqa: BLE001 - starter code, log and retry
            print(f"worker {target} registration attempt {i}/{attempts} failed: {e}", flush=True)
            if i < attempts:
                time.sleep(delay_secs)
    return False


READY_TIMEOUT_SECS = 10.0
TERMINATE_TIMEOUT_SECS = 5.0


def terminate_and_reap(proc: multiprocessing.Process) -> None:
    """terminate() (SIGTERM) doesn't guarantee the process actually exits --
    a worker that ignores or is slow to handle it would otherwise hang this
    join() forever. Escalate to kill() (SIGKILL, not ignorable) if it's
    still alive after a bounded wait, then join again to actually reap it."""
    proc.terminate()
    proc.join(timeout=TERMINATE_TIMEOUT_SECS)
    if proc.is_alive():
        proc.kill()
        proc.join()


def supervise(port: int) -> None:
    while True:
        ready = multiprocessing.Event()
        proc = multiprocessing.Process(target=run_worker, args=(port, ready))
        proc.start()
        # Blocks (bounded) until the worker has actually bound its socket --
        # without this, /upstreams/add could register a port Conduit can
        # route to before anything is listening on it. Bounded so a worker
        # that fails *before* HTTPServer(...) ever succeeds (import error,
        # port already in use, permission error) doesn't hang this slot
        # forever -- the other workers keep running unaffected either way.
        if not ready.wait(timeout=READY_TIMEOUT_SECS):
            print(
                f"worker {port}: not ready after {READY_TIMEOUT_SECS}s "
                f"(exitcode={proc.exitcode}), killing and respawning",
                flush=True,
            )
            terminate_and_reap(proc)
            time.sleep(0.5)
            continue
        target = f"http://127.0.0.1:{port}"
        if not register_with_retry(target):
            print(f"worker {port}: giving up on registration, killing and respawning", flush=True)
            terminate_and_reap(proc)
            time.sleep(0.5)
            continue
        print(f"worker {port} registered with Conduit", flush=True)

        # Re-register on a timer so a `conduit reload` (which clears every
        # dynamic registration, even for an unrelated config change) doesn't
        # silently drop this worker until it next crashes and respawns.
        # Idempotent -- see the Admin API section above.
        while proc.is_alive():
            proc.join(timeout=30)
            if proc.is_alive():
                try:
                    call_admin("/upstreams/add", {"route": ROUTE, "target": target, "weight": 1})
                except Exception as e:  # noqa: BLE001 - starter code, log and keep going
                    print(f"worker {port} periodic re-registration failed: {e}", flush=True)

        try:
            call_admin("/upstreams/remove", {"route": ROUTE, "target": target})
        except Exception as e:  # noqa: BLE001 - starter code, log and keep respawning
            print(f"worker {port} deregistration failed: {e}", flush=True)
        print(f"worker {port} exited, respawning...", flush=True)
        time.sleep(0.5)


if __name__ == "__main__":
    procs = [
        multiprocessing.Process(target=supervise, args=(BASE_PORT + i,))
        for i in range(NUM_WORKERS)
    ]
    for p in procs:
        p.start()
    for p in procs:
        p.join()
# worker.py — your application, one instance per worker process
from http.server import BaseHTTPRequestHandler, HTTPServer
import json
import multiprocessing


class Handler(BaseHTTPRequestHandler):
    def do_GET(self):
        # self.headers["X-Request-ID"] is already set by Conduit's
        # XRequestIdGuard — propagate it into your own logs.
        body = json.dumps({
            "handledBy": f"worker-{self.server.server_port}",
            "requestId": self.headers.get("X-Request-ID"),
        }).encode()
        self.send_response(200)
        self.send_header("Content-Type", "application/json")
        self.end_headers()
        self.wfile.write(body)


def serve(port: int, ready: multiprocessing.synchronize.Event | None = None) -> None:
    server = HTTPServer(("127.0.0.1", port), Handler)  # binds + listens synchronously
    if ready is not None:
        ready.set()  # signal the supervisor only after the socket is actually listening
    server.serve_forever()
# conduit.yaml — see the Node.js example above for why the global:/sites:
# shape is required here rather than the flat single-site shorthand.
global:
  admin:
    bind: "127.0.0.1:2019"
sites:
  - port: 8080
    proxy:
      "/api":
        strategy: round-robin
        targets:
          - http://127.0.0.1:5001 # worker 0 — static seed target

What this deliberately does not do

Open questions before this becomes a published package

Tracked in #290 / #291: the package name (placeholder only, not decided), whether to standardize on child_process/ multiprocessing or delegate to existing tools (cluster/PM2 for Node, Gunicorn/uWSGI for Python) for the actual process management, and how much of the health-check/backoff logic above is worth generalizing into a real library versus leaving as copy-paste-and-adapt starter code.