Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
29 commits
Select commit Hold shift + click to select a range
48822fc
Add initial error handling
whitehawk Feb 27, 2026
b34be7f
Implement full rollback initial draft
whitehawk Mar 2, 2026
211e6bf
Tests and fixes for rollback
whitehawk Mar 3, 2026
fb1467e
Fix test
whitehawk Mar 3, 2026
bdf7626
Merge branch 'feature/ADBDEV-6608' into ADBDEV-9083-re
whitehawk Mar 3, 2026
35a8f4e
Refactoring and improvements
whitehawk Mar 3, 2026
7549678
Update check_down_segments
whitehawk Mar 3, 2026
d706583
Add tests and fixes
whitehawk Mar 3, 2026
f777e69
Fix tests, drop schema at the end of the rollback
whitehawk Mar 4, 2026
6e4bcd8
Update test descriptions
whitehawk Mar 4, 2026
197c187
Add test
whitehawk Mar 4, 2026
41c5695
Add 6.1.2 test
whitehawk Mar 4, 2026
682b12c
Add test 6.1.3.
whitehawk Mar 4, 2026
50a0033
Add test 6.3.2
whitehawk Mar 4, 2026
f018d2b
Add test 6.4.2.
whitehawk Mar 4, 2026
e971872
Update comments, remove dbg logs
whitehawk Mar 4, 2026
8379e87
Update rollback prepare
whitehawk Mar 4, 2026
07d3dca
Truncate steps table
whitehawk Mar 4, 2026
67c24cd
Improve logging
whitehawk Mar 5, 2026
7deaba5
Add test 6.4.3., update test 7.2
whitehawk Mar 5, 2026
0a5fbb0
Rename wrap_state_func_with_faults() to wrap_func_with_faults()
whitehawk Mar 5, 2026
3e7d4f8
Update basic tests, and fix code for them
whitehawk Mar 5, 2026
43915d0
Add case for 6.2 test, and related fix
whitehawk Mar 5, 2026
934ad5f
Add 6.2.2. test
whitehawk Mar 5, 2026
a58de48
Add test 7.3.2.
whitehawk Mar 5, 2026
45d69d8
Refactor
whitehawk Mar 6, 2026
ee3f199
Split 6.2.2. test
whitehawk Mar 6, 2026
0e3020f
Uncomment some checks
whitehawk Mar 6, 2026
ecf1efc
Minor test updates
whitehawk Mar 6, 2026
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
21 changes: 2 additions & 19 deletions gpMgmt/bin/ggrebalance
Original file line number Diff line number Diff line change
Expand Up @@ -273,28 +273,11 @@ def check_down_segments(logger: Any, options: Any, dburl: dbconn.DbURL):
conn = dbconn.connect(dburl, encoding='UTF8', allowSystemTableMods=True)
dbconn.execSQL(conn, "SELECT gp_request_fts_probe_scan()")
cnt_primaries_down = int(dbconn.queryRow(conn, f"SELECT COUNT(1) FROM gp_segment_configuration WHERE role = 'p' AND status = 'd'")[0])
cnt_mirrors_down = int(dbconn.queryRow(conn, f"SELECT COUNT(1) FROM gp_segment_configuration WHERE role = 'm' AND status = 'd'")[0])
conn.close()

if cnt_primaries_down != 0:
raise Exception('Detected some primary segments are down, please recover manually')

if cnt_mirrors_down != 0:
logger.info("Some mirrors are down, trying to recover them, it may take some time...")
recoverseg_options = "-a -F"
if options.logfile_directory is not None:
recoverseg_options = recoverseg_options + f' -l "{str(options.logfile_directory)}"'
global cmd_recoverseg
try:
cmd_recoverseg = GpRecoverSeg("Running gprecoverseg", options=recoverseg_options)
cmd_recoverseg.run(validateAfter=True)
except Exception as e:
logger.error(str(e))
error_msg = f"Failed to execute 'gprecoverseg {recoverseg_options}'"
raise Exception(error_msg)
finally:
cmd_recoverseg = None

def main(options, args, parser):
conn = None
try:
Expand Down Expand Up @@ -322,14 +305,14 @@ def main(options, args, parser):

dburl = dbconn.DbURL(dbname=DBNAME, port=gpenv.getCoordinatorPort())

check_down_segments(logger, options, dburl)

check_running_gputils(dburl, options.coordinator_data_directory)

create_pid_file(options.coordinator_data_directory)

gparray_dump_filename = options.coordinator_data_directory + '/gparraydump'

check_down_segments(logger, options, dburl)

