Class: Legion::LLM::Inference::Executor
- Inherits:
-
Object
- Object
- Legion::LLM::Inference::Executor
- Includes:
- ContextWindow, Escalation, Routing, ToolInjection, NativeToolLoop, RouteAttempts, Steps::Billing, Steps::Classification, Steps::ConfidenceScoring, Steps::Debate, Steps::GaiaAdvisory, Steps::GutCheck, Steps::KnowledgeCapture, Steps::Logging, Steps::Metering, Steps::PostResponse, Steps::PromptCache, Steps::RagContext, Steps::Rbac, Steps::SkillInjector, Steps::StickyPersist, Steps::StickyRunners, Steps::TokenBudget, Steps::ToolCalls, Steps::ToolDiscovery, Steps::ToolHistory, Steps::TriggerMatch, Legion::Logging::Helper
- Defined in:
- lib/legion/llm/inference/executor.rb,
lib/legion/llm/inference/executor/routing.rb,
lib/legion/llm/inference/executor/escalation.rb,
lib/legion/llm/inference/executor/context_window.rb,
lib/legion/llm/inference/executor/tool_injection.rb,
lib/legion/llm/inference/executor/payload_builder.rb
Defined Under Namespace
Modules: ContextWindow, Escalation, PayloadBuilder, Routing, ToolInjection Classes: ToolResultEvent
Constant Summary collapse
- PRE_PROVIDER_STEPS =
%i[ tracing_init idempotency conversation_uuid context_load rbac classification billing gaia_advisory tier_assignment rag_context trigger_match sticky_runners skill_injector tool_history_inject tool_discovery routing request_normalization token_budget ].freeze
- POST_PROVIDER_STEPS =
%i[ response_normalization post_response gut_check metering debate confidence_scoring tool_calls sticky_persist context_store knowledge_capture response_return ].freeze
- STEPS =
(PRE_PROVIDER_STEPS + %i[provider_call] + POST_PROVIDER_STEPS).freeze
- ASYNC_SAFE_STEPS =
%i[post_response knowledge_capture].freeze
- THINKING_TAG_PAIRS =
[ ['<thinking>', '</thinking>'], ['<think>', '</think>'], ['<thought>', '</thought>'] ].freeze
- CONFIG_ERROR_PATTERNS =
[ /AccessDeniedException/, /InvalidModel/i, /model.*not found/i, /not authorized/i, /AWS Marketplace/i ].freeze
- REQUEST_PAYLOAD_ERROR_PATTERNS =
[ /input_schema/i, /tools\.\d+/, /messages\.\d+/, /Field required/i, /ValidationException/ ].freeze
- CONTEXT_OVERFLOW_ERROR_PATTERNS =
[ /maximum context length/i, /context length.*input_tokens/i, /prompt contains at least \d+ input tokens/i ].freeze
- ASYNC_THREAD_POOL =
Concurrent::FixedThreadPool.new(4, fallback_policy: :caller_runs)
Constants included from Steps::StickyPersist
Steps::StickyPersist::SENSITIVE_PARAM_NAMES
Constants included from Steps::Debate
Steps::Debate::CHALLENGER_PROMPT, Steps::Debate::JUDGE_PROMPT, Steps::Debate::REBUTTAL_PROMPT
Constants included from Steps::KnowledgeCapture
Steps::KnowledgeCapture::EMBED_MAX_CHARS
Constants included from Steps::TriggerMatch
Steps::TriggerMatch::HARNESS_PREFIXES
Constants included from Steps::Classification
Steps::Classification::EMAIL_PATTERN, Steps::Classification::LEVELS, Steps::Classification::PHI_KEYWORDS, Steps::Classification::PII_PATTERNS, Steps::Classification::PII_PATTERNS_CORE, Steps::Classification::PII_PATTERNS_EXTENDED
Constants included from ToolInjection
ToolInjection::GENERATION_PARAMS
Constants included from NativeToolLoop
NativeToolLoop::LEAKED_TOKEN_ARG_RE, NativeToolLoop::LEAKED_TOKEN_RE, NativeToolLoop::QWEN_PARAM_RE, NativeToolLoop::QWEN_TOOL_USE_RE
Instance Attribute Summary collapse
-
#audit ⇒ Object
readonly
Returns the value of attribute audit.
-
#confidence_score ⇒ Object
readonly
Returns the value of attribute confidence_score.
-
#discovered_tools ⇒ Object
readonly
Returns the value of attribute discovered_tools.
-
#enrichments ⇒ Object
readonly
Returns the value of attribute enrichments.
-
#profile ⇒ Object
readonly
Returns the value of attribute profile.
-
#request ⇒ Object
readonly
Returns the value of attribute request.
-
#timeline ⇒ Object
readonly
Returns the value of attribute timeline.
-
#tool_event_handler ⇒ Object
Returns the value of attribute tool_event_handler.
-
#tracing ⇒ Object
readonly
Returns the value of attribute tracing.
-
#warnings ⇒ Object
readonly
Returns the value of attribute warnings.
Instance Method Summary collapse
- #call ⇒ Object
-
#call_responses(stream: false, stream_observer: nil, &block) ⇒ Object
N×N: Delegates to the canonical execution path.
- #call_stream(stream_observer: nil, &block) ⇒ Object
- #context_accounting ⇒ Object
-
#initialize(request) ⇒ Executor
constructor
A new instance of Executor.
-
#release_preflight_lease ⇒ Object
Release the streaming preflight DispatchLease (SSOT v3 §15.2).
-
#stream_preflight! ⇒ Object
SSOT v3 §19 streaming preflight.
Methods included from Steps::GutCheck
Methods included from Steps::StickyPersist
Methods included from Steps::ToolHistory
Methods included from Steps::StickyRunners
Methods included from Steps::Metering
build_event, content_fields, flush_spool, identity_fields, operational_fields, publish_event, publish_or_spool, timing_and_context, token_fields
Methods included from Steps::Debate
#debate_enabled?, #gaia_debate_trigger?, #run_debate, #step_debate
Methods included from Steps::PromptCache
#apply_cache_control, #apply_conversation_breakpoint, #sort_tools_deterministically
Methods included from Steps::TokenBudget
Methods included from Steps::ConfidenceScoring
Methods included from Steps::KnowledgeCapture
Methods included from Steps::ToolCalls
Methods included from Steps::ToolDiscovery
Methods included from Steps::SkillInjector
Methods included from Steps::TriggerMatch
#extract_recent_text, #extract_session_text, #format_trigger_log_value, #harness_message?, #log_trigger_match, #normalize_message_words, #rank_and_cap, #record_trigger_match_timeline, #sanitize_trigger_text, #settings_extension_tool_matches, #settings_extensions, #step_trigger_match, #strip_tagged_blocks, #subtract_always_loaded, #trigger_scan_depth, #trigger_tool_limit, #trigger_tool_name, #trigger_tool_name_for_word_match, #trigger_words_for_entry
Methods included from Steps::RagContext
Methods included from Steps::PostResponse
Methods included from Steps::GaiaAdvisory
#build_partner_context, #step_gaia_advisory
Methods included from Steps::Billing
Methods included from Steps::Classification
Methods included from Steps::Rbac
Methods included from ToolInjection
#add_native_tool_definition, #add_pinned_special_tool_definitions, #add_registry_tool_definitions, #add_requested_deferred_tool_definitions_from_settings, #add_settings_extensions_tool_definitions, #apply_generation_params!, #client_tool_passthrough_allowed?, #client_tool_passthrough_enabled?, #client_tool_passthrough_list, #client_tool_passthrough_name_variants, #client_tool_policy_variants, #native_dispatch_chat_options, #native_dispatch_options, #native_dispatch_thinking, #native_dispatch_tools, #native_tool_definition_duplicate?, #native_tool_definition_name_variants, #native_tool_definitions, #native_tool_loop_continuation_prompt, #native_tool_loop_system, #non_executable_client_tool?, #record_system_accounting, #record_tool_accounting, #registry_tool_injection_requested?, #request_tool_names, #request_tool_source, #resolve_registry_tool_source
Methods included from ContextWindow
#compact_to_fit, #empty_assistant_message?, #enforce_context_window, #estimate_message_tokens, #estimate_tool_token_budget, #last_user_message_index, #native_dispatch_messages, #reduce_messages_for_dispatch, #resolved_context_window, #strip_leading_thinking_block, #strip_thinking_from_history, #strip_thinking_pure, #tool_result_message?, #trim_oversized_tool_results, #trim_oversized_tool_results_pure
Methods included from Escalation
#attempts_exhausted_rejection, #client_stream_error?, #execute_provider_request, #execute_provider_request_native, #execute_provider_request_stream, #execute_provider_request_stream_native, #internal_error?, #non_provider_failure?, #populate_ssot_v3_resolved_state, #record_provider_response, #run_provider_call_engine, #run_provider_call_ssot_v3_stream, #run_provider_call_ssot_v3_stream_loop, #sse_translation_error?, #ssot_v3_execute_attempt, #ssot_v3_stream_failover_outcome, #ssot_v3_stream_handle_failure, #ssot_v3_stream_lane_hash, #ssot_v3_stream_preflight, #ssot_v3_stream_session_and_attempt, #step_provider_call, #step_provider_call_stream
Methods included from Routing
#build_ssot_v3_routing_requirements, #local_provider?, #merge_response_offering_metadata, #normalize_offering_metadata, #required_output_tokens_for_request, #step_request_normalization, #step_routing, #step_tier_assignment, #use_native_dispatch?
Constructor Details
#initialize(request) ⇒ Executor
Returns a new instance of Executor.
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 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 |
# File 'lib/legion/llm/inference/executor.rb', line 106 def initialize(request) @request = request @profile = Profile.derive(request.caller) @timeline = Timeline.new @tracing = nil @enrichments = {} @audit = {} @warnings = [] @timestamps = { received: Time.now } @raw_response = nil @exchange_id = nil @discovered_tools = [] @triggered_tools = [] @resolved_provider = nil @resolved_instance = nil @resolved_model = nil @resolved_tier = nil @resolved_offering_id = nil @resolved_offering_metadata = {} @confidence_score = nil @escalation_history = [] @route_attempts = [] @current_escalation_context = nil @routing_requirements = nil @current_attempt_context = nil @pre_provider_steps_done = false @stream_session = nil @preflight_lease = nil @last_ssot_dispatch_outcome = nil @proactive_tier_assignment = nil @tool_event_handler = nil @sticky_turn_snapshot = nil @pending_tool_history = Concurrent::Array.new @pending_tool_history_mutex = Mutex.new @deferred_tool_audits = [] @injected_tool_map = {} @native_tool_source_map = {} @freshly_triggered_keys = [] @applied_signals = { advisory_id: nil, behavioral_synapse_ids: [], trace_ids: [], advisory_types: [], envelope_keys: [], prediction_id: nil, response_stats: {} } @context_accounting = ContextAccounting.empty end |
Instance Attribute Details
#audit ⇒ Object (readonly)
Returns the value of attribute audit.
34 35 36 |
# File 'lib/legion/llm/inference/executor.rb', line 34 def audit @audit end |
#confidence_score ⇒ Object (readonly)
Returns the value of attribute confidence_score.
34 35 36 |
# File 'lib/legion/llm/inference/executor.rb', line 34 def confidence_score @confidence_score end |
#discovered_tools ⇒ Object (readonly)
Returns the value of attribute discovered_tools.
34 35 36 |
# File 'lib/legion/llm/inference/executor.rb', line 34 def discovered_tools @discovered_tools end |
#enrichments ⇒ Object (readonly)
Returns the value of attribute enrichments.
34 35 36 |
# File 'lib/legion/llm/inference/executor.rb', line 34 def enrichments @enrichments end |
#profile ⇒ Object (readonly)
Returns the value of attribute profile.
34 35 36 |
# File 'lib/legion/llm/inference/executor.rb', line 34 def profile @profile end |
#request ⇒ Object (readonly)
Returns the value of attribute request.
34 35 36 |
# File 'lib/legion/llm/inference/executor.rb', line 34 def request @request end |
#timeline ⇒ Object (readonly)
Returns the value of attribute timeline.
34 35 36 |
# File 'lib/legion/llm/inference/executor.rb', line 34 def timeline @timeline end |
#tool_event_handler ⇒ Object
Returns the value of attribute tool_event_handler.
36 37 38 |
# File 'lib/legion/llm/inference/executor.rb', line 36 def tool_event_handler @tool_event_handler end |
#tracing ⇒ Object (readonly)
Returns the value of attribute tracing.
34 35 36 |
# File 'lib/legion/llm/inference/executor.rb', line 34 def tracing @tracing end |
#warnings ⇒ Object (readonly)
Returns the value of attribute warnings.
34 35 36 |
# File 'lib/legion/llm/inference/executor.rb', line 34 def warnings @warnings end |
Instance Method Details
#call ⇒ Object
156 157 158 159 160 161 162 163 164 165 |
# File 'lib/legion/llm/inference/executor.rb', line 156 def call set_log_context Thread.current[:legion_llm_in_pipeline] = true log.debug "[llm][executor] action=call request_id=#{@request.id} profile=#{@profile}" execute_steps build_response ensure Thread.current[:legion_llm_in_pipeline] = nil clear_log_context end |
#call_responses(stream: false, stream_observer: nil, &block) ⇒ Object
N×N: Delegates to the canonical execution path. The API namespace translator has already parsed the Responses API format into canonical form. The provider adapter decides how to wire canonical requests internally — the executor is format-agnostic.
228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 |
# File 'lib/legion/llm/inference/executor.rb', line 228 def call_responses(stream: false, stream_observer: nil, **, &block) @stream_observer = stream_observer set_log_context Thread.current[:legion_llm_in_pipeline] = true log.debug "[llm][executor] action=call_responses->canonical request_id=#{@request.id} profile=#{@profile} stream=#{stream}" execute_pre_provider_steps if stream && block step_provider_call_stream(&block) else step_provider_call end execute_post_provider_steps build_response ensure Thread.current[:legion_llm_in_pipeline] = nil @stream_observer = nil clear_log_context end |
#call_stream(stream_observer: nil, &block) ⇒ Object
192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 |
# File 'lib/legion/llm/inference/executor.rb', line 192 def call_stream(stream_observer: nil, &block) @stream_observer = stream_observer return call unless block set_log_context Thread.current[:legion_llm_in_pipeline] = true log.debug "[llm][executor] action=call_stream request_id=#{@request.id} profile=#{@profile}" execute_pre_provider_steps unless @pre_provider_steps_done step_provider_call_stream(&block) execute_post_provider_steps build_response ensure release_preflight_lease @stream_observer = nil Thread.current[:legion_llm_in_pipeline] = nil clear_log_context end |
#context_accounting ⇒ Object
38 39 40 |
# File 'lib/legion/llm/inference/executor.rb', line 38 def context_accounting @context_accounting ||= ContextAccounting.empty end |
#release_preflight_lease ⇒ Object
Release the streaming preflight DispatchLease (SSOT v3 §15.2). Always called from call_stream's ensure so every streaming exit — success, provider failure, cancellation, client disconnect, thread interruption — retires the exact lease acquired during preflight exactly once.
214 215 216 217 218 219 220 221 222 |
# File 'lib/legion/llm/inference/executor.rb', line 214 def release_preflight_lease return unless @preflight_lease @preflight_lease.release unless @preflight_lease.released? rescue StandardError => e handle_exception(e, level: :warn, handled: true, operation: 'llm.pipeline.release_preflight_lease') ensure @preflight_lease = nil end |
#stream_preflight! ⇒ Object
SSOT v3 §19 streaming preflight. Runs the pre-provider steps and, when the SSOT inventory path is active, selects AND acquires the exact lane (RoutingSession#next_attempt! + a DispatchLease) BEFORE the route opens the SSE event-stream. A routing rejection therefore surfaces as Errors::RoutingRejected here — mapped to a proper HTTP status by the route's RoutingErrorMapper rescue — instead of an SSE server_error emitted after the response headers are already committed.
Returns the selected lane Hash (for StreamAssembler#initial_lane) when SSOT selection ran, or nil when the SSOT path is not active (fallback preserves the prior behavior: the route opens SSE and call_stream selects inline). The selected AttemptContext, RoutingSession, and lease are retained on the executor and reused by the subsequent call_stream.
180 181 182 183 184 185 186 187 188 189 190 |
# File 'lib/legion/llm/inference/executor.rb', line 180 def stream_preflight! set_log_context Thread.current[:legion_llm_in_pipeline] = true log.debug "[llm][executor] action=stream_preflight request_id=#{@request.id} profile=#{@profile}" execute_pre_provider_steps @pre_provider_steps_done = true ssot_v3_stream_preflight ensure Thread.current[:legion_llm_in_pipeline] = nil clear_log_context end |