use std::collections::HashMap;
use std::sync::{Arc, OnceLock};
use parking_lot::Mutex;
use strop_core::worker::CancelToken;
use strop_remote::worker_transport::RemoteWorker;
use strop_worker_deploy::deploy::{deploy, ArtifactSupply, Consent, DeployOutcome, DeployRequest};
use strop_worker_deploy::provider::DeployProvider;
use strop_worker_deploy::VerifiedObject;
use strop_workspace::operation::{FsFailure, FsFailureKind};
use strop_workspace::RemoteEndpoint;
use crate::editor::namespace::map_client;
use crate::editor::worker_catalog::{catalog_for, map_refusal, select_binary};
pub(crate) struct WorkerReady {
pub worker: RemoteWorker,
pub target: String,
pub artifact: VerifiedObject,
}
#[derive(Default)]
struct EndpointLease {
ready: OnceLock<Arc<WorkerReady>>,
deploying: Mutex<()>,
}
#[derive(Default)]
struct EndpointTable {
closed: bool,
endpoints: HashMap<RemoteEndpoint, Arc<EndpointLease>>,
}
#[derive(Clone, Default)]
pub(crate) struct RemoteWorkers {
inner: Arc<Mutex<EndpointTable>>,
}
fn failure(kind: FsFailureKind, detail: impl Into<String>) -> FsFailure {
FsFailure::new(kind, detail)
}
impl RemoteWorkers {
pub(crate) fn get(&self, endpoint: &RemoteEndpoint) -> Option<RemoteWorker> {
self.get_ready(endpoint).map(|ready| ready.worker.clone())
}
pub(crate) fn get_ready(&self, endpoint: &RemoteEndpoint) -> Option<Arc<WorkerReady>> {
let table = self.inner.lock();
if table.closed {
return None;
}
table
.endpoints
.get(endpoint)
.and_then(|lease| lease.ready.get())
.cloned()
}
pub(crate) fn admitting(&self, endpoint: &RemoteEndpoint) -> bool {
let table = self.inner.lock();
!table.closed
&& table
.endpoints
.get(endpoint)
.is_some_and(|lease| lease.deploying.try_lock().is_none())
}
pub(crate) fn admit(
&self,
endpoint: &RemoteEndpoint,
action: &str,
token: &CancelToken,
) -> Result<RemoteWorker, FsFailure> {
let lease = {
let mut table = self.inner.lock();
if table.closed {
return Err(failure(
FsFailureKind::Cancelled,
"editor worker owner closed",
));
}
Arc::clone(
table
.endpoints
.entry(endpoint.clone())
.or_insert_with(|| Arc::new(EndpointLease::default())),
)
};
if let Some(ready) = lease.ready.get() {
return Ok(ready.worker.clone());
}
let _deploying = lease.deploying.lock();
{
let table = self.inner.lock();
if table.closed {
return Err(failure(
FsFailureKind::Cancelled,
"editor worker owner closed",
));
}
if let Some(ready) = lease.ready.get() {
return Ok(ready.worker.clone());
}
}
let facts = strop_remote::bootstrap::discover(endpoint, token).map_err(|error| {
failure(
FsFailureKind::Permission,
format!("worker unavailable on {endpoint}: {error}"),
)
})?;
let Some(target) = facts.local_binary_target() else {
return Err(failure(
FsFailureKind::Unsupported,
format!(
"no worker artifact for {}/{} on this build; obtain the matching artifact through the release pipeline",
facts.system, facts.machine
),
));
};
let binary = select_binary(target)?;
let catalog = catalog_for(target, &binary)?;
let provider =
strop_remote::deploy_provider::SftpDeployProvider::connect(endpoint, &facts, token)
.map_err(|error| {
failure(
FsFailureKind::Io,
format!("deployment channel to {endpoint}: {error}"),
)
})?;
let report = deploy(
&provider,
&DeployRequest {
catalog,
editor_version: env!("CARGO_PKG_VERSION").to_string(),
consent: Consent::Granted {
action: action.to_string(),
},
supply: ArtifactSupply::LocalBinary { path: binary },
online: false,
},
);
match report.outcome {
DeployOutcome::Probed(ready) => {
let worker = RemoteWorker::connect(endpoint, &ready.object.path, target);
worker.worker().capabilities().map_err(|error| {
failure(
FsFailureKind::Io,
format!("live worker at {endpoint} could not start: {error}"),
)
})?;
let maintenance = worker
.worker()
.collect_cache(token, &provider.endpoint().context)
.map_err(map_client)?;
if let Some(error) = maintenance.failure {
return Err(failure(
error.kind,
format!(
"cache maintenance after {} retirements: {}",
maintenance.report.removed_objects.len(),
error.detail
),
));
}
let info = Arc::new(WorkerReady {
worker: worker.clone(),
target: target.to_owned(),
artifact: ready.object,
});
let published = {
let table = self.inner.lock();
if table.closed {
false
} else {
let published = lease.ready.set(info);
debug_assert!(published.is_ok(), "one deployment per endpoint lock");
true
}
};
if !published {
worker.worker().close();
return Err(failure(
FsFailureKind::Cancelled,
"editor worker owner closed",
));
}
Ok(worker)
}
DeployOutcome::Fallback(fallback) => {
Err(failure(FsFailureKind::Unsupported, fallback.to_string()))
}
DeployOutcome::Refused(refusal) => Err(map_refusal(refusal)),
DeployOutcome::PublishedNotReady { reason, .. } => Err(failure(
FsFailureKind::Io,
format!("worker published on {endpoint} but not ready: {reason}"),
)),
}
}
pub(crate) fn close_all(&self) {
let leases = {
let mut table = self.inner.lock();
if table.closed {
return;
}
table.closed = true;
std::mem::take(&mut table.endpoints)
};
for (_, lease) in leases {
if let Some(ready) = lease.ready.get() {
ready.worker.worker().close();
}
}
}
}
#[cfg(test)]
mod tests {
#[test]
fn synthesized_catalog_parses_and_binds_this_build() {
let directory = tempfile::tempdir().expect("tempdir");
let binary = directory.path().join("strop");
std::fs::write(&binary, b"worker bytes").expect("write");
let target = strop_remote::bootstrap::local_target().expect("a shipping target");
let catalog = super::catalog_for(target, &binary).expect("catalog parses");
assert_eq!(catalog.version, env!("CARGO_PKG_VERSION"));
assert_eq!(
catalog.worker.min_editor,
strop_worker_deploy::MIN_EDITOR_VERSION
);
match catalog.compatibility(
env!("CARGO_PKG_VERSION"),
strop_worker_protocol::PROTOCOL_VERSION,
target,
) {
strop_worker_deploy::Compatibility::Compatible(artifact) => {
assert_eq!(artifact.target, target);
assert_eq!(artifact.bytes, 12);
}
other => panic!("the catalog must bind this build, got {other:?}"),
}
assert!(matches!(
catalog.compatibility("0.0.0", strop_worker_protocol::PROTOCOL_VERSION, target),
strop_worker_deploy::Compatibility::Fallback(_)
));
}
#[test]
fn lookup_does_not_park_behind_an_endpoint_deployment() {
use super::*;
let workers = RemoteWorkers::default();
let endpoint = RemoteEndpoint::parse("ssh://test.example").unwrap();
let lease = Arc::new(EndpointLease::default());
workers
.inner
.lock()
.endpoints
.insert(endpoint.clone(), Arc::clone(&lease));
let (held_tx, held_rx) = std::sync::mpsc::channel();
let (release_tx, release_rx) = std::sync::mpsc::channel();
let deployment = std::thread::spawn(move || {
let _guard = lease.deploying.lock();
held_tx.send(()).unwrap();
release_rx.recv().unwrap();
});
held_rx.recv().unwrap();
let (reply_tx, reply_rx) = std::sync::mpsc::channel();
let lookup = std::thread::spawn(move || reply_tx.send(workers.get(&endpoint)).unwrap());
assert!(reply_rx
.recv_timeout(std::time::Duration::from_secs(2))
.expect("lookup waited for the deployment")
.is_none());
release_tx.send(()).unwrap();
deployment.join().unwrap();
lookup.join().unwrap();
}
}