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
34 changes: 20 additions & 14 deletions conformance/harness/runner.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ import (
"testing"
"time"

"github.com/dagucloud/dagu/v2/internal/cmn/cmdutil"
"github.com/stretchr/testify/require"
)

Expand All @@ -36,13 +37,13 @@ type Result struct {

// Process is a running Dagu command.
type Process struct {
t *testing.T
cancel context.CancelFunc
done chan struct{}
stdout *bytes.Buffer
stderr *bytes.Buffer
err error
cancelOnce sync.Once
t *testing.T
proc *cmdutil.ManagedProcess
done chan struct{}
stdout *bytes.Buffer
stderr *bytes.Buffer
err error
stopOnce sync.Once
}

// ExitCode returns the command's process exit code.
Expand Down Expand Up @@ -99,29 +100,29 @@ func (r *Runner) RunWithEnv(env []string, args ...string) *Result {
func (r *Runner) StartWithEnv(env []string, args ...string) *Process {
r.t.Helper()

ctx, cancel := context.WithCancel(context.Background())
stdout := &bytes.Buffer{}
stderr := &bytes.Buffer{}
// Binary-level tests intentionally execute the configured Dagu binary.
cmd := exec.CommandContext(ctx, daguBinary(r.t), args...) //nolint:gosec
cmd := exec.Command(daguBinary(r.t), args...) //nolint:gosec
cmd.Dir = r.dir
cmd.Env = appendEnv(append(isolatedEnv(r.t), "PWD="+r.dir), env...)
cmd.Stdout = stdout
cmd.Stderr = stderr
if err := cmd.Start(); err != nil {
cancel()
proc, err := cmdutil.StartManagedProcess(cmd)
if err != nil {
r.t.Fatalf("starting dagu %s: %v", strings.Join(args, " "), err)
}

process := &Process{
t: r.t,
cancel: cancel,
proc: proc,
done: make(chan struct{}),
stdout: stdout,
stderr: stderr,
}
go func() {
process.err = cmd.Wait()
process.err = proc.Wait()
_ = proc.Release()
close(process.done)
}()
r.t.Cleanup(func() {
Expand Down Expand Up @@ -153,7 +154,12 @@ func (p *Process) Done() <-chan struct{} {
// Stop terminates the command and waits for it to exit.
func (p *Process) Stop() {
p.t.Helper()
p.cancelOnce.Do(p.cancel)
p.stopOnce.Do(func() {
_, _ = p.proc.Stop(cmdutil.StopRequest{
Intent: cmdutil.ForceTermination(),
Reason: cmdutil.StopReasonShutdown,
})
})
<-p.done
}

Expand Down
2 changes: 1 addition & 1 deletion internal/incident/model.go
Original file line number Diff line number Diff line change
Expand Up @@ -169,7 +169,7 @@ type Store interface {

SaveState(ctx context.Context, state *IncidentState) error
GetState(ctx context.Context, providerID, dedupKey string) (*IncidentState, error)
ListOpenStatesByDAG(ctx context.Context, dagName string) ([]*IncidentState, error)
ListOpenStates(ctx context.Context) ([]*IncidentState, error)
DeleteState(ctx context.Context, providerID, dedupKey string) error
}

Expand Down
14 changes: 12 additions & 2 deletions internal/intg/reschedule_source_file_api_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ import (

"github.com/dagucloud/dagu/v2/api/v1"
"github.com/dagucloud/dagu/v2/internal/cmn/config"
"github.com/dagucloud/dagu/v2/internal/cmn/fileutil"
"github.com/dagucloud/dagu/v2/internal/test"
"github.com/stretchr/testify/require"
)
Expand Down Expand Up @@ -48,7 +49,7 @@ steps:
attempt, dag := test.WaitForAttemptSnapshotWithDAG(t, server, dagName, enqBody.DagRunId)
require.NotNil(t, attempt)
require.Empty(t, dag.Location)
require.Equal(t, dagPath, dag.SourceFile)
requireSameFile(t, dagPath, dag.SourceFile)

assertQueuedRunSpecFromFile(t, server, dagName, enqBody.DagRunId, true)

Expand All @@ -73,7 +74,16 @@ steps:

_, rescheduledDAG := test.WaitForAttemptSnapshotWithDAG(t, server, dagName, rescheduleBody.DagRunId)
require.Contains(t, string(rescheduledDAG.YamlData), "echo current file")
require.Equal(t, dagPath, rescheduledDAG.SourceFile)
requireSameFile(t, dagPath, rescheduledDAG.SourceFile)
}

func requireSameFile(t *testing.T, expected, actual string) {
t.Helper()
expectedInfo, err := fileutil.Stat(expected)
require.NoError(t, err)
actualInfo, err := fileutil.Stat(actual)
require.NoError(t, err)
require.True(t, os.SameFile(expectedInfo, actualInfo), "%q and %q do not identify the same file", expected, actual)
}

func assertQueuedRunSpecFromFile(t *testing.T, server test.Server, dagName, dagRunID string, want bool) {
Expand Down
8 changes: 8 additions & 0 deletions internal/persis/file/dagrun/attempt.go
Original file line number Diff line number Diff line change
Expand Up @@ -172,6 +172,14 @@ func (att *Attempt) Open(ctx context.Context) error {
if err := fileutil.WriteFileAtomic(filepath.Join(dir, DAGDefinition), dagJSON, 0600); err != nil {
return fmt.Errorf("failed to write DAG definition: %w", err)
}
} else {
dag, err := att.ReadDAG(ctx)
switch {
case err == nil:
att.dag = dag
case !errors.Is(err, os.ErrNotExist):
return fmt.Errorf("failed to restore DAG definition: %w", err)
}
}

// Create the per-run work directory so steps can use DAG_RUN_WORK_DIR immediately
Expand Down
30 changes: 28 additions & 2 deletions internal/persis/file/dagrun/attempt_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,23 @@ func TestAttempt_Open(t *testing.T) {
assert.NoError(t, err)
}

func TestAttempt_OpenRejectsCorruptDAGDefinition(t *testing.T) {
dir := createTempDir(t)
file := filepath.Join(dir, "status.dat")
ctx := context.Background()

att, err := NewAttempt(file, nil, WithDAG(&core.DAG{Name: "test"}))
require.NoError(t, err)
require.NoError(t, att.Open(ctx))
require.NoError(t, att.Close(ctx))
require.NoError(t, os.WriteFile(filepath.Join(dir, DAGDefinition), []byte("{"), 0600))

reopened, err := NewAttempt(file, nil)
require.NoError(t, err)
err = reopened.Open(ctx)
require.ErrorContains(t, err, "failed to restore DAG definition")
}

func TestAttempt_Write(t *testing.T) {
dir := createTempDir(t)
file := filepath.Join(dir, "status.dat")
Expand Down Expand Up @@ -906,14 +923,19 @@ func TestAttempt_WriteEmitsLifecycleTransitionsAndStatusUpdates(t *testing.T) {
service := eventstore.New(store)
ctx := eventstore.WithContext(context.Background(), service, eventstore.Source{Service: eventstore.SourceServiceServer})

dag := &core.DAG{Name: "TestDAG", Location: filepath.Join(dir, "test-dag.yaml")}
dag := &core.DAG{
Name: "TestDAG",
Location: filepath.Join(dir, "test-dag.yaml"),
Labels: core.NewLabels([]string{"workspace=ops"}),
}
att, err := NewAttempt(file, nil, WithDAG(dag))
require.NoError(t, err)
require.NoError(t, att.Open(ctx))

queued := createTestStatus(core.Queued)
queued.AttemptID = "attempt-1"
queued.QueuedAt = time.Now().UTC().Format(time.RFC3339)
queued.Labels = dag.Labels.Strings()
require.NoError(t, att.Write(ctx, queued))
require.NoError(t, att.Write(ctx, queued))

Expand Down Expand Up @@ -942,6 +964,7 @@ func TestAttempt_WriteEmitsLifecycleTransitionsAndStatusUpdates(t *testing.T) {
snapshot, err := eventstore.DAGRunSnapshotFromEvent(store.events[0])
require.NoError(t, err)
assert.Equal(t, "test-dag", snapshot.DAGFile)
assert.Equal(t, []string{"workspace=ops"}, snapshot.Labels)
assert.Equal(t, core.Queued, snapshot.Status)
}

Expand All @@ -966,7 +989,7 @@ func TestAttempt_OpenRestoresLastEmittedLifecycleState(t *testing.T) {
require.NoError(t, att.Close(ctx))
require.Len(t, store.events, 1)

reopened, err := NewAttempt(file, nil, WithDAG(dag))
reopened, err := NewAttempt(file, nil)
require.NoError(t, err)
require.NoError(t, reopened.Open(ctx))
require.NoError(t, reopened.Write(ctx, queued))
Expand All @@ -982,6 +1005,9 @@ func TestAttempt_OpenRestoresLastEmittedLifecycleState(t *testing.T) {
eventstore.TypeDAGRunUpdated,
eventstore.TypeDAGRunRunning,
}, captureEventTypes(store.events))
snapshot, err := eventstore.DAGRunSnapshotFromEvent(store.events[2])
require.NoError(t, err)
assert.Equal(t, "test-dag", snapshot.DAGFile)
}

type captureEventStore struct {
Expand Down
4 changes: 2 additions & 2 deletions internal/persis/file/incident/store.go
Original file line number Diff line number Diff line change
Expand Up @@ -243,7 +243,7 @@ func (s *Store) GetState(_ context.Context, providerID, dedupKey string) (*incid
return state, nil
}

func (s *Store) ListOpenStatesByDAG(_ context.Context, dagName string) ([]*incident.IncidentState, error) {
func (s *Store) ListOpenStates(_ context.Context) ([]*incident.IncidentState, error) {
s.mu.RLock()
defer s.mu.RUnlock()
entries, err := os.ReadDir(s.stateDir())
Expand All @@ -262,7 +262,7 @@ func (s *Store) ListOpenStatesByDAG(_ context.Context, dagName string) ([]*incid
if err != nil {
return nil, err
}
if state.DAGName != dagName || state.Status != incident.IncidentStatusOpen {
if state.Status != incident.IncidentStatusOpen {
continue
}
states = append(states, state)
Expand Down
2 changes: 1 addition & 1 deletion internal/persis/file/incident/store_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -97,7 +97,7 @@ func TestStorePersistsPolicySetAndState(t *testing.T) {
assert.Equal(t, incident.IncidentStatusOpen, loadedState.Status)
assert.Equal(t, "daily", loadedState.DAGName)

openStates, err := store.ListOpenStatesByDAG(context.Background(), "daily")
openStates, err := store.ListOpenStates(context.Background())
require.NoError(t, err)
require.Len(t, openStates, 1)
assert.Equal(t, state.DedupKey, openStates[0].DedupKey)
Expand Down
18 changes: 5 additions & 13 deletions internal/service/chatbridge/monitor.go
Original file line number Diff line number Diff line change
Expand Up @@ -419,14 +419,15 @@ func (m *NotificationMonitor) pollSource(ctx context.Context) {
if event == nil || !m.isInterestedEventType(event.Type) {
continue
}
status, err := eventstore.DAGRunStatusFromEvent(event)
snapshot, err := eventstore.DAGRunSnapshotFromEvent(event)
if err != nil {
m.logger.Warn("Failed to decode DAG-run event payload",
slog.String("event_id", event.ID),
slog.String("error", err.Error()),
)
continue
}
status := snapshot.DAGRunStatus()
observedAt := event.RecordedAt
if observedAt.IsZero() {
observedAt = time.Now().UTC()
Expand All @@ -435,6 +436,7 @@ func (m *NotificationMonitor) pollSource(ctx context.Context) {
Key: NotificationSeenKey(status),
Type: event.Type,
Status: status,
DAGFile: snapshot.DAGFile,
ObservedAt: observedAt.UTC(),
})
}
Expand Down Expand Up @@ -1176,12 +1178,7 @@ func cloneNotificationMonitorState(state notificationMonitorState) notificationM
}
pending := make(map[string]NotificationEvent, len(destState.Pending))
for key, event := range destState.Pending {
pending[key] = NotificationEvent{
Key: event.Key,
Type: event.Type,
Status: cloneNotificationStatus(event.Status),
ObservedAt: event.ObservedAt,
}
pending[key] = cloneNotificationEvent(event)
}
clone.Destinations[destination] = &notificationDestinationState{
Pending: pending,
Expand Down Expand Up @@ -1303,12 +1300,7 @@ func enqueueNotifications(state *notificationMonitorState, destinations []string
delete(destState.Pending, pendingKey)
}

destState.Pending[event.Key] = NotificationEvent{
Key: event.Key,
Type: event.Type,
Status: cloneNotificationStatus(event.Status),
ObservedAt: event.ObservedAt,
}
destState.Pending[event.Key] = cloneNotificationEvent(event)
queued = append(queued, queuedNotification{
destination: destination,
event: event,
Expand Down
20 changes: 16 additions & 4 deletions internal/service/chatbridge/monitor_state_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -75,7 +75,7 @@ func TestNotificationMonitor_BootstrapsFromCurrentHeadAndOnlyDeliversFutureEvent

var (
mu sync.Mutex
delivered []string
delivered []NotificationEvent
)
transport := &fakeNotificationTransport{
destinations: []string{"dest-1"},
Expand All @@ -84,7 +84,7 @@ func TestNotificationMonitor_BootstrapsFromCurrentHeadAndOnlyDeliversFutureEvent
defer mu.Unlock()
for _, event := range batch.Events {
if event.Status != nil {
delivered = append(delivered, event.Status.DAGRunID)
delivered = append(delivered, cloneNotificationEvent(event))
}
}
return true
Expand All @@ -108,20 +108,28 @@ func TestNotificationMonitor_BootstrapsFromCurrentHeadAndOnlyDeliversFutureEvent
AttemptID: "attempt-new",
Status: core.Failed,
Error: "new failure",
Labels: []string{"workspace=ops", "team=platform"},
FinishedAt: time.Now().UTC().Format(time.RFC3339),
}
require.NoError(t, service.Emit(context.Background(), eventstore.NewDAGRunEvent(
eventstore.Source{Service: eventstore.SourceServiceServer, Instance: "test"},
eventstore.TypeDAGRunFailed,
newStatus,
nil,
map[string]any{eventstore.DAGFileNameDataKey: "briefing-file"},
)))

require.Eventually(t, func() bool {
mu.Lock()
defer mu.Unlock()
return len(delivered) == 1 && delivered[0] == "run-new"
return len(delivered) == 1
}, notificationMonitorEventuallyTimeout(time.Second), 10*time.Millisecond)
mu.Lock()
deliveredEvent := cloneNotificationEvent(delivered[0])
mu.Unlock()
require.NotNil(t, deliveredEvent.Status)
assert.Equal(t, "run-new", deliveredEvent.Status.DAGRunID)
assert.Equal(t, "briefing-file", deliveredEvent.DAGFile)
assert.Equal(t, []string{"workspace=ops", "team=platform"}, deliveredEvent.Status.Labels)

require.Eventually(t, func() bool {
return !monitor.IsDelivered("dest-1", oldStatus) && monitor.IsDelivered("dest-1", newStatus)
Expand All @@ -137,6 +145,7 @@ func TestNotificationMonitor_RestartRequeuesPersistedPending(t *testing.T) {

status := &exec.DAGRunStatus{
Name: "briefing",
Labels: []string{"workspace=ops"},
Status: core.Failed,
DAGRunID: "run-1",
AttemptID: "attempt-1",
Expand All @@ -149,6 +158,7 @@ func TestNotificationMonitor_RestartRequeuesPersistedPending(t *testing.T) {
NotificationSeenKey(status): {
Key: NotificationSeenKey(status),
Status: cloneNotificationStatus(status),
DAGFile: "briefing-file",
ObservedAt: time.Now().UTC(),
},
},
Expand All @@ -168,6 +178,8 @@ func TestNotificationMonitor_RestartRequeuesPersistedPending(t *testing.T) {
assert.Equal(t, "dest-1", destination)
require.Len(t, batch.Events, 1)
assert.Equal(t, "run-1", batch.Events[0].Status.DAGRunID)
assert.Equal(t, "briefing-file", batch.Events[0].DAGFile)
assert.Equal(t, []string{"workspace=ops"}, batch.Events[0].Status.Labels)
calls++
return true
},
Expand Down
Loading