logger.info('Init gparray from catalog')
try:
gparray = GpArray.initFromCatalog(dburl, utility=True)
Expand Down
9 changes: 6 additions & 3 deletions gpMgmt/bin/gpmovemirrors
Original file line number Diff line number Diff line change
Expand Up @@ -85,6 +85,8 @@ def parseargs():
help='show this help message and exit.')
parser.add_option('-a', dest="interactive", action='store_false', default=True,
help="quiet mode, do not require user input for confirmations")
parser.add_option('--skip-resource-estimation', dest='skip_resource_estimation', metavar='<skip_resource_estimation>',
action='store_true', default=False, help='Skip resource estimation (storage)')
parser.add_option('--usage', action="briefhelp")

parser.set_defaults(verbose=False, filters=[], slice=(None, None))
Expand Down Expand Up @@ -376,9 +378,10 @@ try:
pairs.append(pair)

""" Validating Disk Space requirement """
disk_usage = RelocateDiskUsage(pairs, options.batch_size, options)
if not disk_usage.validate_disk_space():
raise InsufficientDiskSpaceError("Insufficient disk space on target mirror hosts.")
if not options.skip_resource_estimation:
disk_usage = RelocateDiskUsage(pairs, options.batch_size, options)
if not disk_usage.validate_disk_space():
raise InsufficientDiskSpaceError("Insufficient disk space on target mirror hosts.")

""" Prepare common execution steps for running commands on segments """
oldMirrorsToMove = [mirror for mirror in newConfig.oldMirrorList if not mirror.inPlace]
Expand Down
41 changes: 41 additions & 0 deletions gpMgmt/bin/gppylib/commands/gp.py
Original file line number Diff line number Diff line change
Expand Up @@ -325,6 +325,28 @@ def __init__(self, datadir, mode, wait, timeout):
self.append("stop")


class PgCtlStatusArgs(CmdArgs):
"""
Used by CoordinatorStop, SegmentStop to format the pg_ctl command
to stop a backend postmaster

>>> str(PgCtlStatusArgs("/data1/coordinator/gpseg-1"))
'$GPHOME/bin/pg_ctl -D /data1/coordinator/gpseg-1 status'

"""

def __init__(self, datadir):
"""
@param datadir: database data directory
"""
CmdArgs.__init__(self, [
"$GPHOME/bin/pg_ctl",
"-D", str(datadir),
])
self.append("status")



