Skip to content
A Bekkour
Writing

Deploys that survive a dropped stream or a control-plane restart

11 min read

On September 24 I shipped a small change to Railyard’s control plane, and the release restarted the job worker in the middle of a Chatwoot build. The build died on the server. The deploy row in Postgres stayed running forever, because the code that would have marked it finished was the code that had just been killed.

My fix that day was defensive. A recurring job now fails any deploy that has been silent for 30 minutes, and the control-plane release script waits for in-flight deploys to finish before it restarts the worker. That meant shipping my own code waited on whatever happened to be building, and a Rails build on a small server can take most of an hour. It also did nothing for the other ways a build died: a gRPC keepalive ping that went unanswered while the server was pegged by the build itself, or a short network blip between the agent and the control plane.

Yesterday I fixed the cause.

Stream lifetime and work lifetime

Railyard is two programs (the split is in two programs and a contract). The Rails control plane decides what to deploy; a Go agent on each server does it. A deploy was one gRPC call, Agent.Deploy, with a streamed response: the control plane sends a DeployRequest, and the agent streams back DeployEvent messages, one per line of build output, until a final KIND_COMPLETED with the exit code.

In Go, every gRPC handler gets a context, and the stream’s context is cancelled when the call ends for any reason: the client hung up, the deadline passed, the connection broke. Cancellation is how Go tells running work that nobody is waiting for it anymore. The agent ran every Docker command through that context with exec.CommandContext, which kills the process when the context is cancelled. That is usually what you want. If a user closes their browser, you stop the query you were running for them.

Here it glued two lifetimes together. The stream lives as long as a TCP connection and the processes on both ends stay healthy. The build should live as long as someone wants the deploy. A control-plane restart, a ping timeout and a network blip all looked the same to the agent: context cancelled, kill docker build.

A session that outlives the stream

I moved the build into a session that lives on the agent, keyed by an ID the control plane chooses. Three fields went into DeployRequest (session_id, resume, resume_after) and one into DeployEvent (seq).

Every event the build produces gets a sequence number, starting at 1, and goes into a buffer on the session. The stream that started the deploy becomes one follower of that buffer. If it drops, the build keeps running. When the control plane comes back, it calls Deploy again with the same session_id, resume = true, and resume_after set to the last sequence number it stored. The agent replays everything after that number and then follows live output.

Here is the agent’s handler, trimmed:

func (s *Server) deployFollowed(req *pb.DeployRequest, stream pb.Agent_DeployServer) error {
	d, fresh := registry.open(req.SessionId, !req.Resume)
	if d == nil {
		return status.Errorf(codes.NotFound,
			"deploy session %s is not on this agent (the agent restarted mid-deploy)", req.SessionId)
	}
	if fresh {
		ctx, cancel := context.WithCancel(context.WithoutCancel(stream.Context()))
		d.cancel = cancel
		go func() {
			defer cancel()
			d.finish(s.runDeploy(req, &detachedStream{Agent_DeployServer: stream, ctx: ctx, session: d}))
			time.AfterFunc(sessionRetention, func() { registry.forget(req.SessionId, d) })
		}()
	}
	return d.follow(stream, req.ResumeAfter)
}

The first detail is context.WithoutCancel, added in Go 1.21. It returns a context that keeps the parent’s values but drops its cancellation and its deadline. In a gRPC handler, the values include the incoming metadata, which is where the control plane’s auth token travels. So the build gets a context that still knows who asked for it, but no longer dies when the stream does. It gets its own cancel instead, which the session holds.

The second is detachedStream. The existing deploy code, which clones, builds, releases and swaps containers, takes a stream and calls Send on it. detachedStream wraps the original stream and overrides two methods: Context() returns the detached context, and Send appends the event to the session buffer instead of writing to the network. The deploy code itself did not change.

That leaves the question of who cancels a build nobody wants. When the last follower detaches, the session starts a timer. If no one reattaches within the grace period (10 minutes; 5 in the first version), the agent cancels the build. A finished session is kept for 15 minutes so a late reattach can still read the result, then forgotten. The buffer is capped at 50,000 events, oldest dropped first.

The control plane’s half

On the Rails side, a deploy can have several steps that run on agents: building on a dedicated build server, running on the app’s server, extra placements on other servers. Each step has a name like build-7 or app-3. I added a JSONB column, step_progress, to the deploys table, holding each step’s session ID and the last sequence number stored.

When the worker restarts, the deploy job is still a row in Solid Queue’s tables, so it runs again. DeployJob sees that step_progress is already filled in, logs “Control plane restarted, picking this deploy back up”, skips the start-of-deploy work like the “deploy started” webhook, and walks the steps again. A step with a saved result returns it straight away. A step with a saved session resumes:

def stream_step(channel, request, step:)
  saved = @deploy.step_progress[step] || {}
  return saved["result"] if saved["result"]

  seq = saved["seq"].to_i
  request.session_id = saved["session"] || new_session_id
  request.resume = saved["session"].present?
  request.resume_after = seq
  supersede_steps(step) unless request.resume
  begin
    channel.stub.deploy(request, metadata: channel.auth_metadata).each do |event|
      record_event(event)
      next unless event.seq.positive?
      seq = event.seq
      save_step(step, "session" => request.session_id, "seq" => seq)
    end
  rescue GRPC::NotFound
    raise unless request.resume
    log_system("The agent no longer has this step running, starting it again")
    request.session_id, request.resume, request.resume_after = new_session_id, false, 0
    retry
  end
  # ... save and return [succeeded, exit_code, stderr_tail] ...
end

