use crate::error::{EngineError, Result};
use crate::types::{ProvisionedPreview, RemoteWorkspaceConfig};
use crate::workspace_contract::PreviewSpec;
use crate::workspace_provider::{
PreviewPlaceholder, ProgressSink, ProvisionSpec, ReadinessOutcome, TeardownMode,
WorkspaceHandle, WorkspaceProvider, WorkspaceProviderKind,
};
use std::sync::Arc;
use std::time::Duration;
pub const ADAPTER_VERSION: &str = "coder-v1";
const WORKSPACE_NAME_PREFIX: &str = "kranz-remote-";
const REQUEST_TIMEOUT: Duration = Duration::from_secs(30);
const READY_TIMEOUT: Duration = Duration::from_secs(300);
const READY_POLL_INTERVAL: Duration = Duration::from_secs(2);
pub(crate) const PROVIDER_REASON_PREFIX: &str = "workspace provider:";
pub(crate) fn provider_block_reason(detail: &str) -> String {
crate::scrub::scrub(&format!(
"{PROVIDER_REASON_PREFIX} {detail} (owner: provider — fix the substrate, then unblock \
and re-run)"
))
}
#[derive(Debug, Clone, PartialEq)]
pub struct SubstrateWorkspaceSpec {
pub template: String,
pub name: String,
pub env_names: Vec<String>,
pub idle_after_hours: Option<f64>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SubstrateUrl {
pub name: String,
pub url: String,
pub auth: Option<bool>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SubstrateWorkspace {
pub id: String,
pub urls: Vec<SubstrateUrl>,
pub takeover: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum SubstrateStatus {
Ready,
Pending,
Failed { reason: String },
}
#[async_trait::async_trait]
pub trait SubstrateClient: Send + Sync {
async fn create_workspace(&self, spec: &SubstrateWorkspaceSpec) -> Result<SubstrateWorkspace>;
async fn workspace_status(&self, id: &str) -> Result<SubstrateStatus>;
async fn delete_workspace(&self, id: &str) -> Result<()>;
async fn stop_workspace(&self, id: &str) -> Result<()>;
}
#[derive(Debug, Clone, PartialEq)]
pub struct RemoteConfig {
pub base_url: String,
pub template: String,
pub token_env: String,
pub idle_after_hours: Option<f64>,
}
impl RemoteConfig {
pub fn require(config: Option<&RemoteWorkspaceConfig>) -> Result<RemoteConfig> {
fn missing(key: &str) -> EngineError {
EngineError::Config(format!(
"workspace.provider \"remote\" needs {key} configured (owner: operator — set the \
mission config key); refusing rather than silently falling back to local"
))
}
let Some(config) = config else {
return Err(missing("workspace.remote.baseUrl"));
};
Ok(RemoteConfig {
base_url: config
.base_url
.clone()
.ok_or_else(|| missing("workspace.remote.baseUrl"))?,
template: config
.template
.clone()
.ok_or_else(|| missing("workspace.remote.template"))?,
token_env: config
.token_env
.clone()
.ok_or_else(|| missing("workspace.remote.tokenEnv"))?,
idle_after_hours: config.idle_after_hours,
})
}
}
type ClientHook = Arc<dyn Fn(&RemoteConfig) -> Result<Arc<dyn SubstrateClient>> + Send + Sync>;
fn coder_client_from_env(config: &RemoteConfig) -> Result<Arc<dyn SubstrateClient>> {
let token = std::env::var(&config.token_env)
.ok()
.filter(|value| !value.is_empty())
.ok_or_else(|| {
EngineError::Config(format!(
"workspace.provider \"remote\": the substrate token env var {} (from \
workspace.remote.tokenEnv) is not set (owner: operator — export it before \
`kranz work`); refusing rather than silently falling back to local",
config.token_env
))
})?;
Ok(Arc::new(CoderHttpClient::new(&config.base_url, &token)?))
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum PollOutcome {
Ready,
Failed { reason: String },
TimedOut { waited_secs: u64 },
}
#[derive(Debug, Clone)]
pub struct RemoteWorkspace {
pub name: String,
pub id: String,
pub takeover: Option<String>,
pub previews: Vec<ProvisionedPreview>,
pub injected_env_names: Vec<String>,
pub poll: PollOutcome,
}
pub struct RemoteWorkspaceProvider {
config: RemoteConfig,
client_hook: ClientHook,
ready_timeout: Duration,
poll_interval: Duration,
}
impl RemoteWorkspaceProvider {
pub(crate) fn from_config(config: Option<&RemoteWorkspaceConfig>) -> Result<Self> {
Ok(Self {
config: RemoteConfig::require(config)?,
client_hook: Arc::new(coder_client_from_env),
ready_timeout: READY_TIMEOUT,
poll_interval: READY_POLL_INTERVAL,
})
}
#[cfg(test)]
pub(crate) fn with_client(config: RemoteConfig, client: Arc<dyn SubstrateClient>) -> Self {
Self {
config,
client_hook: Arc::new(move |_config| Ok(Arc::clone(&client))),
ready_timeout: Duration::from_millis(120),
poll_interval: Duration::from_millis(5),
}
}
async fn poll_until_ready(&self, client: &dyn SubstrateClient, id: &str) -> PollOutcome {
let deadline = std::time::Instant::now() + self.ready_timeout;
loop {
match client.workspace_status(id).await {
Ok(SubstrateStatus::Ready) => return PollOutcome::Ready,
Ok(SubstrateStatus::Failed { reason }) => return PollOutcome::Failed { reason },
Ok(SubstrateStatus::Pending) => {}
Err(e) => {
return PollOutcome::Failed {
reason: format!("status poll failed: {e}"),
};
}
}
if std::time::Instant::now() >= deadline {
return PollOutcome::TimedOut {
waited_secs: self.ready_timeout.as_secs(),
};
}
tokio::time::sleep(self.poll_interval).await;
}
}
}
fn workspace_name(mission_id: &str) -> String {
let sanitized: String = mission_id
.chars()
.map(|c| {
if c.is_ascii_alphanumeric() {
c.to_ascii_lowercase()
} else if c == '-' || c == '_' {
c
} else {
'-'
}
})
.collect();
if sanitized.is_empty() {
format!("{WORKSPACE_NAME_PREFIX}mission")
} else {
format!("{WORKSPACE_NAME_PREFIX}{sanitized}")
}
}
fn preview_placeholders(
previews: &[PreviewSpec],
urls: &[SubstrateUrl],
) -> Vec<PreviewPlaceholder> {
previews
.iter()
.map(|preview| {
let url_template = urls
.iter()
.find(|url| url.name == preview.name)
.map(|url| url.url.clone())
.unwrap_or_else(|| preview.url_template.clone());
PreviewPlaceholder {
name: preview.name.clone(),
url_template,
}
})
.collect()
}
fn map_previews(previews: &[PreviewSpec], urls: &[SubstrateUrl]) -> Vec<ProvisionedPreview> {
previews
.iter()
.filter_map(|preview| {
urls.iter()
.find(|url| url.name == preview.name)
.map(|url| ProvisionedPreview {
name: preview.name.clone(),
url: url.url.clone(),
auth: url.auth,
})
})
.collect()
}
fn provision_detail(
name: &str,
id: &str,
env_names: &[String],
idle_after_hours: Option<f64>,
) -> String {
let injected = if env_names.is_empty() {
"none".to_string()
} else {
env_names.join(",")
};
let idle = match idle_after_hours {
Some(hours) => format!("; idle policy: hibernate after {hours}h (substrate-owned)"),
None => String::new(),
};
if id.is_empty() {
format!("substrate workspace {name}: create failed (injected env names: {injected}){idle}")
} else {
format!("substrate workspace {name} (id {id}); injected env names: {injected}{idle}")
}
}
#[async_trait::async_trait]
impl WorkspaceProvider for RemoteWorkspaceProvider {
fn kind(&self) -> WorkspaceProviderKind {
WorkspaceProviderKind::Remote
}
async fn provision(&self, spec: &ProvisionSpec) -> Result<WorkspaceHandle> {
if !spec.repo_root.is_dir() {
return Err(EngineError::InvalidState(format!(
"remote provision: execution cwd {} does not exist",
spec.repo_root.display()
)));
}
let env = crate::runner::contract_env(spec.base_sha.as_deref());
let Some(contract) = &spec.contract else {
return Ok(WorkspaceHandle {
cwd: spec.repo_root.clone(),
env,
previews: Vec::new(),
contract: None,
detail: None,
container: None,
remote: None,
gate_env: spec.gate_env.clone(),
});
};
let client = (self.client_hook)(&self.config)?;
let name = workspace_name(&spec.mission_id);
let env_names = contract.secrets.clone();
let (id, urls, takeover, poll) = match client
.create_workspace(&SubstrateWorkspaceSpec {
template: self.config.template.clone(),
name: name.clone(),
env_names: env_names.clone(),
idle_after_hours: self.config.idle_after_hours,
})
.await
{
Ok(created) => {
let poll = self.poll_until_ready(&*client, &created.id).await;
(created.id, created.urls, created.takeover, poll)
}
Err(e) => (
String::new(),
Vec::new(),
None,
PollOutcome::Failed {
reason: format!("create_workspace failed: {e}"),
},
),
};
Ok(WorkspaceHandle {
cwd: spec.repo_root.clone(),
env,
previews: preview_placeholders(&contract.previews, &urls),
contract: Some(contract.clone()),
detail: Some(provision_detail(
&name,
&id,
&env_names,
self.config.idle_after_hours,
)),
container: None,
remote: Some(RemoteWorkspace {
name,
id,
takeover,
previews: map_previews(&contract.previews, &urls),
injected_env_names: env_names,
poll,
}),
gate_env: spec.gate_env.clone(),
})
}
async fn readiness(
&self,
handle: &WorkspaceHandle,
progress: &mut ProgressSink<'_>,
) -> Result<ReadinessOutcome> {
let Some(remote) = &handle.remote else {
return Ok(ReadinessOutcome::Ready);
};
match &remote.poll {
PollOutcome::Ready => {
progress(
"workspace remote: substrate reports ready (substrate-reported readiness \
only — contract bootstrap/readiness commands do not execute on the remote \
substrate in v1)",
None,
)?;
Ok(ReadinessOutcome::Ready)
}
PollOutcome::Failed { reason } => Ok(ReadinessOutcome::ProviderFailed {
detail: format!(
"substrate workspace {} failed to provision: {reason}",
remote.name
),
}),
PollOutcome::TimedOut { waited_secs } => Ok(ReadinessOutcome::ProviderFailed {
detail: format!(
"substrate workspace {} did not become ready within {waited_secs}s",
remote.name
),
}),
}
}
async fn teardown(&self, handle: WorkspaceHandle, mode: TeardownMode) -> Result<()> {
let Some(remote) = &handle.remote else {
return Ok(()); };
if remote.id.is_empty() {
return Ok(()); }
match mode {
TeardownMode::Keep => Ok(()),
TeardownMode::Hibernate => {
let client = (self.client_hook)(&self.config)?;
client.stop_workspace(&remote.id).await.map_err(|e| {
EngineError::InvalidState(format!(
"remote teardown (hibernate): stop_workspace failed (owner: provider): {}",
crate::scrub::scrub(&e.to_string())
))
})
}
TeardownMode::Destroy => {
let client = (self.client_hook)(&self.config)?;
client.delete_workspace(&remote.id).await.map_err(|e| {
EngineError::InvalidState(format!(
"remote teardown (destroy): delete_workspace failed (owner: provider): {}",
crate::scrub::scrub(&e.to_string())
))
})
}
}
}
}
pub struct CoderHttpClient {
base_url: String,
token: String,
client: reqwest::Client,
}
impl CoderHttpClient {
pub fn new(base_url: &str, token: &str) -> Result<Self> {
let parsed = reqwest::Url::parse(base_url).map_err(|e| {
EngineError::Config(format!(
"workspace.remote.baseUrl {base_url:?} is not a valid URL: {e} (owner: operator)"
))
})?;
if !matches!(parsed.scheme(), "http" | "https") || parsed.host_str().is_none() {
return Err(EngineError::Config(format!(
"workspace.remote.baseUrl {base_url:?} needs an http(s) URL with a host (owner: \
operator)"
)));
}
let client = reqwest::Client::builder()
.timeout(REQUEST_TIMEOUT)
.build()
.map_err(|e| {
EngineError::Config(format!("could not build the substrate HTTP client: {e}"))
})?;
Ok(Self {
base_url: base_url.trim_end_matches('/').to_string(),
token: token.to_string(),
client,
})
}
async fn send(
&self,
request: reqwest::RequestBuilder,
op: &'static str,
) -> Result<reqwest::Response> {
let response = request
.header("Coder-Session-Token", &self.token)
.send()
.await
.map_err(|e| {
EngineError::InvalidState(format!(
"coder substrate {op}: request failed: {}",
crate::scrub::scrub(&e.to_string())
))
})?;
let status = response.status();
if !status.is_success() {
let body = response.text().await.unwrap_or_default();
let tail = crate::command_exec::last_chars_local(&body, 500);
return Err(EngineError::InvalidState(format!(
"coder substrate {op}: HTTP {status}: {}",
crate::scrub::scrub(tail.trim())
)));
}
Ok(response)
}
async fn json(response: reqwest::Response, op: &'static str) -> Result<serde_json::Value> {
response.json().await.map_err(|e| {
EngineError::InvalidState(format!("coder substrate {op}: response was not JSON: {e}"))
})
}
async fn transition(&self, id: &str, transition: &str, op: &'static str) -> Result<()> {
let url = format!("{}/api/v2/workspaces/{id}/builds", self.base_url);
self.send(
self.client
.post(&url)
.json(&serde_json::json!({ "transition": transition })),
op,
)
.await?;
Ok(())
}
}
fn parse_substrate_url(value: &serde_json::Value) -> Option<SubstrateUrl> {
Some(SubstrateUrl {
name: value.get("name")?.as_str()?.to_string(),
url: value.get("url")?.as_str()?.to_string(),
auth: value.get("auth").and_then(serde_json::Value::as_bool),
})
}
fn parse_status(value: &serde_json::Value) -> SubstrateStatus {
let status = value
.get("status")
.and_then(|v| v.as_str())
.or_else(|| {
value
.get("latest_build")
.and_then(|build| build.get("status"))
.and_then(|v| v.as_str())
})
.unwrap_or("");
match status {
"running" => SubstrateStatus::Ready,
"failed" => SubstrateStatus::Failed {
reason: "substrate reported workspace status \"failed\"".to_string(),
},
_ => SubstrateStatus::Pending,
}
}
#[async_trait::async_trait]
impl SubstrateClient for CoderHttpClient {
async fn create_workspace(&self, spec: &SubstrateWorkspaceSpec) -> Result<SubstrateWorkspace> {
let url = format!("{}/api/v2/users/me/workspaces", self.base_url);
let mut body = serde_json::json!({
"template_id": spec.template,
"name": spec.name,
"env_names": spec.env_names,
});
if let Some(hours) = spec.idle_after_hours {
body["idle_after_hours"] = serde_json::json!(hours);
}
let response = self
.send(self.client.post(&url).json(&body), "create_workspace")
.await?;
let value = Self::json(response, "create_workspace").await?;
let id = value
.get("id")
.and_then(|v| v.as_str())
.filter(|id| !id.is_empty())
.ok_or_else(|| {
EngineError::InvalidState(
"coder substrate create_workspace: response carried no workspace id"
.to_string(),
)
})?;
let urls = value
.get("urls")
.and_then(|v| v.as_array())
.map(|entries| entries.iter().filter_map(parse_substrate_url).collect())
.unwrap_or_default();
let takeover = value
.get("takeover")
.and_then(|v| v.as_str())
.map(str::to_string);
Ok(SubstrateWorkspace {
id: id.to_string(),
urls,
takeover,
})
}
async fn workspace_status(&self, id: &str) -> Result<SubstrateStatus> {
let url = format!("{}/api/v2/workspaces/{id}", self.base_url);
let response = self.send(self.client.get(&url), "workspace_status").await?;
let value = Self::json(response, "workspace_status").await?;
Ok(parse_status(&value))
}
async fn delete_workspace(&self, id: &str) -> Result<()> {
self.transition(id, "delete", "delete_workspace").await
}
async fn stop_workspace(&self, id: &str) -> Result<()> {
self.transition(id, "stop", "stop_workspace").await
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::workspace_contract::{parse_workspace_contract, WorkspaceContract};
use std::collections::VecDeque;
use std::sync::Mutex;
fn config() -> RemoteConfig {
RemoteConfig {
base_url: "https://coder.internal.example.com".to_string(),
template: "tmpl-baked-ami".to_string(),
token_env: "CODER_SESSION_TOKEN".to_string(),
idle_after_hours: None,
}
}
fn spec(root: &std::path::Path, contract: Option<WorkspaceContract>) -> ProvisionSpec {
let runtime_dir = root.join(".kranz").join("missions").join("m-test");
ProvisionSpec {
mission_id: "m-test".to_string(),
repo_root: root.to_path_buf(),
gate_env: crate::workspace_provider::GateEnvPolicy::for_mission(&runtime_dir, &[]),
runtime_dir,
base_sha: Some("deadbeefcafe".to_string()),
contract,
}
}
fn contract(json: &[u8]) -> WorkspaceContract {
parse_workspace_contract(json).expect("valid contract")
}
fn remote_contract() -> WorkspaceContract {
contract(
br#"{
"schemaVersion": 1,
"bootstrap": ["echo never-run-remotely"],
"readiness": ["true"],
"previews": [
{ "name": "app", "urlTemplate": "http://localhost:{port}/" },
{ "name": "db", "urlTemplate": "postgres://localhost:{port}/" }
],
"secrets": ["DATABASE_URL", "STRIPE_API_KEY"]
}"#,
)
}
#[derive(Default)]
struct Progress(Vec<String>);
impl Progress {
fn sink(&mut self) -> impl FnMut(&str, Option<String>) -> Result<()> + Send + use<'_> {
|summary, _detail| {
self.0.push(summary.to_string());
Ok(())
}
}
}
struct FakeSubstrateClient {
created: Mutex<Vec<SubstrateWorkspaceSpec>>,
ops: Mutex<Vec<String>>,
statuses: Mutex<VecDeque<SubstrateStatus>>,
fail_create: Mutex<Option<String>>,
workspace: SubstrateWorkspace,
}
impl FakeSubstrateClient {
fn new(statuses: Vec<SubstrateStatus>) -> Arc<Self> {
Arc::new(Self {
created: Mutex::new(Vec::new()),
ops: Mutex::new(Vec::new()),
statuses: Mutex::new(statuses.into()),
fail_create: Mutex::new(None),
workspace: SubstrateWorkspace {
id: "ws-abc123".to_string(),
urls: vec![SubstrateUrl {
name: "app".to_string(),
url: "https://app--m-test.coder.internal.example.com".to_string(),
auth: Some(true),
}],
takeover: Some(
"ssh://coder.internal.example.com/kranz-remote-m-test".to_string(),
),
},
})
}
fn failing_create(message: &str) -> Arc<Self> {
let client = Self::new(vec![]);
client
.fail_create
.lock()
.unwrap()
.replace(message.to_string());
client
}
fn ops(&self) -> Vec<String> {
self.ops.lock().unwrap().clone()
}
}
#[async_trait::async_trait]
impl SubstrateClient for FakeSubstrateClient {
async fn create_workspace(
&self,
spec: &SubstrateWorkspaceSpec,
) -> Result<SubstrateWorkspace> {
self.created.lock().unwrap().push(spec.clone());
if let Some(message) = self.fail_create.lock().unwrap().as_ref() {
return Err(EngineError::InvalidState(message.clone()));
}
Ok(self.workspace.clone())
}
async fn workspace_status(&self, id: &str) -> Result<SubstrateStatus> {
self.ops.lock().unwrap().push(format!("status:{id}"));
let mut statuses = self.statuses.lock().unwrap();
if statuses.len() > 1 {
Ok(statuses.pop_front().expect("len > 1"))
} else {
Ok(statuses.front().cloned().unwrap_or(SubstrateStatus::Ready))
}
}
async fn delete_workspace(&self, id: &str) -> Result<()> {
self.ops.lock().unwrap().push(format!("delete:{id}"));
Ok(())
}
async fn stop_workspace(&self, id: &str) -> Result<()> {
self.ops.lock().unwrap().push(format!("stop:{id}"));
Ok(())
}
}
fn provider_with(client: Arc<FakeSubstrateClient>) -> RemoteWorkspaceProvider {
RemoteWorkspaceProvider::with_client(config(), client)
}
async fn provisioned(
provider: &RemoteWorkspaceProvider,
root: &std::path::Path,
) -> WorkspaceHandle {
provider
.provision(&spec(root, Some(remote_contract())))
.await
.expect("provision")
}
#[tokio::test]
async fn remote_workspace_provision_creates_with_template_name_and_secret_names() {
let dir = tempfile::tempdir().expect("tempdir");
let client = FakeSubstrateClient::new(vec![SubstrateStatus::Ready]);
let provider = provider_with(Arc::clone(&client));
let handle = provisioned(&provider, dir.path()).await;
let created = client.created.lock().unwrap().clone();
assert_eq!(created.len(), 1, "exactly one create_workspace call");
assert_eq!(created[0].template, "tmpl-baked-ami");
assert_eq!(created[0].name, "kranz-remote-m-test");
assert_eq!(
created[0].env_names,
vec!["DATABASE_URL".to_string(), "STRIPE_API_KEY".to_string()],
"the contract's secret NAMES — and only names — cross to the substrate"
);
assert_eq!(
handle.cwd,
dir.path(),
"sessions stay in the local worktree in v1"
);
assert_eq!(
handle.env.get("KRANZ_BASE_SHA").map(String::as_str),
Some("deadbeefcafe")
);
let remote = handle.remote.as_ref().expect("remote state");
assert_eq!(remote.id, "ws-abc123");
assert_eq!(remote.injected_env_names, created[0].env_names);
assert!(
handle.detail.as_deref().unwrap().contains("ws-abc123")
&& handle
.detail
.as_deref()
.unwrap()
.contains("DATABASE_URL,STRIPE_API_KEY"),
"the provisioned detail records the workspace and injected names: {:?}",
handle.detail
);
}
#[tokio::test]
async fn remote_workspace_ready_poll_lands_previews_and_takeover() {
let dir = tempfile::tempdir().expect("tempdir");
let client =
FakeSubstrateClient::new(vec![SubstrateStatus::Pending, SubstrateStatus::Ready]);
let provider = provider_with(Arc::clone(&client));
let handle = provisioned(&provider, dir.path()).await;
assert!(
client
.ops()
.iter()
.filter(|op| op.starts_with("status:"))
.count()
>= 2,
"the poll looped past the pending status: {:?}",
client.ops()
);
let remote = handle.remote.as_ref().expect("remote state");
assert_eq!(remote.poll, PollOutcome::Ready);
assert_eq!(
remote.takeover.as_deref(),
Some("ssh://coder.internal.example.com/kranz-remote-m-test")
);
assert_eq!(
remote.previews,
vec![ProvisionedPreview {
name: "app".to_string(),
url: "https://app--m-test.coder.internal.example.com".to_string(),
auth: Some(true),
}],
"only the substrate-reported URL is recorded — db stays unfabricated"
);
assert_eq!(
handle.previews,
vec![
PreviewPlaceholder {
name: "app".to_string(),
url_template: "https://app--m-test.coder.internal.example.com".to_string(),
},
PreviewPlaceholder {
name: "db".to_string(),
url_template: "postgres://localhost:{port}/".to_string(),
},
]
);
let mut progress = Progress::default();
let outcome = provider
.readiness(&handle, &mut progress.sink())
.await
.expect("readiness");
assert!(matches!(outcome, ReadinessOutcome::Ready), "{outcome:?}");
assert_eq!(progress.0.len(), 1);
assert!(
progress.0[0].starts_with("workspace remote: substrate reports ready")
&& progress.0[0].contains("substrate-reported readiness only"),
"honest readiness wording, no gate-outcome prefix: {}",
progress.0[0]
);
}
#[tokio::test]
async fn remote_workspace_failed_status_is_a_provider_owned_block() {
let dir = tempfile::tempdir().expect("tempdir");
let client = FakeSubstrateClient::new(vec![SubstrateStatus::Failed {
reason: "template build exited 1".to_string(),
}]);
let provider = provider_with(Arc::clone(&client));
let handle = provisioned(&provider, dir.path()).await;
let mut progress = Progress::default();
let outcome = provider
.readiness(&handle, &mut progress.sink())
.await
.expect("readiness");
let ReadinessOutcome::ProviderFailed { detail } = outcome else {
panic!("a failed substrate status must be ProviderFailed, got {outcome:?}");
};
assert!(detail.contains("kranz-remote-m-test"), "{detail}");
assert!(detail.contains("template build exited 1"), "{detail}");
let reason = provider_block_reason(&detail);
assert!(reason.starts_with("workspace provider:"), "{reason}");
assert!(reason.contains("owner: provider"), "{reason}");
assert!(!reason.contains("repo-setup"), "{reason}");
assert!(
!reason.starts_with(crate::workspace_gate::GATE_REASON_PREFIX),
"not gate-owned ⇒ the pass path never auto-lifts it: {reason}"
);
assert!(
progress.0.is_empty(),
"a failed workspace reports no ready line"
);
}
#[tokio::test]
async fn remote_workspace_poll_timeout_is_provider_owned_and_bounded() {
let dir = tempfile::tempdir().expect("tempdir");
let client = FakeSubstrateClient::new(vec![SubstrateStatus::Pending]);
let provider = provider_with(Arc::clone(&client));
let handle =
tokio::time::timeout(Duration::from_secs(10), provisioned(&provider, dir.path()))
.await
.expect("the provision poll must terminate inside its bound");
assert!(matches!(
handle.remote.as_ref().expect("remote").poll,
PollOutcome::TimedOut { .. }
));
let mut progress = Progress::default();
let outcome = provider
.readiness(&handle, &mut progress.sink())
.await
.expect("readiness");
let ReadinessOutcome::ProviderFailed { detail } = outcome else {
panic!("a poll timeout must be ProviderFailed, got {outcome:?}");
};
assert!(detail.contains("did not become ready within"), "{detail}");
}
#[tokio::test]
async fn remote_workspace_create_failure_is_recorded_and_teardown_noops() {
let dir = tempfile::tempdir().expect("tempdir");
let client = FakeSubstrateClient::failing_create("HTTP 401: bad token");
let provider = provider_with(Arc::clone(&client));
let handle = provisioned(&provider, dir.path()).await;
let remote = handle.remote.as_ref().expect("remote state");
assert!(remote.id.is_empty(), "no id when create failed");
assert!(
matches!(&remote.poll, PollOutcome::Failed { reason } if reason.contains("HTTP 401")),
"{:?}",
remote.poll
);
let mut progress = Progress::default();
let outcome = provider
.readiness(&handle, &mut progress.sink())
.await
.expect("readiness");
assert!(
matches!(outcome, ReadinessOutcome::ProviderFailed { .. }),
"{outcome:?}"
);
for mode in [
TeardownMode::Keep,
TeardownMode::Hibernate,
TeardownMode::Destroy,
] {
provider
.teardown(handle.clone(), mode)
.await
.expect("teardown with no workspace no-ops");
}
assert!(
!client
.ops()
.iter()
.any(|op| op.starts_with("stop:") || op.starts_with("delete:")),
"nothing exists ⇒ no substrate teardown calls: {:?}",
client.ops()
);
}
#[tokio::test]
async fn remote_workspace_teardown_modes_map_to_delete_stop_keep() {
let dir = tempfile::tempdir().expect("tempdir");
let client = FakeSubstrateClient::new(vec![SubstrateStatus::Ready]);
let provider = provider_with(Arc::clone(&client));
let handle = provisioned(&provider, dir.path()).await;
provider
.teardown(handle, TeardownMode::Keep)
.await
.expect("keep");
assert!(
!client
.ops()
.iter()
.any(|op| op.starts_with("stop:") || op.starts_with("delete:")),
"Keep leaves the workspace running — no teardown call: {:?}",
client.ops()
);
let client = FakeSubstrateClient::new(vec![SubstrateStatus::Ready]);
let provider = provider_with(Arc::clone(&client));
let handle = provisioned(&provider, dir.path()).await;
provider
.teardown(handle, TeardownMode::Hibernate)
.await
.expect("hibernate");
assert!(
client.ops().contains(&"stop:ws-abc123".to_string())
&& !client.ops().iter().any(|op| op.starts_with("delete:")),
"Hibernate is stop_workspace: {:?}",
client.ops()
);
let client = FakeSubstrateClient::new(vec![SubstrateStatus::Ready]);
let provider = provider_with(Arc::clone(&client));
let handle = provisioned(&provider, dir.path()).await;
provider
.teardown(handle, TeardownMode::Destroy)
.await
.expect("destroy");
assert!(
client.ops().contains(&"delete:ws-abc123".to_string())
&& !client.ops().iter().any(|op| op.starts_with("stop:")),
"Destroy is delete_workspace: {:?}",
client.ops()
);
}
#[tokio::test]
async fn idle_after_hours_passes_through_to_create_and_the_provisioned_detail() {
let validated = RemoteConfig::require(Some(&RemoteWorkspaceConfig {
base_url: Some("https://coder.internal.example.com".to_string()),
template: Some("tmpl-baked-ami".to_string()),
token_env: Some("CODER_SESSION_TOKEN".to_string()),
idle_after_hours: Some(24.0),
}))
.expect("complete remote config validates with an idle policy");
assert_eq!(validated.idle_after_hours, Some(24.0));
let without = RemoteConfig::require(Some(&RemoteWorkspaceConfig {
base_url: Some("https://coder.internal.example.com".to_string()),
template: Some("tmpl-baked-ami".to_string()),
token_env: Some("CODER_SESSION_TOKEN".to_string()),
idle_after_hours: None,
}))
.expect("complete remote config validates without an idle policy");
assert_eq!(without.idle_after_hours, None);
let dir = tempfile::tempdir().expect("tempdir");
let client = FakeSubstrateClient::new(vec![SubstrateStatus::Ready]);
let provider = RemoteWorkspaceProvider::with_client(
validated,
Arc::clone(&client) as Arc<dyn SubstrateClient>,
);
let handle = provisioned(&provider, dir.path()).await;
let created = client.created.lock().unwrap().clone();
assert_eq!(created.len(), 1);
assert_eq!(
created[0].idle_after_hours,
Some(24.0),
"the idle policy VALUE crosses to the substrate verbatim"
);
assert!(
handle
.detail
.as_deref()
.unwrap()
.contains("idle policy: hibernate after 24h (substrate-owned)"),
"the provisioned event records the substrate-owned policy: {:?}",
handle.detail
);
}
#[tokio::test]
async fn remote_workspace_contract_less_provision_never_touches_the_substrate() {
let dir = tempfile::tempdir().expect("tempdir");
let client = FakeSubstrateClient::new(vec![SubstrateStatus::Ready]);
let provider = provider_with(Arc::clone(&client));
let handle = provider
.provision(&spec(dir.path(), None))
.await
.expect("contract-less provision");
assert!(handle.remote.is_none() && handle.contract.is_none());
assert!(
client.created.lock().unwrap().is_empty() && client.ops().is_empty(),
"no contract ⇒ no substrate workspace (D-H)"
);
let mut progress = Progress::default();
let outcome = provider
.readiness(&handle, &mut progress.sink())
.await
.expect("readiness");
assert!(matches!(outcome, ReadinessOutcome::Ready));
assert!(progress.0.is_empty(), "contract-less ⇒ silent");
}
#[tokio::test]
async fn remote_workspace_missing_creds_fail_closed_naming_the_env_var() {
let dir = tempfile::tempdir().expect("tempdir");
let var = format!("KRANZ_TEST_REMOTE_TOKEN_UNSET_{}", std::process::id());
std::env::remove_var(&var); let mut cfg = config();
cfg.token_env = var.clone();
let provider = RemoteWorkspaceProvider {
config: cfg,
client_hook: Arc::new(coder_client_from_env),
ready_timeout: Duration::from_millis(50),
poll_interval: Duration::from_millis(5),
};
let err = provider
.provision(&spec(dir.path(), Some(remote_contract())))
.await
.expect_err("missing creds must fail closed at provision");
let msg = err.to_string();
assert!(msg.contains(&var), "names the env var NAME: {msg}");
assert!(msg.contains("workspace.remote.tokenEnv"), "{msg}");
assert!(msg.contains("owner: operator"), "{msg}");
assert!(
msg.contains("refusing rather than silently falling back"),
"{msg}"
);
}
#[test]
fn remote_workspace_name_is_mission_owned_and_coder_shaped() {
assert_eq!(workspace_name("m-test"), "kranz-remote-m-test");
assert_eq!(workspace_name("M.Test X"), "kranz-remote-m-test-x");
assert_eq!(
workspace_name("..."),
"kranz-remote----",
"every non-alnum sanitizes to '-'"
);
assert_eq!(
workspace_name(""),
"kranz-remote-mission",
"the prefix guarantees a valid leading character"
);
}
#[test]
fn coder_http_client_rejects_a_bad_base_url() {
let err = CoderHttpClient::new("not a url", "tok")
.err()
.expect("invalid URL fails closed");
assert!(
err.to_string().contains("workspace.remote.baseUrl"),
"{err}"
);
let err = CoderHttpClient::new("file:///etc/passwd", "tok")
.err()
.expect("non-http(s) schemes fail closed");
assert!(err.to_string().contains("http(s)"), "{err}");
}
#[test]
fn coder_http_status_mapping_is_conservative() {
assert_eq!(
parse_status(&serde_json::json!({"status": "running"})),
SubstrateStatus::Ready
);
assert_eq!(
parse_status(&serde_json::json!({"latest_build": {"status": "running"}})),
SubstrateStatus::Ready,
"stock Coder nests the state under latest_build"
);
assert!(matches!(
parse_status(&serde_json::json!({"latest_build": {"status": "failed"}})),
SubstrateStatus::Failed { .. }
));
for unknown in [
serde_json::json!({"latest_build": {"status": "starting"}}),
serde_json::json!({"status": "stopping"}),
serde_json::json!({}),
] {
assert_eq!(
parse_status(&unknown),
SubstrateStatus::Pending,
"unknown/absent states keep polling, never guess ready: {unknown}"
);
}
}
fn find_subslice(haystack: &[u8], needle: &[u8]) -> Option<usize> {
haystack
.windows(needle.len())
.position(|window| window == needle)
}
async fn serve_once(
listener: tokio::net::TcpListener,
body: &str,
recorded: Arc<Mutex<Vec<String>>>,
) {
use tokio::io::{AsyncReadExt, AsyncWriteExt};
let (mut socket, _) = listener.accept().await.expect("accept");
let mut buf = Vec::new();
let mut chunk = [0u8; 4096];
loop {
let n = socket.read(&mut chunk).await.expect("read request");
assert!(n > 0, "connection closed before the full request arrived");
buf.extend_from_slice(&chunk[..n]);
if let Some(pos) = find_subslice(&buf, b"\r\n\r\n") {
let headers = String::from_utf8_lossy(&buf[..pos]).to_string();
let content_length = headers
.lines()
.find_map(|line| {
line.to_ascii_lowercase()
.strip_prefix("content-length:")
.and_then(|value| value.trim().parse::<usize>().ok())
})
.unwrap_or(0);
if buf.len() >= pos + 4 + content_length {
break;
}
}
}
recorded
.lock()
.unwrap()
.push(String::from_utf8_lossy(&buf).to_string());
let response = format!(
"HTTP/1.1 200 OK\r\ncontent-type: application/json\r\ncontent-length: {}\r\nconnection: close\r\n\r\n{body}",
body.len()
);
socket
.write_all(response.as_bytes())
.await
.expect("write response");
}
#[tokio::test]
async fn coder_http_client_maps_the_coder_shaped_wire_over_loopback() {
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("bind loopback");
let base_url = format!("http://{}", listener.local_addr().expect("addr"));
let recorded = Arc::new(Mutex::new(Vec::new()));
let create_body = serde_json::json!({
"id": "ws-loop1",
"urls": [{"name": "app", "url": "https://app.example.com", "auth": true}],
"takeover": "https://coder.example.com/@me/ws-loop1",
})
.to_string();
let client = CoderHttpClient::new(&base_url, "test-session-token").expect("client");
let create_spec = SubstrateWorkspaceSpec {
template: "tmpl-1".to_string(),
name: "kranz-remote-m-1".to_string(),
env_names: vec!["DATABASE_URL".to_string()],
idle_after_hours: None,
};
let ((), workspace) = tokio::join!(
serve_once(listener, &create_body, Arc::clone(&recorded)),
tokio::time::timeout(
Duration::from_secs(10),
client.create_workspace(&create_spec)
)
);
let workspace = workspace.expect("bounded").expect("create_workspace");
assert_eq!(workspace.id, "ws-loop1");
assert_eq!(
workspace.urls,
vec![SubstrateUrl {
name: "app".to_string(),
url: "https://app.example.com".to_string(),
auth: Some(true),
}]
);
assert_eq!(
workspace.takeover.as_deref(),
Some("https://coder.example.com/@me/ws-loop1")
);
let requests = recorded.lock().unwrap().clone();
assert_eq!(requests.len(), 1);
let request = &requests[0];
assert!(
request.starts_with("POST /api/v2/users/me/workspaces "),
"the create verb+path: {}",
request.lines().next().unwrap_or("")
);
assert!(
request.contains("coder-session-token: test-session-token"),
"Coder's auth header carries the token (and the token appears NOWHERE else): {request}"
);
let body = request.split("\r\n\r\n").nth(1).expect("a JSON body");
let body: serde_json::Value = serde_json::from_str(body).expect("body is JSON");
assert_eq!(body["template_id"], "tmpl-1");
assert_eq!(body["name"], "kranz-remote-m-1");
assert_eq!(
body["env_names"],
serde_json::json!(["DATABASE_URL"]),
"secret NAMES on the wire — never values"
);
}
#[tokio::test]
async fn coder_http_create_body_carries_idle_after_hours_only_when_configured() {
for (idle_after_hours, expected) in
[(Some(24.0), Some(serde_json::json!(24.0))), (None, None)]
{
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("bind loopback");
let base_url = format!("http://{}", listener.local_addr().expect("addr"));
let recorded = Arc::new(Mutex::new(Vec::new()));
let client = CoderHttpClient::new(&base_url, "tok").expect("client");
let spec = SubstrateWorkspaceSpec {
template: "tmpl-1".to_string(),
name: "kranz-remote-m-1".to_string(),
env_names: vec![],
idle_after_hours,
};
let ((), created) = tokio::join!(
serve_once(listener, r#"{"id":"ws-1"}"#, Arc::clone(&recorded)),
tokio::time::timeout(Duration::from_secs(10), client.create_workspace(&spec))
);
created.expect("bounded").expect("create_workspace");
let request = recorded.lock().unwrap()[0].clone();
let body = request.split("\r\n\r\n").nth(1).expect("a JSON body");
let body: serde_json::Value = serde_json::from_str(body).expect("body is JSON");
assert_eq!(
body.get("idle_after_hours").cloned(),
expected,
"idle_after_hours rides the wire only when configured: {body}"
);
}
}
#[tokio::test]
async fn coder_http_client_status_and_transitions_over_loopback() {
let recorded = Arc::new(Mutex::new(Vec::new()));
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("bind loopback");
let base_url = format!("http://{}", listener.local_addr().expect("addr"));
let client = CoderHttpClient::new(&base_url, "tok").expect("client");
let ((), status) = tokio::join!(
serve_once(
listener,
r#"{"latest_build":{"status":"running"}}"#,
Arc::clone(&recorded)
),
tokio::time::timeout(Duration::from_secs(10), client.workspace_status("ws-9"))
);
let status = status.expect("bounded").expect("status");
assert_eq!(status, SubstrateStatus::Ready);
assert!(
recorded.lock().unwrap()[0].starts_with("GET /api/v2/workspaces/ws-9 "),
"status path: {:?}",
recorded.lock().unwrap()[0].lines().next()
);
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("bind loopback");
let base_url = format!("http://{}", listener.local_addr().expect("addr"));
let client = CoderHttpClient::new(&base_url, "tok").expect("client");
let ((), stopped) = tokio::join!(
serve_once(listener, "{}", Arc::clone(&recorded)),
tokio::time::timeout(Duration::from_secs(10), client.stop_workspace("ws-9"))
);
stopped.expect("bounded").expect("stop");
let stop_request = recorded.lock().unwrap().last().unwrap().clone();
assert!(
stop_request.starts_with("POST /api/v2/workspaces/ws-9/builds "),
"{stop_request}"
);
assert!(
stop_request.contains(r#""transition":"stop""#),
"the stop transition body: {stop_request}"
);
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("bind loopback");
let base_url = format!("http://{}", listener.local_addr().expect("addr"));
let client = CoderHttpClient::new(&base_url, "tok").expect("client");
let ((), deleted) = tokio::join!(
serve_once(listener, "{}", Arc::clone(&recorded)),
tokio::time::timeout(Duration::from_secs(10), client.delete_workspace("ws-9"))
);
deleted.expect("bounded").expect("delete");
let delete_request = recorded.lock().unwrap().last().unwrap().clone();
assert!(
delete_request.contains(r#""transition":"delete""#),
"the delete transition body: {delete_request}"
);
}
}