1use std::collections::HashMap;
2use std::path::PathBuf;
3use std::sync::Arc;
4
5use anyhow::Context;
6use futures::StreamExt;
7use futures::future::{AbortHandle, Abortable, BoxFuture, try_join_all};
8use roder_api::catalog::{
9 EDIT_TOOL_EDIT, EDIT_TOOL_PATCH, PROVIDER_GEMINI, REASONING_NONE, built_in_model_profile,
10 built_in_model_profile_for_provider, lookup_model,
11};
12use roder_api::context::PolicyGate;
13use roder_api::events::*;
14use roder_api::extension::ExtensionRegistry;
15use roder_api::inference::{
16 AgentInferenceRequest, HostedWebSearchConfig, HostedWebSearchMode, InferenceEngine,
17 InferenceEvent, InferenceTurnContext, InstructionBundle, ModelHarnessProfile,
18 ModelSchemaPolicy, ModelSelection, OutputConfig, ReasoningConfig, RuntimeHints, RuntimeProfile,
19 TokenUsage, ToolCallCompleted, ToolSearchConfig, ToolSearchConfigOverlay,
20 finish_reason_from_stop_reason,
21};
22use roder_api::inference_routing::{InferenceRoutingOutcome, ModelSelectionMode};
23use roder_api::policy_mode::{PolicyDecision, PolicyMode};
24use roder_api::reliability::{
25 ReliabilityContext, ReliabilityDetails, ReliabilityErrorClass, ReliabilityLimitRecorded,
26 ReliabilityRequestPolicy, ReliabilityRetryDecision, ReliabilityRetryRecorded,
27 provider_retry_delay_ms,
28};
29use roder_api::remote_runner::{
30 RemoteRunnerSession, RemoteWorkspace, RunnerDestination, ThreadRunnerBinding,
31};
32use roder_api::subagents::SubagentDefinition;
33use roder_api::teams::TeamMemberStatus;
34use roder_api::thread::{
35 ThreadItemEvent, ThreadItemEventKind, ThreadMetadata, ThreadSnapshot, ThreadStore,
36 ThreadUsageMetadata, is_synthetic_event_thread_id, validate_thread_workspace,
37};
38use roder_api::tools::{ToolCall, ToolChoice, ToolExecutionContext, ToolRegistry, ToolResult};
39use roder_api::transcript::{
40 AssistantMessage, ErrorRecord, InputImage, ReasoningSummary, ToolCallRecord, ToolResultRecord,
41 TranscriptItem, UserMessage,
42};
43use roder_sandbox::ScopedFilesystem;
44use roder_sandbox::process::LocalProcessRunner;
45use roder_skills::{SkillRegistry, SkillRegistryOptions};
46use time::{Duration, OffsetDateTime};
47use tokio::sync::{Mutex, RwLock, oneshot};
48
49
50use crate::artifacts::{
51 ContextArtifactStore as FilesystemContextArtifactStore, default_context_artifact_dir,
52};
53use crate::bus::EventBus;
54use crate::dynamic_workflows::{
55 DynamicWorkflowEffortProfile, RuntimeDynamicWorkflowConfig, WorkflowTriggerDecision,
56 classify_workflow_trigger, ultracode_reasoning_level_for_model,
57};
58use crate::fake_provider::FakeInferenceEngine;
59use crate::goals::RuntimeGoalController;
60use crate::inference_routing::{
61 InferenceRoutingRequest, RuntimeInferenceRouterConfig, collect_inference_routing_candidates,
62 route_inference_selection, transcript_failure_count_since,
63};
64use crate::instructions::{
65 apply_model_instruction_overlay, apply_plan_mode, apply_runtime_profile,
66 apply_task_ledger_required, apply_thread_developer_instructions, apply_turn_developer_context,
67};
68use crate::policy_gate::DefaultPolicyGate;
69use crate::reliability::{
70 ReliabilityLimitHit, RuntimeReliabilityConfig, TurnReliabilityState,
71 provider_stream_retry_cause,
72};
73pub use crate::speed_policy::RuntimeSpeedPolicyConfig;
74use crate::speed_policy::{SpeedPolicyState, reasoning_from_decision};
75use crate::subagent_traces::RuntimeSubagentTraceSink;
76use crate::teams::{TeamManager, TeamMemberStartRequest, TeamStartRequest, TeamState};
77use crate::thread_item_cache::{ThreadItemCache, ThreadItemCacheEntry};
78use crate::verification_gate::VerificationGateState;
79
80const MAX_TOOL_ROUNDS_PER_TURN: usize = 1024;
81const FINAL_ANSWER_PHASE: &str = "final_answer";
82pub(crate) const TASK_LEDGER_TOOL_NAME: &str = "task_ledger.update";
83const TASK_LEDGER_COMPLETION_REMINDER_LIMIT: u8 = 2;
84const TASK_LEDGER_SCOREABLE_CHECKPOINT_SECONDS: u64 = 180;
85const TASK_LEDGER_SCOREABLE_CHECKPOINT_LIMIT: u8 = 1;
86pub(crate) const MIN_CHILD_DEADLINE_SECONDS: u64 = 2;
87const MODEL_PROFILE_TRACE_KIND: &str = "model_profile_segment";
88const MODEL_SWITCH_SUMMARY_PREFIX: &str = "Model switch summary:";
89
90#[derive(Clone, Copy, Debug, Eq, PartialEq)]
91enum InferenceTimeoutAction {
92 ScoreableCheckpoint,
93 Finalization,
94}
95
96#[derive(Debug, Clone)]
97pub struct RuntimeConfig {
98 pub default_provider: String,
99 pub default_model: String,
100 pub reasoning: Option<String>,
101 pub auto_compact_token_limit: Option<u32>,
102 pub file_backed_dynamic_context: bool,
103 pub hosted_web_search: HostedWebSearchConfig,
104 pub tool_search: ToolSearchConfig,
105 pub provider_tool_search: HashMap<String, ToolSearchConfigOverlay>,
106 pub model_tool_search: HashMap<String, ToolSearchConfigOverlay>,
107 pub model_edit_tools: HashMap<String, String>,
108 pub model_parallel_tool_calls: HashMap<String, bool>,
109 pub model_profiles: HashMap<String, ModelHarnessProfile>,
110 pub tool_allowlist: Vec<String>,
111 pub external_tool_timeout_seconds: u64,
113 pub command_shell: String,
114 pub workspace: Option<String>,
115 pub policy_mode: PolicyMode,
116 pub runtime_profile: RuntimeProfile,
117 pub inference_router: RuntimeInferenceRouterConfig,
118 pub speed_policy: RuntimeSpeedPolicyConfig,
119 pub dynamic_workflows: RuntimeDynamicWorkflowConfig,
120 pub reliability: RuntimeReliabilityConfig,
121 pub turn_deadline_seconds: Option<u64>,
122 pub remote_runner_destination: Option<RunnerDestination>,
123 pub team_data_dir: Option<PathBuf>,
124 pub roadmap_data_dir: Option<PathBuf>,
125 pub media_generation: crate::media_generation::RuntimeMediaGenerationConfig,
126}
127
128impl Default for RuntimeConfig {
129 fn default() -> Self {
130 Self {
131 default_provider: roder_api::catalog::PROVIDER_MOCK.to_string(),
132 default_model: "mock".to_string(),
133 reasoning: None,
134 auto_compact_token_limit: None,
135 file_backed_dynamic_context: true,
136 hosted_web_search: HostedWebSearchConfig::cached(),
137 tool_search: ToolSearchConfig::default(),
138 provider_tool_search: HashMap::new(),
139 model_tool_search: HashMap::new(),
140 model_edit_tools: HashMap::new(),
141 model_parallel_tool_calls: HashMap::new(),
142 model_profiles: HashMap::new(),
143 tool_allowlist: Vec::new(),
144 external_tool_timeout_seconds: DEFAULT_EXTERNAL_TOOL_TIMEOUT_SECONDS,
145 command_shell: roder_api::command_shell::default_command_shell(),
146 workspace: None,
147 policy_mode: PolicyMode::Default,
148 runtime_profile: RuntimeProfile::Interactive,
149 inference_router: RuntimeInferenceRouterConfig::default(),
150 speed_policy: RuntimeSpeedPolicyConfig::default(),
151 dynamic_workflows: RuntimeDynamicWorkflowConfig::default(),
152 reliability: RuntimeReliabilityConfig::default(),
153 turn_deadline_seconds: None,
154 remote_runner_destination: None,
155 team_data_dir: None,
156 roadmap_data_dir: None,
157 media_generation: crate::media_generation::RuntimeMediaGenerationConfig::default(),
158 }
159 }
160}
161
162#[derive(Debug, Clone)]
163pub struct StartTurnRequest {
164 pub thread_id: ThreadId,
165 pub message: String,
166 pub images: Vec<InputImage>,
167 pub provider_override: Option<String>,
168 pub model_override: Option<String>,
169 pub reasoning_override: Option<String>,
170 pub workspace: String,
171 pub instructions: InstructionBundle,
172 pub developer_context: Option<String>,
178 pub task_ledger_required: bool,
179}
180
181#[derive(Debug, Clone)]
182pub struct CreateThreadRequest {
183 pub title: Option<String>,
184 pub workspace: String,
185 pub workspace_id: Option<String>,
186 pub root_id: Option<String>,
187 pub provider: Option<String>,
188 pub model: Option<String>,
189 pub selection_mode: Option<ModelSelectionMode>,
190 pub tool_allowlist: Vec<String>,
192 pub developer_instructions: Option<String>,
194 pub external_tools: Vec<roder_api::tools::ToolSpec>,
196 pub runner: Option<ThreadRunnerSelection>,
198}
199
200#[derive(Debug, Clone)]
206pub struct ThreadRunnerSelection {
207 pub provider_id: String,
208 pub config: serde_json::Value,
209 pub workspace: String,
211 pub read_roots: Vec<String>,
217}
218
219#[derive(Debug, Clone, PartialEq, Eq)]
220pub struct PendingPlanExit {
221 pub thread_id: ThreadId,
222 pub turn_id: TurnId,
223 pub request_id: String,
224 pub target_mode: PolicyMode,
225 pub plan_summary: Option<String>,
226 pub next_steps: Vec<String>,
227 pub requested_at: OffsetDateTime,
228 pub expires_at: Option<OffsetDateTime>,
229}
230
231pub(crate) struct PendingToolApproval {
232 pub(crate) thread_id: ThreadId,
233 pub(crate) turn_id: TurnId,
234 pub(crate) tool_id: String,
235 pub(crate) tool_name: String,
236 pub(crate) call: roder_api::tools::ToolCall,
237 pub(crate) tx: oneshot::Sender<bool>,
238}
239
240pub(crate) struct PendingUserInput {
241 pub(crate) thread_id: ThreadId,
242 pub(crate) turn_id: TurnId,
243 pub(crate) tx: oneshot::Sender<serde_json::Value>,
244}
245
246#[derive(Debug, Clone, PartialEq, Eq)]
248pub struct ExternalToolResolution {
249 pub output: String,
250 pub is_error: bool,
251}
252
253pub(crate) struct PendingExternalToolCall {
254 pub(crate) thread_id: ThreadId,
255 pub(crate) turn_id: TurnId,
256 pub(crate) tool_id: String,
257 pub(crate) tool_name: String,
258 pub(crate) tx: oneshot::Sender<ExternalToolResolution>,
259}
260
261#[derive(Clone)]
262struct ActiveTurnHandle {
263 thread_id: ThreadId,
264 abort: AbortHandle,
265 steers: Arc<Mutex<Vec<UserMessage>>>,
266}
267
268#[derive(Debug, Clone, Default, PartialEq, Eq)]
269pub struct ThreadActivity {
270 pub active_turn_id: Option<TurnId>,
271 pub active_flags: Vec<String>,
272}
273
274#[derive(Debug, Clone, Default)]
276pub(crate) struct ThreadTurnOverrides {
277 pub(crate) tool_allowlist: Vec<String>,
278 pub(crate) developer_instructions: Option<String>,
279 pub(crate) external_tools: Vec<roder_api::tools::ToolSpec>,
280}
281
282#[derive(Debug, Clone, Copy, PartialEq, Eq)]
283pub(crate) enum TurnRunOutcome {
284 Completed,
285 Stopped,
286}
287
288impl PendingPlanExit {
289 pub fn new(
290 thread_id: ThreadId,
291 turn_id: TurnId,
292 request_id: String,
293 target_mode: PolicyMode,
294 plan_summary: Option<String>,
295 next_steps: Vec<String>,
296 ) -> Self {
297 let requested_at = OffsetDateTime::now_utc();
298 Self {
299 thread_id,
300 turn_id,
301 request_id,
302 target_mode,
303 plan_summary,
304 next_steps,
305 requested_at,
306 expires_at: Some(requested_at + default_plan_exit_timeout()),
307 }
308 }
309
310 pub fn is_expired(&self, now: OffsetDateTime) -> bool {
311 self.expires_at.is_some_and(|expires_at| now >= expires_at)
312 }
313}
314
315pub fn default_plan_exit_timeout() -> Duration {
316 Duration::minutes(10)
317}
318
319pub const DEFAULT_EXTERNAL_TOOL_TIMEOUT_SECONDS: u64 = 300;
320
321pub struct Runtime {
322 pub bus: EventBus,
323 pub registry: ExtensionRegistry,
324 config: RwLock<RuntimeConfig>,
325 pending_plan_exit: RwLock<Option<PendingPlanExit>>,
326 pub(crate) pending_tool_approvals: Mutex<HashMap<String, PendingToolApproval>>,
327 pub(crate) pending_user_inputs: Mutex<HashMap<String, PendingUserInput>>,
328 pub(crate) pending_external_tool_calls: Mutex<HashMap<String, PendingExternalToolCall>>,
329 active_turns: RwLock<HashMap<TurnId, ActiveTurnHandle>>,
330 workspace: PathBuf,
331 teams: TeamManager,
332 pub(crate) roadmaps: Mutex<roder_roadmap::RoadmapRuntime>,
333 pub(crate) goals: Arc<RuntimeGoalController>,
334 context_artifacts: roder_api::artifacts::ContextArtifactStore,
335 pub(crate) thread_store: Option<Arc<dyn ThreadStore>>,
336 thread_item_cache: Mutex<ThreadItemCache>,
337 pub(crate) tool_registry: ToolRegistry,
338 media_generation: Arc<crate::media_generation::MediaGenerationService>,
339 pub(crate) skills: RwLock<SkillRegistry>,
340 event_sink_dispatcher: tokio::sync::OnceCell<crate::event_sink_dispatch::EventSinkDispatcher>,
343 pub(crate) compaction_hysteresis: std::sync::Mutex<HashMap<ThreadId, u32>>,
344}
345
346impl Runtime {
347 pub fn new(registry: ExtensionRegistry, config: RuntimeConfig) -> anyhow::Result<Self> {
348 if registry.inference_engines.is_empty() {
349 anyhow::bail!("at least one inference engine must be registered");
350 }
351 validate_runtime_config_reasoning(&config)?;
352 validate_runtime_inference_router_config(®istry, &config)?;
353
354 let bus = EventBus::new(1024);
355 let thread_store = registry
356 .thread_stores
357 .first()
358 .map(|factory| factory.create());
359 let mut tool_registry = ToolRegistry::default();
360 for contributor in ®istry.tools {
361 contributor
362 .contribute(&mut tool_registry)
363 .with_context(|| format!("tool contributor {} failed", contributor.id()))?;
364 }
365 crate::agent_control_tools::contribute_agent_control_tools(&mut tool_registry)?;
366
367 let media_generation = Arc::new(crate::media_generation::MediaGenerationService::new(
368 registry.media_generator_providers.clone(),
369 config.media_generation.clone(),
370 ));
371 tool_registry.replace(Arc::new(
372 crate::media_generation::MediaGenerateImageTool::new(media_generation.clone()),
373 ));
374
375 let team_data_dir = config.team_data_dir.clone();
376 let workspace = config
377 .workspace
378 .clone()
379 .map(PathBuf::from)
380 .unwrap_or(std::env::current_dir()?);
381 let roadmap_data_dir = config
382 .roadmap_data_dir
383 .clone()
384 .unwrap_or_else(|| workspace.join(".roder"));
385 let context_artifacts = thread_store
386 .as_ref()
387 .and_then(|store| store.context_artifact_store())
388 .or_else(|| {
389 thread_store
390 .as_ref()
391 .and_then(|store| store.local_thread_root())
392 .map(FilesystemContextArtifactStore::shared_thread_scoped)
393 })
394 .unwrap_or_else(|| {
395 FilesystemContextArtifactStore::shared_legacy(default_context_artifact_dir())
396 });
397 let goals = Arc::new(RuntimeGoalController::new(
398 bus.clone(),
399 thread_store.clone(),
400 ));
401 let runtime = Self {
402 bus,
403 registry,
404 config: RwLock::new(config),
405 pending_plan_exit: RwLock::new(None),
406 pending_tool_approvals: Mutex::new(HashMap::new()),
407 pending_user_inputs: Mutex::new(HashMap::new()),
408 pending_external_tool_calls: Mutex::new(HashMap::new()),
409 active_turns: RwLock::new(HashMap::new()),
410 workspace: workspace.clone(),
411 teams: TeamManager::new(
412 team_data_dir.unwrap_or_else(crate::teams::default_team_data_dir),
413 ),
414 roadmaps: Mutex::new(roder_roadmap::RoadmapRuntime::new(
415 workspace,
416 roadmap_data_dir,
417 )),
418 goals,
419 context_artifacts,
420 thread_store,
421 thread_item_cache: Mutex::new(ThreadItemCache::default()),
422 tool_registry,
423 media_generation,
424 skills: RwLock::new(SkillRegistry::load(SkillRegistryOptions::new(
425 PathBuf::new(),
426 ))),
427 event_sink_dispatcher: tokio::sync::OnceCell::new(),
428 compaction_hysteresis: crate::compaction_runtime::compaction_hysteresis_state(),
429 };
430 runtime.bus.emit(RoderEvent::RuntimeStarted(RuntimeStarted {
431 timestamp: OffsetDateTime::now_utc(),
432 }));
433 for manifest in &runtime.registry.manifests {
434 runtime
435 .bus
436 .emit(RoderEvent::ExtensionRegistered(ExtensionRegistered {
437 extension_id: manifest.id.clone(),
438 timestamp: OffsetDateTime::now_utc(),
439 }));
440 }
441 Ok(runtime)
442 }
443
444 pub fn from_engine(engine: Arc<dyn InferenceEngine>) -> anyhow::Result<Self> {
445 let mut builder = roder_api::extension::ExtensionRegistryBuilder::new();
446 builder.inference_engine(engine);
447 Self::new(builder.build()?, RuntimeConfig::default())
448 }
449
450 pub fn fake() -> anyhow::Result<Self> {
451 Self::from_engine(Arc::new(FakeInferenceEngine))
452 }
453
454 pub fn subscribe_events(&self) -> tokio::sync::broadcast::Receiver<EventEnvelope> {
455 self.bus.subscribe()
456 }
457
458 pub fn registry(&self) -> &ExtensionRegistry {
459 &self.registry
460 }
461
462 pub fn media_generation(&self) -> Arc<crate::media_generation::MediaGenerationService> {
463 self.media_generation.clone()
464 }
465
466 pub fn context_artifacts(&self) -> roder_api::artifacts::ContextArtifactStore {
467 self.context_artifacts.clone()
468 }
469
470 pub async fn execute_workflow_tool(
471 &self,
472 thread_id: ThreadId,
473 tool_name: &str,
474 arguments: serde_json::Value,
475 ) -> anyhow::Result<ToolResult> {
476 let Some(executor) = self.tool_registry.get(tool_name) else {
477 anyhow::bail!("tool not found: {tool_name}");
478 };
479 let tool_call = ToolCall {
480 id: format!("slash-{tool_name}"),
481 name: tool_name.to_string(),
482 raw_arguments: serde_json::to_string(&arguments)?,
483 arguments,
484 thread_id: thread_id.clone(),
485 turn_id: "slash-command".to_string(),
486 };
487 let runtime_config = self.status().await;
488 let ctx = self.tool_execution_context(
489 thread_id,
490 "slash-command".to_string(),
491 runtime_config.policy_mode,
492 runtime_config.workspace.as_deref(),
493 Some(&runtime_config.command_shell),
494 );
495 executor.execute(ctx, tool_call).await
496 }
497
498 pub(crate) fn tool_execution_context(
499 &self,
500 thread_id: ThreadId,
501 turn_id: TurnId,
502 mode: PolicyMode,
503 workspace: Option<&str>,
504 command_shell: Option<&str>,
505 ) -> ToolExecutionContext {
506 let mut ctx = ToolExecutionContext::new(thread_id, turn_id, mode)
507 .with_command_shell(command_shell.unwrap_or_default())
508 .with_process_runner(Arc::new(LocalProcessRunner))
509 .with_context_artifacts(self.context_artifacts.backend())
510 .with_goal_controller(self.goals.clone())
511 .with_subagent_trace_sink(Arc::new(RuntimeSubagentTraceSink::new(
512 self.bus.clone(),
513 self.thread_store.clone(),
514 )));
515 if let Some(workspace) = workspace {
516 ctx = ctx.with_workspace_handle(Arc::new(ScopedFilesystem::new(workspace)));
517 }
518 ctx
519 }
520
521 pub async fn status(&self) -> RuntimeConfig {
522 self.config.read().await.clone()
523 }
524
525 pub async fn set_skills(&self, skills: SkillRegistry) {
526 *self.skills.write().await = skills;
527 }
528
529 pub async fn skills_snapshot(&self) -> SkillRegistry {
530 self.skills.read().await.clone()
531 }
532
533 pub fn workspace(&self) -> PathBuf {
534 self.workspace.clone()
535 }
536
537 pub async fn set_remote_runner_destination(&self, destination: Option<RunnerDestination>) {
538 let lifecycle = destination.as_ref().map(|destination| RunnerLifecycle {
539 destination_id: destination.id.clone(),
540 provider_id: destination.provider_id.clone(),
541 state: "configured".to_string(),
542 session_id: None,
543 timestamp: OffsetDateTime::now_utc(),
544 });
545 self.config.write().await.remote_runner_destination = destination;
546 if let Some(lifecycle) = lifecycle {
547 self.emit(RoderEvent::RunnerLifecycle(lifecycle)).await;
548 } else {
549 self.emit(RoderEvent::RunnerLifecycle(RunnerLifecycle {
550 destination_id: "local".to_string(),
551 provider_id: "local".to_string(),
552 state: "local_fallback".to_string(),
553 session_id: None,
554 timestamp: OffsetDateTime::now_utc(),
555 }))
556 .await;
557 }
558 }
559
560 pub async fn set_file_backed_dynamic_context(&self, enabled: bool) -> RuntimeConfig {
561 let mut cfg = self.config.write().await;
562 cfg.file_backed_dynamic_context = enabled;
563 cfg.clone()
564 }
565
566 pub async fn set_command_shell(&self, shell: String) -> RuntimeConfig {
567 let mut cfg = self.config.write().await;
568 cfg.command_shell = shell;
569 cfg.clone()
570 }
571
572 pub async fn pending_plan_exit(&self) -> Option<PendingPlanExit> {
573 let mut pending = self.pending_plan_exit.write().await;
574 let current = pending.clone()?;
575 if !current.is_expired(OffsetDateTime::now_utc()) {
576 return Some(current);
577 }
578 *pending = None;
579 drop(pending);
580 self.emit_plan_exit_resolved(¤t, false, self.status().await.policy_mode)
581 .await;
582 None
583 }
584
585 pub async fn set_policy_mode(
586 &self,
587 mode: PolicyMode,
588 reason: Option<String>,
589 ) -> anyhow::Result<RuntimeConfig> {
590 let mut cfg = self.config.write().await;
591 let previous_mode = cfg.policy_mode;
592 cfg.policy_mode = mode;
593 let next = cfg.clone();
594 drop(cfg);
595 self.emit(RoderEvent::PolicyModeChanged(PolicyModeChanged {
596 thread_id: "runtime".to_string(),
597 turn_id: None,
598 previous_mode,
599 new_mode: mode,
600 reason,
601 timestamp: OffsetDateTime::now_utc(),
602 }))
603 .await;
604 self.auto_resolve_pending_tool_approvals_for_mode(mode)
605 .await;
606 Ok(next)
607 }
608
609 pub async fn set_hosted_web_search(
610 &self,
611 mode: HostedWebSearchMode,
612 ) -> anyhow::Result<RuntimeConfig> {
613 let mut cfg = self.config.write().await;
614 cfg.hosted_web_search = HostedWebSearchConfig { mode };
615 Ok(cfg.clone())
616 }
617
618 async fn auto_resolve_pending_tool_approvals_for_mode(&self, mode: PolicyMode) {
619 let gate = DefaultPolicyGate::new();
620 let mut pending = self.pending_tool_approvals.lock().await;
621 let approval_ids = pending
622 .iter()
623 .filter_map(|(approval_id, approval)| {
624 let ctx = ToolExecutionContext::new(
625 approval.thread_id.clone(),
626 approval.turn_id.clone(),
627 mode,
628 );
629 matches!(
630 gate.decide(&approval.call, mode, &ctx),
631 PolicyDecision::AutoApproved { .. }
632 )
633 .then_some(approval_id.clone())
634 })
635 .collect::<Vec<_>>();
636 let approvals = approval_ids
637 .into_iter()
638 .filter_map(|approval_id| {
639 pending
640 .remove(&approval_id)
641 .map(|approval| (approval_id, approval))
642 })
643 .collect::<Vec<_>>();
644 drop(pending);
645
646 for (approval_id, approval) in approvals {
647 let ctx = ToolExecutionContext::new(
648 approval.thread_id.clone(),
649 approval.turn_id.clone(),
650 mode,
651 );
652 let decision = gate.decide(&approval.call, mode, &ctx);
653 self.emit(RoderEvent::PolicyDecisionRecorded(PolicyDecisionRecorded {
654 thread_id: approval.thread_id.clone(),
655 turn_id: approval.turn_id.clone(),
656 tool_id: approval.tool_id.clone(),
657 tool_name: approval.tool_name.clone(),
658 mode,
659 decision,
660 timestamp: OffsetDateTime::now_utc(),
661 }))
662 .await;
663 if mode == PolicyMode::Bypass {
664 self.emit(RoderEvent::PolicyBypassActive(PolicyBypassActive {
665 thread_id: approval.thread_id.clone(),
666 turn_id: approval.turn_id.clone(),
667 tool_id: approval.tool_id.clone(),
668 tool_name: approval.tool_name.clone(),
669 timestamp: OffsetDateTime::now_utc(),
670 }))
671 .await;
672 }
673 self.emit(RoderEvent::ApprovalResolved(ApprovalResolved {
674 thread_id: approval.thread_id,
675 turn_id: approval.turn_id,
676 approval_id,
677 tool_id: approval.tool_id,
678 tool_name: approval.tool_name,
679 approved: true,
680 timestamp: OffsetDateTime::now_utc(),
681 }))
682 .await;
683 let _ = approval.tx.send(true);
684 }
685 }
686
687 pub async fn record_pending_plan_exit(&self, pending: PendingPlanExit) {
688 *self.pending_plan_exit.write().await = Some(pending.clone());
689 self.emit(RoderEvent::PolicyExitPlanRequested(
690 PolicyExitPlanRequested {
691 thread_id: pending.thread_id,
692 turn_id: pending.turn_id,
693 request_id: pending.request_id,
694 target_mode: pending.target_mode,
695 plan_summary: pending.plan_summary,
696 next_steps: pending.next_steps,
697 timestamp: OffsetDateTime::now_utc(),
698 },
699 ))
700 .await;
701 }
702
703 pub async fn resolve_pending_plan_exit(
704 &self,
705 request_id: &str,
706 approved: bool,
707 ) -> anyhow::Result<Option<PendingPlanExit>> {
708 let mut pending = self.pending_plan_exit.write().await;
709 let Some(current) = pending.clone() else {
710 return Ok(None);
711 };
712 if current.request_id != request_id {
713 anyhow::bail!("pending plan exit request {request_id:?} was not found");
714 }
715 *pending = None;
716 drop(pending);
717
718 let approved = approved && !current.is_expired(OffsetDateTime::now_utc());
719 let resolved_mode = if approved {
720 let mut cfg = self.config.write().await;
721 let previous_mode = cfg.policy_mode;
722 cfg.policy_mode = current.target_mode;
723 drop(cfg);
724 self.emit(RoderEvent::PolicyModeChanged(PolicyModeChanged {
725 thread_id: current.thread_id.clone(),
726 turn_id: Some(current.turn_id.clone()),
727 previous_mode,
728 new_mode: current.target_mode,
729 reason: Some("approved plan exit".to_string()),
730 timestamp: OffsetDateTime::now_utc(),
731 }))
732 .await;
733 self.auto_resolve_pending_tool_approvals_for_mode(current.target_mode)
734 .await;
735 current.target_mode
736 } else {
737 self.status().await.policy_mode
738 };
739 self.emit_plan_exit_resolved(¤t, approved, resolved_mode)
740 .await;
741 Ok(Some(current))
742 }
743
744 pub async fn resolve_tool_approval(
745 &self,
746 approval_id: &str,
747 approved: bool,
748 ) -> anyhow::Result<bool> {
749 let pending = self.pending_tool_approvals.lock().await.remove(approval_id);
750 let Some(pending) = pending else {
751 return Ok(false);
752 };
753 self.emit(RoderEvent::ApprovalResolved(ApprovalResolved {
754 thread_id: pending.thread_id,
755 turn_id: pending.turn_id,
756 approval_id: approval_id.to_string(),
757 tool_id: pending.tool_id,
758 tool_name: pending.tool_name,
759 approved,
760 timestamp: OffsetDateTime::now_utc(),
761 }))
762 .await;
763 let _ = pending.tx.send(approved);
764 Ok(true)
765 }
766
767 pub async fn request_app_server_tool_approval(
768 &self,
769 call: ToolCall,
770 reason: Option<String>,
771 ) -> anyhow::Result<bool> {
772 let approval_id = call.id.clone();
773 let (tx, rx) = oneshot::channel();
774 self.pending_tool_approvals.lock().await.insert(
775 approval_id.clone(),
776 PendingToolApproval {
777 thread_id: call.thread_id.clone(),
778 turn_id: call.turn_id.clone(),
779 tool_id: call.id.clone(),
780 tool_name: call.name.clone(),
781 call: call.clone(),
782 tx,
783 },
784 );
785 self.emit(RoderEvent::ApprovalRequested(ApprovalRequested {
786 thread_id: call.thread_id.clone(),
787 turn_id: call.turn_id.clone(),
788 approval_id,
789 tool_id: call.id.clone(),
790 tool_name: call.name.clone(),
791 reason,
792 timestamp: OffsetDateTime::now_utc(),
793 }))
794 .await;
795 Ok(rx.await.unwrap_or(false))
796 }
797
798 pub async fn resolve_external_tool_call(
801 &self,
802 request_id: &str,
803 resolution: ExternalToolResolution,
804 ) -> anyhow::Result<bool> {
805 let pending = self
806 .pending_external_tool_calls
807 .lock()
808 .await
809 .remove(request_id);
810 let Some(pending) = pending else {
811 return Ok(false);
812 };
813 self.emit(RoderEvent::ExternalToolCallResolved(
814 ExternalToolCallResolved {
815 thread_id: pending.thread_id,
816 turn_id: pending.turn_id,
817 request_id: request_id.to_string(),
818 tool_id: pending.tool_id,
819 tool_name: pending.tool_name,
820 outcome: ExternalToolCallOutcome::Resolved,
821 is_error: resolution.is_error,
822 timestamp: OffsetDateTime::now_utc(),
823 },
824 ))
825 .await;
826 let _ = pending.tx.send(resolution);
827 Ok(true)
828 }
829
830 async fn cancel_pending_external_tool_calls_for_turn(&self, turn_id: &TurnId) {
833 let cancelled = {
834 let mut pending = self.pending_external_tool_calls.lock().await;
835 let request_ids = pending
836 .iter()
837 .filter(|(_, call)| &call.turn_id == turn_id)
838 .map(|(request_id, _)| request_id.clone())
839 .collect::<Vec<_>>();
840 request_ids
841 .into_iter()
842 .filter_map(|request_id| pending.remove(&request_id).map(|call| (request_id, call)))
843 .collect::<Vec<_>>()
844 };
845 for (request_id, call) in cancelled {
846 self.emit(RoderEvent::ExternalToolCallResolved(
847 ExternalToolCallResolved {
848 thread_id: call.thread_id,
849 turn_id: call.turn_id,
850 request_id,
851 tool_id: call.tool_id,
852 tool_name: call.tool_name,
853 outcome: ExternalToolCallOutcome::Cancelled,
854 is_error: true,
855 timestamp: OffsetDateTime::now_utc(),
856 },
857 ))
858 .await;
859 }
860 }
861
862 pub async fn resolve_user_input(
863 &self,
864 request_id: &str,
865 answers: serde_json::Value,
866 ) -> anyhow::Result<bool> {
867 let pending = self.pending_user_inputs.lock().await.remove(request_id);
868 let Some(pending) = pending else {
869 return Ok(false);
870 };
871 self.emit(RoderEvent::UserInputResolved(UserInputResolved {
872 thread_id: pending.thread_id,
873 turn_id: pending.turn_id,
874 request_id: request_id.to_string(),
875 answers: answers.clone(),
876 timestamp: OffsetDateTime::now_utc(),
877 }))
878 .await;
879 let _ = pending.tx.send(answers);
880 Ok(true)
881 }
882
883 async fn emit_plan_exit_resolved(
884 &self,
885 current: &PendingPlanExit,
886 approved: bool,
887 resolved_mode: PolicyMode,
888 ) {
889 self.emit(RoderEvent::PolicyExitPlanResolved(PolicyExitPlanResolved {
890 thread_id: current.thread_id.clone(),
891 turn_id: current.turn_id.clone(),
892 request_id: current.request_id.clone(),
893 approved,
894 target_mode: current.target_mode,
895 resolved_mode,
896 timestamp: OffsetDateTime::now_utc(),
897 }))
898 .await;
899 }
900
901 pub async fn select_provider(
902 &self,
903 provider: String,
904 model: Option<String>,
905 reasoning: Option<String>,
906 ) -> anyhow::Result<RuntimeConfig> {
907 let next = self
908 .preview_provider_selection(provider, model, reasoning)
909 .await?;
910 let mut cfg = self.config.write().await;
911 *cfg = next;
912 Ok(cfg.clone())
913 }
914
915 pub async fn preview_provider_selection(
916 &self,
917 provider: String,
918 model: Option<String>,
919 reasoning: Option<String>,
920 ) -> anyhow::Result<RuntimeConfig> {
921 self.engine_for(&provider)?;
922 let mut cfg = self.config.read().await.clone();
923 cfg.default_provider = provider;
924 if let Some(model) = model {
925 cfg.default_model = model;
926 }
927 if let Some(reasoning) = reasoning {
928 if reasoning == REASONING_NONE
929 && !model_supports_reasoning(&cfg.default_model, &reasoning)
930 {
931 return Ok(cfg.clone());
932 }
933 validate_reasoning_effort(&cfg.default_model, &reasoning)?;
934 cfg.reasoning = Some(reasoning);
935 }
936 Ok(cfg)
937 }
938
939 pub async fn effective_reasoning(&self) -> String {
940 let cfg = self.config.read().await;
941 effective_reasoning_for_model(&cfg, &cfg.default_model)
942 }
943
944 pub fn effective_reasoning_for_config(cfg: &RuntimeConfig) -> String {
945 effective_reasoning_for_model(cfg, &cfg.default_model)
946 }
947
948 pub async fn set_dynamic_workflow_effort(
949 &self,
950 effort_profile: DynamicWorkflowEffortProfile,
951 ) -> RuntimeConfig {
952 let mut cfg = self.config.write().await;
953 cfg.dynamic_workflows.effort_profile = effort_profile;
954 cfg.clone()
955 }
956
957 pub async fn dynamic_workflow_trigger_decision(
958 &self,
959 message: &str,
960 ) -> WorkflowTriggerDecision {
961 let cfg = self.config.read().await;
962 classify_workflow_trigger(message, &cfg.dynamic_workflows)
963 }
964
965 pub async fn create_thread(&self, title: Option<String>) -> anyhow::Result<ThreadMetadata> {
966 self.create_thread_with(CreateThreadRequest {
967 title,
968 workspace: self.workspace.display().to_string(),
969 workspace_id: None,
970 root_id: None,
971 provider: None,
972 model: None,
973 selection_mode: None,
974 tool_allowlist: Vec::new(),
975 developer_instructions: None,
976 external_tools: Vec::new(),
977 runner: None,
978 })
979 .await
980 }
981
982 async fn resolve_thread_runner_binding(
988 &self,
989 thread_id: &str,
990 selection: ThreadRunnerSelection,
991 ) -> anyhow::Result<ThreadRunnerBinding> {
992 let provider = self
993 .registry
994 .remote_runner_providers
995 .iter()
996 .find(|provider| provider.id() == selection.provider_id)
997 .cloned()
998 .ok_or_else(|| {
999 anyhow::anyhow!(
1000 "remote runner provider {:?} is not installed",
1001 selection.provider_id
1002 )
1003 })?;
1004 let workspace = selection.workspace.trim();
1005 anyhow::ensure!(
1006 std::path::Path::new(workspace).is_absolute(),
1007 "runner workspace must be an absolute path on the runner: {workspace:?}"
1008 );
1009 let mut read_roots = Vec::with_capacity(selection.read_roots.len());
1010 for read_root in &selection.read_roots {
1011 let trimmed = read_root.trim();
1012 anyhow::ensure!(
1013 std::path::Path::new(trimmed).is_absolute(),
1014 "runner read root must be an absolute path on the runner: {trimmed:?}"
1015 );
1016 read_roots.push(PathBuf::from(trimmed));
1017 }
1018 let destination = RunnerDestination {
1019 id: format!("thread-{thread_id}"),
1020 provider_id: selection.provider_id,
1021 config: selection.config,
1022 default_manifest: roder_api::remote_runner::RunnerManifest::default(),
1023 };
1024 provider.validate_destination(&destination).await?;
1025 Ok(ThreadRunnerBinding {
1026 destination,
1027 workspace: PathBuf::from(workspace),
1028 read_roots,
1029 })
1030 }
1031
1032 pub async fn validate_thread_runner_selection(
1034 &self,
1035 selection: ThreadRunnerSelection,
1036 ) -> anyhow::Result<()> {
1037 self.resolve_thread_runner_binding("validate", selection)
1038 .await
1039 .map(|_| ())
1040 }
1041
1042 pub async fn create_thread_with(
1043 &self,
1044 req: CreateThreadRequest,
1045 ) -> anyhow::Result<ThreadMetadata> {
1046 let cfg = self.config.read().await.clone();
1047 let now = OffsetDateTime::now_utc();
1048 let workspace = validate_thread_workspace(&req.workspace)?;
1049 let provider = req.provider.unwrap_or(cfg.default_provider);
1050 let model = req.model.unwrap_or(cfg.default_model);
1051 let selection_mode = req
1052 .selection_mode
1053 .unwrap_or_else(|| ModelSelectionMode::manual(provider.clone(), model.clone(), None));
1054 let thread_id = uuid::Uuid::new_v4().to_string();
1055 let runner_binding = match req.runner {
1056 Some(selection) => Some(
1057 self.resolve_thread_runner_binding(&thread_id, selection)
1058 .await?,
1059 ),
1060 None => None,
1061 };
1062 let runner_destination = runner_binding
1063 .as_ref()
1064 .map(|binding| binding.destination.clone())
1065 .or_else(|| cfg.remote_runner_destination.clone());
1066 let metadata = ThreadMetadata {
1067 thread_id,
1068 title: req.title,
1069 workspace,
1070 workspace_id: req.workspace_id,
1071 root_id: req.root_id,
1072 provider: Some(provider),
1073 model: Some(model),
1074 selection_mode: Some(selection_mode),
1075 tool_allowlist: req.tool_allowlist,
1076 developer_instructions: req.developer_instructions,
1077 external_tools: req.external_tools,
1078 runner_destination,
1079 runner_state: None,
1080 runner_binding,
1081 created_at: now,
1082 updated_at: now,
1083 message_count: 0,
1084 usage: None,
1085 parent_thread_id: None,
1086 forked_from_turn_id: None,
1087 workspace_fork: None,
1088 };
1089
1090 let metadata = if let Some(store) = &self.thread_store {
1091 store.create_thread(metadata).await?
1092 } else {
1093 metadata
1094 };
1095 self.emit(RoderEvent::ThreadCreated(ThreadCreated {
1096 thread_id: metadata.thread_id.clone(),
1097 timestamp: OffsetDateTime::now_utc(),
1098 }))
1099 .await;
1100 Ok(metadata)
1101 }
1102
1103 pub async fn list_threads(&self) -> anyhow::Result<Vec<ThreadMetadata>> {
1104 if let Some(store) = &self.thread_store {
1105 return store.list_threads().await;
1106 }
1107 Ok(Vec::new())
1108 }
1109
1110 pub async fn list_threads_page(
1111 &self,
1112 options: roder_api::thread::ThreadListOptions,
1113 ) -> anyhow::Result<roder_api::thread::ThreadListPage> {
1114 if let Some(store) = &self.thread_store {
1115 return store.list_threads_page(options).await;
1116 }
1117 Ok(roder_api::thread::ThreadListPage::default())
1118 }
1119
1120 pub async fn load_thread_metadata(
1121 &self,
1122 thread_id: &str,
1123 ) -> anyhow::Result<Option<roder_api::thread::ThreadMetadata>> {
1124 if let Some(store) = &self.thread_store {
1125 return store.load_thread_metadata(&thread_id.to_string()).await;
1126 }
1127 Ok(None)
1128 }
1129
1130 pub async fn archive_thread(&self, thread_id: &str) -> anyhow::Result<bool> {
1131 let archived = if let Some(store) = &self.thread_store {
1132 store.archive_thread(&thread_id.to_string()).await?
1133 } else {
1134 false
1135 };
1136 if archived {
1137 self.thread_item_cache
1138 .lock()
1139 .await
1140 .remove_thread(&thread_id.to_string());
1141 }
1142 Ok(archived)
1143 }
1144
1145 pub async fn start_team(&self, req: TeamStartRequest) -> anyhow::Result<TeamState> {
1146 let cfg = self.config.read().await.clone();
1147 let workspace = self.workspace.display().to_string();
1148 let lead_thread_id = match req.lead_thread_id {
1149 Some(thread_id) => thread_id,
1150 None => {
1151 self.create_thread_with(CreateThreadRequest {
1152 title: Some("Team lead".to_string()),
1153 workspace: workspace.clone(),
1154 workspace_id: None,
1155 root_id: None,
1156 provider: None,
1157 model: None,
1158 selection_mode: None,
1159 tool_allowlist: Vec::new(),
1160 developer_instructions: None,
1161 external_tools: Vec::new(),
1162 runner: None,
1163 })
1164 .await?
1165 .thread_id
1166 }
1167 };
1168 let team_id = uuid::Uuid::new_v4().to_string();
1169 let mut members = vec![crate::teams::lead_member(
1170 lead_thread_id.clone(),
1171 Some(cfg.default_provider.clone()),
1172 Some(cfg.default_model.clone()),
1173 cfg.policy_mode,
1174 )];
1175
1176 for (index, member) in req.members.into_iter().enumerate() {
1177 let thread = self
1178 .create_thread_with(CreateThreadRequest {
1179 title: Some(member.name.clone()),
1180 workspace: workspace.clone(),
1181 workspace_id: None,
1182 root_id: None,
1183 provider: member.model_provider.clone(),
1184 model: member.model.clone(),
1185 selection_mode: None,
1186 tool_allowlist: Vec::new(),
1187 developer_instructions: None,
1188 external_tools: Vec::new(),
1189 runner: None,
1190 })
1191 .await?;
1192 let member_id = format!("member-{}", index + 1);
1193 let descriptor = crate::teams::teammate_member(
1194 member_id.clone(),
1195 member.name,
1196 thread.thread_id.clone(),
1197 member.model_provider.or(thread.provider),
1198 member.model.or(thread.model),
1199 cfg.policy_mode,
1200 );
1201 self.emit(RoderEvent::TeamMemberStarted(TeamMemberStarted {
1202 team_id: team_id.clone(),
1203 member_id,
1204 member_thread_id: thread.thread_id,
1205 role: descriptor.role,
1206 name: descriptor.name.clone(),
1207 timestamp: OffsetDateTime::now_utc(),
1208 }))
1209 .await;
1210 members.push(descriptor);
1211 }
1212
1213 let now = OffsetDateTime::now_utc();
1214 let team = self
1215 .teams
1216 .insert(TeamState {
1217 id: team_id.clone(),
1218 lead_thread_id: lead_thread_id.clone(),
1219 display_mode: req.display_mode,
1220 members,
1221 mailbox: Vec::new(),
1222 tasks: Vec::new(),
1223 created_at: now,
1224 updated_at: now,
1225 })
1226 .await?;
1227 self.emit(RoderEvent::TeamStarted(TeamStarted {
1228 team_id,
1229 lead_thread_id,
1230 display_mode: team.display_mode,
1231 timestamp: OffsetDateTime::now_utc(),
1232 }))
1233 .await;
1234 Ok(team)
1235 }
1236
1237 pub async fn list_teams(&self) -> Vec<TeamState> {
1238 self.teams.list().await
1239 }
1240
1241 pub async fn read_team(&self, team_id: &str) -> Option<TeamState> {
1242 self.teams.get(team_id).await
1243 }
1244
1245 pub async fn start_team_member(
1246 &self,
1247 team_id: &str,
1248 req: TeamMemberStartRequest,
1249 ) -> anyhow::Result<TeamState> {
1250 let cfg = self.config.read().await.clone();
1251 let thread = self
1252 .create_thread_with(CreateThreadRequest {
1253 title: Some(req.name.clone()),
1254 workspace: self.workspace.display().to_string(),
1255 workspace_id: None,
1256 root_id: None,
1257 provider: req.model_provider.clone(),
1258 model: req.model.clone(),
1259 selection_mode: None,
1260 tool_allowlist: Vec::new(),
1261 developer_instructions: None,
1262 external_tools: Vec::new(),
1263 runner: None,
1264 })
1265 .await?;
1266 let team = self
1267 .read_team(team_id)
1268 .await
1269 .ok_or_else(|| anyhow::anyhow!("unknown team {team_id:?}"))?;
1270 let member_id = format!("member-{}", team.members.len());
1271 let descriptor = crate::teams::teammate_member(
1272 member_id.clone(),
1273 req.name,
1274 thread.thread_id.clone(),
1275 req.model_provider.or(thread.provider),
1276 req.model.or(thread.model),
1277 cfg.policy_mode,
1278 );
1279 let mut next = team;
1280 next.members.push(descriptor.clone());
1281 next.updated_at = OffsetDateTime::now_utc();
1282 let next = self.teams.insert(next).await?;
1283 self.emit(RoderEvent::TeamMemberStarted(TeamMemberStarted {
1284 team_id: next.id.clone(),
1285 member_id,
1286 member_thread_id: descriptor.thread_id,
1287 role: descriptor.role,
1288 name: descriptor.name,
1289 timestamp: OffsetDateTime::now_utc(),
1290 }))
1291 .await;
1292 Ok(next)
1293 }
1294
1295 pub async fn message_team_member(
1296 self: &Arc<Self>,
1297 team_id: &str,
1298 member_id: &str,
1299 message: String,
1300 ) -> anyhow::Result<TurnId> {
1301 let team = self
1302 .read_team(team_id)
1303 .await
1304 .ok_or_else(|| anyhow::anyhow!("unknown team {team_id:?}"))?;
1305 let member = team
1306 .members
1307 .iter()
1308 .find(|member| member.id == member_id)
1309 .ok_or_else(|| anyhow::anyhow!("unknown team member {member_id:?}"))?
1310 .clone();
1311 if member.status == TeamMemberStatus::Closed {
1312 anyhow::bail!("subagent {} is closed", member.name);
1313 }
1314 self.teams
1315 .append_mailbox_message(team_id, None, member_id.to_string(), message.clone())
1316 .await?;
1317 let workspace = self.workspace.display().to_string();
1318 let turn_id = if member.status == TeamMemberStatus::Running {
1319 if let Some(turn_id) = member.current_turn_id.clone() {
1320 self.steer_turn(
1321 member.thread_id.clone(),
1322 turn_id.clone(),
1323 message,
1324 Vec::new(),
1325 )
1326 .await?;
1327 turn_id
1328 } else {
1329 self.start_turn(StartTurnRequest {
1330 thread_id: member.thread_id.clone(),
1331 message,
1332 images: Vec::new(),
1333 provider_override: member.model_provider.clone(),
1334 model_override: member.model.clone(),
1335 reasoning_override: None,
1336 workspace: workspace.clone(),
1337 instructions: crate::default_instructions(),
1338 developer_context: None,
1339 task_ledger_required: false,
1340 })
1341 .await?
1342 }
1343 } else {
1344 self.start_turn(StartTurnRequest {
1345 thread_id: member.thread_id.clone(),
1346 message,
1347 images: Vec::new(),
1348 provider_override: member.model_provider.clone(),
1349 model_override: member.model.clone(),
1350 reasoning_override: None,
1351 workspace,
1352 instructions: crate::default_instructions(),
1353 developer_context: None,
1354 task_ledger_required: false,
1355 })
1356 .await?
1357 };
1358 let is_active = self.active_turns.read().await.contains_key(&turn_id);
1359 self.teams
1360 .update_member(team_id, member_id, |member| {
1361 if is_active {
1362 member.current_turn_id = Some(turn_id.clone());
1363 member.status = TeamMemberStatus::Running;
1364 } else {
1365 member.current_turn_id = None;
1366 member.status = TeamMemberStatus::Completed;
1367 }
1368 })
1369 .await?;
1370 if is_active {
1371 self.emit(RoderEvent::TeamMemberStatusChanged(
1372 TeamMemberStatusChanged {
1373 team_id: team_id.to_string(),
1374 member_id: member_id.to_string(),
1375 member_thread_id: member.thread_id,
1376 status: TeamMemberStatus::Running,
1377 timestamp: OffsetDateTime::now_utc(),
1378 },
1379 ))
1380 .await;
1381 } else {
1382 self.emit(RoderEvent::TeamMemberCompleted(TeamMemberCompleted {
1383 team_id: team_id.to_string(),
1384 member_id: member_id.to_string(),
1385 member_thread_id: member.thread_id,
1386 turn_id: Some(turn_id.clone()),
1387 status: TeamMemberStatus::Completed,
1388 timestamp: OffsetDateTime::now_utc(),
1389 }))
1390 .await;
1391 }
1392 Ok(turn_id)
1393 }
1394
1395 pub async fn set_team_member_policy_mode(
1396 &self,
1397 team_id: &str,
1398 member_id: &str,
1399 policy_mode: PolicyMode,
1400 ) -> anyhow::Result<TeamState> {
1401 self.teams
1402 .set_member_policy_mode(team_id, member_id, policy_mode)
1403 .await
1404 }
1405
1406 pub async fn interrupt_team_member(
1407 &self,
1408 team_id: &str,
1409 member_id: &str,
1410 ) -> anyhow::Result<Option<TurnId>> {
1411 let team = self
1412 .read_team(team_id)
1413 .await
1414 .ok_or_else(|| anyhow::anyhow!("unknown team {team_id:?}"))?;
1415 let member = team
1416 .members
1417 .iter()
1418 .find(|member| member.id == member_id)
1419 .ok_or_else(|| anyhow::anyhow!("unknown team member {member_id:?}"))?
1420 .clone();
1421 let Some(turn_id) = member.current_turn_id.clone() else {
1422 return Ok(None);
1423 };
1424 self.interrupt_turn(member.thread_id.clone(), turn_id.clone())
1425 .await?;
1426 self.teams
1427 .update_member(team_id, member_id, |member| {
1428 member.status = TeamMemberStatus::Interrupted;
1429 member.current_turn_id = None;
1430 })
1431 .await?;
1432 self.emit(RoderEvent::TeamMemberCompleted(TeamMemberCompleted {
1433 team_id: team_id.to_string(),
1434 member_id: member_id.to_string(),
1435 member_thread_id: member.thread_id,
1436 turn_id: Some(turn_id.clone()),
1437 status: TeamMemberStatus::Interrupted,
1438 timestamp: OffsetDateTime::now_utc(),
1439 }))
1440 .await;
1441 Ok(Some(turn_id))
1442 }
1443
1444 pub async fn close_team_member(
1445 &self,
1446 team_id: &str,
1447 member_id: &str,
1448 ) -> anyhow::Result<roder_api::teams::TeamMemberDescriptor> {
1449 let team = self
1450 .read_team(team_id)
1451 .await
1452 .ok_or_else(|| anyhow::anyhow!("unknown team {team_id:?}"))?;
1453 let member = team
1454 .members
1455 .iter()
1456 .find(|member| member.id == member_id)
1457 .ok_or_else(|| anyhow::anyhow!("unknown team member {member_id:?}"))?
1458 .clone();
1459 if member.role == roder_api::teams::TeamMemberRole::Lead {
1460 anyhow::bail!("team lead cannot be closed as a subagent");
1461 }
1462 let interrupted_turn_id = if member.status == TeamMemberStatus::Running {
1463 if let Some(turn_id) = member.current_turn_id.clone() {
1464 self.interrupt_turn(member.thread_id.clone(), turn_id.clone())
1465 .await?;
1466 Some(turn_id)
1467 } else {
1468 None
1469 }
1470 } else {
1471 member.current_turn_id.clone()
1472 };
1473 let updated = self
1474 .teams
1475 .update_member(team_id, member_id, |member| {
1476 member.status = TeamMemberStatus::Closed;
1477 member.current_turn_id = None;
1478 })
1479 .await?;
1480 let closed = updated
1481 .members
1482 .iter()
1483 .find(|member| member.id == member_id)
1484 .cloned()
1485 .ok_or_else(|| anyhow::anyhow!("closed team member disappeared"))?;
1486 self.emit(crate::agent_control_tools::closed_member_event(
1487 team_id.to_string(),
1488 &closed,
1489 interrupted_turn_id,
1490 ))
1491 .await;
1492 Ok(closed)
1493 }
1494
1495 pub async fn cleanup_team(&self, team_id: &str, force: bool) -> anyhow::Result<bool> {
1496 let Some(team) = self.read_team(team_id).await else {
1497 return Ok(false);
1498 };
1499 if !force
1500 && team
1501 .members
1502 .iter()
1503 .any(|member| member.status == TeamMemberStatus::Running)
1504 {
1505 anyhow::bail!("team {team_id:?} has active teammates; use forced cleanup");
1506 }
1507 let removed = self.teams.remove(team_id).await?.is_some();
1508 if removed {
1509 self.emit(RoderEvent::TeamCleanupCompleted(TeamCleanupCompleted {
1510 team_id: team_id.to_string(),
1511 forced: force,
1512 timestamp: OffsetDateTime::now_utc(),
1513 }))
1514 .await;
1515 }
1516 Ok(removed)
1517 }
1518
1519 pub async fn effective_policy_mode_for_thread(&self, thread_id: &str) -> PolicyMode {
1520 if let Some(mode) = self.teams.policy_mode_for_thread(thread_id).await {
1521 return mode;
1522 }
1523 self.status().await.policy_mode
1524 }
1525
1526 async fn complete_team_member_turn(
1527 &self,
1528 thread_id: &ThreadId,
1529 turn_id: &TurnId,
1530 status: TeamMemberStatus,
1531 ) -> anyhow::Result<()> {
1532 let Some((team_id, member)) = self
1533 .teams
1534 .complete_member_turn(thread_id, turn_id, status)
1535 .await?
1536 else {
1537 return Ok(());
1538 };
1539 self.emit(RoderEvent::TeamMemberCompleted(TeamMemberCompleted {
1540 team_id,
1541 member_id: member.id,
1542 member_thread_id: member.thread_id,
1543 turn_id: Some(turn_id.clone()),
1544 status,
1545 timestamp: OffsetDateTime::now_utc(),
1546 }))
1547 .await;
1548 Ok(())
1549 }
1550
1551 pub async fn load_thread(
1552 &self,
1553 thread_id: &ThreadId,
1554 ) -> anyhow::Result<Option<ThreadSnapshot>> {
1555 let loaded = if let Some(store) = &self.thread_store {
1556 store.load_thread(thread_id).await?
1557 } else {
1558 None
1559 };
1560 if loaded.is_some() {
1561 self.emit(RoderEvent::ThreadLoaded(ThreadLoaded {
1562 thread_id: thread_id.clone(),
1563 timestamp: OffsetDateTime::now_utc(),
1564 }))
1565 .await;
1566 }
1567 Ok(loaded)
1568 }
1569
1570 pub async fn workspace_for_thread(&self, thread_id: &ThreadId) -> anyhow::Result<String> {
1571 if let Some(store) = &self.thread_store {
1572 let snapshot = store
1573 .load_thread(thread_id)
1574 .await?
1575 .ok_or_else(|| anyhow::anyhow!("thread not found: {thread_id}"));
1576 match snapshot {
1577 Ok(snapshot) => {
1578 if let Some(metadata) = snapshot.metadata {
1579 if let Some(fork) = &metadata.workspace_fork
1582 && fork.status == roder_api::forks::ForkStatus::Active
1583 && !std::path::Path::new(&metadata.workspace).is_dir()
1584 {
1585 anyhow::bail!(
1586 "workspace fork {} is missing its workspace at {}; restore it or \
1587 remove the fork before running turns in this thread",
1588 fork.id,
1589 metadata.workspace
1590 );
1591 }
1592 return Ok(metadata.workspace);
1593 }
1594 eprintln!(
1595 "thread metadata missing while resolving workspace for {thread_id}; falling back to runtime workspace"
1596 );
1597 }
1598 Err(err) => {
1599 eprintln!(
1600 "thread missing while resolving workspace for {thread_id}: {err}; falling back to runtime workspace"
1601 );
1602 }
1603 }
1604 }
1605 Ok(self.workspace.display().to_string())
1606 }
1607
1608 async fn selection_mode_for_thread(
1609 &self,
1610 thread_id: &ThreadId,
1611 ) -> anyhow::Result<Option<ModelSelectionMode>> {
1612 let Some(store) = &self.thread_store else {
1613 return Ok(None);
1614 };
1615 Ok(store
1616 .load_thread(thread_id)
1617 .await?
1618 .and_then(|snapshot| snapshot.metadata)
1619 .and_then(|metadata| {
1620 metadata
1621 .selection_mode
1622 .or_else(|| match (metadata.provider, metadata.model) {
1623 (Some(provider), Some(model)) => {
1624 Some(ModelSelectionMode::manual(provider, model, None))
1625 }
1626 _ => None,
1627 })
1628 }))
1629 }
1630
1631 pub(crate) async fn thread_turn_overrides(
1633 &self,
1634 thread_id: &ThreadId,
1635 ) -> anyhow::Result<ThreadTurnOverrides> {
1636 let Some(store) = &self.thread_store else {
1637 return Ok(ThreadTurnOverrides::default());
1638 };
1639 Ok(store
1640 .load_thread_metadata(thread_id)
1641 .await?
1642 .map(|metadata| ThreadTurnOverrides {
1643 tool_allowlist: metadata.tool_allowlist,
1644 developer_instructions: metadata.developer_instructions,
1645 external_tools: metadata.external_tools,
1646 })
1647 .unwrap_or_default())
1648 }
1649
1650 pub async fn set_thread_selection_mode(
1651 &self,
1652 thread_id: &ThreadId,
1653 selection_mode: ModelSelectionMode,
1654 ) -> anyhow::Result<()> {
1655 let Some(store) = &self.thread_store else {
1656 return Ok(());
1657 };
1658 let Some(snapshot) = store.load_thread(thread_id).await? else {
1659 anyhow::bail!("thread not found: {thread_id}");
1660 };
1661 let Some(mut metadata) = snapshot.metadata else {
1662 return Ok(());
1663 };
1664 let concrete = selection_mode.concrete_selection();
1665 metadata.provider = Some(concrete.provider);
1666 metadata.model = Some(concrete.model);
1667 metadata.selection_mode = Some(selection_mode);
1668 metadata.updated_at = OffsetDateTime::now_utc();
1669 store.update_thread_metadata(metadata).await?;
1670 Ok(())
1671 }
1672
1673 async fn runner_session_for_thread(
1674 &self,
1675 thread_id: &ThreadId,
1676 ) -> anyhow::Result<Option<(RunnerDestination, Arc<dyn RemoteRunnerSession>)>> {
1677 let metadata = if let Some(store) = &self.thread_store {
1678 store.load_thread_metadata(thread_id).await?
1679 } else {
1680 None
1681 };
1682 let destination = metadata
1684 .as_ref()
1685 .and_then(|metadata| metadata.runner_binding.as_ref())
1686 .map(|binding| binding.destination.clone())
1687 .or(self.config.read().await.remote_runner_destination.clone());
1688 let Some(destination) = destination else {
1689 return Ok(None);
1690 };
1691 let provider = self
1692 .registry
1693 .remote_runner_providers
1694 .iter()
1695 .find(|provider| provider.id() == destination.provider_id)
1696 .cloned()
1697 .ok_or_else(|| {
1698 anyhow::anyhow!(
1699 "remote runner provider {:?} is not installed",
1700 destination.provider_id
1701 )
1702 })?;
1703 let persisted_state = metadata.and_then(|metadata| metadata.runner_state);
1704 let session = if let Some(state) = persisted_state
1705 && state.provider_id == destination.provider_id
1706 && state.destination_id == destination.id
1707 {
1708 match provider.resume_session(state).await {
1709 Ok(session) => session,
1710 Err(_) => provider.create_session(destination.clone()).await?,
1711 }
1712 } else {
1713 provider.create_session(destination.clone()).await?
1714 };
1715 Ok(Some((destination, session)))
1716 }
1717
1718 pub(crate) async fn remote_workspace_for_thread(
1724 &self,
1725 thread_id: &ThreadId,
1726 ) -> anyhow::Result<Option<Arc<RemoteWorkspace>>> {
1727 let Some(store) = &self.thread_store else {
1728 return Ok(None);
1729 };
1730 let binding = store
1731 .load_thread_metadata(thread_id)
1732 .await?
1733 .and_then(|metadata| metadata.runner_binding);
1734 let Some(binding) = binding else {
1735 return Ok(None);
1736 };
1737 let session = self
1738 .runner_session_for_thread(thread_id)
1739 .await?
1740 .map(|(_, session)| session)
1741 .ok_or_else(|| {
1742 anyhow::anyhow!("runner-bound thread {thread_id} has no runner session")
1743 })?;
1744 Ok(Some(Arc::new(RemoteWorkspace {
1745 session,
1746 root: binding.workspace,
1747 read_roots: binding.read_roots,
1748 })))
1749 }
1750
1751 async fn persist_runner_state(
1752 &self,
1753 thread_id: &ThreadId,
1754 runner: Option<&(RunnerDestination, Arc<dyn RemoteRunnerSession>)>,
1755 ) -> anyhow::Result<()> {
1756 let Some((destination, session)) = runner else {
1757 return Ok(());
1758 };
1759 let Some(store) = &self.thread_store else {
1760 return Ok(());
1761 };
1762 let Some(snapshot) = store.load_thread(thread_id).await? else {
1763 return Ok(());
1764 };
1765 let Some(mut metadata) = snapshot.metadata else {
1766 return Ok(());
1767 };
1768 metadata.runner_destination = Some(destination.clone());
1769 metadata.runner_state = Some(session.state());
1770 metadata.updated_at = OffsetDateTime::now_utc();
1771 store.update_thread_metadata(metadata).await?;
1772 Ok(())
1773 }
1774
1775 async fn record_thread_usage_metadata(
1776 &self,
1777 thread_id: &ThreadId,
1778 usage: &TokenUsage,
1779 ) -> anyhow::Result<()> {
1780 if usage.is_empty() {
1781 return Ok(());
1782 }
1783 let Some(store) = &self.thread_store else {
1784 return Ok(());
1785 };
1786 let Some(snapshot) = store.load_thread(thread_id).await? else {
1787 return Ok(());
1788 };
1789 let Some(mut metadata) = snapshot.metadata else {
1790 return Ok(());
1791 };
1792 metadata
1793 .usage
1794 .get_or_insert_with(ThreadUsageMetadata::default)
1795 .add_token_usage(usage);
1796 metadata.updated_at = OffsetDateTime::now_utc();
1797 store.update_thread_metadata(metadata).await?;
1798 Ok(())
1799 }
1800
1801 pub fn start_turn(
1802 self: &Arc<Self>,
1803 mut req: StartTurnRequest,
1804 ) -> BoxFuture<'_, anyhow::Result<TurnId>> {
1805 Box::pin(async move {
1806 req.workspace = validate_thread_workspace(&req.workspace)?;
1807 let cfg = self.config.read().await.clone();
1808 let provider = req
1809 .provider_override
1810 .clone()
1811 .unwrap_or_else(|| cfg.default_provider.clone());
1812 self.engine_for(&provider)?;
1813 let turn_id = uuid::Uuid::new_v4().to_string();
1814 let (abort_handle, abort_registration) = AbortHandle::new_pair();
1815 let active = ActiveTurnHandle {
1816 thread_id: req.thread_id.clone(),
1817 abort: abort_handle,
1818 steers: Arc::new(Mutex::new(Vec::new())),
1819 };
1820 self.active_turns
1821 .write()
1822 .await
1823 .insert(turn_id.clone(), active);
1824 let runtime = Arc::clone(self);
1825 let turn_req = req;
1826 let thread_id_for_task = turn_req.thread_id.clone();
1827 let turn_id_for_task = turn_id.clone();
1828 tokio::spawn(async move {
1829 let result = Abortable::new(
1830 runtime.run_turn(turn_req, turn_id_for_task.clone()),
1831 abort_registration,
1832 )
1833 .await;
1834 runtime
1842 .cancel_pending_external_tool_calls_for_turn(&turn_id_for_task)
1843 .await;
1844 let completed = matches!(&result, Ok(Ok(TurnRunOutcome::Completed)));
1845 if let Ok(Err(err)) = &result {
1846 runtime
1848 .emit(RoderEvent::TurnFailed(TurnFailed {
1849 thread_id: thread_id_for_task.clone(),
1850 turn_id: turn_id_for_task.clone(),
1851 error: err.to_string(),
1852 error_kind: None,
1853 usage: None,
1854 timestamp: OffsetDateTime::now_utc(),
1855 }))
1856 .await;
1857 }
1858 runtime.active_turns.write().await.remove(&turn_id_for_task);
1859 if completed {
1860 let _ = runtime
1861 .continue_active_goal_after_turn(thread_id_for_task)
1862 .await;
1863 }
1864 });
1865 Ok(turn_id)
1866 })
1867 }
1868
1869 pub(crate) async fn has_active_turn_for_thread(&self, thread_id: &ThreadId) -> bool {
1870 self.active_turns
1871 .read()
1872 .await
1873 .values()
1874 .any(|handle| &handle.thread_id == thread_id)
1875 }
1876
1877 pub async fn active_turn_count(&self) -> usize {
1880 self.active_turns.read().await.len()
1881 }
1882
1883 pub async fn active_turn_for_thread(&self, thread_id: &ThreadId) -> Option<TurnId> {
1884 self.active_turns
1885 .read()
1886 .await
1887 .iter()
1888 .find_map(|(turn_id, handle)| (&handle.thread_id == thread_id).then(|| turn_id.clone()))
1889 }
1890
1891 pub async fn thread_activity(&self, thread_id: &ThreadId) -> ThreadActivity {
1892 let Some(active_turn_id) = self.active_turn_for_thread(thread_id).await else {
1893 return ThreadActivity::default();
1894 };
1895
1896 let mut active_flags = Vec::new();
1897 {
1898 let pending_approvals = self.pending_tool_approvals.lock().await;
1899 if pending_approvals
1900 .values()
1901 .any(|pending| &pending.thread_id == thread_id && pending.turn_id == active_turn_id)
1902 {
1903 active_flags.push("approvalRequired".to_string());
1904 }
1905 }
1906 {
1907 let pending_inputs = self.pending_user_inputs.lock().await;
1908 if pending_inputs
1909 .values()
1910 .any(|pending| &pending.thread_id == thread_id && pending.turn_id == active_turn_id)
1911 {
1912 active_flags.push("userInputRequired".to_string());
1913 }
1914 }
1915 {
1916 let pending_external = self.pending_external_tool_calls.lock().await;
1917 if pending_external
1918 .values()
1919 .any(|pending| &pending.thread_id == thread_id && pending.turn_id == active_turn_id)
1920 {
1921 active_flags.push("externalToolPending".to_string());
1922 }
1923 }
1924 if self.pending_plan_exit().await.is_some_and(|pending| {
1925 &pending.thread_id == thread_id && pending.turn_id == active_turn_id
1926 }) {
1927 active_flags.push("planExitRequired".to_string());
1928 }
1929
1930 ThreadActivity {
1931 active_turn_id: Some(active_turn_id),
1932 active_flags,
1933 }
1934 }
1935
1936 pub async fn interrupt_turn(&self, thread_id: ThreadId, turn_id: TurnId) -> anyhow::Result<()> {
1937 if let Some(handle) = self.active_turns.write().await.remove(&turn_id) {
1938 handle.abort.abort();
1939 }
1940 self.cancel_pending_external_tool_calls_for_turn(&turn_id)
1941 .await;
1942 self.emit(RoderEvent::TurnInterrupted(TurnInterrupted {
1943 thread_id,
1944 turn_id,
1945 timestamp: OffsetDateTime::now_utc(),
1946 }))
1947 .await;
1948 Ok(())
1949 }
1950
1951 pub async fn steer_turn(
1952 &self,
1953 thread_id: ThreadId,
1954 turn_id: TurnId,
1955 message: String,
1956 images: Vec<InputImage>,
1957 ) -> anyhow::Result<()> {
1958 let message = message.trim().to_string();
1959 if message.is_empty() && images.is_empty() {
1960 return Ok(());
1961 }
1962
1963 let Some(active) = self.active_turns.read().await.get(&turn_id).cloned() else {
1964 anyhow::bail!("no active turn to steer");
1965 };
1966 active
1967 .steers
1968 .lock()
1969 .await
1970 .push(UserMessage::with_images(message.clone(), images));
1971 self.emit(RoderEvent::TurnSteered(TurnSteered {
1972 thread_id,
1973 turn_id,
1974 message,
1975 timestamp: OffsetDateTime::now_utc(),
1976 }))
1977 .await;
1978 Ok(())
1979 }
1980
1981 pub async fn tool_specs(&self) -> Vec<roder_api::tools::ToolSpec> {
1982 let cfg = self.config.read().await;
1983 let model_profile =
1984 model_profile_for_provider_model(&cfg, &cfg.default_provider, &cfg.default_model);
1985 self.filtered_tool_specs(&cfg, &cfg.default_model, model_profile.as_ref(), &[], &[])
1986 }
1987
1988 pub fn subagent_definitions(&self) -> Vec<SubagentDefinition> {
1989 self.registry
1990 .subagent_dispatchers
1991 .iter()
1992 .flat_map(|dispatcher| dispatcher.definitions())
1993 .collect()
1994 }
1995
1996 async fn run_turn(
1997 self: &Arc<Self>,
1998 req: StartTurnRequest,
1999 turn_id: TurnId,
2000 ) -> anyhow::Result<TurnRunOutcome> {
2001 let turn_started_at = OffsetDateTime::now_utc();
2002 self.emit(RoderEvent::TurnStarted(TurnStarted {
2003 thread_id: req.thread_id.clone(),
2004 turn_id: turn_id.clone(),
2005 runtime_profile: self.config.read().await.runtime_profile,
2006 timestamp: turn_started_at,
2007 }))
2008 .await;
2009 self.persist_turn_item(
2010 &req.thread_id,
2011 &turn_id,
2012 &TranscriptItem::UserMessage(UserMessage::with_images(
2013 req.message.clone(),
2014 req.images.clone(),
2015 )),
2016 )
2017 .await?;
2018
2019 let mut cfg = self.config.read().await.clone();
2020 let runtime_profile = cfg.runtime_profile;
2021 let turn_deadline = turn_deadline_for_config(&cfg);
2022 let deadline_finalization_reserve =
2023 crate::deadline_policy::finalization_reserve_seconds(cfg.turn_deadline_seconds);
2024 let selection_mode = self.selection_mode_for_thread(&req.thread_id).await?;
2025 let concrete_selection = selection_mode
2026 .as_ref()
2027 .map(ModelSelectionMode::concrete_selection);
2028 let default_provider = req
2029 .provider_override
2030 .clone()
2031 .or_else(|| {
2032 concrete_selection
2033 .as_ref()
2034 .map(|selection| selection.provider.clone())
2035 })
2036 .unwrap_or(cfg.default_provider.clone());
2037 let default_model = req
2038 .model_override
2039 .clone()
2040 .or_else(|| {
2041 concrete_selection
2042 .as_ref()
2043 .map(|selection| selection.model.clone())
2044 })
2045 .unwrap_or(cfg.default_model.clone());
2046 if let Some(reasoning) = req.reasoning_override.as_deref().or_else(|| {
2047 selection_mode
2048 .as_ref()
2049 .and_then(ModelSelectionMode::reasoning)
2050 }) {
2051 validate_reasoning_effort(&default_model, reasoning)?;
2052 cfg.reasoning = Some(reasoning.to_string());
2053 }
2054 let turn_has_concrete_model_override =
2055 req.provider_override.is_some() || req.model_override.is_some();
2056 let (turn_inference_router, turn_inference_router_profile) = match &selection_mode {
2057 Some(ModelSelectionMode::Auto {
2058 router_id, profile, ..
2059 }) if !turn_has_concrete_model_override => (
2060 RuntimeInferenceRouterConfig {
2061 enabled: true,
2062 router_id: Some(router_id.clone()),
2063 },
2064 profile.clone(),
2065 ),
2066 _ => (RuntimeInferenceRouterConfig::disabled(), None),
2067 };
2068 let mut provider = default_provider.clone();
2069 let mut model = default_model.clone();
2070 let mut model_profile = model_profile_for_provider_model(&cfg, &provider, &model);
2071 let workspace = req.workspace.clone();
2072 let mut transcript = self.transcript_for_turn(&req, &turn_id, &model).await?;
2073 let mut compacted_this_turn = transcript
2074 .iter()
2075 .any(|item| matches!(item, TranscriptItem::ContextCompaction(_)));
2076 let runner_session = self.runner_session_for_thread(&req.thread_id).await?;
2077 let effective_policy_mode = self.effective_policy_mode_for_thread(&req.thread_id).await;
2078 let thread_overrides = self.thread_turn_overrides(&req.thread_id).await?;
2079 let mut final_assistant_text = String::new();
2080 let mut final_phase_messages = Vec::<AssistantMessage>::new();
2081 let mut final_reasoning_text = String::new();
2082 let mut final_provider_metadata = None;
2083 let mut exhausted_tool_rounds = true;
2084 let mut verification_gate =
2085 VerificationGateState::new(req.message.clone(), runtime_profile);
2086 let mut speed_policy = SpeedPolicyState::default();
2087 let mut reliability = TurnReliabilityState::default();
2088 let mut turn_usage = TokenUsage::default();
2089 let mut turn_finish_reason: Option<String> = None;
2093 let mut deadline_finalization_requested = false;
2094 let mut deadline_scoreable_completion_requested = false;
2095 let mut task_ledger_completion_reminders = 0_u8;
2096 let mut task_ledger_scoreable_checkpoints = 0_u8;
2097 let mut provider_stream_retry_attempts = 0_u32;
2098 let mut routing_candidates = None;
2099 let routing_transcript_start = transcript.len().saturating_sub(1);
2100 let mut routing_escalations = 0_u32;
2101 let mut model_switch_summary_selection = None::<ModelSelection>;
2102
2103 'tool_rounds: for round_index in 0..MAX_TOOL_ROUNDS_PER_TURN {
2104 if let Some(deadline) = turn_deadline
2105 && deadline_expired(deadline)
2106 {
2107 self.fail_turn_due_to_deadline(&req.thread_id, &turn_id, deadline, &transcript)
2108 .await?;
2109 return Ok(TurnRunOutcome::Stopped);
2110 }
2111 let steers = self.drain_turn_steers(&turn_id).await;
2112 self.append_steers(&req, &turn_id, &mut transcript, steers)
2113 .await?;
2114 if runtime_profile == RuntimeProfile::Eval
2115 && let Some(remaining) = crate::deadline_policy::should_start_finalization(
2116 turn_deadline,
2117 deadline_finalization_reserve,
2118 deadline_finalization_requested || deadline_scoreable_completion_requested,
2119 )
2120 {
2121 if req.task_ledger_required
2122 && task_ledger_completion_reminders < TASK_LEDGER_COMPLETION_REMINDER_LIMIT
2123 && let Some(prompt) = task_ledger_completion_prompt(&transcript)
2124 {
2125 task_ledger_completion_reminders += 1;
2126 deadline_scoreable_completion_requested = true;
2127 let item = TranscriptItem::UserMessage(UserMessage::text(
2128 task_ledger_deadline_completion_prompt(
2129 remaining,
2130 deadline_finalization_reserve,
2131 &prompt,
2132 ),
2133 ));
2134 self.persist_turn_item(&req.thread_id, &turn_id, &item)
2135 .await?;
2136 transcript.push(item);
2137 continue 'tool_rounds;
2138 } else {
2139 self.start_deadline_finalization(
2140 &req.thread_id,
2141 &turn_id,
2142 &mut transcript,
2143 remaining,
2144 )
2145 .await?;
2146 deadline_finalization_requested = true;
2147 }
2148 }
2149 if runtime_profile == RuntimeProfile::Eval
2150 && req.task_ledger_required
2151 && !deadline_finalization_requested
2152 && task_ledger_scoreable_checkpoints < TASK_LEDGER_SCOREABLE_CHECKPOINT_LIMIT
2153 && let Some(remaining) = deadline_remaining_seconds(turn_deadline)
2154 && remaining <= TASK_LEDGER_SCOREABLE_CHECKPOINT_SECONDS
2155 && remaining > deadline_finalization_reserve
2156 && let Some(prompt) = task_ledger_completion_prompt(&transcript)
2157 {
2158 task_ledger_scoreable_checkpoints += 1;
2159 let item = TranscriptItem::UserMessage(UserMessage::text(
2160 task_ledger_scoreable_checkpoint_prompt(remaining, &prompt),
2161 ));
2162 self.persist_turn_item(&req.thread_id, &turn_id, &item)
2163 .await?;
2164 transcript.push(item);
2165 continue 'tool_rounds;
2166 }
2167 if turn_inference_router.is_active() && routing_candidates.is_none() {
2168 routing_candidates =
2169 Some(collect_inference_routing_candidates(&self.registry).await);
2170 }
2171 let routing_tools_model = model.clone();
2172 let routing_tools = self.filtered_tool_specs(
2173 &cfg,
2174 &model,
2175 model_profile.as_ref(),
2176 &thread_overrides.tool_allowlist,
2177 &thread_overrides.external_tools,
2178 );
2179 let prior_failures =
2180 transcript_failure_count_since(&transcript, routing_transcript_start)
2181 .max(reliability.tool_failure_count())
2182 .saturating_add(provider_stream_retry_attempts);
2183 let routing_selection = route_inference_selection(
2184 &self.registry,
2185 &turn_inference_router,
2186 InferenceRoutingRequest {
2187 thread_id: &req.thread_id,
2188 turn_id: &turn_id,
2189 round_index: round_index as u32,
2190 runtime_profile,
2191 phase: speed_policy.phase(),
2192 profile: turn_inference_router_profile.as_deref(),
2193 default_selection: ModelSelection {
2194 provider: default_provider.clone(),
2195 model: default_model.clone(),
2196 },
2197 transcript: &transcript,
2198 tools: &routing_tools,
2199 candidates: routing_candidates.as_deref(),
2200 prior_failures,
2201 prior_escalations: routing_escalations,
2202 },
2203 )
2204 .await;
2205 if let Some(decision) = routing_selection.decision.clone() {
2206 if matches!(decision.outcome, InferenceRoutingOutcome::Escalated) {
2207 routing_escalations = routing_escalations.saturating_add(1);
2208 }
2209 self.emit(RoderEvent::InferenceRoutingDecision(
2210 InferenceRoutingDecisionEvent {
2211 thread_id: req.thread_id.clone(),
2212 turn_id: turn_id.clone(),
2213 round_index: round_index as u32,
2214 default_selection: ModelSelection {
2215 provider: default_provider.clone(),
2216 model: default_model.clone(),
2217 },
2218 selected_selection: routing_selection.selection.clone(),
2219 decision,
2220 timestamp: OffsetDateTime::now_utc(),
2221 },
2222 ))
2223 .await;
2224 }
2225 provider = routing_selection.selection.provider.clone();
2226 model = routing_selection.selection.model.clone();
2227 let engine = self.engine_for(&provider)?;
2228 let capabilities = engine.capabilities();
2229 model_profile = model_profile_for_provider_model(&cfg, &provider, &model);
2230 let tools = if capabilities.tool_calls {
2231 if model == routing_tools_model {
2232 routing_tools.clone()
2233 } else {
2234 self.filtered_tool_specs(
2235 &cfg,
2236 &model,
2237 model_profile.as_ref(),
2238 &thread_overrides.tool_allowlist,
2239 &thread_overrides.external_tools,
2240 )
2241 }
2242 } else {
2243 Vec::new()
2244 };
2245 let parallel_tool_calls = parallel_tool_calls_for_model(&cfg, &model);
2246 let tool_choice = if tools.is_empty() {
2247 ToolChoice::None
2248 } else {
2249 ToolChoice::Auto
2250 };
2251 let summary_selection = ModelSelection {
2252 provider: provider.clone(),
2253 model: model.clone(),
2254 };
2255 if model_switch_summary_selection.as_ref() != Some(&summary_selection) {
2256 if let Some(summary) = model_switch_summary(
2257 &transcript,
2258 model_profile.as_ref(),
2259 &provider,
2260 &model,
2261 &tools,
2262 ) {
2263 let item = TranscriptItem::UserMessage(UserMessage::text(summary));
2264 self.persist_turn_item(&req.thread_id, &turn_id, &item)
2265 .await?;
2266 transcript.push(item);
2267 }
2268 model_switch_summary_selection = Some(summary_selection);
2269 }
2270
2271 if !capabilities.image_input && transcript_has_images(&transcript) {
2272 self.fail_turn_with_error(
2273 &req.thread_id,
2274 &turn_id,
2275 format!("provider {provider} does not support image input"),
2276 )
2277 .await?;
2278 return Ok(TurnRunOutcome::Stopped);
2279 }
2280 transcript = self
2281 .compact_transcript_if_needed(
2282 &req.thread_id,
2283 &turn_id,
2284 &provider,
2285 &model,
2286 transcript,
2287 self.compaction_options_for_turn(&req.thread_id, !compacted_this_turn),
2288 )
2289 .await?;
2290 compacted_this_turn = compacted_this_turn
2291 || transcript
2292 .iter()
2293 .any(|item| matches!(item, TranscriptItem::ContextCompaction(_)));
2294
2295 let speed_policy_decision =
2296 speed_policy.decision(runtime_profile, &model, &cfg.speed_policy);
2297 let request_reasoning = reasoning_from_decision(
2298 speed_policy_decision.as_ref(),
2299 routing_selection
2300 .reasoning
2301 .clone()
2302 .unwrap_or_else(|| reasoning_for_model(&cfg, &model)),
2303 );
2304 if let Some(limit) = reliability.record_model_call(
2305 &cfg.reliability,
2306 runtime_profile == RuntimeProfile::Interactive,
2307 ) {
2308 self.fail_turn_due_to_reliability_limit(
2309 &req.thread_id,
2310 &turn_id,
2311 &provider,
2312 &model,
2313 limit,
2314 &transcript,
2315 )
2316 .await?;
2317 return Ok(TurnRunOutcome::Stopped);
2318 }
2319 self.emit(RoderEvent::InferenceStarted(InferenceStarted {
2320 thread_id: req.thread_id.clone(),
2321 turn_id: turn_id.clone(),
2322 engine_id: engine.id(),
2323 model: ModelSelection {
2324 provider: provider.clone(),
2325 model: model.clone(),
2326 },
2327 reasoning: request_reasoning.clone(),
2328 speed_policy: speed_policy_decision.clone(),
2329 deadline_remaining_seconds: deadline_remaining_seconds(turn_deadline),
2330 timestamp: OffsetDateTime::now_utc(),
2331 }))
2332 .await;
2333
2334 let mut instructions = req.instructions.clone();
2335 if let Some(extra) = &thread_overrides.developer_instructions {
2336 instructions = apply_thread_developer_instructions(instructions, extra);
2337 }
2338 if let Some(context) = req.developer_context.as_deref() {
2339 instructions = apply_turn_developer_context(instructions, context);
2340 }
2341 let mut instructions = apply_runtime_profile(instructions, runtime_profile);
2342 if let Some(profile) = &model_profile {
2343 instructions = apply_model_instruction_overlay(instructions, profile);
2344 }
2345 if req.task_ledger_required
2346 && runtime_profile == RuntimeProfile::Eval
2347 && !transcript_has_task_ledger(&transcript)
2348 {
2349 instructions = apply_task_ledger_required(instructions);
2350 }
2351 if effective_policy_mode == PolicyMode::Plan {
2352 instructions = apply_plan_mode(instructions);
2353 }
2354 instructions = self
2355 .goals
2356 .apply_goal_instructions(&req.thread_id, instructions)
2357 .await?;
2358 let mut request_metadata = serde_json::json!({});
2359 if let Some(decision) = &speed_policy_decision {
2360 request_metadata["speedPolicy"] = serde_json::json!(decision);
2361 }
2362 if let Some(decision) = routing_selection.decision.as_ref() {
2363 request_metadata["inferenceRouting"] = serde_json::json!(decision);
2364 }
2365 if let Some(remaining) = deadline_remaining_seconds(turn_deadline) {
2366 request_metadata["deadlineRemainingSeconds"] = serde_json::json!(remaining);
2367 }
2368 if let Some(profile) = &model_profile {
2369 request_metadata["modelProfile"] = serde_json::json!({
2370 "model": profile.model,
2371 "providerFamily": profile.provider_family,
2372 "editTool": profile.edit_tool,
2373 "schemaPolicy": profile.schema_policy,
2374 "instructionOverlay": profile.instruction_overlay,
2375 "parallelToolCalls": profile.parallel_tool_calls,
2376 "autoCompactTokenLimit": profile.auto_compact_token_limit,
2377 });
2378 }
2379 let task_ledger_required_this_round = req.task_ledger_required
2380 && runtime_profile == RuntimeProfile::Eval
2381 && !deadline_finalization_requested
2382 && !transcript_has_task_ledger(&transcript);
2383 let task_ledger_tools = (capabilities.tool_calls && task_ledger_required_this_round)
2384 .then(|| self.task_ledger_tool_specs(model_profile.as_ref()))
2385 .filter(|tools| !tools.is_empty());
2386 let request_tools = if deadline_finalization_requested {
2387 Vec::new()
2388 } else if let Some(ledger_tools) = &task_ledger_tools {
2389 ledger_tools.clone()
2390 } else {
2391 tools.clone()
2392 };
2393 let request_tool_choice = if deadline_finalization_requested {
2394 ToolChoice::None
2395 } else if task_ledger_tools.is_some() {
2396 ToolChoice::Specific(TASK_LEDGER_TOOL_NAME.to_string())
2397 } else {
2398 tool_choice.clone()
2399 };
2400 if deadline_finalization_requested {
2401 request_metadata["deadlineFinalization"] = serde_json::json!({
2402 "reserveSeconds": deadline_finalization_reserve,
2403 "remainingSeconds": deadline_remaining_seconds(turn_deadline),
2404 });
2405 } else if deadline_scoreable_completion_requested {
2406 request_metadata["deadlineScoreableCompletion"] = serde_json::json!({
2407 "reserveSeconds": deadline_finalization_reserve,
2408 "remainingSeconds": deadline_remaining_seconds(turn_deadline),
2409 });
2410 }
2411 let request = AgentInferenceRequest {
2412 model: ModelSelection {
2413 provider: provider.clone(),
2414 model: model.clone(),
2415 },
2416 instructions,
2417 transcript: transcript.clone(),
2418 tools: request_tools,
2419 tool_choice: request_tool_choice,
2420 reasoning: request_reasoning,
2421 output: OutputConfig::default(),
2422 runtime: RuntimeHints {
2423 auto_compact_token_limit: server_side_compaction_threshold(&cfg, &model),
2424 profile: runtime_profile,
2425 parallel_tool_calls: Some(parallel_tool_calls),
2426 hosted_web_search: cfg.hosted_web_search.clone(),
2427 tool_search: tool_search_for_provider_model(&cfg, &provider, &model),
2428 speed_policy: speed_policy_decision,
2429 reliability: Some(cfg.reliability.clone().into()),
2430 deadline_remaining_seconds: deadline_remaining_seconds(turn_deadline),
2431 ..RuntimeHints::default()
2432 },
2433 metadata: request_metadata,
2434 };
2435
2436 let ctx = InferenceTurnContext {
2437 thread_id: &req.thread_id,
2438 turn_id: &turn_id,
2439 tool_executor: Some(std::sync::Arc::new(
2440 crate::tool_execution::RuntimeTurnToolExecutor {
2441 runtime: Arc::clone(self),
2442 thread_id: req.thread_id.clone(),
2443 turn_id: turn_id.clone(),
2444 workspace: Some(workspace.clone()),
2445 deadline: turn_deadline,
2446 },
2447 )),
2448 };
2449 let stream_future = engine.stream_turn(ctx, request);
2450 let mut stream = if let Some((deadline, timeout_action)) = inference_timeout_deadline(
2451 turn_deadline,
2452 runtime_profile,
2453 req.task_ledger_required,
2454 deadline_finalization_reserve,
2455 deadline_finalization_requested || deadline_scoreable_completion_requested,
2456 task_ledger_scoreable_checkpoints,
2457 &transcript,
2458 ) {
2459 match tokio::time::timeout_at(deadline_instant(deadline), stream_future).await {
2460 Ok(stream) => stream?,
2461 Err(_) => {
2462 if runtime_profile == RuntimeProfile::Eval
2463 && !deadline_finalization_requested
2464 {
2465 let remaining = deadline_remaining_seconds(turn_deadline).unwrap_or(0);
2466 if timeout_action == InferenceTimeoutAction::ScoreableCheckpoint
2467 && task_ledger_scoreable_checkpoints
2468 < TASK_LEDGER_SCOREABLE_CHECKPOINT_LIMIT
2469 && let Some(prompt) = task_ledger_completion_prompt(&transcript)
2470 {
2471 task_ledger_scoreable_checkpoints += 1;
2472 let item = TranscriptItem::UserMessage(UserMessage::text(
2473 task_ledger_scoreable_checkpoint_prompt(remaining, &prompt),
2474 ));
2475 self.persist_turn_item(&req.thread_id, &turn_id, &item)
2476 .await?;
2477 transcript.push(item);
2478 continue 'tool_rounds;
2479 }
2480 self.start_deadline_finalization(
2481 &req.thread_id,
2482 &turn_id,
2483 &mut transcript,
2484 remaining,
2485 )
2486 .await?;
2487 deadline_finalization_requested = true;
2488 continue 'tool_rounds;
2489 }
2490 self.fail_turn_due_to_deadline(
2491 &req.thread_id,
2492 &turn_id,
2493 deadline,
2494 &transcript,
2495 )
2496 .await?;
2497 return Ok(TurnRunOutcome::Stopped);
2498 }
2499 }
2500 } else {
2501 stream_future.await?
2502 };
2503 let mut assistant_text = String::new();
2504 let mut phase_messages = Vec::<AssistantMessage>::new();
2505 let mut reasoning_text = String::new();
2506 let mut tool_calls = Vec::new();
2507 let mut provider_metadata = None;
2508
2509 loop {
2510 let next = if let Some((deadline, timeout_action)) = inference_timeout_deadline(
2511 turn_deadline,
2512 runtime_profile,
2513 req.task_ledger_required,
2514 deadline_finalization_reserve,
2515 deadline_finalization_requested || deadline_scoreable_completion_requested,
2516 task_ledger_scoreable_checkpoints,
2517 &transcript,
2518 ) {
2519 match tokio::time::timeout_at(deadline_instant(deadline), stream.next()).await {
2520 Ok(next) => next,
2521 Err(_) => {
2522 if runtime_profile == RuntimeProfile::Eval
2523 && !deadline_finalization_requested
2524 {
2525 let remaining =
2526 deadline_remaining_seconds(turn_deadline).unwrap_or(0);
2527 if timeout_action == InferenceTimeoutAction::ScoreableCheckpoint
2528 && task_ledger_scoreable_checkpoints
2529 < TASK_LEDGER_SCOREABLE_CHECKPOINT_LIMIT
2530 && let Some(prompt) = task_ledger_completion_prompt(&transcript)
2531 {
2532 task_ledger_scoreable_checkpoints += 1;
2533 let item = TranscriptItem::UserMessage(UserMessage::text(
2534 task_ledger_scoreable_checkpoint_prompt(remaining, &prompt),
2535 ));
2536 self.persist_turn_item(&req.thread_id, &turn_id, &item)
2537 .await?;
2538 transcript.push(item);
2539 continue 'tool_rounds;
2540 }
2541 self.start_deadline_finalization(
2542 &req.thread_id,
2543 &turn_id,
2544 &mut transcript,
2545 remaining,
2546 )
2547 .await?;
2548 deadline_finalization_requested = true;
2549 continue 'tool_rounds;
2550 }
2551 self.fail_turn_due_to_deadline(
2552 &req.thread_id,
2553 &turn_id,
2554 deadline,
2555 &transcript,
2556 )
2557 .await?;
2558 return Ok(TurnRunOutcome::Stopped);
2559 }
2560 }
2561 } else {
2562 stream.next().await
2563 };
2564 let Some(res) = next else {
2565 break;
2566 };
2567 let event = match res {
2568 Ok(event) => event,
2569 Err(err) => {
2570 let error = err.to_string();
2571 if runtime_profile == RuntimeProfile::Eval
2572 && !deadline_finalization_requested
2573 && let Some(cause) = provider_stream_retry_cause(&error)
2574 {
2575 let retry_attempt = provider_stream_retry_attempts.saturating_add(1);
2576 let policy: ReliabilityRequestPolicy = cfg.reliability.clone().into();
2577 if retry_attempt < policy.provider_retry_max_attempts {
2578 provider_stream_retry_attempts = retry_attempt;
2579 let delay_ms = provider_retry_delay_ms(&policy, retry_attempt);
2580 self.emit(RoderEvent::ReliabilityRetryRecorded(
2581 ReliabilityRetryRecorded {
2582 context: ReliabilityContext {
2583 thread_id: req.thread_id.clone(),
2584 turn_id: turn_id.clone(),
2585 provider: Some(provider.clone()),
2586 model: Some(model.clone()),
2587 ..ReliabilityContext::default()
2588 },
2589 error_class: ReliabilityErrorClass::ProviderError,
2590 decision: ReliabilityRetryDecision::Retry,
2591 attempt: retry_attempt,
2592 max_attempts: policy.provider_retry_max_attempts,
2593 delay_ms: Some(delay_ms),
2594 details: ReliabilityDetails::redacted(format!(
2595 "{cause}: {error}"
2596 )),
2597 timestamp: OffsetDateTime::now_utc(),
2598 },
2599 ))
2600 .await;
2601 if delay_ms > 0 {
2602 tokio::time::sleep(std::time::Duration::from_millis(delay_ms))
2603 .await;
2604 }
2605 continue 'tool_rounds;
2606 }
2607 }
2608 self.emit(RoderEvent::TurnFailed(TurnFailed {
2609 thread_id: req.thread_id.clone(),
2610 turn_id: turn_id.clone(),
2611 error,
2612 error_kind: None,
2613 usage: None,
2614 timestamp: OffsetDateTime::now_utc(),
2615 }))
2616 .await;
2617 self.complete_team_member_turn(
2618 &req.thread_id,
2619 &turn_id,
2620 TeamMemberStatus::Failed,
2621 )
2622 .await?;
2623 return Err(err);
2624 }
2625 };
2626
2627 let inference_timestamp = OffsetDateTime::now_utc();
2628 self.emit(RoderEvent::InferenceEventReceived(InferenceEventReceived {
2629 thread_id: req.thread_id.clone(),
2630 turn_id: turn_id.clone(),
2631 event: event.clone(),
2632 timestamp: inference_timestamp,
2633 }))
2634 .await;
2635
2636 match event {
2637 InferenceEvent::MessageDelta(delta) => {
2638 if let Some((team_id, member)) =
2639 self.teams.member_for_thread(&req.thread_id).await
2640 {
2641 self.emit(RoderEvent::TeamMemberMessageDelta(TeamMemberMessageDelta {
2642 team_id,
2643 member_id: member.id,
2644 member_thread_id: req.thread_id.clone(),
2645 turn_id: turn_id.clone(),
2646 delta: delta.text.clone(),
2647 timestamp: OffsetDateTime::now_utc(),
2648 }))
2649 .await;
2650 }
2651 if is_final_answer_phase(delta.phase.as_deref()) {
2652 assistant_text.push_str(&delta.text);
2653 } else if let Some(last) = phase_messages.last_mut()
2654 && last.phase == delta.phase
2655 {
2656 last.text.push_str(&delta.text);
2657 } else {
2658 phase_messages.push(AssistantMessage {
2659 text: delta.text,
2660 phase: delta.phase,
2661 });
2662 }
2663 }
2664 InferenceEvent::ReasoningDelta(delta) => reasoning_text.push_str(&delta.text),
2665 InferenceEvent::ToolCallCompleted(call) => tool_calls.push(call),
2666 InferenceEvent::Failed(failure) => {
2667 speed_policy.record_failure();
2668 self.persist_turn_item(
2669 &req.thread_id,
2670 &turn_id,
2671 &TranscriptItem::Error(ErrorRecord {
2672 message: failure.message.clone(),
2673 }),
2674 )
2675 .await?;
2676 self.emit(RoderEvent::TurnFailed(TurnFailed {
2677 thread_id: req.thread_id.clone(),
2678 turn_id: turn_id.clone(),
2679 error: failure.message,
2680 error_kind: None,
2681 usage: None,
2682 timestamp: OffsetDateTime::now_utc(),
2683 }))
2684 .await;
2685 self.complete_team_member_turn(
2686 &req.thread_id,
2687 &turn_id,
2688 TeamMemberStatus::Failed,
2689 )
2690 .await?;
2691 return Ok(TurnRunOutcome::Stopped);
2692 }
2693 InferenceEvent::Usage(usage) => {
2694 turn_usage.add_assign(&usage);
2695 }
2696 InferenceEvent::Completed(metadata) => {
2697 turn_finish_reason = metadata
2698 .stop_reason
2699 .as_deref()
2700 .map(finish_reason_from_stop_reason);
2701 }
2702 InferenceEvent::Compaction(_)
2703 | InferenceEvent::HostedToolCallStarted(_)
2704 | InferenceEvent::HostedToolCallCompleted(_)
2705 | InferenceEvent::ToolCallStarted(_)
2706 | InferenceEvent::ToolCallDelta(_) => {}
2707 InferenceEvent::ProviderMetadata(metadata) => {
2708 provider_metadata = Some(metadata);
2709 }
2710 }
2711 }
2712
2713 speed_policy.record_model_output(
2714 !assistant_text.is_empty() || !phase_messages.is_empty(),
2715 tool_calls.len(),
2716 );
2717 if tool_calls.is_empty() {
2718 let steers = self.drain_turn_steers(&turn_id).await;
2719 if !steers.is_empty() {
2720 for message in phase_messages {
2721 let item = TranscriptItem::AssistantMessage(message);
2722 self.persist_turn_item(&req.thread_id, &turn_id, &item)
2723 .await?;
2724 transcript.push(item);
2725 self.persist_model_profile_segment(
2726 &req.thread_id,
2727 &turn_id,
2728 model_profile.as_ref(),
2729 &provider,
2730 &model,
2731 "assistant",
2732 )
2733 .await?;
2734 }
2735 if !assistant_text.is_empty() {
2736 let assistant = TranscriptItem::AssistantMessage(AssistantMessage {
2737 text: assistant_text,
2738 phase: Some(FINAL_ANSWER_PHASE.to_string()),
2739 });
2740 self.persist_turn_item(&req.thread_id, &turn_id, &assistant)
2741 .await?;
2742 transcript.push(assistant);
2743 self.persist_model_profile_segment(
2744 &req.thread_id,
2745 &turn_id,
2746 model_profile.as_ref(),
2747 &provider,
2748 &model,
2749 "assistant",
2750 )
2751 .await?;
2752 }
2753 if let Some(metadata) = provider_metadata {
2754 let item = TranscriptItem::ProviderMetadata(metadata);
2755 self.persist_turn_item(&req.thread_id, &turn_id, &item)
2756 .await?;
2757 transcript.push(item);
2758 }
2759 self.append_steers(&req, &turn_id, &mut transcript, steers)
2760 .await?;
2761 continue;
2762 }
2763 if !deadline_finalization_requested
2764 && req.task_ledger_required
2765 && runtime_profile == RuntimeProfile::Eval
2766 && task_ledger_completion_reminders < TASK_LEDGER_COMPLETION_REMINDER_LIMIT
2767 && (!assistant_text.trim().is_empty() || !phase_messages.is_empty())
2768 && let Some(prompt) = task_ledger_completion_prompt(&transcript)
2769 {
2770 task_ledger_completion_reminders += 1;
2771 let item = TranscriptItem::UserMessage(UserMessage::text(prompt));
2772 self.persist_turn_item(&req.thread_id, &turn_id, &item)
2773 .await?;
2774 transcript.push(item);
2775 continue;
2776 }
2777 if !deadline_finalization_requested
2778 && let Some(prompt) = verification_gate.blocking_prompt()
2779 {
2780 speed_policy.record_verification_required();
2781 self.emit(RoderEvent::VerificationRequired(VerificationRequired {
2782 thread_id: req.thread_id.clone(),
2783 turn_id: turn_id.clone(),
2784 reason: verification_gate.reason(),
2785 changed_files: verification_gate.changed_files(),
2786 tool_evidence: verification_gate.tool_evidence.clone(),
2787 tests_run: verification_gate.tests_run.clone(),
2788 open_gaps: verification_gate.open_gaps.clone(),
2789 timestamp: OffsetDateTime::now_utc(),
2790 }))
2791 .await;
2792 let item = TranscriptItem::UserMessage(UserMessage::text(prompt));
2793 self.persist_turn_item(&req.thread_id, &turn_id, &item)
2794 .await?;
2795 transcript.push(item);
2796 continue;
2797 }
2798 if deadline_finalization_requested
2799 && assistant_text.trim().is_empty()
2800 && phase_messages.is_empty()
2801 {
2802 assistant_text = format!(
2803 "Deadline finalization completed without model text. {}",
2804 turn_partial_result(&transcript)
2805 );
2806 }
2807 final_phase_messages = phase_messages;
2808 final_assistant_text = assistant_text;
2809 final_reasoning_text = reasoning_text;
2810 final_provider_metadata = provider_metadata;
2811 exhausted_tool_rounds = false;
2812 break;
2813 }
2814
2815 for message in phase_messages {
2816 let item = TranscriptItem::AssistantMessage(message);
2817 self.persist_turn_item(&req.thread_id, &turn_id, &item)
2818 .await?;
2819 transcript.push(item);
2820 self.persist_model_profile_segment(
2821 &req.thread_id,
2822 &turn_id,
2823 model_profile.as_ref(),
2824 &provider,
2825 &model,
2826 "assistant",
2827 )
2828 .await?;
2829 }
2830 if !assistant_text.is_empty() {
2831 transcript.push(TranscriptItem::AssistantMessage(AssistantMessage {
2832 text: assistant_text,
2833 phase: Some(FINAL_ANSWER_PHASE.to_string()),
2834 }));
2835 self.persist_model_profile_segment(
2836 &req.thread_id,
2837 &turn_id,
2838 model_profile.as_ref(),
2839 &provider,
2840 &model,
2841 "assistant",
2842 )
2843 .await?;
2844 }
2845 if let Some(metadata) = provider_metadata {
2846 let item = TranscriptItem::ProviderMetadata(metadata);
2847 self.persist_turn_item(&req.thread_id, &turn_id, &item)
2848 .await?;
2849 transcript.push(item);
2850 }
2851 for call in &tool_calls {
2852 let tool_item = TranscriptItem::ToolCall(ToolCallRecord {
2853 id: call.id.clone(),
2854 name: call.name.clone(),
2855 arguments: call.arguments.clone(),
2856 });
2857 self.persist_turn_item(&req.thread_id, &turn_id, &tool_item)
2858 .await?;
2859 transcript.push(tool_item);
2860 self.persist_model_profile_segment(
2861 &req.thread_id,
2862 &turn_id,
2863 model_profile.as_ref(),
2864 &provider,
2865 &model,
2866 "tool_call",
2867 )
2868 .await?;
2869 }
2870 if let Some(deadline) = turn_deadline
2871 && deadline_expired(deadline)
2872 {
2873 self.fail_turn_due_to_deadline(&req.thread_id, &turn_id, deadline, &transcript)
2874 .await?;
2875 return Ok(TurnRunOutcome::Stopped);
2876 }
2877 let results = self
2878 .route_tool_calls(
2879 &req.thread_id,
2880 &turn_id,
2881 tool_calls,
2882 parallel_tool_calls,
2883 Some(workspace.as_str()),
2884 turn_deadline,
2885 )
2886 .await?;
2887 if let Some(limit) = reliability.record_tool_results(
2888 &cfg.reliability,
2889 &results,
2890 runtime_profile == RuntimeProfile::Interactive,
2891 ) {
2892 self.fail_turn_due_to_reliability_limit(
2893 &req.thread_id,
2894 &turn_id,
2895 &provider,
2896 &model,
2897 limit,
2898 &transcript,
2899 )
2900 .await?;
2901 return Ok(TurnRunOutcome::Stopped);
2902 }
2903 for result in results {
2904 verification_gate.record_tool_result(&result);
2905 transcript.push(TranscriptItem::ToolResult(result));
2906 self.persist_model_profile_segment(
2907 &req.thread_id,
2908 &turn_id,
2909 model_profile.as_ref(),
2910 &provider,
2911 &model,
2912 "tool_result",
2913 )
2914 .await?;
2915 }
2916 transcript = self
2917 .compact_transcript_if_needed(
2918 &req.thread_id,
2919 &turn_id,
2920 &provider,
2921 &model,
2922 transcript,
2923 self.compaction_options_for_turn(&req.thread_id, !compacted_this_turn),
2924 )
2925 .await?;
2926 compacted_this_turn = compacted_this_turn
2927 || transcript
2928 .iter()
2929 .any(|item| matches!(item, TranscriptItem::ContextCompaction(_)));
2930 }
2931
2932 if exhausted_tool_rounds {
2933 let message =
2934 format!("tool call limit reached after {MAX_TOOL_ROUNDS_PER_TURN} rounds");
2935 self.persist_turn_item(
2936 &req.thread_id,
2937 &turn_id,
2938 &TranscriptItem::Error(ErrorRecord {
2939 message: message.clone(),
2940 }),
2941 )
2942 .await?;
2943 self.emit(RoderEvent::TurnFailed(TurnFailed {
2944 thread_id: req.thread_id.clone(),
2945 turn_id: turn_id.clone(),
2946 error: message,
2947 error_kind: None,
2948 usage: None,
2949 timestamp: OffsetDateTime::now_utc(),
2950 }))
2951 .await;
2952 self.complete_team_member_turn(&req.thread_id, &turn_id, TeamMemberStatus::Failed)
2953 .await?;
2954 return Ok(TurnRunOutcome::Stopped);
2955 }
2956
2957 if !final_reasoning_text.is_empty() {
2958 self.persist_turn_item(
2959 &req.thread_id,
2960 &turn_id,
2961 &TranscriptItem::ReasoningSummary(ReasoningSummary {
2962 text: final_reasoning_text,
2963 }),
2964 )
2965 .await?;
2966 }
2967 for message in final_phase_messages {
2968 self.persist_turn_item(
2969 &req.thread_id,
2970 &turn_id,
2971 &TranscriptItem::AssistantMessage(message),
2972 )
2973 .await?;
2974 self.persist_model_profile_segment(
2975 &req.thread_id,
2976 &turn_id,
2977 model_profile.as_ref(),
2978 &provider,
2979 &model,
2980 "assistant",
2981 )
2982 .await?;
2983 }
2984 if !final_assistant_text.is_empty() {
2985 self.persist_turn_item(
2986 &req.thread_id,
2987 &turn_id,
2988 &TranscriptItem::AssistantMessage(AssistantMessage {
2989 text: final_assistant_text,
2990 phase: Some(FINAL_ANSWER_PHASE.to_string()),
2991 }),
2992 )
2993 .await?;
2994 self.persist_model_profile_segment(
2995 &req.thread_id,
2996 &turn_id,
2997 model_profile.as_ref(),
2998 &provider,
2999 &model,
3000 "assistant",
3001 )
3002 .await?;
3003 }
3004 if let Some(metadata) = final_provider_metadata {
3005 self.persist_turn_item(
3006 &req.thread_id,
3007 &turn_id,
3008 &TranscriptItem::ProviderMetadata(metadata),
3009 )
3010 .await?;
3011 }
3012
3013 let turn_usage_tokens = turn_usage.total_tokens as i64;
3014 let completed_usage = (!turn_usage.is_empty()).then_some(turn_usage.clone());
3015 self.record_thread_usage_metadata(&req.thread_id, &turn_usage)
3016 .await?;
3017 self.goals
3018 .account_turn_usage(
3019 &req.thread_id,
3020 turn_usage_tokens,
3021 OffsetDateTime::now_utc() - turn_started_at,
3022 )
3023 .await?;
3024 self.emit(RoderEvent::TurnCompleted(TurnCompleted {
3025 thread_id: req.thread_id.clone(),
3026 turn_id: turn_id.clone(),
3027 usage: completed_usage,
3028 finish_reason: turn_finish_reason,
3029 timestamp: OffsetDateTime::now_utc(),
3030 }))
3031 .await;
3032 self.complete_team_member_turn(&req.thread_id, &turn_id, TeamMemberStatus::Completed)
3033 .await?;
3034 self.persist_runner_state(&req.thread_id, runner_session.as_ref())
3035 .await?;
3036 Ok(TurnRunOutcome::Completed)
3037 }
3038
3039 async fn drain_turn_steers(&self, turn_id: &TurnId) -> Vec<UserMessage> {
3040 let Some(active) = self.active_turns.read().await.get(turn_id).cloned() else {
3041 return Vec::new();
3042 };
3043 let mut steers = active.steers.lock().await;
3044 std::mem::take(&mut *steers)
3045 }
3046
3047 async fn route_tool_calls(
3048 self: &Arc<Self>,
3049 thread_id: &ThreadId,
3050 turn_id: &TurnId,
3051 calls: Vec<ToolCallCompleted>,
3052 parallel: bool,
3053 workspace: Option<&str>,
3054 deadline: Option<OffsetDateTime>,
3055 ) -> anyhow::Result<Vec<ToolResultRecord>> {
3056 let force_sequential = calls
3057 .iter()
3058 .any(|call| crate::agent_control_tools::is_agent_control_tool(&call.name));
3059 if parallel && !force_sequential {
3060 try_join_all(
3061 calls.into_iter().map(|call| {
3062 self.route_tool_call(thread_id, turn_id, call, workspace, deadline)
3063 }),
3064 )
3065 .await
3066 } else {
3067 let mut results = Vec::with_capacity(calls.len());
3068 for call in calls {
3069 results.push(
3070 self.route_tool_call(thread_id, turn_id, call, workspace, deadline)
3071 .await?,
3072 );
3073 }
3074 Ok(results)
3075 }
3076 }
3077
3078 async fn fail_turn_with_error(
3079 &self,
3080 thread_id: &ThreadId,
3081 turn_id: &TurnId,
3082 message: String,
3083 ) -> anyhow::Result<()> {
3084 self.persist_turn_item(
3085 thread_id,
3086 turn_id,
3087 &TranscriptItem::Error(ErrorRecord {
3088 message: message.clone(),
3089 }),
3090 )
3091 .await?;
3092 self.emit(RoderEvent::TurnFailed(TurnFailed {
3093 thread_id: thread_id.clone(),
3094 turn_id: turn_id.clone(),
3095 error: message,
3096 error_kind: None,
3097 usage: None,
3098 timestamp: OffsetDateTime::now_utc(),
3099 }))
3100 .await;
3101 self.complete_team_member_turn(thread_id, turn_id, TeamMemberStatus::Failed)
3102 .await?;
3103 Ok(())
3104 }
3105
3106 async fn fail_turn_due_to_deadline(
3107 &self,
3108 thread_id: &ThreadId,
3109 turn_id: &TurnId,
3110 deadline: OffsetDateTime,
3111 transcript: &[TranscriptItem],
3112 ) -> anyhow::Result<()> {
3113 let partial_result = turn_partial_result(transcript);
3114 self.emit(RoderEvent::TurnPartialResult(TurnPartialResult {
3115 thread_id: thread_id.clone(),
3116 turn_id: turn_id.clone(),
3117 summary: partial_result.clone(),
3118 timestamp: OffsetDateTime::now_utc(),
3119 }))
3120 .await;
3121 self.emit(RoderEvent::TurnDeadlineExceeded(TurnDeadlineExceeded {
3122 thread_id: thread_id.clone(),
3123 turn_id: turn_id.clone(),
3124 deadline,
3125 partial_result: partial_result.clone(),
3126 timestamp: OffsetDateTime::now_utc(),
3127 }))
3128 .await;
3129 let message = "turn deadline expired".to_string();
3130 self.persist_turn_item(
3131 thread_id,
3132 turn_id,
3133 &TranscriptItem::Error(ErrorRecord {
3134 message: format!("{message}: {partial_result}"),
3135 }),
3136 )
3137 .await?;
3138 self.emit(RoderEvent::TurnFailed(TurnFailed {
3139 thread_id: thread_id.clone(),
3140 turn_id: turn_id.clone(),
3141 error: message,
3142 error_kind: Some("deadline_timeout".to_string()),
3143 usage: None,
3144 timestamp: OffsetDateTime::now_utc(),
3145 }))
3146 .await;
3147 self.complete_team_member_turn(thread_id, turn_id, TeamMemberStatus::Failed)
3148 .await?;
3149 Ok(())
3150 }
3151
3152 async fn start_deadline_finalization(
3153 &self,
3154 thread_id: &ThreadId,
3155 turn_id: &TurnId,
3156 transcript: &mut Vec<TranscriptItem>,
3157 remaining_seconds: u64,
3158 ) -> anyhow::Result<()> {
3159 let item = TranscriptItem::UserMessage(crate::deadline_policy::finalization_message(
3160 remaining_seconds,
3161 ));
3162 self.persist_turn_item(thread_id, turn_id, &item).await?;
3163 transcript.push(item);
3164 self.emit(RoderEvent::TurnPartialResult(TurnPartialResult {
3165 thread_id: thread_id.clone(),
3166 turn_id: turn_id.clone(),
3167 summary: turn_partial_result(transcript),
3168 timestamp: OffsetDateTime::now_utc(),
3169 }))
3170 .await;
3171 Ok(())
3172 }
3173
3174 async fn fail_turn_due_to_reliability_limit(
3175 &self,
3176 thread_id: &ThreadId,
3177 turn_id: &TurnId,
3178 provider: &str,
3179 model: &str,
3180 limit: ReliabilityLimitHit,
3181 transcript: &[TranscriptItem],
3182 ) -> anyhow::Result<()> {
3183 self.emit(RoderEvent::ReliabilityLimitRecorded(
3184 ReliabilityLimitRecorded {
3185 context: ReliabilityContext {
3186 thread_id: thread_id.clone(),
3187 turn_id: turn_id.clone(),
3188 tool_id: None,
3189 tool_name: None,
3190 provider: Some(provider.to_string()),
3191 model: Some(model.to_string()),
3192 },
3193 error_class: limit.error_class,
3194 limit_kind: limit.limit_kind,
3195 decision: limit.decision,
3196 current: limit.current,
3197 limit: limit.limit,
3198 details: ReliabilityDetails::redacted(&limit.message),
3199 timestamp: OffsetDateTime::now_utc(),
3200 },
3201 ))
3202 .await;
3203 let partial_result = turn_partial_result(transcript);
3204 self.emit(RoderEvent::TurnPartialResult(TurnPartialResult {
3205 thread_id: thread_id.clone(),
3206 turn_id: turn_id.clone(),
3207 summary: partial_result.clone(),
3208 timestamp: OffsetDateTime::now_utc(),
3209 }))
3210 .await;
3211 let message = format!("reliability limit reached: {}", limit.message);
3212 self.persist_turn_item(
3213 thread_id,
3214 turn_id,
3215 &TranscriptItem::Error(ErrorRecord {
3216 message: format!("{message}: {partial_result}"),
3217 }),
3218 )
3219 .await?;
3220 self.emit(RoderEvent::TurnFailed(TurnFailed {
3221 thread_id: thread_id.clone(),
3222 turn_id: turn_id.clone(),
3223 error: message,
3224 error_kind: Some("reliability_limit".to_string()),
3225 usage: None,
3226 timestamp: OffsetDateTime::now_utc(),
3227 }))
3228 .await;
3229 self.complete_team_member_turn(thread_id, turn_id, TeamMemberStatus::Failed)
3230 .await?;
3231 Ok(())
3232 }
3233
3234 async fn append_steers(
3235 &self,
3236 req: &StartTurnRequest,
3237 turn_id: &TurnId,
3238 transcript: &mut Vec<TranscriptItem>,
3239 steers: Vec<UserMessage>,
3240 ) -> anyhow::Result<()> {
3241 for mut steer in steers {
3242 steer.text = steer.text.trim().to_string();
3243 if steer.text.is_empty() && steer.images.is_empty() {
3244 continue;
3245 }
3246 let item = TranscriptItem::UserMessage(steer);
3247 self.persist_turn_item(&req.thread_id, turn_id, &item)
3248 .await?;
3249 transcript.push(item);
3250 }
3251 Ok(())
3252 }
3253
3254 async fn persist_model_profile_segment(
3255 &self,
3256 thread_id: &ThreadId,
3257 turn_id: &TurnId,
3258 profile: Option<&ModelHarnessProfile>,
3259 provider: &str,
3260 model: &str,
3261 segment: &str,
3262 ) -> anyhow::Result<()> {
3263 let item = TranscriptItem::ProviderMetadata(model_profile_segment_metadata(
3264 profile, provider, model, segment,
3265 ));
3266 self.persist_turn_item(thread_id, turn_id, &item).await
3267 }
3268
3269 fn filtered_tool_specs(
3275 &self,
3276 cfg: &RuntimeConfig,
3277 model: &str,
3278 profile: Option<&ModelHarnessProfile>,
3279 thread_allowlist: &[String],
3280 external_tools: &[roder_api::tools::ToolSpec],
3281 ) -> Vec<roder_api::tools::ToolSpec> {
3282 let mut specs = self
3283 .tool_registry
3284 .specs_for_edit_tool_with_schema_policy(
3285 edit_tool_for_model(cfg, model),
3286 schema_policy_for_model(profile),
3287 )
3288 .into_iter()
3289 .filter(|spec| {
3290 allowlist_permits(&cfg.tool_allowlist, &spec.name)
3291 && allowlist_permits(thread_allowlist, &spec.name)
3292 && !external_tools.iter().any(|tool| tool.name == spec.name)
3293 })
3294 .collect::<Vec<_>>();
3295 specs.extend(external_tools.iter().cloned());
3296 specs
3297 }
3298
3299 fn task_ledger_tool_specs(
3300 &self,
3301 profile: Option<&ModelHarnessProfile>,
3302 ) -> Vec<roder_api::tools::ToolSpec> {
3303 self.tool_registry
3304 .get(TASK_LEDGER_TOOL_NAME)
3305 .map(|tool| {
3306 tool.spec()
3307 .normalized_for_model_profile(schema_policy_for_model(profile))
3308 })
3309 .into_iter()
3310 .collect()
3311 }
3312
3313 pub(crate) fn engine_for(&self, provider: &str) -> anyhow::Result<Arc<dyn InferenceEngine>> {
3314 self.registry
3315 .inference_engine(provider)
3316 .or_else(|| {
3317 self.registry
3318 .default_inference_engine()
3319 .filter(|engine| provider.is_empty() || engine.id() == provider)
3320 })
3321 .ok_or_else(|| anyhow::anyhow!("inference provider {provider:?} is not registered"))
3322 }
3323
3324 pub async fn emit(&self, event: RoderEvent) -> EventEnvelope {
3325 let envelope = self.bus.emit(event);
3326 if let (Some(store), Some(thread_id)) = (&self.thread_store, envelope.thread_id.as_ref())
3327 && should_persist_thread_event(thread_id)
3328 {
3329 let _ = store.append_event(thread_id, &envelope).await;
3330 }
3331 let dispatcher = self
3335 .event_sink_dispatcher
3336 .get_or_init(|| async {
3337 crate::event_sink_dispatch::EventSinkDispatcher::start(
3338 &self.registry.event_sinks,
3339 self.bus.clone(),
3340 )
3341 })
3342 .await;
3343 if !dispatcher.is_empty() {
3344 dispatcher.dispatch(&envelope, &self.bus);
3345 }
3346 envelope
3347 }
3348
3349 pub async fn record_thread_item_event_kind(
3355 &self,
3356 thread_id: &ThreadId,
3357 turn_id: &TurnId,
3358 timestamp: OffsetDateTime,
3359 kind: ThreadItemEventKind,
3360 ) -> anyhow::Result<ThreadItemEvent> {
3361 let seq = self.next_thread_item_event_seq(thread_id).await?;
3362 let item_event = ThreadItemEvent {
3363 seq,
3364 event_id: format!("{turn_id}-item-event-{seq}"),
3365 thread_id: thread_id.clone(),
3366 turn_id: turn_id.clone(),
3367 timestamp,
3368 event: kind,
3369 };
3370 if let Some(store) = &self.thread_store {
3371 store.append_item_event(thread_id, &item_event).await?;
3372 }
3373 self.remember_thread_item_event(&item_event).await?;
3374 Ok(item_event)
3375 }
3376
3377 async fn next_thread_item_event_seq(&self, thread_id: &ThreadId) -> anyhow::Result<u64> {
3378 self.ensure_thread_item_cache(thread_id).await?;
3379 Ok(self
3380 .thread_item_cache
3381 .lock()
3382 .await
3383 .next_item_event_seq(thread_id))
3384 }
3385
3386 pub async fn thread_item_exists(
3387 &self,
3388 thread_id: &ThreadId,
3389 turn_id: &TurnId,
3390 item_id: &str,
3391 ) -> anyhow::Result<bool> {
3392 self.ensure_thread_item_cache(thread_id).await?;
3393 Ok(self
3394 .thread_item_cache
3395 .lock()
3396 .await
3397 .thread_item_exists(thread_id, turn_id, item_id))
3398 }
3399
3400 pub async fn current_reasoning_item_id(
3401 &self,
3402 thread_id: &ThreadId,
3403 turn_id: &TurnId,
3404 ) -> anyhow::Result<Option<String>> {
3405 self.ensure_thread_item_cache(thread_id).await?;
3406 Ok(self
3407 .thread_item_cache
3408 .lock()
3409 .await
3410 .current_reasoning_item_id(thread_id, turn_id))
3411 }
3412
3413 async fn remember_thread_item_event(&self, item_event: &ThreadItemEvent) -> anyhow::Result<()> {
3414 self.ensure_thread_item_cache(&item_event.thread_id).await?;
3415 self.thread_item_cache
3416 .lock()
3417 .await
3418 .remember_item_event(item_event);
3419 Ok(())
3420 }
3421
3422 pub async fn latest_transcript_item_index(
3423 &self,
3424 thread_id: &ThreadId,
3425 turn_id: &TurnId,
3426 ) -> anyhow::Result<Option<usize>> {
3427 self.ensure_thread_item_cache(thread_id).await?;
3428 Ok(self
3429 .thread_item_cache
3430 .lock()
3431 .await
3432 .latest_transcript_item_index(thread_id, turn_id))
3433 }
3434
3435 async fn next_transcript_item_index(
3436 &self,
3437 thread_id: &ThreadId,
3438 turn_id: &TurnId,
3439 ) -> anyhow::Result<usize> {
3440 self.ensure_thread_item_cache(thread_id).await?;
3441 Ok(self
3442 .thread_item_cache
3443 .lock()
3444 .await
3445 .next_transcript_item_index(thread_id, turn_id))
3446 }
3447
3448 async fn remember_transcript_item_index(
3449 &self,
3450 thread_id: &ThreadId,
3451 turn_id: &TurnId,
3452 item_index: usize,
3453 ) -> anyhow::Result<()> {
3454 self.ensure_thread_item_cache(thread_id).await?;
3455 self.thread_item_cache
3456 .lock()
3457 .await
3458 .remember_transcript_item_index(thread_id, turn_id, item_index);
3459 Ok(())
3460 }
3461
3462 async fn ensure_thread_item_cache(&self, thread_id: &ThreadId) -> anyhow::Result<()> {
3463 if self
3464 .thread_item_cache
3465 .lock()
3466 .await
3467 .contains_thread(thread_id)
3468 {
3469 return Ok(());
3470 }
3471
3472 let snapshot = if let Some(store) = &self.thread_store {
3473 store.load_thread(thread_id).await?
3474 } else {
3475 None
3476 };
3477 self.thread_item_cache.lock().await.ensure_thread(
3478 thread_id,
3479 ThreadItemCacheEntry::from_snapshot(snapshot.as_ref()),
3480 );
3481 Ok(())
3482 }
3483
3484 pub(crate) async fn persist_turn_item(
3485 &self,
3486 thread_id: &ThreadId,
3487 turn_id: &TurnId,
3488 item: &TranscriptItem,
3489 ) -> anyhow::Result<()> {
3490 let item_index = self.next_transcript_item_index(thread_id, turn_id).await?;
3491 let timestamp = OffsetDateTime::now_utc();
3492 self.emit(RoderEvent::TranscriptItemAppended(TranscriptItemAppended {
3493 thread_id: thread_id.clone(),
3494 turn_id: turn_id.clone(),
3495 item_type: match item {
3496 TranscriptItem::UserMessage(_) => "user_message",
3497 TranscriptItem::AssistantMessage(_) => "assistant_message",
3498 TranscriptItem::ReasoningSummary(_) => "reasoning_summary",
3499 TranscriptItem::ToolCall(_) => "tool_call",
3500 TranscriptItem::ToolResult(_) => "tool_result",
3501 TranscriptItem::FileChange(_) => "file_change",
3502 TranscriptItem::ContextCompaction(_) => "context_compaction",
3503 TranscriptItem::Error(_) => "error",
3504 TranscriptItem::ProviderMetadata(_) => "provider_metadata",
3505 }
3506 .to_string(),
3507 item_index: Some(item_index),
3508 item: Some(item.clone()),
3509 timestamp,
3510 }))
3511 .await;
3512 self.remember_transcript_item_index(thread_id, turn_id, item_index)
3513 .await?;
3514 Ok(())
3515 }
3516}
3517
3518fn transcript_has_images(transcript: &[TranscriptItem]) -> bool {
3519 transcript.iter().any(|item| {
3520 matches!(
3521 item,
3522 TranscriptItem::UserMessage(message) if !message.images.is_empty()
3523 )
3524 })
3525}
3526
3527fn transcript_has_task_ledger(transcript: &[TranscriptItem]) -> bool {
3528 transcript.iter().any(|item| {
3529 matches!(
3530 item,
3531 TranscriptItem::ToolResult(result)
3532 if result.name.as_deref() == Some(TASK_LEDGER_TOOL_NAME) && !result.is_error
3533 )
3534 })
3535}
3536
3537fn task_ledger_completion_prompt(transcript: &[TranscriptItem]) -> Option<String> {
3538 let latest = transcript.iter().rev().find_map(|item| match item {
3539 TranscriptItem::ToolResult(result)
3540 if result.name.as_deref() == Some(TASK_LEDGER_TOOL_NAME) && !result.is_error =>
3541 {
3542 Some(result.result.as_str())
3543 }
3544 _ => None,
3545 })?;
3546 if !task_ledger_has_open_items(latest) {
3547 return None;
3548 }
3549
3550 let mut ledger = latest.chars().take(1500).collect::<String>();
3551 if latest.chars().nth(1500).is_some() {
3552 ledger.push_str("...");
3553 }
3554 Some(format!(
3555 "Task Ledger Completion Required: the latest task ledger still has pending or in-progress items. Do not provide a final answer yet. Use tools to complete the remaining scoreable work, create or update any required output files, then call `{TASK_LEDGER_TOOL_NAME}` with every task completed and evidence before finalizing.\n\nLatest ledger:\n{ledger}"
3556 ))
3557}
3558
3559fn task_ledger_deadline_completion_prompt(
3560 remaining_seconds: u64,
3561 reserve_seconds: u64,
3562 completion_prompt: &str,
3563) -> String {
3564 format!(
3565 "Eval deadline scoreable completion: {remaining_seconds} seconds remain in the {reserve_seconds}-second finalization reserve. Do not browse, search, or start slow work. Use the available tools now to create or update the required scoreable output files, run only a quick local check if needed, then update the task ledger to completed before finalizing.\n\n{completion_prompt}"
3566 )
3567}
3568
3569fn task_ledger_scoreable_checkpoint_prompt(
3570 remaining_seconds: u64,
3571 completion_prompt: &str,
3572) -> String {
3573 format!(
3574 "Scoreable Output Checkpoint: {remaining_seconds} seconds remain before the eval deadline. Before any further research, browsing, or long commands, use tools now to ensure the required output file(s) exist with the best evidence-backed answer, even if provisional. If a scoreable file already exists, read it and preserve that candidate unless you have stronger task-specific evidence for a replacement. Do not overwrite a plausible dated, historical, or local-evidence candidate with a current live-page, partial-coverage, or weaker guess merely to refresh the checkpoint. You may continue refining afterward, but do not apologize or finalize until the scoreable file exists and the task ledger is updated.\n\n{completion_prompt}"
3575 )
3576}
3577
3578fn task_ledger_has_open_items(ledger: &str) -> bool {
3579 ledger.lines().any(|line| {
3580 let line = line.trim_start();
3581 line.starts_with("- pending:") || line.starts_with("- in_progress:")
3582 })
3583}
3584
3585fn turn_deadline_for_config(cfg: &RuntimeConfig) -> Option<OffsetDateTime> {
3586 if !cfg.runtime_profile.is_non_interactive() {
3587 return None;
3588 }
3589 cfg.turn_deadline_seconds
3590 .filter(|seconds| *seconds > 0)
3591 .map(|seconds| OffsetDateTime::now_utc() + Duration::seconds(seconds as i64))
3592}
3593
3594fn deadline_expired(deadline: OffsetDateTime) -> bool {
3595 OffsetDateTime::now_utc() >= deadline
3596}
3597
3598pub(crate) fn deadline_remaining_seconds(deadline: Option<OffsetDateTime>) -> Option<u64> {
3599 let deadline = deadline?;
3600 if deadline <= OffsetDateTime::now_utc() {
3601 return Some(0);
3602 }
3603 Some(
3604 (deadline - OffsetDateTime::now_utc())
3605 .unsigned_abs()
3606 .as_secs()
3607 .max(1),
3608 )
3609}
3610
3611fn deadline_instant(deadline: OffsetDateTime) -> tokio::time::Instant {
3612 let now = OffsetDateTime::now_utc();
3613 if deadline <= now {
3614 return tokio::time::Instant::now();
3615 }
3616 tokio::time::Instant::now() + (deadline - now).unsigned_abs()
3617}
3618
3619fn inference_timeout_deadline(
3620 deadline: Option<OffsetDateTime>,
3621 runtime_profile: RuntimeProfile,
3622 task_ledger_required: bool,
3623 reserve_seconds: u64,
3624 finalization_requested: bool,
3625 task_ledger_scoreable_checkpoints: u8,
3626 transcript: &[TranscriptItem],
3627) -> Option<(OffsetDateTime, InferenceTimeoutAction)> {
3628 let deadline = deadline?;
3629 if runtime_profile == RuntimeProfile::Eval
3630 && task_ledger_required
3631 && !finalization_requested
3632 && task_ledger_scoreable_checkpoints < TASK_LEDGER_SCOREABLE_CHECKPOINT_LIMIT
3633 && TASK_LEDGER_SCOREABLE_CHECKPOINT_SECONDS > reserve_seconds
3634 && task_ledger_completion_prompt(transcript).is_some()
3635 {
3636 let checkpoint_deadline =
3637 deadline - Duration::seconds(TASK_LEDGER_SCOREABLE_CHECKPOINT_SECONDS as i64);
3638 if checkpoint_deadline > OffsetDateTime::now_utc() {
3639 return Some((
3640 checkpoint_deadline,
3641 InferenceTimeoutAction::ScoreableCheckpoint,
3642 ));
3643 }
3644 }
3645 if runtime_profile == RuntimeProfile::Eval && !finalization_requested {
3646 return Some((
3647 deadline - Duration::seconds(reserve_seconds as i64),
3648 InferenceTimeoutAction::Finalization,
3649 ));
3650 }
3651 Some((deadline, InferenceTimeoutAction::Finalization))
3652}
3653
3654fn turn_partial_result(transcript: &[TranscriptItem]) -> String {
3655 let tool_results = transcript
3656 .iter()
3657 .filter(|item| matches!(item, TranscriptItem::ToolResult(_)))
3658 .count();
3659 let assistant_messages = transcript
3660 .iter()
3661 .filter(|item| matches!(item, TranscriptItem::AssistantMessage(_)))
3662 .count();
3663 format!(
3664 "partial turn state: {} transcript items, {assistant_messages} assistant messages, {tool_results} tool results",
3665 transcript.len()
3666 )
3667}
3668
3669fn reasoning_for_model(cfg: &RuntimeConfig, model: &str) -> ReasoningConfig {
3670 let level = effective_reasoning_for_model(cfg, model);
3671 match level.as_str() {
3672 "" | REASONING_NONE => ReasoningConfig::default(),
3673 level => ReasoningConfig {
3674 enabled: true,
3675 level: Some(level.to_string()),
3676 },
3677 }
3678}
3679
3680fn server_side_compaction_threshold(cfg: &RuntimeConfig, model: &str) -> Option<u32> {
3681 let entry = lookup_model(model)?;
3682 if !entry.supports_compaction {
3683 return None;
3684 }
3685 cfg.auto_compact_token_limit
3686 .or_else(|| {
3687 model_profile_for_model(cfg, model).and_then(|profile| profile.auto_compact_token_limit)
3688 })
3689 .or(Some(entry.auto_compact_token_limit))
3690 .filter(|threshold| *threshold > 0)
3691}
3692
3693pub(crate) fn tool_search_for_provider_model(
3694 cfg: &RuntimeConfig,
3695 provider: &str,
3696 model: &str,
3697) -> ToolSearchConfig {
3698 let mut resolved = cfg.tool_search.clone();
3699 if let Some(provider_config) = cfg.provider_tool_search.get(provider) {
3700 provider_config.apply_to(&mut resolved);
3701 }
3702 if let Some(model_config) = cfg.model_tool_search.get(model) {
3703 model_config.apply_to(&mut resolved);
3704 }
3705 resolved
3706}
3707
3708fn parallel_tool_calls_for_model(cfg: &RuntimeConfig, model: &str) -> bool {
3709 cfg.model_parallel_tool_calls
3710 .get(model)
3711 .copied()
3712 .or_else(|| {
3713 model_profile_for_model(cfg, model).and_then(|profile| profile.parallel_tool_calls)
3714 })
3715 .unwrap_or(true)
3716}
3717
3718fn effective_reasoning_for_model(cfg: &RuntimeConfig, model: &str) -> String {
3719 let base_reasoning = default_effective_reasoning_for_model(cfg, model);
3720 if cfg.dynamic_workflows.effort_profile == DynamicWorkflowEffortProfile::Ultracode {
3721 return ultracode_reasoning_level_for_model(
3722 model,
3723 &cfg.speed_policy.ultracode_reasoning,
3724 &base_reasoning,
3725 );
3726 }
3727 base_reasoning
3728}
3729
3730fn default_effective_reasoning_for_model(cfg: &RuntimeConfig, model: &str) -> String {
3731 let Some(entry) = lookup_model(model) else {
3732 return cfg
3733 .reasoning
3734 .clone()
3735 .unwrap_or_else(|| REASONING_NONE.to_string());
3736 };
3737 if entry.supported_reasoning.is_empty() {
3738 return REASONING_NONE.to_string();
3739 }
3740 cfg.reasoning
3741 .as_deref()
3742 .filter(|reasoning| {
3743 entry
3744 .supported_reasoning
3745 .iter()
3746 .any(|option| option.effort == *reasoning)
3747 })
3748 .map(str::to_string)
3749 .or_else(|| {
3750 model_profile_for_model(cfg, model)
3751 .and_then(|profile| profile.reasoning.orientation)
3752 .filter(|reasoning| {
3753 entry
3754 .supported_reasoning
3755 .iter()
3756 .any(|option| option.effort == reasoning)
3757 })
3758 })
3759 .unwrap_or_else(|| entry.default_reasoning.to_string())
3760}
3761
3762fn validate_reasoning_effort(model: &str, effort: &str) -> anyhow::Result<()> {
3763 if effort == REASONING_NONE && !model_supports_reasoning(model, effort) {
3764 return Ok(());
3765 }
3766 let Some(entry) = lookup_model(model) else {
3767 return Ok(());
3768 };
3769 if entry
3770 .supported_reasoning
3771 .iter()
3772 .any(|option| option.effort == effort)
3773 {
3774 Ok(())
3775 } else {
3776 anyhow::bail!("model {model} does not support reasoning effort {effort}")
3777 }
3778}
3779
3780fn validate_runtime_config_reasoning(cfg: &RuntimeConfig) -> anyhow::Result<()> {
3781 let Some(reasoning) = cfg.reasoning.as_deref() else {
3782 return Ok(());
3783 };
3784 let Some(entry) = lookup_model(&cfg.default_model) else {
3785 return Ok(());
3786 };
3787 if entry.provider != PROVIDER_GEMINI {
3788 return Ok(());
3789 }
3790 validate_reasoning_effort(&cfg.default_model, reasoning)
3791}
3792
3793fn validate_runtime_inference_router_config(
3794 registry: &ExtensionRegistry,
3795 cfg: &RuntimeConfig,
3796) -> anyhow::Result<()> {
3797 if !cfg.inference_router.enabled {
3798 return Ok(());
3799 }
3800 let Some(router_id) = cfg.inference_router.router_id.as_deref() else {
3801 anyhow::bail!("inference_router.enabled requires inference_router.router");
3802 };
3803 if registry.inference_router(router_id).is_some() {
3804 return Ok(());
3805 }
3806 let available = registry
3807 .inference_routers
3808 .iter()
3809 .map(|router| router.id())
3810 .collect::<Vec<_>>()
3811 .join(", ");
3812 if available.is_empty() {
3813 anyhow::bail!("inference router {router_id:?} is not registered");
3814 }
3815 anyhow::bail!(
3816 "inference router {router_id:?} is not registered; available routers: {available}"
3817 );
3818}
3819
3820fn model_supports_reasoning(model: &str, effort: &str) -> bool {
3821 lookup_model(model)
3822 .map(|entry| {
3823 entry
3824 .supported_reasoning
3825 .iter()
3826 .any(|option| option.effort == effort)
3827 })
3828 .unwrap_or(false)
3829}
3830
3831fn is_final_answer_phase(phase: Option<&str>) -> bool {
3832 phase.is_none_or(|phase| phase.is_empty() || phase == FINAL_ANSWER_PHASE)
3833}
3834
3835fn edit_tool_for_model<'a>(cfg: &'a RuntimeConfig, model: &'a str) -> Option<&'a str> {
3836 cfg.model_edit_tools
3837 .get(model)
3838 .map(String::as_str)
3839 .or_else(|| {
3840 cfg.model_profiles
3841 .get(model)
3842 .and_then(|profile| profile.edit_tool.as_deref())
3843 })
3844 .or_else(|| lookup_model(model).and_then(|entry| entry.edit_tool))
3845 .or(Some(EDIT_TOOL_EDIT))
3846}
3847
3848fn model_profile_for_model(cfg: &RuntimeConfig, model: &str) -> Option<ModelHarnessProfile> {
3849 cfg.model_profiles
3850 .get(model)
3851 .cloned()
3852 .or_else(|| built_in_model_profile(model))
3853}
3854
3855pub(crate) fn allowlist_permits(allowlist: &[String], tool_name: &str) -> bool {
3856 allowlist.is_empty() || allowlist.iter().any(|allowed| allowed == tool_name)
3857}
3858
3859fn model_profile_for_provider_model(
3869 cfg: &RuntimeConfig,
3870 provider: &str,
3871 model: &str,
3872) -> Option<ModelHarnessProfile> {
3873 cfg.model_profiles
3874 .get(model)
3875 .cloned()
3876 .or_else(|| built_in_model_profile_for_provider(provider, model))
3877}
3878
3879fn schema_policy_for_model(profile: Option<&ModelHarnessProfile>) -> ModelSchemaPolicy {
3880 profile
3881 .map(|profile| profile.schema_policy)
3882 .unwrap_or_default()
3883}
3884
3885fn model_profile_segment_metadata(
3886 profile: Option<&ModelHarnessProfile>,
3887 provider: &str,
3888 model: &str,
3889 segment: &str,
3890) -> serde_json::Value {
3891 serde_json::json!({
3892 "kind": MODEL_PROFILE_TRACE_KIND,
3893 "segment": segment,
3894 "provider": provider,
3895 "model": model,
3896 "profileModel": profile.map(|profile| profile.model.as_str()).unwrap_or(model),
3897 "providerFamily": profile.map(|profile| profile.provider_family),
3898 "editTool": profile.and_then(|profile| profile.edit_tool.as_deref()),
3899 "schemaPolicy": profile.map(|profile| profile.schema_policy),
3900 "instructionOverlay": profile.map(|profile| profile.instruction_overlay),
3901 "parallelToolCalls": profile.and_then(|profile| profile.parallel_tool_calls),
3902 "autoCompactTokenLimit": profile.and_then(|profile| profile.auto_compact_token_limit),
3903 })
3904}
3905
3906fn model_switch_summary(
3907 transcript: &[TranscriptItem],
3908 profile: Option<&ModelHarnessProfile>,
3909 provider: &str,
3910 model: &str,
3911 tools: &[roder_api::tools::ToolSpec],
3912) -> Option<String> {
3913 let previous = latest_model_profile_segment(transcript)?;
3914 let previous_model = previous
3915 .get("model")
3916 .and_then(serde_json::Value::as_str)
3917 .unwrap_or_default();
3918 let previous_provider = previous
3919 .get("provider")
3920 .and_then(serde_json::Value::as_str)
3921 .unwrap_or_default();
3922 if previous_model == model && previous_provider == provider {
3923 return None;
3924 }
3925
3926 let previous_profile = previous
3927 .get("profileModel")
3928 .and_then(serde_json::Value::as_str)
3929 .unwrap_or(previous_model);
3930 let current_profile = profile
3931 .map(|profile| profile.model.as_str())
3932 .unwrap_or(model);
3933 let previous_edit_tool = previous
3934 .get("editTool")
3935 .and_then(serde_json::Value::as_str)
3936 .unwrap_or("none");
3937 let current_edit_tool = profile
3938 .and_then(|profile| profile.edit_tool.as_deref())
3939 .unwrap_or("none");
3940 let tool_names = tools
3941 .iter()
3942 .map(|tool| tool.name.as_str())
3943 .take(12)
3944 .collect::<Vec<_>>()
3945 .join(", ");
3946 Some(format!(
3947 "{MODEL_SWITCH_SUMMARY_PREFIX} previous profile {previous_provider}/{previous_profile} used edit tool {previous_edit_tool}. Current profile {provider}/{current_profile} uses edit tool {current_edit_tool}. Available tools now: {}.",
3948 if tool_names.is_empty() {
3949 "none"
3950 } else {
3951 &tool_names
3952 }
3953 ))
3954}
3955
3956fn latest_model_profile_segment(transcript: &[TranscriptItem]) -> Option<&serde_json::Value> {
3957 transcript.iter().rev().find_map(|item| {
3958 let TranscriptItem::ProviderMetadata(value) = item else {
3959 return None;
3960 };
3961 (value.get("kind").and_then(serde_json::Value::as_str) == Some(MODEL_PROFILE_TRACE_KIND))
3962 .then_some(value)
3963 })
3964}
3965
3966pub fn validate_edit_tool(value: &str) -> anyhow::Result<()> {
3967 match value.trim() {
3968 EDIT_TOOL_PATCH | EDIT_TOOL_EDIT => Ok(()),
3969 _ => anyhow::bail!(
3970 "unsupported edit_tool {value:?}; allowed values: {EDIT_TOOL_PATCH}, {EDIT_TOOL_EDIT}"
3971 ),
3972 }
3973}
3974
3975fn should_persist_thread_event(thread_id: &str) -> bool {
3976 !is_synthetic_event_thread_id(thread_id)
3977}
3978
3979#[cfg(test)]
3980mod tests {
3981 use super::*;
3982 use futures::stream;
3983 use roder_api::catalog::{
3984 PROVIDER_MOCK, REASONING_HIGH, REASONING_LOW, REASONING_MEDIUM, REASONING_MINIMAL,
3985 REASONING_NONE, REASONING_XHIGH,
3986 };
3987 use roder_api::extension::ExtensionRegistryBuilder;
3988 use roder_api::inference::{
3989 CompletionMetadata, InferenceCapabilities, InferenceEngine, InferenceEventStream,
3990 InferenceProviderContext, InferenceTurnContext, MessageDelta, ModelDescriptor,
3991 ModelInstructionOverlay, ModelProfileReasoning, ModelSchemaPolicy, ProviderFamily,
3992 ReasoningEffortDescriptor,
3993 };
3994 use roder_api::inference_routing::{
3995 InferenceRouter, InferenceRoutingContext, InferenceRoutingDecision, InferenceRoutingOutcome,
3996 };
3997 use roder_api::thread::ThreadStoreFactory;
3998 use roder_api::tools::{ToolContributor, ToolExecutor, ToolSpec};
3999 use roder_ext_jsonl_thread_store::store::JsonlThreadStoreFactory;
4000 use std::sync::Mutex as StdMutex;
4001
4002 fn test_workspace() -> String {
4003 std::env::current_dir().unwrap().display().to_string()
4004 }
4005
4006 struct MetadataMissingStore;
4007
4008 #[async_trait::async_trait]
4009 impl ThreadStore for MetadataMissingStore {
4010 fn id(&self) -> roder_api::thread::ThreadStoreId {
4011 "metadata-missing-store".to_string()
4012 }
4013
4014 async fn create_thread(&self, metadata: ThreadMetadata) -> anyhow::Result<ThreadMetadata> {
4015 Ok(metadata)
4016 }
4017
4018 async fn list_threads(&self) -> anyhow::Result<Vec<ThreadMetadata>> {
4019 Ok(Vec::new())
4020 }
4021
4022 async fn load_thread(
4023 &self,
4024 _thread_id: &ThreadId,
4025 ) -> anyhow::Result<Option<ThreadSnapshot>> {
4026 Ok(Some(ThreadSnapshot {
4027 metadata: None,
4028 ..ThreadSnapshot::default()
4029 }))
4030 }
4031
4032 async fn append_event(
4033 &self,
4034 _thread_id: &ThreadId,
4035 _envelope: &EventEnvelope,
4036 ) -> anyhow::Result<()> {
4037 Ok(())
4038 }
4039 }
4040
4041 struct MetadataMissingStoreFactory;
4042
4043 impl ThreadStoreFactory for MetadataMissingStoreFactory {
4044 fn id(&self) -> roder_api::thread::ThreadStoreId {
4045 "metadata-missing-store".to_string()
4046 }
4047
4048 fn create(&self) -> Arc<dyn ThreadStore> {
4049 Arc::new(MetadataMissingStore)
4050 }
4051 }
4052
4053 #[test]
4054 fn synthetic_app_server_events_are_not_thread_events() {
4055 for thread_id in ["app-server", "runtime", "thread-workflow"] {
4056 assert!(!should_persist_thread_event(thread_id));
4057 }
4058 assert!(should_persist_thread_event("thread-discovery"));
4059 assert!(should_persist_thread_event("thread-plan"));
4060 assert!(should_persist_thread_event("thread-process"));
4061 assert!(should_persist_thread_event("thread-1"));
4062 }
4063
4064 #[test]
4065 fn server_side_compaction_uses_catalog_ninety_percent_default() {
4066 assert_eq!(
4067 server_side_compaction_threshold(&RuntimeConfig::default(), "gpt-5.5"),
4068 Some(945_000)
4069 );
4070 assert_eq!(
4071 server_side_compaction_threshold(&RuntimeConfig::default(), "gpt-5.3-codex-spark"),
4072 Some(115_200)
4073 );
4074 }
4075
4076 #[test]
4077 fn server_side_compaction_respects_explicit_config_override() {
4078 let cfg = RuntimeConfig {
4079 auto_compact_token_limit: Some(123_456),
4080 ..RuntimeConfig::default()
4081 };
4082
4083 assert_eq!(
4084 server_side_compaction_threshold(&cfg, "gpt-5.5"),
4085 Some(123_456)
4086 );
4087 }
4088
4089 #[tokio::test]
4090 async fn pre_request_compaction_runs_when_server_side_model_is_at_context_window() {
4091 let captured = Arc::new(StdMutex::new(None));
4092 let mut builder = ExtensionRegistryBuilder::new();
4093 builder.inference_engine(Arc::new(CapturingEngine {
4094 request: captured.clone(),
4095 }));
4096 let thread_root = std::env::temp_dir().join(format!(
4097 "roder-pre-request-compaction-{}",
4098 uuid::Uuid::new_v4()
4099 ));
4100 builder.thread_store_factory(Arc::new(JsonlThreadStoreFactory {
4101 base_path: thread_root.clone(),
4102 }));
4103 let runtime = Arc::new(
4104 Runtime::new(
4105 builder.build().unwrap(),
4106 RuntimeConfig {
4107 default_provider: PROVIDER_MOCK.to_string(),
4108 default_model: "gpt-5.5".to_string(),
4109 file_backed_dynamic_context: true,
4110 ..RuntimeConfig::default()
4111 },
4112 )
4113 .unwrap(),
4114 );
4115 let thread_id = runtime
4116 .create_thread(Some("Pre-request compaction".to_string()))
4117 .await
4118 .unwrap()
4119 .thread_id;
4120 let old_turn = "old-turn".to_string();
4121 runtime
4122 .persist_turn_item(
4123 &thread_id,
4124 &old_turn,
4125 &TranscriptItem::UserMessage(UserMessage::text("old context ".repeat(4_300_000))),
4126 )
4127 .await
4128 .unwrap();
4129
4130 let mut events = runtime.subscribe_events();
4131 runtime
4132 .start_turn(StartTurnRequest {
4133 thread_id: thread_id.clone(),
4134 message: "continue".to_string(),
4135 images: Vec::new(),
4136 provider_override: None,
4137 model_override: None,
4138 reasoning_override: None,
4139 workspace: test_workspace(),
4140 instructions: InstructionBundle::default(),
4141 developer_context: None,
4142 task_ledger_required: false,
4143 })
4144 .await
4145 .unwrap();
4146 loop {
4147 let envelope = tokio::time::timeout(std::time::Duration::from_secs(5), events.recv())
4148 .await
4149 .unwrap()
4150 .unwrap();
4151 if envelope.thread_id.as_deref() == Some(&thread_id)
4152 && matches!(envelope.event, RoderEvent::TurnCompleted(_))
4153 {
4154 break;
4155 }
4156 }
4157
4158 let request = captured.lock().unwrap().clone().unwrap();
4159 assert!(
4160 matches!(
4161 request.transcript.first(),
4162 Some(TranscriptItem::ContextCompaction(_))
4163 ),
4164 "provider request should start with a local emergency compaction item"
4165 );
4166 assert!(
4167 request.transcript.len() < 4,
4168 "provider request should not replay the full oversized prior transcript: {:?}",
4169 request.transcript
4170 );
4171
4172 let _ = std::fs::remove_dir_all(thread_root);
4173 }
4174
4175 #[tokio::test]
4176 async fn continue_after_context_window_failure_compacts_before_provider_request() {
4177 let captured = Arc::new(StdMutex::new(None));
4178 let mut builder = ExtensionRegistryBuilder::new();
4179 builder.inference_engine(Arc::new(CapturingEngine {
4180 request: captured.clone(),
4181 }));
4182 let thread_root = std::env::temp_dir().join(format!(
4183 "roder-context-failure-continue-{}",
4184 uuid::Uuid::new_v4()
4185 ));
4186 builder.thread_store_factory(Arc::new(JsonlThreadStoreFactory {
4187 base_path: thread_root.clone(),
4188 }));
4189 let runtime = Arc::new(
4190 Runtime::new(
4191 builder.build().unwrap(),
4192 RuntimeConfig {
4193 default_provider: PROVIDER_MOCK.to_string(),
4194 default_model: "gpt-5.5".to_string(),
4195 file_backed_dynamic_context: true,
4196 ..RuntimeConfig::default()
4197 },
4198 )
4199 .unwrap(),
4200 );
4201 let thread_id = runtime
4202 .create_thread(Some("Context failure continue".to_string()))
4203 .await
4204 .unwrap()
4205 .thread_id;
4206 let failed_turn = "failed-turn".to_string();
4207 runtime
4208 .persist_turn_item(
4209 &thread_id,
4210 &failed_turn,
4211 &TranscriptItem::UserMessage(UserMessage::text("old work ".repeat(10_000))),
4212 )
4213 .await
4214 .unwrap();
4215 runtime
4216 .persist_turn_item(
4217 &thread_id,
4218 &failed_turn,
4219 &TranscriptItem::Error(ErrorRecord {
4220 message: "Your input exceeds the context window of this model. Please adjust your input and try again."
4221 .to_string(),
4222 }),
4223 )
4224 .await
4225 .unwrap();
4226
4227 let mut events = runtime.subscribe_events();
4228 runtime
4229 .start_turn(StartTurnRequest {
4230 thread_id: thread_id.clone(),
4231 message: "continue".to_string(),
4232 images: Vec::new(),
4233 provider_override: None,
4234 model_override: None,
4235 reasoning_override: None,
4236 workspace: test_workspace(),
4237 instructions: InstructionBundle::default(),
4238 developer_context: None,
4239 task_ledger_required: false,
4240 })
4241 .await
4242 .unwrap();
4243 loop {
4244 let envelope = tokio::time::timeout(std::time::Duration::from_secs(5), events.recv())
4245 .await
4246 .unwrap()
4247 .unwrap();
4248 if envelope.thread_id.as_deref() == Some(&thread_id)
4249 && matches!(envelope.event, RoderEvent::TurnCompleted(_))
4250 {
4251 break;
4252 }
4253 }
4254
4255 let request = captured.lock().unwrap().clone().unwrap();
4256 assert!(
4257 matches!(
4258 request.transcript.first(),
4259 Some(TranscriptItem::ContextCompaction(_))
4260 ),
4261 "provider request after context-window failure should start with local compaction"
4262 );
4263 assert!(
4264 request
4265 .transcript
4266 .iter()
4267 .any(|item| matches!(item, TranscriptItem::UserMessage(message) if message.text == "continue")),
4268 "current continue prompt must be preserved: {:?}",
4269 request.transcript
4270 );
4271 assert!(
4272 !request.transcript.iter().any(
4273 |item| matches!(item, TranscriptItem::Error(error) if error.message.contains("context window"))
4274 ),
4275 "raw prior context-window error should be summarized, not replayed: {:?}",
4276 request.transcript
4277 );
4278
4279 let _ = std::fs::remove_dir_all(thread_root);
4280 }
4281
4282 #[tokio::test]
4283 async fn workspace_for_thread_falls_back_when_metadata_is_missing() {
4284 let workspace = test_workspace();
4285 let mut builder = ExtensionRegistryBuilder::new();
4286 builder.inference_engine(Arc::new(FakeInferenceEngine));
4287 builder.thread_store_factory(Arc::new(MetadataMissingStoreFactory));
4288 let runtime = Runtime::new(
4289 builder.build().unwrap(),
4290 RuntimeConfig {
4291 workspace: Some(workspace.clone()),
4292 ..RuntimeConfig::default()
4293 },
4294 )
4295 .unwrap();
4296
4297 let resolved = runtime
4298 .workspace_for_thread(&ThreadId::from("thread-workflow"))
4299 .await
4300 .unwrap();
4301
4302 assert_eq!(resolved, workspace);
4303 }
4304
4305 #[tokio::test]
4306 async fn automations_can_create_project_thread_with_model_overrides() {
4307 let runtime = Runtime::fake().unwrap();
4308 let workspace = std::env::temp_dir().join("project");
4309 let metadata = runtime
4310 .create_thread_with(CreateThreadRequest {
4311 title: Some("Automation: nightly status".to_string()),
4312 workspace: workspace.display().to_string(),
4313 workspace_id: None,
4314 root_id: None,
4315 provider: Some("mock".to_string()),
4316 model: Some("mock".to_string()),
4317 selection_mode: None,
4318 tool_allowlist: Vec::new(),
4319 developer_instructions: None,
4320 external_tools: Vec::new(),
4321 runner: None,
4322 })
4323 .await
4324 .unwrap();
4325
4326 assert_eq!(
4327 metadata.title.as_deref(),
4328 Some("Automation: nightly status")
4329 );
4330 assert_eq!(metadata.workspace, workspace.display().to_string());
4331 assert_eq!(metadata.provider.as_deref(), Some("mock"));
4332 assert_eq!(metadata.model.as_deref(), Some("mock"));
4333 }
4334
4335 #[test]
4336 fn server_side_compaction_is_only_enabled_for_supported_models() {
4337 let cfg = RuntimeConfig {
4338 auto_compact_token_limit: Some(123_456),
4339 ..RuntimeConfig::default()
4340 };
4341
4342 assert_eq!(server_side_compaction_threshold(&cfg, "mock"), None);
4343 assert_eq!(
4344 server_side_compaction_threshold(&cfg, "codex-auto-review"),
4345 None
4346 );
4347 }
4348
4349 #[test]
4350 fn reasoning_is_disabled_for_models_without_reasoning_support() {
4351 let cfg = RuntimeConfig {
4352 reasoning: Some(REASONING_HIGH.to_string()),
4353 ..RuntimeConfig::default()
4354 };
4355
4356 assert_eq!(
4357 effective_reasoning_for_model(&cfg, "claude-haiku-4-5-20251001"),
4358 REASONING_NONE
4359 );
4360 assert_eq!(
4361 reasoning_for_model(&cfg, "claude-haiku-4-5-20251001"),
4362 ReasoningConfig::default()
4363 );
4364 }
4365
4366 #[test]
4367 fn unsupported_configured_reasoning_falls_back_to_model_default() {
4368 let cfg = RuntimeConfig {
4369 reasoning: Some(REASONING_MINIMAL.to_string()),
4370 ..RuntimeConfig::default()
4371 };
4372
4373 assert_eq!(
4374 effective_reasoning_for_model(&cfg, "gpt-5.5"),
4375 REASONING_MEDIUM
4376 );
4377 }
4378
4379 #[test]
4380 fn unsupported_configured_gemini_reasoning_is_rejected() {
4381 let mut builder = ExtensionRegistryBuilder::new();
4382 builder.inference_engine(std::sync::Arc::new(FakeInferenceEngine));
4383
4384 let err = match Runtime::new(
4385 builder.build().unwrap(),
4386 RuntimeConfig {
4387 default_model: "gemini-3.5-flash".to_string(),
4388 reasoning: Some(REASONING_XHIGH.to_string()),
4389 ..RuntimeConfig::default()
4390 },
4391 ) {
4392 Ok(_) => panic!("expected unsupported Gemini reasoning to be rejected"),
4393 Err(err) => err,
4394 };
4395
4396 assert!(
4397 err.to_string()
4398 .contains("model gemini-3.5-flash does not support reasoning effort xhigh")
4399 );
4400 }
4401
4402 #[tokio::test]
4403 async fn selecting_none_for_non_reasoning_model_preserves_stored_preference() {
4404 let runtime = Runtime::new(
4405 Runtime::fake().unwrap().registry,
4406 RuntimeConfig {
4407 reasoning: Some(REASONING_HIGH.to_string()),
4408 ..RuntimeConfig::default()
4409 },
4410 )
4411 .unwrap();
4412
4413 let cfg = runtime
4414 .select_provider(
4415 roder_api::catalog::PROVIDER_MOCK.to_string(),
4416 Some("claude-haiku-4-5-20251001".to_string()),
4417 Some(REASONING_NONE.to_string()),
4418 )
4419 .await
4420 .unwrap();
4421
4422 assert_eq!(cfg.reasoning.as_deref(), Some(REASONING_HIGH));
4423 assert_eq!(runtime.effective_reasoning().await, REASONING_NONE);
4424 }
4425
4426 #[tokio::test]
4427 async fn selecting_none_for_model_that_supports_none_updates_preference() {
4428 let runtime = Runtime::new(
4429 Runtime::fake().unwrap().registry,
4430 RuntimeConfig {
4431 reasoning: Some(REASONING_HIGH.to_string()),
4432 ..RuntimeConfig::default()
4433 },
4434 )
4435 .unwrap();
4436
4437 let cfg = runtime
4438 .select_provider(
4439 roder_api::catalog::PROVIDER_MOCK.to_string(),
4440 Some("mock".to_string()),
4441 Some(REASONING_NONE.to_string()),
4442 )
4443 .await
4444 .unwrap();
4445
4446 assert_eq!(cfg.reasoning.as_deref(), Some(REASONING_NONE));
4447 }
4448
4449 #[test]
4450 fn parallel_tool_calls_default_on_with_model_override() {
4451 assert!(parallel_tool_calls_for_model(
4452 &RuntimeConfig::default(),
4453 "custom-model"
4454 ));
4455
4456 let cfg = RuntimeConfig {
4457 model_parallel_tool_calls: std::collections::HashMap::from([(
4458 "custom-model".to_string(),
4459 false,
4460 )]),
4461 ..RuntimeConfig::default()
4462 };
4463
4464 assert!(!parallel_tool_calls_for_model(&cfg, "custom-model"));
4465 assert!(parallel_tool_calls_for_model(&cfg, "other-model"));
4466 }
4467
4468 #[test]
4469 fn profile_parallel_tool_calls_applies_between_config_and_default() {
4470 let cfg = RuntimeConfig {
4471 model_profiles: std::collections::HashMap::from([(
4472 "gpt-5.5".to_string(),
4473 test_model_profile("gpt-5.5"),
4474 )]),
4475 ..RuntimeConfig::default()
4476 };
4477
4478 assert!(!parallel_tool_calls_for_model(&cfg, "gpt-5.5"));
4479
4480 let cfg = RuntimeConfig {
4481 model_parallel_tool_calls: std::collections::HashMap::from([(
4482 "gpt-5.5".to_string(),
4483 true,
4484 )]),
4485 ..cfg
4486 };
4487
4488 assert!(parallel_tool_calls_for_model(&cfg, "gpt-5.5"));
4489 }
4490
4491 struct CapturingEngine {
4492 request: Arc<StdMutex<Option<AgentInferenceRequest>>>,
4493 }
4494
4495 #[async_trait::async_trait]
4496 impl InferenceEngine for CapturingEngine {
4497 fn id(&self) -> String {
4498 roder_api::catalog::PROVIDER_MOCK.to_string()
4499 }
4500
4501 fn capabilities(&self) -> InferenceCapabilities {
4502 InferenceCapabilities::coding_agent_default()
4503 }
4504
4505 async fn list_models(
4506 &self,
4507 _ctx: InferenceProviderContext<'_>,
4508 ) -> anyhow::Result<Vec<roder_api::inference::ModelDescriptor>> {
4509 Ok(roder_api::catalog::models_for_provider(
4510 roder_api::catalog::PROVIDER_MOCK,
4511 true,
4512 ))
4513 }
4514
4515 async fn stream_turn(
4516 &self,
4517 _ctx: InferenceTurnContext<'_>,
4518 request: AgentInferenceRequest,
4519 ) -> anyhow::Result<InferenceEventStream> {
4520 *self.request.lock().unwrap() = Some(request);
4521 Ok(Box::pin(stream::iter(vec![
4522 Ok(InferenceEvent::MessageDelta(MessageDelta {
4523 text: "done".to_string(),
4524 phase: None,
4525 })),
4526 Ok(InferenceEvent::Completed(CompletionMetadata {
4527 stop_reason: Some("stop".to_string()),
4528 provider_response_id: None,
4529 })),
4530 ])))
4531 }
4532 }
4533
4534 struct RoutingCaptureEngine {
4535 id: &'static str,
4536 models: Vec<ModelDescriptor>,
4537 requests: Arc<StdMutex<Vec<AgentInferenceRequest>>>,
4538 }
4539
4540 #[async_trait::async_trait]
4541 impl InferenceEngine for RoutingCaptureEngine {
4542 fn id(&self) -> String {
4543 self.id.to_string()
4544 }
4545
4546 fn capabilities(&self) -> InferenceCapabilities {
4547 InferenceCapabilities::coding_agent_default()
4548 }
4549
4550 async fn list_models(
4551 &self,
4552 _ctx: InferenceProviderContext<'_>,
4553 ) -> anyhow::Result<Vec<ModelDescriptor>> {
4554 Ok(self.models.clone())
4555 }
4556
4557 async fn stream_turn(
4558 &self,
4559 _ctx: InferenceTurnContext<'_>,
4560 request: AgentInferenceRequest,
4561 ) -> anyhow::Result<InferenceEventStream> {
4562 self.requests.lock().unwrap().push(request);
4563 Ok(Box::pin(stream::iter(vec![
4564 Ok(InferenceEvent::MessageDelta(MessageDelta {
4565 text: "routed".to_string(),
4566 phase: None,
4567 })),
4568 Ok(InferenceEvent::Completed(CompletionMetadata {
4569 stop_reason: Some("stop".to_string()),
4570 provider_response_id: None,
4571 })),
4572 ])))
4573 }
4574 }
4575
4576 struct StaticRouter {
4577 id: &'static str,
4578 decision: InferenceRoutingDecision,
4579 contexts: Arc<StdMutex<Vec<InferenceRoutingContext>>>,
4580 }
4581
4582 #[async_trait::async_trait]
4583 impl InferenceRouter for StaticRouter {
4584 fn id(&self) -> String {
4585 self.id.to_string()
4586 }
4587
4588 async fn route(
4589 &self,
4590 context: InferenceRoutingContext,
4591 ) -> anyhow::Result<InferenceRoutingDecision> {
4592 self.contexts.lock().unwrap().push(context);
4593 Ok(self.decision.clone())
4594 }
4595 }
4596
4597 fn routing_test_model(id: &str, supported_reasoning: &[&str]) -> ModelDescriptor {
4598 ModelDescriptor {
4599 id: id.to_string(),
4600 name: id.to_string(),
4601 context_window: Some(128_000),
4602 default_reasoning: supported_reasoning
4603 .first()
4604 .map(|effort| (*effort).to_string()),
4605 supported_reasoning: supported_reasoning
4606 .iter()
4607 .map(|effort| ReasoningEffortDescriptor {
4608 effort: (*effort).to_string(),
4609 description: format!("{effort} reasoning"),
4610 })
4611 .collect(),
4612 }
4613 }
4614
4615 struct TaskLedgerCompletionGateEngine {
4616 calls: StdMutex<u32>,
4617 requests: Arc<StdMutex<Vec<AgentInferenceRequest>>>,
4618 }
4619
4620 #[async_trait::async_trait]
4621 impl InferenceEngine for TaskLedgerCompletionGateEngine {
4622 fn id(&self) -> String {
4623 roder_api::catalog::PROVIDER_MOCK.to_string()
4624 }
4625
4626 fn capabilities(&self) -> InferenceCapabilities {
4627 InferenceCapabilities::coding_agent_default()
4628 }
4629
4630 async fn list_models(
4631 &self,
4632 _ctx: InferenceProviderContext<'_>,
4633 ) -> anyhow::Result<Vec<roder_api::inference::ModelDescriptor>> {
4634 Ok(roder_api::catalog::models_for_provider(
4635 roder_api::catalog::PROVIDER_MOCK,
4636 true,
4637 ))
4638 }
4639
4640 async fn stream_turn(
4641 &self,
4642 _ctx: InferenceTurnContext<'_>,
4643 request: AgentInferenceRequest,
4644 ) -> anyhow::Result<InferenceEventStream> {
4645 self.requests.lock().unwrap().push(request);
4646 let mut calls = self.calls.lock().unwrap();
4647 *calls += 1;
4648 let events = match *calls {
4649 1 => vec![Ok(InferenceEvent::ToolCallCompleted(ToolCallCompleted {
4650 id: "ledger-open".to_string(),
4651 name: TASK_LEDGER_TOOL_NAME.to_string(),
4652 arguments: serde_json::json!({
4653 "tasks": [
4654 {
4655 "id": "inspect",
4656 "content": "Inspect local assets",
4657 "status": "completed",
4658 "evidence": "listed workspace"
4659 },
4660 {
4661 "id": "write",
4662 "content": "Write /app/result.txt",
4663 "status": "pending"
4664 }
4665 ],
4666 "requireCompletionEvidence": true
4667 })
4668 .to_string(),
4669 }))],
4670 3 => vec![Ok(InferenceEvent::ToolCallCompleted(ToolCallCompleted {
4671 id: "ledger-complete".to_string(),
4672 name: TASK_LEDGER_TOOL_NAME.to_string(),
4673 arguments: serde_json::json!({
4674 "tasks": [
4675 {
4676 "id": "inspect",
4677 "content": "Inspect local assets",
4678 "status": "completed",
4679 "evidence": "listed workspace"
4680 },
4681 {
4682 "id": "write",
4683 "content": "Write /app/result.txt",
4684 "status": "completed",
4685 "evidence": "wrote answer"
4686 }
4687 ],
4688 "requireCompletionEvidence": true
4689 })
4690 .to_string(),
4691 }))],
4692 _ => vec![Ok(InferenceEvent::MessageDelta(MessageDelta {
4693 text: "final".to_string(),
4694 phase: None,
4695 }))],
4696 };
4697 Ok(Box::pin(stream::iter(events.into_iter().chain(
4698 std::iter::once(Ok(InferenceEvent::Completed(CompletionMetadata {
4699 stop_reason: Some("stop".to_string()),
4700 provider_response_id: None,
4701 }))),
4702 ))))
4703 }
4704 }
4705
4706 struct VerificationGateEngine {
4707 calls: StdMutex<u32>,
4708 }
4709
4710 #[async_trait::async_trait]
4711 impl InferenceEngine for VerificationGateEngine {
4712 fn id(&self) -> String {
4713 roder_api::catalog::PROVIDER_MOCK.to_string()
4714 }
4715
4716 fn capabilities(&self) -> InferenceCapabilities {
4717 InferenceCapabilities::coding_agent_default()
4718 }
4719
4720 async fn list_models(
4721 &self,
4722 _ctx: InferenceProviderContext<'_>,
4723 ) -> anyhow::Result<Vec<roder_api::inference::ModelDescriptor>> {
4724 Ok(roder_api::catalog::models_for_provider(
4725 roder_api::catalog::PROVIDER_MOCK,
4726 true,
4727 ))
4728 }
4729
4730 async fn stream_turn(
4731 &self,
4732 _ctx: InferenceTurnContext<'_>,
4733 request: AgentInferenceRequest,
4734 ) -> anyhow::Result<InferenceEventStream> {
4735 let mut calls = self.calls.lock().unwrap();
4736 *calls += 1;
4737 let events = match *calls {
4738 1 => vec![Ok(InferenceEvent::ToolCallCompleted(ToolCallCompleted {
4739 id: "write-1".to_string(),
4740 name: "write_file".to_string(),
4741 arguments: serde_json::json!({
4742 "path": "src/lib.rs",
4743 "content": "pub fn answer() -> u8 { 42 }\n"
4744 })
4745 .to_string(),
4746 }))],
4747 2 => vec![Ok(InferenceEvent::MessageDelta(MessageDelta {
4748 text: "done too early".to_string(),
4749 phase: None,
4750 }))],
4751 3 if request.transcript.iter().any(|item| {
4752 matches!(
4753 item,
4754 TranscriptItem::UserMessage(message)
4755 if message.text.contains("Verification gate blocked final completion")
4756 )
4757 }) =>
4758 {
4759 vec![Ok(InferenceEvent::ToolCallCompleted(ToolCallCompleted {
4760 id: "verify-1".to_string(),
4761 name: crate::verification_gate::VERIFICATION_TOOL_NAME.to_string(),
4762 arguments: serde_json::json!({
4763 "originalTask": "write code",
4764 "changedFiles": ["src/lib.rs"],
4765 "toolEvidence": ["write_file wrote src/lib.rs"],
4766 "testsRun": ["cargo test -p roder-core verification_gate"],
4767 "openGaps": [],
4768 "status": "completed"
4769 })
4770 .to_string(),
4771 }))]
4772 }
4773 _ => vec![Ok(InferenceEvent::MessageDelta(MessageDelta {
4774 text: "verified final".to_string(),
4775 phase: None,
4776 }))],
4777 };
4778 Ok(Box::pin(stream::iter(events.into_iter().chain(
4779 std::iter::once(Ok(InferenceEvent::Completed(CompletionMetadata {
4780 stop_reason: Some("stop".to_string()),
4781 provider_response_id: None,
4782 }))),
4783 ))))
4784 }
4785 }
4786
4787 struct SpeedPolicyEngine {
4788 calls: StdMutex<u32>,
4789 requests: Arc<StdMutex<Vec<AgentInferenceRequest>>>,
4790 }
4791
4792 #[async_trait::async_trait]
4793 impl InferenceEngine for SpeedPolicyEngine {
4794 fn id(&self) -> String {
4795 roder_api::catalog::PROVIDER_MOCK.to_string()
4796 }
4797
4798 fn capabilities(&self) -> InferenceCapabilities {
4799 InferenceCapabilities::coding_agent_default()
4800 }
4801
4802 async fn list_models(
4803 &self,
4804 _ctx: InferenceProviderContext<'_>,
4805 ) -> anyhow::Result<Vec<roder_api::inference::ModelDescriptor>> {
4806 Ok(roder_api::catalog::models_for_provider(
4807 roder_api::catalog::PROVIDER_MOCK,
4808 true,
4809 ))
4810 }
4811
4812 async fn stream_turn(
4813 &self,
4814 _ctx: InferenceTurnContext<'_>,
4815 request: AgentInferenceRequest,
4816 ) -> anyhow::Result<InferenceEventStream> {
4817 self.requests.lock().unwrap().push(request.clone());
4818 let mut calls = self.calls.lock().unwrap();
4819 *calls += 1;
4820 let events = match *calls {
4821 1 => vec![Ok(InferenceEvent::ToolCallCompleted(ToolCallCompleted {
4822 id: "write-1".to_string(),
4823 name: "write_file".to_string(),
4824 arguments: serde_json::json!({
4825 "path": "src/lib.rs",
4826 "content": "pub fn answer() -> u8 { 42 }\n"
4827 })
4828 .to_string(),
4829 }))],
4830 3 if request.transcript.iter().any(|item| {
4831 matches!(
4832 item,
4833 TranscriptItem::UserMessage(message)
4834 if message.text.contains("Verification gate blocked final completion")
4835 )
4836 }) =>
4837 {
4838 vec![Ok(InferenceEvent::ToolCallCompleted(ToolCallCompleted {
4839 id: "verify-1".to_string(),
4840 name: crate::verification_gate::VERIFICATION_TOOL_NAME.to_string(),
4841 arguments: serde_json::json!({
4842 "originalTask": "write code",
4843 "changedFiles": ["src/lib.rs"],
4844 "toolEvidence": ["write_file wrote src/lib.rs"],
4845 "testsRun": ["cargo test -p roder-core speed_policy"],
4846 "openGaps": [],
4847 "status": "completed"
4848 })
4849 .to_string(),
4850 }))]
4851 }
4852 _ => vec![Ok(InferenceEvent::MessageDelta(MessageDelta {
4853 text: "done".to_string(),
4854 phase: None,
4855 }))],
4856 };
4857 Ok(Box::pin(stream::iter(events.into_iter().chain(
4858 std::iter::once(Ok(InferenceEvent::Completed(CompletionMetadata {
4859 stop_reason: Some("stop".to_string()),
4860 provider_response_id: None,
4861 }))),
4862 ))))
4863 }
4864 }
4865
4866 struct SwitchCaptureEngine {
4867 requests: Arc<StdMutex<Vec<AgentInferenceRequest>>>,
4868 }
4869
4870 #[async_trait::async_trait]
4871 impl InferenceEngine for SwitchCaptureEngine {
4872 fn id(&self) -> String {
4873 roder_api::catalog::PROVIDER_MOCK.to_string()
4874 }
4875
4876 fn capabilities(&self) -> InferenceCapabilities {
4877 InferenceCapabilities::coding_agent_default()
4878 }
4879
4880 async fn list_models(
4881 &self,
4882 _ctx: InferenceProviderContext<'_>,
4883 ) -> anyhow::Result<Vec<roder_api::inference::ModelDescriptor>> {
4884 Ok(roder_api::catalog::models_for_provider(
4885 roder_api::catalog::PROVIDER_MOCK,
4886 true,
4887 ))
4888 }
4889
4890 async fn stream_turn(
4891 &self,
4892 _ctx: InferenceTurnContext<'_>,
4893 request: AgentInferenceRequest,
4894 ) -> anyhow::Result<InferenceEventStream> {
4895 self.requests.lock().unwrap().push(request);
4896 Ok(Box::pin(stream::iter(vec![
4897 Ok(InferenceEvent::MessageDelta(MessageDelta {
4898 text: "done".to_string(),
4899 phase: None,
4900 })),
4901 Ok(InferenceEvent::Completed(CompletionMetadata {
4902 stop_reason: Some("stop".to_string()),
4903 provider_response_id: None,
4904 })),
4905 ])))
4906 }
4907 }
4908
4909 struct DeadlineEngine;
4910
4911 #[async_trait::async_trait]
4912 impl InferenceEngine for DeadlineEngine {
4913 fn id(&self) -> String {
4914 roder_api::catalog::PROVIDER_MOCK.to_string()
4915 }
4916
4917 fn capabilities(&self) -> InferenceCapabilities {
4918 InferenceCapabilities::coding_agent_default()
4919 }
4920
4921 async fn list_models(
4922 &self,
4923 _ctx: InferenceProviderContext<'_>,
4924 ) -> anyhow::Result<Vec<roder_api::inference::ModelDescriptor>> {
4925 Ok(Vec::new())
4926 }
4927
4928 async fn stream_turn(
4929 &self,
4930 _ctx: InferenceTurnContext<'_>,
4931 _request: AgentInferenceRequest,
4932 ) -> anyhow::Result<InferenceEventStream> {
4933 Ok(Box::pin(stream::once(async {
4934 tokio::time::sleep(std::time::Duration::from_secs(60)).await;
4935 Ok(InferenceEvent::MessageDelta(MessageDelta {
4936 text: "too late".to_string(),
4937 phase: None,
4938 }))
4939 })))
4940 }
4941 }
4942
4943 struct WriteFileContributor;
4944
4945 impl ToolContributor for WriteFileContributor {
4946 fn id(&self) -> String {
4947 "test-write".to_string()
4948 }
4949
4950 fn contribute(&self, registry: &mut ToolRegistry) -> anyhow::Result<()> {
4951 registry.register(Arc::new(WriteFileTool))
4952 }
4953 }
4954
4955 struct WriteFileTool;
4956
4957 #[async_trait::async_trait]
4958 impl ToolExecutor for WriteFileTool {
4959 fn spec(&self) -> ToolSpec {
4960 ToolSpec {
4961 name: "write_file".to_string(),
4962 description: "Write a test file.".to_string(),
4963 parameters: serde_json::json!({
4964 "type": "object",
4965 "properties": {
4966 "path": { "type": "string" },
4967 "content": { "type": "string" }
4968 },
4969 "required": ["path", "content"],
4970 "additionalProperties": false
4971 }),
4972 }
4973 }
4974
4975 async fn execute(
4976 &self,
4977 _ctx: ToolExecutionContext,
4978 call: ToolCall,
4979 ) -> anyhow::Result<ToolResult> {
4980 let path = call
4981 .arguments
4982 .get("path")
4983 .and_then(serde_json::Value::as_str)
4984 .unwrap_or("src/lib.rs");
4985 Ok(ToolResult {
4986 id: call.id,
4987 name: call.name,
4988 text: format!("wrote {path}"),
4989 data: serde_json::json!({ "path": path }),
4990 is_error: false,
4991 })
4992 }
4993 }
4994
4995 struct ProfileToolContributor;
4996
4997 impl ToolContributor for ProfileToolContributor {
4998 fn id(&self) -> String {
4999 "profile-tools".to_string()
5000 }
5001
5002 fn contribute(&self, registry: &mut ToolRegistry) -> anyhow::Result<()> {
5003 for name in ["apply_patch", "edit", "multi_edit", "write_file"] {
5004 registry.register(Arc::new(ProfileTool {
5005 name: name.to_string(),
5006 }))?;
5007 }
5008 Ok(())
5009 }
5010 }
5011
5012 struct ProfileTool {
5013 name: String,
5014 }
5015
5016 #[async_trait::async_trait]
5017 impl ToolExecutor for ProfileTool {
5018 fn spec(&self) -> ToolSpec {
5019 ToolSpec {
5020 name: self.name.clone(),
5021 description: format!("{} test tool", self.name),
5022 parameters: serde_json::json!({
5023 "type": "object",
5024 "properties": {
5025 "path": { "type": "string" },
5026 "content": { "type": "string" }
5027 },
5028 "required": ["path", "content"],
5029 "additionalProperties": false
5030 }),
5031 }
5032 }
5033
5034 async fn execute(
5035 &self,
5036 _ctx: ToolExecutionContext,
5037 call: ToolCall,
5038 ) -> anyhow::Result<ToolResult> {
5039 Ok(ToolResult {
5040 id: call.id,
5041 name: call.name,
5042 text: "ok".to_string(),
5043 data: serde_json::json!({}),
5044 is_error: false,
5045 })
5046 }
5047 }
5048
5049 fn test_model_profile(model: &str) -> ModelHarnessProfile {
5050 ModelHarnessProfile {
5051 model: model.to_string(),
5052 provider: roder_api::catalog::PROVIDER_OPENAI.to_string(),
5053 provider_family: ProviderFamily::OpenAi,
5054 edit_tool: Some(EDIT_TOOL_EDIT.to_string()),
5055 schema_policy: ModelSchemaPolicy::StandardRequiredFirst,
5056 instruction_overlay: ModelInstructionOverlay::IntuitiveContext,
5057 reasoning: ModelProfileReasoning {
5058 orientation: Some(REASONING_LOW.to_string()),
5059 execution: Some(REASONING_LOW.to_string()),
5060 verification: Some(REASONING_LOW.to_string()),
5061 recovery: Some(REASONING_LOW.to_string()),
5062 },
5063 parallel_tool_calls: Some(false),
5064 auto_compact_token_limit: Some(123_000),
5065 }
5066 }
5067
5068 async fn captured_profile_request(cfg: RuntimeConfig) -> AgentInferenceRequest {
5069 let captured = Arc::new(StdMutex::new(None));
5070 let mut builder = ExtensionRegistryBuilder::new();
5071 builder.inference_engine(Arc::new(CapturingEngine {
5072 request: captured.clone(),
5073 }));
5074 builder.tool_contributor(Arc::new(ProfileToolContributor));
5075 let runtime = Arc::new(Runtime::new(builder.build().unwrap(), cfg).unwrap());
5076 let mut rx = runtime.subscribe_events();
5077 let turn_id = runtime
5078 .start_turn(StartTurnRequest {
5079 thread_id: "thread-model-profile".to_string(),
5080 message: "use profile knobs".to_string(),
5081 images: Vec::new(),
5082 provider_override: None,
5083 model_override: None,
5084 reasoning_override: None,
5085 workspace: test_workspace(),
5086 instructions: InstructionBundle {
5087 system: None,
5088 developer: Some("base developer".to_string()),
5089 developer_context: None,
5090 },
5091 developer_context: None,
5092 task_ledger_required: false,
5093 })
5094 .await
5095 .unwrap();
5096
5097 tokio::time::timeout(std::time::Duration::from_secs(5), async {
5098 loop {
5099 let envelope = rx.recv().await.unwrap();
5100 if envelope.turn_id.as_deref() != Some(&turn_id) {
5101 continue;
5102 }
5103 match envelope.event {
5104 RoderEvent::TurnCompleted(_) => break,
5105 RoderEvent::TurnFailed(event) => panic!("turn failed: {}", event.error),
5106 _ => {}
5107 }
5108 }
5109 })
5110 .await
5111 .unwrap();
5112
5113 captured.lock().unwrap().clone().unwrap()
5114 }
5115
5116 struct ToolThenStopEngine {
5117 calls: StdMutex<u32>,
5118 }
5119
5120 #[async_trait::async_trait]
5121 impl InferenceEngine for ToolThenStopEngine {
5122 fn id(&self) -> String {
5123 roder_api::catalog::PROVIDER_MOCK.to_string()
5124 }
5125
5126 fn capabilities(&self) -> InferenceCapabilities {
5127 InferenceCapabilities::coding_agent_default()
5128 }
5129
5130 async fn list_models(
5131 &self,
5132 _ctx: InferenceProviderContext<'_>,
5133 ) -> anyhow::Result<Vec<roder_api::inference::ModelDescriptor>> {
5134 Ok(roder_api::catalog::models_for_provider(
5135 roder_api::catalog::PROVIDER_MOCK,
5136 true,
5137 ))
5138 }
5139
5140 async fn stream_turn(
5141 &self,
5142 _ctx: InferenceTurnContext<'_>,
5143 _request: AgentInferenceRequest,
5144 ) -> anyhow::Result<InferenceEventStream> {
5145 let mut calls = self.calls.lock().unwrap();
5146 *calls += 1;
5147 let events = match *calls {
5148 1 => vec![
5149 Ok(InferenceEvent::ToolCallCompleted(ToolCallCompleted {
5150 id: "write-1".to_string(),
5151 name: "write_file".to_string(),
5152 arguments: serde_json::json!({
5153 "path": "src/lib.rs",
5154 "content": "pub fn answer() -> u8 { 42 }\n"
5155 })
5156 .to_string(),
5157 })),
5158 Ok(InferenceEvent::Completed(CompletionMetadata {
5159 stop_reason: Some("tool_use".to_string()),
5160 provider_response_id: None,
5161 })),
5162 ],
5163 _ => vec![
5164 Ok(InferenceEvent::MessageDelta(MessageDelta {
5165 text: "final".to_string(),
5166 phase: None,
5167 })),
5168 Ok(InferenceEvent::Completed(CompletionMetadata {
5169 stop_reason: Some("end_turn".to_string()),
5170 provider_response_id: None,
5171 })),
5172 ],
5173 };
5174 Ok(Box::pin(stream::iter(events)))
5175 }
5176 }
5177
5178 #[tokio::test]
5179 async fn turn_completed_reports_terminal_step_finish_reason() {
5180 let mut builder = ExtensionRegistryBuilder::new();
5181 builder.inference_engine(Arc::new(ToolThenStopEngine {
5182 calls: StdMutex::new(0),
5183 }));
5184 builder.tool_contributor(Arc::new(WriteFileContributor));
5185 let runtime = Arc::new(
5186 Runtime::new(
5187 builder.build().unwrap(),
5188 RuntimeConfig {
5189 policy_mode: PolicyMode::Bypass,
5190 ..RuntimeConfig::default()
5191 },
5192 )
5193 .unwrap(),
5194 );
5195 let mut rx = runtime.subscribe_events();
5196 let turn_id = runtime
5197 .start_turn(StartTurnRequest {
5198 thread_id: "thread-finish-reason".to_string(),
5199 message: "write then finish".to_string(),
5200 images: Vec::new(),
5201 provider_override: None,
5202 model_override: None,
5203 reasoning_override: None,
5204 workspace: test_workspace(),
5205 instructions: InstructionBundle {
5206 system: None,
5207 developer: None,
5208 developer_context: None,
5209 },
5210 developer_context: None,
5211 task_ledger_required: false,
5212 })
5213 .await
5214 .unwrap();
5215
5216 let completed = tokio::time::timeout(std::time::Duration::from_secs(5), async {
5217 loop {
5218 let envelope = rx.recv().await.unwrap();
5219 if envelope.turn_id.as_deref() != Some(&turn_id) {
5220 continue;
5221 }
5222 match envelope.event {
5223 RoderEvent::TurnCompleted(event) => break event,
5224 RoderEvent::TurnFailed(event) => panic!("turn failed: {}", event.error),
5225 _ => {}
5226 }
5227 }
5228 })
5229 .await
5230 .unwrap();
5231
5232 assert_eq!(completed.finish_reason.as_deref(), Some("stop"));
5235 }
5236
5237 #[tokio::test]
5238 async fn inference_router_selection_changes_request_model_and_records_event() {
5239 let default_requests = Arc::new(StdMutex::new(Vec::<AgentInferenceRequest>::new()));
5240 let routed_requests = Arc::new(StdMutex::new(Vec::<AgentInferenceRequest>::new()));
5241 let contexts = Arc::new(StdMutex::new(Vec::<InferenceRoutingContext>::new()));
5242 let selected = ModelSelection {
5243 provider: "routed-provider".to_string(),
5244 model: "routed-model".to_string(),
5245 };
5246 let default = ModelSelection {
5247 provider: roder_api::catalog::PROVIDER_MOCK.to_string(),
5248 model: "mock".to_string(),
5249 };
5250 let decision = InferenceRoutingDecision {
5251 reasoning: Some(ReasoningConfig {
5252 enabled: true,
5253 level: Some(REASONING_LOW.to_string()),
5254 }),
5255 confidence: Some(0.91),
5256 baseline: Some(default.clone()),
5257 matched_signals: vec![roder_api::inference_routing::InferenceRoutingSignal::new(
5258 "intent", "routine",
5259 )],
5260 ..InferenceRoutingDecision::selected("test-router", selected.clone(), "routine request")
5261 };
5262
5263 let mut builder = ExtensionRegistryBuilder::new();
5264 builder.inference_engine(Arc::new(RoutingCaptureEngine {
5265 id: roder_api::catalog::PROVIDER_MOCK,
5266 models: vec![routing_test_model("mock", &[REASONING_LOW])],
5267 requests: default_requests.clone(),
5268 }));
5269 builder.inference_engine(Arc::new(RoutingCaptureEngine {
5270 id: "routed-provider",
5271 models: vec![routing_test_model(
5272 "routed-model",
5273 &[REASONING_LOW, REASONING_MEDIUM],
5274 )],
5275 requests: routed_requests.clone(),
5276 }));
5277 builder.inference_router(Arc::new(StaticRouter {
5278 id: "test-router",
5279 decision,
5280 contexts: contexts.clone(),
5281 }));
5282 let thread_root =
5283 std::env::temp_dir().join(format!("roder-routing-auto-{}", uuid::Uuid::new_v4()));
5284 builder.thread_store_factory(Arc::new(JsonlThreadStoreFactory {
5285 base_path: thread_root.clone(),
5286 }));
5287 let runtime = Arc::new(
5288 Runtime::new(
5289 builder.build().unwrap(),
5290 RuntimeConfig {
5291 default_provider: default.provider.clone(),
5292 default_model: default.model.clone(),
5293 ..RuntimeConfig::default()
5294 },
5295 )
5296 .unwrap(),
5297 );
5298 let thread_id = runtime
5299 .create_thread_with(CreateThreadRequest {
5300 title: Some("Routing auto".to_string()),
5301 workspace: test_workspace(),
5302 workspace_id: None,
5303 root_id: None,
5304 provider: Some(default.provider.clone()),
5305 model: Some(default.model.clone()),
5306 tool_allowlist: Vec::new(),
5307 developer_instructions: None,
5308 external_tools: Vec::new(),
5309 selection_mode: Some(ModelSelectionMode::auto(
5310 "test-router:coding",
5311 "test-router",
5312 "Auto: Coding",
5313 default.clone(),
5314 Some("coding".to_string()),
5315 None,
5316 )),
5317 runner: None,
5318 })
5319 .await
5320 .unwrap()
5321 .thread_id;
5322 let mut rx = runtime.subscribe_events();
5323 let turn_id = runtime
5324 .start_turn(StartTurnRequest {
5325 thread_id: thread_id.clone(),
5326 message: "small cleanup".to_string(),
5327 images: Vec::new(),
5328 provider_override: None,
5329 model_override: None,
5330 reasoning_override: None,
5331 workspace: test_workspace(),
5332 instructions: InstructionBundle::default(),
5333 developer_context: None,
5334 task_ledger_required: false,
5335 })
5336 .await
5337 .unwrap();
5338
5339 let mut routing_event = None;
5340 let mut inference_started = None;
5341 tokio::time::timeout(std::time::Duration::from_secs(5), async {
5342 loop {
5343 let envelope = rx.recv().await.unwrap();
5344 if envelope.turn_id.as_deref() != Some(&turn_id) {
5345 continue;
5346 }
5347 match envelope.event {
5348 RoderEvent::InferenceRoutingDecision(event) => {
5349 routing_event = Some(event);
5350 }
5351 RoderEvent::InferenceStarted(event) => {
5352 inference_started = Some(event);
5353 }
5354 RoderEvent::TurnCompleted(_) => break,
5355 RoderEvent::TurnFailed(event) => panic!("turn failed: {}", event.error),
5356 _ => {}
5357 }
5358 }
5359 })
5360 .await
5361 .unwrap();
5362
5363 assert!(default_requests.lock().unwrap().is_empty());
5364 let routed_requests = routed_requests.lock().unwrap();
5365 assert_eq!(routed_requests.len(), 1);
5366 assert_eq!(routed_requests[0].model, selected);
5367 assert_eq!(
5368 routed_requests[0].reasoning.level.as_deref(),
5369 Some(REASONING_LOW)
5370 );
5371 assert_eq!(
5372 routed_requests[0].metadata["inferenceRouting"]["outcome"],
5373 "selected"
5374 );
5375
5376 let routing_event = routing_event.expect("routing decision event");
5377 assert_eq!(routing_event.default_selection, default);
5378 assert_eq!(routing_event.selected_selection, selected);
5379 assert_eq!(
5380 routing_event.decision.outcome,
5381 InferenceRoutingOutcome::Selected
5382 );
5383 assert_eq!(
5384 inference_started.expect("inference started event").model,
5385 selected
5386 );
5387
5388 let contexts = contexts.lock().unwrap();
5389 assert_eq!(contexts.len(), 1);
5390 assert_eq!(contexts[0].default_selection, default);
5391 assert_eq!(contexts[0].candidates.len(), 2);
5392 assert!(
5393 contexts[0]
5394 .signals
5395 .iter()
5396 .any(|signal| signal.key == "profile" && signal.value == "coding")
5397 );
5398 let _ = std::fs::remove_dir_all(thread_root);
5399 }
5400
5401 #[tokio::test]
5402 async fn inference_router_is_bypassed_for_explicit_selection() {
5403 let requests = Arc::new(StdMutex::new(Vec::<AgentInferenceRequest>::new()));
5404 let contexts = Arc::new(StdMutex::new(Vec::<InferenceRoutingContext>::new()));
5405 let selected = ModelSelection {
5406 provider: roder_api::catalog::PROVIDER_MOCK.to_string(),
5407 model: "mock".to_string(),
5408 };
5409
5410 let mut builder = ExtensionRegistryBuilder::new();
5411 builder.inference_engine(Arc::new(RoutingCaptureEngine {
5412 id: roder_api::catalog::PROVIDER_MOCK,
5413 models: vec![routing_test_model("mock", &[REASONING_LOW])],
5414 requests: requests.clone(),
5415 }));
5416 builder.inference_router(Arc::new(StaticRouter {
5417 id: "test-router",
5418 decision: InferenceRoutingDecision::selected(
5419 "test-router",
5420 ModelSelection {
5421 provider: "missing".to_string(),
5422 model: "missing".to_string(),
5423 },
5424 "would route if called",
5425 ),
5426 contexts: contexts.clone(),
5427 }));
5428 let thread_root =
5429 std::env::temp_dir().join(format!("roder-routing-explicit-{}", uuid::Uuid::new_v4()));
5430 builder.thread_store_factory(Arc::new(JsonlThreadStoreFactory {
5431 base_path: thread_root.clone(),
5432 }));
5433 let runtime = Arc::new(
5434 Runtime::new(
5435 builder.build().unwrap(),
5436 RuntimeConfig {
5437 default_provider: selected.provider.clone(),
5438 default_model: selected.model.clone(),
5439 inference_router: RuntimeInferenceRouterConfig {
5440 enabled: true,
5441 router_id: Some("test-router".to_string()),
5442 },
5443 ..RuntimeConfig::default()
5444 },
5445 )
5446 .unwrap(),
5447 );
5448 let thread_id = runtime
5449 .create_thread_with(CreateThreadRequest {
5450 title: Some("Routing explicit".to_string()),
5451 workspace: test_workspace(),
5452 workspace_id: None,
5453 root_id: None,
5454 provider: Some(selected.provider.clone()),
5455 model: Some(selected.model.clone()),
5456 tool_allowlist: Vec::new(),
5457 developer_instructions: None,
5458 external_tools: Vec::new(),
5459 selection_mode: Some(ModelSelectionMode::auto(
5460 "test-router:default",
5461 "test-router",
5462 "Auto",
5463 selected.clone(),
5464 None,
5465 None,
5466 )),
5467 runner: None,
5468 })
5469 .await
5470 .unwrap()
5471 .thread_id;
5472 let mut rx = runtime.subscribe_events();
5473 let turn_id = runtime
5474 .start_turn(StartTurnRequest {
5475 thread_id,
5476 message: "use explicit selection".to_string(),
5477 images: Vec::new(),
5478 provider_override: Some(selected.provider.clone()),
5479 model_override: Some(selected.model.clone()),
5480 reasoning_override: None,
5481 workspace: test_workspace(),
5482 instructions: InstructionBundle::default(),
5483 developer_context: None,
5484 task_ledger_required: false,
5485 })
5486 .await
5487 .unwrap();
5488
5489 let mut saw_routing_event = false;
5490 tokio::time::timeout(std::time::Duration::from_secs(5), async {
5491 loop {
5492 let envelope = rx.recv().await.unwrap();
5493 if envelope.turn_id.as_deref() != Some(&turn_id) {
5494 continue;
5495 }
5496 match envelope.event {
5497 RoderEvent::InferenceRoutingDecision(_) => {
5498 saw_routing_event = true;
5499 }
5500 RoderEvent::TurnCompleted(_) => break,
5501 RoderEvent::TurnFailed(event) => panic!("turn failed: {}", event.error),
5502 _ => {}
5503 }
5504 }
5505 })
5506 .await
5507 .unwrap();
5508
5509 assert!(!saw_routing_event);
5510 assert!(contexts.lock().unwrap().is_empty());
5511 let requests = requests.lock().unwrap();
5512 assert_eq!(requests.len(), 1);
5513 assert_eq!(requests[0].model, selected);
5514 let _ = std::fs::remove_dir_all(thread_root);
5515 }
5516
5517 #[tokio::test]
5518 async fn inference_router_is_bypassed_for_manual_selection_mode() {
5519 let requests = Arc::new(StdMutex::new(Vec::<AgentInferenceRequest>::new()));
5520 let contexts = Arc::new(StdMutex::new(Vec::<InferenceRoutingContext>::new()));
5521 let selected = ModelSelection {
5522 provider: roder_api::catalog::PROVIDER_MOCK.to_string(),
5523 model: "mock".to_string(),
5524 };
5525
5526 let mut builder = ExtensionRegistryBuilder::new();
5527 builder.inference_engine(Arc::new(RoutingCaptureEngine {
5528 id: roder_api::catalog::PROVIDER_MOCK,
5529 models: vec![routing_test_model("mock", &[REASONING_LOW])],
5530 requests: requests.clone(),
5531 }));
5532 builder.inference_router(Arc::new(StaticRouter {
5533 id: "test-router",
5534 decision: InferenceRoutingDecision::selected(
5535 "test-router",
5536 ModelSelection {
5537 provider: "missing".to_string(),
5538 model: "missing".to_string(),
5539 },
5540 "would route if called",
5541 ),
5542 contexts: contexts.clone(),
5543 }));
5544 let thread_root =
5545 std::env::temp_dir().join(format!("roder-routing-manual-{}", uuid::Uuid::new_v4()));
5546 builder.thread_store_factory(Arc::new(JsonlThreadStoreFactory {
5547 base_path: thread_root.clone(),
5548 }));
5549 let runtime = Arc::new(
5550 Runtime::new(
5551 builder.build().unwrap(),
5552 RuntimeConfig {
5553 default_provider: selected.provider.clone(),
5554 default_model: selected.model.clone(),
5555 inference_router: RuntimeInferenceRouterConfig {
5556 enabled: true,
5557 router_id: Some("test-router".to_string()),
5558 },
5559 ..RuntimeConfig::default()
5560 },
5561 )
5562 .unwrap(),
5563 );
5564 let thread_id = runtime
5565 .create_thread_with(CreateThreadRequest {
5566 title: Some("Routing manual".to_string()),
5567 workspace: test_workspace(),
5568 workspace_id: None,
5569 root_id: None,
5570 provider: Some(selected.provider.clone()),
5571 model: Some(selected.model.clone()),
5572 tool_allowlist: Vec::new(),
5573 developer_instructions: None,
5574 external_tools: Vec::new(),
5575 selection_mode: Some(ModelSelectionMode::manual(
5576 selected.provider.clone(),
5577 selected.model.clone(),
5578 None,
5579 )),
5580 runner: None,
5581 })
5582 .await
5583 .unwrap()
5584 .thread_id;
5585 let mut rx = runtime.subscribe_events();
5586 let turn_id = runtime
5587 .start_turn(StartTurnRequest {
5588 thread_id,
5589 message: "use selected manual model".to_string(),
5590 images: Vec::new(),
5591 provider_override: None,
5592 model_override: None,
5593 reasoning_override: None,
5594 workspace: test_workspace(),
5595 instructions: InstructionBundle::default(),
5596 developer_context: None,
5597 task_ledger_required: false,
5598 })
5599 .await
5600 .unwrap();
5601
5602 let mut saw_routing_event = false;
5603 tokio::time::timeout(std::time::Duration::from_secs(5), async {
5604 loop {
5605 let envelope = rx.recv().await.unwrap();
5606 if envelope.turn_id.as_deref() != Some(&turn_id) {
5607 continue;
5608 }
5609 match envelope.event {
5610 RoderEvent::InferenceRoutingDecision(_) => {
5611 saw_routing_event = true;
5612 }
5613 RoderEvent::TurnCompleted(_) => break,
5614 RoderEvent::TurnFailed(event) => panic!("turn failed: {}", event.error),
5615 _ => {}
5616 }
5617 }
5618 })
5619 .await
5620 .unwrap();
5621
5622 assert!(!saw_routing_event);
5623 assert!(contexts.lock().unwrap().is_empty());
5624 let requests = requests.lock().unwrap();
5625 assert_eq!(requests.len(), 1);
5626 assert_eq!(requests[0].model, selected);
5627 let _ = std::fs::remove_dir_all(thread_root);
5628 }
5629
5630 #[test]
5631 fn enabled_inference_router_requires_registered_router() {
5632 let mut builder = ExtensionRegistryBuilder::new();
5633 builder.inference_engine(Arc::new(FakeInferenceEngine));
5634
5635 let err = match Runtime::new(
5636 builder.build().unwrap(),
5637 RuntimeConfig {
5638 inference_router: RuntimeInferenceRouterConfig {
5639 enabled: true,
5640 router_id: Some("missing-router".to_string()),
5641 },
5642 ..RuntimeConfig::default()
5643 },
5644 ) {
5645 Ok(_) => panic!("runtime should reject unknown inference router"),
5646 Err(err) => err,
5647 };
5648
5649 assert!(
5650 err.to_string()
5651 .contains("inference router \"missing-router\" is not registered")
5652 );
5653 }
5654
5655 #[tokio::test]
5656 async fn model_profile_routes_request_knobs_to_next_inference() {
5657 let request = captured_profile_request(RuntimeConfig {
5658 default_provider: roder_api::catalog::PROVIDER_MOCK.to_string(),
5659 default_model: "gpt-5.5".to_string(),
5660 model_profiles: std::collections::HashMap::from([(
5661 "gpt-5.5".to_string(),
5662 test_model_profile("gpt-5.5"),
5663 )]),
5664 ..RuntimeConfig::default()
5665 })
5666 .await;
5667
5668 let tool_names = request
5669 .tools
5670 .iter()
5671 .map(|tool| tool.name.as_str())
5672 .collect::<Vec<_>>();
5673 assert!(tool_names.contains(&"apply_patch"));
5674 assert!(tool_names.contains(&"edit"));
5675 assert!(tool_names.contains(&"multi_edit"));
5676 assert!(tool_names.contains(&"write_file"));
5677 assert_eq!(request.reasoning.level.as_deref(), Some(REASONING_LOW));
5678 assert_eq!(request.runtime.parallel_tool_calls, Some(false));
5679 assert_eq!(request.runtime.auto_compact_token_limit, Some(123_000));
5680 assert!(
5681 request
5682 .instructions
5683 .developer
5684 .as_deref()
5685 .unwrap_or_default()
5686 .contains("Use the provided context as the current working set")
5687 );
5688 assert_eq!(
5689 request
5690 .metadata
5691 .pointer("/modelProfile/schemaPolicy")
5692 .and_then(serde_json::Value::as_str),
5693 Some("standard_required_first")
5694 );
5695 }
5696
5697 #[tokio::test]
5698 async fn turn_developer_context_reaches_inference_and_does_not_persist() {
5699 let captured = Arc::new(StdMutex::new(None));
5700 let mut builder = ExtensionRegistryBuilder::new();
5701 builder.inference_engine(Arc::new(CapturingEngine {
5702 request: captured.clone(),
5703 }));
5704 let runtime = Arc::new(
5705 Runtime::new(
5706 builder.build().unwrap(),
5707 RuntimeConfig {
5708 default_provider: roder_api::catalog::PROVIDER_MOCK.to_string(),
5709 default_model: "gpt-5.5".to_string(),
5710 ..RuntimeConfig::default()
5711 },
5712 )
5713 .unwrap(),
5714 );
5715
5716 async fn run_turn(runtime: &Arc<Runtime>, developer_context: Option<String>) {
5717 let mut rx = runtime.subscribe_events();
5718 let turn_id = runtime
5719 .start_turn(StartTurnRequest {
5720 thread_id: "thread-turn-context".to_string(),
5721 message: "hello".to_string(),
5722 images: Vec::new(),
5723 provider_override: None,
5724 model_override: None,
5725 reasoning_override: None,
5726 workspace: test_workspace(),
5727 instructions: InstructionBundle::default(),
5728 developer_context,
5729 task_ledger_required: false,
5730 })
5731 .await
5732 .unwrap();
5733 tokio::time::timeout(std::time::Duration::from_secs(5), async {
5734 loop {
5735 let envelope = rx.recv().await.unwrap();
5736 if envelope.turn_id.as_deref() != Some(&turn_id) {
5737 continue;
5738 }
5739 match envelope.event {
5740 RoderEvent::TurnCompleted(_) => break,
5741 RoderEvent::TurnFailed(event) => panic!("turn failed: {}", event.error),
5742 _ => {}
5743 }
5744 }
5745 })
5746 .await
5747 .unwrap();
5748 }
5749
5750 run_turn(
5751 &runtime,
5752 Some("Connected accounts: example-service.".to_string()),
5753 )
5754 .await;
5755 let request = captured.lock().unwrap().clone().unwrap();
5756 assert_eq!(
5757 request.instructions.developer_context.as_deref(),
5758 Some("Connected accounts: example-service.")
5759 );
5760
5761 run_turn(&runtime, None).await;
5764 let request = captured.lock().unwrap().clone().unwrap();
5765 assert_eq!(request.instructions.developer_context, None);
5766 }
5767
5768 #[tokio::test]
5769 async fn tool_search_overrides_route_to_next_inference_request() {
5770 let request = captured_profile_request(RuntimeConfig {
5771 default_provider: roder_api::catalog::PROVIDER_MOCK.to_string(),
5772 default_model: "gpt-5.4".to_string(),
5773 tool_search: ToolSearchConfig {
5774 mode: roder_api::inference::ToolSearchMode::Auto,
5775 max_catalog_items: Some(100),
5776 ..ToolSearchConfig::default()
5777 },
5778 provider_tool_search: std::collections::HashMap::from([(
5779 roder_api::catalog::PROVIDER_MOCK.to_string(),
5780 roder_api::inference::ToolSearchConfigOverlay {
5781 include_skills: Some(false),
5782 provider_variant: Some(roder_api::inference::ToolSearchProviderVariant::Regex),
5783 ..Default::default()
5784 },
5785 )]),
5786 model_tool_search: std::collections::HashMap::from([(
5787 "gpt-5.4".to_string(),
5788 roder_api::inference::ToolSearchConfigOverlay {
5789 mode: Some(roder_api::inference::ToolSearchMode::ProviderNative),
5790 max_catalog_items: Some(25),
5791 provider_variant: Some(roder_api::inference::ToolSearchProviderVariant::Bm25),
5792 ..Default::default()
5793 },
5794 )]),
5795 ..RuntimeConfig::default()
5796 })
5797 .await;
5798
5799 assert_eq!(
5800 request.runtime.tool_search.mode,
5801 roder_api::inference::ToolSearchMode::ProviderNative
5802 );
5803 assert_eq!(request.runtime.tool_search.max_catalog_items, Some(25));
5804 assert_eq!(
5805 request.runtime.tool_search.provider_variant,
5806 roder_api::inference::ToolSearchProviderVariant::Bm25
5807 );
5808 }
5809
5810 #[tokio::test]
5811 async fn context_entrypoint_hints_use_turn_workspace() {
5812 let process_workspace = runtime_test_workspace("entrypoint-process");
5813 let thread_workspace = runtime_test_workspace("entrypoint-thread");
5814 std::fs::create_dir_all(process_workspace.join("src")).unwrap();
5815 std::fs::create_dir_all(thread_workspace.join("src")).unwrap();
5816 std::fs::write(
5817 process_workspace.join("src/sidebar-thread-groups.ts"),
5818 "export const desktopLeak = true;\n",
5819 )
5820 .unwrap();
5821 std::fs::write(
5822 thread_workspace.join("src/voice-plan-feedback.ts"),
5823 "export const voicePlanFeedback = true;\n",
5824 )
5825 .unwrap();
5826
5827 let captured = Arc::new(StdMutex::new(None));
5828 let mut builder = ExtensionRegistryBuilder::new();
5829 builder.inference_engine(Arc::new(CapturingEngine {
5830 request: captured.clone(),
5831 }));
5832 builder.context_planner(Arc::new(roder_context::EntrypointContextPlanner::new(
5833 process_workspace.clone(),
5834 )));
5835 let runtime = Arc::new(
5836 Runtime::new(
5837 builder.build().unwrap(),
5838 RuntimeConfig {
5839 workspace: Some(process_workspace.display().to_string()),
5840 ..RuntimeConfig::default()
5841 },
5842 )
5843 .unwrap(),
5844 );
5845 let mut rx = runtime.subscribe_events();
5846
5847 let turn_id = runtime
5848 .start_turn(StartTurnRequest {
5849 thread_id: "thread-workspace-entrypoint".to_string(),
5850 message: "investigate voice plan feedback".to_string(),
5851 images: Vec::new(),
5852 provider_override: None,
5853 model_override: None,
5854 reasoning_override: None,
5855 workspace: thread_workspace.display().to_string(),
5856 instructions: crate::instructions::default_instructions(),
5857 developer_context: None,
5858 task_ledger_required: false,
5859 })
5860 .await
5861 .unwrap();
5862
5863 tokio::time::timeout(std::time::Duration::from_secs(5), async {
5864 loop {
5865 let envelope = rx.recv().await.unwrap();
5866 if envelope.turn_id.as_deref() != Some(&turn_id) {
5867 continue;
5868 }
5869 match envelope.event {
5870 RoderEvent::TurnCompleted(_) => break,
5871 RoderEvent::TurnFailed(event) => panic!("turn failed: {}", event.error),
5872 _ => {}
5873 }
5874 }
5875 })
5876 .await
5877 .unwrap();
5878
5879 let request = captured.lock().unwrap().clone().expect("captured request");
5880 let transcript_text = request
5881 .transcript
5882 .iter()
5883 .map(|item| match item {
5884 TranscriptItem::UserMessage(message) => message.text.as_str(),
5885 _ => "",
5886 })
5887 .collect::<Vec<_>>()
5888 .join("\n");
5889 assert!(transcript_text.contains("src/voice-plan-feedback.ts"));
5890 assert!(!transcript_text.contains("src/sidebar-thread-groups.ts"));
5891
5892 let _ = std::fs::remove_dir_all(process_workspace);
5893 let _ = std::fs::remove_dir_all(thread_workspace);
5894 }
5895
5896 fn runtime_test_workspace(name: &str) -> std::path::PathBuf {
5897 let path =
5898 std::env::temp_dir().join(format!("roder-runtime-{name}-{}", uuid::Uuid::new_v4()));
5899 let _ = std::fs::remove_dir_all(&path);
5900 std::fs::create_dir_all(&path).unwrap();
5901 path
5902 }
5903
5904 #[tokio::test]
5905 async fn model_profile_user_model_knobs_override_profile_defaults() {
5906 let request = captured_profile_request(RuntimeConfig {
5907 default_provider: roder_api::catalog::PROVIDER_MOCK.to_string(),
5908 default_model: "gpt-5.5".to_string(),
5909 reasoning: Some(REASONING_HIGH.to_string()),
5910 auto_compact_token_limit: Some(999),
5911 model_edit_tools: std::collections::HashMap::from([(
5912 "gpt-5.5".to_string(),
5913 EDIT_TOOL_PATCH.to_string(),
5914 )]),
5915 model_parallel_tool_calls: std::collections::HashMap::from([(
5916 "gpt-5.5".to_string(),
5917 true,
5918 )]),
5919 model_profiles: std::collections::HashMap::from([(
5920 "gpt-5.5".to_string(),
5921 test_model_profile("gpt-5.5"),
5922 )]),
5923 ..RuntimeConfig::default()
5924 })
5925 .await;
5926
5927 let tool_names = request
5928 .tools
5929 .iter()
5930 .map(|tool| tool.name.as_str())
5931 .collect::<Vec<_>>();
5932 assert!(tool_names.contains(&"apply_patch"));
5933 assert!(!tool_names.contains(&"edit"));
5934 assert_eq!(request.reasoning.level.as_deref(), Some(REASONING_HIGH));
5935 assert_eq!(request.runtime.parallel_tool_calls, Some(true));
5936 assert_eq!(request.runtime.auto_compact_token_limit, Some(999));
5937 }
5938
5939 #[tokio::test]
5940 async fn runtime_tool_allowlist_filters_advertised_tools() {
5941 let request = captured_profile_request(RuntimeConfig {
5942 default_provider: roder_api::catalog::PROVIDER_MOCK.to_string(),
5943 default_model: "gpt-5.5".to_string(),
5944 tool_allowlist: vec!["edit".to_string()],
5945 model_profiles: std::collections::HashMap::from([(
5946 "gpt-5.5".to_string(),
5947 test_model_profile("gpt-5.5"),
5948 )]),
5949 ..RuntimeConfig::default()
5950 })
5951 .await;
5952
5953 let tool_names = request
5954 .tools
5955 .iter()
5956 .map(|tool| tool.name.as_str())
5957 .collect::<Vec<_>>();
5958 assert_eq!(tool_names, vec!["edit"]);
5959 }
5960
5961 async fn captured_thread_override_request(
5964 runtime: &Arc<Runtime>,
5965 requests: &Arc<StdMutex<Vec<AgentInferenceRequest>>>,
5966 tool_allowlist: Vec<String>,
5967 developer_instructions: Option<String>,
5968 external_tools: Vec<ToolSpec>,
5969 ) -> AgentInferenceRequest {
5970 let thread_id = runtime
5971 .create_thread_with(CreateThreadRequest {
5972 title: Some("Thread overrides".to_string()),
5973 workspace: test_workspace(),
5974 workspace_id: None,
5975 root_id: None,
5976 provider: None,
5977 model: None,
5978 selection_mode: None,
5979 tool_allowlist,
5980 developer_instructions,
5981 external_tools,
5982 runner: None,
5983 })
5984 .await
5985 .unwrap()
5986 .thread_id;
5987 let mut rx = runtime.subscribe_events();
5988 let turn_id = runtime
5989 .start_turn(StartTurnRequest {
5990 thread_id,
5991 message: "hello".to_string(),
5992 images: Vec::new(),
5993 provider_override: None,
5994 model_override: None,
5995 reasoning_override: None,
5996 workspace: test_workspace(),
5997 instructions: crate::default_instructions(),
5998 developer_context: None,
5999 task_ledger_required: false,
6000 })
6001 .await
6002 .unwrap();
6003 tokio::time::timeout(std::time::Duration::from_secs(5), async {
6004 loop {
6005 let envelope = rx.recv().await.unwrap();
6006 if envelope.turn_id.as_deref() != Some(&turn_id) {
6007 continue;
6008 }
6009 match envelope.event {
6010 RoderEvent::TurnCompleted(_) => break,
6011 RoderEvent::TurnFailed(event) => panic!("turn failed: {}", event.error),
6012 _ => {}
6013 }
6014 }
6015 })
6016 .await
6017 .unwrap();
6018 requests.lock().unwrap().pop().expect("captured request")
6019 }
6020
6021 #[tokio::test]
6022 async fn thread_tool_allowlist_filters_only_that_thread() {
6023 let requests = Arc::new(StdMutex::new(Vec::new()));
6024 let thread_root =
6025 std::env::temp_dir().join(format!("roder-thread-allowlist-{}", uuid::Uuid::new_v4()));
6026 let mut builder = ExtensionRegistryBuilder::new();
6027 builder.inference_engine(Arc::new(SwitchCaptureEngine {
6028 requests: requests.clone(),
6029 }));
6030 builder.thread_store_factory(Arc::new(JsonlThreadStoreFactory {
6031 base_path: thread_root.clone(),
6032 }));
6033 builder.tool_contributor(Arc::new(ProfileToolContributor));
6034 let runtime = Arc::new(
6035 Runtime::new(
6036 builder.build().unwrap(),
6037 RuntimeConfig {
6038 default_provider: roder_api::catalog::PROVIDER_MOCK.to_string(),
6039 default_model: "gpt-5.5".to_string(),
6040 model_profiles: std::collections::HashMap::from([(
6041 "gpt-5.5".to_string(),
6042 test_model_profile("gpt-5.5"),
6043 )]),
6044 ..RuntimeConfig::default()
6045 },
6046 )
6047 .unwrap(),
6048 );
6049
6050 let allowlisted = captured_thread_override_request(
6051 &runtime,
6052 &requests,
6053 vec!["edit".to_string()],
6054 None,
6055 Vec::new(),
6056 )
6057 .await;
6058 let unrestricted =
6059 captured_thread_override_request(&runtime, &requests, Vec::new(), None, Vec::new())
6060 .await;
6061
6062 let allowlisted_names = allowlisted
6063 .tools
6064 .iter()
6065 .map(|tool| tool.name.as_str())
6066 .collect::<Vec<_>>();
6067 assert_eq!(allowlisted_names, vec!["edit"]);
6068 let unrestricted_names = unrestricted
6069 .tools
6070 .iter()
6071 .map(|tool| tool.name.as_str())
6072 .collect::<Vec<_>>();
6073 assert!(unrestricted_names.contains(&"edit"));
6074 assert!(unrestricted_names.len() > 1);
6075
6076 let _ = std::fs::remove_dir_all(thread_root);
6077 }
6078
6079 fn runtime_with_edit_allowlist(
6081 requests: &Arc<StdMutex<Vec<AgentInferenceRequest>>>,
6082 thread_root: &std::path::Path,
6083 ) -> Arc<Runtime> {
6084 let mut builder = ExtensionRegistryBuilder::new();
6085 builder.inference_engine(Arc::new(SwitchCaptureEngine {
6086 requests: requests.clone(),
6087 }));
6088 builder.thread_store_factory(Arc::new(JsonlThreadStoreFactory {
6089 base_path: thread_root.to_path_buf(),
6090 }));
6091 builder.tool_contributor(Arc::new(ProfileToolContributor));
6092 Arc::new(
6093 Runtime::new(
6094 builder.build().unwrap(),
6095 RuntimeConfig {
6096 default_provider: roder_api::catalog::PROVIDER_MOCK.to_string(),
6097 default_model: "gpt-5.5".to_string(),
6098 tool_allowlist: vec!["edit".to_string()],
6099 model_profiles: std::collections::HashMap::from([(
6100 "gpt-5.5".to_string(),
6101 test_model_profile("gpt-5.5"),
6102 )]),
6103 ..RuntimeConfig::default()
6104 },
6105 )
6106 .unwrap(),
6107 )
6108 }
6109
6110 #[tokio::test]
6111 async fn runtime_and_thread_allowlists_intersect() {
6112 let requests = Arc::new(StdMutex::new(Vec::new()));
6113 let thread_root = std::env::temp_dir().join(format!(
6114 "roder-allowlist-intersect-{}",
6115 uuid::Uuid::new_v4()
6116 ));
6117 let runtime = runtime_with_edit_allowlist(&requests, &thread_root);
6118
6119 let request = captured_thread_override_request(
6120 &runtime,
6121 &requests,
6122 vec!["edit".to_string(), "write_file".to_string()],
6123 None,
6124 Vec::new(),
6125 )
6126 .await;
6127
6128 let names = request
6130 .tools
6131 .iter()
6132 .map(|tool| tool.name.as_str())
6133 .collect::<Vec<_>>();
6134 assert_eq!(names, vec!["edit"]);
6135
6136 let _ = std::fs::remove_dir_all(thread_root);
6137 }
6138
6139 #[tokio::test]
6140 async fn route_tool_call_denies_tools_outside_allowlists() {
6141 let requests = Arc::new(StdMutex::new(Vec::new()));
6142 let thread_root =
6143 std::env::temp_dir().join(format!("roder-allowlist-dispatch-{}", uuid::Uuid::new_v4()));
6144 let runtime = runtime_with_edit_allowlist(&requests, &thread_root);
6145 let thread_id = runtime
6146 .create_thread_with(CreateThreadRequest {
6147 title: Some("Dispatch allowlist".to_string()),
6148 workspace: test_workspace(),
6149 workspace_id: None,
6150 root_id: None,
6151 provider: None,
6152 model: None,
6153 selection_mode: None,
6154 tool_allowlist: vec!["edit".to_string(), "write_file".to_string()],
6155 developer_instructions: None,
6156 external_tools: Vec::new(),
6157 runner: None,
6158 })
6159 .await
6160 .unwrap()
6161 .thread_id;
6162
6163 let result = runtime
6165 .route_tool_call(
6166 &thread_id,
6167 &"turn-allowlist-dispatch".to_string(),
6168 roder_api::inference::ToolCallCompleted {
6169 id: "call-1".to_string(),
6170 name: "write_file".to_string(),
6171 arguments: r#"{"path":"a.txt","content":"hi"}"#.to_string(),
6172 },
6173 None,
6174 None,
6175 )
6176 .await
6177 .unwrap();
6178
6179 assert!(result.is_error);
6180 assert!(
6181 result
6182 .result
6183 .contains("not permitted by the tool allowlist"),
6184 "unexpected result: {}",
6185 result.result
6186 );
6187
6188 let _ = std::fs::remove_dir_all(thread_root);
6189 }
6190
6191 struct SignalledFailureEngine {
6193 started: tokio::sync::mpsc::UnboundedSender<()>,
6194 proceed: Arc<tokio::sync::Notify>,
6195 }
6196
6197 #[async_trait::async_trait]
6198 impl InferenceEngine for SignalledFailureEngine {
6199 fn id(&self) -> String {
6200 roder_api::catalog::PROVIDER_MOCK.to_string()
6201 }
6202
6203 fn capabilities(&self) -> InferenceCapabilities {
6204 InferenceCapabilities::coding_agent_default()
6205 }
6206
6207 async fn list_models(
6208 &self,
6209 _ctx: InferenceProviderContext<'_>,
6210 ) -> anyhow::Result<Vec<roder_api::inference::ModelDescriptor>> {
6211 Ok(roder_api::catalog::models_for_provider(
6212 roder_api::catalog::PROVIDER_MOCK,
6213 true,
6214 ))
6215 }
6216
6217 async fn stream_turn(
6218 &self,
6219 _ctx: InferenceTurnContext<'_>,
6220 _request: AgentInferenceRequest,
6221 ) -> anyhow::Result<InferenceEventStream> {
6222 let _ = self.started.send(());
6223 self.proceed.notified().await;
6224 anyhow::bail!("engine failed mid-turn")
6225 }
6226 }
6227
6228 #[tokio::test]
6229 async fn failed_turn_sweeps_pending_external_tool_calls() {
6230 let (started_tx, mut started_rx) = tokio::sync::mpsc::unbounded_channel();
6231 let proceed = Arc::new(tokio::sync::Notify::new());
6232 let mut builder = ExtensionRegistryBuilder::new();
6233 builder.inference_engine(Arc::new(SignalledFailureEngine {
6234 started: started_tx,
6235 proceed: proceed.clone(),
6236 }));
6237 let runtime =
6238 Arc::new(Runtime::new(builder.build().unwrap(), RuntimeConfig::default()).unwrap());
6239 let mut rx = runtime.subscribe_events();
6240 let turn_id = runtime
6241 .start_turn(StartTurnRequest {
6242 thread_id: "thread-sweep".to_string(),
6243 message: "go".to_string(),
6244 images: Vec::new(),
6245 provider_override: None,
6246 model_override: None,
6247 reasoning_override: None,
6248 workspace: test_workspace(),
6249 instructions: InstructionBundle {
6250 system: None,
6251 developer: None,
6252 developer_context: None,
6253 },
6254 developer_context: None,
6255 task_ledger_required: false,
6256 })
6257 .await
6258 .unwrap();
6259 tokio::time::timeout(std::time::Duration::from_secs(5), started_rx.recv())
6260 .await
6261 .unwrap()
6262 .unwrap();
6263
6264 let (tx, _pending_rx) = oneshot::channel();
6265 runtime.pending_external_tool_calls.lock().await.insert(
6266 "exttool-sweep-test".to_string(),
6267 PendingExternalToolCall {
6268 thread_id: "thread-sweep".to_string(),
6269 turn_id: turn_id.clone(),
6270 tool_id: "call-1".to_string(),
6271 tool_name: "acme_lookup".to_string(),
6272 tx,
6273 },
6274 );
6275 proceed.notify_one();
6276
6277 let outcome = tokio::time::timeout(std::time::Duration::from_secs(5), async {
6278 loop {
6279 let envelope = rx.recv().await.unwrap();
6280 if let RoderEvent::ExternalToolCallResolved(event) = envelope.event
6281 && event.request_id == "exttool-sweep-test"
6282 {
6283 break event.outcome;
6284 }
6285 }
6286 })
6287 .await
6288 .expect("turn failure must resolve pending external tool calls");
6289 assert_eq!(outcome, ExternalToolCallOutcome::Cancelled);
6290 assert!(runtime.pending_external_tool_calls.lock().await.is_empty());
6291 }
6292
6293 #[tokio::test]
6294 async fn thread_developer_instructions_layer_under_harness_prompt() {
6295 let requests = Arc::new(StdMutex::new(Vec::new()));
6296 let thread_root = std::env::temp_dir().join(format!(
6297 "roder-thread-instructions-{}",
6298 uuid::Uuid::new_v4()
6299 ));
6300 let mut builder = ExtensionRegistryBuilder::new();
6301 builder.inference_engine(Arc::new(SwitchCaptureEngine {
6302 requests: requests.clone(),
6303 }));
6304 builder.thread_store_factory(Arc::new(JsonlThreadStoreFactory {
6305 base_path: thread_root.clone(),
6306 }));
6307 builder.tool_contributor(Arc::new(ProfileToolContributor));
6308 let runtime =
6309 Arc::new(Runtime::new(builder.build().unwrap(), RuntimeConfig::default()).unwrap());
6310
6311 let request = captured_thread_override_request(
6312 &runtime,
6313 &requests,
6314 Vec::new(),
6315 Some("You are embedded in a host app.".to_string()),
6316 Vec::new(),
6317 )
6318 .await;
6319
6320 let system = request.instructions.system.expect("system instructions");
6321 assert!(system.starts_with("You are Roder"));
6322 let developer = request
6323 .instructions
6324 .developer
6325 .expect("developer instructions");
6326 assert!(developer.starts_with("You are embedded in a host app."));
6327
6328 let plain =
6329 captured_thread_override_request(&runtime, &requests, Vec::new(), None, Vec::new())
6330 .await;
6331 assert_eq!(plain.instructions.developer, None);
6332
6333 let _ = std::fs::remove_dir_all(thread_root);
6334 }
6335
6336 #[tokio::test]
6337 async fn thread_external_tools_are_advertised_and_shadow_builtins() {
6338 let requests = Arc::new(StdMutex::new(Vec::new()));
6339 let thread_root = std::env::temp_dir().join(format!(
6340 "roder-thread-external-tools-{}",
6341 uuid::Uuid::new_v4()
6342 ));
6343 let mut builder = ExtensionRegistryBuilder::new();
6344 builder.inference_engine(Arc::new(SwitchCaptureEngine {
6345 requests: requests.clone(),
6346 }));
6347 builder.thread_store_factory(Arc::new(JsonlThreadStoreFactory {
6348 base_path: thread_root.clone(),
6349 }));
6350 builder.tool_contributor(Arc::new(ProfileToolContributor));
6351 let runtime =
6352 Arc::new(Runtime::new(builder.build().unwrap(), RuntimeConfig::default()).unwrap());
6353
6354 let external_tools = vec![
6355 ToolSpec {
6356 name: "acme_lookup".to_string(),
6357 description: "Look up Acme workspace state.".to_string(),
6358 parameters: serde_json::json!({
6359 "type": "object",
6360 "properties": { "query": { "type": "string" } },
6361 "required": ["query"]
6362 }),
6363 },
6364 ToolSpec {
6365 name: "edit".to_string(),
6366 description: "Host-managed edit.".to_string(),
6367 parameters: serde_json::json!({ "type": "object" }),
6368 },
6369 ];
6370 let request =
6371 captured_thread_override_request(&runtime, &requests, Vec::new(), None, external_tools)
6372 .await;
6373
6374 let acme = request
6375 .tools
6376 .iter()
6377 .find(|tool| tool.name == "acme_lookup")
6378 .expect("external tool advertised");
6379 assert_eq!(acme.description, "Look up Acme workspace state.");
6380 assert_eq!(acme.parameters["required"][0], "query");
6381 let edits = request
6382 .tools
6383 .iter()
6384 .filter(|tool| tool.name == "edit")
6385 .collect::<Vec<_>>();
6386 assert_eq!(edits.len(), 1, "external edit shadows the builtin");
6387 assert_eq!(edits[0].description, "Host-managed edit.");
6388
6389 let plain =
6390 captured_thread_override_request(&runtime, &requests, Vec::new(), None, Vec::new())
6391 .await;
6392 assert!(plain.tools.iter().all(|tool| tool.name != "acme_lookup"));
6393 let plain_edit = plain
6394 .tools
6395 .iter()
6396 .find(|tool| tool.name == "edit")
6397 .expect("builtin edit advertised on plain thread");
6398 assert_eq!(plain_edit.description, "edit test tool");
6399
6400 let _ = std::fs::remove_dir_all(thread_root);
6401 }
6402
6403 #[tokio::test]
6404 async fn model_switch_injects_summary_and_records_profile_segments() {
6405 let requests = Arc::new(StdMutex::new(Vec::new()));
6406 let thread_root = std::env::temp_dir().join(format!(
6407 "roder-model-switch-thread-{}",
6408 uuid::Uuid::new_v4()
6409 ));
6410 let mut builder = ExtensionRegistryBuilder::new();
6411 builder.inference_engine(Arc::new(SwitchCaptureEngine {
6412 requests: requests.clone(),
6413 }));
6414 builder.thread_store_factory(Arc::new(JsonlThreadStoreFactory {
6415 base_path: thread_root.clone(),
6416 }));
6417 builder.tool_contributor(Arc::new(ProfileToolContributor));
6418 let mut claude_profile = test_model_profile("claude-haiku-4-5-20251001");
6419 claude_profile.provider_family = ProviderFamily::Anthropic;
6420 claude_profile.edit_tool = Some(EDIT_TOOL_EDIT.to_string());
6421 let runtime = Arc::new(
6422 Runtime::new(
6423 builder.build().unwrap(),
6424 RuntimeConfig {
6425 default_provider: roder_api::catalog::PROVIDER_MOCK.to_string(),
6426 default_model: "gpt-5.5".to_string(),
6427 model_profiles: std::collections::HashMap::from([
6428 ("gpt-5.5".to_string(), test_model_profile("gpt-5.5")),
6429 ("claude-haiku-4-5-20251001".to_string(), claude_profile),
6430 ]),
6431 ..RuntimeConfig::default()
6432 },
6433 )
6434 .unwrap(),
6435 );
6436 let thread_id = runtime
6437 .create_thread_with(CreateThreadRequest {
6438 title: Some("Model switch".to_string()),
6439 workspace: test_workspace(),
6440 workspace_id: None,
6441 root_id: None,
6442 provider: None,
6443 model: None,
6444 selection_mode: None,
6445 tool_allowlist: Vec::new(),
6446 developer_instructions: None,
6447 external_tools: Vec::new(),
6448 runner: None,
6449 })
6450 .await
6451 .unwrap()
6452 .thread_id;
6453 let mut rx = runtime.subscribe_events();
6454 for (message, model_override) in [
6455 ("first turn", None),
6456 ("second turn", Some("claude-haiku-4-5-20251001".to_string())),
6457 ] {
6458 let turn_id = runtime
6459 .start_turn(StartTurnRequest {
6460 thread_id: thread_id.clone(),
6461 message: message.to_string(),
6462 images: Vec::new(),
6463 provider_override: None,
6464 model_override,
6465 reasoning_override: None,
6466 workspace: test_workspace(),
6467 instructions: InstructionBundle::default(),
6468 developer_context: None,
6469 task_ledger_required: false,
6470 })
6471 .await
6472 .unwrap();
6473 tokio::time::timeout(std::time::Duration::from_secs(5), async {
6474 loop {
6475 let envelope = rx.recv().await.unwrap();
6476 if envelope.turn_id.as_deref() != Some(&turn_id) {
6477 continue;
6478 }
6479 match envelope.event {
6480 RoderEvent::TurnCompleted(_) => break,
6481 RoderEvent::TurnFailed(event) => panic!("turn failed: {}", event.error),
6482 _ => {}
6483 }
6484 }
6485 })
6486 .await
6487 .unwrap();
6488 }
6489
6490 let captured = requests.lock().unwrap().clone();
6491 assert_eq!(captured.len(), 2);
6492 assert!(captured[1].transcript.iter().any(|item| {
6493 matches!(
6494 item,
6495 TranscriptItem::UserMessage(message)
6496 if message.text.starts_with(MODEL_SWITCH_SUMMARY_PREFIX)
6497 && message.text.contains("previous profile mock/gpt-5.5")
6498 && message.text.contains("Current profile mock/claude-haiku-4-5-20251001")
6499 && message.text.contains("Available tools now:")
6500 )
6501 }));
6502
6503 let snapshot = runtime
6504 .thread_store
6505 .as_ref()
6506 .unwrap()
6507 .load_thread(&thread_id)
6508 .await
6509 .unwrap()
6510 .unwrap();
6511 let trace_segments = snapshot
6512 .turns
6513 .iter()
6514 .flat_map(|turn| &turn.items)
6515 .filter(|item| {
6516 matches!(
6517 item,
6518 TranscriptItem::ProviderMetadata(value)
6519 if value.get("kind").and_then(serde_json::Value::as_str)
6520 == Some(MODEL_PROFILE_TRACE_KIND)
6521 && value.get("segment").and_then(serde_json::Value::as_str)
6522 == Some("assistant")
6523 )
6524 })
6525 .count();
6526 assert!(trace_segments >= 2);
6527 let _ = std::fs::remove_dir_all(thread_root);
6528 }
6529
6530 struct CountingTaskTool {
6531 calls: Arc<StdMutex<u32>>,
6532 }
6533
6534 #[async_trait::async_trait]
6535 impl ToolExecutor for CountingTaskTool {
6536 fn spec(&self) -> ToolSpec {
6537 ToolSpec {
6538 name: "task".to_string(),
6539 description: "Dispatch a test subagent.".to_string(),
6540 parameters: serde_json::json!({
6541 "type": "object",
6542 "properties": {
6543 "description": { "type": "string" },
6544 "prompt": { "type": "string" },
6545 "parent_deadline_seconds": { "type": "integer" }
6546 },
6547 "required": ["description", "prompt"],
6548 "additionalProperties": false
6549 }),
6550 }
6551 }
6552
6553 async fn execute(
6554 &self,
6555 _ctx: ToolExecutionContext,
6556 call: ToolCall,
6557 ) -> anyhow::Result<ToolResult> {
6558 *self.calls.lock().unwrap() += 1;
6559 Ok(ToolResult {
6560 id: call.id,
6561 name: call.name,
6562 text: "started child".to_string(),
6563 data: serde_json::json!({}),
6564 is_error: false,
6565 })
6566 }
6567 }
6568
6569 #[tokio::test]
6570 async fn runtime_profile_reaches_inference_request_and_turn_metadata() {
6571 let captured = Arc::new(StdMutex::new(None));
6572 let mut builder = ExtensionRegistryBuilder::new();
6573 builder.inference_engine(Arc::new(CapturingEngine {
6574 request: captured.clone(),
6575 }));
6576 let runtime = Arc::new(
6577 Runtime::new(
6578 builder.build().unwrap(),
6579 RuntimeConfig {
6580 runtime_profile: RuntimeProfile::NonInteractive,
6581 ..RuntimeConfig::default()
6582 },
6583 )
6584 .unwrap(),
6585 );
6586 let mut rx = runtime.subscribe_events();
6587 let turn_id = runtime
6588 .start_turn(StartTurnRequest {
6589 thread_id: "thread-profile".to_string(),
6590 message: "work unattended".to_string(),
6591 images: Vec::new(),
6592 provider_override: None,
6593 model_override: None,
6594 reasoning_override: None,
6595 workspace: test_workspace(),
6596 instructions: InstructionBundle {
6597 system: None,
6598 developer: Some("base developer".to_string()),
6599 developer_context: None,
6600 },
6601 developer_context: None,
6602 task_ledger_required: false,
6603 })
6604 .await
6605 .unwrap();
6606
6607 let mut observed_profile = None;
6608 tokio::time::timeout(std::time::Duration::from_secs(5), async {
6609 loop {
6610 let envelope = rx.recv().await.unwrap();
6611 if envelope.turn_id.as_deref() != Some(&turn_id) {
6612 continue;
6613 }
6614 match envelope.event {
6615 RoderEvent::TurnStarted(event) => {
6616 observed_profile = Some(event.runtime_profile);
6617 }
6618 RoderEvent::TurnCompleted(_) => break,
6619 RoderEvent::TurnFailed(event) => panic!("turn failed: {}", event.error),
6620 _ => {}
6621 }
6622 }
6623 })
6624 .await
6625 .unwrap();
6626
6627 assert_eq!(observed_profile, Some(RuntimeProfile::NonInteractive));
6628 let request = captured.lock().unwrap().clone().unwrap();
6629 assert_eq!(request.runtime.profile, RuntimeProfile::NonInteractive);
6630 let developer = request.instructions.developer.unwrap();
6631 assert!(developer.contains("base developer"));
6632 assert!(developer.contains("non-interactive profile"));
6633 }
6634
6635 #[tokio::test]
6636 async fn global_policy_mode_changes_do_not_create_runtime_thread_directory() {
6637 let workspace = runtime_test_workspace("global-policy-mode");
6638 let thread_root = workspace.join("threads");
6639 let mut builder = ExtensionRegistryBuilder::new();
6640 builder.inference_engine(Arc::new(FakeInferenceEngine));
6641 builder.thread_store_factory(Arc::new(JsonlThreadStoreFactory {
6642 base_path: thread_root.clone(),
6643 }));
6644 let runtime = Runtime::new(
6645 builder.build().unwrap(),
6646 RuntimeConfig {
6647 workspace: Some(workspace.display().to_string()),
6648 ..Default::default()
6649 },
6650 )
6651 .unwrap();
6652
6653 runtime
6654 .set_policy_mode(PolicyMode::AcceptAll, Some("test".to_string()))
6655 .await
6656 .unwrap();
6657
6658 assert!(!thread_root.join("runtime").exists());
6659 let _ = std::fs::remove_dir_all(workspace);
6660 }
6661
6662 #[tokio::test]
6663 async fn task_ledger_enforcement_injects_eval_reminder_before_work() {
6664 let captured = Arc::new(StdMutex::new(None));
6665 let mut builder = ExtensionRegistryBuilder::new();
6666 builder.inference_engine(Arc::new(CapturingEngine {
6667 request: captured.clone(),
6668 }));
6669 builder.tool_contributor(Arc::new(
6670 roder_ext_task_ledger::TaskLedgerToolContributor::default(),
6671 ));
6672 let runtime = Arc::new(
6673 Runtime::new(
6674 builder.build().unwrap(),
6675 RuntimeConfig {
6676 runtime_profile: RuntimeProfile::Eval,
6677 policy_mode: PolicyMode::Bypass,
6678 ..RuntimeConfig::default()
6679 },
6680 )
6681 .unwrap(),
6682 );
6683 let mut rx = runtime.subscribe_events();
6684 let turn_id = runtime
6685 .start_turn(StartTurnRequest {
6686 thread_id: "thread-ledger".to_string(),
6687 message: "decomposed work".to_string(),
6688 images: Vec::new(),
6689 provider_override: None,
6690 model_override: None,
6691 reasoning_override: None,
6692 workspace: test_workspace(),
6693 instructions: InstructionBundle::default(),
6694 developer_context: None,
6695 task_ledger_required: true,
6696 })
6697 .await
6698 .unwrap();
6699
6700 tokio::time::timeout(std::time::Duration::from_secs(5), async {
6701 loop {
6702 let envelope = rx.recv().await.unwrap();
6703 if envelope.turn_id.as_deref() == Some(&turn_id)
6704 && matches!(envelope.event, RoderEvent::TurnCompleted(_))
6705 {
6706 break;
6707 }
6708 }
6709 })
6710 .await
6711 .unwrap();
6712
6713 let request = captured.lock().unwrap().clone().unwrap();
6714 let developer = request.instructions.developer.unwrap();
6715 assert!(developer.contains("Task Ledger Required"));
6716 assert!(developer.contains("task_ledger.update"));
6717 let tool_names: Vec<_> = request
6718 .tools
6719 .iter()
6720 .map(|tool| tool.name.as_str())
6721 .collect();
6722 assert!(
6723 tool_names.contains(&TASK_LEDGER_TOOL_NAME),
6724 "tool names: {tool_names:?}"
6725 );
6726 assert_eq!(
6727 request.tool_choice,
6728 ToolChoice::Specific(TASK_LEDGER_TOOL_NAME.to_string())
6729 );
6730 assert_eq!(request.tools.len(), 1);
6731 assert_eq!(request.tools[0].name, TASK_LEDGER_TOOL_NAME);
6732 }
6733
6734 #[tokio::test]
6735 async fn eval_task_ledger_blocks_final_answer_until_open_items_are_completed() {
6736 let requests = Arc::new(StdMutex::new(Vec::new()));
6737 let mut builder = ExtensionRegistryBuilder::new();
6738 builder.inference_engine(Arc::new(TaskLedgerCompletionGateEngine {
6739 calls: StdMutex::new(0),
6740 requests: requests.clone(),
6741 }));
6742 builder.tool_contributor(Arc::new(
6743 roder_ext_task_ledger::TaskLedgerToolContributor::default(),
6744 ));
6745 let runtime = Arc::new(
6746 Runtime::new(
6747 builder.build().unwrap(),
6748 RuntimeConfig {
6749 runtime_profile: RuntimeProfile::Eval,
6750 policy_mode: PolicyMode::Bypass,
6751 ..RuntimeConfig::default()
6752 },
6753 )
6754 .unwrap(),
6755 );
6756 let mut rx = runtime.subscribe_events();
6757 let turn_id = runtime
6758 .start_turn(StartTurnRequest {
6759 thread_id: "thread-ledger-completion".to_string(),
6760 message: "write the answer file".to_string(),
6761 images: Vec::new(),
6762 provider_override: None,
6763 model_override: None,
6764 reasoning_override: None,
6765 workspace: test_workspace(),
6766 instructions: InstructionBundle::default(),
6767 developer_context: None,
6768 task_ledger_required: true,
6769 })
6770 .await
6771 .unwrap();
6772
6773 tokio::time::timeout(std::time::Duration::from_secs(5), async {
6774 loop {
6775 let envelope = rx.recv().await.unwrap();
6776 if envelope.turn_id.as_deref() != Some(&turn_id) {
6777 continue;
6778 }
6779 match envelope.event {
6780 RoderEvent::TurnCompleted(_) => break,
6781 RoderEvent::TurnFailed(event) => panic!("turn failed: {}", event.error),
6782 _ => {}
6783 }
6784 }
6785 })
6786 .await
6787 .unwrap();
6788
6789 let requests = requests.lock().unwrap().clone();
6790 assert_eq!(requests.len(), 4);
6791 assert!(requests[2].transcript.iter().any(|item| {
6792 matches!(
6793 item,
6794 TranscriptItem::UserMessage(message)
6795 if message.text.contains("Task Ledger Completion Required")
6796 && message.text.contains("Write /app/result.txt")
6797 )
6798 }));
6799 assert!(requests[3].transcript.iter().any(|item| {
6800 matches!(
6801 item,
6802 TranscriptItem::ToolResult(result)
6803 if result.name.as_deref() == Some(TASK_LEDGER_TOOL_NAME)
6804 && result.result.contains("Task ledger: 2/2 completed")
6805 )
6806 }));
6807 }
6808
6809 #[tokio::test]
6810 async fn eval_task_ledger_checkpoint_requests_scoreable_file_before_final_reserve() {
6811 let requests = Arc::new(StdMutex::new(Vec::new()));
6812 let mut builder = ExtensionRegistryBuilder::new();
6813 builder.inference_engine(Arc::new(TaskLedgerCompletionGateEngine {
6814 calls: StdMutex::new(0),
6815 requests: requests.clone(),
6816 }));
6817 builder.tool_contributor(Arc::new(
6818 roder_ext_task_ledger::TaskLedgerToolContributor::default(),
6819 ));
6820 let runtime = Arc::new(
6821 Runtime::new(
6822 builder.build().unwrap(),
6823 RuntimeConfig {
6824 runtime_profile: RuntimeProfile::Eval,
6825 policy_mode: PolicyMode::Bypass,
6826 turn_deadline_seconds: Some(120),
6827 ..RuntimeConfig::default()
6828 },
6829 )
6830 .unwrap(),
6831 );
6832 let mut rx = runtime.subscribe_events();
6833 let turn_id = runtime
6834 .start_turn(StartTurnRequest {
6835 thread_id: "thread-ledger-checkpoint".to_string(),
6836 message: "write the answer file".to_string(),
6837 images: Vec::new(),
6838 provider_override: None,
6839 model_override: None,
6840 reasoning_override: None,
6841 workspace: test_workspace(),
6842 instructions: InstructionBundle::default(),
6843 developer_context: None,
6844 task_ledger_required: true,
6845 })
6846 .await
6847 .unwrap();
6848
6849 tokio::time::timeout(std::time::Duration::from_secs(5), async {
6850 loop {
6851 let envelope = rx.recv().await.unwrap();
6852 if envelope.turn_id.as_deref() != Some(&turn_id) {
6853 continue;
6854 }
6855 match envelope.event {
6856 RoderEvent::TurnCompleted(_) => break,
6857 RoderEvent::TurnFailed(event) => panic!("turn failed: {}", event.error),
6858 _ => {}
6859 }
6860 }
6861 })
6862 .await
6863 .unwrap();
6864
6865 let requests = requests.lock().unwrap().clone();
6866 assert!(requests.len() >= 2);
6867 assert!(requests[1].transcript.iter().any(|item| {
6868 matches!(
6869 item,
6870 TranscriptItem::UserMessage(message)
6871 if message.text.contains("Scoreable Output Checkpoint")
6872 && message.text.contains("ensure the required output file(s) exist")
6873 && message.text.contains("Write /app/result.txt")
6874 )
6875 }));
6876 }
6877
6878 #[test]
6879 fn deadline_task_ledger_prompt_preserves_scoreable_work_instruction() {
6880 let prompt = task_ledger_deadline_completion_prompt(
6881 12,
6882 30,
6883 "Task Ledger Completion Required: write /app/result.txt, then call task_ledger.update",
6884 );
6885
6886 assert!(prompt.contains("12 seconds remain"));
6887 assert!(prompt.contains("create or update the required scoreable output files"));
6888 assert!(prompt.contains("write /app/result.txt"));
6889 assert!(prompt.contains(TASK_LEDGER_TOOL_NAME));
6890 }
6891
6892 #[test]
6893 fn scoreable_checkpoint_prompt_preserves_provisional_file_instruction() {
6894 let prompt = task_ledger_scoreable_checkpoint_prompt(
6895 120,
6896 "Task Ledger Completion Required: write /app/result.txt, then call task_ledger.update",
6897 );
6898
6899 assert!(prompt.contains("120 seconds remain"));
6900 assert!(prompt.contains("best evidence-backed answer"));
6901 assert!(prompt.contains("even if provisional"));
6902 assert!(prompt.contains("preserve that candidate"));
6903 assert!(prompt.contains("partial-coverage"));
6904 assert!(prompt.contains("write /app/result.txt"));
6905 assert!(prompt.contains(TASK_LEDGER_TOOL_NAME));
6906 }
6907
6908 #[test]
6909 fn open_task_ledger_moves_inference_timeout_to_scoreable_checkpoint() {
6910 let deadline = Some(OffsetDateTime::now_utc() + Duration::seconds(870));
6911 let transcript = vec![TranscriptItem::ToolResult(ToolResultRecord {
6912 id: "ledger-open".to_string(),
6913 name: Some(TASK_LEDGER_TOOL_NAME.to_string()),
6914 result: "Task ledger: 0/1 completed\n- pending: Write /app/result.txt [write]"
6915 .to_string(),
6916 display_payload: None,
6917 is_error: false,
6918 })];
6919
6920 let (_, action) = inference_timeout_deadline(
6921 deadline,
6922 RuntimeProfile::Eval,
6923 true,
6924 30,
6925 false,
6926 0,
6927 &transcript,
6928 )
6929 .unwrap();
6930
6931 assert_eq!(action, InferenceTimeoutAction::ScoreableCheckpoint);
6932 }
6933
6934 #[tokio::test]
6935 async fn verification_gate_forces_eval_code_changes_through_review() {
6936 let mut builder = ExtensionRegistryBuilder::new();
6937 builder.inference_engine(Arc::new(VerificationGateEngine {
6938 calls: StdMutex::new(0),
6939 }));
6940 builder.tool_contributor(Arc::new(WriteFileContributor));
6941 builder.tool_contributor(Arc::new(
6942 roder_ext_verification::VerificationToolContributor,
6943 ));
6944 let runtime = Arc::new(
6945 Runtime::new(
6946 builder.build().unwrap(),
6947 RuntimeConfig {
6948 default_provider: roder_api::catalog::PROVIDER_MOCK.to_string(),
6949 default_model: "mock".to_string(),
6950 runtime_profile: RuntimeProfile::Eval,
6951 policy_mode: PolicyMode::Bypass,
6952 ..RuntimeConfig::default()
6953 },
6954 )
6955 .unwrap(),
6956 );
6957 let mut rx = runtime.subscribe_events();
6958 let turn_id = runtime
6959 .start_turn(StartTurnRequest {
6960 thread_id: "thread-verification".to_string(),
6961 message: "write code".to_string(),
6962 images: Vec::new(),
6963 provider_override: None,
6964 model_override: None,
6965 reasoning_override: None,
6966 workspace: test_workspace(),
6967 instructions: InstructionBundle::default(),
6968 developer_context: None,
6969 task_ledger_required: false,
6970 })
6971 .await
6972 .unwrap();
6973
6974 let mut saw_required = false;
6975 let mut saw_completed = false;
6976 let mut final_text = String::new();
6977 tokio::time::timeout(std::time::Duration::from_secs(5), async {
6978 loop {
6979 let envelope = rx.recv().await.unwrap();
6980 if envelope.turn_id.as_deref() != Some(&turn_id) {
6981 continue;
6982 }
6983 match envelope.event {
6984 RoderEvent::VerificationRequired(event) => {
6985 saw_required = true;
6986 assert_eq!(event.changed_files, vec!["src/lib.rs"]);
6987 }
6988 RoderEvent::VerificationCompleted(event) => {
6989 saw_completed = true;
6990 assert!(event.passed);
6991 }
6992 RoderEvent::InferenceEventReceived(event) => {
6993 if let InferenceEvent::MessageDelta(delta) = event.event {
6994 final_text.push_str(&delta.text);
6995 }
6996 }
6997 RoderEvent::TurnCompleted(_) => break,
6998 _ => {}
6999 }
7000 }
7001 })
7002 .await
7003 .unwrap();
7004
7005 assert!(saw_required);
7006 assert!(saw_completed);
7007 assert!(final_text.contains("verified final"));
7008 }
7009
7010 #[tokio::test]
7011 async fn speed_policy_changes_reasoning_across_eval_model_calls_without_model_switch() {
7012 let requests = Arc::new(StdMutex::new(Vec::new()));
7013 let mut builder = ExtensionRegistryBuilder::new();
7014 builder.inference_engine(Arc::new(SpeedPolicyEngine {
7015 calls: StdMutex::new(0),
7016 requests: requests.clone(),
7017 }));
7018 builder.tool_contributor(Arc::new(WriteFileContributor));
7019 builder.tool_contributor(Arc::new(
7020 roder_ext_verification::VerificationToolContributor,
7021 ));
7022 let runtime = Arc::new(
7023 Runtime::new(
7024 builder.build().unwrap(),
7025 RuntimeConfig {
7026 default_provider: roder_api::catalog::PROVIDER_MOCK.to_string(),
7027 default_model: "gpt-5.5".to_string(),
7028 runtime_profile: RuntimeProfile::Eval,
7029 policy_mode: PolicyMode::Bypass,
7030 ..RuntimeConfig::default()
7031 },
7032 )
7033 .unwrap(),
7034 );
7035 let mut rx = runtime.subscribe_events();
7036 let turn_id = runtime
7037 .start_turn(StartTurnRequest {
7038 thread_id: "thread-speed-policy".to_string(),
7039 message: "write code".to_string(),
7040 images: Vec::new(),
7041 provider_override: None,
7042 model_override: None,
7043 reasoning_override: None,
7044 workspace: test_workspace(),
7045 instructions: InstructionBundle::default(),
7046 developer_context: None,
7047 task_ledger_required: false,
7048 })
7049 .await
7050 .unwrap();
7051
7052 let mut saw_speed_policy_event = false;
7053 tokio::time::timeout(std::time::Duration::from_secs(5), async {
7054 loop {
7055 let envelope = rx.recv().await.unwrap();
7056 if envelope.turn_id.as_deref() != Some(&turn_id) {
7057 continue;
7058 }
7059 match envelope.event {
7060 RoderEvent::InferenceStarted(event) => {
7061 if event.speed_policy.is_some() {
7062 saw_speed_policy_event = true;
7063 }
7064 }
7065 RoderEvent::TurnCompleted(_) => break,
7066 RoderEvent::TurnFailed(event) => panic!("turn failed: {}", event.error),
7067 _ => {}
7068 }
7069 }
7070 })
7071 .await
7072 .unwrap();
7073
7074 let requests = requests.lock().unwrap().clone();
7075 assert!(saw_speed_policy_event);
7076 assert!(requests.len() >= 4);
7077 assert!(requests.iter().all(|request| {
7078 request.model.provider == roder_api::catalog::PROVIDER_MOCK
7079 && request.model.model == "gpt-5.5"
7080 }));
7081 assert_eq!(
7082 requests[0].runtime.speed_policy.as_ref().map(|d| d.phase),
7083 Some(roder_api::inference::SpeedPolicyPhase::Orientation)
7084 );
7085 assert_eq!(requests[0].reasoning.level.as_deref(), Some(REASONING_HIGH));
7086 assert_eq!(
7087 requests[1].runtime.speed_policy.as_ref().map(|d| d.phase),
7088 Some(roder_api::inference::SpeedPolicyPhase::Execution)
7089 );
7090 assert_eq!(requests[1].reasoning.level.as_deref(), Some(REASONING_LOW));
7091 assert_eq!(
7092 requests[2].runtime.speed_policy.as_ref().map(|d| d.phase),
7093 Some(roder_api::inference::SpeedPolicyPhase::Verification)
7094 );
7095 assert_eq!(requests[2].reasoning.level.as_deref(), Some(REASONING_HIGH));
7096 assert_eq!(
7097 requests[2]
7098 .metadata
7099 .pointer("/speedPolicy/phase")
7100 .and_then(serde_json::Value::as_str),
7101 Some("verification")
7102 );
7103 }
7104
7105 #[tokio::test]
7106 async fn deadline_turn_timeout_emits_partial_result_and_clears_active_turn() {
7107 let mut builder = ExtensionRegistryBuilder::new();
7108 builder.inference_engine(Arc::new(DeadlineEngine));
7109 let runtime = Arc::new(
7110 Runtime::new(
7111 builder.build().unwrap(),
7112 RuntimeConfig {
7113 default_provider: roder_api::catalog::PROVIDER_MOCK.to_string(),
7114 default_model: "mock".to_string(),
7115 runtime_profile: RuntimeProfile::Eval,
7116 turn_deadline_seconds: Some(1),
7117 ..RuntimeConfig::default()
7118 },
7119 )
7120 .unwrap(),
7121 );
7122 let mut rx = runtime.subscribe_events();
7123 let turn_id = runtime
7124 .start_turn(StartTurnRequest {
7125 thread_id: "thread-deadline".to_string(),
7126 message: "slow work".to_string(),
7127 images: Vec::new(),
7128 provider_override: None,
7129 model_override: None,
7130 reasoning_override: None,
7131 workspace: test_workspace(),
7132 instructions: InstructionBundle::default(),
7133 developer_context: None,
7134 task_ledger_required: false,
7135 })
7136 .await
7137 .unwrap();
7138
7139 let mut saw_partial = false;
7140 let mut saw_deadline = false;
7141 let mut failed_kind = None;
7142 tokio::time::timeout(std::time::Duration::from_secs(5), async {
7143 loop {
7144 let envelope = rx.recv().await.unwrap();
7145 if envelope.turn_id.as_deref() != Some(&turn_id) {
7146 continue;
7147 }
7148 match envelope.event {
7149 RoderEvent::TurnPartialResult(event) => {
7150 saw_partial = event.summary.contains("partial turn state");
7151 }
7152 RoderEvent::TurnDeadlineExceeded(event) => {
7153 saw_deadline = event.partial_result.contains("transcript items");
7154 }
7155 RoderEvent::TurnFailed(event) => {
7156 failed_kind = event.error_kind;
7157 break;
7158 }
7159 _ => {}
7160 }
7161 }
7162 })
7163 .await
7164 .unwrap();
7165
7166 for _ in 0..20 {
7167 if !runtime.active_turns.read().await.contains_key(&turn_id) {
7168 break;
7169 }
7170 tokio::time::sleep(std::time::Duration::from_millis(10)).await;
7171 }
7172 assert!(saw_partial);
7173 assert!(saw_deadline);
7174 assert_eq!(failed_kind.as_deref(), Some("deadline_timeout"));
7175 assert!(!runtime.active_turns.read().await.contains_key(&turn_id));
7176 }
7177
7178 #[tokio::test]
7179 async fn deadline_skips_subagent_task_when_remaining_budget_is_too_low() {
7180 let calls = Arc::new(StdMutex::new(0));
7181 let mut builder = ExtensionRegistryBuilder::new();
7182 builder.inference_engine(Arc::new(CapturingEngine {
7183 request: Arc::new(StdMutex::new(None)),
7184 }));
7185 let task_tool = Arc::new(CountingTaskTool {
7186 calls: calls.clone(),
7187 });
7188 builder.tool_contributor(Arc::new(TestToolContributor { tool: task_tool }));
7189 let runtime = Arc::new(
7190 Runtime::new(
7191 builder.build().unwrap(),
7192 RuntimeConfig {
7193 policy_mode: PolicyMode::Bypass,
7194 ..RuntimeConfig::default()
7195 },
7196 )
7197 .unwrap(),
7198 );
7199
7200 let result = runtime
7201 .route_tool_call(
7202 &"thread-deadline-task".to_string(),
7203 &"turn-deadline-task".to_string(),
7204 ToolCallCompleted {
7205 id: "task-1".to_string(),
7206 name: "task".to_string(),
7207 arguments: serde_json::json!({
7208 "description": "inspect",
7209 "prompt": "read"
7210 })
7211 .to_string(),
7212 },
7213 None,
7214 Some(OffsetDateTime::now_utc() + Duration::seconds(1)),
7215 )
7216 .await
7217 .unwrap();
7218
7219 assert!(result.is_error);
7220 assert!(result.result.contains("deadline policy skipped"));
7221 assert_eq!(*calls.lock().unwrap(), 0);
7222 }
7223
7224 struct TestToolContributor {
7225 tool: Arc<dyn ToolExecutor>,
7226 }
7227
7228 impl ToolContributor for TestToolContributor {
7229 fn id(&self) -> String {
7230 "test-tool".to_string()
7231 }
7232
7233 fn contribute(&self, registry: &mut ToolRegistry) -> anyhow::Result<()> {
7234 registry.register(self.tool.clone())
7235 }
7236 }
7237}