systemprompt_cli/runner/routing/
mod.rs1use std::io::{self, Write};
11
12use anyhow::{Context, Result};
13use systemprompt_client::{OutputSink, RemoteCliExecutor, RemoteCliRequest};
14use systemprompt_cloud::{SessionKey, SessionStore, StoredTenant, TenantStore};
15use systemprompt_config::ProfileBootstrap;
16use systemprompt_identifiers::{ContextId, SessionToken};
17use systemprompt_logging::CliService;
18
19use crate::paths::ResolvedPaths;
20
21#[doc(hidden)]
22#[derive(Debug)]
23pub enum ExecutionTarget {
24 Local,
25 Remote {
26 hostname: String,
27 token: SessionToken,
28 context: ContextId,
29 },
30}
31
32pub fn determine_execution_target() -> Result<ExecutionTarget> {
33 let Ok(profile) = ProfileBootstrap::get() else {
34 tracing::debug!("No profile loaded, routing to local execution");
35 return Ok(ExecutionTarget::Local);
36 };
37
38 if profile.target.is_local() {
39 tracing::debug!(
40 profile_name = %profile.name,
41 "Profile target is local, routing to local execution"
42 );
43 return Ok(ExecutionTarget::Local);
44 }
45
46 let Some(tenant_id) = profile.cloud.as_ref().and_then(|c| c.tenant_id.as_ref()) else {
47 tracing::debug!(
48 profile_name = %profile.name,
49 "Profile has no tenant_id, routing to local execution"
50 );
51 return Ok(ExecutionTarget::Local);
52 };
53
54 tracing::debug!(
55 profile_name = %profile.name,
56 tenant_id = %tenant_id,
57 "Profile has tenant_id, resolving remote execution target"
58 );
59
60 let tenant = resolve_tenant(profile, tenant_id)?;
61 let hostname = tenant
62 .hostname
63 .as_ref()
64 .context("Tenant has no hostname configured")?
65 .clone();
66
67 let session_key = SessionKey::Tenant(tenant_id.clone());
68 let session = load_session_for_key(profile, &session_key, &profile.security.issuer)?;
69
70 tracing::info!(
71 hostname = %hostname,
72 tenant_id = %tenant_id,
73 "Routing to remote execution"
74 );
75
76 Ok(ExecutionTarget::Remote {
77 hostname,
78 token: session.session_token,
79 context: session.context_id,
80 })
81}
82
83pub fn resolve_tenant(
84 profile: &systemprompt_manifest::Profile,
85 tenant: &systemprompt_identifiers::TenantId,
86) -> Result<StoredTenant> {
87 let tenants_path = ResolvedPaths::from_profile(profile).tenants_path();
88
89 let store = TenantStore::load_from_path(&tenants_path).with_context(|| {
90 format!(
91 "Failed to load tenants from {}. Run 'systemprompt cloud tenant list' to sync.",
92 tenants_path.display()
93 )
94 })?;
95
96 store.find_tenant(tenant).cloned().with_context(|| {
97 format!(
98 "Tenant '{tenant}' not found in local tenant store. Run 'systemprompt cloud tenant \
99 list' to sync."
100 )
101 })
102}
103
104pub fn load_session_for_key(
105 profile: &systemprompt_manifest::Profile,
106 session_key: &SessionKey,
107 issuer: &str,
108) -> Result<systemprompt_cloud::CliSession> {
109 let sessions_dir = ResolvedPaths::from_profile(profile).sessions_dir();
110
111 let store = SessionStore::load_or_create(&sessions_dir)?;
112
113 store
114 .get_valid_session(session_key, issuer)
115 .cloned()
116 .context("No active session. Run 'systemprompt admin session login'.")
117}
118
119struct StdioSink {
120 stdout: io::Stdout,
121 stderr: io::Stderr,
122}
123
124impl OutputSink for StdioSink {
125 fn stdout_chunk(&mut self, data: &str) -> io::Result<()> {
126 write!(self.stdout, "{}", data)?;
127 self.stdout.flush()
128 }
129
130 fn stderr_chunk(&mut self, data: &str) -> io::Result<()> {
131 write!(self.stderr, "{}", data)?;
132 self.stderr.flush()
133 }
134
135 fn error_message(&mut self, message: &str) {
136 CliService::error(message);
137 }
138}
139
140pub async fn execute_remote(
141 hostname: &str,
142 token: &SessionToken,
143 context: &ContextId,
144 args: &[String],
145 timeout_secs: u64,
146) -> Result<i32> {
147 let executor = RemoteCliExecutor::new(&format!("https://{hostname}"), timeout_secs)
148 .context("Failed to create HTTP client")?;
149 let mut sink = StdioSink {
150 stdout: io::stdout(),
151 stderr: io::stderr(),
152 };
153 let request = RemoteCliRequest {
154 token,
155 context: Some(context),
156 args,
157 };
158 Ok(executor.execute(request, &mut sink).await?)
159}