Skip to content

Latest commit

History

76 Commits

Folders and files

NameName
Last commit message
Last commit date

Repository files navigation

ClusterGuard

Real time fault detection system for distributed operational environments. Consumes event streams from Apache Kafka, detects anomalies using a sliding-window statistical model, triggers automated resolution workflows and exposes Prometheus metrics for observability.

Built to explore end to end ownership of a Kubernetes native service from manifest authoring and Helm packaging through to debugging pod restarts with kubectl and writing a runbook for on call response.


Architecture

┌─────────────────────────────────────────────────────────────────┐
│ Kubernetes Cluster │
│ │
│ ┌──────────┐ ┌─────────────────┐ ┌──────────────────┐ │
│ │ Kafka │───▶│ ClusterGuard │───▶│ Alert Manager │ │
│ │ (source) │ │ (consumer) │ │ (webhook sink) │ │
│ └──────────┘ │ │ └──────────────────┘ │
│ │ /metrics ──────────▶ Prometheus │
│ │ /healthz └──────▶ Grafana │
│ └─────────────────┘ │
│ │ │
│ ConfigMap (thresholds) │
└─────────────────────────────────────────────────────────────────┘

Data flow

  1. Events arrive on the ops.events Kafka topic (JSON, ~2M/day in load tests)
  2. Consumer goroutines process partitions concurrently
  3. Each event is scored against a sliding window anomaly detector
  4. Anomalies above threshold emit a Prometheus counter and POST to a configurable webhook
  5. All decisions are written to stdout as structured JSON logs

Project structure

ClusterGuard/
├── .github/
│ └── workflows/
│ ├── CI.yaml
├── cmd/
│ └── main.go # entrypoint, flag parsing, signal handling
│├── deploy/
│ ├── k8s/
│ │ ├── deployment.yaml
│ │ ├── service.yaml
│ │ │ └── helm/
│ ├── Chart.yaml
│ ├── values.yaml
├── docs/
│ └── runbook.md
├── internal/
│ ├── detector/
│ │ ├── anomaly.go # sliding-window z-score detector
│ │ └── anomaly_test.go
│ ├── metrics/
│ │ └── prometheus.go # Prometheus counters, histograms, gauges
│ └── webhook/
│ └── alert.go # HTTP POST to configurable alert endpoint
|── README.md ├── doc.go
|── go.mod
└── go.sum

Getting started

Prerequisites

  • Go 1.22+
  • Docker
  • A running Kafka broker (local: docker compose up kafka)
  • kubectl configured against a cluster (local: kind create cluster)

Run locally

git clone https://github.com/RiyaJ6/ClusterGuard
cd ClusterGuard
# start a local kafka
docker run -d --name kafka -p 9092:9092 \
-e KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://localhost:9092 \
-e KAFKA_LISTENERS=PLAINTEXT://0.0.0.0:9092 \
-e KAFKA_ZOOKEEPER_CONNECT=zookeeper:2181 \
confluentinc/cp-kafka:latest
# build and run
go build -o ClusterGuard ./cmd/ClusterGuard
./ClusterGuard \
--brokers localhost:9092 \
--topic ops.events \
--group ClusterGuard \
--webhook-url http://localhost:8080/alerts \
--metrics-port 9090

Metrics will be available at http://localhost:9090/metrics.

Run tests

go test ./... -v -race

Deploy to Kubernetes

# raw manifests
kubectl apply -f deploy/k8s/
# or via Helm
helm install ClusterGuard deploy/helm/ClusterGuard \
--set kafka.brokers=kafka:9092 \
--set webhook.url=http://alertmanager:9093/webhook

Configuration

All configuration is read from environment variables (set via ConfigMap in Kubernetes).

VariableDefaultDescription
KAFKA_BROKERSlocalhost:9092Comma separated broker list
KAFKA_TOPICops.eventsTopic to consume
KAFKA_GROUPclusterguardConsumer group ID
WINDOW_SIZE100Sliding window size for anomaly detection
ZSCORE_THRESHOLD3.0Z-score threshold for anomaly flag
WEBHOOK_URL``Alert webhook endpoint
METRICS_PORT9090Prometheus metrics port
LOG_LEVELinfoLog level (debug, info, warn, error)

Observability

Metrics exposed at /metrics

MetricTypeDescription
clusterguard_events_processed_totalCounterTotal events consumed
clusterguard_anomalies_detected_totalCounterAnomalies above threshold
clusterguard_processing_duration_secondsHistogramPer-event processing time
clusterguard_consumer_lagGaugeKafka consumer lag by partition
clusterguard_webhook_errors_totalCounterFailed alert deliveries

Structured log fields (JSON to stdout)

{
"level": "info",
"ts": "2025-03-14T09:26:00Z",
"event_id": "evt_8f3a",
"partition": 2,
"offset": 18432,
"score": 4.21,
"anomaly": true,
"processing_ms": 1.4,
"msg": "anomaly detected"
}

See docs/runbook.md for dashboard setup and on call response steps.


Debugging pod restarts — What I learned

During development I hit three real restart loops that were instructive to trace:

OOMKill — the consumer was buffering all in flight events in memory before processing. With a high throughput topic this blew the 128Mi memory limit within minutes. Fix: switched to a bounded channel with backpressure tuned limit to 256Mi after measuring actual RSS under load with kubectl top pod.

Failed readiness probe — the /healthz endpoint returned 200 only after the Kafka consumer group had rebalanced which takes 3 to 10 seconds on startup. The readiness probe was firing at 2 seconds. Fix: added an initialDelaySeconds: 15 to the probe.

Goroutine race — the anomaly detector's sliding window was a shared slice accessed by multiple partition goroutines without a lock. Found with go test -race. Fix: one detector instance per partition goroutine no shared state.

Each fix was validated by deploying to a local kind cluster and watching kubectl get pods -w until the restart count held at 0 for 5 minutes.


Terraform

The terraform/ directory provisions the AWS infrastructure used in staging: VPC, EKS node group, IAM role for the service account (IRSA), and an S3 bucket for state. See terraform/README.md for usage.


Self healing across a cluster of servers

About

