mirror of
https://github.com/temporalio/temporal.git
synced 2026-08-30 18:41:49 -07:00
fix: [Scheduler] clarify version override diagnostics
This commit is contained in:
@@ -3603,7 +3603,7 @@ is non-fatal: the search continues past this threshold.`,
|
||||
-1,
|
||||
`SchedulerV1VersionCeiling caps the workflow version the V1 scheduler records into history, so histories written on this cluster stay replayable on peer clusters that do not support newer versions. Set it to the highest scheduler version supported by the lowest peer. Intended for multi-cluster failover and rollback. The supported floor is version 1 (OSS v1.20). A negative value (the default) disables the cap.
|
||||
The ceiling is reread on every tweakables evaluation. The version never decreases within a run, but raising or removing a ceiling can advance it on the next evaluation. A lower ceiling is recorded immediately; if it is below the version already recorded for the run, that version is retained.
|
||||
Operational notes: (1) A ceiling below 12 holds the version below CHASM migration support, so it pauses all V1->V2 CHASM migrations for the namespace until the ceiling is lifted (deferred, not dropped). (2) A ceiling below 6 skips custom search-attribute updates on schedule edits. (3) This caps V1 scheduler histories only; schedules already migrated to CHASM V2 are not made rollback-safe by it.`,
|
||||
Operational notes: (1) A ceiling below 12 holds fresh or not-yet-advanced runs below CHASM migration support, so it pauses their V1->V2 CHASM migrations until the ceiling is lifted (deferred, not dropped). It cannot downgrade a version already recorded in an existing run. (2) A ceiling below 6 skips custom search-attribute updates on schedule edits in fresh or not-yet-advanced runs. (3) This caps V1 scheduler histories only; schedules already migrated to CHASM V2 are not made rollback-safe by it.`,
|
||||
)
|
||||
SchedulerV1VersionOverride = NewNamespaceIntSetting(
|
||||
"worker.schedulerV1VersionOverride",
|
||||
|
||||
@@ -8,6 +8,7 @@ import (
|
||||
"github.com/stretchr/testify/mock"
|
||||
"github.com/stretchr/testify/require"
|
||||
schedulepb "go.temporal.io/api/schedule/v1"
|
||||
"go.temporal.io/sdk/log"
|
||||
"go.temporal.io/sdk/workflow"
|
||||
schedulespb "go.temporal.io/server/api/schedule/v1"
|
||||
schedulerpb "go.temporal.io/server/chasm/lib/scheduler/gen/schedulerpb/v1"
|
||||
@@ -15,6 +16,19 @@ import (
|
||||
"google.golang.org/protobuf/types/known/durationpb"
|
||||
)
|
||||
|
||||
type warningLogger struct {
|
||||
warnings []string
|
||||
}
|
||||
|
||||
var _ log.Logger = (*warningLogger)(nil)
|
||||
|
||||
func (*warningLogger) Debug(string, ...any) {}
|
||||
func (*warningLogger) Info(string, ...any) {}
|
||||
func (l *warningLogger) Warn(msg string, _ ...any) {
|
||||
l.warnings = append(l.warnings, msg)
|
||||
}
|
||||
func (*warningLogger) Error(string, ...any) {}
|
||||
|
||||
func TestDetermineVersionTransitions(t *testing.T) {
|
||||
for _, tc := range []struct {
|
||||
name string
|
||||
@@ -117,6 +131,36 @@ func TestDetermineVersionTransitions(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestDetermineVersionDiagnostics(t *testing.T) {
|
||||
t.Run("deduplicates invalid override warnings", func(t *testing.T) {
|
||||
logger := &warningLogger{}
|
||||
s := &scheduler{
|
||||
logger: logger,
|
||||
versionCeiling: func() int { return -1 },
|
||||
versionOverride: func() int { return int(LatestSchedulerWorkflowVersion) + 1 },
|
||||
}
|
||||
|
||||
s.determineVersion(TriggerImmediatelyTimestamp)
|
||||
s.determineVersion(TriggerImmediatelyTimestamp)
|
||||
|
||||
require.Equal(t, []string{"worker.schedulerV1VersionOverride is outside the supported range; ignored"}, logger.warnings)
|
||||
})
|
||||
|
||||
t.Run("does not report a ceiling that caps an override as ineffective", func(t *testing.T) {
|
||||
logger := &warningLogger{}
|
||||
s := &scheduler{
|
||||
logger: logger,
|
||||
versionCeiling: func() int { return int(MigrationHandoffFixes) },
|
||||
versionOverride: func() int { return int(LatestSchedulerWorkflowVersion) },
|
||||
}
|
||||
|
||||
version, _ := s.determineVersion(TriggerImmediatelyTimestamp)
|
||||
|
||||
require.Equal(t, SchedulerWorkflowVersion(MigrationHandoffFixes), version)
|
||||
require.Empty(t, logger.warnings)
|
||||
})
|
||||
}
|
||||
|
||||
// TestVersionCeilingDefersCHASMMigration verifies that a clamp below the CHASM gate keeps
|
||||
// migration markers out of history. The pending migration is retained for a later fresh run.
|
||||
func (s *workflowSuite) TestVersionCeilingDefersCHASMMigration() {
|
||||
|
||||
@@ -135,6 +135,8 @@ type (
|
||||
versionOverride func() int
|
||||
lastVersionCeiling int
|
||||
hasLastVersionCeiling bool
|
||||
lastVersionOverride int
|
||||
hasLastVersionOverride bool
|
||||
|
||||
tweakables TweakablePolicies
|
||||
|
||||
@@ -1145,6 +1147,12 @@ func (s *scheduler) handleMigrateSignal(ch workflow.ReceiveChannel, _ bool) {
|
||||
"namespace", s.State.Namespace,
|
||||
"schedule-id", s.State.ScheduleId,
|
||||
)
|
||||
if !s.tweakables.EnableCHASMMigration {
|
||||
s.logger.Error("failed assertion: received migrate signal while CHASM migration is disabled",
|
||||
"namespace", s.State.Namespace,
|
||||
"schedule-id", s.State.ScheduleId,
|
||||
)
|
||||
}
|
||||
s.State.PendingMigration = true
|
||||
}
|
||||
|
||||
@@ -1880,18 +1888,23 @@ func (s *scheduler) hasMinVersion(version SchedulerWorkflowVersion) bool {
|
||||
// the latest ceiling is applied on each iteration.
|
||||
func (s *scheduler) determineVersion(defaultVersion SchedulerWorkflowVersion) (SchedulerWorkflowVersion, int) {
|
||||
ceiling := s.versionCeiling()
|
||||
override := s.versionOverride()
|
||||
validOverride := override >= int(defaultVersion) && override <= int(LatestSchedulerWorkflowVersion)
|
||||
if ceiling != s.lastVersionCeiling || !s.hasLastVersionCeiling {
|
||||
if ceiling > int(defaultVersion) {
|
||||
if ceiling > int(resolveVersionBeforeCeiling(defaultVersion, override)) {
|
||||
s.logger.Warn("worker.schedulerV1VersionCeiling above the version this binary records; no effect",
|
||||
"ceiling", ceiling, "recordedVersion", defaultVersion)
|
||||
}
|
||||
s.lastVersionCeiling = ceiling
|
||||
s.hasLastVersionCeiling = true
|
||||
}
|
||||
override := s.versionOverride()
|
||||
if override > int(LatestSchedulerWorkflowVersion) {
|
||||
s.logger.Warn("worker.schedulerV1VersionOverride above the latest supported version; ignored",
|
||||
"override", override, "latestSupportedVersion", LatestSchedulerWorkflowVersion)
|
||||
if override != s.lastVersionOverride || !s.hasLastVersionOverride {
|
||||
if override >= 0 && !validOverride {
|
||||
s.logger.Warn("worker.schedulerV1VersionOverride is outside the supported range; ignored",
|
||||
"override", override, "defaultVersion", defaultVersion, "latestSupportedVersion", LatestSchedulerWorkflowVersion)
|
||||
}
|
||||
s.lastVersionOverride = override
|
||||
s.hasLastVersionOverride = true
|
||||
}
|
||||
return determineVersionTransition(defaultVersion, s.tweakables.Version, ceiling, override)
|
||||
}
|
||||
@@ -1911,10 +1924,14 @@ func clampVersion(v SchedulerWorkflowVersion, ceiling int) SchedulerWorkflowVers
|
||||
|
||||
// resolveVersion applies a valid override, then lowers the result to the ceiling.
|
||||
func resolveVersion(v SchedulerWorkflowVersion, ceiling, override int) SchedulerWorkflowVersion {
|
||||
return clampVersion(resolveVersionBeforeCeiling(v, override), ceiling)
|
||||
}
|
||||
|
||||
func resolveVersionBeforeCeiling(v SchedulerWorkflowVersion, override int) SchedulerWorkflowVersion {
|
||||
if override >= int(v) && override <= int(LatestSchedulerWorkflowVersion) {
|
||||
v = SchedulerWorkflowVersion(override)
|
||||
return SchedulerWorkflowVersion(override)
|
||||
}
|
||||
return clampVersion(v, ceiling)
|
||||
return v
|
||||
}
|
||||
|
||||
func panicIfErr(err error) {
|
||||
|
||||
Reference in New Issue
Block a user