diff --git a/scripts/logs_loop/start.py b/scripts/logs_loop/start.py index e4bd1f5..0a9150e 100644 --- a/scripts/logs_loop/start.py +++ b/scripts/logs_loop/start.py @@ -75,117 +75,124 @@ api_instance = client.CoreV1Api() batch_api = client.BatchV1Api() namespace = "kubeinvaders" -#for key in r.scan_iter("log:*"): -# r.delete(key) +for key in r.scan_iter("logs_enabled:*"): + if r.get(key) == "1": -while True: - if not r.exists("log_cleaner"): - if pathlib.Path("/var/www/html/chaoslogs.html").exists(): - os.remove("/var/www/html/chaoslogs.html") - r.set(f"log_cleaner", "foobar") - r.expire(f"log_cleaner", 30) + logid = key.split(":")[1] - logging.info("Loop iteration...") - file = pathlib.Path('/var/www/html/chaoslogs.html') - if not file.exists(): - for key in r.scan_iter("log:*"): - r.delete(key) + while True: + if not r.exists("log_cleaner:{logid}"): + if pathlib.Path(f"/var/www/html/chaoslogs{logid}.html").exists(): + os.remove(f"/var/www/html/chaoslogs{logid}.html") + r.set(f"log_cleaner:{logid}", "1") + r.expire(f"log_cleaner:{logid}", 30) - webtail_pods = [] - final_pod_list = [] - if r.exists("log_pod_regex") and r.exists('logs_enabled'): - logging.info("Found Redis keys for log tail...") - if r.get("logs_enabled") == "1": - logging.info("Found regex log_pod_regex in Redis. Logs from all pods should be collected") - log_pod_regex = r.get("log_pod_regex") - logging.info(f"log_pod_regex is => |{log_pod_regex}|") - try: - api_response = api_instance.list_pod_for_all_namespaces() - except ApiException as e: - logging.info(e) - logging.info(f"Going to search pod compliant with the regex on {len(api_response.items)} pods") + logging.info(f"Loop iteration for log id {logid}") - json_re = json.loads(log_pod_regex) - pod_re = json_re["pod"] - namespace_re = json_re["namespace"] - annotations_re = json_re["annotations"] - labels_re = json_re["labels"] + file = pathlib.Path(f"/var/www/html/chaoslogs{logid}.html") - logging.info(f"Gobal Json Regex is #{json_re}") - logging.info(f"Regex for pod name is #{pod_re}") - logging.info(f"Regex namespace name #{namespace_re}") - logging.info(f"Regex for labels is #{labels_re}") - logging.info(f"Regex for annotation is #{annotations_re}") + if not file.exists(): + for key in r.scan_iter(f"log:{logid}:*"): + r.delete(key) - for pod in api_response.items: - if re.search(f"{pod_re}", pod.metadata.name) and re.search(f"{namespace_re}", pod.metadata.namespace) and re.search(f"{labels_re}", str(pod.metadata.labels)) and re.search(f"{annotations_re}", str(pod.metadata.annotations)): - webtail_pods.append(pod) - logging.info(f"Taking log of {pod.metadata.name} because it is compliant with the regex {log_pod_regex}") + webtail_pods = [] + final_pod_list = [] + if r.exists(f"log_pod_regex:{logid}") and r.exists(f"logs_enabled:{logid}"): + logging.info(f"[logid:{logid}] Found Redis keys for log tail") + if r.get(f"logs_enabled:{logid}") == "1": + logging.info(f"[logid:{logid}] Found regex log_pod_regex in Redis. Logs from all pods should be collected") - try: - api_response = api_instance.list_namespaced_pod(namespace="kubeinvaders") - except ApiException as e: - logging.info(e) + log_pod_regex = r.get(f"log_pod_regex:{logid]}") - webtail_switch = False + logging.info(f"[logid:{logid}] log_pod_regex is => |{log_pod_regex}|") - final_pod_list = webtail_pods + api_response.items - - if len(webtail_pods) > 0: - webtail_switch = True - - for pod in final_pod_list: - if webtail_switch or (pod.metadata.labels.get('approle') != None and pod.metadata.labels['approle'] == 'chaosnode' and pod.status.phase != "Pending"): - try: - latest_log_tail = r.get(f"log_time:{pod.metadata.name}") - logging.info(f"Reading logs of {pod.metadata.name} on {pod.metadata.namespace}") - - if r.exists(f"log_time:{pod.metadata.name}"): - latest_log_tail_time = r.get(f"log_time:{pod.metadata.name}") - else: - latest_log_tail_time = time.time() - logging.info(f"Latest latest_log_tail for {pod.metadata.name} is {latest_log_tail_time}. Current Unix Time is {time.time()}") - - since = int(time.time() - float(latest_log_tail_time)) - - logging.info(f"Diff from time.time() and latest_log_tail_time for {pod.metadata.name} is {since}") - - if since == 0: - since = 1 - - api_response = api_instance.read_namespaced_pod_log(name=pod.metadata.name, namespace=pod.metadata.namespace, tail_lines=1, since_seconds=since) - - if api_response == "": - continue - logging.info(f"API Response: {api_response}") - - logrow = f"
[namespace:{pod.metadata.namespace}][pod:{pod.metadata.name}]
>>>{api_response}
" - - store = False - sha256log = sha256(logrow.encode('utf-8')).hexdigest() - - if r.exists(f"log:{pod.metadata.name}:{sha256log}"): - current_row = r.get(f"log:{pod.metadata.name}:{sha256log}") - if current_row != logrow: - store = True - - if not r.exists(f"log:{pod.metadata.name}:{sha256log}") or store: - file = pathlib.Path('/var/www/html') - if file.exists(): - log_html_file = pathlib.Path('/var/www/html/chaoslogs.html') - line_prepender(log_html_file, logrow) - - r.set(f"log:{pod.metadata.name}:{sha256log}", logrow) - r.set(f"log_time:{pod.metadata.name}", time.time()) - r.expire(f"log:{pod.metadata.name}:{sha256log}", 30) - - logging.info(f"Phase of {pod.metadata.name} is {pod.status.phase}") - if pod.status.phase == "Succeeded" and pod.metadata.labels['approle'] == 'chaosnode': try: - api_response = api_instance.delete_namespaced_pod(pod.metadata.name, namespace = pod.metadata.namespace) - logging.info(f"Deleted pod {pod.metadata.name}") + api_response = api_instance.list_pod_for_all_namespaces() except ApiException as e: logging.info(e) + logging.info(f"[logid:{logid}] Going to search pod compliant with the regex on {len(api_response.items)} pods") + + json_re = json.loads(log_pod_regex) + pod_re = json_re["pod"] + namespace_re = json_re["namespace"] + annotations_re = json_re["annotations"] + labels_re = json_re["labels"] + + logging.info(f"Gobal Json Regex is #{json_re}") + logging.info(f"Regex for pod name is #{pod_re}") + logging.info(f"Regex namespace name #{namespace_re}") + logging.info(f"Regex for labels is #{labels_re}") + logging.info(f"Regex for annotation is #{annotations_re}") + + for pod in api_response.items: + if re.search(f"{pod_re}", pod.metadata.name) and re.search(f"{namespace_re}", pod.metadata.namespace) and re.search(f"{labels_re}", str(pod.metadata.labels)) and re.search(f"{annotations_re}", str(pod.metadata.annotations)): + webtail_pods.append(pod) + logging.info(f"[logid:{logid}] Taking log of {pod.metadata.name} because it is compliant with the regex {log_pod_regex}") + + try: + api_response = api_instance.list_namespaced_pod(namespace="kubeinvaders") except ApiException as e: logging.info(e) - time.sleep(0.5) + + webtail_switch = False + + final_pod_list = webtail_pods + api_response.items + + if len(webtail_pods) > 0: + webtail_switch = True + + for pod in final_pod_list: + if webtail_switch or (pod.metadata.labels.get('approle') != None and pod.metadata.labels['approle'] == 'chaosnode' and pod.status.phase != "Pending"): + try: + latest_log_tail = r.get(f"log_time:{pod.metadata.name}") + logging.info(f"[logid:{logid}] Reading logs of {pod.metadata.name} on {pod.metadata.namespace}") + + if r.exists(f"log_time:{logid}:{pod.metadata.name}"): + latest_log_tail_time = r.get(f"log_time:{logid}:{pod.metadata.name}") + else: + latest_log_tail_time = time.time() + logging.info(f"[logid:{logid}] Latest latest_log_tail for {pod.metadata.name} is {latest_log_tail_time}. Current Unix Time is {time.time()}") + + since = int(time.time() - float(latest_log_tail_time)) + + logging.info(f"[logid:{logid}] Diff from time.time() and latest_log_tail_time for {pod.metadata.name} is {since}") + + if since == 0: + since = 1 + + api_response = api_instance.read_namespaced_pod_log(name=pod.metadata.name, namespace=pod.metadata.namespace, tail_lines=1, since_seconds=since) + + if api_response == "": + continue + logging.info(f"[logid:{logid}] API Response: {api_response}") + + logrow = f"
[namespace:{pod.metadata.namespace}][pod:{pod.metadata.name}]
>>>{api_response}
" + + store = False + sha256log = sha256(logrow.encode('utf-8')).hexdigest() + + if r.exists(f"log:{logid}:{pod.metadata.name}:{sha256log}"): + current_row = r.get(f"log:{logid}:{pod.metadata.name}:{sha256log}") + if current_row != logrow: + store = True + + if not r.exists(f"log:{logid}:{pod.metadata.name}:{sha256log}") or store: + file = pathlib.Path('/var/www/html') + if file.exists(): + log_html_file = pathlib.Path(f"/var/www/html/chaoslogs{logid}.html") + line_prepender(log_html_file, logrow) + + r.set(f"log:{logid}{pod.metadata.name}:{sha256log}", logrow) + r.set(f"log_time:{logid}:{pod.metadata.name}", time.time()) + r.expire(f"log:{logid}:{pod.metadata.name}:{sha256log}", 30) + + logging.info(f"[logid:{logid}] Phase of {pod.metadata.name} is {pod.status.phase}") + if pod.status.phase == "Succeeded" and pod.metadata.labels['approle'] == 'chaosnode': + try: + api_response = 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) + except ApiException as e: + logging.info(e) + time.sleep(0.5)