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
53 changes: 28 additions & 25 deletions internal/advancer/advancer.go
Original file line number Diff line number Diff line change
Expand Up @@ -313,7 +313,16 @@ func (s *Service) processInputs(
// Store the result in the database
err = s.repository.StoreAdvanceResult(ctx, input.EpochApplicationID, result)
if err != nil {
if errors.Is(err, repository.ErrApplicationNotRunnable) {
var errCause string
switch {
case errors.Is(err, context.Canceled) && errors.Is(ctx.Err(), context.Canceled):
errCause = "canceled advance result persistence"
// Shutdown interrupted persistence after the machine advanced.
// Discard that runtime; restart must recover from persisted state.
s.Logger.Debug("Advance result persistence canceled during shutdown; closing machine",
"application", app.Name, "index", input.Index, "error", err)
case errors.Is(err, repository.ErrApplicationNotRunnable):
errCause = "application status race"
// Another service durably fenced this application after the
// machine began its advance. Discard only this stale runtime; the
// database already contains the authoritative app-local outcome.
Expand All @@ -322,41 +331,35 @@ func (s *Service) processInputs(
"epoch", input.EpochIndex,
"index", input.Index,
"error", err)
closeErr := machine.Close()
if closeErr != nil {
s.Logger.Error("Could not close stale machine after application status race",
"application", app.Name,
"error", closeErr)
}
return processed, false, errors.Join(err, closeErr)
default:
errCause = "its advance result was not confirmed saved; service shutdown is still required"
// Advance has already changed the live machine, but the transaction
// did not confirm that its result was saved. The database may still
// show this input as pending. Reusing this machine could then execute
// the input again from the wrong state, so StoreAdvanceResult is not
// retried against this live machine.
s.Logger.Error(
"Could not confirm that the advance result was saved; "+
"the live machine has already advanced, so services will stop; "+
"after the node is restarted, execution will use persisted state",
"application", app.Name,
"epoch", input.EpochIndex,
"index", input.Index,
"error", err)
s.supervisor.Fatal(fmt.Errorf("unconfirmed advance result for %s: %w", app.Name, err))
}

// Advance has already changed the live machine, but the transaction
// did not confirm that its result was saved. The database may still
// show this input as pending. Reusing this machine could then execute
// the input again from the wrong state, so StoreAdvanceResult is not
// retried against this live machine.
s.Logger.Error(
"Could not confirm that the advance result was saved; "+
"the live machine has already advanced, so services will stop; "+
"after the node is restarted, execution will use persisted state",
"application", app.Name,
"epoch", input.EpochIndex,
"index", input.Index,
"error", err)

// Try to close the machine now so the already-advanced runtime cannot
// be used again. Cancel services even if Close fails. After the node
// is restarted, the machine is rebuilt from persisted state, and the
// database decides whether this input is still pending and needs a
// safe retry.
closeErr := machine.Close()
s.supervisor.Fatal(fmt.Errorf("unconfirmed advance result for %s: %w", app.Name, errors.Join(err, closeErr)))
if closeErr != nil {
s.Logger.Error("Could not close the machine after its advance result "+
"was not confirmed saved; service shutdown is still required",
s.Logger.Error("Could not close the machine after "+errCause,
"application", app.Name,
"error", closeErr)
s.supervisor.Fatal(fmt.Errorf("close machine for %s after %s: %w", app.Name, errCause, closeErr))
}
return processed, false, errors.Join(err, closeErr)
}
Expand Down
66 changes: 65 additions & 1 deletion internal/advancer/advancer_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -2336,6 +2336,7 @@ type MockMachineInstance struct {
machineImpl *MockMachineImpl
createSnapshotError error
destroyAfterSnapshotError bool
closeError error
closeCalls int
advanceCalls int
}
Expand Down Expand Up @@ -2386,7 +2387,7 @@ func (m *MockMachineInstance) Hash(ctx context.Context) ([32]byte, error) {
// Close implements the MachineInstance interface for testing
func (m *MockMachineInstance) Close() error {
m.closeCalls++
return nil
return m.closeError
}

// ------------------------------------------------------------------------------------------------
Expand All @@ -2398,6 +2399,7 @@ type MockRepository struct {
GetInputsReturn map[common.Address][]*Input
GetInputsError error
GetInputsBlock bool
StoreAdvanceHook func(context.Context) error
StoreAdvanceError error
StoreAdvanceCommitError error
StoreAdvanceFailCount int
Expand Down Expand Up @@ -2502,6 +2504,9 @@ func (mock *MockRepository) StoreAdvanceResult(
appID int64,
res *AdvanceResult,
) error {
if mock.StoreAdvanceHook != nil {
return mock.StoreAdvanceHook(ctx)
}
// Check for context cancellation
if ctx.Err() != nil {
return ctx.Err()
Expand Down Expand Up @@ -2750,3 +2755,62 @@ func marshal(res *AdvanceResult) []byte {
}
return data
}

func (s *AdvancerSuite) TestStoreAdvanceShutdownClassification() {
storageErr := errors.New("storage failed")
closeErr := errors.New("machine close failed")
for _, tc := range []struct {
name string
cancel bool
storeErr error
closeErr error
fatalErr error
}{
{"shutdown cancellation", true, fmt.Errorf("store: %w", context.Canceled), nil, nil},
{"independent cancellation", false, context.Canceled, nil, context.Canceled},
{"storage failure", false, storageErr, nil, storageErr},
{"storage failure during shutdown", true, storageErr, nil, storageErr},
{"close failure during shutdown", true, context.Canceled, closeErr, closeErr},
} {
s.Run(tc.name, func() {
require := s.Require()
env := s.setupOneApp()
ctx, cancel := context.WithCancel(s.T().Context())
defer cancel()
machine := env.mm.Map[env.app.Application.ID]
machine.closeError = tc.closeErr
env.repo.StoreAdvanceHook = func(storeCtx context.Context) error {
require.NoError(storeCtx.Err(), "cancellation must happen during persistence")
require.Equal(1, machine.advanceCalls)
if tc.cancel {
cancel()
}
return tc.storeErr
}
pending := newInput(env.app.Application.ID, 0, 0, marshal(randomAdvanceResult(0)))
address := env.app.Application.IApplicationAddress
env.repo.GetInputsReturn = map[common.Address][]*Input{address: {pending}}
processed, stopped, err := env.service.processInputs(ctx, env.app.Application, []*Input{
pending, newInput(env.app.Application.ID, 0, 1, []byte("unreachable")),
})
require.ErrorIs(err, tc.storeErr)
require.Zero(processed)
require.False(stopped)
require.Equal(1, machine.closeCalls)
require.Equal(1, machine.advanceCalls)
require.Empty(env.repo.StoredResults)
require.Equal([]*Input{pending}, env.repo.GetInputsReturn[address])
fatal := env.supervisor.FatalError.Load()
if tc.fatalErr == nil {
require.Nil(fatal, "shutdown cancellation must not become fatal")
} else {
require.NotNil(fatal)
require.ErrorIs(*fatal, tc.fatalErr)
if tc.closeErr != nil {
require.ErrorIs(err, tc.closeErr)
require.NotErrorIs(*fatal, context.Canceled)
}
}
})
}
}
2 changes: 1 addition & 1 deletion internal/evmreader/error_paths_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -561,7 +561,7 @@ func (s *EvmReaderSuite) TestEpochLengthZeroSetsAppCorrupted() {
s.evmReader.repository = repo

err := s.evmReader.readAndStoreInputs(s.ctx, 100, 110, apps)
s.Require().ErrorIs(err, errScanIncomplete)
s.Require().NoError(err)

// App must be set inoperable
repo.AssertNumberOfCalls(s.T(), "UpdateApplicationStatus", 1)
Expand Down
7 changes: 4 additions & 3 deletions internal/evmreader/evmreader.go
Original file line number Diff line number Diff line change
Expand Up @@ -177,9 +177,10 @@ func (r *Service) processBlockHead(
return r.runBlockScanners(ctx, apps, blockNumber) && len(apps) == len(observableApps)
}

// runBlockScanners reports whether all scheduled observation work completed without
// errors. An idle scan is healthy. Always run the remaining scanners after a
// failure so one application's stall does not prevent progress elsewhere.
// runBlockScanners reports whether observation completed without shared failures.
// Known input corruption with a persisted integrity status is application-local
// degradation; observation still runs on later ticks. An idle scan is healthy.
// Always run remaining scanners so one application's stall cannot block others.
func (r *Service) runBlockScanners(
ctx context.Context,
apps []appContracts,
Expand Down
62 changes: 40 additions & 22 deletions internal/evmreader/input.go
Original file line number Diff line number Diff line change
Expand Up @@ -300,6 +300,28 @@ func indexInputsIntoEpochs(
return epochInputMap, nil
}

// recordInputCorruption separates a known application-local integrity failure
// from failure to persist its status. Status helpers update app.Status only after
// a successful write and preserve previously persisted integrity terminals.
// Observation must continue on later ticks, including for terminal applications.
func (r *Service) recordInputCorruption(
ctx context.Context,
app *Application,
reasonFmt string,
args ...any,
) bool {
reason := fmt.Sprintf(reasonFmt, args...)
// setApplicationCorrupted always returns non-nil (the reason text itself).
// The DB error case is already logged inside setApplicationStatus.
// TODO: consider returning only DB errors instead of always returning an error.
_ = r.setApplicationCorrupted(ctx, app, "%s", reason)
if app.Status == ApplicationStatus_Corrupted || app.Status == ApplicationStatus_Diverged {
r.Logger.Warn("Input observation degraded for application", "application", app.Name, "reason", reason)
return true
}
return false
}

// readAndStoreInputs reads, inputs from the InputSource given specific filter options, indexes
// them into epochs and store the indexed inputs and epochs
func (r *Service) readAndStoreInputs(
Expand All @@ -308,8 +330,6 @@ func (r *Service) readAndStoreInputs(
mostRecentBlockNumber uint64,
apps []appContracts,
) error {
var scanErr error

if len(apps) == 0 {
r.Logger.Warn("No valid running applications")
return nil
Expand All @@ -325,8 +345,10 @@ func (r *Service) readAndStoreInputs(
err)
}

var scanIncomplete bool

if len(appInputsMap) != len(apps) {
scanErr = errScanIncomplete
scanIncomplete = true
}
addrToApp := mapAddressToApp(apps)

Expand All @@ -337,19 +359,15 @@ func (r *Service) readAndStoreInputs(
if !exists {
r.Logger.Error("Application address on input not found",
"address", address)
scanErr = errScanIncomplete
scanIncomplete = true
continue
}

epochLength := app.application.EpochLength
if epochLength == 0 {
// setApplicationCorrupted always returns non-nil (the reason text itself).
// The DB error case is already logged inside setApplicationStatus.
// On DB success the app is marked inoperable and won't reappear next tick.
// On DB failure the app reappears as Enabled next tick, retrying this path.
_ = r.setApplicationCorrupted(ctx, app.application,
ok := r.recordInputCorruption(ctx, app.application,
"Application has epoch length of zero")
scanErr = errScanIncomplete
scanIncomplete = scanIncomplete || !ok
continue
}

Expand All @@ -372,7 +390,7 @@ func (r *Service) readAndStoreInputs(
"error", err,
)
}
scanErr = errScanIncomplete
scanIncomplete = true
continue
}

Expand All @@ -381,8 +399,10 @@ func (r *Service) readAndStoreInputs(
epochLength, currentEpoch, inputs, mostRecentBlockNumber)
if err != nil {
if errors.Is(err, ErrInputForNonOpenEpoch) {
return r.setApplicationCorrupted(ctx, app.application,
ok := r.recordInputCorruption(ctx, app.application,
"Should never happen. %v", err)
scanIncomplete = scanIncomplete || !ok
continue
}
return fmt.Errorf("error indexing inputs: %w", err)
}
Expand Down Expand Up @@ -416,23 +436,18 @@ func (r *Service) readAndStoreInputs(
)
if err != nil {
if errors.Is(err, repository.ErrInputLogIdentityConflict) {
// A stored input's L1 log identity disagrees with rescanned
// chain data. Retrying the same insert every tick cannot
// succeed; without escalation the app would stall silently
// with Status OK. See setApplicationCorrupted contract in
// the epochLength == 0 branch above.
_ = r.setApplicationCorrupted(ctx, app.application,
ok := r.recordInputCorruption(ctx, app.application,
"stored input L1 log identity conflicts with rescanned chain data"+
" (possible reorg past the input cursor); operator reset required. %v", err)
scanErr = errScanIncomplete
scanIncomplete = scanIncomplete || !ok
continue
}
r.Logger.Error("Error storing inputs and epochs",
"application", app.application.Name,
"address", address,
"error", err,
)
scanErr = errScanIncomplete
scanIncomplete = true
continue
}
r.Logger.Debug("Inputs and epochs stored successfully",
Expand Down Expand Up @@ -486,7 +501,7 @@ func (r *Service) readAndStoreInputs(
"error", err,
)
}
scanErr = errScanIncomplete
scanIncomplete = true
} else {
r.Logger.Debug("Updated LastInputCheckBlock for applications without inputs",
"app_ids", appsToUpdate,
Expand All @@ -495,7 +510,10 @@ func (r *Service) readAndStoreInputs(
}
}

return scanErr
if scanIncomplete {
return errScanIncomplete
}
return nil
}

// readInputsFromBlockchain fetches inputs for each application independently.
Expand Down
Loading