Skip to content

Latest commit

Β 

History

28 Commits

Folders and files

NameName
Last commit message
Last commit date
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 

Repository files navigation

🎬 StreamGuard AI

Agentic AI Reliability Supervisor for Cinema & Video Streaming

Google Cloud Gemini ADK Grafana MCP Confluent MCP License: MIT


πŸš€ Quick Links


πŸ“Œ Executive Summary

StreamGuard AI is an autonomous, agentic AI reliability and continuity supervisor built for cinema, OTT, and live video-streaming platforms.

Built for the Google Cloud Agentic Cinema: The Blockbuster Hackathon (Grafana Labs Track), StreamGuard AI replaces manual dashboard hunting during major movie premieres, live cinema events, or esports streams. It autonomously correlates Grafana Cloud observability telemetry with Confluent Kafka event streams using Google Agent Development Kit (ADK) and Gemini 2.5/3.6 to reason over operational evidence in real time.

Instead of requiring engineers to manually inspect dozens of disconnected dashboards, Loki logs, Prometheus metrics, alerts, and event topics under intense pressure, StreamGuard AI executes a multi-source investigation and delivers an evidence-based operational report.

It answers the critical operational questions:

"What is failing? Where is it failing? What evidence supports the diagnosis? What might be causing it? How many viewers could be affected? And what should the operator investigate next?"


🎯 The Problem & The Solution

πŸ’₯ The Problem

Modern OTT and cinema streaming pipelines involve complex, multi-tier infrastructure:

Ingest Pipeline βž” Transcoding Clusters βž” Origin Infrastructure βž” Regional CDN Egress βž” Edge Nodes βž” Viewer Playback

During high-demand movie releases or live broadcasts, infrastructure problems rapidly propagate into viewer-facing incidents (buffering, dropped frames, transcode queue saturation, 5xx gateway errors).

Traditional monitoring tools provide raw data, but SRE teams still have to manually correlate signals across disparate systems under extreme time pressure.

The real operational challenge is not simply:

"Is CPU utilization high?"

It is:

"Is the streaming experience degrading, where is it happening, what evidence explains it, and what exact steps should the SRE team take next?"


