From cd0cf6443690a14dbcbef74477dad953487c20b9 Mon Sep 17 00:00:00 2001 From: Eugenio Marzo Date: Wed, 25 Jan 2023 00:20:33 +0100 Subject: [PATCH] added since for log collector --- html5/index.html | 6 +- nginx/KubeInvaders.conf | 2 +- scripts/logs_loop/start.py | 189 ++++++++++++++++++++++--------------- 3 files changed, 116 insertions(+), 81 deletions(-) diff --git a/html5/index.html b/html5/index.html index 69c7301..2a74e10 100644 --- a/html5/index.html +++ b/html5/index.html @@ -285,7 +285,7 @@ jobs: image: docker.io/luckysideburn/kubeinvaders-stress-ng:latest command: "stress-ng" args: - - --help + - --version mem-attack-job: additional-labels: @@ -295,7 +295,7 @@ jobs: image: docker.io/luckysideburn/kubeinvaders-stress-ng:latest command: "stress-ng" args: - - --help + - --version experiments: - name: cpu-attack-exp @@ -365,7 +365,7 @@ experiments:
- +
diff --git a/nginx/KubeInvaders.conf b/nginx/KubeInvaders.conf index 272b369..1070515 100644 --- a/nginx/KubeInvaders.conf +++ b/nginx/KubeInvaders.conf @@ -183,7 +183,7 @@ server { local logid = args["logid"] red:set("do_not_clean_log:" .. logid, "1") - red:expire("do_not_clean_log:" .. logid, "10") + red:expire("do_not_clean_log:" .. logid, "20") red:set("logs_enabled:" .. logid, "1") red:expire("logs_enabled:" .. logid, "10") diff --git a/scripts/logs_loop/start.py b/scripts/logs_loop/start.py index 7843d5c..a8f4d2a 100644 --- a/scripts/logs_loop/start.py +++ b/scripts/logs_loop/start.py @@ -17,9 +17,11 @@ import re from hashlib import sha256 import time import urllib3 +import datetime urllib3.disable_warnings(urllib3.exceptions.InsecureRequestWarning) -def create_pod_list(logid, api_response_items, current_regex): +def create_pod_list(logid, api_responses, current_regex): + webtail_pods = [] json_re = json.loads(current_regex) regexsha = sha256(current_regex.encode('utf-8')).hexdigest() pod_re = json_re["pod"] @@ -27,50 +29,67 @@ def create_pod_list(logid, api_response_items, current_regex): annotations_re = json_re["annotations"] labels_re = json_re["labels"] containers_re = json_re["containers"] - webtail_pods = [] + + for api_response in api_responses: + pods_pending = 0 + pods_running = 0 + pods_succeeded = 0 - for pod in api_response_items: - if r.exists(f"regex_cmp:{regexsha}:{logid}:{pod.metadata.namespace}:{pod.metadata.name}"): - cached_regex_match = r.get(f"regex_cmp:{regexsha}:{logid}:{pod.metadata.namespace}:{pod.metadata.name}") - if cached_regex_match == "maching": - 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.debug(f"[k-inv][regexmatch][logid:{logid}][{cached_regex_match}] IS CHACHED IN REDIS") + for pod in api_response.items: + if pod.status.phase == "Pending": + pods_pending = pods_pending + 1 + if pod.status.phase == "Running": + pods_running = pods_running + 1 + if pod.status.phase == "Succeeded": + pods_succeeded = pods_succeeded + 1 + + # old_logs = r.get(f"logs:chaoslogs-{logid}") + # r.set(f"logs:chaoslogs-{logid}", f"
[k-inv] pods on Pending phase: {pods_pending}
[k-inv] pods on Succeeded phase: {pods_succeeded}
[k-inv] pods on Running phase: {pods_running}
{old_logs}") + + # if pod.status.phase != "Succedeed" and pod.status.phase != "Running": + # continue + + if r.exists(f"regex_cmp:{regexsha}:{logid}:{pod.metadata.namespace}:{pod.metadata.name}"): + cached_regex_match = r.get(f"regex_cmp:{regexsha}:{logid}:{pod.metadata.namespace}:{pod.metadata.name}") + if cached_regex_match == "maching": + 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.debug(f"[k-inv][regexmatch][logid:{logid}][{cached_regex_match}] IS CHACHED IN REDIS") + + else: + 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.debug(f"[k-inv][regexmatch][logid:{logid}][{cached_regex_match}] IS CHACHED IN REDIS") else: - 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.debug(f"[k-inv][regexmatch][logid:{logid}][{cached_regex_match}] IS CHACHED IN REDIS") + if re.search(f"{pod_re}", pod.metadata.name) or re.search(r"{pod_re}", pod.metadata.name): + #logging.debug(f"[logid:{logid}][k-in][regexmatch] |{pod_re}| |{pod.metadata.name}| MATCHED") + regex_key_name = f"regex_cmp:{regexsha}:{logid}:{pod.metadata.namespace}:{pod.metadata.name}" - else: - if re.search(f"{pod_re}", pod.metadata.name) or re.search(r"{pod_re}", pod.metadata.name): - #logging.debug(f"[logid:{logid}][k-in][regexmatch] |{pod_re}| |{pod.metadata.name}| MATCHED") - regex_key_name = f"regex_cmp:{regexsha}:{logid}:{pod.metadata.namespace}:{pod.metadata.name}" - - if re.search(f"{namespace_re}", pod.metadata.namespace) or re.search(r"{namespace_re}", pod.metadata.namespace): - #logging.debug(f"[logid:{logid}][k-inv][regexmatch] |{namespace_re}| |{pod.metadata.namespace}| MATCHED") - if re.search(f"{labels_re}", str(pod.metadata.labels)) or re.search(r"{labels_re}", str(pod.metadata.labels)): - #logging.debug(f"[logid:{logid}][k-inv][regexmatch] |{labels_re}| |{str(pod.metadata.labels)}| MATCHED") - if re.search(f"{annotations_re}", str(pod.metadata.annotations)) or re.search(r"{annotations_re}", str(pod.metadata.annotations)): - #logging.debug(f"[logid:{logid}][k-inv][regexmatch] |{annotations_re}| |{str(pod.metadata.annotations)}| MATCHED") - webtail_pods.append(pod) - regex_match_info = f"[logid:{logid}] Taking logs from {pod.metadata.name}. It is compliant with the Regex {current_regex}" - r.set(regex_key_name, "maching") - logging.debug(regex_match_info) - r.set(f"log_status:{logid}", regex_match_info) + if re.search(f"{namespace_re}", pod.metadata.namespace) or re.search(r"{namespace_re}", pod.metadata.namespace): + #logging.debug(f"[logid:{logid}][k-inv][regexmatch] |{namespace_re}| |{pod.metadata.namespace}| MATCHED") + if re.search(f"{labels_re}", str(pod.metadata.labels)) or re.search(r"{labels_re}", str(pod.metadata.labels)): + logging.debug(f"[logid:{logid}][k-inv][regexmatch] |{labels_re}| |{str(pod.metadata.labels)}| MATCHED") + if re.search(f"{annotations_re}", str(pod.metadata.annotations)) or re.search(r"{annotations_re}", str(pod.metadata.annotations)): + logging.debug(f"[logid:{logid}][k-inv][regexmatch] |{annotations_re}| |{str(pod.metadata.annotations)}| MATCHED") + webtail_pods.append(pod) + regex_match_info = f"[logid:{logid}] Taking logs from {pod.metadata.name}. It is compliant with the Regex {current_regex}" + r.set(regex_key_name, "maching") + logging.debug(regex_match_info) + r.set(f"log_status:{logid}", regex_match_info) + else: + logging.debug(f"[logid:{logid}][k-inv][regexmatch] |{annotations_re}| |{str(pod.metadata.annotations)}| FAILED") + r.set(regex_key_name, "not_maching") else: - logging.debug(f"[logid:{logid}][k-inv][regexmatch] |{annotations_re}| |{str(pod.metadata.annotations)}| FAILED") + logging.debug(f"[logid:{logid}][k-inv][regexmatch] |{labels_re}| |{str(pod.metadata.labels)}| FAILED") r.set(regex_key_name, "not_maching") else: - logging.debug(f"[logid:{logid}][k-inv][regexmatch] |{labels_re}| |{str(pod.metadata.labels)}| FAILED") + logging.debug(f"[logid:{logid}][k-inv][regexmatch] |{namespace_re}| |{pod.metadata.namespace}| FAILED") r.set(regex_key_name, "not_maching") else: - logging.debug(f"[logid:{logid}][k-inv][regexmatch] |{namespace_re}| |{pod.metadata.namespace}| FAILED") + logging.debug(f"[logid:{logid}][k-inv][regexmatch] |{pod_re}| |{pod.metadata.name}| FAILED") r.set(regex_key_name, "not_maching") - else: - logging.debug(f"[logid:{logid}][k-inv][regexmatch] |{pod_re}| |{pod.metadata.name}| FAILED") - r.set(regex_key_name, "not_maching") return webtail_pods def log_cleaner(logid): @@ -81,7 +100,8 @@ def log_cleaner(logid): for key in r.scan_iter(f"log:{logid}:*"): r.delete(key) - r.set(f"logs:chaoslogs-{logid}", "
[k-inv] Logs has been cleaned...
") + # old_logs = r.get(f"logs:chaoslogs-{logid}") + # r.set(f"logs:chaoslogs-{logid}", f"
[k-inv] Logs has been cleaned...
{old_logs}") r.set(f"do_not_clean_log:{logid}", "1") r.expire(f"do_not_clean_log:{logid}", 60) @@ -136,7 +156,7 @@ def compute_line(api_response_line, container): r.set(f"log:{logid}:{pod.metadata.name}:{container}:{sha256log}", logrow) r.set(f"log_time:{logid}:{pod.metadata.name}:{container}", time.time()) - r.expire(f"log:{logid}:{pod.metadata.name}:{container}:{sha256log}", 30) + r.expire(f"log:{logid}:{pod.metadata.name}:{container}:{sha256log}", 60) logging.basicConfig(level=os.environ.get("LOGLEVEL", "DEBUG")) logging.getLogger('kubernetes').setLevel(logging.ERROR) @@ -152,7 +172,7 @@ else: if os.environ.get("DEV"): logging.debug("Setting env var for dev...") - r.set("log_pod_regex", '{"pod":".*", "namespace":"namespace1", "labels":".*", "annotations":".*", "containers": ".*"}') + r.set("log_pod_regex", '{"since": 60, "pod":".*", "namespace":"namespace1", "labels":".*", "annotations":".*", "containers": ".*"}') r.set("logs_enabled:aaaa", 1) r.expire("logs_enabled:aaaa", 10) r.set("programming_mode", 0) @@ -189,43 +209,42 @@ while True: r.set(f"log_status:{logid}", f"[k-inv][logs-loop] {key} is using this regex: {current_regex}") logging.debug(f"[logid:{logid}] Checking do_not_clean_log Redis key") - log_cleaner(logid) try: - api_response = api_instance.list_pod_for_all_namespaces() + json_re = json.loads(current_regex) + namespace_re = json_re["namespace"] + logging.debug(f"[logid:{logid}][k-inv][logs-loop] Taking list of namespaces") + namespaces_list = api_instance.list_namespace() + api_responses = [] + + for namespace in namespaces_list.items: + #logging.debug(f"[logid:{logid}][k-inv][logs-loop] Found namespace {namespace.metadata.name}") + if re.search(f"{namespace_re}", namespace.metadata.name): + logging.debug(f"[logid:{logid}][k-inv][logs-loop][NAMESPACE-MATCHING] {namespace.metadata.name}") + logging.debug(f"[logid:{logid}][k-inv][logs-loop] Taking pods from namespace {namespace.metadata.name}") + api_responses.append(api_instance.list_namespaced_pod(namespace.metadata.name)) + except ApiException as e: logging.debug(e) - pods_found_info = f"[logid:{logid}][k-inv][logs-loop] Looking for pods compliant with the current regex. Scanning {len(api_response.items)} pods" - r.set(f"log_status:{logid}", pods_found_info) + #pods_found_info = f"[logid:{logid}][k-inv][logs-loop] Looking for pods compliant with the current regex. Scanning {len(api_response.items)} pods" + #r.set(f"log_status:{logid}", pods_found_info) - webtail_pods = create_pod_list(logid, api_response.items, current_regex) + webtail_pods = create_pod_list(logid, api_responses, current_regex) json_re = json.loads(current_regex) containers_re = json_re["containers"] webtail_pods_len = len(webtail_pods) old_logs = r.get(f"logs:chaoslogs-{logid}") + user_since = int(json_re["since"]) - if r.exists(f"logs:webtail_pods_len:{logid}") and str(r.get(f"logs:webtail_pods_len:{logid}")) != str(webtail_pods_len): - r.set(f"logs:chaoslogs-{logid}", f"
[k-inv] K-inv found {webtail_pods_len} pods to read logs from
{old_logs}") + + # if r.exists(f"logs:webtail_pods_len:{logid}") and str(r.get(f"logs:webtail_pods_len:{logid}")) != str(webtail_pods_len): + # r.set(f"logs:chaoslogs-{logid}", f"
[k-inv] K-inv found {webtail_pods_len} pods to read logs from
{old_logs}") r.set(f"logs:webtail_pods_len:{logid}", webtail_pods_len) r.set(f"pods_match_regex:{logid}", webtail_pods_len) logging.debug(f"[logid:{logid}][k-inv][logs-loop] Current Regex: {current_regex}") - - pods_pending = 0 - pods_running = 0 - pods_succeeded = 0 - - for pod in webtail_pods: - if pod.status.phase == "Pending": - pods_pending = pods_pending + 1 - if pod.status.phase == "Running": - pods_running = pods_running + 1 - if pod.status.phase == "Succeeded": - pods_succeeded = pods_succeeded + 1 - - r.set(f"logs:chaoslogs-{logid}", f"
[k-inv] pods on Pending phase: {pods_pending}
[k-inv] pods on Succeeded phase: {pods_succeeded}
[k-inv] pods on Running phase: {pods_running}
") for pod in webtail_pods: if pod.status.phase == "Unknown" and pod.status.phase == "Pending": @@ -236,6 +255,11 @@ while True: # if "containers_re" in locals() or "containers_re" in globals(): # if re.search(f"{containers_re}", container.name): container_list.append(container.name) + + if len(container_list) == 1: + only_one_container = True + else: + only_one_container = False for container in container_list: logging.debug(f"[logid:{logid}][k-inv][logs-loop] Listing containers of {pod.metadata.name}. Computing {container} phase: {pod.status.phase}") @@ -244,39 +268,50 @@ while True: try: if r.exists(f"log_time:{logid}:{pod.metadata.name}:{container}"): latest_log_tail_time = float(r.get(f"log_time:{logid}:{pod.metadata.name}:{container}")) + since = int(time.time() - float(latest_log_tail_time)) + 1 + else: latest_log_tail_time = time.time() + pod_start_time = int(datetime.datetime.timestamp(pod.status.start_time)) + logging.debug(f"[logid:{logid}][k-inv][logs-loop] POD's start time {pod_start_time}") + since = int(time.time() - pod_start_time) + 1 - since = int(time.time() - float(latest_log_tail_time)) + 1 - logging.debug(f"[logid:{logid}][k-inv][logs-loop] Time types: {type(latest_log_tail_time)} {type(time.time())} {type(since)} since={since}") + logging.debug(f"[logid:{logid}][k-inv][logs-loop] user_since is {user_since}") - if since == 0: - since = 1 - + if since > user_since: + continue + + logging.debug(f"[logid:{logid}][k-inv][logs-loop] Time types: {type(latest_log_tail_time)} {type(time.time())} {type(since)} since={since}") logging.debug(f"[logid:{logid}][k-inv][logs-loop] Calling K8s API for reading logs of {pod.metadata.name} container {container} in namespace {pod.metadata.namespace} since {since} seconds - phase {pod.status.phase}") - api_response = api_instance.read_namespaced_pod_log(name=pod.metadata.name, namespace=pod.metadata.namespace, since_seconds=since, container=container, tail_lines=since) + if only_one_container: + api_response = api_instance.read_namespaced_pod_log(name=pod.metadata.name, namespace=pod.metadata.namespace, since_seconds=since) + else: + api_response = api_instance.read_namespaced_pod_log(name=pod.metadata.name, namespace=pod.metadata.namespace, since_seconds=since, container=container) logging.debug(f"[logid:{logid}][k-inv][logs-loop] Computing K8s API response for reading logs of {pod.metadata.name} in namespace {pod.metadata.namespace} - phase {pod.status.phase}") logging.debug(f"[logid:{logid}][k-inv][logs-loop] {type(api_response)} {api_response}") r.set(f"log_time:{logid}:{pod.metadata.name}:{container}", time.time()) - k = 5 + # k = 5 - regex_return = re.search(r'[\w]+', api_response) - logging.debug(f"[logid:{logid}][k-inv][logs-loop] Regex on api_response: {regex_return}") + # regex_return = re.search(r'[\w]+', api_response) + # logging.debug(f"[logid:{logid}][k-inv][logs-loop] Regex on api_response: {regex_return}") - while not re.search(r'[\w]+', api_response) and k < 60: - logging.debug(f"[logid:{logid}][k-inv][logs-loop][logs collector attempt {k}] Calling K8s API for reading logs of {pod.metadata.name} container {container} in namespace {pod.metadata.namespace} since {since} seconds - phase {pod.status.phase}") - api_response = api_instance.read_namespaced_pod_log(name=pod.metadata.name, namespace=pod.metadata.namespace, since_seconds=since, container=container, tail_lines=since) - logging.debug(f"[logid:{logid}][k-inv][logs-loop] API Response: {api_response}") - # regex_return = re.search(r'[\w]+', api_response) - # logging.debug(f"[logid:{logid}][k-inv][logs-loop] Regex on api_response: {regex_return}") + # while not re.search(r'[\w]+', api_response) and k < 60: + # logging.debug(f"[logid:{logid}][k-inv][logs-loop][logs collector attempt {k}] Calling K8s API for reading logs of {pod.metadata.name} container {container} in namespace {pod.metadata.namespace} since {since} seconds - phase {pod.status.phase}") + # if only_one_container: + # api_response = api_instance.read_namespaced_pod_log(name=pod.metadata.name, namespace=pod.metadata.namespace, since_seconds=since) + # else: + # api_response = api_instance.read_namespaced_pod_log(name=pod.metadata.name, namespace=pod.metadata.namespace, since_seconds=since, container=container) + # logging.debug(f"[logid:{logid}][k-inv][logs-loop] API Response: {api_response}") + # # regex_return = re.search(r'[\w]+', api_response) + # # logging.debug(f"[logid:{logid}][k-inv][logs-loop] Regex on api_response: {regex_return}") - since = since + 1 - k = k + 1 - time.sleep(0.5) + # since = since + 1 + # k = k + 1 + # time.sleep(0.5) if not re.search(r'[\w]+', api_response): # regex_return = re.search(r'[\w]+', api_response)