Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 9 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@ server.

<picture>
<source media="(prefers-color-scheme: dark)" srcset="docs/media/deploy-dark.svg">
<img src="docs/media/deploy-light.svg" width="760" alt="An example Onebox session: ob plan prints a sealed diff of two image changes, a new workload and a migration whose data effect is unknown, then ob deploy rolls the workloads, verifies, and finishes with release r-0042 serving.">
<img src="docs/media/deploy-light.svg" width="760" alt="An example Onebox session: ob plan prints a sealed diff of two image changes, a new workload and a migration whose data effect is unknown, then ob deploy reports completed steps with durations and finishes with release r-0042 deployed.">
</picture>

<sub>A rendering of an example session, not a recording of one.</sub>
Expand Down Expand Up @@ -201,6 +201,14 @@ contract.

## Development

Preview the CLI's spinner, replica progress, and nested health/drain waits with
`go run ./scripts/ui-demo`; add `-fail` for a failed-healthcheck example. The
demo simulates a deployment locally and never connects to a server.

Interactive terminals show a live progress line. Pipes, CI, and `TERM=dumb`
receive static progress updates. Set `ONEBOX_NO_ANIMATION=1` to use static
output in a terminal; `NO_COLOR` disables colors independently.

```sh
just check # local pre-commit gate
just e2e # opt-in Docker end-to-end suite
Expand Down
13 changes: 9 additions & 4 deletions docs/media/deploy-dark.svg
Loading
Sorry, something went wrong. Reload?
Sorry, we cannot display this file.
Sorry, this file is invalid so it cannot be displayed.
13 changes: 9 additions & 4 deletions docs/media/deploy-light.svg
Loading
Sorry, something went wrong. Reload?
Sorry, we cannot display this file.
Sorry, this file is invalid so it cannot be displayed.
2 changes: 1 addition & 1 deletion go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ go 1.27.0

require (
github.com/charmbracelet/lipgloss v1.1.0
github.com/charmbracelet/x/ansi v0.11.8
github.com/compose-spec/compose-go/v2 v2.14.0
github.com/distribution/reference v0.6.0
github.com/opencontainers/go-digest v1.0.0
Expand All @@ -20,7 +21,6 @@ require (
require (
github.com/aymanbagabas/go-osc52/v2 v2.0.1 // indirect
github.com/charmbracelet/colorprofile v0.4.3 // indirect
github.com/charmbracelet/x/ansi v0.11.8 // indirect
github.com/charmbracelet/x/cellbuf v0.0.15 // indirect
github.com/charmbracelet/x/term v0.2.2 // indirect
github.com/clipperhouse/displaywidth v0.11.0 // indirect
Expand Down
8 changes: 5 additions & 3 deletions internal/engine/deploy.go
Original file line number Diff line number Diff line change
Expand Up @@ -274,7 +274,7 @@ func (e *Engine) runPhases(ctx context.Context, jw *journal.Writer, releaseID, l
if done["transfer"] {
e.logf("transfer: already complete (resume)")
} else {
tr := e.ui.Step("transfer", false)
tr := e.ui.Step("transfer", true)
if localStagingDir == "" {
res, err := e.T.Run(ctx, "test -d "+q(remoteDir))
if err != nil || res.ExitCode != 0 {
Expand Down Expand Up @@ -414,7 +414,7 @@ func (e *Engine) runPhases(ctx context.Context, jw *journal.Writer, releaseID, l
}

e.progress("verification", "started", "")
vf := e.ui.Step("verify", false)
vf := e.ui.Step("verify", true)
if err := e.Verify(ctx); err != nil {
vf(err)
e.progress("verification", "failed", "verification failed; inspect journal evidence")
Expand All @@ -430,7 +430,7 @@ func (e *Engine) runPhases(ctx context.Context, jw *journal.Writer, releaseID, l
e.progress("verification", "succeeded", "")

e.progress("activation", "started", "")
fin := e.ui.Step("activate", false)
fin := e.ui.Step("activate", true)
if err := jw.Append(ctx, journal.Record{
Phase: "activation", Event: "intent", Detail: "release=" + releaseID,
}); err != nil {
Expand Down Expand Up @@ -719,12 +719,14 @@ func (e *Engine) releaseRoles(ctx context.Context, remoteCompose string) error {
for _, roleName := range e.Spec.ReleaseOrder() {
role := e.Spec.Workloads[roleName]
e.logf("release %s (%s)", roleName, role.Mode())
step := e.ui.Step(roleName+" "+role.Mode(), true)
var err error
if role.Mode() == "rolling" {
err = e.RollRole(ctx, roleName, remoteCompose)
} else {
err = e.RecreateRole(ctx, roleName, remoteCompose)
}
step(err)
if err != nil {
return fmt.Errorf("%s: %w", roleName, err)
}
Expand Down
2 changes: 2 additions & 0 deletions internal/engine/recovery.go
Original file line number Diff line number Diff line change
Expand Up @@ -250,12 +250,14 @@ func (e *Engine) restoreReleaseRoles(ctx context.Context, previous *Engine, prev
continue
}
}
step := previous.ui.Step(roleName+" "+role.Mode(), true)
var err error
if role.Mode() == "rolling" {
err = previous.RollRole(ctx, roleName, composePath)
} else {
err = previous.RecreateRole(ctx, roleName, composePath)
}
step(err)
if err != nil {
return fmt.Errorf("restore %s: %w", roleName, err)
}
Expand Down
6 changes: 5 additions & 1 deletion internal/engine/recovery_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -81,8 +81,9 @@ func TestRecoveryRetryKeepsCheckpointUntilHealthyAndSweepsStaleRoles(t *testing.
}
return base(command)
}
var output bytes.Buffer
engine := New(testConfig(), testProject(t), target, Options{
Out: &bytes.Buffer{}, Sleep: noSleep,
Out: &output, Sleep: noSleep,
Now: func() time.Time { return time.Date(2026, 8, 9, 12, 0, 0, 0, time.UTC) },
})
seedInterruptedRecoveryState(t, engine)
Expand All @@ -99,6 +100,9 @@ func TestRecoveryRetryKeepsCheckpointUntilHealthyAndSweepsStaleRoles(t *testing.
if !errors.As(err, &incomplete) || incomplete.Code() != "recovery_incomplete" || incomplete.Phase != "verify" {
t.Fatalf("first recovery error = %#v", err)
}
if strings.Count(output.String(), "✓ web rolling") != 1 {
t.Fatalf("restored workload must leave one completed step before verification: %s", output.String())
}
if _, err := release.ReadActivationCheckpoint(context.Background(), target, engine.Names()); err != nil {
t.Fatalf("failed recovery cleared checkpoint: %v", err)
}
Expand Down
18 changes: 17 additions & 1 deletion internal/engine/roll.go
Original file line number Diff line number Diff line change
Expand Up @@ -135,6 +135,8 @@ func (e *Engine) rollRoleForRelease(ctx context.Context, roleName, remoteCompose
cc := e.composeCmdForProject(remoteComposePath, remoteProjectDir)
desired := role.Count()
within, pollEvery := role.ReadyTiming()
update, stop := e.ui.Progress(roleName+" rolling", desired)
defer stop()

pulled := false
// Each pass converges by one step: add a missing new replica, or retire a
Expand All @@ -149,6 +151,11 @@ func (e *Engine) rollRoleForRelease(ctx context.Context, roleName, remoteCompose
return err
}
olds := subtract(cur, news)
// A running newcomer alone is not a completed replacement: its old
// replica must have drained and retired too. Reserve the final segment
// until slot assignment has succeeded, including when resuming a roll.
completed := min(len(news), max(0, desired-len(olds)), max(0, desired-1))
update(completed, "inspecting replicas")
if guard > 4*(desired+len(olds))+8 {
return fmt.Errorf("role %s: roll did not converge (news=%d olds=%d)", roleName, len(news), len(olds))
}
Expand All @@ -158,6 +165,7 @@ func (e *Engine) rollRoleForRelease(ctx context.Context, roleName, remoteCompose
break
}
// surplus new replicas (count reduced) — drain one down to desired
update(completed, "retiring surplus replica")
if err := e.retireContainer(ctx, role, news[desired], pollEvery); err != nil {
return err
}
Expand All @@ -167,6 +175,7 @@ func (e *Engine) rollRoleForRelease(ctx context.Context, roleName, remoteCompose
// surge one new replica if we still need more of the new release
if len(news) < desired {
if !pulled {
update(completed, "checking image")
if err := e.pullBeforeRelease(ctx, svc, cc); err != nil {
return err
}
Expand All @@ -193,6 +202,7 @@ func (e *Engine) rollRoleForRelease(ctx context.Context, roleName, remoteCompose
}
known := idSet(news)
scale := len(cur) + 1
update(completed, fmt.Sprintf("starting replica %d", len(news)+1))
if res, err := e.mutate(ctx, fmt.Sprintf("%s up -d --no-deps --no-recreate --scale %s=%d %s", cc, svc, scale, svc)); err != nil {
return err
} else if res.ExitCode != 0 {
Expand All @@ -218,6 +228,7 @@ func (e *Engine) rollRoleForRelease(ctx context.Context, roleName, remoteCompose
return err
}
// join: the newcomer becomes a routable endpoint via its healthcheck
update(completed, fmt.Sprintf("replica %d healthcheck", len(news)+1))
if err := e.waitHealth(ctx, newID, "healthy", within, pollEvery); err != nil {
e.logf("join failed for %s — removing new container, existing keep serving", roleName)
cleanupErr := e.mutateChecked(ctx, "remove unhealthy newcomer "+newID, "docker rm -f "+newID)
Expand All @@ -227,19 +238,22 @@ func (e *Engine) rollRoleForRelease(ctx context.Context, roleName, remoteCompose

// retire one old, freeing its slot for the newcomer just added
if len(olds) > 0 {
update(completed, "retiring old replica")
if err := e.retireContainer(ctx, role, olds[0], pollEvery); err != nil {
return err
}
}

// hand clean slot names to any new container that doesn't have one yet
update(completed, "assigning replica slots")
if err := e.reslot(ctx, svc, releaseID, generation, desired); err != nil {
return err
}
}
if err := e.reslot(ctx, svc, releaseID, generation, desired); err != nil {
return err
}
update(desired, "replicas converged")
return nil
}

Expand Down Expand Up @@ -477,7 +491,8 @@ func subtract(all, remove []string) []string {

func (e *Engine) waitHealth(ctx context.Context, id, want string, budget, interval time.Duration) error {
deadline := e.Opts.Now().Add(budget)
_, stop := e.ui.Busy(fmt.Sprintf("waiting %.12s → %s", id, want))
label := fmt.Sprintf("waiting %.12s → %s", id, want)
update, stop := e.ui.Busy(label)
defer stop()
for {
h, err := e.healthOf(ctx, id)
Expand All @@ -487,6 +502,7 @@ func (e *Engine) waitHealth(ctx context.Context, id, want string, budget, interv
if h == want {
return nil
}
update(label + " · " + h)
// A container that has exited will not report a different health
// status later, so waiting out the budget only delays the same answer
// with a worse message.
Expand Down
42 changes: 42 additions & 0 deletions internal/engine/roll_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -269,6 +269,48 @@ func TestRollRoleTwoReplicasCleanSlots(t *testing.T) {
}
}

func TestRollProgressCompletesOnlyAfterRetirementAndSlots(t *testing.T) {
f := rollFake()
var out bytes.Buffer
base := f.Dynamic
f.Dynamic = func(cmd string) (transport.Result, bool) {
if strings.HasPrefix(cmd, "docker stop ") || strings.HasPrefix(cmd, "docker rename ") {
if strings.Contains(out.String(), "web rolling · 1/1") {
t.Fatalf("progress completed before retirement/slot assignment: %s", out.String())
}
}
return base(cmd)
}
e := New(testConfig(), testProject(t), f, Options{Out: &out, Sleep: noSleep})
if err := e.RollRole(context.Background(), "web", "/var/lib/onebox/app/releases/R1/compose.yaml"); err != nil {
t.Fatal(err)
}
if !strings.Contains(out.String(), "web rolling · 1/1 · replicas converged") {
t.Fatalf("successful roll must finish measured progress: %s", out.String())
}
}

func TestRollProgressDoesNotCompleteOnFailedHealthcheck(t *testing.T) {
f := rollFake()
base := f.Dynamic
f.Dynamic = func(cmd string) (transport.Result, bool) {
if strings.Contains(cmd, "State.Health") && strings.Contains(cmd, "NEW1") {
return transport.Result{Stdout: "starting\n"}, true
}
return base(cmd)
}
cfg := testConfig()
cfg.Workloads["web"] = withinMillis(cfg.Workloads["web"], 1)
var out bytes.Buffer
e := New(cfg, testProject(t), f, Options{Out: &out, Sleep: noSleep})
if err := e.RollRole(context.Background(), "web", "F"); err == nil {
t.Fatal("expected failed healthcheck")
}
if strings.Contains(out.String(), "web rolling · 1/1") || strings.Contains(out.String(), "replicas converged") {
t.Fatalf("failed roll must not show full progress: %s", out.String())
}
}

// drain.grace sets the docker stop -t timeout when retiring a drained container;
// absent it stays at the conservative 30s (asserted by the sequence tests).
func TestRollRoleDrainGraceConfigurable(t *testing.T) {
Expand Down
3 changes: 3 additions & 0 deletions internal/engine/rollback_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -82,6 +82,9 @@ func TestRollbackReplaysSnapshotChoreography(t *testing.T) {
if err := e.Rollback(context.Background()); err != nil {
t.Fatalf("rollback: %v\n%s", err, strings.Join(f.Commands, "\n"))
}
if strings.Count(out.String(), "✓ worker recreate") != 1 {
t.Fatalf("rollback must leave one completed workload step: %s", out.String())
}
seq := strings.Join(f.Commands, "\n")
if !strings.Contains(seq, "--force-recreate --timeout 30 worker") {
t.Fatalf("snapshot choreography (worker recreate) not replayed:\n%s", seq)
Expand Down
6 changes: 6 additions & 0 deletions internal/engine/secret_generation_rolling_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -226,6 +226,9 @@ func TestSecretGenerationRollsRollingWorkloads(t *testing.T) {
if _, err := engine.SecretsPushBatch(context.Background(), generationPayloads()); err != nil {
t.Fatalf("push: %v\n%s", err, strings.Join(fake.Commands, "\n"))
}
if strings.Count(output.String(), "✓ web rolling") != 1 {
t.Fatalf("secret roll must leave one completed workload step: %s", output.String())
}
commands := strings.Join(fake.Commands, "\n")
webUp := upCommandsFor(fake.Commands, "web")
if len(webUp) == 0 {
Expand Down Expand Up @@ -263,6 +266,9 @@ func TestSecretGenerationRollingKeepsServingReplicaWhenNewcomerNeverHealthy(t *t
if _, err := engine.SecretsPushBatch(context.Background(), generationPayloads()); err == nil {
t.Fatalf("unhealthy newcomer must fail the push:\n%s", strings.Join(fake.Commands, "\n"))
}
if !strings.Contains(output.String(), "✗ web rolling") || strings.Contains(output.String(), "✓ web rolling") {
t.Fatalf("failed secret roll must leave a failed workload step: %s", output.String())
}
commands := strings.Join(fake.Commands, "\n")
if !strings.Contains(commands, "docker rm -f W2") {
t.Fatalf("unhealthy newcomer was not removed:\n%s", commands)
Expand Down
4 changes: 3 additions & 1 deletion internal/engine/secretspush.go
Original file line number Diff line number Diff line change
Expand Up @@ -740,7 +740,7 @@ func (e *Engine) cleanupOrphanSecretGenerations(ctx context.Context, releaseID s
return e.mutateChecked(ctx, "clean orphaned secret generations", command)
}

func (e *Engine) forceSecretGeneration(ctx context.Context, checkpoint release.SecretCheckpoint, workload, generation string) error {
func (e *Engine) forceSecretGeneration(ctx context.Context, checkpoint release.SecretCheckpoint, workload, generation string) (err error) {
// Already converged. Two paths reach here that way: a resume after a crash,
// and recovery after a roll whose unhealthy newcomer was removed without
// any old replica being touched. Neither has anything to replace, and
Expand All @@ -760,6 +760,8 @@ func (e *Engine) forceSecretGeneration(ctx context.Context, checkpoint release.S
return err
}
composePath := generationDir + "/compose.yaml"
step := e.ui.Step(workload+" "+e.Spec.Workloads[workload].Mode(), true)
defer func() { step(err) }()
// A rolling workload rotates its secret the way it takes a release: surge
// one replica on the new generation, gate it healthy, retire one old.
// Recreating instead destroyed every serving replica before the first
Expand Down
Loading
Loading