From 15aa749e340799e1aedae39662d10cb65b4be51e Mon Sep 17 00:00:00 2001 From: Eugenio Marzo Date: Thu, 12 Jan 2023 17:12:39 +0100 Subject: [PATCH] fix logs and metrics --- scripts/chaos-node.lua | 16 +++++++------- scripts/logs_loop/start.py | 14 +++---------- scripts/logs_loop/start.sh | 4 ++-- scripts/metrics_loop/start.py | 35 +++++++++++++++++++++++-------- scripts/metrics_loop/start.sh | 9 +++++++- scripts/pod.lua | 2 +- scripts/programming_mode/start.py | 4 ++-- 7 files changed, 49 insertions(+), 35 deletions(-) diff --git a/scripts/chaos-node.lua b/scripts/chaos-node.lua index a8fdf66..9a88654 100644 --- a/scripts/chaos-node.lua +++ b/scripts/chaos-node.lua @@ -11,7 +11,7 @@ local chaos_container = "" local red = redis:new() local okredis, errredis = red:connect("unix:/tmp/redis.sock") -function read_all(file) +local function read_all(file) local f = assert(io.open(file, "rb")) local content = f:read("*all") f:close() @@ -89,23 +89,21 @@ else chaos_container = config["default_chaos_container"] end -body = [[ +local body = [[ { "apiVersion": "batch/v1", "kind": "Job", "metadata": { "name": "kubeinvaders-chaos-]] .. rand .. [[", "labels": { - "app": "kubeinvaders", - "approle": "chaosnode" + "chaos-controller": "kubeinvaders" } }, "spec": { "template": { "metadata": { "labels": { - "app": "kubeinvaders", - "approle": "chaosnode" + "chaos-controller": "kubeinvaders" } }, "spec": { @@ -125,7 +123,7 @@ local headers2 = { ["Content-Length"] = string.len(body) } -url = k8s_url .. "/apis/batch/v1/namespaces/" .. namespace .. "/jobs" +local url = k8s_url .. "/apis/batch/v1/namespaces/" .. namespace .. "/jobs" ngx.log(ngx.INFO, "Creating chaos_node job kubeinvaders-chaos-" ..rand) local ok, statusCode, headers, statusText = https.request{ @@ -158,7 +156,7 @@ for k,v in ipairs(resp) do decoded = json.decode(v) if decoded["kind"] == "JobList" then for k2,v2 in ipairs(decoded["items"]) do - if v2["status"]["succeeded"] == 1 and v2["metadata"]["labels"]["approle"] == "chaosnode" then + if v2["status"]["succeeded"] == 1 and v2["metadata"]["labels"]["chaos-controller"] == "kubeinvaders" then delete_job = "kubectl delete job " .. v2["metadata"]["name"] .. " --token=" .. token .. " --server=" .. k8s_url .. " --insecure-skip-tls-verify=true -n " .. namespace ngx.log(ngx.INFO, delete_pod) end @@ -184,7 +182,7 @@ for k,v in ipairs(resp) do decoded = json.decode(v) if decoded["kind"] == "PodList" then for k2,v2 in ipairs(decoded["items"]) do - if v2["status"]["phase"] == "Succeeded" and v2["metadata"]["labels"]["approle"] == "chaosnode" then + if v2["status"]["phase"] == "Succeeded" and v2["metadata"]["labels"]["chaos-controller"] == "kubeinvaders" then delete_pod = "kubectl delete pod " .. v2["metadata"]["name"] .. " --token=" .. token .. " --server=" .. k8s_url .. " --insecure-skip-tls-verify=true -n " .. namespace ngx.log(ngx.INFO, delete_pod) end diff --git a/scripts/logs_loop/start.py b/scripts/logs_loop/start.py index a36d13c..ec4057a 100644 --- a/scripts/logs_loop/start.py +++ b/scripts/logs_loop/start.py @@ -37,12 +37,12 @@ def create_pod_list(logid, api_response_items, current_regex): webtail_pods.append(pod) regex_match_info = f"[logid:{logid}][k-inv][logs-loop] Taking logs of {pod.metadata.name}. Redis has cached that {current_regex} is good for {pod.metadata.name}" r.set(f"log_status:{logid}", regex_match_info) - logging.info(f"[logid:{logid}][k-in][regexmatch] |{namespace_re}| |{pod.metadata.namespace}| IS CHACHED IN REDIS") + logging.info(f"[k-inv][regexmatch][logid:{logid}][{cached_regex_match}] IS CHACHED IN REDIS") else: - regex_match_info = f"[logid:{logid}][logs-loop] Skipping logs of {pod.metadata.name}. Redis has cached that {current_regex} is not good for {pod.metadata.name}" + regex_match_info = f"[logs-loop][logid:{logid}] Skipping logs of {pod.metadata.name}. Redis has cached that {current_regex} is not good for {pod.metadata.name}" logging.debug(regex_match_info) - logging.info(f"[logid:{logid}][k-inv][regexmatch] |{namespace_re}| |{pod.metadata.namespace}| IS CHACHED IN REDIS") + logging.info(f"[k-inv][regexmatch][logid:{logid}][{cached_regex_match}] IS CHACHED IN REDIS") else: if re.search(f"{pod_re}", pod.metadata.name) or re.search(r"{pod_re}", pod.metadata.name): @@ -139,14 +139,6 @@ def compute_line(api_response_line, container): r.set(f"log_time:{logid}:{pod.metadata.name}:{container}", time.time()) r.expire(f"log:{logid}:{pod.metadata.name}:{container}:{sha256log}", 30) - # logging.info(f"[logid:{logid}] Phase of {pod.metadata.name} is {pod.status.phase}") - # if pod.status.phase == "Succeeded" and 'approle' in pod.metadata.labels and pod.metadata.labels['approle'] == 'chaosnode': - # try: - # api_instance.delete_namespaced_pod(pod.metadata.name, namespace = pod.metadata.namespace) - # logging.info(f"[logid:{logid}] Deleted pod {pod.metadata.name}") - # except ApiException as e: - # logging.info(e) - logging.basicConfig(level=logging.INFO) logging.info('Starting script for KubeInvaders taking logs from pods...') diff --git a/scripts/logs_loop/start.sh b/scripts/logs_loop/start.sh index df77a07..018edd7 100755 --- a/scripts/logs_loop/start.sh +++ b/scripts/logs_loop/start.sh @@ -13,7 +13,7 @@ python3 /opt/logs_loop/start.py https://${KUBERNETES_SERVICE_HOST}:${KUBERNETES_ while true do - ( ps -ef | grep start.py | grep logs_loop | grep -v grep ) || echo "Error. Logs Loop is down..." - ( ps -ef | grep start.py | grep logs_loop | grep -v grep ) || ( python3 /opt/logs_loop/start.py https://${KUBERNETES_SERVICE_HOST}:${KUBERNETES_SERVICE_PORT_HTTPS} & ) + ( ps -ef | grep start.py | grep logs_loop | grep -v grep &> /dev/null) || echo "Error. Logs Loop is down..." + ( ps -ef | grep start.py | grep logs_loop | grep -v grep &> /dev/null) || ( python3 /opt/logs_loop/start.py https://${KUBERNETES_SERVICE_HOST}:${KUBERNETES_SERVICE_PORT_HTTPS} & ) sleep 2 done diff --git a/scripts/metrics_loop/start.py b/scripts/metrics_loop/start.py index d6c2e0b..a069b09 100644 --- a/scripts/metrics_loop/start.py +++ b/scripts/metrics_loop/start.py @@ -12,6 +12,7 @@ import random import redis import time import urllib3 +import time urllib3.disable_warnings(urllib3.exceptions.InsecureRequestWarning) def create_container(image, name, command, args): @@ -33,13 +34,13 @@ def create_container(image, name, command, args): def create_pod_template(pod_name, container, job_name): pod_template = client.V1PodTemplateSpec( spec=client.V1PodSpec(restart_policy="Never", containers=[container]), - metadata=client.V1ObjectMeta(name=pod_name, labels={"pod-name-prefix": pod_name, "approle": "chaosnode", "job-name": job_name}), + metadata=client.V1ObjectMeta(name=pod_name, labels={"chaos-controller": "kubeinvaders", "job-name": job_name}), ) return pod_template def create_job(job_name, pod_template): - metadata = client.V1ObjectMeta(name=job_name, labels={"job-name": job_name, "approle": "chaosnode"}) + metadata = client.V1ObjectMeta(name=job_name, labels={"chaos-controller": "kubeinvaders"}) job = client.V1Job( api_version="batch/v1", @@ -54,7 +55,7 @@ def create_job(job_name, pod_template): r = redis.Redis(unix_socket_path='/tmp/redis.sock') # create logger -logging.basicConfig(level=os.environ.get("LOGLEVEL", "INFO")) +logging.basicConfig(level=os.environ.get("LOGLEVEL", "DEBUG")) logging.info('Starting script for KubeInvaders programming mode') configuration = client.Configuration() @@ -70,24 +71,40 @@ client.Configuration.set_default(configuration) api_instance = client.CoreV1Api() batch_api = client.BatchV1Api() -namespace = "kubeinvaders" while True: try: - api_response = api_instance.list_namespaced_pod(namespace="kubeinvaders") + label_selector="chaos-controller=kubeinvaders" + api_response = api_instance.list_pod_for_all_namespaces(label_selector=label_selector) except ApiException as e: logging.info(e) r.set("current_chaos_job_pod", 0) for pod in api_response.items: - if pod.metadata.labels.get('approle') != None and pod.metadata.labels['approle'] == 'chaosnode': - if pod.status.phase == "Pending" or pod.status.phase == "Running": - r.incr('current_chaos_job_pod') + if pod.status.phase == "Pending" or pod.status.phase == "Running": + logging.info(f"[k-inv][metrics_loop] Found pod {pod.metadata.name}. It is in {pod.status.phase} phase. Incrementing current_chaos_job_pod Redis key") + r.incr('current_chaos_job_pod') + + if pod.status.phase != "Pending" and pod.status.phase != "Running" and not r.exists(f"pod:time:{pod.metadata.namespace}:{pod.metadata.name}"): + logging.info(f"[k-inv][metrics_loop] Found pod {pod.metadata.name}. It is in {pod.status.phase} phase. Tracking time in pod:time:{pod.metadata.namespace}:{pod.metadata.name} Redis key") + r.set(f"pod:time:{pod.metadata.namespace}:{pod.metadata.name}", int(time.time())) + + elif pod.status.phase != "Pending" and pod.status.phase != "Running" and r.exists(f"pod:time:{pod.metadata.namespace}:{pod.metadata.name}"): + logging.info(f"[k-inv][metrics_loop] Found pod {pod.metadata.name}. It is in {pod.status.phase} phase. Comparing time in pod:time:{pod.metadata.namespace}:{pod.metadata.name} Redis key with now") + now = int(time.time()) + pod_time = int(r.get(f"pod:time:{pod.metadata.namespace}:{pod.metadata.name}")) + logging.info(f"[k-inv][metrics_loop] For {pod.metadata.name} comparing now:{now} with pod_time:{pod_time}") + if (now - pod_time > 30): + try: + api_instance.delete_namespaced_pod(pod.metadata.name, namespace = pod.metadata.namespace) + logging.info(f"[k-inv][metrics_loop] Deleting pod {pod.metadata.name}") + r.delete(f"pod:time:{pod.metadata.namespace}:{pod.metadata.name}") + except ApiException as e: + logging.info(e) if pod.metadata.labels.get('chaos-codename') != None: codename = pod.metadata.labels.get('chaos-codename') job_name = pod.metadata.labels.get('job-name') exp_name = pod.metadata.labels.get('experiment-name') r.set(f"chaos_jobs_status:{codename}:{exp_name}:{job_name}", pod.status.phase) - time.sleep(1) diff --git a/scripts/metrics_loop/start.sh b/scripts/metrics_loop/start.sh index 00177e7..9163e44 100755 --- a/scripts/metrics_loop/start.sh +++ b/scripts/metrics_loop/start.sh @@ -8,4 +8,11 @@ else export TOKEN="$(cat /var/run/secrets/kubernetes.io/serviceaccount/token)" fi -python3 /opt/metrics_loop/start.py https://${KUBERNETES_SERVICE_HOST}:${KUBERNETES_SERVICE_PORT_HTTPS} +python3 /opt/metrics_loop/start.py https://${KUBERNETES_SERVICE_HOST}:${KUBERNETES_SERVICE_PORT_HTTPS} & + +while true +do + ( ps -ef | grep start.py | grep metrics_loop | grep -v grep &> /dev/null) || echo "Error. Metrics Loop is down..." + ( ps -ef | grep start.py | grep metrics_loop | grep -v grep &> /dev/null) || (python3 /opt/metrics_loop/start.py https://${KUBERNETES_SERVICE_HOST}:${KUBERNETES_SERVICE_PORT_HTTPS} & ) + sleep 2 +done \ No newline at end of file diff --git a/scripts/pod.lua b/scripts/pod.lua index 5024f91..1986753 100644 --- a/scripts/pod.lua +++ b/scripts/pod.lua @@ -107,7 +107,7 @@ if action == "list" then decoded = json.decode(v) if decoded["kind"] == "PodList" then for k2,v2 in ipairs(decoded["items"]) do - if v2["status"]["phase"] == "Running" and v2["metadata"]["labels"]["approle"] ~= "chaosnode" then + if v2["status"]["phase"] == "Running" and v2["metadata"]["labels"]["chaos-controller"] ~= "kubeinvaders" then ngx.log(ngx.INFO, "found pod " .. v2["metadata"]["name"]) pods["items"][i] = v2["metadata"]["name"] i = i + 1 diff --git a/scripts/programming_mode/start.py b/scripts/programming_mode/start.py index 3a3d4ee..2fc9dd9 100644 --- a/scripts/programming_mode/start.py +++ b/scripts/programming_mode/start.py @@ -28,7 +28,7 @@ def create_container(image, name, command, args): return container def create_pod_template(pod_name, additional_labels, container, exp_name): - pod_labels = {"app": "kubeinvaders", "approle": "chaosnode", "experiment-name": exp_name} + pod_labels = {"chaos-controller": "kubeinvaders", "experiment-name": exp_name} pod_labels.update(additional_labels) pod_template = client.V1PodTemplateSpec( spec=client.V1PodSpec(restart_policy="Never", containers=[container]), @@ -38,7 +38,7 @@ def create_pod_template(pod_name, additional_labels, container, exp_name): return pod_template def create_job(job_name, pod_template): - metadata = client.V1ObjectMeta(name=job_name, labels={"job-name": job_name, "approle": "chaosnode"}) + metadata = client.V1ObjectMeta(name=job_name, labels={"chaos-controller": "kubeinvaders"}) job = client.V1Job( api_version="batch/v1",