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
74 changes: 16 additions & 58 deletions internal/cmd/enqueue.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,17 +4,14 @@
package cmd

import (
"errors"
"fmt"
"log/slog"
"time"

"github.com/dagucloud/dagu/internal/cmn/logger"
"github.com/dagucloud/dagu/internal/cmn/logger/tag"
"github.com/dagucloud/dagu/internal/cmn/stringutil"
"github.com/dagucloud/dagu/internal/core"
"github.com/dagucloud/dagu/internal/core/exec"
"github.com/dagucloud/dagu/internal/runtime/transform"
"github.com/dagucloud/dagu/internal/dagrun/intake"
"github.com/spf13/cobra"
)

Expand Down Expand Up @@ -96,68 +93,29 @@ func enqueueDAGRun(ctx *Context, dag *core.DAG, dagRunID string, triggerType cor
return fmt.Errorf("queues are disabled in configuration")
}

logFile, err := ctx.GenLogFileName(dag, dagRunID)
if err != nil {
return fmt.Errorf("failed to generate log file name: %w", err)
}

dagRun := exec.NewDAGRunRef(dag.Name, dagRunID)

if _, err = ctx.DAGRunStore.FindAttempt(ctx, dagRun); err == nil {
if _, err := ctx.DAGRunStore.FindAttempt(ctx, dagRun); err == nil {
return fmt.Errorf("DAG %q with ID %q already exists", dag.Name, dagRunID)
}
artifactDir, err := ctx.GenArtifactDir(dag, dagRunID)
if err != nil {
return fmt.Errorf("failed to generate artifact directory: %w", err)
}

att, err := ctx.DAGRunStore.CreateAttempt(ctx.Context, dag, time.Now(), dagRunID, exec.NewDAGRunAttemptOptions{})
queued, err := intake.EnqueueRun(ctx.Context, intake.QueueRequest{
DAGRunStore: ctx.DAGRunStore,
QueueStore: ctx.QueueStore,
DAG: dag,
DAGRunID: dagRunID,
LogBaseDir: ctx.Config.Paths.LogDir,
ArtifactBaseDir: ctx.Config.Paths.ArtifactDir,
TriggerType: triggerType,
ScheduleTime: scheduleTime,
ProceedOnStatusCloseErr: true,
})
if err != nil {
return fmt.Errorf("failed to create run: %w", err)
}

opts := []transform.StatusOption{
transform.WithLogFilePath(logFile),
transform.WithArchiveDir(artifactDir),
transform.WithAttemptID(att.ID()),
transform.WithPreconditions(dag.Preconditions),
transform.WithQueuedAt(stringutil.FormatTime(time.Now())),
transform.WithHierarchyRefs(
exec.NewDAGRunRef(dag.Name, dagRunID),
exec.DAGRunRef{},
),
transform.WithTriggerType(triggerType),
}

if scheduleTime != "" {
opts = append(opts, transform.WithScheduleTime(scheduleTime))
}

dagStatus := transform.NewStatusBuilder(dag).Create(dagRunID, core.Queued, 0, time.Time{}, opts...)

if err := att.Open(ctx.Context); err != nil {
return fmt.Errorf("failed to open run: %w", err)
}

if err := att.Write(ctx.Context, dagStatus); err != nil {
_ = att.Close(ctx.Context)
return fmt.Errorf("failed to save status: %w", err)
return err
}

closeErr := att.Close(ctx.Context)
if closeErr != nil {
if queued.StatusCloseErr != nil {
logger.Warn(ctx.Context, "Failed to close queued status before enqueue",
tag.Error(closeErr))
}

if err := ctx.QueueStore.Enqueue(ctx.Context, dag.ProcGroup(), exec.QueuePriorityLow, dagRun); err != nil {
if closeErr != nil {
return errors.Join(
fmt.Errorf("failed to close run: %w", closeErr),
fmt.Errorf("failed to enqueue dag-run: %w", err),
)
}
return fmt.Errorf("failed to enqueue dag-run: %w", err)
tag.Error(queued.StatusCloseErr))
}

logger.Info(ctx.Context, "Enqueued dag-run",
Expand Down
161 changes: 13 additions & 148 deletions internal/cmd/local_execution.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,151 +5,14 @@ package cmd

import (
"context"
"errors"
"fmt"
"time"

"github.com/dagucloud/dagu/internal/cmn/logger"
"github.com/dagucloud/dagu/internal/cmn/logger/tag"
"github.com/dagucloud/dagu/internal/core"
"github.com/dagucloud/dagu/internal/core/exec"
"github.com/dagucloud/dagu/internal/runtime/transform"
"github.com/dagucloud/dagu/internal/dagrun/intake"
)

var errLocalExecutionAlreadyExists = errors.New("local execution already exists")

type localExecutionPreparation struct {
Attempt exec.DAGRunAttempt
Proc exec.ProcHandle
}

func prepareLocalExecution(
ctx *Context,
dag *core.DAG,
dagRunID string,
root exec.DAGRunRef,
parent exec.DAGRunRef,
triggerType core.TriggerType,
scheduleTime string,
buildAttempt func(context.Context) (exec.DAGRunAttempt, error),
) (*localExecutionPreparation, error) {
if dag == nil {
return nil, fmt.Errorf("dag is required")
}
if dagRunID == "" {
return nil, fmt.Errorf("dag-run ID is required")
}
if buildAttempt == nil {
return nil, fmt.Errorf("attempt builder is required")
}
if root.Zero() {
root = exec.NewDAGRunRef(dag.Name, dagRunID)
}

if err := ctx.ProcStore.Lock(ctx, dag.ProcGroup()); err != nil {
return nil, fmt.Errorf("failed to lock process group: %w", err)
}
defer ctx.ProcStore.Unlock(ctx, dag.ProcGroup())

attempt, err := buildAttempt(ctx.Context)
if err != nil {
if errors.Is(err, exec.ErrDAGRunAlreadyExists) {
return nil, fmt.Errorf("%w: dag-run ID %s already exists for DAG %s", errLocalExecutionAlreadyExists, dagRunID, dag.Name)
}
return nil, fmt.Errorf("failed to prepare execution attempt: %w", err)
}
if attempt == nil {
return nil, fmt.Errorf("attempt builder returned nil attempt")
}
attempt.SetDAG(dag)

proc, err := ctx.ProcStore.Acquire(ctx, dag.ProcGroup(), exec.ProcMeta{
StartedAt: time.Now().Unix(),
Name: dag.Name,
DAGRunID: dagRunID,
AttemptID: attempt.ID(),
RootName: root.Name,
RootDAGRunID: root.ID,
})
if err != nil {
_ = recordPreparedAttemptFailure(ctx, attempt, dag, dagRunID, root, parent, triggerType, scheduleTime, err)
return nil, fmt.Errorf("%w: %w", errProcAcquisitionFailed, err)
}

return &localExecutionPreparation{
Attempt: attempt,
Proc: proc,
}, nil
}

func recordPreparedAttemptFailure(
ctx *Context,
attempt exec.DAGRunAttempt,
dag *core.DAG,
dagRunID string,
root exec.DAGRunRef,
parent exec.DAGRunRef,
triggerType core.TriggerType,
scheduleTime string,
runErr error,
) error {
if attempt == nil {
return fmt.Errorf("attempt is required")
}
if dag == nil {
return fmt.Errorf("dag is required")
}
if dagRunID == "" {
return fmt.Errorf("dag-run ID is required")
}
if root.Zero() {
root = exec.NewDAGRunRef(dag.Name, dagRunID)
}

logPath, logPathErr := ctx.GenLogFileName(dag, dagRunID)
if logPathErr != nil {
logger.Warn(ctx, "Failed to generate log file path for prepared local execution failure",
tag.Error(logPathErr),
tag.DAG(dag.Name),
tag.RunID(dagRunID),
)
}
artifactDir, artifactDirErr := ctx.GenArtifactDir(dag, dagRunID)
if artifactDirErr != nil {
logger.Warn(ctx, "Failed to generate artifact directory for prepared local execution failure",
tag.Error(artifactDirErr),
tag.DAG(dag.Name),
tag.RunID(dagRunID),
)
}
opts := []transform.StatusOption{
transform.WithAttemptID(attempt.ID()),
transform.WithHierarchyRefs(root, parent),
transform.WithLogFilePath(logPath),
transform.WithArchiveDir(artifactDir),
transform.WithFinishedAt(time.Now()),
transform.WithError(runErr.Error()),
transform.WithWorkerID("local"),
transform.WithTriggerType(triggerType),
}
if scheduleTime != "" {
opts = append(opts, transform.WithScheduleTime(scheduleTime))
}
status := transform.NewStatusBuilder(dag).Create(dagRunID, core.Failed, 0, time.Now(), opts...)

if err := attempt.Open(ctx.Context); err != nil {
return fmt.Errorf("failed to open attempt for failure recording: %w", err)
}
defer func() {
_ = attempt.Close(ctx.Context)
}()

if err := attempt.Write(ctx.Context, status); err != nil {
return fmt.Errorf("failed to write failed status: %w", err)
}
return nil
}

func withPreparedLocalExecution(
ctx *Context,
dag *core.DAG,
Expand All @@ -161,16 +24,18 @@ func withPreparedLocalExecution(
buildAttempt func(context.Context) (exec.DAGRunAttempt, error),
run func(exec.DAGRunAttempt) error,
) error {
prepared, err := prepareLocalExecution(
ctx,
dag,
dagRunID,
root,
parent,
triggerType,
scheduleTime,
buildAttempt,
)
prepared, err := intake.PrepareLocalExecution(ctx.Context, intake.LocalRequest{
ProcStore: ctx.ProcStore,
DAG: dag,
DAGRunID: dagRunID,
Root: root,
Parent: parent,
TriggerType: triggerType,
ScheduleTime: scheduleTime,
LogBaseDir: ctx.Config.Paths.LogDir,
ArtifactBaseDir: ctx.Config.Paths.ArtifactDir,
BuildAttempt: buildAttempt,
})
if err != nil {
logger.Debug(ctx, "Failed to prepare local execution", tag.Error(err))
return err
Expand Down
2 changes: 0 additions & 2 deletions internal/cmd/start.go
Original file line number Diff line number Diff line change
Expand Up @@ -224,8 +224,6 @@ func runStart(ctx *Context, args []string) error {
return tryExecuteDAG(ctx, dag, dagRunID, root, workerID, attemptID, triggerType, scheduleTime)
}

var errProcAcquisitionFailed = errors.New("failed to acquire process handle")

// tryExecuteDAG acquires a process handle and executes the DAG.
func tryExecuteDAG(ctx *Context, dag *core.DAG, dagRunID string, root exec.DAGRunRef, workerID, attemptID string, triggerType core.TriggerType, scheduleTime string) error {
// Check for dispatch to coordinator for distributed execution.
Expand Down
Loading
Loading