Skip to content

Run background processes

Last updated View as MarkdownAgent setup

Start a process that keeps running after the request ends, such as a build, a test suite, or an agent task. The process object that exec() returns belongs to the request that started it, so later requests manage the process through files in the sandbox. The process writes its process ID, output, and exit code to those files.

For a process that runs until the container stops, such as a development server or a tunnel, refer to Run a server in the background.

Prerequisites

Run a process in the background

  1. Add shell scripts that run a process and report its state:

    src/index.jsjs
    const ROOT = "/var/lib/processes";
    const INACTIVITY_TIMEOUT_MS = 10 * 60 * 1000;
    
    // Runs the command in its own process group. Records its process ID when
    // it starts and its exit code when it ends.
    const RUN = `dir=$1; shift
    setsid sh -c 'echo "$$ $(cat /proc/sys/kernel/random/boot_id)" >"$0/pid"; exec "$@"' \\
    	"$dir" "$@" >"$dir/stdout.log" 2>"$dir/stderr.log"
    echo "$?" >"$dir/exit-code.tmp" && mv "$dir/exit-code.tmp" "$dir/exit-code"`;
    
    // Defines current(), which reads the process ID in a directory into $pid
    // when the process started in this instance.
    const CURRENT = `current() {
    	read -r pid boot 2>/dev/null <"$1/pid" &&
    		[ "$boot" = "$(cat /proc/sys/kernel/random/boot_id)" ]
    }`;
    
    // Prints the state of the process that owns the directory.
    const STATUS = `${CURRENT}
    dir=$1
    if [ ! -d "$dir" ]; then echo missing
    elif [ -e "$dir/exit-code" ]; then echo "exited $(cat "$dir/exit-code")"
    elif [ ! -e "$dir/pid" ]; then echo starting
    elif current "$dir" && kill -0 "$pid" 2>/dev/null; then echo "running $pid"
    elif [ -e "$dir/exit-code" ]; then echo "exited $(cat "$dir/exit-code")"
    else echo lost
    fi`;
    
    async function run(container, argv) {
    	const process = await container.exec(argv);
    	const output = await process.output();
    
    	return {
    		exitCode: output.exitCode,
    		stdout: new TextDecoder().decode(output.stdout),
    	};
    }
    src/index.tsts
    const ROOT = "/var/lib/processes";
    const INACTIVITY_TIMEOUT_MS = 10 * 60 * 1000;
    
    // Runs the command in its own process group. Records its process ID when
    // it starts and its exit code when it ends.
    const RUN = `dir=$1; shift
    setsid sh -c 'echo "$$ $(cat /proc/sys/kernel/random/boot_id)" >"$0/pid"; exec "$@"' \\
    	"$dir" "$@" >"$dir/stdout.log" 2>"$dir/stderr.log"
    echo "$?" >"$dir/exit-code.tmp" && mv "$dir/exit-code.tmp" "$dir/exit-code"`;
    
    // Defines current(), which reads the process ID in a directory into $pid
    // when the process started in this instance.
    const CURRENT = `current() {
    	read -r pid boot 2>/dev/null <"$1/pid" &&
    		[ "$boot" = "$(cat /proc/sys/kernel/random/boot_id)" ]
    }`;
    
    // Prints the state of the process that owns the directory.
    const STATUS = `${CURRENT}
    dir=$1
    if [ ! -d "$dir" ]; then echo missing
    elif [ -e "$dir/exit-code" ]; then echo "exited $(cat "$dir/exit-code")"
    elif [ ! -e "$dir/pid" ]; then echo starting
    elif current "$dir" && kill -0 "$pid" 2>/dev/null; then echo "running $pid"
    elif [ -e "$dir/exit-code" ]; then echo "exited $(cat "$dir/exit-code")"
    else echo lost
    fi`;
    
    type Status =
    	| { state: "starting" | "lost" }
    	| { state: "running"; pid: number }
    	| { state: "exited"; exitCode: number };
    
    async function run(container: Container, argv: string[]) {
    	const process = await container.exec(argv);
    	const output = await process.output();
    
    	return {
    		exitCode: output.exitCode,
    		stdout: new TextDecoder().decode(output.stdout),
    	};
    }

    setsid starts the command in a new process group. Stopping the group also stops the processes the command starts, such as the node process behind npm run dev. The script writes the exit code to a temporary file and then renames it, so a reader never sees a partial file. Do not start the command with & instead: a shell starts background jobs with SIGINT ignored.

    Each process keeps its files in its own directory under /var/lib/processes on the instance disk. The logs grow until the process ends. When the disk fills, writes from the process fail, so delete old process directories, or send output that you want to keep to a mounted bucket.

  2. Add a constructor to your Durable Object. It sets the inactivity timeout again when a restarted Durable Object finds the container running:

    src/index.tsts
    export class MyContainer extends DurableObject<Env> {
    	constructor(ctx: DurableObjectState, env: Env) {
    		super(ctx, env);
    		const container = ctx.container;
    		if (container?.running) {
    			void ctx.blockConcurrencyWhile(() =>
    				container.setInactivityTimeout(INACTIVITY_TIMEOUT_MS),
    			);
    		}
    	}
    }

    Then add a method that starts a process:

    src/index.tsts
    export class MyContainer extends DurableObject<Env> {
    	// ...
    
    	async startProcess(id: string, argv: string[]): Promise<boolean> {
    		const container = this.ctx.container;
    
    		if (!container) {
    			throw new Error("The container binding is not configured");
    		}
    
    		if (!container.running) {
    			container.start({
    				image: "cloudflare/debian-trixie",
    				entrypoint: ["sleep", "infinity"],
    				enableInternet: false,
    			});
    			await container.setInactivityTimeout(INACTIVITY_TIMEOUT_MS);
    		}
    
    		// mkdir fails when another process already uses the ID.
    		const dir = `${ROOT}/${id}`;
    		await run(container, ["mkdir", "-p", ROOT]);
    		const created = await run(container, ["mkdir", dir]);
    
    		if (created.exitCode !== 0) {
    			return false;
    		}
    
    		// Ignoring the output lets the process outlive this request.
    		await container.exec(["sh", "-c", RUN, "sh", dir, ...argv], {
    			stdout: "ignore",
    			stderr: "ignore",
    		});
    
    		await this.ctx.storage.put(`process:${id}`, argv);
    		await this.ctx.storage.setAlarm(Date.now() + 60_000);
    		return true;
    	}
    }
  3. Add a method that reports the state of a process:

    src/index.tsts
    export class MyContainer extends DurableObject<Env> {
    	// ...
    
    	async status(id: string): Promise<Status | undefined> {
    		const container = this.ctx.container;
    
    		if (!container?.running) {
    			return undefined;
    		}
    
    		const { stdout } = await run(container, [
    			"sh",
    			"-c",
    			STATUS,
    			"sh",
    			`${ROOT}/${id}`,
    		]);
    		const [state, value] = stdout.trim().split(" ");
    
    		if (state === "running") {
    			return { state, pid: Number(value) };
    		}
    
    		if (state === "exited") {
    			return { state, exitCode: Number(value) };
    		}
    
    		if (state === "starting" || state === "lost") {
    			return { state };
    		}
    
    		return undefined;
    	}
    }

    kill -0 checks that the process exists without sending it a signal. A process is lost when it ended without recording an exit code, for example because something sent SIGKILL to the script that records it.

    The process ID file also records the boot ID, which changes in every instance. A snapshot restores the file, and a new instance reuses the same process IDs, so a restored process ID can belong to an unrelated process. current() ignores a process ID from another instance, so a process that was running when the snapshot was taken is lost.

  4. Add methods that return the output of a process so far, and follow it while the process runs:

    src/index.tsts
    export class MyContainer extends DurableObject<Env> {
    	// ...
    
    	async logs(id: string, stream: "stdout" | "stderr"): Promise<string | undefined> {
    		if (!(await this.status(id))) {
    			return undefined;
    		}
    
    		const path = `${ROOT}/${id}/${stream}.log`;
    		return (await run(this.ctx.container!, ["cat", path])).stdout;
    	}
    
    	async follow(id: string, stream: "stdout" | "stderr"): Promise<Response | undefined> {
    		const status = await this.status(id);
    
    		if (!status) {
    			return undefined;
    		}
    
    		const path = `${ROOT}/${id}/${stream}.log`;
    		const argv =
    			status.state === "running"
    				? ["tail", "-n", "+1", "-F", "--pid", `${status.pid}`, path]
    				: ["cat", path];
    		const process = await this.ctx.container!.exec(argv, { stderr: "ignore" });
    		const { readable, writable } = new TransformStream<Uint8Array, Uint8Array>();
    		const writer = writable.getWriter();
    		const encoder = new TextEncoder();
    
    		const send = (event: string, data: unknown) =>
    			writer.write(
    				encoder.encode(`event: ${event}\ndata: ${JSON.stringify(data)}\n\n`),
    			);
    
    		let exited = false;
    		const markExited = () => {
    			exited = true;
    		};
    		process.exitCode.then(markExited, markExited);
    
    		// Signaling a process that has exited records an internal error.
    		const stop = () => {
    			if (!exited) process.kill();
    		};
    
    		// Stop following when the client disconnects.
    		writer.closed.catch(stop);
    
    		// Detect a disconnected client while the log is quiet.
    		const heartbeat = setInterval(() => {
    			writer.write(encoder.encode(": keep-alive\n\n")).catch(() => {});
    		}, 5_000);
    
    		const pump = async () => {
    			try {
    				for await (const text of process.stdout!.pipeThrough(
    					new TextDecoderStream(),
    				)) {
    					await send(stream, text);
    				}
    				const ended = await this.status(id);
    				const exitCode = ended?.state === "exited" ? ended.exitCode : null;
    				await send("exit", { exitCode });
    				await writer.close();
    			} catch {
    				stop();
    			} finally {
    				clearInterval(heartbeat);
    			}
    		};
    
    		void pump();
    
    		return new Response(readable, {
    			headers: {
    				"Content-Type": "text/event-stream",
    				"Cache-Control": "no-cache",
    			},
    		});
    	}
    }

    follow() sends the output so far, then new output until the process exits and tail --pid stops. Each event is named after the stream, and holds a JSON-encoded chunk of output, which can be part of a line or several lines. When the output ends, follow() sends an exit event with the exit code, or null when the process recorded none. A page can read the stream with EventSource, as in Stream command output.

    follow() notices a disconnected client only when a write fails. The keep-alive comment every five seconds makes that write happen while the log is quiet, so follow() stops tail. Without the comment, tail keeps running until the process exits.

  5. Add a method that stops a process:

    src/index.tsts
    export class MyContainer extends DurableObject<Env> {
    	// ...
    
    	async stop(id: string): Promise<boolean> {
    		const status = await this.status(id);
    
    		if (status?.state !== "running") {
    			return false;
    		}
    
    		// `kill` is a shell built-in. A negative ID signals the process group.
    		const { exitCode } = await run(this.ctx.container!, [
    			"sh",
    			"-c",
    			'kill -s TERM -- "$1"',
    			"sh",
    			`-${status.pid}`,
    		]);
    		return exitCode === 0;
    	}
    }

    A process stopped by a signal exits with 128 plus the signal number, such as 143 for SIGTERM.

  6. Add an alarm handler that keeps the sandbox running while any process runs:

    src/index.tsts
    export class MyContainer extends DurableObject<Env> {
    	// ...
    
    	async alarm(): Promise<void> {
    		let running = false;
    
    		for (const [key, argv] of await this.ctx.storage.list<string[]>({
    			prefix: "process:",
    		})) {
    			const id = key.slice("process:".length);
    			const status = await this.status(id);
    
    			if (status?.state === "running" || status?.state === "starting") {
    				running = true;
    			} else {
    				console.log({ event: "process.ended", id, argv, status });
    				await this.ctx.storage.delete(key);
    			}
    		}
    
    		if (running) {
    			await this.ctx.storage.setAlarm(Date.now() + 60_000);
    		}
    	}
    }

    A running process does not keep the container running. The alarm checks every process each minute, and each alarm keeps the container running. A Durable Object that restarted between alarms sets the 10-minute inactivity timeout again in its constructor. The container stops 10 minutes after the alarm finds that the last process has ended. For more information, refer to Sandbox lifetime.

    A Durable Object has one alarm. If your class already has an alarm() handler, merge the process checks into it, and set the alarm to the earliest time that any check needs. For more information, refer to Alarms.

    When a process ends, the alarm writes a process.ended log to Workers Logs within a minute. Replace the console.log() call with what your application does next, such as sending a notification.

  7. Add routes to your Worker that call these methods:

    src/index.jsjs
    const command = [
    	"sh",
    	"-c",
    	"echo ready; for i in $(seq 1 300); do echo tick $i; sleep 1; done",
    ];
    
    export default {
    	async fetch(request, env) {
    		const url = new URL(request.url);
    		const match = /^\/processes\/([a-z0-9-]{1,63})(\/logs)?$/.exec(
    			url.pathname,
    		);
    
    		if (!match) {
    			return new Response("Not found", { status: 404 });
    		}
    
    		const [, id, logs] = match;
    		const sandbox = env.MY_CONTAINER.getByName("sandbox");
    
    		if (logs) {
    			const stream =
    				url.searchParams.get("stream") === "stderr" ? "stderr" : "stdout";
    
    			if (url.searchParams.has("follow")) {
    				const response = await sandbox.follow(id, stream);
    				return response ?? new Response("Not found", { status: 404 });
    			}
    
    			const output = await sandbox.logs(id, stream);
    			return output === undefined
    				? new Response("Not found", { status: 404 })
    				: new Response(output);
    		}
    
    		if (request.method === "POST") {
    			return (await sandbox.startProcess(id, command))
    				? new Response(null, { status: 201 })
    				: new Response("The process ID is in use", { status: 409 });
    		}
    
    		if (request.method === "DELETE") {
    			return (await sandbox.stop(id))
    				? new Response(null, { status: 202 })
    				: new Response("The process is not running", { status: 409 });
    		}
    
    		const status = await sandbox.status(id);
    		return status
    			? Response.json(status)
    			: new Response("Not found", { status: 404 });
    	},
    };
    src/index.tsts
    const command = [
    	"sh",
    	"-c",
    	"echo ready; for i in $(seq 1 300); do echo tick $i; sleep 1; done",
    ];
    
    export default {
    	async fetch(request: Request, env: Env): Promise<Response> {
    		const url = new URL(request.url);
    		const match = /^\/processes\/([a-z0-9-]{1,63})(\/logs)?$/.exec(
    			url.pathname,
    		);
    
    		if (!match) {
    			return new Response("Not found", { status: 404 });
    		}
    
    		const [, id, logs] = match;
    		const sandbox = env.MY_CONTAINER.getByName("sandbox");
    
    		if (logs) {
    			const stream =
    				url.searchParams.get("stream") === "stderr" ? "stderr" : "stdout";
    
    			if (url.searchParams.has("follow")) {
    				const response = await sandbox.follow(id, stream);
    				return response ?? new Response("Not found", { status: 404 });
    			}
    
    			const output = await sandbox.logs(id, stream);
    			return output === undefined
    				? new Response("Not found", { status: 404 })
    				: new Response(output);
    		}
    
    		if (request.method === "POST") {
    			return (await sandbox.startProcess(id, command))
    				? new Response(null, { status: 201 })
    				: new Response("The process ID is in use", { status: 409 });
    		}
    
    		if (request.method === "DELETE") {
    			return (await sandbox.stop(id))
    				? new Response(null, { status: 202 })
    				: new Response("The process is not running", { status: 409 });
    		}
    
    		const status = await sandbox.status(id);
    		return status
    			? Response.json(status)
    			: new Response("Not found", { status: 404 });
    	},
    } satisfies ExportedHandler<Env>;

    This example command prints ready, then prints a line each second for five minutes. Replace it with your own command. The route accepts only process IDs that are safe in a path, because each ID becomes a directory name.

  8. Deploy your Worker:

    npx wrangler deploy
  9. Start a process named ticker, then check it. Replace the example hostname with the workers.dev URL that Wrangler prints:

    curl https://<YOUR_WORKER>.<YOUR_SUBDOMAIN>.workers.dev/processes/ticker --request POST
    curl https://<YOUR_WORKER>.<YOUR_SUBDOMAIN>.workers.dev/processes/ticker

    The first request responds with 201, and a second POST with the same ID responds with 409. The status includes the process ID:

    { "state": "running", "pid": 354 }
  10. Read the output so far, then follow it:

    curl https://<YOUR_WORKER>.<YOUR_SUBDOMAIN>.workers.dev/processes/ticker/logs
    curl "https://<YOUR_WORKER>.<YOUR_SUBDOMAIN>.workers.dev/processes/ticker/logs?follow" --no-buffer

    The first request prints the lines so far and ends. The second prints them as events, then a new event each second:

    event: stdout
    data: "ready\ntick 1\ntick 2\ntick 3\ntick 4\n"
    
    event: stdout
    data: "tick 5\n"
    
    event: stdout
    data: "tick 6\n"

    Add stream=stderr to the query to read standard error.

  11. In another terminal, stop the process, then check it again:

    curl https://<YOUR_WORKER>.<YOUR_SUBDOMAIN>.workers.dev/processes/ticker --request DELETE
    curl https://<YOUR_WORKER>.<YOUR_SUBDOMAIN>.workers.dev/processes/ticker

    The DELETE request responds with 202. The followed output ends with an exit event:

    event: exit
    data: {"exitCode":143}

    The status shows the same exit code:

    { "state": "exited", "exitCode": 143 }

