id: ai-multi-provider-fallback-router
namespace: company.ai
inputs:
- id: prompt
type: STRING
displayName: User Prompt
description: Text prompt to submit to the large language model.
defaults: Explain the architecture of an event-driven workflow orchestrator in 3
concise bullet points.
- id: system_prompt
type: STRING
displayName: System Instructions
description: System persona instructions guiding the model behavior.
defaults: You are an expert distributed systems engineer and technical architect.
- id: temperature
type: FLOAT
displayName: Model Temperature
description: Sampling temperature for generation randomness (0.0 to 1.0).
defaults: 0.3
- id: timeout_seconds
type: INT
displayName: Primary Provider Timeout (Seconds)
description: Maximum seconds to wait for primary provider response before
triggering circuit breaker failover.
defaults: 8
- id: notify_on_failover
type: BOOL
displayName: Notify on Failover
description: Send an incident notification card to Slack when automatic provider
failover is triggered.
defaults: true
concurrency:
limit: 10
triggers:
- id: llm_router_webhook
type: io.kestra.plugin.core.trigger.Webhook
description: Event-driven Webhook endpoint ready to receive inference requests
from downstream applications and agents.
key: llm-router-endpoint
tasks:
- id: log_request
type: io.kestra.plugin.core.log.Log
description: Log incoming inference request metadata for observability.
message: >-
Inference request received: Prompt length: {{ inputs.prompt | length }}
characters, Primary Timeout: {{ inputs.timeout_seconds }}s, Temperature:
{{ inputs.temperature }}.
- id: route_inference
type: io.kestra.plugin.scripts.python.Script
description: Attempt LLM generation against primary provider (OpenAI) with
automated circuit breaker failover to secondary provider (Google Gemini)
on rate limits, timeouts, or errors.
containerImage: python:3.11-slim
taskRunner:
type: io.kestra.plugin.scripts.runner.docker.Docker
beforeCommands:
- pip install --no-cache-dir kestra requests
env:
OPENAI_API_KEY: "{{ secret('OPENAI_API_KEY') }}"
GEMINI_API_KEY: "{{ secret('GEMINI_API_KEY') }}"
script: |
import os
import time
import json
import requests
from kestra import Kestra
prompt = """{{ inputs.prompt }}""".strip()
system_prompt = """{{ inputs.system_prompt }}""".strip()
temperature = float("{{ inputs.temperature }}")
timeout_sec = int("{{ inputs.timeout_seconds }}")
openai_key = os.environ.get("OPENAI_API_KEY", "").strip()
gemini_key = os.environ.get("GEMINI_API_KEY", "").strip()
response_text = ""
active_provider = ""
fallback_triggered = False
primary_error = ""
start_time = time.time()
# Tier 1: Primary Provider - OpenAI (gpt-4o-mini)
if openai_key:
print("Tier 1: Attempting inference via primary provider (OpenAI gpt-4o-mini)...")
try:
headers = {
"Authorization": f"Bearer {openai_key}",
"Content-Type": "application/json"
}
payload = {
"model": "gpt-4o-mini",
"messages": [
{"role": "system", "content": system_prompt},
{"role": "user", "content": prompt}
],
"temperature": temperature
}
resp = requests.post(
"https://api.openai.com/v1/chat/completions",
headers=headers,
json=payload,
timeout=timeout_sec
)
if resp.status_code == 200:
data = resp.json()
response_text = data["choices"][0]["message"]["content"]
active_provider = "OpenAI (gpt-4o-mini)"
fallback_triggered = False
print("Primary provider (OpenAI) responded successfully.")
else:
primary_error = f"HTTP {resp.status_code}: {resp.text[:200]}"
print(f"Primary provider failed with status {resp.status_code}. Initiating failover...")
fallback_triggered = True
except Exception as e:
primary_error = f"Exception: {str(e)}"
print(f"Primary provider connection failed/timed out: {e}. Initiating failover...")
fallback_triggered = True
else:
primary_error = "OPENAI_API_KEY missing or empty."
fallback_triggered = True
# Tier 2: Fallback Provider - Google Gemini (gemini-1.5-flash)
if fallback_triggered:
print("Tier 2: Engaging circuit breaker fallback to secondary provider (Google Gemini 1.5 Flash)...")
if not gemini_key:
raise RuntimeError(f"Primary failed ({primary_error}) and GEMINI_API_KEY is not configured for fallback.")
gemini_url = f"https://generativelanguage.googleapis.com/v1beta/models/gemini-1.5-flash:generateContent?key={gemini_key}"
gemini_payload = {
"contents": [
{
"parts": [
{"text": f"System context: {system_prompt}\n\nUser request: {prompt}"}
]
}
],
"generationConfig": {
"temperature": temperature
}
}
resp_gemini = requests.post(gemini_url, json=gemini_payload, timeout=15)
if resp_gemini.status_code == 200:
g_data = resp_gemini.json()
response_text = g_data["candidates"][0]["content"]["parts"][0]["text"]
active_provider = "Google Gemini (gemini-1.5-flash)"
print("Fallback provider (Google Gemini) responded successfully.")
else:
raise RuntimeError(f"Both providers failed! Primary error: {primary_error}, Gemini error: HTTP {resp_gemini.status_code} {resp_gemini.text[:200]}")
elapsed_ms = int((time.time() - start_time) * 1000)
print(f"Inference completed in {elapsed_ms}ms via {active_provider}. Failover: {fallback_triggered}")
Kestra.outputs({
"response_text": response_text,
"active_provider": active_provider,
"fallback_triggered": fallback_triggered,
"latency_ms": elapsed_ms,
"primary_error": primary_error or "None"
})
- id: check_failover
type: io.kestra.plugin.core.flow.If
description: Branch based on whether an automatic provider failover was triggered.
condition: "{{ outputs.route_inference.vars.fallback_triggered }}"
then:
- id: notify_slack_failover
type: io.kestra.plugin.core.flow.If
description: Post failover alert to Slack if notifications are enabled.
condition: "{{ inputs.notify_on_failover }}"
then:
- id: slack_failover_card
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
description: Send failover incident card to Slack AI operations channel.
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"text": "LLM Circuit Breaker Failover Engaged!\n• Primary Provider Failed: OpenAI (`{{ outputs.route_inference.vars.primary_error }}`)\n• Fallback Provider Activated: `{{ outputs.route_inference.vars.active_provider }}`\n• Total Latency: {{ outputs.route_inference.vars.latency_ms }}ms\n• Request Status: Downstream request served successfully without interruption.\nExecution: {{ execution.id }}"
}
else:
- id: log_primary_success
type: io.kestra.plugin.core.log.Log
description: Record healthy primary provider generation telemetry in execution logs.
message: >-
Inference fulfilled by primary provider: {{
outputs.route_inference.vars.active_provider }}. Latency: {{
outputs.route_inference.vars.latency_ms }}ms. Fallback was not needed.
errors:
- id: on_pipeline_error
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
description: Alert the AI platform team if all configured LLM providers fail or
credentials are misconfigured.
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"text": "CRITICAL: LLM Multi-Provider Gateway completely failed in flow `{{ flow.id }}`! Both primary and fallback providers failed to generate completion. Execution: {{ execution.id }}."
}
outputs:
- id: response_text
type: STRING
description: Normalized text completion generated by the active model provider.
value: "{{ outputs.route_inference.vars.response_text }}"
- id: active_provider
type: STRING
description: Name of the model provider that fulfilled the generation request.
value: "{{ outputs.route_inference.vars.active_provider }}"
- id: fallback_triggered
type: BOOL
description: Indicates whether automatic circuit breaker failover to the
secondary provider occurred.
value: "{{ outputs.route_inference.vars.fallback_triggered }}"
- id: latency_ms
type: INT
description: End-to-end inference turnaround time in milliseconds.
value: "{{ outputs.route_inference.vars.latency_ms }}"