Generate HandlerAggregatedStatus based on staus and heartbeat reported by handler.

This commit is contained in:
eridai
2014-06-06 17:02:08 +08:00
parent 58e70fb000
commit fd0aa07e0a
+85 -83
View File
@@ -100,6 +100,7 @@ global provisioned
provisioned=False
global provisionError
provisionError=None
HandlerStatusToAggStatus = {"installed":"Installing", "enabled":"Ready", "unintalled":"NotReady", "disabled":"NotReady"}
WaagentConf = """\
#
@@ -2916,7 +2917,7 @@ class ExtensionsConfig(object):
Log("No RuntimeSettings for " + name + " V " + version)
SimpleLog(p.plugin_log,"No RuntimeSettings for " + name + " V " + version)
SetFileContents(root +"/config/" + incarnation +".settings", config )
SetFileContents(root +"/config/" + seqNo +".settings", config )
#create HandlerEnvironment.json
handler_env='[{ "name": "'+name+'", "seqNo": "'+seqNo+'", "version": 1.0, "handlerEnvironment": { "logFolder": "'+os.path.dirname(p.plugin_log)+'", "configFolder": "' + root + '/config", "statusFolder": "' + root + '/status", "heartbeatFile": "'+ root + '/heartbeat.log"}}]'
SetFileContents(root+'/HandlerEnvironment.json',handler_env)
@@ -3009,7 +3010,7 @@ class ExtensionsConfig(object):
Error("No RuntimeSettings for " + name + " V " + version)
SimpleLog(p.plugin_log,"No RuntimeSettings for " + name + " V " + version)
SetFileContents(root +"/config/" + incarnation +".settings", config )
SetFileContents(root +"/config/" + seqNo +".settings", config )
# state is still enable
if (self.GetHandlerState(handler) == 'NotInstalled'): # run install first if true
@@ -3160,7 +3161,6 @@ class ExtensionsConfig(object):
return None
status=''
statuses=''
sent_suffix = '_sent'
for p in self.Plugins:
if p.getAttribute("state") == 'uninstall' or p.getAttribute("restricted") == 'true' :
continue
@@ -3169,24 +3169,14 @@ class ExtensionsConfig(object):
if p.getAttribute("isJson") != 'true':
LogIfVerbose("Plugin " + name+" version: " +version+" is not a JSON Extension. Skipping.")
continue
status_file=LibDir+'/'+name+'-'+version+'/status/'+incarnation+'.status'
if os.path.exists(status_file) !=True and os.path.exists(status_file+sent_suffix) != True:
if p.getAttribute("state") == 'disabled' :
LogIfVerbose(name+'-'+version+' is disabled. No status to report.')
continue
Error("Unable to locate " + status_file)
Error('Status report '+status_file+' Not sent!')
continue
elif os.path.exists(status_file) !=True and os.path.exists(status_file+sent_suffix) == True:
status_file = status_file+sent_suffix
reportHeartbeat = False
if len(p.getAttribute("manifestdata"))<1:
Error("Failed to get manifestdata.")
else:
reportHeartbeat = json.loads(p.getAttribute("manifestdata"))[0]['handlerManifest']['reportHeartbeat']
if len(statuses)>0:
statuses+=','
statuses+=GetFileContents(status_file)
if status_file.find(sent_suffix) < 0:
if os.path.exists(status_file + sent_suffix):
os.remove(status_file + sent_suffix)
os.rename(status_file,status_file+sent_suffix)
statuses+=self.GenerateAggStatus(name, version, reportHeartbeat)
if len(statuses)<1:
LogIfVerbose('No Handler status to report')
return None
@@ -3212,72 +3202,85 @@ class ExtensionsConfig(object):
return None
self.Util.Endpoint=uri.split('/')[2]
self.Util.HttpPutBlockBlob(uri, status)
Log('Status report '+status+' sent to ' + uri)
LogIfVerbose('Status report '+status+' sent to ' + uri)
return True
def CheckHeartbeat(self):
try:
incarnation=self.Extensions[0].getAttribute("goalStateIncarnation")
uri=GetNodeTextData(self.Extensions[0].getElementsByTagName("StatusUploadBlob")[0])
except:
Error('Error parsing ExtensionsConfig. Unable to check hearbeat.')
return None
for p in self.Plugins:
if p.getAttribute("state") == 'disabled' or p.getAttribute("state") == 'uninstall' or p.getAttribute("restricted") == 'true' :
continue
version=p.getAttribute("version")
name=p.getAttribute("name")
if p.getAttribute("isJson") != 'true':
Error("Plugin " + name+" version: " +version+" is not a JSON Extension. Skipping.")
continue
try:
if len(p.getAttribute("manifestdata"))<1 or json.loads(p.getAttribute("manifestdata"))[0]['handlerManifest']['reportHeartbeat']!=True :
Error("JSON error, unable to process manifestdata")
def GetCurrentSequenceNumber(self, plugin_base_dir):
"""
Get the settings file with biggest file number in config folder
"""
config_dir = os.path.join(plugin_base_dir, 'config')
seq_no = 0
for subdir, dirs, files in os.walk(config_dir):
for file in files:
try:
cur_seq_no = int(os.path.basename(file).split('.')[0])
if cur_seq_no > seq_no:
seq_no = cur_seq_no
except ValueError:
continue
except:
continue
heartbeat_file=LibDir+'/'+name+'-'+version+'/heartbeat.log'
status_file=LibDir+'/'+name+'-'+version+'/status/'+incarnation+'.status'
if not os.path.exists(heartbeat_file):
Error('Missing '+ heartbeat_file)
continue
return str(seq_no)
def GenerateAggStatus(self, name, version, reportHeartbeat = False):
"""
Generate the status which Azure can understand by the status and heartbeat reported by extension
"""
plugin_base_dir = LibDir+'/'+name+'-'+version+'/'
current_seq_no = self.GetCurrentSequenceNumber(plugin_base_dir)
status_file=os.path.join(plugin_base_dir, 'status/', current_seq_no +'.status')
heartbeat_file = os.path.join(plugin_base_dir, 'heartbeat.log')
handler_state_file = os.path.join(plugin_base_dir, 'config', 'HandlerState')
agg_state = 'NotReady'
handler_state = None
status_obj = None
status_code = None
formatted_message = None
localized_message = None
if os.path.exists(handler_state_file):
handler_state = GetFileContents(handler_state_file).lower()
if HandlerStatusToAggStatus.has_key(handler_state):
agg_state = HandlerStatusToAggStatus[handler_state]
if reportHeartbeat:
if os.path.exists(heartbeat_file):
d=int(time.time()-os.stat(heartbeat_file).st_mtime)
if d > 600 : # not updated for more than 10 min
agg_state = 'Unresponsive'
else:
try:
heartbeat = json.loads(GetFileContents(heartbeat_file))[0]["heartbeat"]
agg_state = heartbeat.get("status")
status_code = heartbeat.get("code")
formatted_message = heartbeat.get("formattedMessage")
localized_message = heartbeat.get("message")
except:
Error("Incorrect heartbeat file. Ignore it. ")
else:
heartbeat=GetFileContents(heartbeat_file)
try:
hb=json.loads(heartbeat)
except:
Error("JSON error, unable to process " + heartbeat_file)
try:
d=int(time.time()-os.stat(heartbeat_file).st_mtime)
except:
Error("Unable to stat " + heartbeat_file)
continue
if d >= 700: # stop sending heartbeats
return 'NotReady'
if d < 120: # within 2 mins considered active
return 'Ready'
if d < 600: # less than 10 mins unknown
state='Unknown'
else: # more than 10 mins with no update considered notready
state='NotReady'
try:
stat_rept='{"handlerName":"' + name + '","handlerVersion":"'+version+ '","status":"' +hb[0]['heartbeat']['status'] + '","code":' + hb[0]['heartbeat']['code'] + ',"formattedMessage":{"lang":"en-US","message":"' + hb[0]['heartbeat']['Message'] + '"}}'
cur_file=status_file+'_current'
with open(cur_file,'w+') as f:
f.write(stat_rept)
# if inc.status exists, rename the inc.status to inc.status_sent
if os.path.exists(status_file) == True:
os.rename(status_file,status_file+'_sent')
# rename inc.status_current to inc.status
os.rename(cur_file,status_file)
# remove inc.status_sent
if os.path.exists(status_file+'_sent') == True:
os.unlink(status_file+'_sent')
except:
Error("Unable to create " + status_file)
continue
agg_state = 'Unresponsive'
#get status file reported by extension
if os.path.exists(status_file):
# raw status generated by extension is an array, get the first item and remove the unnecessary element
try:
status_obj = json.loads(GetFileContents(status_file))[0]
del status_obj["version"]
except:
Error("Incorrect status file. Will NOT settingsStatus in settings. ")
agg_status_obj = {"handlerName": name, "handlerVersion": version, "status": agg_state, "runtimeSettingsStatus" :
{"sequenceNumber": current_seq_no}}
if status_obj:
agg_status_obj["runtimeSettingsStatus"]["settingsStatus"] = status_obj
if status_code != None:
agg_status_obj["code"] = status_code
if formatted_message:
agg_status_obj["formattedMessage"] = formatted_message
if localized_message:
agg_status_obj["message"] = localized_message
agg_status_string = json.dumps(agg_status_obj)
LogIfVerbose("Handler Aggregated Status:" + agg_status_string)
return agg_status_string
def SetHandlerState(self, handler, state=''):
zip_dir=LibDir+"/" + handler
@@ -4625,7 +4628,6 @@ class Agent(Util):
# report the status/heartbeat results of extension processing
if goalState.ExtensionsConfig != None :
goalState.ExtensionsConfig.ReportHandlerStatus()
goalState.ExtensionsConfig.CheckHeartbeat()
time.sleep(25 - sleepToReduceAccessDenied)