Skip to content
Draft
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
88 changes: 88 additions & 0 deletions chasm/lib/scheduler/buffered_start_lifecycle_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,88 @@
package scheduler_test

import (
"testing"

"github.com/stretchr/testify/require"
enumspb "go.temporal.io/api/enums/v1"
"go.temporal.io/api/serviceerror"
schedulespb "go.temporal.io/server/api/schedule/v1"
"go.temporal.io/server/chasm"
"go.temporal.io/server/chasm/lib/scheduler"
schedulerpb "go.temporal.io/server/chasm/lib/scheduler/gen/schedulerpb/v1"
"go.uber.org/mock/gomock"
"google.golang.org/protobuf/types/known/timestamppb"
)

func TestBufferedStartLifecycleTransitionsEngine(t *testing.T) {
t.Run("unprocessed to ready", func(t *testing.T) {
env := newSchedulerTestEngine(t, defaultSchedule())
now := env.timeSource.Now()

err := env.updateScheduler(
func(s *scheduler.Scheduler, ctx chasm.MutableContext) error {
invoker := s.Invoker.Get(ctx)
invoker.BufferedStarts = []*schedulespb.BufferedStart{{
NominalTime: timestamppb.New(now),
ActualTime: timestamppb.New(now),
RequestId: "ready",
WorkflowId: "ready-workflow",
OverlapPolicy: enumspb.SCHEDULE_OVERLAP_POLICY_ALLOW_ALL,
}}
ctx.AddTask(invoker, chasm.TaskAttributes{}, &schedulerpb.InvokerProcessBufferTask{})
return nil
})
require.NoError(t, err)

err = env.readScheduler(
func(s *scheduler.Scheduler, ctx chasm.Context) error {
start := s.Invoker.Get(ctx).GetBufferedStarts()[0]
require.Equal(t, int64(1), start.GetAttempt())
require.Empty(t, start.GetRunId())
require.Nil(t, start.GetBackoffTime())
return nil
})
require.NoError(t, err)
})

t.Run("ready to retrying", func(t *testing.T) {
env := newInvokerExecuteEngine(t)
now := env.timeSource.Now()

err := env.updateScheduler(
func(s *scheduler.Scheduler, ctx chasm.MutableContext) error {
invoker := s.Invoker.Get(ctx)
invoker.LastProcessedTime = timestamppb.New(now)
invoker.BufferedStarts = []*schedulespb.BufferedStart{{
NominalTime: timestamppb.New(now),
ActualTime: timestamppb.New(now),
DesiredTime: timestamppb.New(now),
RequestId: "retrying",
WorkflowId: "retrying-workflow",
OverlapPolicy: enumspb.SCHEDULE_OVERLAP_POLICY_ALLOW_ALL,
Attempt: 1,
}}
ctx.AddTask(invoker, chasm.TaskAttributes{}, &schedulerpb.InvokerExecuteTask{})
return nil
})
require.NoError(t, err)

env.frontendClient.EXPECT().
StartWorkflowExecution(gomock.Any(), startWorkflowExecutionRequestIDMatches("retrying")).
Return(nil, serviceerror.NewDeadlineExceeded("retry"))

executed, err := env.engine.FireSideEffectTasks(env.rootRef, now)
require.NoError(t, err)
require.Equal(t, 1, executed)

err = env.readScheduler(
func(s *scheduler.Scheduler, ctx chasm.Context) error {
start := s.Invoker.Get(ctx).GetBufferedStarts()[0]
require.Equal(t, int64(2), start.GetAttempt())
require.Empty(t, start.GetRunId())
require.True(t, start.GetBackoffTime().AsTime().After(now))
return nil
})
require.NoError(t, err)
})
}
108 changes: 108 additions & 0 deletions chasm/lib/scheduler/internal/buffered_start.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,108 @@
package internal

import (
"time"

schedulespb "go.temporal.io/server/api/schedule/v1"
"google.golang.org/protobuf/types/known/timestamppb"
)

// BufferedStartState identifies the lifecycle state encoded by a BufferedStart's fields.
type BufferedStartState int

