<?xml version="1.0" encoding="UTF-8"?>
<rss version="2.0" xmlns:atom="http://www.w3.org/2005/Atom" xmlns:dc="http://purl.org/dc/elements/1.1/">
  <channel>
    <title>DEV Community: Sachin Singhal</title>
    <description>The latest articles on DEV Community by Sachin Singhal (@sachin_xp).</description>
    <link>https://dev.to/sachin_xp</link>
    <image>
      <url>https://media2.dev.to/dynamic/image/width=90,height=90,fit=cover,gravity=auto,format=auto/https:%2F%2Fdev-to-uploads.s3.us-east-2.amazonaws.com%2Fuploads%2Fuser%2Fprofile_image%2F2613642%2Fa4d59948-1f87-47dd-bf4c-7a7dd9b039ca.png</url>
      <title>DEV Community: Sachin Singhal</title>
      <link>https://dev.to/sachin_xp</link>
    </image>
    <atom:link rel="self" type="application/rss+xml" href="https://dev.to/feed/sachin_xp"/>
    <language>en</language>
    <item>
      <title>Fampay Solution#1</title>
      <dc:creator>Sachin Singhal</dc:creator>
      <pubDate>Wed, 25 Dec 2024 14:21:46 +0000</pubDate>
      <link>https://dev.to/sachin_xp/fampay-solution1-2m7g</link>
      <guid>https://dev.to/sachin_xp/fampay-solution1-2m7g</guid>
      <description>&lt;ul&gt;
&lt;li&gt;AI-Powered Severity Analysis: Classifies logs using BERT.&lt;/li&gt;
&lt;li&gt;Anomaly Detection: Identifies unusual patterns with KMeans clustering.&lt;/li&gt;
&lt;li&gt;Self-Healing Automation: Executes scripts to mitigate critical issues.&lt;/li&gt;
&lt;li&gt;Slack Integration: Sends real-time notifications for critical logs.&lt;/li&gt;
&lt;li&gt;Prometheus Monitoring: Tracks bugs and RCA metrics.&lt;/li&gt;
&lt;/ul&gt;

&lt;p&gt;import logging&lt;br&gt;
import os&lt;br&gt;
import json&lt;br&gt;
import pandas as pd&lt;br&gt;
import numpy as np&lt;br&gt;
from kafka import KafkaConsumer&lt;br&gt;
from elasticsearch import Elasticsearch&lt;br&gt;
from transformers import AutoTokenizer, TFAutoModel&lt;br&gt;
from tensorflow.keras.models import Model&lt;br&gt;
from sklearn.cluster import KMeans&lt;br&gt;
from causalinference import CausalModel&lt;br&gt;
from prometheus_client import start_http_server, Counter&lt;br&gt;
from slack_sdk import WebClient&lt;br&gt;
from slack_sdk.errors import SlackApiError&lt;br&gt;
import subprocess&lt;br&gt;
import requests&lt;/p&gt;

