Merge branch 'work/companion' into develop

This commit is contained in:
2026-08-21 07:21:09 +09:00
25 changed files with 1642 additions and 396 deletions
+88 -45
View File
@@ -59,7 +59,6 @@ pub struct WorkerRetentionPolicyUpdate {
pub struct WorkerRemovalPlanRequest {
pub workspace_id: String,
pub worker: RuntimeWorkerRef,
pub expected_worker_revision: String,
pub reason: String,
}
@@ -154,8 +153,6 @@ pub enum WorkerRetentionError {
WorkerNotFound,
#[error("Worker belongs to a different Workspace")]
CrossWorkspace,
#[error("Worker revision changed: expected {expected}, current {actual}")]
WorkerRevisionConflict { expected: String, actual: String },
#[error("Worker removal is blocked: {0:?}")]
Blocked(Vec<WorkerRemovalBlocker>),
#[error("Worker removal plan {plan_id} is stale: {reason}")]
@@ -308,17 +305,16 @@ impl SqliteWorkspaceStore {
return Err(StoreError::InvalidInput(if other{"cross-workspace".into()}else{"worker-missing".into()}));
}
};
if worker.updated_at!=req.expected_worker_revision { return Err(StoreError::InvalidInput(format!("worker-conflict:{}:{}",req.expected_worker_revision,worker.updated_at))); }
let mut blockers=Vec::new();
if worker.retention_state=="pinned" { blockers.push(WorkerRemovalBlocker::Hold); }
if let Some((assignment_id,ticket_id))=tx.query_row("SELECT a.assignment_id,a.ticket_id FROM ticket_current_worker_assignments c JOIN ticket_worker_assignments a ON a.workspace_id=c.workspace_id AND a.ticket_id=c.ticket_id AND a.assignment_id=c.assignment_id WHERE a.workspace_id=?1 AND a.runtime_id=?2 AND a.worker_id=?3",params![req.workspace_id,req.worker.runtime_id,req.worker.worker_id],|r|Ok((r.get(0)?,r.get(1)?))).optional()? {
blockers.push(WorkerRemovalBlocker::CurrentAssignment{assignment_id,ticket_id});
}
let fp=fingerprint(req,inv,&policy,&blockers)?;
let fp=fingerprint(req,&worker.updated_at,inv,&policy,&blockers)?;
let plan_id=stable("wrp",&fp); let operation_id=stable("wro",&fp);
let archive_id=(policy.session_disposition==SessionDisposition::Archive).then(||stable("wra",&fp));
let state=if blockers.is_empty(){WorkerRemovalPlanState::Planned}else{WorkerRemovalPlanState::Blocked};
tx.execute("INSERT OR IGNORE INTO worker_removal_operations(operation_id,plan_id,input_fingerprint,workspace_id,runtime_id,worker_id,worker_revision,run_generation,policy_id,policy_revision,session_disposition,metadata_disposition,archive_retention_kind,archive_retention_seconds,diagnostics_disposition,diagnostics_retention_seconds,archive_id,blockers_json,state,reason,created_at,updated_at) VALUES(?1,?2,?3,?4,?5,?6,?7,?8,?9,?10,?11,?12,?13,?14,?15,?16,?17,?18,?19,?20,?21,?21)",params![operation_id,plan_id,fp,req.workspace_id,req.worker.runtime_id,req.worker.worker_id,req.expected_worker_revision,inv.run_generation,policy.policy_id,policy.revision,sess(policy.session_disposition),meta(policy.metadata_disposition),archive_kind(policy.archive_retention),archive_seconds(policy.archive_retention),diag(policy.diagnostics_disposition),policy.diagnostics_retention_seconds,archive_id,serde_json::to_string(&blockers).map_err(|e|StoreError::InvalidInput(e.to_string()))?,state_s(state),req.reason,now])?;
tx.execute("INSERT OR IGNORE INTO worker_removal_operations(operation_id,plan_id,input_fingerprint,workspace_id,runtime_id,worker_id,worker_revision,run_generation,policy_id,policy_revision,session_disposition,metadata_disposition,archive_retention_kind,archive_retention_seconds,diagnostics_disposition,diagnostics_retention_seconds,archive_id,blockers_json,state,reason,created_at,updated_at) VALUES(?1,?2,?3,?4,?5,?6,?7,?8,?9,?10,?11,?12,?13,?14,?15,?16,?17,?18,?19,?20,?21,?21)",params![operation_id,plan_id,fp,req.workspace_id,req.worker.runtime_id,req.worker.worker_id,worker.updated_at,inv.run_generation,policy.policy_id,policy.revision,sess(policy.session_disposition),meta(policy.metadata_disposition),archive_kind(policy.archive_retention),archive_seconds(policy.archive_retention),diag(policy.diagnostics_disposition),policy.diagnostics_retention_seconds,archive_id,serde_json::to_string(&blockers).map_err(|e|StoreError::InvalidInput(e.to_string()))?,state_s(state),req.reason,now])?;
let plan=load_plan(&tx,&plan_id)?.ok_or_else(||StoreError::InvalidInput("plan missing".into()))?;
if plan.input_fingerprint!=fp{return Err(StoreError::InvalidInput(format!("fingerprint:{}",plan.operation_id)));}
tx.commit()?; Ok(plan)
@@ -422,27 +418,22 @@ impl SqliteWorkspaceStore {
&self,
workspace_id: &str,
worker: &RuntimeWorkerRef,
expected_worker_revision: &str,
reason: &str,
) -> Result<Option<PreparedWorkerRemoval>, WorkerRetentionError> {
bounded("workspace", workspace_id, 160)?;
bounded("revision", expected_worker_revision, 256)?;
bounded("reason", reason, 512)?;
let plan = self.with_conn(|conn| {
conn.query_row(
"SELECT plan_id FROM worker_removal_operations
WHERE workspace_id=?1 AND runtime_id=?2 AND worker_id=?3
AND worker_revision=?4 AND reason=?5
AND state IN ('executing','failed','succeeded')
AND state IN ('planned','executing','failed','succeeded')
AND (
state='succeeded' OR worker_revision=(
SELECT updated_at FROM worker_registry
WHERE workspace_id=?1 AND runtime_id=?2 AND worker_id=?3
)
)
ORDER BY CASE state WHEN 'succeeded' THEN 0 ELSE 1 END,
created_at DESC LIMIT 1",
params![
workspace_id,
worker.runtime_id,
worker.worker_id,
expected_worker_revision,
reason,
],
params![workspace_id, worker.runtime_id, worker.worker_id],
|row| row.get::<_, String>(0),
)
.optional()
@@ -848,6 +839,7 @@ fn stale_error(plan: &WorkerRemovalPlan, reason: &str) -> StoreError {
}
fn fingerprint(
r: &WorkerRemovalPlanRequest,
worker_revision: &str,
i: &WorkerRetentionInventory,
p: &WorkerRetentionPolicy,
b: &[WorkerRemovalBlocker],
@@ -856,7 +848,7 @@ fn fingerprint(
r.workspace_id,
r.worker.runtime_id,
r.worker.worker_id,
r.expected_worker_revision,
worker_revision,
i.run_generation,
i.session_id,
i.segment_ids,
@@ -887,7 +879,6 @@ fn validate_plan(
i: &WorkerRetentionInventory,
) -> Result<(), WorkerRetentionError> {
bounded("workspace", &r.workspace_id, 160)?;
bounded("revision", &r.expected_worker_revision, 256)?;
bounded("reason", &r.reason, 2000)?;
if i.workspace_id != r.workspace_id
|| i.runtime_id != r.worker.runtime_id
@@ -947,13 +938,6 @@ fn map_error(e: StoreError) -> WorkerRetentionError {
if m == "worker-missing" {
return WorkerRetentionError::WorkerNotFound;
}
if let Some(x) = m.strip_prefix("worker-conflict:") {
let mut s = x.splitn(2, ':');
return WorkerRetentionError::WorkerRevisionConflict {
expected: s.next().unwrap_or_default().into(),
actual: s.next().unwrap_or_default().into(),
};
}
if let Some(x) = m.strip_prefix("fingerprint:") {
return WorkerRetentionError::OperationFingerprintConflict {
operation_id: x.into(),
@@ -1111,7 +1095,6 @@ mod tests {
runtime_id: "r".into(),
worker_id: worker_id().to_string(),
},
expected_worker_revision: "rev1".into(),
reason: "cleanup".into(),
}
}
@@ -1146,6 +1129,7 @@ mod tests {
let a = s.plan_worker_removal(&req(), &inv()).unwrap();
let b = s.plan_worker_removal(&req(), &inv()).unwrap();
assert_eq!(a.plan_id, b.plan_id);
assert_eq!(a.worker_revision, "rev1");
s.with_conn(|c| {
c.execute(
"UPDATE worker_registry SET retention_state='pinned' WHERE workspace_id='w'",
@@ -1154,9 +1138,7 @@ mod tests {
Ok(())
})
.unwrap();
let mut q = req();
q.expected_worker_revision = "rev1".into();
let p = s.plan_worker_removal(&q, &inv()).unwrap();
let p = s.plan_worker_removal(&req(), &inv()).unwrap();
assert_eq!(p.blockers, vec![WorkerRemovalBlocker::Hold]);
assert!(matches!(
s.begin_worker_removal("w", &p.plan_id, &p.input_fingerprint),
@@ -1580,6 +1562,77 @@ mod tests {
);
}
#[test]
fn worker_removal_recovery_is_target_keyed_and_preserves_original_reason() {
let store = setup();
let request = req();
let plan = store.plan_worker_removal(&request, &inv()).unwrap();
let recovered_planned = store
.recover_worker_removal_execution("w", &request.worker)
.unwrap()
.unwrap();
assert_eq!(
recovered_planned.plan.state,
WorkerRemovalPlanState::Planned
);
assert_eq!(recovered_planned.plan.plan_id, plan.plan_id);
store
.prepare_worker_removal_execution("w", &plan.plan_id, &plan.input_fingerprint)
.unwrap();
store
.fail_worker_removal(
"w",
&plan.operation_id,
&plan.input_fingerprint,
"runtime_remove_failed",
)
.unwrap();
let recovered = store
.recover_worker_removal_execution("w", &request.worker)
.unwrap()
.unwrap();
assert_eq!(recovered.plan.plan_id, plan.plan_id);
assert_eq!(recovered.plan.reason, request.reason);
}
#[test]
fn stale_failed_removal_is_not_recovered_after_worker_authority_changes() {
let store = setup();
let request = req();
let plan = store.plan_worker_removal(&request, &inv()).unwrap();
store
.prepare_worker_removal_execution("w", &plan.plan_id, &plan.input_fingerprint)
.unwrap();
store
.fail_worker_removal(
"w",
&plan.operation_id,
&plan.input_fingerprint,
"runtime_remove_failed",
)
.unwrap();
store
.with_conn(|conn| {
conn.execute(
"UPDATE worker_registry SET updated_at='rev2' WHERE workspace_id='w' AND runtime_id='r' AND worker_id=?1",
[worker_id().to_string()],
)?;
Ok(())
})
.unwrap();
assert!(
store
.recover_worker_removal_execution("w", &request.worker)
.unwrap()
.is_none()
);
let replacement = store.plan_worker_removal(&request, &inv()).unwrap();
assert_eq!(replacement.worker_revision, "rev2");
assert_ne!(replacement.plan_id, plan.plan_id);
}
#[test]
fn succeeded_worker_removal_recovers_after_registry_purge() {
let s = setup();
@@ -1632,18 +1685,13 @@ mod tests {
.is_none()
);
let recovered = s
.recover_worker_removal_execution(
"w",
&request.worker,
&request.expected_worker_revision,
&request.reason,
)
.recover_worker_removal_execution("w", &request.worker)
.unwrap()
.unwrap();
assert_eq!(recovered.plan.state, WorkerRemovalPlanState::Succeeded);
assert_eq!(
recovered.runtime_request.expected_worker_revision,
request.expected_worker_revision
plan.worker_revision
);
}
@@ -1662,12 +1710,7 @@ mod tests {
)
.unwrap();
let recovered = s
.recover_worker_removal_execution(
"w",
&request.worker,
&request.expected_worker_revision,
&request.reason,
)
.recover_worker_removal_execution("w", &request.worker)
.unwrap()
.unwrap();
assert_eq!(
+57 -103
View File
@@ -411,7 +411,6 @@ impl WorkspaceWorkerRemoveExecutor {
source: crate::worker_source::VerifiedWorkerMutationSource,
target_runtime_id: &str,
target_worker_id: &str,
expected_worker_revision: &str,
reason: &str,
) -> std::result::Result<worker::WorkspaceResponse, String> {
let reason = reason.trim();
@@ -509,73 +508,63 @@ impl WorkspaceWorkerRemoveExecutor {
let prepared = self
.store
.recover_worker_removal_execution(
&self.workspace_id,
&target,
expected_worker_revision,
reason,
)
.recover_worker_removal_execution(&self.workspace_id, &target)
.map_err(|_| "Worker removal recovery authority is unavailable".to_string())?;
if let Some(prepared) = prepared {
if prepared.plan.state == crate::retention::WorkerRemovalPlanState::Succeeded {
return Ok(worker_remove_success_response(&target));
}
let must_close_session =
prepared.prior_failure_category.as_deref() == Some("workdir_session_close_failed");
let must_release_attachment = must_close_session
|| prepared.prior_failure_category.as_deref()
== Some("workdir_attachment_release_failed");
let prepared =
if prepared.plan.state == crate::retention::WorkerRemovalPlanState::Failed {
match self.store.prepare_worker_removal_execution(
&self.workspace_id,
&prepared.plan.plan_id,
&prepared.plan.input_fingerprint,
) {
Ok(prepared) => prepared,
Err(error) => return Ok(worker_retention_error_response(error)),
}
} else {
prepared
};
if must_close_session {
let session = {
self.workdir_sessions
.lock()
.map_err(|_| "Workdir session registry was poisoned".to_string())?
.get(&target)
.cloned()
};
if let Some(session) = session {
if session.close().await.is_err() {
let _ = self.store.fail_worker_removal(
&self.workspace_id,
&prepared.plan.operation_id,
&prepared.plan.input_fingerprint,
"workdir_session_close_failed",
);
return Ok(worker_remove_error_response(
StatusCode::SERVICE_UNAVAILABLE,
"attachment_close_failed",
"Worker Workdir session could not be closed; removal can be retried",
));
}
self.workdir_sessions
.lock()
.map_err(|_| "Workdir session registry was poisoned".to_string())?
.remove(&target);
let prepared = if matches!(
prepared.plan.state,
crate::retention::WorkerRemovalPlanState::Planned
| crate::retention::WorkerRemovalPlanState::Failed
) {
match self.store.prepare_worker_removal_execution(
&self.workspace_id,
&prepared.plan.plan_id,
&prepared.plan.input_fingerprint,
) {
Ok(prepared) => prepared,
Err(error) => return Ok(worker_retention_error_response(error)),
}
}
if must_release_attachment
&& self
.store
.detach_worker_workdir(
} else {
prepared
};
let session = {
self.workdir_sessions
.lock()
.map_err(|_| "Workdir session registry was poisoned".to_string())?
.get(&target)
.cloned()
};
if let Some(session) = session {
if session.close().await.is_err() {
let _ = self.store.fail_worker_removal(
&self.workspace_id,
&target,
None,
&Utc::now().to_rfc3339_opts(SecondsFormat::Millis, true),
)
.is_err()
&prepared.plan.operation_id,
&prepared.plan.input_fingerprint,
"workdir_session_close_failed",
);
return Ok(worker_remove_error_response(
StatusCode::SERVICE_UNAVAILABLE,
"attachment_close_failed",
"Worker Workdir session could not be closed; removal can be retried",
));
}
self.workdir_sessions
.lock()
.map_err(|_| "Workdir session registry was poisoned".to_string())?
.remove(&target);
}
if self
.store
.detach_worker_workdir(
&self.workspace_id,
&target,
None,
&Utc::now().to_rfc3339_opts(SecondsFormat::Millis, true),
)
.is_err()
{
let _ = self.store.fail_worker_removal(
&self.workspace_id,
@@ -632,7 +621,6 @@ impl WorkspaceWorkerRemoveExecutor {
let request = crate::retention::WorkerRemovalPlanRequest {
workspace_id: self.workspace_id.clone(),
worker: target.clone(),
expected_worker_revision: expected_worker_revision.to_string(),
reason: reason.to_string(),
};
let plan = match self.store.plan_worker_removal(&request, &inventory) {
@@ -737,13 +725,11 @@ impl crate::worker_source::VerifiedWorkerRemoveExecutor for WorkspaceWorkerRemov
source: crate::worker_source::VerifiedWorkerMutationSource,
target_runtime_id: &str,
target_worker_id: &str,
expected_worker_revision: &str,
reason: &str,
) -> std::result::Result<worker::WorkspaceResponse, String> {
let executor = self.clone();
let target_runtime_id = target_runtime_id.to_string();
let target_worker_id = target_worker_id.to_string();
let expected_worker_revision = expected_worker_revision.to_string();
let reason = reason.to_string();
std::thread::spawn(move || {
tokio::runtime::Builder::new_current_thread()
@@ -754,7 +740,6 @@ impl crate::worker_source::VerifiedWorkerRemoveExecutor for WorkspaceWorkerRemov
source,
&target_runtime_id,
&target_worker_id,
&expected_worker_revision,
&reason,
))
})
@@ -6476,14 +6461,13 @@ fn worker_retention_error_response(
"worker_not_found",
"Worker was not found in this Workspace",
),
crate::retention::WorkerRetentionError::WorkerRevisionConflict { .. }
| crate::retention::WorkerRetentionError::PolicyRevisionConflict { .. }
crate::retention::WorkerRetentionError::PolicyRevisionConflict { .. }
| crate::retention::WorkerRetentionError::StalePlan { .. }
| crate::retention::WorkerRetentionError::OperationFingerprintConflict { .. } => {
worker_remove_error_response(
StatusCode::CONFLICT,
"worker_revision_conflict",
"Worker removal state changed; reread the Worker and retry",
"worker_removal_conflict",
"Worker removal state changed; retry the operation",
)
}
crate::retention::WorkerRetentionError::Blocked(_) => worker_remove_error_response(
@@ -6510,7 +6494,6 @@ fn worker_retention_error_response(
struct WorkerRemoveBoundaryRequest {
target_runtime_id: String,
target_worker_id: String,
expected_worker_revision: String,
reason: String,
}
@@ -6548,7 +6531,6 @@ async fn scoped_worker_remove_source_boundary(
source,
&request.target_runtime_id,
&request.target_worker_id,
&request.expected_worker_revision,
&request.reason,
)
.await
@@ -16357,7 +16339,7 @@ mod tests {
let temp = tempfile::tempdir().unwrap();
let app = build_router(test_api(temp.path()).await);
let body = r#"{"target_runtime_id":"runtime-target","target_worker_id":"target-worker","expected_worker_revision":"revision-1","reason":"retire target Worker"}"#;
let body = r#"{"target_runtime_id":"runtime-target","target_worker_id":"target-worker","reason":"retire target Worker"}"#;
let browser = app
.clone()
.oneshot(
@@ -16474,7 +16456,6 @@ mod tests {
fresh_proof,
"runtime-target",
"target-worker",
"revision-1",
"retire target Worker",
)
.unwrap_err();
@@ -16482,7 +16463,7 @@ mod tests {
}
#[tokio::test]
async fn worker_remove_rejects_self_running_and_stale_revision_at_caller_boundary() {
async fn worker_remove_rejects_self_and_running_at_caller_boundary() {
let temp = tempfile::tempdir().unwrap();
let api = test_api(temp.path()).await;
let Json(orchestrator) = scoped_start_workspace_orchestrator(
@@ -16507,7 +16488,6 @@ mod tests {
verified_source(),
&source.runtime_id,
&source.worker_id,
"irrelevant",
"must reject self",
)
.await
@@ -16549,37 +16529,12 @@ mod tests {
verified_source(),
&target.runtime_id,
&target.worker_id,
"irrelevant",
"must reject a live Worker",
)
.await
.unwrap();
assert_eq!(running_response.status, StatusCode::CONFLICT.as_u16());
assert!(running_response.body.contains("worker_not_stopped"));
api.runtime
.stop_worker(
&target,
WorkerLifecycleRequest {
reason: Some("prepare stale revision guard".to_string()),
ticket_assignment: None,
},
)
.unwrap();
let summary = api.runtime.worker(&target).unwrap();
let record = sync_worker_observation(&api, &summary).unwrap();
let stale_response = executor
.execute_async(
verified_source(),
&target.runtime_id,
&target.worker_id,
&format!("{}-stale", record.updated_at),
"must reject stale revision",
)
.await
.unwrap();
assert_eq!(stale_response.status, StatusCode::CONFLICT.as_u16());
assert!(stale_response.body.contains("worker_revision_conflict"));
}
#[tokio::test]
@@ -16653,7 +16608,7 @@ mod tests {
)
.unwrap();
let summary = api.runtime.worker(&target).unwrap();
let record = sync_worker_observation(&api, &summary).unwrap();
sync_worker_observation(&api, &summary).unwrap();
seed_worker_control_grant(&api, &source, &target, "embedded-valid-proof");
let response = WorkspaceWorkerRemoveExecutor::new(&api)
@@ -16667,7 +16622,6 @@ mod tests {
},
&target.runtime_id,
&target.worker_id,
&record.updated_at,
"retire completed Worker",
)
.await
@@ -16878,7 +16832,7 @@ mod tests {
route_token,
)
.body(Body::from(
r#"{"target_runtime_id":"runtime-target","target_worker_id":"target-worker","expected_worker_revision":"revision-1","reason":"retire target Worker"}"#,
r#"{"target_runtime_id":"runtime-target","target_worker_id":"target-worker","reason":"retire target Worker"}"#,
))
.unwrap(),
)
+2 -12
View File
@@ -653,13 +653,11 @@ pub trait ControlPlaneStore: Send + Sync {
&self,
workspace_id: &str,
worker: &RuntimeWorkerRef,
expected_worker_revision: &str,
reason: &str,
) -> std::result::Result<
Option<crate::retention::PreparedWorkerRemoval>,
crate::retention::WorkerRetentionError,
> {
let _ = (workspace_id, worker, expected_worker_revision, reason);
let _ = (workspace_id, worker);
Ok(None)
}
fn fail_worker_removal(
@@ -1650,19 +1648,11 @@ impl ControlPlaneStore for SqliteWorkspaceStore {
&self,
workspace_id: &str,
worker: &RuntimeWorkerRef,
expected_worker_revision: &str,
reason: &str,
) -> std::result::Result<
Option<crate::retention::PreparedWorkerRemoval>,
crate::retention::WorkerRetentionError,
> {
SqliteWorkspaceStore::recover_worker_removal_execution(
self,
workspace_id,
worker,
expected_worker_revision,
reason,
)
SqliteWorkspaceStore::recover_worker_removal_execution(self, workspace_id, worker)
}
fn fail_worker_removal(
+1 -9
View File
@@ -175,7 +175,6 @@ pub(crate) trait VerifiedWorkerRemoveExecutor: Send + Sync {
source: VerifiedWorkerMutationSource,
target_runtime_id: &str,
target_worker_id: &str,
expected_worker_revision: &str,
reason: &str,
) -> Result<worker::WorkspaceResponse, String>;
}
@@ -217,7 +216,6 @@ impl worker_runtime::worker_source::EmbeddedWorkerMutationDispatcher
proof: InProcessWorkerMutationProof,
target_runtime_id: &str,
target_worker_id: &str,
expected_worker_revision: &str,
reason: &str,
) -> Result<
worker::WorkspaceResponse,
@@ -241,13 +239,7 @@ impl worker_runtime::worker_source::EmbeddedWorkerMutationDispatcher
)
})?;
executor
.execute(
source,
target_runtime_id,
target_worker_id,
expected_worker_revision,
reason,
)
.execute(source, target_runtime_id, target_worker_id, reason)
.map_err(worker_runtime::worker_source::RuntimeWorkerMutationForwardError::Embedded)
}
}