use std::collections::HashMap;
use std::sync::{Arc, OnceLock};
use parking_lot::Mutex;
use strop_containers::{ContainerIdentity, EngineRef};
use strop_core::worker::CancelToken;
use strop_worker_client::{ClientError, Worker};
use strop_worker_deploy::provider::{DeployProvider, ProviderError};
use strop_worker_deploy::{
deploy, ArtifactSupply, Consent, ContainerProvider, DeployOrigin, DeployOutcome, DeployRequest,
ShellPolicy, VerifiedObject,
};
use strop_workspace::operation::{FsFailure, FsFailureKind};
use strop_workspace::{ContainerId, DirectorySnapshot, Filesystem, ResourceLocation};
use crate::editor::namespace::map_client;
use crate::editor::worker_catalog::{catalog_for, map_refusal, select_binary};
fn failure(kind: FsFailureKind, detail: impl Into<String>) -> FsFailure {
FsFailure::new(kind, detail)
}
fn provider_failure(error: ProviderError) -> FsFailure {
let kind = match &error {
ProviderError::Permission(_) => FsFailureKind::Permission,
ProviderError::ReadOnly(_) | ProviderError::NoExec(_) | ProviderError::NotFound(_) => {
FsFailureKind::Unsupported
}
ProviderError::Handshake(_) | ProviderError::CorruptCache(_) => FsFailureKind::Protocol,
ProviderError::Transport(_) | ProviderError::DiskFull(_) => FsFailureKind::Io,
};
failure(kind, error.to_string())
}
#[derive(Default)]
struct Entry {
ready: OnceLock<Arc<BoundWorker>>,
admitting: Mutex<()>,
engine: Mutex<Option<EngineRef>>,
}
struct Incarnation {
started_at: String,
entry: Arc<Entry>,
}
impl Incarnation {
fn new(identity: &ContainerIdentity) -> Self {
Self {
started_at: identity.started_at.clone(),
entry: Arc::new(Entry::default()),
}
}
}
#[derive(Default)]
struct IncarnationTable {
closed: bool,
entries: HashMap<String, Incarnation>,
}
pub(crate) struct BoundWorker {
id: ContainerId,
started_at: String,
worker: Worker,
target: String,
artifact: VerifiedObject,
}
impl BoundWorker {
pub(crate) fn worker(&self) -> &Worker {
&self.worker
}
pub(crate) fn target(&self) -> &str {
&self.target
}
pub(crate) fn artifact_path(&self) -> &str {
&self.artifact.path
}
pub(crate) fn artifact_sha256(&self) -> &str {
&self.artifact.sha256
}
pub(crate) fn matches(&self, identity: &ContainerIdentity) -> bool {
identity.id == self.id.as_str() && identity.started_at == self.started_at
}
fn location(&self, location: &ResourceLocation) -> Result<ResourceLocation, ClientError> {
if !matches!(&location.filesystem, Filesystem::Container(id) if id == &self.id) {
return Err(ClientError::Protocol(
strop_worker_protocol::ProtocolError::Unexpected {
message: "container location belongs to another worker namespace".into(),
},
));
}
Ok(ResourceLocation::local(location.path.clone()))
}
pub(crate) fn observe(
&self,
token: &CancelToken,
location: &ResourceLocation,
) -> Result<Option<strop_workspace::Observation>, ClientError> {
let native = self.location(location)?;
let observations = self.worker.observe(token, vec![native.clone()])?;
let mut observations = observations.into_iter();
let Some(observation) = observations.next() else {
return Err(ClientError::Protocol(
strop_worker_protocol::ProtocolError::Unexpected {
message: "container worker returned no observation".into(),
},
));
};
if observation.location != native || observations.next().is_some() {
return Err(ClientError::Protocol(
strop_worker_protocol::ProtocolError::Unexpected {
message: "container worker returned a foreign observation".into(),
},
));
}
Ok(observation.value)
}
pub(crate) fn list(
&self,
token: &CancelToken,
location: &ResourceLocation,
) -> Result<DirectorySnapshot, ClientError> {
let native = self.location(location)?;
let mut snapshot = self.worker.list(token, native.clone())?;
if snapshot.location != native {
return Err(ClientError::Protocol(
strop_worker_protocol::ProtocolError::Unexpected {
message: "container worker returned a foreign listing".into(),
},
));
}
snapshot.location = location.clone();
Ok(snapshot)
}
pub(crate) fn read(
&self,
token: &CancelToken,
location: &ResourceLocation,
length: u64,
) -> Result<strop_worker_client::ReadPayload, ClientError> {
self.worker
.read(token, self.location(location)?, 0, Some(length))
}
}
#[derive(Clone, Default)]
pub(crate) struct ContainerWorkers {
inner: Arc<Mutex<IncarnationTable>>,
}
impl ContainerWorkers {
pub(crate) fn note_inspect(
&self,
identity: &ContainerIdentity,
engine: EngineRef,
) -> Result<(), FsFailure> {
let mut table = self.inner.lock();
if table.closed {
return Err(failure(
FsFailureKind::Cancelled,
"editor worker owner closed",
));
}
let slot = table
.entries
.entry(identity.id.clone())
.or_insert_with(|| Incarnation::new(identity));
if slot.started_at != identity.started_at {
return Err(failure(
FsFailureKind::Conflict,
"container restarted; old attachment cannot bind the new incarnation",
));
}
let mut captured = slot.entry.engine.lock();
if captured
.as_ref()
.is_some_and(|prior| !prior.same_connection(&engine))
{
return Err(failure(
FsFailureKind::Conflict,
"selected Docker engine changed; old attachment cannot bind another engine",
));
}
*captured = Some(engine);
Ok(())
}
pub(crate) fn get(&self, identity: &ContainerIdentity) -> Option<Arc<BoundWorker>> {
self.inner
.lock()
.entries
.get(&identity.id)
.and_then(|slot| {
(slot.started_at == identity.started_at)
.then(|| slot.entry.ready.get().cloned())
.flatten()
})
}
pub(crate) fn get_for(&self, id: &ContainerId, started_at: &str) -> Option<Arc<BoundWorker>> {
self.inner.lock().entries.get(id.as_str()).and_then(|slot| {
(slot.started_at == started_at)
.then(|| slot.entry.ready.get().cloned())
.flatten()
})
}
pub(crate) fn admit(
&self,
identity: &ContainerIdentity,
action: &str,
token: &CancelToken,
) -> Result<Arc<BoundWorker>, FsFailure> {
let id = ContainerId::canonical(identity.id.clone())
.map_err(|error| failure(FsFailureKind::Protocol, error.to_string()))?;
let entry = {
let table = self.inner.lock();
if table.closed {
return Err(failure(
FsFailureKind::Cancelled,
"editor worker owner closed",
));
}
table
.entries
.get(&identity.id)
.filter(|slot| slot.started_at == identity.started_at)
.map(|slot| Arc::clone(&slot.entry))
.ok_or_else(|| {
failure(
FsFailureKind::Conflict,
"attach this container incarnation before admitting its worker",
)
})?
};
if let Some(worker) = entry.ready.get() {
return Ok(Arc::clone(worker));
}
let _admitting = entry.admitting.lock();
if self.inner.lock().closed {
return Err(failure(
FsFailureKind::Cancelled,
"editor worker owner closed",
));
}
if let Some(worker) = entry.ready.get() {
return Ok(Arc::clone(worker));
}
let engine = entry.engine.lock().clone().ok_or_else(|| {
failure(
FsFailureKind::Unsupported,
"attach this container before admitting its worker",
)
})?;
let preinstalled = match std::env::var_os("STROP_CONTAINER_WORKER_PATH") {
Some(path) => Some(path.into_string().map_err(|_| {
failure(
FsFailureKind::InvalidPath,
"STROP_CONTAINER_WORKER_PATH must be UTF-8 container path text",
)
})?),
None => None,
};
let shell = if preinstalled.is_some() {
ShellPolicy::Absent
} else {
ShellPolicy::Required
};
let provider = ContainerProvider::capture(&engine, identity, shell, token)
.map_err(provider_failure)?;
let binary = select_binary(&provider.endpoint().target)?;
let catalog = catalog_for(&provider.endpoint().target, &binary)?;
let supply = match preinstalled {
Some(path) => ArtifactSupply::Preinstalled { path },
None => ArtifactSupply::LocalBinary { path: binary },
};
let report = deploy(
&provider,
&DeployRequest {
catalog,
editor_version: env!("CARGO_PKG_VERSION").into(),
consent: Consent::Granted {
action: action.to_string(),
},
supply,
online: false,
},
);
let ready = match report.outcome {
DeployOutcome::Probed(ready) => ready,
DeployOutcome::Fallback(reason) => {
return Err(failure(FsFailureKind::Unsupported, reason.to_string()));
}
DeployOutcome::Refused(reason) => return Err(map_refusal(reason)),
DeployOutcome::PublishedNotReady { reason, .. } => {
return Err(failure(FsFailureKind::Io, reason));
}
};
let worker_lease = provider.worker(&ready.object.path);
worker_lease.capabilities().map_err(|error| {
failure(
FsFailureKind::Io,
format!("container live worker could not start: {error}"),
)
})?;
if ready.origin != DeployOrigin::Preinstalled {
let maintenance = worker_lease
.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 worker = Arc::new(BoundWorker {
id,
started_at: identity.started_at.clone(),
target: provider.endpoint().target.clone(),
artifact: ready.object,
worker: worker_lease,
});
let published = {
let table = self.inner.lock();
if table.closed {
false
} else {
let published = entry.ready.set(Arc::clone(&worker));
debug_assert!(published.is_ok(), "one container admission per identity");
true
}
};
if !published {
worker.worker().close();
return Err(failure(
FsFailureKind::Cancelled,
"editor worker owner closed",
));
}
Ok(worker)
}
pub(crate) fn close_all(&self) {
let entries = {
let mut table = self.inner.lock();
if table.closed {
return;
}
table.closed = true;
std::mem::take(&mut table.entries)
};
for (_, incarnation) in entries {
if let Some(worker) = incarnation.entry.ready.get() {
worker.worker().close();
}
}
}
}