&lt;p&gt;logging.basicConfig(level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s')&lt;/p&gt;

&lt;p&gt;bug_detected_counter = Counter("bug_detected", "Number of bugs detected in logs")&lt;br&gt;
root_cause_counter = Counter("root_cause_analysis", "Number of root causes identified")&lt;/p&gt;

&lt;p&gt;KAFKA_BROKER = os.getenv("KAFKA_BROKER", "localhost:9092")&lt;br&gt;
TOPIC = os.getenv("KAFKA_TOPIC", "application_logs")&lt;/p&gt;

&lt;p&gt;consumer = KafkaConsumer(&lt;br&gt;
    TOPIC,&lt;br&gt;
    bootstrap_servers=[KAFKA_BROKER],&lt;br&gt;
    value_deserializer=lambda x: json.loads(x.decode('utf-8'))&lt;br&gt;
)&lt;/p&gt;

&lt;p&gt;ELASTIC_HOST = os.getenv("ELASTIC_HOST", "&lt;a href="http://localhost:9200%22" rel="noopener noreferrer"&gt;http://localhost:9200"&lt;/a&gt;)&lt;br&gt;
es = Elasticsearch([ELASTIC_HOST])&lt;/p&gt;

&lt;p&gt;SLACK_TOKEN = os.getenv("SLACK_TOKEN")&lt;br&gt;
SLACK_CHANNEL = os.getenv("SLACK_CHANNEL", "#bug-notifications")&lt;br&gt;
slack_client = WebClient(token=SLACK_TOKEN)&lt;/p&gt;

&lt;p&gt;tokenizer = AutoTokenizer.from_pretrained("bert-base-uncased")&lt;br&gt;
bert_model = TFAutoModel.from_pretrained("bert-base-uncased")&lt;/p&gt;

&lt;p&gt;def root_cause_analysis(data):&lt;br&gt;
    logging.info("Performing Root Cause Analysis...")&lt;br&gt;
    causal = CausalModel(&lt;br&gt;
        Y=data["severity_label"].values, &lt;br&gt;
        D=data["source_system"].values,  # Adjust column names as per your logs&lt;br&gt;
        X=data.drop(columns=["severity_label", "source_system"]).values&lt;br&gt;
    )&lt;br&gt;
    causal.est_via_matching()&lt;br&gt;
    root_cause_counter.inc()&lt;br&gt;
    return causal.estimates&lt;/p&gt;

&lt;p&gt;def cluster_anomalies(logs):&lt;br&gt;
    logging.info("Clustering logs for anomalies...")&lt;br&gt;
    kmeans = KMeans(n_clusters=3, random_state=42)&lt;br&gt;
    clusters = kmeans.fit_predict(logs)&lt;br&gt;
    anomalies = logs[clusters == 2]  # Example: Cluster 2 as anomalies&lt;br&gt;
    return anomalies&lt;/p&gt;

&lt;h1&gt;
  
  
  Notification System
&lt;/h1&gt;

&lt;p&gt;def notify_slack(message):&lt;br&gt;
    try:&lt;br&gt;
        slack_client.chat_postMessage(channel=SLACK_CHANNEL, text=message)&lt;br&gt;
    except SlackApiError as e:&lt;br&gt;
        logging.error(f"Slack API error: {e.response['error']}")&lt;/p&gt;

&lt;p&gt;def run_self_healing_script(script_path):&lt;br&gt;
    logging.info(f"Running self-healing script: {script_path}")&lt;br&gt;
    try:&lt;br&gt;
        subprocess.run(["bash", script_path], check=True)&lt;br&gt;
    except subprocess.CalledProcessError as e:&lt;br&gt;
        logging.error(f"Error executing self-healing script: {e}")&lt;/p&gt;

&lt;p&gt;def process_log(log):&lt;br&gt;
    try:&lt;br&gt;
        message = log['message']&lt;br&gt;
        severity = log['severity']&lt;br&gt;
        source = log['source']&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;    # Tokenize for Model Inference
    tokens = tokenizer([message], padding=True, truncation=True, max_length=128, return_tensors="tf")
    embeddings = bert_model.predict({"input_ids": tokens["input_ids"], "attention_mask": tokens["attention_mask"]})
    severity_label = np.argmax(embeddings[0][0])


    es.index(index="bug_logs", body={
        "message": message,
        "severity": severity,
        "severity_label": int(severity_label),
        "source": source
    })

    # Notification &amp;amp; Healing
    if severity_label &amp;gt;= 2:  # ERROR or CRITICAL
        notify_slack(f"Bug Detected: {message}, Severity: {severity_label}")
        run_self_healing_script("/path/to/healing_script.sh")  # Customize path

    bug_detected_counter.inc()
except Exception as e:
    logging.error(f"Error processing log: {e}")
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;

&lt;p&gt;def process_logs_stream():&lt;br&gt;
    logging.info("Starting log stream processing...")&lt;br&gt;
    for message in consumer:&lt;br&gt;
        log = message.value&lt;br&gt;
        process_log(log)&lt;/p&gt;

&lt;p&gt;def generate_reports():&lt;br&gt;
    logging.info("Generating advanced reports...")&lt;br&gt;
    logs = es.search(index="bug_logs", body={"query": {"match_all": {}}}, size=10000)&lt;br&gt;
    data = pd.DataFrame([log['_source'] for log in logs['hits']['hits']])&lt;br&gt;
    anomaly_logs = cluster_anomalies(data)&lt;br&gt;
    anomaly_logs.to_csv("anomaly_logs.csv", index=False)&lt;/p&gt;

&lt;p&gt;if &lt;strong&gt;name&lt;/strong&gt; == "&lt;strong&gt;main&lt;/strong&gt;":&lt;br&gt;
    # Start Prometheus Monitoring&lt;br&gt;
    start_http_server(8000)&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;# Process Logs from Kafka
process_logs_stream()
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;

</description>
    </item>
  </channel>
</rss>
