Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
21 commits
Select commit Hold shift + click to select a range
791867d
feat(workspace): support async, durable workspace up execution
skevetter Jul 28, 2026
65ce40f
fix(lint): satisfy golangci-lint on the async workspace-up changes
skevetter Jul 28, 2026
5f6e271
fix(lint): drop forbidigo nolint by writing to os.Stdout directly
skevetter Jul 28, 2026
6078bae
fix(e2e): implement the async task protocol in the mock devsy CLI
skevetter Jul 28, 2026
d212862
style: update comments
skevetter Jul 29, 2026
5b95539
fix(task): preserve terminal task state and tighten status reporting
skevetter Jul 29, 2026
480427b
fix(desktop): keep detached up tasks cancellable
skevetter Jul 29, 2026
545ad00
refactor(status): move status package out of devcontainer
skevetter Jul 29, 2026
c7952dc
fix(task): reconcile abandoned tasks and close up/stop races
skevetter Jul 29, 2026
e742a9d
style: cleanup comments
skevetter Jul 30, 2026
eda77d1
style: fix formatting
skevetter Jul 30, 2026
deaeb7c
fix(task): use a worker-held lock for task liveness
skevetter Jul 30, 2026
6807530
fix(task): hold the worker lock across the abandoned-task transition
skevetter Jul 30, 2026
2b4ba67
style: update comments
skevetter Jul 30, 2026
0d202d4
fix(task): fsync task state so a crash can't lose a terminal result
skevetter Jul 30, 2026
a9694b9
style: update comments
skevetter Jul 30, 2026
416d00d
style: update podman variable
skevetter Jul 30, 2026
6bee286
fix: apply testfile naming convention
skevetter Jul 30, 2026
3f63b4d
fix(task): make the worker-lock test seam visible across packages
skevetter Jul 30, 2026
183d5e7
style: trim explanatory comments to essentials
skevetter Jul 30, 2026
4206641
fix(desktop): handle a log follower that fails to start
skevetter Jul 30, 2026
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
6 changes: 4 additions & 2 deletions cmd/internal/agentworkspace/up.go
Original file line number Diff line number Diff line change
Expand Up @@ -150,7 +150,7 @@ func (cmd *UpCmd) up(
workspaceInfo *provider.AgentWorkspaceInfo,
tunnelClient tunnel.TunnelClient,
) error {
result, err := cmd.devsyUp(ctx, workspaceInfo)
result, err := cmd.devsyUp(ctx, workspaceInfo, tunnelClient)
if err != nil {
errResult := &config2.Result{
Error: err.Error(),
Expand Down Expand Up @@ -200,16 +200,18 @@ func (cmd *UpCmd) sendResult(
func (cmd *UpCmd) devsyUp(
ctx context.Context,
workspaceInfo *provider.AgentWorkspaceInfo,
tunnelClient tunnel.TunnelClient,
) (*config2.Result, error) {
runner, err := CreateRunner(ctx, workspaceInfo)
if err != nil {
return nil, err
}

reporter := tunnelserver.NewTunnelStatusReporter(ctx, tunnelClient)
return runner.Up(ctx, devcontainer.UpOptions{
CLIOptions: workspaceInfo.CLIOptions,
RegistryCache: workspaceInfo.RegistryCache,
}, workspaceInfo.InjectTimeout)
}, workspaceInfo.InjectTimeout, reporter)
}

func CreateRunner(
Expand Down
2 changes: 2 additions & 0 deletions cmd/internal/container_tunnel.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ import (
"github.com/devsy-org/devsy/pkg/flags/names"
"github.com/devsy-org/devsy/pkg/log"
provider2 "github.com/devsy-org/devsy/pkg/provider"
"github.com/devsy-org/devsy/pkg/status"
"github.com/spf13/cobra"
)

Expand Down Expand Up @@ -161,6 +162,7 @@ func StartContainer(
ctx,
devcontainer.UpOptions{NoBuild: true},
workspaceConfig.InjectTimeout,
status.Nop(),
)
if err != nil {
return result, err
Expand Down
341 changes: 341 additions & 0 deletions cmd/workspace/task.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,341 @@
package workspace

import (
"context"
"encoding/json"
"fmt"
"os"
"time"

"github.com/devsy-org/devsy/cmd/flags"
config2 "github.com/devsy-org/devsy/pkg/devcontainer/config"
cliflags "github.com/devsy-org/devsy/pkg/flags"
"github.com/devsy-org/devsy/pkg/flags/names"
"github.com/devsy-org/devsy/pkg/output"
"github.com/devsy-org/devsy/pkg/status"
"github.com/devsy-org/devsy/pkg/task"
"github.com/spf13/cobra"
)

// NewTaskCmd builds the workspace task parent command.
func NewTaskCmd(globalFlags *flags.GlobalFlags) *cobra.Command {
taskCmd := &cobra.Command{
Use: "task",
Short: "Manage background tasks (e.g. from 'up --detach')",
}
taskCmd.AddCommand(newTaskListCmd(globalFlags))
taskCmd.AddCommand(newTaskGetCmd(globalFlags))
taskCmd.AddCommand(newTaskLogsCmd(globalFlags))
taskCmd.AddCommand(newTaskCancelCmd(globalFlags))
taskCmd.AddCommand(newTaskRmCmd(globalFlags))
return taskCmd
}

type taskListCmd struct {
*flags.GlobalFlags
}

func newTaskListCmd(globalFlags *flags.GlobalFlags) *cobra.Command {
cmd := &taskListCmd{GlobalFlags: globalFlags}
return &cobra.Command{
Use: "list",
Aliases: []string{"ls"},
Short: "List background tasks, most recently started first",
Args: cobra.NoArgs,
RunE: func(*cobra.Command, []string) error {
return cmd.run()
},
}
}

func (cmd *taskListCmd) run() error {
emitJSON, err := resolveEmitJSON(cmd.GlobalFlags)
if err != nil {
return err
}

store, err := task.NewStore()
if err != nil {
return err
}
states, err := store.List()
if err != nil {
return err
}

if emitJSON {
return json.NewEncoder(os.Stdout).Encode(states)
}
for _, s := range states {
_, _ = fmt.Fprintf(os.Stdout, "%s\t%s\t%s\t%s\n", s.ID, s.Status, s.Command, s.WorkspaceID)
}
return nil
}

type taskGetCmd struct {
*flags.GlobalFlags
}

func newTaskGetCmd(globalFlags *flags.GlobalFlags) *cobra.Command {
cmd := &taskGetCmd{GlobalFlags: globalFlags}
return &cobra.Command{
Use: "get <task-id>",
Aliases: []string{"describe", "show"},
Short: "Show a task's current status",
Args: cobra.ExactArgs(1),
RunE: func(_ *cobra.Command, args []string) error {
return cmd.run(args[0])
},
}
}

func (cmd *taskGetCmd) run(id string) error {
emitJSON, err := resolveEmitJSON(cmd.GlobalFlags)
if err != nil {
return err
}
store, err := task.NewStore()
if err != nil {
return err
}
state, err := store.Get(id)
if err != nil {
return err
}
return reportTaskState(state, emitJSON)
}

type taskLogsCmd struct {
*flags.GlobalFlags

Follow bool
Interval string
}

func newTaskLogsCmd(globalFlags *flags.GlobalFlags) *cobra.Command {
cmd := &taskLogsCmd{GlobalFlags: globalFlags}
logsCmd := &cobra.Command{
Use: "logs <task-id>",
Aliases: []string{"attach"},
Short: "Show a task's status, or follow it until it finishes",
Args: cobra.ExactArgs(1),
RunE: func(_ *cobra.Command, args []string) error {
return cmd.run(args[0])
},
}
cliflags.Add(logsCmd,
cliflags.Bool(&cmd.Follow, names.Follow, false,
"Poll until the task reaches a terminal state instead of reporting once").
Shorthand("f"),
cliflags.String(&cmd.Interval, names.Interval, "500ms",
"Poll interval when --follow is set"),
)
return logsCmd
}

func (cmd *taskLogsCmd) run(id string) error {
emitJSON, err := resolveEmitJSON(cmd.GlobalFlags)
if err != nil {
return err
}
store, err := task.NewStore()
if err != nil {
return err
}

if !cmd.Follow {
state, err := store.Get(id)
if err != nil {
return err
}
return reportTaskState(state, emitJSON)
}

interval, err := time.ParseDuration(cmd.Interval)
if err != nil {
return fmt.Errorf("parse --interval: %w", err)
}
if interval <= 0 {
return fmt.Errorf("--interval must be positive, got %q", cmd.Interval)
}
return followTask(context.Background(), store, followTaskOptions{
id: id,
interval: interval,
emitJSON: emitJSON,
})
}

type followTaskOptions struct {
id string
interval time.Duration
emitJSON bool
}

func followTask(ctx context.Context, store *task.Store, opts followTaskOptions) error {
var last *task.State
ticker := time.NewTicker(opts.interval)
defer ticker.Stop()

for {
state, err := store.Get(opts.id)
if err != nil {
return err
}
emitTaskTransition(last, state, opts.emitJSON)
last = state

if state.Status.Terminal() {
return reportTaskState(state, opts.emitJSON)
}

select {
case <-ctx.Done():
return ctx.Err()
case <-ticker.C:
}
}
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.

func emitTaskTransition(last, current *task.State, emitJSON bool) {
if last != nil && last.Phase == current.Phase && last.Step == current.Step {
return
}
if current.Phase == "" {
return
}
event := status.Event{Phase: status.Phase(current.Phase), Step: current.Step, Started: true}
if emitJSON {
_ = config2.WriteStatusJSON(os.Stdout, event)
return
}
_, _ = fmt.Fprintf(os.Stdout, "task %s: %s\n", current.ID, current.Phase)
}

type taskCancelCmd struct {
*flags.GlobalFlags
}

func newTaskCancelCmd(globalFlags *flags.GlobalFlags) *cobra.Command {
cmd := &taskCancelCmd{GlobalFlags: globalFlags}
return &cobra.Command{
Use: "cancel <task-id>",
Aliases: []string{"stop"},
Short: "Stop a task's process and mark it failed",
Args: cobra.ExactArgs(1),
RunE: func(_ *cobra.Command, args []string) error {
return cmd.run(args[0])
},
}
}

func (cmd *taskCancelCmd) run(id string) error {
emitJSON, err := resolveEmitJSON(cmd.GlobalFlags)
if err != nil {
return err
}
store, err := task.NewStore()
if err != nil {
return err
}
if err := store.Open(id).Cancel(); err != nil {
return err
}
state, err := store.Get(id)
if err != nil {
return err
}

// Report the raw state. Canceling is not itself a failure.
if emitJSON {
return json.NewEncoder(os.Stdout).Encode(state)
}
_, _ = fmt.Fprintf(os.Stdout, "task %s: canceled\n", state.ID)
return nil
}

type taskRmCmd struct {
*flags.GlobalFlags

Force bool
}

func newTaskRmCmd(globalFlags *flags.GlobalFlags) *cobra.Command {
cmd := &taskRmCmd{GlobalFlags: globalFlags}
rmCmd := &cobra.Command{
Use: "rm <task-id>",
Aliases: []string{"delete", "remove"},
Short: "Delete a finished task's state",
Args: cobra.ExactArgs(1),
RunE: func(_ *cobra.Command, args []string) error {
return cmd.run(args[0])
},
}
cliflags.Add(rmCmd,
cliflags.Bool(&cmd.Force, names.Force, false,
"Stop the task first if it's still pending or running, then delete"),
)
return rmCmd
}

func (cmd *taskRmCmd) run(id string) error {
store, err := task.NewStore()
if err != nil {
return err
}
if cmd.Force {
if err := store.Open(id).Cancel(); err != nil {
return err
}
}
return store.Delete(id, cmd.Force)
}

func resolveEmitJSON(g *flags.GlobalFlags) (bool, error) {
mode, err := output.ResolveMode(g.ResultFormat)
if err != nil {
return false, err
}
return mode == output.ModeJSON, nil
}

func reportTaskState(state *task.State, emitJSON bool) error {
if emitJSON {
return reportTaskStateJSON(state)
}

_, _ = fmt.Fprintf(os.Stdout, "task %s: %s\n", state.ID, state.Status)
if state.Status == task.StatusFailed {
return fmt.Errorf("%s", state.Error)
}
return nil
}

func reportTaskStateJSON(state *task.State) error {
switch state.Status {
case task.StatusFailed:
if err := config2.WriteErrorJSON(os.Stdout, state.Error); err != nil {
return err
}
return fmt.Errorf("%s", state.Error)
case task.StatusSucceeded:
return config2.WriteResultJSON(os.Stdout, resultEnvelopeFrom(state))
default:
return config2.WriteStatusJSON(os.Stdout, status.Event{
Phase: status.Phase(state.Phase),
Step: state.Step,
Started: true,
})
}
}

func resultEnvelopeFrom(state *task.State) config2.ResultEnvelope {
if state.Result == nil {
return config2.ResultEnvelope{}
}
return config2.ResultEnvelope{
ContainerID: config2.GetContainerID(state.Result),
RemoteUser: config2.GetRemoteUser(state.Result),
Warnings: state.Result.HostWarnings,
Recovery: state.Result.RecoveryContainer,
}
}
Loading
Loading