const (
BufferedStartStateInvalid BufferedStartState = iota
BufferedStartStateUnprocessed
BufferedStartStateDeferred
BufferedStartStateReady
BufferedStartStateBackingOff
BufferedStartStateStarted
BufferedStartStateCompleted
)

// ClassifyBufferedStart returns the lifecycle state encoded by start at retryEvaluationTime.
func ClassifyBufferedStart(
start *schedulespb.BufferedStart,
retryEvaluationTime time.Time,
) BufferedStartState {
if start == nil || start.GetAttempt() < -1 {
return BufferedStartStateInvalid
}

hasRun := start.GetRunId() != ""
hasCompletion := start.GetCompleted() != nil
hasBackoff := start.GetBackoffTime() != nil

switch start.GetAttempt() {
case -1:
if hasRun || hasCompletion || hasBackoff {
return BufferedStartStateInvalid
}
return BufferedStartStateDeferred
case 0:
if hasRun || hasCompletion || hasBackoff {
return BufferedStartStateInvalid
}
return BufferedStartStateUnprocessed
default:
if hasRun {
if hasCompletion {
return BufferedStartStateCompleted
}
return BufferedStartStateStarted
}
if hasCompletion {
return BufferedStartStateInvalid
}
if hasBackoff && start.GetBackoffTime().AsTime().After(retryEvaluationTime) {
return BufferedStartStateBackingOff
}
return BufferedStartStateReady
}
}

// MarkStartUnprocessed marks a start for its initial overlap-policy pass.
func MarkStartUnprocessed(start *schedulespb.BufferedStart) {
start.Attempt = 0
}

// MarkStartMigratedUnprocessed clears V2 retry state from a pending V1 start.
func MarkStartMigratedUnprocessed(start *schedulespb.BufferedStart) {
MarkStartUnprocessed(start)
start.BackoffTime = nil
}

// MarkStartDeferred marks a start as waiting for an overlapping action to complete.
func MarkStartDeferred(start *schedulespb.BufferedStart) {
start.Attempt = -1
}

// MarkStartReady marks a start as ready for its first execution attempt.
func MarkStartReady(start *schedulespb.BufferedStart) {
start.Attempt = 1
}

// MarkStartRetrying advances a start to its next attempt and backoff deadline.
func MarkStartRetrying(
start *schedulespb.BufferedStart,
nextAttempt int64,
backoffTime *timestamppb.Timestamp,
) {
start.Attempt = nextAttempt
start.BackoffTime = backoffTime
}

// MarkStartStarted records the execution created for a start.
func MarkStartStarted(
start *schedulespb.BufferedStart,
runID string,
startTime *timestamppb.Timestamp,
) {
start.RunId = runID
start.StartTime = startTime
}

// MarkStartCompleted records a start's terminal workflow result.
func MarkStartCompleted(start *schedulespb.BufferedStart, result *schedulespb.CompletedResult) {
start.Completed = result
}
150 changes: 150 additions & 0 deletions chasm/lib/scheduler/internal/buffered_start_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,150 @@
package internal

import (
"testing"
"time"

"github.com/stretchr/testify/require"
schedulespb "go.temporal.io/server/api/schedule/v1"
"google.golang.org/protobuf/types/known/timestamppb"
)

func TestClassifyBufferedStart(t *testing.T) {
now := time.Now()

tests := []struct {
name string
start *schedulespb.BufferedStart
want BufferedStartState
}{
{name: "nil", want: BufferedStartStateInvalid},
{name: "unknown negative attempt", start: &schedulespb.BufferedStart{Attempt: -2}, want: BufferedStartStateInvalid},
{name: "unprocessed", start: unprocessedStart(), want: BufferedStartStateUnprocessed},
{name: "deferred", start: deferredStart(), want: BufferedStartStateDeferred},
{name: "ready without backoff", start: readyStart(), want: BufferedStartStateReady},
{
name: "ready before backoff boundary",
start: func() *schedulespb.BufferedStart {
start := readyStart()
start.BackoffTime = timestamppb.New(now.Add(-time.Nanosecond))
return start
}(),
want: BufferedStartStateReady,
},
{
name: "ready at backoff boundary",
start: func() *schedulespb.BufferedStart {
start := readyStart()
start.BackoffTime = timestamppb.New(now)
return start
}(),
want: BufferedStartStateReady,
},
{name: "backing off", start: backingOffStart(now), want: BufferedStartStateBackingOff},
{name: "started", start: runningStart(now), want: BufferedStartStateStarted},
{name: "completed", start: completedStart(now), want: BufferedStartStateCompleted},
{
name: "completion before start",
start: &schedulespb.BufferedStart{
Attempt: 1,
Completed: &schedulespb.CompletedResult{},
},
want: BufferedStartStateInvalid,
},
{
name: "unprocessed with run ID",
start: &schedulespb.BufferedStart{
RunId: "run-id",
},
want: BufferedStartStateInvalid,
},
{
name: "deferred with backoff",
start: &schedulespb.BufferedStart{
Attempt: -1,
BackoffTime: timestamppb.New(now),
},
want: BufferedStartStateInvalid,
},
}

for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
require.Equal(t, tt.want, ClassifyBufferedStart(tt.start, now))
})
}
}

func TestBufferedStartTransitions(t *testing.T) {
now := time.Now()
backoffTime := timestamppb.New(now.Add(time.Minute))
startTime := timestamppb.New(now.Add(2 * time.Minute))
completion := &schedulespb.CompletedResult{CloseTime: timestamppb.New(now.Add(3 * time.Minute))}
start := &schedulespb.BufferedStart{RequestId: "request-id"}

MarkStartDeferred(start)
require.Equal(t, int64(-1), start.GetAttempt())
require.Equal(t, "request-id", start.GetRequestId())
require.Equal(t, BufferedStartStateDeferred, ClassifyBufferedStart(start, now))

MarkStartUnprocessed(start)
require.Zero(t, start.GetAttempt())
require.Equal(t, BufferedStartStateUnprocessed, ClassifyBufferedStart(start, now))

MarkStartReady(start)
require.Equal(t, int64(1), start.GetAttempt())
require.Equal(t, BufferedStartStateReady, ClassifyBufferedStart(start, now))

MarkStartRetrying(start, 2, backoffTime)
require.Equal(t, int64(2), start.GetAttempt())
require.Same(t, backoffTime, start.GetBackoffTime())
require.Equal(t, BufferedStartStateBackingOff, ClassifyBufferedStart(start, now))

MarkStartStarted(start, "run-id", startTime)
require.Equal(t, "run-id", start.GetRunId())
require.Same(t, startTime, start.GetStartTime())
require.Equal(t, int64(2), start.GetAttempt())
require.Same(t, backoffTime, start.GetBackoffTime())
require.Equal(t, BufferedStartStateStarted, ClassifyBufferedStart(start, now))

MarkStartCompleted(start, completion)
require.Same(t, completion, start.GetCompleted())
require.Equal(t, "run-id", start.GetRunId())
require.Equal(t, BufferedStartStateCompleted, ClassifyBufferedStart(start, now))
}

func unprocessedStart() *schedulespb.BufferedStart {
return &schedulespb.BufferedStart{RequestId: "unprocessed"}
}

func deferredStart() *schedulespb.BufferedStart {
return &schedulespb.BufferedStart{RequestId: "deferred", Attempt: -1}
}

func readyStart() *schedulespb.BufferedStart {
return &schedulespb.BufferedStart{RequestId: "ready", Attempt: 1}
}

func backingOffStart(now time.Time) *schedulespb.BufferedStart {
return &schedulespb.BufferedStart{
RequestId: "backing-off",
Attempt: 2,
BackoffTime: timestamppb.New(now.Add(time.Nanosecond)),
}
}

func runningStart(now time.Time) *schedulespb.BufferedStart {
return &schedulespb.BufferedStart{
RequestId: "started",
Attempt: 1,
RunId: "run-id",
StartTime: timestamppb.New(now),
}
}

func completedStart(now time.Time) *schedulespb.BufferedStart {
start := runningStart(now)
start.RequestId = "completed"
start.Completed = &schedulespb.CompletedResult{CloseTime: timestamppb.New(now.Add(time.Minute))}
return start
}
Loading
Loading