class CoordinatorStart(Command):
def __init__(self, name, dataDir, port, era,
wrapper, wrapper_args, specialMode=None, restrictedMode=False, timeout=SEGMENT_TIMEOUT_DEFAULT,
Expand Down Expand Up @@ -441,6 +463,25 @@ def remote(name, hostname, dataDir, mode='smart'):
cmd.run(validateAfter=True)
return cmd

#-----------------------------------------------
class SegmentStatus(Command):
def __init__(self, name, dataDir, ctxt=LOCAL, remoteHost=None):

self.cmdStr = str( PgCtlStatusArgs(dataDir) )
Command.__init__(self, name, self.cmdStr, ctxt, remoteHost)

@staticmethod
def local(name, dataDir):
cmd=SegmentStatus(name, dataDir)
cmd.run(validateAfter=False)
return cmd

@staticmethod
def remote(name, hostname, dataDir):
cmd=SegmentStatus(name, dataDir, ctxt=REMOTE, remoteHost=hostname)
cmd.run(validateAfter=False)
return cmd

#-----------------------------------------------
class SegmentIsShutDown(Command):
"""
Expand Down
12 changes: 10 additions & 2 deletions gpMgmt/bin/gppylib/fault_injection.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@
GPMGMT_FAULT_FILE_FLAG = 'GPMGMT_FAULT_FILE_FLAG'

GPMGMT_FAULT_TYPE_SYSPEND = 'suspend'
GPMGMT_FAULT_TYPE_VALUE = 'value'

def inject_fault(fault_point):
if GPMGMT_FAULT_POINT in os.environ and fault_point == os.environ[GPMGMT_FAULT_POINT]:
Expand All @@ -32,10 +33,17 @@ def raise_exception(delay: int):
else:
raise Exception('Fault Injection %s' % os.environ[GPMGMT_FAULT_POINT])

def inject_fault_get_value() -> str:
if GPMGMT_FAULT_TYPE in os.environ and os.environ[GPMGMT_FAULT_TYPE] == GPMGMT_FAULT_TYPE_VALUE:
if GPMGMT_FAULT_POINT in os.environ:
return os.environ[GPMGMT_FAULT_POINT]
return ''

# decorator for test purposes
def wrap_state_func_with_faults(func):
def wrap_func_with_faults(func):
def func_with_faults(*args):
inject_fault(f'{func.__name__}_begin')
func(*args)
result = func(*args)
inject_fault(f'{func.__name__}_end')
return result
return func_with_faults
3 changes: 3 additions & 0 deletions gpMgmt/bin/gppylib/operations/buildMirrorSegments.py
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@
from gppylib.commands.gp import is_pid_postmaster, get_pid_from_remotehost
from gppylib.commands.unix import check_pid_on_remotehost
from gppylib.programs.clsRecoverSegment_triples import RecoveryTriplet
from gppylib.fault_injection import *

logger = gplog.get_default_logger()

Expand Down Expand Up @@ -313,6 +314,7 @@ def _trigger_fts_probe(self, port=0):
dbconn.execSQL(conn,"SELECT gp_request_fts_probe_scan()")
conn.close()

@wrap_func_with_faults
def _update_config(self, recovery_info_by_host, gpArray):
# should use mainUtils.getProgramName but I can't make it work!
programName = os.path.split(sys.argv[0])[-1]
Expand Down Expand Up @@ -587,6 +589,7 @@ def _run_recovery(self, action_name, recovery_info_by_host, gpEnv):
self._remove_progress_files(recovery_info_by_host, recovery_results)
return recovery_results

@wrap_func_with_faults
def _do_recovery(self, recovery_info_by_host, gpEnv):
"""
# Recover and start segments using gpsegrecovery, which will internally call either
Expand Down
3 changes: 3 additions & 0 deletions gpMgmt/bin/gppylib/operations/rebalanceSegments.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@
from gppylib.commands.gp import GpSegStopCmd
from gppylib.commands import base
from gppylib import gplog
from gppylib.fault_injection import *

from gppylib.operations.segment_reconfigurer import SegmentReconfigurer

Expand Down Expand Up @@ -116,6 +117,8 @@ def rebalance(self):
pool.addCommand(cmd)

base.join_and_indicate_progress(pool)

inject_fault('GpSegmentRebalanceOperation_rebalance_at_seg_stop')

failed_count = 0
completed = pool.getCompletedItems()
Expand Down
71 changes: 44 additions & 27 deletions gpMgmt/bin/gprebalance_modules/ggrebalance_main_sm.py
Original file line number Diff line number Diff line change
Expand Up @@ -145,6 +145,11 @@ def __init__(self, conn: dbconn.Connection, logger: Any, dburl: dbconn.DbURL, op
self.plan = None
self.main_state_from_prev_run = self.rebalance_schema.getMainStateFromPreviousRun()

self.shrink_state_from_prev_run = self.rebalance_schema.getShrinkStateFromPreviousRun()
self.is_shrink_rollback_in_progress = self.gg_shrink.state_is_from_rollback_flow(self.shrink_state_from_prev_run)
self.prev_shrink_run_was_complete = self.gg_shrink.state_is_final(self.shrink_state_from_prev_run)


def on_every_state(self) -> None:
if self.state in self.states_logged:
self.rebalance_schema.storeMainState(self.state)
Expand All @@ -168,7 +173,7 @@ def shutdown(self) -> None:

# state callbacks start here

@wrap_state_func_with_faults
@wrap_func_with_faults
def on_enter_STATE_OPTIONS_VALIDATION(self) -> None:
if self.options.clean_required:
self.trigger('move_to_STATE_CLEANUP')
Expand All @@ -177,25 +182,37 @@ def on_enter_STATE_OPTIONS_VALIDATION(self) -> None:
else:
self.trigger('move_to_STATE_PLANNING_STARTED')

@wrap_state_func_with_faults
@wrap_func_with_faults
def on_enter_STATE_CLEANUP(self) -> None:
if not self.rebalance_schema.schemaExists():
self.logger.info(f"Rebalance schema doesn't exist. Cleanup is not required.")
else:
prev_run_was_complete = (self.main_state_from_prev_run == 'STATE_EXECUTOR_DONE' or
self.main_state_from_prev_run == 'STATE_ROLLBACK')
self.gg_shrink.cleanup(prev_run_was_complete)
self.plan = self.rebalance_schema.retrieveSavedPlan()
if isinstance(self.plan, ShrinkPlan):
self.gg_shrink.cleanup(self.prev_shrink_run_was_complete)
self.rebalance_schema.dropSchema()
self.logger.info('Cleanup is complete')
self.trigger('move_to_STATE_END')

@wrap_state_func_with_faults
@wrap_func_with_faults
def on_enter_STATE_ROLLBACK(self) -> None:
self.plan = self.rebalance_schema.retrieveSavedPlan()
self.gg_shrink.rollback(self.plan)
self.trigger('move_to_STATE_END')

@wrap_state_func_with_faults
try:
if self.main_state_from_prev_run == 'STATE_EXECUTOR_DONE':
self.logger.info("Previous run was completed successfully. Can't perform rollback.")
return
self.plan = self.rebalance_schema.retrieveSavedPlan()
if isinstance(self.plan, ShrinkPlan):
if self.is_shrink_rollback_in_progress:
self.logger.info("Rollback is already in progress, and was interrupted. Execute 'ggrebalance' without '-r' flag.")
return
if not self.prev_shrink_run_was_complete:
self.gg_shrink.rollback(self.plan)
return
self.gg_rebalance.rollback()
finally:
self.trigger('move_to_STATE_END')

@wrap_func_with_faults
def on_enter_STATE_PLANNING_STARTED(self) -> None:
if self.options.target_segment_count != None:
self.plan = Planner(self.logger, self.dburl, self.gparray, self.options).plan()
Expand All @@ -205,11 +222,11 @@ def on_enter_STATE_PLANNING_STARTED(self) -> None:

self.trigger('move_to_STATE_PLANNING_DONE')

@wrap_state_func_with_faults
@wrap_func_with_faults
def on_enter_STATE_PLANNING_DONE(self) -> None:
self.trigger('move_to_STATE_CHECK_PREVIOUS_RUN')

@wrap_state_func_with_faults
@wrap_func_with_faults
def on_enter_STATE_CHECK_PREVIOUS_RUN(self) -> None:
if not self.rebalance_schema.schemaExists():
if self.plan == None:
Expand Down Expand Up @@ -243,57 +260,57 @@ def on_enter_STATE_CHECK_PREVIOUS_RUN(self) -> None:

self.trigger('move_to_STATE_EXECUTOR_STARTED')

@wrap_state_func_with_faults
@wrap_func_with_faults
def on_enter_STATE_SETUP_SCHEMA_STARTED(self) -> None:
# Create schema and status tables.
# It will also save plan in order to use it for recovering after interruption
self.rebalance_schema.createSchema(self.plan)
self.trigger('move_to_STATE_SETUP_SCHEMA_DONE')

@wrap_state_func_with_faults
@wrap_func_with_faults
def on_enter_STATE_SETUP_SCHEMA_DONE(self) -> None:
self.logger.info(f'Created "{self.rebalance_schema.getSchemaName()}" schema')
self.trigger('move_to_STATE_EXECUTOR_STARTED')

@wrap_state_func_with_faults
@wrap_func_with_faults
def on_enter_STATE_EXECUTOR_STARTED(self) -> None:
if isinstance(self.plan, ShrinkPlan):
shrink_state_from_prev_run = self.rebalance_schema.getShrinkStateFromPreviousRun()
if not self.gg_shrink.state_is_final(shrink_state_from_prev_run):
if not self.prev_shrink_run_was_complete:
self.trigger('move_to_STATE_SHRINK_STARTED')
return
self.trigger('move_to_STATE_REBALANCE_STARTED')

@wrap_state_func_with_faults
@wrap_func_with_faults
def on_enter_STATE_EXECUTOR_DONE(self) -> None:
self.trigger('move_to_STATE_END')

@wrap_state_func_with_faults
@wrap_func_with_faults
def on_enter_STATE_SHRINK_STARTED(self) -> None:
self.gg_shrink.run(self.plan)
self.trigger('move_to_STATE_SHRINK_DONE')

@wrap_state_func_with_faults
@wrap_func_with_faults
def on_enter_STATE_SHRINK_DONE(self) -> None:
self.trigger('move_to_STATE_REBALANCE_STARTED')

@wrap_state_func_with_faults
@wrap_func_with_faults
def on_enter_STATE_REBALANCE_STARTED(self) -> None:
if self.plan is not None and self.plan.getMoves() is not None:
if (self.plan is not None and
self.plan.getMoves() is not None and
not self.is_shrink_rollback_in_progress):
self.gg_rebalance.run(self.plan)
self.logger.info('Rebalance is complete')

self.trigger('move_to_STATE_REBALANCE_DONE')

@wrap_state_func_with_faults
@wrap_func_with_faults
def on_enter_STATE_REBALANCE_DONE(self) -> None:
self.trigger('move_to_STATE_EXECUTOR_DONE')

@wrap_state_func_with_faults
@wrap_func_with_faults
def on_enter_STATE_END(self) -> None:
pass

@wrap_state_func_with_faults
@wrap_func_with_faults
def on_enter_STATE_ERROR(self) -> None:
raise Exception('Main SM entered STATE_ERROR')

Expand Down
Loading