πŸ’‘ The Solution: StreamGuard AI Workflow

               β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
               β”‚    Operator Request / Autonomous Goal     β”‚
               β”‚   "Investigate current streaming health"  β”‚
               β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
                                     β”‚
                                     β–Ό
               β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
               β”‚         StreamGuard AI Root Agent         β”‚
               β”‚            (Google ADK + Gemini)          β”‚
               β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
                                     β”‚
                 β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”΄β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
                 β–Ό                                       β–Ό
    β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”           β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
    β”‚ Broadcast Monitoring Agentβ”‚           β”‚   Event Streaming Agent   β”‚
    β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜           β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
                 β”‚                                       β”‚
                 β–Ό                                       β–Ό
    β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”           β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
    β”‚     Grafana Cloud MCP     β”‚           β”‚       Confluent MCP       β”‚
    β”‚   (grafana/mcp-grafana)   β”‚           β”‚   (Kafka Stream Health)   β”‚
    β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜           β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
                 β”‚                                       β”‚
         β”Œβ”€β”€β”€β”€β”€β”€β”€β”΄β”€β”€β”€β”€β”€β”€β”€β”                       β”Œβ”€β”€β”€β”€β”€β”€β”€β”΄β”€β”€β”€β”€β”€β”€β”€β”
         β–Ό               β–Ό                       β–Ό               β–Ό
  Prometheus Metrics  Loki Logs             Kafka Topics   Stream Events
         β”‚               β”‚                       β”‚               β”‚
         β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”΄β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”΄β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
                                     β”‚
                                     β–Ό
                         Multi-Source Correlation
                                     β”‚
                                     β–Ό
                        3-Tier Evidence Engine
                     (Verified | Derived | Hypotheses)
                                     β”‚
                                     β–Ό
                     Closed-Loop Grafana Write-Back
                 (Dashboard Event Annotations #ANN)

🧠 Core System Capabilities

1. πŸ€– Multi-Agent Investigation Architecture (Google ADK)

StreamGuard AI uses specialized Google ADK agents to isolate operational concerns:

  • Root Agent: Accepts the operator goal, orchestrates subagent handoffs, synthesizes evidence, and generates final operational reports.
  • Broadcast Monitoring Agent: Specializes in Grafana Cloud observability via grafana/mcp-grafana.
  • Event Streaming Agent: Specializes in Confluent Kafka event topic discovery and message stream analysis.

2. πŸ“Š Grafana Cloud MCP Integration (grafana/mcp-grafana)

Demonstrates active runtime usage of the official Grafana MCP Server:

  • PromQL Metrics (query_prometheus): Queries ingest CPU load, transcode queue depth, CDN egress bandwidth, and viewer buffer ratios.
  • LogQL Logs (query_loki_logs): Inspects 5xx HTTP gateway timeouts and transcode drop-frame stack traces.
  • Incidents API (list_incidents, create_grafana_incident): Inspects and opens official Grafana incidents.

3. 🌊 Confluent Kafka MCP Integration

Inspects the event-streaming layer:

  • Discovers Kafka topics (stream-health, continuity-alerts, broadcast-incidents).
  • Inspects schema subjects and consumes real-time stream-health events.
  • Correlates Kafka stream drops with Grafana telemetry.

4. πŸ”Ž 3-Tier Evidence-Based Reasoning

Prevents assumptions from being presented as confirmed telemetry by categorizing findings into:

  • βœ… Verified Evidence: Direct telemetry returned by Grafana or Confluent tools (e.g. Transcode CPU > 94%).
  • πŸ”Ž Derived Observations: Cross-system correlated patterns (e.g. CPU spikes coincide with dropped frames 2m later).
  • ⚠️ Hypotheses: Possible explanations requiring further profiling (e.g. Background job competing for CPU).

5. 🌎 Regional Correlation Safety

Enforces strict tag validation: if telemetry lacks explicit region_id tags, the agent never invents regional claims, explicitly outputting a disclaimer instead.

6. ✍️ Bidirectional Closed-Loop Grafana Write-Back

When operational context is established, StreamGuard AI invokes annotate_grafana_dashboard to place event markers directly onto live Grafana SRE dashboards (#ANN-8924).

7. πŸ“ˆ Grafana Agent Observability (Observing the AI)

Makes the AI agent itself observable, tracking model calls, execution latency, token counts, and estimated cost in real-time.


πŸ—οΈ System Architecture Diagrams

High-Level System Architecture Flowchart

flowchart TD
    classDef operatorStyle fill:#0284c7,color:#fff,stroke:#0369a1,stroke-width:2px;
    classDef agentStyle fill:#7c3aed,color:#fff,stroke:#5b21b6,stroke-width:2px;
    classDef mcpStyle fill:#0ea5e9,color:#fff,stroke:#0369a1,stroke-width:2px;
    classDef dataStyle fill:#d97706,color:#fff,stroke:#92400e,stroke-width:2px;
    classDef engineStyle fill:#059669,color:#fff,stroke:#047857,stroke-width:2px;
    classDef outputStyle fill:#db2777,color:#fff,stroke:#9d174d,stroke-width:2px;

    Op["πŸ‘€ Streaming Operator / Technical Director"]:::operatorStyle
    RootAgent["🧠 StreamGuard AI Root Agent"]:::agentStyle
    ADK["⚑ Google ADK + Gemini (2.5 Flash / 3.6 Pro)"]:::agentStyle

    subgraph AgentLayer ["πŸ€– Multi-Agent Delegation Layer"]
        BMAgent["πŸŽ₯ Broadcast Monitoring Agent\n(Grafana Telemetry & Incidents)"]:::agentStyle
        ESAgent["πŸ“‘ Event Streaming Agent\n(Kafka Topics & Event Streams)"]:::agentStyle
    end

    subgraph MCPLayer ["πŸ”Œ Model Context Protocol (MCP) Integration"]
        GrafanaMCP["πŸ“Š Grafana Cloud MCP Server\n(grafana/mcp-grafana)"]:::mcpStyle
        ConfluentMCP["🌊 Confluent Kafka MCP Server"]:::mcpStyle
    end

    subgraph DataLayer ["πŸ—„οΈ Multi-System Observability Data"]
        subgraph GrafanaCloud ["Grafana Cloud Platform"]
            PromData["πŸ“ˆ Prometheus Metrics\n(CPU, Latency, Buffer Ratio, Frame Drops)"]:::dataStyle
            LokiData["πŸ“‹ Loki Logs\n(HTTP 502/504, Stack Traces)"]:::dataStyle
            AlertData["🚨 Alert Groups & Incidents\n(Active Alerts, Severity)"]:::dataStyle
        end

        subgraph ConfluentCloud ["Confluent Cloud Kafka"]
            KafkaTopics["πŸ’¬ Stream-Health Topics\n(stream-health, broadcast-incidents)"]:::dataStyle
            KafkaSchema["πŸ“œ Schema Registry & Subjects"]:::dataStyle
            KafkaMsgs["βœ‰οΈ Real-Time Kafka Event Messages"]:::dataStyle
        end
    end

    subgraph ReasoningEngine ["🧠 Evidence Correlation & RCA Engine"]
        MultiCorr["πŸ”— Cross-System Signal Correlation\n(CPU βž” Latency βž” Frame Drop βž” Buffering)"]:::engineStyle
        EvidClass["πŸ”Ž Evidence Classifier\n(Verified Evidence | Derived Obs | Hypotheses)"]:::engineStyle
        RegSafety["🌎 Regional Correlation Safety Evaluator\n(Strict Identifier Check)"]:::engineStyle
        ReportGen["πŸ“ Evidence-Based Operational Report Generator"]:::engineStyle
    end

    subgraph OutputLayer ["πŸ“€ Closed-Loop Actions & Observability"]
        WriteBack["✍️ Grafana Write-Back\n(Dashboard Event Annotations & Incident Logs)"]:::outputStyle
        AgentO11y["πŸ“Š Grafana Agent Observability\n(Token Usage, Latency, Cost, Model Traces)"]:::outputStyle
        OpReport["πŸ“„ Final Incident Investigation Report\n(Root Cause, Impact, Recommended Actions)"]:::outputStyle
    end

    Op --> RootAgent
    RootAgent --> ADK
    ADK --> BMAgent
    ADK --> ESAgent

    BMAgent --> GrafanaMCP
    ESAgent --> ConfluentMCP

    GrafanaMCP --> PromData & LokiData & AlertData
    ConfluentMCP --> KafkaTopics & KafkaSchema & KafkaMsgs

    PromData & LokiData & AlertData --> MultiCorr
    KafkaTopics & KafkaMsgs --> MultiCorr

    MultiCorr --> EvidClass --> RegSafety --> ReportGen
    ReportGen --> OpReport & WriteBack
    ADK --> AgentO11y
Loading

πŸ“ Repository Structure

StreamGuard-AI/
β”‚
β”œβ”€β”€ backend/                        # FastAPI Service Backend & MCP Connectors
β”‚   β”œβ”€β”€ app/
β”‚   β”‚   β”œβ”€β”€ agents/                 # Google ADK Agent definitions & prompts
β”‚   β”‚   β”œβ”€β”€ mcp/                    # Grafana Cloud & Confluent MCP integrations
β”‚   β”‚   └── simulator/              # Live telemetry anomaly injector
β”‚   β”œβ”€β”€ main.py                     # Entry point server
β”‚   └── requirements.txt
β”‚
β”œβ”€β”€ frontend/                       # StreamGuard AI Director Console UI
β”‚   β”œβ”€β”€ index.html                  # Live dashboard interface
β”‚   β”œβ”€β”€ css/
β”‚   └── js/
β”‚
β”œβ”€β”€ diagrams/                       # High-resolution SVG/PNG/PDF architecture diagrams
β”‚   β”œβ”€β”€ StreamGuard_AI_Architecture_Diagrams.pdf
β”‚   β”œβ”€β”€ Architecture_Flowcharts_Interactive.html
β”‚   β”œβ”€β”€ Sequence_Flow_Interactive.html
β”‚   └── StreamGuard_AI_Architecture_Diagram_Clean.png
β”‚
β”œβ”€β”€ streamguard_agent/              # ADK Web Agent definition package
β”‚   └── agent.py
β”‚
β”œβ”€β”€ .gcp/                           # Google Cloud Run deployment scripts
β”‚   └── deploy.sh
β”‚
β”œβ”€β”€ Dockerfile                      # Container build definition
β”œβ”€β”€ requirements.txt                # Python dependencies
β”œβ”€β”€ README.md                       # Documentation
└── LICENSE                         # MIT License

βš™οΈ Quick Start & Local Setup

Prerequisites

  • Python 3.11+
  • Google Cloud Project with Vertex AI / Gemini API enabled
  • Grafana Cloud Account + Service Account Token
  • Confluent Cloud Account + API Credentials

1. Clone & Install Dependencies

git clone https://github.com/Shrushti72/StreamGuard-AI.git
cd StreamGuard-AI

pip install -r requirements.txt

2. Configure Environment Credentials

Create a .env file or export environment variables:

# Gemini / Vertex AI Credentials
export GEMINI_API_KEY="your-gemini-api-key"

# Grafana Cloud Credentials
export GRAFANA_URL="https://your-instance.grafana.net"
export GRAFANA_SERVICE_ACCOUNT_TOKEN="your-grafana-token"

# Confluent Cloud Credentials
export CONFLUENT_API_KEY="your-confluent-key"
export CONFLUENT_API_SECRET="your-confluent-secret"

# Grafana Agent Observability Credentials
export AGENTO11Y_ENDPOINT="https://otlp-gateway.grafana.net"
export AGENTO11Y_AUTH_TOKEN="your-agent-o11y-token"

3. Run Locally via Google ADK CLI

adk web streamguard_agent \
  --host 0.0.0.0 \
  --port 8080

Open http://localhost:8080 to access the StreamGuard AI Director Console.


☁️ Deployment to Google Cloud Run

Deploy serverless to Google Cloud Run using Google Cloud Secret Manager for credentials:

# Submit build to Google Cloud Artifact Registry
gcloud builds submit --tag gcr.io/$PROJECT_ID/streamguard-ai

# Deploy to Cloud Run with injected secrets
gcloud run deploy streamguard-ai \
  --image gcr.io/$PROJECT_ID/streamguard-ai \
  --platform managed \
  --region us-central1 \
  --allow-unauthenticated \
  --set-secrets="GRAFANA_SERVICE_ACCOUNT_TOKEN=grafana-token:latest,CONFLUENT_API_SECRET=confluent-secret:latest"

πŸ§ͺ Hackathon Component Validation

Component / Feature Technology Stack Status
Agent Reasoning Engine Google Gemini 2.5 Flash / 3.6 Pro βœ… Verified
Multi-Agent Orchestration Google Agent Development Kit (ADK) βœ… Verified
Cloud Hosting Google Cloud Run (us-central1) βœ… Live & Deployed
Grafana Metrics (PromQL) Grafana Cloud MCP (query_prometheus) βœ… Verified
Grafana Logs (LogQL) Grafana Cloud MCP (query_loki_logs) βœ… Verified
Grafana Write-Back Grafana Cloud MCP (annotate_grafana_dashboard) βœ… Verified
Kafka Event Streaming Confluent MCP (consume_kafka_messages) βœ… Verified
AI Agent Observability Grafana Agent Observability (OTLP) βœ… Verified
Credential Security Google Cloud Secret Manager βœ… Configured

πŸ“œ License

This project is licensed under the MIT License. See the LICENSE file for details.


πŸ‘©β€πŸ’» Author & Acknowledgments

Shrushti Wakchaure β€” Built for the Google Cloud Agentic Cinema: The Blockbuster Hackathon.

Special thanks to Google Cloud and Grafana Labs for providing the ADK framework and official Grafana Cloud MCP Server integrations.

About

No description, website, or topics provided.

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages