Class: RubyLlmMesh::Router

Inherits:
Object
  • Object
show all
Defined in:
lib/ruby_llm_mesh/router.rb,
sig/ruby_llm_mesh.rbs

Constant Summary collapse

PROVIDER_MAP =
{
  openai: Providers::Openai,
  anthropic: Providers::Anthropic,
  local_node: Providers::LocalNode,
  local_mesh: Providers::LocalNode # alias — multi-peer aware via PeerRegistry
}.freeze

Class Method Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(config: RubyLlmMesh.configuration, circuit_breaker: self.class.circuit_breaker, semantic_cache: nil, budget: nil) ⇒ Router

Returns a new instance of Router.



37
38
39
40
41
42
43
# File 'lib/ruby_llm_mesh/router.rb', line 37

def initialize(config: RubyLlmMesh.configuration, circuit_breaker: self.class.circuit_breaker,
               semantic_cache: nil, budget: nil)
  @config = config
  @circuit_breaker = circuit_breaker
  @semantic_cache = semantic_cache
  @budget = budget
end

Class Method Details

.circuit_breakerObject



13
14
15
16
17
18
# File 'lib/ruby_llm_mesh/router.rb', line 13

def circuit_breaker
  @circuit_breaker ||= CircuitBreaker.new(
    failure_threshold: RubyLlmMesh.configuration.circuit_failure_threshold,
    reset_timeout: RubyLlmMesh.configuration.circuit_reset_timeout
  )
end

.health_monitorObject



32
33
34
# File 'lib/ruby_llm_mesh/router.rb', line 32

def health_monitor
  Mesh::HealthMonitor.instance
end

.peer_registryObject



28
29
30
# File 'lib/ruby_llm_mesh/router.rb', line 28

def peer_registry
  Mesh::PeerRegistry.instance
end

.reset_circuit_breaker!Object



20
21
22
# File 'lib/ruby_llm_mesh/router.rb', line 20

def reset_circuit_breaker!
  @circuit_breaker = nil
end

.semantic_cacheObject



24
25
26
# File 'lib/ruby_llm_mesh/router.rb', line 24

def semantic_cache
  Cache::SemanticCache.instance
end

Instance Method Details

#complete(prompt:, providers: nil, fallback: nil, system: nil, model: nil, **options) ⇒ Object

Raises:

  • (ArgumentError)


45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
# File 'lib/ruby_llm_mesh/router.rb', line 45

def complete(prompt:, providers: nil, fallback: nil, system: nil, model: nil, **options)
  raise ArgumentError, "prompt is required" if prompt.nil? || prompt.to_s.strip.empty?

  skip_cache = options.delete(:skip_cache)
  cache = resolve_cache
  unless skip_cache
    cached = cache&.lookup(prompt, system: system)
    if cached
      log(:info, "Semantic cache hit")
      return cached
    end
  end

  maybe_refresh_peer_health!

  provider_list = Array(providers || @config.default_providers).map(&:to_sym)
  raise ArgumentError, "providers list cannot be empty" if provider_list.empty?

  use_fallback = fallback.nil? ? @config.fallback : fallback
  errors = {}
  attempted = 0

  provider_list.each_with_index do |provider_name, index|
    circuit_key = circuit_key_for(provider_name)

    unless PROVIDER_MAP.key?(provider_name)
      errors[provider_name] = ProviderError.new("Unknown provider: #{provider_name}", provider: provider_name)
      break unless use_fallback
      next
    end

    unless @circuit_breaker.allow?(circuit_key)
      errors[provider_name] = CircuitOpenError.new(
        "Circuit open for #{provider_name}",
        provider: provider_name
      )
      log(:warn, "Skipping #{provider_name} — circuit open")
      break unless use_fallback
      next
    end

    begin
      attempted += 1
      log(:info, "Routing to #{provider_name}")
      budget = resolve_budget
      estimated_tokens = Budget.estimate_tokens(prompt: prompt, system: system, max_tokens: options[:max_tokens])
      estimated_usd = Budget.estimate_usd(tokens: estimated_tokens, model: model, config: @config)
      budget.check!(estimated_tokens: estimated_tokens, estimated_usd: estimated_usd)

      provider = PROVIDER_MAP[provider_name].new(@config)
      response = with_retries(provider_name) do
        provider.complete(prompt: prompt, system: system, model: model, **options)
      end
      budget.consume!(usage: response.usage, provider: provider_name, model: response.model || model)
      @circuit_breaker.record_success(circuit_key)

      result = Response.new(
        content: response.content,
        provider: response.provider,
        model: response.model,
        usage: response.usage,
        raw: response.raw,
        latency_ms: response.latency_ms,
        fallback_used: index.positive?,
        cache_hit: false
      )
      cache&.store_response(prompt, result, system: system) unless skip_cache
      return result
    rescue BudgetExceededError
      raise
    rescue AuthenticationError => e
      @circuit_breaker.record_failure(circuit_key)
      errors[provider_name] = e
      log(:error, "#{provider_name} authentication failed: #{e.message}")
      break unless use_fallback
    rescue RateLimitError => e
      # Trip circuit immediately so subsequent requests skip this provider
      force_open_circuit!(circuit_key)
      errors[provider_name] = e
      log(:error, "#{provider_name} rate limited: #{e.message}")
      break unless use_fallback
    rescue ProviderError => e
      @circuit_breaker.record_failure(circuit_key)
      errors[provider_name] = e
      log(:error, "#{provider_name} failed: #{e.message}")
      break unless use_fallback
    end
  end

  raise AllProvidersFailedError, errors
end