Skip to content
Open
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
35 changes: 35 additions & 0 deletions docs/research/migration-delete-race.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,35 @@
# Delete versus an in-flight forward create

Base: ad7b2298d. Fixed seed: `delete-create`, logical clock 2026-09-04T12:00:00Z.
Invariant: a successful schedule deletion cannot be followed by an active scheduler created from an older migration snapshot.

The deterministic test holds forward creation at a channel barrier. The real frontend `DeleteSchedule` executes its CHASM path against the real component engine, observes NotFound, then terminates the V1 workflow through a mocked History RPC. After deletion succeeds, the barrier releases the real `CreateSchedulerFromMigration` constructor and its CHASM transaction. The destination is active. The native control creates first, then deletes through the same frontend/component paths and observes a closed destination.

This is a component/ingress counterexample, not a full History persistence experiment. V1 termination is represented by a successful History RPC; the test does not run a live V1 worker, ambiguous database commits, or namespace failover. The forward service handler's only creation fence is the CHASM execution key and its existing-execution/sentinel handling. No state from the successful V1 delete is supplied to that creation transaction.

```mermaid
sequenceDiagram
participant Create as Forward create
participant Delete as DeleteSchedule
participant CHASM
participant V1
Note over Create: held before CHASM creation
Delete->>CHASM: delete
CHASM-->>Delete: NotFound
Delete->>V1: terminate
V1-->>Delete: success
Delete-->>Delete: return success
Create->>CHASM: create from old snapshot
CHASM-->>Create: active destination committed
```

![Scrubbed failure](migration-delete-race.svg)

Control: `go test -tags test_dep ./service/frontend -run '^TestMigrationDelete' -count=1`.
Counterexample: `TEMPORAL_RUN_MIGRATION_COUNTEREXAMPLES=1 go test -tags test_dep ./service/frontend -run '^TestMigrationDeleteCounterexample$' -count=1`.

Upstream audit: all public open PR titles and bodies fetched on 2026-09-04; no matching durable delete/incarnation fence found. #11924 changes sentinel classification, not this race.

Disposition: **no local production workaround**. Depend on [PR #44](https://github.com/chaptersix/temporal/pull/44)'s History ingress/incarnation protocol. Delete needs a replicated tombstone or generation advance in the same authoritative ownership protocol checked by migration creation. The generation must survive response loss, task retry, destination deletion/retention, and namespace failover. A second read before creation leaves exactly the same interleaving between that read and the write. A process-local lock or tombstone cannot fence failover or delayed RPCs.

Until that protocol exists, operators must stop forward-migration admission, reconcile ambiguous creates, and drain existing attempts before relying on deletion across backends. The existing API does not establish this fence. This is an explicit fail-closed operational disposition, not a claim that the current implementation fails closed. No new public API, proto field, or workflow behavior is introduced here.
1 change: 1 addition & 0 deletions docs/research/migration-delete-race.svg
Loading
Sorry, something went wrong. Reload?
Sorry, we cannot display this file.
Sorry, this file is invalid so it cannot be displayed.
103 changes: 103 additions & 0 deletions service/frontend/migration_delete_race_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,103 @@
package frontend

import (
"context"
"os"
"testing"
"time"

"github.com/stretchr/testify/require"
schedulepb "go.temporal.io/api/schedule/v1"
"go.temporal.io/api/serviceerror"
"go.temporal.io/api/workflowservice/v1"
"go.temporal.io/server/api/historyservice/v1"
"go.temporal.io/server/api/historyservicemock/v1"
"go.temporal.io/server/chasm"
"go.temporal.io/server/chasm/chasmtest"
chasmscheduler "go.temporal.io/server/chasm/lib/scheduler"
schedulerpb "go.temporal.io/server/chasm/lib/scheduler/gen/schedulerpb/v1"
"go.temporal.io/server/common/clock"
"go.temporal.io/server/common/log"
"go.temporal.io/server/common/metrics"
"go.temporal.io/server/common/namespace"
legacyscheduler "go.temporal.io/server/service/worker/scheduler"
"go.uber.org/mock/gomock"
"google.golang.org/grpc"
"google.golang.org/protobuf/types/known/timestamppb"
)

type migrationDeleteRaceClient struct {
schedulerpb.SchedulerServiceClient
engineCtx context.Context
ref chasm.ComponentRef
}

func (c *migrationDeleteRaceClient) DeleteSchedule(_ context.Context, req *schedulerpb.DeleteScheduleRequest, _ ...grpc.CallOption) (*schedulerpb.DeleteScheduleResponse, error) {
response, _, err := chasm.UpdateComponent(c.engineCtx, c.ref, func(s *chasmscheduler.Scheduler, ctx chasm.MutableContext, _ struct{}) (*schedulerpb.DeleteScheduleResponse, error) {
return s.Delete(ctx, req)
}, struct{}{})
return response, err
}

func TestMigrationDeleteNativeControl(t *testing.T) { testMigrationDeleteRace(t, false) }
func TestMigrationDeleteCounterexample(t *testing.T) {
if os.Getenv("TEMPORAL_RUN_MIGRATION_COUNTEREXAMPLES") != "1" {
t.Skip("set TEMPORAL_RUN_MIGRATION_COUNTEREXAMPLES=1")
}
testMigrationDeleteRace(t, true)
}

func testMigrationDeleteRace(t *testing.T, inFlight bool) {
t.Helper()
logger := log.NewNoopLogger()
config := &chasmscheduler.Config{Tweakables: func(string) chasmscheduler.Tweakables { return chasmscheduler.DefaultTweakables }}
builder := legacyscheduler.NewSpecBuilder(func() int { return 0 }, func() int { return 0 })
processor := chasmscheduler.NewSpecProcessor(config, metrics.NoopMetricsHandler, logger, builder)
registry := chasm.NewRegistry(logger)
require.NoError(t, registry.Register(&chasm.CoreLibrary{}))
require.NoError(t, registry.Register(chasmscheduler.NewLibrary(config, nil,
chasmscheduler.NewSchedulerIdleTaskHandler(chasmscheduler.SchedulerIdleTaskHandlerOptions{Config: config, MetricsHandler: metrics.NoopMetricsHandler, BaseLogger: logger}), nil,
chasmscheduler.NewGeneratorTaskHandler(chasmscheduler.GeneratorTaskHandlerOptions{Config: config, MetricsHandler: metrics.NoopMetricsHandler, BaseLogger: logger, SpecBuilder: builder, SpecProcessor: processor}), nil, nil, nil, nil)))
now := time.Date(2026, 9, 4, 12, 0, 0, 0, time.UTC)
engine := chasmtest.NewEngine(t, registry, chasmtest.WithTimeSource(clock.NewEventTimeSource().Update(now)))
engineCtx := chasm.NewEngineContext(context.Background(), engine)
key := chasm.ExecutionKey{NamespaceID: "ns-id", BusinessID: "schedule-id"}
ref := chasm.NewComponentRef[*chasmscheduler.Scheduler](key)
create := func() error {
_, err := chasm.StartExecution(engineCtx, key, chasmscheduler.CreateSchedulerFromMigration, &schedulerpb.CreateFromMigrationStateRequest{NamespaceId: key.NamespaceID, State: &schedulerpb.SchedulerMigrationState{
SchedulerState: &schedulerpb.SchedulerState{Namespace: "ns", NamespaceId: key.NamespaceID, ScheduleId: key.BusinessID, Schedule: &schedulepb.Schedule{Spec: &schedulepb.ScheduleSpec{}, State: &schedulepb.ScheduleState{Paused: true}}, Info: &schedulepb.ScheduleInfo{}},
GeneratorState: &schedulerpb.GeneratorState{LastProcessedTime: timestamppb.New(now)}, InvokerState: &schedulerpb.InvokerState{},
}})
return err
}
started, release, created := make(chan struct{}), make(chan struct{}), make(chan error, 1)
if inFlight {
go func() { close(started); <-release; created <- create() }()
<-started
} else {
require.NoError(t, create())
}
ctrl := gomock.NewController(t)
namespaces := namespace.NewMockRegistry(ctrl)
namespaces.EXPECT().GetNamespaceID(namespace.Name("ns")).Return(namespace.ID("ns-id"), nil).Times(2)
history := historyservicemock.NewMockHistoryServiceClient(ctrl)
history.EXPECT().TerminateWorkflowExecution(gomock.Any(), gomock.Any()).DoAndReturn(func(_ context.Context, req *historyservice.TerminateWorkflowExecutionRequest, _ ...grpc.CallOption) (*historyservice.TerminateWorkflowExecutionResponse, error) {
require.Equal(t, legacyscheduler.WorkflowIDPrefix+key.BusinessID, req.TerminateRequest.WorkflowExecution.WorkflowId)
return &historyservice.TerminateWorkflowExecutionResponse{}, nil
}).Times(1)
handler := &WorkflowHandler{logger: logger, namespaceRegistry: namespaces, historyClient: history, schedulerClient: &migrationDeleteRaceClient{engineCtx: engineCtx, ref: ref}, config: &Config{EnableSchedules: func(string) bool { return true }, EnableCHASMSchedulerCreation: func(string) bool { return true }}}
_, deleteErr := handler.DeleteSchedule(context.Background(), &workflowservice.DeleteScheduleRequest{Namespace: "ns", ScheduleId: key.BusinessID})
if inFlight {
close(release)
require.NoError(t, <-created)
}
require.NoError(t, deleteErr)
_, err := chasm.ReadComponent(engineCtx, ref, func(s *chasmscheduler.Scheduler, _ chasm.Context, _ struct{}) (struct{}, error) {
require.True(t, s.Closed, "seed=delete-create: successful DeleteSchedule must not be followed by an active migrated scheduler")
return struct{}{}, nil
}, struct{}{})
if err != nil {
var notFound *serviceerror.NotFound
require.ErrorAs(t, err, &notFound)
}
}
Loading