The NotFound branch handles the case the session cannot: the agent itself restarted (an update, a crash, a reboot) and its in-memory sessions are gone. Then the control plane starts that one step again under a new session ID. The new build reuses whatever layers Docker already cached.

With that in place, the release script stopped waiting. It used to poll every 30 seconds, for up to 90 minutes, until no deploy was active before restarting the worker. Now it restarts the worker straight away, and running deploys are picked back up from the agents.

Never run a deploy twice

A resume must never start a new build, because deploys are not safe to repeat blindly. A second build of the same commit burns a build’s worth of CPU on a small server; a second release can run migrations twice or race the first one for the container swap.

On the agent, that is the !req.Resume argument to registry.open. A fresh request may create a session. A resume may only find one. If the session is not there, the answer is NotFound, and starting over is a separate decision the control plane makes on purpose, with a new ID.

There was one subtler case. Protocol Buffers skip fields a program does not know. An agent that had not been updated yet would read a resume request, ignore session_id and resume as unknown, and run a brand-new deploy. So the control plane only reattaches after a dropped stream when it has seen at least one event with a sequence number. An old agent never sends one, so a dropped stream from an old agent fails the way it always did, and nothing runs twice.

This is the at-least-once versus exactly-once distinction in a small form. Over a network, a sender that retries until it gets an answer delivers each message at least once, and sometimes more. You cannot make delivery happen exactly once. You can make processing idempotent: give each message an identity (here, a session ID and a sequence number) and have the receiver skip what it has already applied.

I got the receiver side wrong in the first version. To save writes, the control plane stored the sequence number at most once a second, while it wrote each log line as it arrived. After a crash, the saved position could be up to a second behind the log lines already in the table, so a resume replayed those lines and the deploy log showed them twice. The same night I changed it to save the position with every event. A second bug from that night: a resume could reattach to a step that a newer attempt had already replaced, and follow a build nobody wanted anymore. Starting a step fresh now clears any unfinished steps of the same kind, which is the supersede_steps call above, so a resume can only find the current one.

Not holding a thread for the build

The deploy job still sat in a Solid Queue worker thread, reading the stream, for the length of the build. Production had a handful of job threads serving every queue. A few concurrent deploys meant the recurring jobs (the deploy reconciler, offline-server detection, app health checks, webhooks) waited behind them, for up to an hour.

The next step is built and on a branch, not deployed yet. The control plane sends the request with detach = true. The agent creates the session, answers with a single KIND_ACCEPTED event, and the RPC ends in about a second and a half. The agent then pushes events to a new RPC on the control plane, PushDeployEvents, in batches. It retries each batch until the control plane acknowledges it, and sends an empty batch every minute as a keep-alive during quiet stretches of a build. The deploy job ends right after ACCEPTED. When the batch holding KIND_COMPLETED arrives, the control plane enqueues the job again for the next step. Deploys also get their own worker with their own threads, so long work cannot starve the sweeps.

Since the agent retries, the control plane will see some batches twice. The receiving side does the dedup and the checkpoint under one row lock, in one transaction:

def apply_agent_events!(step_name, session:, events:)
  with_lock do
    data = step(step_name)
    next [ events.map(&:seq).max.to_i, true ] if data["session"] != session || data.key?("result") || finished?

    fresh = events.select { |event| event.seq.to_i.zero? || event.seq > data["seq"].to_i }
    store_agent_events(fresh, data)
    data["seq"] = [ data["seq"].to_i, *fresh.map(&:seq) ].max
    save_step(step_name, data)
    [ data["seq"], data.key?("result") ]
  end
end

The log rows (one insert_all per batch, not one INSERT per line) and the new position commit together or not at all, which closes the gap the once-a-second checkpoint had. The return value is the acknowledged sequence number plus a stop flag, which tells the agent to stop reporting a session the control plane no longer follows.

The context bug

When I ran the built agent inside Docker-in-Docker against the control plane’s real gRPC server, every detached build failed at once with Unauthenticated.

The deploy code checks the caller’s token at the top, reading it from the context’s gRPC metadata. In the refactor that added PushDeployEvents, I had rewritten the line that creates the detached context, and it now started from nothing:

// before: no cancellation, but also no metadata
ctx, cancel := context.WithTimeout(context.Background(), buildTimeout(req)+detachedMargin)

// after: outlives the RPC but keeps its metadata, so the run authorizes again
ctx, cancel := context.WithTimeout(
	context.WithoutCancel(stream.Context()),
	buildTimeout(req)+detachedMargin,
)

context.Background() is an empty root. No cancellation, which was the point, but also no values, so no metadata and no token. The first version had used context.WithoutCancel(stream.Context()) for exactly this reason, and I lost it while moving code around.

A context in Go carries two separate things, a cancellation signal and a bag of request-scoped values. Detaching work from a request means keeping the values and dropping the signal. Background() drops both. The regression test now starts a detached deploy, cancels the caller’s context the way a finished RPC would, and checks that the build still passes its auth check 50 milliseconds later.

After the fix, the same Docker-in-Docker run went end to end: the deploy RPC returned in about 1.5 seconds, the agent pushed 66 events in batches, the step completed, and the control plane picked the deploy up and finished it.

The session buffer lives in the agent’s memory, and that is the remaining trade-off. If the agent restarts mid-build, the session is gone and that step starts over; the deploy page explains it when it happens. A disk-backed buffer would close that case at the cost of more moving parts on every server. For control-plane releases, the problem is gone: the worker restarts the moment new code lands, and the builds running at that moment carry on to the end.

Al Mokhtar Bekkour

Senior Rails & Go engineer in Quebec. I'm building Railyard, deploy software that runs your whole app on servers you own, and writing here about how it works. Open to work.

← All writing