Merge pull request #108 from yuezh/fix_handler_state

Fix handler state
This commit is contained in:
lizzha
2015-09-16 10:18:11 +08:00
7 changed files with 170 additions and 54 deletions
+150 -38
View File
@@ -39,6 +39,17 @@ HANDLER_ENVIRONMENT_VERSION = 1.0
VALID_EXTENSION_STATUS = ['transitioning', 'error', 'success', 'warning']
VALID_HANDLER_STATUS = ['Ready', 'NotReady', "Installing", "Unresponsive"]
def handler_state_to_status(handler_state):
if handler_state == "Enabled":
return "Ready"
elif handler_state in VALID_HANDLER_STATUS:
return handler_state
else:
return "NotReady"
def validate_has_key(obj, key, fullname):
if key not in obj:
raise ExtensionError("Missing: {0}".format(fullname))
@@ -128,6 +139,12 @@ def get_installed_version(target_name):
installed_version = version
return installed_version
class ExtHandlerState(object):
Enabled = "Enabled"
Disabled = "Disabled"
Failed = "Failed"
class ExtHandlersHandler(object):
def process(self):
@@ -140,7 +157,7 @@ class ExtHandlersHandler(object):
if ext_handlers.extHandlers is None or \
len(ext_handlers.extHandlers) == 0:
logger.info("No extensions to handle")
logger.verb("No extensions to handle")
return
vm_status = prot.VMStatus()
@@ -161,7 +178,7 @@ class ExtHandlersHandler(object):
vm_status.vmAgent.extensionHandlers.append(handler_status)
try:
logger.info("Report vm agent status")
logger.verb("Report vm agent status")
protocol.report_vm_status(vm_status)
except prot.ProtocolError as e:
add_event(name="WALA", is_success=False, message = text(e))
@@ -196,8 +213,8 @@ class ExtHandlerInstance(object):
self.update_policy = ext_handler.properties.upgradePolicy
self.curr_version = curr_version
self.enabled = False
self.installed = installed
self.handler_state = None
self.lib_dir = OSUTIL.get_lib_dir()
self.ext_status = prot.ExtensionStatus()
@@ -205,12 +222,14 @@ class ExtHandlerInstance(object):
self.handler_status.name = self.name
self.handler_status.version = self.curr_version
self.ext = None
#Currently, extension will have no more than 1 instance
#Currently, extension settings will have no more than 1 instance
if len(ext_handler.properties.extensions) > 0:
self.ext = ext_handler.properties.extensions[0]
self.ext_status.sequenceNumber = self.ext.sequenceNumber
self.handler_status.extensions = [self.ext.name]
else:
#When no extension settings, set sequenceNumber to 0
self.ext = prot.Extension(sequenceNumber=0)
self.ext_status.sequenceNumber = self.ext.sequenceNumber
prefix = "[{0}]".format(self.get_full_name())
self.logger = logger.Logger(logger.DEFAULT_LOGGER, prefix)
@@ -224,38 +243,80 @@ class ExtHandlerInstance(object):
def handle(self):
self.init_logger()
self.logger.info("Start processing extension handler")
self.logger.verb("Start processing extension handler")
try:
self.handle_state()
self.collect_ext_status()
self.collect_handler_status()
if self.installed:
self.collect_ext_status()
self.collect_handler_status()
except ExtensionError as e:
self.report_event(is_success=False, message=text(e))
self.logger.info("Finished processing extension handler")
self.logger.verb("Finished processing extension handler")
def handle_state(self):
if self.installed:
self.logger.info("Installed version:{0}", self.curr_version)
handler_state = self.get_state()
self.enabled = (handler_state == "Ready")
self.handler_state = self.get_state()
self.handler_status.status = handler_state_to_status(self.handler_state)
self.logger.verb("Handler state: {0}", self.handler_state)
self.logger.verb("Sequence number: {0}", self.ext.sequenceNumber)
if self.state == 'enabled':
self.handle_enable()
if self.handler_state == ExtHandlerState.Failed:
self.logger.verb("Found previous failure, quit handle_enable")
return
if self.handler_state == ExtHandlerState.Enabled:
self.logger.verb("Already enabled with sequenceNumber: {0}",
self.ext.sequenceNumber)
self.logger.verb("Quit handle_enable")
return
try:
new = self.handle_enable()
if new is not None:
#Upgrade happened
new.set_state(ExtHandlerState.Enabled)
else:
self.set_state(ExtHandlerState.Enabled)
except ExtensionError as e:
self.set_state(ExtHandlerState.Failed)
raise e
elif self.state == 'disabled':
self.handle_disable()
if self.handler_state == ExtHandlerState.Failed:
self.logger.verb("Found previous failure, quit handle_disable")
return
if self.handler_state == ExtHandlerState.Disabled:
self.logger.verb("Already disabled with sequenceNumber: {0}",
self.ext.sequenceNumber)
self.logger.verb("Quit handle_disable")
return
try:
self.handle_disable()
self.set_state(ExtHandlerState.Disabled)
except ExtensionError as e:
self.set_state(ExtHandlerState.Failed)
raise e
elif self.state == 'uninstall':
self.handle_disable()
self.handle_uninstall()
try:
self.handle_uninstall()
except ExtensionError as e:
self.set_state(ExtHandlerState.Failed)
raise e
else:
raise ExtensionError("Unknown state:{0}".format(self.state))
def handle_enable(self):
target_version = self.get_target_version()
logger.info("Target version: {0}", target_version)
self.logger.info("Target version: {0}", target_version)
if self.installed:
if Version(target_version) > Version(self.curr_version):
self.upgrade(target_version)
return self.upgrade(target_version)
elif Version(target_version) == Version(self.curr_version):
self.enable()
else:
@@ -272,13 +333,19 @@ class ExtHandlerInstance(object):
self.enable()
def handle_disable(self):
if not self.installed or not self.enabled:
if not self.installed:
self.logger.verb("Not installed, quit disable")
return
self.disable()
def handle_uninstall(self):
if not self.installed:
self.logger.verb("Not installed, quit unistall")
self.handler_status = None
self.ext_status = None
return
self.disable()
self.uninstall()
def report_event(self, is_success=True, message=""):
@@ -288,6 +355,7 @@ class ExtHandlerInstance(object):
self.ext_status.code = -1
if self.handler_status is not None:
self.handler_status.message = message
self.set_state_message(message)
if not is_success:
self.handler_status.status = "NotReady"
add_event(name=self.name, op=self.ext_status.operation,
@@ -321,6 +389,7 @@ class ExtHandlerInstance(object):
new.install()
self.logger.info("Enable new extension")
new.enable()
return new
def download(self):
self.logger.info("Download extension package")
@@ -361,9 +430,8 @@ class ExtHandlerInstance(object):
fileutil.mkdir(status_dir, mode=0o700)
conf_dir = self.get_conf_dir()
fileutil.mkdir(conf_dir, mode=0o700)
#Init handler state to uninstall
self.set_state("NotReady")
self.make_handler_state_dir()
#Save HandlerEnvironment.json
self.create_handler_env()
@@ -374,7 +442,6 @@ class ExtHandlerInstance(object):
man = self.load_manifest()
self.launch_command(man.get_enable_command())
self.set_state("Ready")
def disable(self):
self.logger.info("Disable extension.")
@@ -382,45 +449,47 @@ class ExtHandlerInstance(object):
man = self.load_manifest()
self.launch_command(man.get_disable_command(), timeout=900)
self.set_state("NotReady")
def install(self):
self.logger.info("Install extension.")
self.set_operation(WALAEventOperation.Install)
man = self.load_manifest()
self.set_state("Installing")
self.launch_command(man.get_install_command(), timeout=900)
self.set_state("Ready")
self.installed = True
def uninstall(self):
self.logger.info("Uninstall extension.")
self.set_operation(WALAEventOperation.UnInstall)
man = self.load_manifest()
self.launch_command(man.get_uninstall_command())
self.logger.info("Remove ext handler dir: {0}", self.get_base_dir())
try:
shutil.rmtree(self.get_base_dir())
except IOError as e:
raise ExtensionError("Failed to rm ext handler dir: {0}".format(e))
self.installed = False
self.handler_status = None
self.ext_status = None
def update(self):
self.logger.info("Update extension.")
self.set_operation(WALAEventOperation.Update)
man = self.load_manifest()
self.launch_command(man.get_update_command(), timeout=900)
def collect_handler_status(self):
self.logger.info("Collect extension handler status")
self.logger.verb("Collect extension handler status")
if self.handler_status is None:
return
self.handler_status.status = self.get_state()
handler_state = self.get_state()
self.handler_status.status = handler_state_to_status(handler_state)
self.handler_status.message = self.get_state_messsage()
man = self.load_manifest()
if man.is_report_heartbeat():
heartbeat = self.collect_heartbeat()
@@ -428,7 +497,7 @@ class ExtHandlerInstance(object):
self.handler_status.status = heartbeat['status']
def collect_ext_status(self):
self.logger.info("Collect extension status")
self.logger.verb("Collect extension status")
if self.handler_status is None:
return
@@ -444,21 +513,54 @@ class ExtHandlerInstance(object):
raise ExtensionError("Failed to get status file: {0}".format(e))
except ValueError as e:
raise ExtensionError("Malformed status file: {0}".format(e))
def make_handler_state_dir(self):
handler_state_dir = self.get_handler_state_dir()
fileutil.mkdir(handler_state_dir, 0o600)
if not os.path.exists(handler_state_dir):
os.makedirs(handler_state_dir)
def get_state(self):
handler_state_file = self.get_handler_state_file()
if not os.path.isfile(handler_state_file):
return None
try:
handler_state = fileutil.read_file(handler_state_file)
if handler_state is not None:
handler_state = handler_state.rstrip()
return handler_state
except IOError as e:
raise ExtensionError("Failed to get handler status: {0}".format(e))
raise ExtensionError("Failed to get handler state: {0}".format(e))
def set_state(self, state):
handler_state_file = self.get_handler_state_file()
if not os.path.isfile(handler_state_file):
self.make_handler_state_dir()
try:
fileutil.write_file(handler_state_file, state)
except IOError as e:
raise ExtensionError("Failed to set handler status: {0}".format(e))
raise ExtensionError("Failed to set handler state: {0}".format(e))
def get_state_messsage(self):
handler_state_message_file= self.get_handler_state_message_file()
if not os.path.isfile(handler_state_message_file):
return None
try:
message = fileutil.read_file(handler_state_message_file)
return message
except IOError as e:
raise ExtensionError(("Failed to get handler state message: {0}"
"").format(e))
def set_state_message(self, message):
handler_state_message_file = self.get_handler_state_message_file()
if not os.path.isfile(handler_state_message_file):
self.make_handler_state_dir()
try:
fileutil.write_file(handler_state_message_file, message)
except IOError as e:
raise ExtensionError(("Failed to set handler status message: {0}"
"").format(e))
def collect_heartbeat(self):
self.logger.info("Collect heart beat")
@@ -523,7 +625,7 @@ class ExtHandlerInstance(object):
def update_settings(self):
if self.ext is None:
self.logger.verbose("Extension has no settings")
self.logger.verb("Extension has no settings")
return
settings = {
@@ -602,9 +704,19 @@ class ExtHandlerInstance(object):
def get_settings_file(self):
return os.path.join(self.get_conf_dir(),
"{0}.settings".format(self.ext.sequenceNumber))
def get_handler_state_dir(self):
return os.path.join(OSUTIL.get_lib_dir(), "handler_state",
self.get_full_name())
def get_handler_state_file(self):
return os.path.join(self.get_conf_dir(), 'HandlerState')
return os.path.join(self.get_handler_state_dir(),
'{0}.state'.format(self.ext.sequenceNumber))
def get_handler_state_message_file(self):
return os.path.join(self.get_handler_state_dir(),
'{0}.message'.format(self.ext.sequenceNumber))
def get_heartbeat_file(self):
return os.path.join(self.get_base_dir(), 'heartbeat.log')
+3 -3
View File
@@ -82,11 +82,11 @@ class EventMonitor(object):
def collect_event(self, evt_file_name):
try:
logger.info("Found event file: {0}", evt_file_name)
logger.verb("Found event file: {0}", evt_file_name)
with open(evt_file_name, "rb") as evt_file:
#if fail to open or delete the file, throw exception
json_str = evt_file.read().decode("utf-8",'ignore')
logger.info("Processed event file: {0}", evt_file_name)
logger.verb("Processed event file: {0}", evt_file_name)
os.remove(evt_file_name)
return json_str
except IOError as e:
@@ -109,7 +109,7 @@ class EventMonitor(object):
data = json.loads(data_str)
except ValueError as e:
logger.verb(data_str)
logger.error("Failed to decode json event file{0}", e)
logger.error("Failed to decode json event file: {0}", e)
continue
event = prot.TelemetryEvent()
+2 -2
View File
@@ -35,7 +35,7 @@ class Logger(object):
self.appenders.extend(logger.appenders)
self.prefix = prefix
def verbose(self, msg_format, *args):
def verb(self, msg_format, *args):
self.log(LogLevel.VERBOSE, msg_format, *args)
def info(self, msg_format, *args):
@@ -135,7 +135,7 @@ def add_logger_appender(appender_type, level=LogLevel.INFO, path=None):
DEFAULT_LOGGER.add_appender(appender_type, level, path)
def verb(msg_format, *args):
DEFAULT_LOGGER.verbose(msg_format, *args)
DEFAULT_LOGGER.verb(msg_format, *args)
def info(msg_format, *args):
DEFAULT_LOGGER.info(msg_format, *args)
+1 -1
View File
@@ -103,7 +103,7 @@ class CertList(DataContract):
self.certificates = DataContractList(Cert)
class Extension(DataContract):
def __init__(self, name=None, sequenceNumber=None, publicSettings={},
def __init__(self, name=None, sequenceNumber=None, publicSettings=None,
privateSettings=None, certificateThumbprint=None):
self.name = name
self.sequenceNumber = sequenceNumber
+6 -3
View File
@@ -123,10 +123,15 @@ def _fetch_uri(uri, headers, chk_proxy=False):
except restutil.HttpError as e:
raise ProtocolError(text(e))
if(resp.status == httpclient.FORBIDDEN):
logger.info("Sleep to prevent throttling.")
time.sleep(10)
if(resp.status == httpclient.GONE):
raise WireProtocolResourceGone(uri)
if(resp.status != httpclient.OK):
raise ProtocolError("{0} - {1}".format(resp.status, uri))
data = resp.read()
if data is None:
return None
@@ -260,8 +265,6 @@ def ext_handler_status_to_v1(handler_status, ext_statuses, timestamp):
"lang":"en-US",
"message": handler_status.message
},
'runtimeSettingsStatus' : {
}
}
if len(handler_status.extensions) > 0:
@@ -320,7 +323,7 @@ class StatusBlob(object):
def upload(self, url):
#TODO upload extension only if content has changed
logger.info("Upload status blob")
logger.verb("Upload status blob")
blob_type = self.get_blob_type(url)
data = self.to_json()
+3 -2
View File
@@ -112,8 +112,10 @@ class TestExtensions(unittest.TestCase):
self.assertEqual("/tmp/TestExt-2.0/status", test_ext.get_status_dir())
self.assertEqual("/tmp/TestExt-2.0/status/0.status",
test_ext.get_status_file())
self.assertEqual("/tmp/TestExt-2.0/config/HandlerState",
self.assertEqual("/tmp/handler_state/TestExt-2.0/0.state",
test_ext.get_handler_state_file())
self.assertEqual("/tmp/handler_state/TestExt-2.0/0.message",
test_ext.get_handler_state_message_file())
self.assertEqual("/tmp/TestExt-2.0/config", test_ext.get_conf_dir())
self.assertEqual("/tmp/TestExt-2.0/config/0.settings",
test_ext.get_settings_file())
@@ -143,7 +145,6 @@ class TestExtensions(unittest.TestCase):
test_ext.handle_uninstall()
self.assertEqual(None, mock_launch_command.args)
self.assertEqual(None, mock_set_state.args)
self.assertEqual(None, test_ext.ext_status.operation)
test_ext = ext.ExtHandlerInstance(ext_sample, pkg_list_sample,
ext_sample.properties.version, True)
+5 -5
View File
@@ -31,7 +31,7 @@ class TestLogger(unittest.TestCase):
def test_no_appender(self):
#The logger won't throw exception even if no appender.
_logger = logger.Logger()
_logger.verbose("Assert no exception")
_logger.verb("Assert no exception")
_logger.info("Assert no exception")
_logger.warn("Assert no exception")
_logger.error("Assert no exception")
@@ -41,8 +41,8 @@ class TestLogger(unittest.TestCase):
_logger.info("This is an exception {0}", Exception("Test"))
_logger.info("This is an number {0}", 0)
_logger.info("This is an boolean {0}", True)
_logger.verbose("{0}")
_logger.verbose("{0} {1}", 0, 1)
_logger.verb("{0}")
_logger.verb("{0} {1}", 0, 1)
_logger.info("{0} {1}", 0, 1)
_logger.warn("{0} {1}", 0, 1)
_logger.error("{0} {1}", 0, 1)
@@ -61,7 +61,7 @@ class TestLogger(unittest.TestCase):
self.assertTrue(tools.simple_file_grep('/tmp/testlog', msg))
msg = text(uuid.uuid4())
_logger.verbose("Verbose should not be logged: {0}", msg)
_logger.verb("Verbose should not be logged: {0}", msg)
self.assertFalse(tools.simple_file_grep('/tmp/testlog', msg))
@@ -76,7 +76,7 @@ class TestLogger(unittest.TestCase):
self.assertTrue(tools.simple_file_grep('/tmp/testlog', msg))
msg = text(uuid.uuid4())
_logger.verbose("Test logger: {0}", msg)
_logger.verb("Test logger: {0}", msg)
self.assertFalse(tools.simple_file_grep('/tmp/testlog', msg))