Wait for a log line

A development server prints a line, such as ready, when it accepts requests. To wait for that line before you use the server, check the log files until a line matches or the process exits:

  1. Add a script that checks the logs of a process:

    src/index.jsjs
    // Prints the first matching line and exits with 0, or exits with 3 when
    // the process ends without printing one.
    const WAIT = `dir=$1; pattern=$2
    while :; do
      exited=false; [ -e "$dir/exit-code" ] && exited=true
      line=$(grep -h -m 1 -E -e "$pattern" "$dir/stdout.log" "$dir/stderr.log" 2>/dev/null | head -n 1)
      [ -n "$line" ] && { printf '%s\\n' "$line"; exit 0; }
      $exited && exit 3
      sleep 0.2
    done`;
    src/index.tsts
    // Prints the first matching line and exits with 0, or exits with 3 when
    // the process ends without printing one.
    const WAIT = `dir=$1; pattern=$2
    while :; do
      exited=false; [ -e "$dir/exit-code" ] && exited=true
      line=$(grep -h -m 1 -E -e "$pattern" "$dir/stdout.log" "$dir/stderr.log" 2>/dev/null | head -n 1)
      [ -n "$line" ] && { printf '%s\\n' "$line"; exit 0; }
      $exited && exit 3
      sleep 0.2
    done`;

    Do not pipe tail -F into grep -m 1 instead: tail keeps running after grep exits, until the next line arrives.

  2. Add a method that runs the script with a time limit:

    src/index.tsts
    export class MyContainer extends DurableObject<Env> {
    	// ...
    
    	async waitForLog(id: string, pattern: string, timeoutMs: number) {
    		const container = this.ctx.container;
    
    		if (!container?.running) {
    			return undefined;
    		}
    
    		const result = await run(container, [
    			"timeout",
    			"--kill-after=5",
    			`${timeoutMs / 1000}`,
    			"sh",
    			"-c",
    			WAIT,
    			"sh",
    			`${ROOT}/${id}`,
    			pattern,
    		]);
    
    		if (result.exitCode === 0) {
    			return { state: "matched", line: result.stdout.trimEnd() };
    		}
    
    		if (result.exitCode === 3) {
    			return { state: "exited" };
    		}
    
    		// 124 when the time limit passes, or 137 if timeout also sends SIGKILL.
    		if (result.exitCode === 124 || result.exitCode === 137) {
    			return { state: "timed-out" };
    		}
    
    		throw new Error(`The log check failed with exit code ${result.exitCode}`);
    	}
    }

    timeout stops the script when the time limit passes. If the script is still running five seconds later, --kill-after=5 stops it with SIGKILL.

  3. Call it after startProcess(), for example with the pattern ^ready:

    src/index.tsts
    const ready = await sandbox.waitForLog(id, "^ready", 30_000);
    { "state": "matched", "line": "ready" }

    The pattern is an extended regular expression. The state is exited when the process ends without printing a matching line, and timed-out when the time limit passes first.

Was this helpful?