Distributed fault detection on kubernetes.

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { // Add copy buttons to all
 blocks
(function() {
function addCopyButtons() {
document.querySelectorAll('pre code').forEach(function(codeBlock) {
if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;
codeBlock.parentElement.setAttribute('data-copy-added', 'true');
var btn = document.createElement('button');
btn.textContent = 'Copy';
btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';
btn.onmouseover = function() { this.style.opacity = '1'; };
btn.onmouseout = function() { this.style.opacity = '0.7'; };
btn.onclick = function() {
navigator.clipboard.writeText(codeBlock.textContent).then(function() {
btn.textContent = 'Copied!';
setTimeout(function() { btn.textContent = 'Copy'; }, 1500);
});
};
codeBlock.parentElement.style.position = 'relative';
codeBlock.parentElement.appendChild(btn);
});
}
addCopyButtons();
// Re-run on dynamic content
var observer = new MutationObserver(addCopyButtons);
observer.observe(document.body, { childList: true, subtree: true });
})();
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
GitHub - RiyaJ6/ClusterGuard: Distributed fault detection on kubernetes. · GitHub
Skip to content

Latest commit

History

76 Commits

Folders and files

NameName
Last commit message
Last commit date

Repository files navigation

ClusterGuard

Real time fault detection system for distributed operational environments. Consumes event streams from Apache Kafka, detects anomalies using a sliding-window statistical model, triggers automated resolution workflows and exposes Prometheus metrics for observability.

Built to explore end to end ownership of a Kubernetes native service from manifest authoring and Helm packaging through to debugging pod restarts with kubectl and writing a runbook for on call response.


Architecture

┌─────────────────────────────────────────────────────────────────┐
│ Kubernetes Cluster │
│ │
│ ┌──────────┐ ┌─────────────────┐ ┌──────────────────┐ │
│ │ Kafka │───▶│ ClusterGuard │───▶│ Alert Manager │ │
│ │ (source) │ │ (consumer) │ │ (webhook sink) │ │
│ └──────────┘ │ │ └──────────────────┘ │
│ │ /metrics ──────────▶ Prometheus │
│ │ /healthz └──────▶ Grafana │
│ └─────────────────┘ │
│ │ │
│ ConfigMap (thresholds) │
└─────────────────────────────────────────────────────────────────┘

Data flow

  1. Events arrive on the ops.events Kafka topic (JSON, ~2M/day in load tests)
  2. Consumer goroutines process partitions concurrently
  3. Each event is scored against a sliding window anomaly detector
  4. Anomalies above threshold emit a Prometheus counter and POST to a configurable webhook
  5. All decisions are written to stdout as structured JSON logs

Project structure

ClusterGuard/
├── .github/
│ └── workflows/
│ ├── CI.yaml
├── cmd/
│ └── main.go # entrypoint, flag parsing, signal handling
│├── deploy/
│ ├── k8s/
│ │ ├── deployment.yaml
│ │ ├── service.yaml
│ │ │ └── helm/
│ ├── Chart.yaml
│ ├── values.yaml
├── docs/
│ └── runbook.md
├── internal/
│ ├── detector/
│ │ ├── anomaly.go # sliding-window z-score detector
│ │ └── anomaly_test.go
│ ├── metrics/
│ │ └── prometheus.go # Prometheus counters, histograms, gauges
│ └── webhook/
│ └── alert.go # HTTP POST to configurable alert endpoint
|── README.md ├── doc.go
|── go.mod
└── go.sum

Getting started

Prerequisites

  • Go 1.22+
  • Docker
  • A running Kafka broker (local: docker compose up kafka)
  • kubectl configured against a cluster (local: kind create cluster)

Run locally

git clone https://github.com/RiyaJ6/ClusterGuard
cd ClusterGuard
# start a local kafka
docker run -d --name kafka -p 9092:9092 \
-e KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://localhost:9092 \
-e KAFKA_LISTENERS=PLAINTEXT://0.0.0.0:9092 \
-e KAFKA_ZOOKEEPER_CONNECT=zookeeper:2181 \
confluentinc/cp-kafka:latest
# build and run
go build -o ClusterGuard ./cmd/ClusterGuard
./ClusterGuard \
--brokers localhost:9092 \
--topic ops.events \
--group ClusterGuard \
--webhook-url http://localhost:8080/alerts \
--metrics-port 9090

Metrics will be available at http://localhost:9090/metrics.

Run tests

go test ./... -v -race

Deploy to Kubernetes

# raw manifests
kubectl apply -f deploy/k8s/
# or via Helm
helm install ClusterGuard deploy/helm/ClusterGuard \
--set kafka.brokers=kafka:9092 \
--set webhook.url=http://alertmanager:9093/webhook

Configuration

All configuration is read from environment variables (set via ConfigMap in Kubernetes).

VariableDefaultDescription
KAFKA_BROKERSlocalhost:9092Comma separated broker list
KAFKA_TOPICops.eventsTopic to consume
KAFKA_GROUPclusterguardConsumer group ID
WINDOW_SIZE100Sliding window size for anomaly detection
ZSCORE_THRESHOLD3.0Z-score threshold for anomaly flag
WEBHOOK_URL``Alert webhook endpoint
METRICS_PORT9090Prometheus metrics port
LOG_LEVELinfoLog level (debug, info, warn, error)

Observability

Metrics exposed at /metrics

MetricTypeDescription
clusterguard_events_processed_totalCounterTotal events consumed
clusterguard_anomalies_detected_totalCounterAnomalies above threshold
clusterguard_processing_duration_secondsHistogramPer-event processing time
clusterguard_consumer_lagGaugeKafka consumer lag by partition
clusterguard_webhook_errors_totalCounterFailed alert deliveries

Structured log fields (JSON to stdout)

{
"level": "info",
"ts": "2025-03-14T09:26:00Z",
"event_id": "evt_8f3a",
"partition": 2,
"offset": 18432,
"score": 4.21,
"anomaly": true,
"processing_ms": 1.4,
"msg": "anomaly detected"
}

See docs/runbook.md for dashboard setup and on call response steps.


Debugging pod restarts — What I learned

During development I hit three real restart loops that were instructive to trace:

OOMKill — the consumer was buffering all in flight events in memory before processing. With a high throughput topic this blew the 128Mi memory limit within minutes. Fix: switched to a bounded channel with backpressure tuned limit to 256Mi after measuring actual RSS under load with kubectl top pod.

Failed readiness probe — the /healthz endpoint returned 200 only after the Kafka consumer group had rebalanced which takes 3 to 10 seconds on startup. The readiness probe was firing at 2 seconds. Fix: added an initialDelaySeconds: 15 to the probe.

Goroutine race — the anomaly detector's sliding window was a shared slice accessed by multiple partition goroutines without a lock. Found with go test -race. Fix: one detector instance per partition goroutine no shared state.

Each fix was validated by deploying to a local kind cluster and watching kubectl get pods -w until the restart count held at 0 for 5 minutes.


Terraform

The terraform/ directory provisions the AWS infrastructure used in staging: VPC, EKS node group, IAM role for the service account (IRSA), and an S3 bucket for state. See terraform/README.md for usage.


Self healing across a cluster of servers

About

Distributed fault detection on kubernetes.

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { // Force GitHub README to respect dark mode (function() { var style = document.createElement('style'); style.textContent = ' .markdown-body { color-scheme: dark light; } .markdown-body pre { background: #161b22 !important; } .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; } .markdown-body table th, .markdown-body table td { border-color: #30363d !important; } .markdown-body img { background: #0d1117; } .markdown-body blockquote { border-left-color: #8b949e; } .markdown-body hr { border-color: #30363d; } '; document.head.appendChild(style); })(); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' GitHub - RiyaJ6/ClusterGuard: Distributed fault detection on kubernetes. · GitHub
Skip to content

Latest commit

History

76 Commits

Folders and files

NameName
Last commit message
Last commit date

Repository files navigation

ClusterGuard

Real time fault detection system for distributed operational environments. Consumes event streams from Apache Kafka, detects anomalies using a sliding-window statistical model, triggers automated resolution workflows and exposes Prometheus metrics for observability.

Built to explore end to end ownership of a Kubernetes native service from manifest authoring and Helm packaging through to debugging pod restarts with kubectl and writing a runbook for on call response.


Architecture

┌─────────────────────────────────────────────────────────────────┐
│ Kubernetes Cluster │
│ │
│ ┌──────────┐ ┌─────────────────┐ ┌──────────────────┐ │
│ │ Kafka │───▶│ ClusterGuard │───▶│ Alert Manager │ │
│ │ (source) │ │ (consumer) │ │ (webhook sink) │ │
│ └──────────┘ │ │ └──────────────────┘ │
│ │ /metrics ──────────▶ Prometheus │
│ │ /healthz └──────▶ Grafana │
│ └─────────────────┘ │
│ │ │
│ ConfigMap (thresholds) │
└─────────────────────────────────────────────────────────────────┘

Data flow

  1. Events arrive on the ops.events Kafka topic (JSON, ~2M/day in load tests)
  2. Consumer goroutines process partitions concurrently
  3. Each event is scored against a sliding window anomaly detector
  4. Anomalies above threshold emit a Prometheus counter and POST to a configurable webhook
  5. All decisions are written to stdout as structured JSON logs

Project structure

ClusterGuard/
├── .github/
│ └── workflows/
│ ├── CI.yaml
├── cmd/
│ └── main.go # entrypoint, flag parsing, signal handling
│├── deploy/
│ ├── k8s/
│ │ ├── deployment.yaml
│ │ ├── service.yaml
│ │ │ └── helm/
│ ├── Chart.yaml
│ ├── values.yaml
├── docs/
│ └── runbook.md
├── internal/
│ ├── detector/
│ │ ├── anomaly.go # sliding-window z-score detector
│ │ └── anomaly_test.go
│ ├── metrics/
│ │ └── prometheus.go # Prometheus counters, histograms, gauges
│ └── webhook/
│ └── alert.go # HTTP POST to configurable alert endpoint
|── README.md ├── doc.go
|── go.mod
└── go.sum

Getting started

Prerequisites

  • Go 1.22+
  • Docker
  • A running Kafka broker (local: docker compose up kafka)
  • kubectl configured against a cluster (local: kind create cluster)

Run locally

git clone https://github.com/RiyaJ6/ClusterGuard
cd ClusterGuard
# start a local kafka
docker run -d --name kafka -p 9092:9092 \
-e KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://localhost:9092 \
-e KAFKA_LISTENERS=PLAINTEXT://0.0.0.0:9092 \
-e KAFKA_ZOOKEEPER_CONNECT=zookeeper:2181 \
confluentinc/cp-kafka:latest
# build and run
go build -o ClusterGuard ./cmd/ClusterGuard
./ClusterGuard \
--brokers localhost:9092 \
--topic ops.events \
--group ClusterGuard \
--webhook-url http://localhost:8080/alerts \
--metrics-port 9090

Metrics will be available at http://localhost:9090/metrics.

Run tests

go test ./... -v -race

Deploy to Kubernetes

# raw manifests
kubectl apply -f deploy/k8s/
# or via Helm
helm install ClusterGuard deploy/helm/ClusterGuard \
--set kafka.brokers=kafka:9092 \
--set webhook.url=http://alertmanager:9093/webhook

Configuration

All configuration is read from environment variables (set via ConfigMap in Kubernetes).

VariableDefaultDescription
KAFKA_BROKERSlocalhost:9092Comma separated broker list
KAFKA_TOPICops.eventsTopic to consume
KAFKA_GROUPclusterguardConsumer group ID
WINDOW_SIZE100Sliding window size for anomaly detection
ZSCORE_THRESHOLD3.0Z-score threshold for anomaly flag
WEBHOOK_URL``Alert webhook endpoint
METRICS_PORT9090Prometheus metrics port
LOG_LEVELinfoLog level (debug, info, warn, error)

Observability

Metrics exposed at /metrics

MetricTypeDescription
clusterguard_events_processed_totalCounterTotal events consumed
clusterguard_anomalies_detected_totalCounterAnomalies above threshold
clusterguard_processing_duration_secondsHistogramPer-event processing time
clusterguard_consumer_lagGaugeKafka consumer lag by partition
clusterguard_webhook_errors_totalCounterFailed alert deliveries

Structured log fields (JSON to stdout)

{
"level": "info",
"ts": "2025-03-14T09:26:00Z",
"event_id": "evt_8f3a",
"partition": 2,
"offset": 18432,
"score": 4.21,
"anomaly": true,
"processing_ms": 1.4,
"msg": "anomaly detected"
}

See docs/runbook.md for dashboard setup and on call response steps.


Debugging pod restarts — What I learned

During development I hit three real restart loops that were instructive to trace:

OOMKill — the consumer was buffering all in flight events in memory before processing. With a high throughput topic this blew the 128Mi memory limit within minutes. Fix: switched to a bounded channel with backpressure tuned limit to 256Mi after measuring actual RSS under load with kubectl top pod.

Failed readiness probe — the /healthz endpoint returned 200 only after the Kafka consumer group had rebalanced which takes 3 to 10 seconds on startup. The readiness probe was firing at 2 seconds. Fix: added an initialDelaySeconds: 15 to the probe.

Goroutine race — the anomaly detector's sliding window was a shared slice accessed by multiple partition goroutines without a lock. Found with go test -race. Fix: one detector instance per partition goroutine no shared state.

Each fix was validated by deploying to a local kind cluster and watching kubectl get pods -w until the restart count held at 0 for 5 minutes.


Terraform

The terraform/ directory provisions the AWS infrastructure used in staging: VPC, EKS node group, IAM role for the service account (IRSA), and an S3 bucket for state. See terraform/README.md for usage.


Self healing across a cluster of servers

About

Distributed fault detection on kubernetes.

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { // Highlight search terms from Google/DuckDuckGo/Bing referrer (function() { var ref = document.referrer; var terms = []; if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) { var url = new URL(ref); var q = url.searchParams.get('q') || url.searchParams.get('p'); if (q) { terms = q.split(/\s+/).filter(function(t) { return t.length > 2; }); } } if (terms.length === 0) return; var style = document.createElement('style'); style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }'; document.head.appendChild(style); function highlight(node) { if (node.nodeType === 3) { // text node var text = node.textContent; var found = false; terms.forEach(function(term) { var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\]\\]/g, '\\') + ')', 'gi'); if (regex.test(text)) { found = true; var frag = document.createDocumentFragment(); var parts = text.split(regex); parts.forEach(function(part, i) { if (i % 2 === 0) { frag.appendChild(document.createTextNode(part)); } else { var span = document.createElement('span'); span.className = 'userscript-highlight'; span.textContent = part; frag.appendChild(span); } }); node.parentNode.replaceChild(frag, node); } }); } else if (node.nodeType === 1 && node.childNodes) { // element var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT']; if (!skipTags.includes(node.tagName)) { Array.from(node.childNodes).forEach(highlight); } } } highlight(document.body); // Re-highlight on dynamic content var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1 || node.nodeType === 3) highlight(node); }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' GitHub - RiyaJ6/ClusterGuard: Distributed fault detection on kubernetes. · GitHub
Skip to content

Latest commit

History

76 Commits

Folders and files

NameName
Last commit message
Last commit date

Repository files navigation

ClusterGuard

Real time fault detection system for distributed operational environments. Consumes event streams from Apache Kafka, detects anomalies using a sliding-window statistical model, triggers automated resolution workflows and exposes Prometheus metrics for observability.

Built to explore end to end ownership of a Kubernetes native service from manifest authoring and Helm packaging through to debugging pod restarts with kubectl and writing a runbook for on call response.


Architecture

┌─────────────────────────────────────────────────────────────────┐
│ Kubernetes Cluster │
│ │
│ ┌──────────┐ ┌─────────────────┐ ┌──────────────────┐ │
│ │ Kafka │───▶│ ClusterGuard │───▶│ Alert Manager │ │
│ │ (source) │ │ (consumer) │ │ (webhook sink) │ │
│ └──────────┘ │ │ └──────────────────┘ │
│ │ /metrics ──────────▶ Prometheus │
│ │ /healthz └──────▶ Grafana │
│ └─────────────────┘ │
│ │ │
│ ConfigMap (thresholds) │
└─────────────────────────────────────────────────────────────────┘

Data flow

  1. Events arrive on the ops.events Kafka topic (JSON, ~2M/day in load tests)
  2. Consumer goroutines process partitions concurrently
  3. Each event is scored against a sliding window anomaly detector
  4. Anomalies above threshold emit a Prometheus counter and POST to a configurable webhook
  5. All decisions are written to stdout as structured JSON logs

Project structure

ClusterGuard/
├── .github/
│ └── workflows/
│ ├── CI.yaml
├── cmd/
│ └── main.go # entrypoint, flag parsing, signal handling
│├── deploy/
│ ├── k8s/
│ │ ├── deployment.yaml
│ │ ├── service.yaml
│ │ │ └── helm/
│ ├── Chart.yaml
│ ├── values.yaml
├── docs/
│ └── runbook.md
├── internal/
│ ├── detector/
│ │ ├── anomaly.go # sliding-window z-score detector
│ │ └── anomaly_test.go
│ ├── metrics/
│ │ └── prometheus.go # Prometheus counters, histograms, gauges
│ └── webhook/
│ └── alert.go # HTTP POST to configurable alert endpoint
|── README.md ├── doc.go
|── go.mod
└── go.sum

Getting started

Prerequisites

  • Go 1.22+
  • Docker
  • A running Kafka broker (local: docker compose up kafka)
  • kubectl configured against a cluster (local: kind create cluster)

Run locally

git clone https://github.com/RiyaJ6/ClusterGuard
cd ClusterGuard
# start a local kafka
docker run -d --name kafka -p 9092:9092 \
-e KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://localhost:9092 \
-e KAFKA_LISTENERS=PLAINTEXT://0.0.0.0:9092 \
-e KAFKA_ZOOKEEPER_CONNECT=zookeeper:2181 \
confluentinc/cp-kafka:latest
# build and run
go build -o ClusterGuard ./cmd/ClusterGuard
./ClusterGuard \
--brokers localhost:9092 \
--topic ops.events \
--group ClusterGuard \
--webhook-url http://localhost:8080/alerts \
--metrics-port 9090

Metrics will be available at http://localhost:9090/metrics.

Run tests

go test ./... -v -race

Deploy to Kubernetes

# raw manifests
kubectl apply -f deploy/k8s/
# or via Helm
helm install ClusterGuard deploy/helm/ClusterGuard \
--set kafka.brokers=kafka:9092 \
--set webhook.url=http://alertmanager:9093/webhook

Configuration

All configuration is read from environment variables (set via ConfigMap in Kubernetes).

VariableDefaultDescription
KAFKA_BROKERSlocalhost:9092Comma separated broker list
KAFKA_TOPICops.eventsTopic to consume
KAFKA_GROUPclusterguardConsumer group ID
WINDOW_SIZE100Sliding window size for anomaly detection
ZSCORE_THRESHOLD3.0Z-score threshold for anomaly flag
WEBHOOK_URL``Alert webhook endpoint
METRICS_PORT9090Prometheus metrics port
LOG_LEVELinfoLog level (debug, info, warn, error)

Observability

Metrics exposed at /metrics

MetricTypeDescription
clusterguard_events_processed_totalCounterTotal events consumed
clusterguard_anomalies_detected_totalCounterAnomalies above threshold
clusterguard_processing_duration_secondsHistogramPer-event processing time
clusterguard_consumer_lagGaugeKafka consumer lag by partition
clusterguard_webhook_errors_totalCounterFailed alert deliveries

Structured log fields (JSON to stdout)

{
"level": "info",
"ts": "2025-03-14T09:26:00Z",
"event_id": "evt_8f3a",
"partition": 2,
"offset": 18432,
"score": 4.21,
"anomaly": true,
"processing_ms": 1.4,
"msg": "anomaly detected"
}

See docs/runbook.md for dashboard setup and on call response steps.


Debugging pod restarts — What I learned

During development I hit three real restart loops that were instructive to trace:

OOMKill — the consumer was buffering all in flight events in memory before processing. With a high throughput topic this blew the 128Mi memory limit within minutes. Fix: switched to a bounded channel with backpressure tuned limit to 256Mi after measuring actual RSS under load with kubectl top pod.

Failed readiness probe — the /healthz endpoint returned 200 only after the Kafka consumer group had rebalanced which takes 3 to 10 seconds on startup. The readiness probe was firing at 2 seconds. Fix: added an initialDelaySeconds: 15 to the probe.

Goroutine race — the anomaly detector's sliding window was a shared slice accessed by multiple partition goroutines without a lock. Found with go test -race. Fix: one detector instance per partition goroutine no shared state.

Each fix was validated by deploying to a local kind cluster and watching kubectl get pods -w until the restart count held at 0 for 5 minutes.


Terraform

The terraform/ directory provisions the AWS infrastructure used in staging: VPC, EKS node group, IAM role for the service account (IRSA), and an S3 bucket for state. See terraform/README.md for usage.


Self healing across a cluster of servers

About

Distributed fault detection on kubernetes.

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { // Strip utm_, fbclid, gclid, etc. from all links on page (function() { var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content', 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid', 'ref', 'ref_src', 'source', 'medium', 'campaign']; function cleanUrl(url) { try { var u = new URL(url, window.location.origin); var changed = false; trackingParams.forEach(function(p) { if (u.searchParams.has(p)) { u.searchParams.delete(p); changed = true; } }); return changed ? u.toString() : url; } catch (e) { return url; } } function cleanLinks() { document.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } cleanLinks(); var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1) { if (node.tagName === 'A') cleanLinks(); node.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + ' GitHub - RiyaJ6/ClusterGuard: Distributed fault detection on kubernetes. · GitHub
Skip to content

Latest commit

History

76 Commits

Folders and files

NameName
Last commit message
Last commit date

Repository files navigation

ClusterGuard

Real time fault detection system for distributed operational environments. Consumes event streams from Apache Kafka, detects anomalies using a sliding-window statistical model, triggers automated resolution workflows and exposes Prometheus metrics for observability.

Built to explore end to end ownership of a Kubernetes native service from manifest authoring and Helm packaging through to debugging pod restarts with kubectl and writing a runbook for on call response.


Architecture

┌─────────────────────────────────────────────────────────────────┐
│ Kubernetes Cluster │
│ │
│ ┌──────────┐ ┌─────────────────┐ ┌──────────────────┐ │
│ │ Kafka │───▶│ ClusterGuard │───▶│ Alert Manager │ │
│ │ (source) │ │ (consumer) │ │ (webhook sink) │ │
│ └──────────┘ │ │ └──────────────────┘ │
│ │ /metrics ──────────▶ Prometheus │
│ │ /healthz └──────▶ Grafana │
│ └─────────────────┘ │
│ │ │
│ ConfigMap (thresholds) │
└─────────────────────────────────────────────────────────────────┘

Data flow

  1. Events arrive on the ops.events Kafka topic (JSON, ~2M/day in load tests)
  2. Consumer goroutines process partitions concurrently
  3. Each event is scored against a sliding window anomaly detector
  4. Anomalies above threshold emit a Prometheus counter and POST to a configurable webhook
  5. All decisions are written to stdout as structured JSON logs

Project structure

ClusterGuard/
├── .github/
│ └── workflows/
│ ├── CI.yaml
├── cmd/
│ └── main.go # entrypoint, flag parsing, signal handling
│├── deploy/
│ ├── k8s/
│ │ ├── deployment.yaml
│ │ ├── service.yaml
│ │ │ └── helm/
│ ├── Chart.yaml
│ ├── values.yaml
├── docs/
│ └── runbook.md
├── internal/
│ ├── detector/
│ │ ├── anomaly.go # sliding-window z-score detector
│ │ └── anomaly_test.go
│ ├── metrics/
│ │ └── prometheus.go # Prometheus counters, histograms, gauges
│ └── webhook/
│ └── alert.go # HTTP POST to configurable alert endpoint
|── README.md ├── doc.go
|── go.mod
└── go.sum

Getting started

Prerequisites

  • Go 1.22+
  • Docker
  • A running Kafka broker (local: docker compose up kafka)
  • kubectl configured against a cluster (local: kind create cluster)

Run locally

git clone https://github.com/RiyaJ6/ClusterGuard
cd ClusterGuard
# start a local kafka
docker run -d --name kafka -p 9092:9092 \
-e KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://localhost:9092 \
-e KAFKA_LISTENERS=PLAINTEXT://0.0.0.0:9092 \
-e KAFKA_ZOOKEEPER_CONNECT=zookeeper:2181 \
confluentinc/cp-kafka:latest
# build and run
go build -o ClusterGuard ./cmd/ClusterGuard
./ClusterGuard \
--brokers localhost:9092 \
--topic ops.events \
--group ClusterGuard \
--webhook-url http://localhost:8080/alerts \
--metrics-port 9090

Metrics will be available at http://localhost:9090/metrics.

Run tests

go test ./... -v -race

Deploy to Kubernetes

# raw manifests
kubectl apply -f deploy/k8s/
# or via Helm
helm install ClusterGuard deploy/helm/ClusterGuard \
--set kafka.brokers=kafka:9092 \
--set webhook.url=http://alertmanager:9093/webhook

Configuration

All configuration is read from environment variables (set via ConfigMap in Kubernetes).

VariableDefaultDescription
KAFKA_BROKERSlocalhost:9092Comma separated broker list
KAFKA_TOPICops.eventsTopic to consume
KAFKA_GROUPclusterguardConsumer group ID
WINDOW_SIZE100Sliding window size for anomaly detection
ZSCORE_THRESHOLD3.0Z-score threshold for anomaly flag
WEBHOOK_URL``Alert webhook endpoint
METRICS_PORT9090Prometheus metrics port
LOG_LEVELinfoLog level (debug, info, warn, error)

Observability

Metrics exposed at /metrics

MetricTypeDescription
clusterguard_events_processed_totalCounterTotal events consumed
clusterguard_anomalies_detected_totalCounterAnomalies above threshold
clusterguard_processing_duration_secondsHistogramPer-event processing time
clusterguard_consumer_lagGaugeKafka consumer lag by partition
clusterguard_webhook_errors_totalCounterFailed alert deliveries

Structured log fields (JSON to stdout)

{
"level": "info",
"ts": "2025-03-14T09:26:00Z",
"event_id": "evt_8f3a",
"partition": 2,
"offset": 18432,
"score": 4.21,
"anomaly": true,
"processing_ms": 1.4,
"msg": "anomaly detected"
}

See docs/runbook.md for dashboard setup and on call response steps.


Debugging pod restarts — What I learned

During development I hit three real restart loops that were instructive to trace:

OOMKill — the consumer was buffering all in flight events in memory before processing. With a high throughput topic this blew the 128Mi memory limit within minutes. Fix: switched to a bounded channel with backpressure tuned limit to 256Mi after measuring actual RSS under load with kubectl top pod.

Failed readiness probe — the /healthz endpoint returned 200 only after the Kafka consumer group had rebalanced which takes 3 to 10 seconds on startup. The readiness probe was firing at 2 seconds. Fix: added an initialDelaySeconds: 15 to the probe.

Goroutine race — the anomaly detector's sliding window was a shared slice accessed by multiple partition goroutines without a lock. Found with go test -race. Fix: one detector instance per partition goroutine no shared state.

Each fix was validated by deploying to a local kind cluster and watching kubectl get pods -w until the restart count held at 0 for 5 minutes.


Terraform

The terraform/ directory provisions the AWS infrastructure used in staging: VPC, EKS node group, IAM role for the service account (IRSA), and an S3 bucket for state. See terraform/README.md for usage.


Self healing across a cluster of servers

About

Distributed fault detection on kubernetes.

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { // Auto-enable theater mode on YouTube (function() { function tryTheater() { var btn = document.querySelector('button[aria-label="Theater mode"], ytd-player #player button[title="Theater mode"]'); if (btn && !btn.classList.contains('activated')) { btn.click(); } } // Try immediately tryTheater(); // Try after navigation (SPA) var lastUrl = location.href; setInterval(function() { if (location.href !== lastUrl) { lastUrl = location.href; setTimeout(tryTheater, 500); } }, 1000); // Also try on player load var observer = new MutationObserver(tryTheater); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' GitHub - RiyaJ6/ClusterGuard: Distributed fault detection on kubernetes. · GitHub
Skip to content

Latest commit

History

76 Commits

Folders and files

NameName
Last commit message
Last commit date

Repository files navigation

ClusterGuard

Real time fault detection system for distributed operational environments. Consumes event streams from Apache Kafka, detects anomalies using a sliding-window statistical model, triggers automated resolution workflows and exposes Prometheus metrics for observability.

Built to explore end to end ownership of a Kubernetes native service from manifest authoring and Helm packaging through to debugging pod restarts with kubectl and writing a runbook for on call response.


Architecture

┌─────────────────────────────────────────────────────────────────┐
│ Kubernetes Cluster │
│ │
│ ┌──────────┐ ┌─────────────────┐ ┌──────────────────┐ │
│ │ Kafka │───▶│ ClusterGuard │───▶│ Alert Manager │ │
│ │ (source) │ │ (consumer) │ │ (webhook sink) │ │
│ └──────────┘ │ │ └──────────────────┘ │
│ │ /metrics ──────────▶ Prometheus │
│ │ /healthz └──────▶ Grafana │
│ └─────────────────┘ │
│ │ │
│ ConfigMap (thresholds) │
└─────────────────────────────────────────────────────────────────┘

Data flow

  1. Events arrive on the ops.events Kafka topic (JSON, ~2M/day in load tests)
  2. Consumer goroutines process partitions concurrently
  3. Each event is scored against a sliding window anomaly detector
  4. Anomalies above threshold emit a Prometheus counter and POST to a configurable webhook
  5. All decisions are written to stdout as structured JSON logs

Project structure

ClusterGuard/
├── .github/
│ └── workflows/
│ ├── CI.yaml
├── cmd/
│ └── main.go # entrypoint, flag parsing, signal handling
│├── deploy/
│ ├── k8s/
│ │ ├── deployment.yaml
│ │ ├── service.yaml
│ │ │ └── helm/
│ ├── Chart.yaml
│ ├── values.yaml
├── docs/
│ └── runbook.md
├── internal/
│ ├── detector/
│ │ ├── anomaly.go # sliding-window z-score detector
│ │ └── anomaly_test.go
│ ├── metrics/
│ │ └── prometheus.go # Prometheus counters, histograms, gauges
│ └── webhook/
│ └── alert.go # HTTP POST to configurable alert endpoint
|── README.md ├── doc.go
|── go.mod
└── go.sum

Getting started

Prerequisites

  • Go 1.22+
  • Docker
  • A running Kafka broker (local: docker compose up kafka)
  • kubectl configured against a cluster (local: kind create cluster)

Run locally

git clone https://github.com/RiyaJ6/ClusterGuard
cd ClusterGuard
# start a local kafka
docker run -d --name kafka -p 9092:9092 \
-e KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://localhost:9092 \
-e KAFKA_LISTENERS=PLAINTEXT://0.0.0.0:9092 \
-e KAFKA_ZOOKEEPER_CONNECT=zookeeper:2181 \
confluentinc/cp-kafka:latest
# build and run
go build -o ClusterGuard ./cmd/ClusterGuard
./ClusterGuard \
--brokers localhost:9092 \
--topic ops.events \
--group ClusterGuard \
--webhook-url http://localhost:8080/alerts \
--metrics-port 9090

Metrics will be available at http://localhost:9090/metrics.

Run tests

go test ./... -v -race

Deploy to Kubernetes

# raw manifests
kubectl apply -f deploy/k8s/
# or via Helm
helm install ClusterGuard deploy/helm/ClusterGuard \
--set kafka.brokers=kafka:9092 \
--set webhook.url=http://alertmanager:9093/webhook

Configuration

All configuration is read from environment variables (set via ConfigMap in Kubernetes).

VariableDefaultDescription
KAFKA_BROKERSlocalhost:9092Comma separated broker list
KAFKA_TOPICops.eventsTopic to consume
KAFKA_GROUPclusterguardConsumer group ID
WINDOW_SIZE100Sliding window size for anomaly detection
ZSCORE_THRESHOLD3.0Z-score threshold for anomaly flag
WEBHOOK_URL``Alert webhook endpoint
METRICS_PORT9090Prometheus metrics port
LOG_LEVELinfoLog level (debug, info, warn, error)

Observability

Metrics exposed at /metrics

MetricTypeDescription
clusterguard_events_processed_totalCounterTotal events consumed
clusterguard_anomalies_detected_totalCounterAnomalies above threshold
clusterguard_processing_duration_secondsHistogramPer-event processing time
clusterguard_consumer_lagGaugeKafka consumer lag by partition
clusterguard_webhook_errors_totalCounterFailed alert deliveries

Structured log fields (JSON to stdout)

{
"level": "info",
"ts": "2025-03-14T09:26:00Z",
"event_id": "evt_8f3a",
"partition": 2,
"offset": 18432,
"score": 4.21,
"anomaly": true,
"processing_ms": 1.4,
"msg": "anomaly detected"
}

See docs/runbook.md for dashboard setup and on call response steps.


Debugging pod restarts — What I learned

During development I hit three real restart loops that were instructive to trace:

OOMKill — the consumer was buffering all in flight events in memory before processing. With a high throughput topic this blew the 128Mi memory limit within minutes. Fix: switched to a bounded channel with backpressure tuned limit to 256Mi after measuring actual RSS under load with kubectl top pod.

Failed readiness probe — the /healthz endpoint returned 200 only after the Kafka consumer group had rebalanced which takes 3 to 10 seconds on startup. The readiness probe was firing at 2 seconds. Fix: added an initialDelaySeconds: 15 to the probe.

Goroutine race — the anomaly detector's sliding window was a shared slice accessed by multiple partition goroutines without a lock. Found with go test -race. Fix: one detector instance per partition goroutine no shared state.

Each fix was validated by deploying to a local kind cluster and watching kubectl get pods -w until the restart count held at 0 for 5 minutes.


Terraform

The terraform/ directory provisions the AWS infrastructure used in staging: VPC, EKS node group, IAM role for the service account (IRSA), and an S3 bucket for state. See terraform/README.md for usage.


Self healing across a cluster of servers

About

Distributed fault detection on kubernetes.

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { // Remove or un-stick sticky/fixed headers that block content (function() { function unstick() { document.querySelectorAll('header, nav, [role="banner"], .header, .navbar, .sticky, .fixed-top, [style*="position: fixed"], [style*="position:sticky"]').forEach(function(el) { if (el.style.position === 'fixed' || el.style.position === 'sticky' || getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') { el.style.position = 'static'; el.style.top = 'auto'; el.style.zIndex = 'auto'; } }); } unstick(); var observer = new MutationObserver(unstick); observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] }); })(); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' GitHub - RiyaJ6/ClusterGuard: Distributed fault detection on kubernetes. · GitHub
Skip to content

Latest commit

History

76 Commits

Folders and files

NameName
Last commit message
Last commit date

Repository files navigation

ClusterGuard

Real time fault detection system for distributed operational environments. Consumes event streams from Apache Kafka, detects anomalies using a sliding-window statistical model, triggers automated resolution workflows and exposes Prometheus metrics for observability.

Built to explore end to end ownership of a Kubernetes native service from manifest authoring and Helm packaging through to debugging pod restarts with kubectl and writing a runbook for on call response.


Architecture

┌─────────────────────────────────────────────────────────────────┐
│ Kubernetes Cluster │
│ │
│ ┌──────────┐ ┌─────────────────┐ ┌──────────────────┐ │
│ │ Kafka │───▶│ ClusterGuard │───▶│ Alert Manager │ │
│ │ (source) │ │ (consumer) │ │ (webhook sink) │ │
│ └──────────┘ │ │ └──────────────────┘ │
│ │ /metrics ──────────▶ Prometheus │
│ │ /healthz └──────▶ Grafana │
│ └─────────────────┘ │
│ │ │
│ ConfigMap (thresholds) │
└─────────────────────────────────────────────────────────────────┘

Data flow

  1. Events arrive on the ops.events Kafka topic (JSON, ~2M/day in load tests)
  2. Consumer goroutines process partitions concurrently
  3. Each event is scored against a sliding window anomaly detector
  4. Anomalies above threshold emit a Prometheus counter and POST to a configurable webhook
  5. All decisions are written to stdout as structured JSON logs

Project structure

ClusterGuard/
├── .github/
│ └── workflows/
│ ├── CI.yaml
├── cmd/
│ └── main.go # entrypoint, flag parsing, signal handling
│├── deploy/
│ ├── k8s/
│ │ ├── deployment.yaml
│ │ ├── service.yaml
│ │ │ └── helm/
│ ├── Chart.yaml
│ ├── values.yaml
├── docs/
│ └── runbook.md
├── internal/
│ ├── detector/
│ │ ├── anomaly.go # sliding-window z-score detector
│ │ └── anomaly_test.go
│ ├── metrics/
│ │ └── prometheus.go # Prometheus counters, histograms, gauges
│ └── webhook/
│ └── alert.go # HTTP POST to configurable alert endpoint
|── README.md ├── doc.go
|── go.mod
└── go.sum

Getting started

Prerequisites

  • Go 1.22+
  • Docker
  • A running Kafka broker (local: docker compose up kafka)
  • kubectl configured against a cluster (local: kind create cluster)

Run locally

git clone https://github.com/RiyaJ6/ClusterGuard
cd ClusterGuard
# start a local kafka
docker run -d --name kafka -p 9092:9092 \
-e KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://localhost:9092 \
-e KAFKA_LISTENERS=PLAINTEXT://0.0.0.0:9092 \
-e KAFKA_ZOOKEEPER_CONNECT=zookeeper:2181 \
confluentinc/cp-kafka:latest
# build and run
go build -o ClusterGuard ./cmd/ClusterGuard
./ClusterGuard \
--brokers localhost:9092 \
--topic ops.events \
--group ClusterGuard \
--webhook-url http://localhost:8080/alerts \
--metrics-port 9090

Metrics will be available at http://localhost:9090/metrics.

Run tests

go test ./... -v -race

Deploy to Kubernetes

# raw manifests
kubectl apply -f deploy/k8s/
# or via Helm
helm install ClusterGuard deploy/helm/ClusterGuard \
--set kafka.brokers=kafka:9092 \
--set webhook.url=http://alertmanager:9093/webhook

Configuration

All configuration is read from environment variables (set via ConfigMap in Kubernetes).

VariableDefaultDescription
KAFKA_BROKERSlocalhost:9092Comma separated broker list
KAFKA_TOPICops.eventsTopic to consume
KAFKA_GROUPclusterguardConsumer group ID
WINDOW_SIZE100Sliding window size for anomaly detection
ZSCORE_THRESHOLD3.0Z-score threshold for anomaly flag
WEBHOOK_URL``Alert webhook endpoint
METRICS_PORT9090Prometheus metrics port
LOG_LEVELinfoLog level (debug, info, warn, error)

Observability

Metrics exposed at /metrics

MetricTypeDescription
clusterguard_events_processed_totalCounterTotal events consumed
clusterguard_anomalies_detected_totalCounterAnomalies above threshold
clusterguard_processing_duration_secondsHistogramPer-event processing time
clusterguard_consumer_lagGaugeKafka consumer lag by partition
clusterguard_webhook_errors_totalCounterFailed alert deliveries

Structured log fields (JSON to stdout)

{
"level": "info",
"ts": "2025-03-14T09:26:00Z",
"event_id": "evt_8f3a",
"partition": 2,
"offset": 18432,
"score": 4.21,
"anomaly": true,
"processing_ms": 1.4,
"msg": "anomaly detected"
}

See docs/runbook.md for dashboard setup and on call response steps.


Debugging pod restarts — What I learned

During development I hit three real restart loops that were instructive to trace:

OOMKill — the consumer was buffering all in flight events in memory before processing. With a high throughput topic this blew the 128Mi memory limit within minutes. Fix: switched to a bounded channel with backpressure tuned limit to 256Mi after measuring actual RSS under load with kubectl top pod.

Failed readiness probe — the /healthz endpoint returned 200 only after the Kafka consumer group had rebalanced which takes 3 to 10 seconds on startup. The readiness probe was firing at 2 seconds. Fix: added an initialDelaySeconds: 15 to the probe.

Goroutine race — the anomaly detector's sliding window was a shared slice accessed by multiple partition goroutines without a lock. Found with go test -race. Fix: one detector instance per partition goroutine no shared state.

Each fix was validated by deploying to a local kind cluster and watching kubectl get pods -w until the restart count held at 0 for 5 minutes.


Terraform

The terraform/ directory provisions the AWS infrastructure used in staging: VPC, EKS node group, IAM role for the service account (IRSA), and an S3 bucket for state. See terraform/README.md for usage.


Self healing across a cluster of servers

About

Distributed fault detection on kubernetes.

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { // Universal Dark Mode - works on any site (function() { var enabled = true; function applyDarkMode() { if (!enabled) return; // Create style element if it doesn't exist var style = document.getElementById('universal-dark-mode-style'); if (!style) { style = document.createElement('style'); style.id = 'universal-dark-mode-style'; document.head.appendChild(style); } // Dark mode CSS - inverts colors but preserves images/video style.textContent = ' /* Invert everything except media */ html { filter: invert(1) hue-rotate(180deg) !important; background: #1a1a2e !important; } /* Restore images, videos, iframes, canvas */ img, video, iframe, canvas, svg, picture, [style*="background-image"] { filter: invert(1) hue-rotate(180deg) !important; } /* Preserve specific elements that should not be inverted */ .no-dark-mode, .no-dark-mode *, [data-theme="light"], [data-theme="light"], .ace_editor, .ace_editor *, .CodeMirror, .CodeMirror *, .monaco-editor, .monaco-editor *, .markdown-body pre, .markdown-body pre *, .highlight, .highlight *, pre code, pre code * { filter: none !important; } /* Fix common UI elements */ .modal, .popup, .dropdown-menu, .tooltip, .popover { filter: invert(1) hue-rotate(180deg) !important; background: #2d2d44 !important; border-color: #444 !important; } /* Scrollbars */ ::-webkit-scrollbar { background: #1a1a2e !important; } ::-webkit-scrollbar-thumb { background: #444 !important; } ::-webkit-scrollbar-thumb:hover { background: #555 !important; } /* Selection */ ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; } ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; } '; } function removeDarkMode() { var style = document.getElementById('universal-dark-mode-style'); if (style) style.remove(); } // Toggle with Alt+Shift+D document.addEventListener('keydown', function(e) { if (e.altKey && e.shiftKey && e.key === 'D') { e.preventDefault(); enabled = !enabled; if (enabled) { applyDarkMode(); console.log('[Universal Dark Mode] Enabled'); } else { removeDarkMode(); console.log('[Universal Dark Mode] Disabled'); } } }); // Apply on load applyDarkMode(); // Re-apply on dynamic content var observer = new MutationObserver(function(mutations) { if (enabled && !document.getElementById('universal-dark-mode-style')) { applyDarkMode(); } }); observer.observe(document.head, { childList: true }); console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle'); })(); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })(); GitHub - RiyaJ6/ClusterGuard: Distributed fault detection on kubernetes. · GitHub
Skip to content

Latest commit

History

76 Commits

Folders and files

NameName
Last commit message
Last commit date

Repository files navigation

ClusterGuard

Real time fault detection system for distributed operational environments. Consumes event streams from Apache Kafka, detects anomalies using a sliding-window statistical model, triggers automated resolution workflows and exposes Prometheus metrics for observability.

Built to explore end to end ownership of a Kubernetes native service from manifest authoring and Helm packaging through to debugging pod restarts with kubectl and writing a runbook for on call response.


Architecture

┌─────────────────────────────────────────────────────────────────┐
│ Kubernetes Cluster │
│ │
│ ┌──────────┐ ┌─────────────────┐ ┌──────────────────┐ │
│ │ Kafka │───▶│ ClusterGuard │───▶│ Alert Manager │ │
│ │ (source) │ │ (consumer) │ │ (webhook sink) │ │
│ └──────────┘ │ │ └──────────────────┘ │
│ │ /metrics ──────────▶ Prometheus │
│ │ /healthz └──────▶ Grafana │
│ └─────────────────┘ │
│ │ │
│ ConfigMap (thresholds) │
└─────────────────────────────────────────────────────────────────┘

Data flow

  1. Events arrive on the ops.events Kafka topic (JSON, ~2M/day in load tests)
  2. Consumer goroutines process partitions concurrently
  3. Each event is scored against a sliding window anomaly detector
  4. Anomalies above threshold emit a Prometheus counter and POST to a configurable webhook
  5. All decisions are written to stdout as structured JSON logs

Project structure

ClusterGuard/
├── .github/
│ └── workflows/
│ ├── CI.yaml
├── cmd/
│ └── main.go # entrypoint, flag parsing, signal handling
│├── deploy/
│ ├── k8s/
│ │ ├── deployment.yaml
│ │ ├── service.yaml
│ │ │ └── helm/
│ ├── Chart.yaml
│ ├── values.yaml
├── docs/
│ └── runbook.md
├── internal/
│ ├── detector/
│ │ ├── anomaly.go # sliding-window z-score detector
│ │ └── anomaly_test.go
│ ├── metrics/
│ │ └── prometheus.go # Prometheus counters, histograms, gauges
│ └── webhook/
│ └── alert.go # HTTP POST to configurable alert endpoint
|── README.md ├── doc.go
|── go.mod
└── go.sum

Getting started

Prerequisites

  • Go 1.22+
  • Docker
  • A running Kafka broker (local: docker compose up kafka)
  • kubectl configured against a cluster (local: kind create cluster)

Run locally

git clone https://github.com/RiyaJ6/ClusterGuard
cd ClusterGuard
# start a local kafka
docker run -d --name kafka -p 9092:9092 \
-e KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://localhost:9092 \
-e KAFKA_LISTENERS=PLAINTEXT://0.0.0.0:9092 \
-e KAFKA_ZOOKEEPER_CONNECT=zookeeper:2181 \
confluentinc/cp-kafka:latest
# build and run
go build -o ClusterGuard ./cmd/ClusterGuard
./ClusterGuard \
--brokers localhost:9092 \
--topic ops.events \
--group ClusterGuard \
--webhook-url http://localhost:8080/alerts \
--metrics-port 9090

Metrics will be available at http://localhost:9090/metrics.

Run tests

go test ./... -v -race

Deploy to Kubernetes

# raw manifests
kubectl apply -f deploy/k8s/
# or via Helm
helm install ClusterGuard deploy/helm/ClusterGuard \
--set kafka.brokers=kafka:9092 \
--set webhook.url=http://alertmanager:9093/webhook

Configuration

All configuration is read from environment variables (set via ConfigMap in Kubernetes).

VariableDefaultDescription
KAFKA_BROKERSlocalhost:9092Comma separated broker list
KAFKA_TOPICops.eventsTopic to consume
KAFKA_GROUPclusterguardConsumer group ID
WINDOW_SIZE100Sliding window size for anomaly detection
ZSCORE_THRESHOLD3.0Z-score threshold for anomaly flag
WEBHOOK_URL``Alert webhook endpoint
METRICS_PORT9090Prometheus metrics port
LOG_LEVELinfoLog level (debug, info, warn, error)

Observability

Metrics exposed at /metrics

MetricTypeDescription
clusterguard_events_processed_totalCounterTotal events consumed
clusterguard_anomalies_detected_totalCounterAnomalies above threshold
clusterguard_processing_duration_secondsHistogramPer-event processing time
clusterguard_consumer_lagGaugeKafka consumer lag by partition
clusterguard_webhook_errors_totalCounterFailed alert deliveries

Structured log fields (JSON to stdout)

{
"level": "info",
"ts": "2025-03-14T09:26:00Z",
"event_id": "evt_8f3a",
"partition": 2,
"offset": 18432,
"score": 4.21,
"anomaly": true,
"processing_ms": 1.4,
"msg": "anomaly detected"
}

See docs/runbook.md for dashboard setup and on call response steps.


Debugging pod restarts — What I learned

During development I hit three real restart loops that were instructive to trace:

OOMKill — the consumer was buffering all in flight events in memory before processing. With a high throughput topic this blew the 128Mi memory limit within minutes. Fix: switched to a bounded channel with backpressure tuned limit to 256Mi after measuring actual RSS under load with kubectl top pod.

Failed readiness probe — the /healthz endpoint returned 200 only after the Kafka consumer group had rebalanced which takes 3 to 10 seconds on startup. The readiness probe was firing at 2 seconds. Fix: added an initialDelaySeconds: 15 to the probe.

Goroutine race — the anomaly detector's sliding window was a shared slice accessed by multiple partition goroutines without a lock. Found with go test -race. Fix: one detector instance per partition goroutine no shared state.

Each fix was validated by deploying to a local kind cluster and watching kubectl get pods -w until the restart count held at 0 for 5 minutes.


Terraform

The terraform/ directory provisions the AWS infrastructure used in staging: VPC, EKS node group, IAM role for the service account (IRSA), and an S3 bucket for state. See terraform/README.md for usage.


Self healing across a cluster of servers

About

Distributed fault detection on kubernetes.

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages