13 Commits
27 changed files with 3600 additions and 239 deletions
+29
View File
@@ -530,6 +530,15 @@ pub enum Event {
revision: u64,
event: Box<Event>,
},
/// Terminal removal fence for one parent-owned Internal Worker session.
///
/// Clients discard the matching child and descendants, then ignore later
/// nested events for this identity until an authoritative snapshot replaces
/// the projection.
InternalWorkerRemoved {
worker: InternalWorkerRef,
revision: u64,
},
/// Server-side segment log rotated to a fresh `SegmentStart`.
///
/// Fires on compaction and on auto-fork when the store head drifts
@@ -1932,6 +1941,26 @@ mod tests {
}
}
#[test]
fn internal_worker_removal_roundtrip_preserves_terminal_fence() {
let event = Event::InternalWorkerRemoved {
worker: InternalWorkerRef {
session_id: "session-1".into(),
name: "research".into(),
parent_session_id: Some("parent-session".into()),
kind: InternalWorkerKind::SubWorker,
},
revision: 8,
};
let json = serde_json::to_string(&event).unwrap();
let decoded: Event = serde_json::from_str(&json).unwrap();
assert!(matches!(
decoded,
Event::InternalWorkerRemoved { worker, revision }
if worker.session_id == "session-1" && revision == 8
));
}
#[test]
fn legacy_snapshot_defaults_internal_workers_to_empty() {
let snapshot: Event = serde_json::from_value(serde_json::json!({
+4 -6
View File
@@ -848,18 +848,16 @@ fn collect_foreign_key_diagnostics(
)
})
.collect::<BTreeSet<_>>();
// The Ticket component owns its required foreign keys, while an integrated host may
// strengthen Workspace/domain boundaries with additional references to host-owned
// tables. Reject missing component constraints, but do not treat those host extensions
// as Ticket schema drift.
for missing in expected.difference(&actual) {
push_diagnostic(
diagnostics,
format!("table {table:?} is missing foreign key {missing:?}"),
);
}
for unexpected in actual.difference(&expected) {
push_diagnostic(
diagnostics,
format!("table {table:?} has unexpected foreign key {unexpected:?}"),
);
}
}
fn collect_foreign_key_check_diagnostics(
+10 -3
View File
@@ -72,13 +72,20 @@ impl Tool for EditTool {
})
.await
.map_err(ToolsError::from)?;
self.tracker.record_workdir_hash(&path, result.content_hash);
let replacements = result.replacements;
self.tracker.record_workdir_edit(
&path,
result.content_hash,
replacements,
params.new_string.lines().count(),
params.old_string.lines().count(),
);
let summary = format!(
"Edited {} ({} replacement{})",
path,
result.replacements,
if result.replacements == 1 { "" } else { "s" }
replacements,
if replacements == 1 { "" } else { "s" }
);
let preview = make_preview(&params.new_string, &params.new_string);
+1 -1
View File
@@ -27,7 +27,7 @@ pub use error::ToolsError;
pub use glob::glob_tool;
pub use grep::grep_tool;
pub use read::read_tool;
pub use tracker::Tracker;
pub use tracker::{ChangeStat, Tracker};
pub use view_image::view_image_tool;
pub use web::{web_fetch_tool, web_search_tool};
pub use write::write_tool;
+90 -2
View File
@@ -119,12 +119,22 @@ fn normalize_path_lexically(path: &Path) -> PathBuf {
normalized
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct ChangeStat {
pub added: u64,
pub deleted: u64,
}
#[derive(Debug, Default)]
struct Inner {
/// Hash of each file's last observed contents, keyed by canonical path.
hashes: HashMap<PathBuf, ContentHash>,
/// Line count paired with observations that included the file content.
line_counts: HashMap<PathBuf, usize>,
/// LRU list of touched files. Front = most recently touched.
recency: VecDeque<PathBuf>,
/// Successful Write/Edit mutations attributed to this session's tools.
change_stat: ChangeStat,
}
/// Canonical-path keyed tracker of file observations and their recency.
@@ -187,8 +197,27 @@ impl Tracker {
}
}
pub fn record_workdir_content(&self, path: &workdir::WorkdirPath, bytes: &[u8]) {
self.record_workdir_hash(path, hash_bytes(bytes));
pub fn record_workdir_content(&self, path: &workdir::WorkdirPath, content: &[u8]) {
let key = PathBuf::from(path.as_str());
let hash = hash_bytes(content);
let line_count = String::from_utf8_lossy(content).lines().count();
let mut inner = self.inner.lock().unwrap_or_else(|e| e.into_inner());
inner.line_counts.insert(key.clone(), line_count);
inner.hashes.insert(key.clone(), hash);
inner.recency.retain(|candidate| candidate != &key);
inner.recency.push_front(key);
if inner.recency.len() > RECENCY_CAPACITY {
inner.recency.pop_back();
}
}
pub fn observed_workdir_line_count(&self, path: &workdir::WorkdirPath) -> Option<usize> {
self.inner
.lock()
.unwrap_or_else(|e| e.into_inner())
.line_counts
.get(Path::new(path.as_str()))
.copied()
}
pub fn record_workdir_hash(&self, path: &workdir::WorkdirPath, hash: workdir::ContentHash) {
@@ -202,6 +231,50 @@ impl Tracker {
}
}
/// Record a successful, session-attributable source mutation.
///
/// Callers supply line counts derived from the exact replacement accepted
/// by a Write/Edit tool. Bash and external process mutations are excluded
/// because this tracker cannot attribute them to one tool operation
/// authoritatively.
pub fn record_change(&self, added: usize, deleted: usize) {
let mut inner = self.inner.lock().unwrap_or_else(|e| e.into_inner());
inner.change_stat.added = inner.change_stat.added.saturating_add(added as u64);
inner.change_stat.deleted = inner.change_stat.deleted.saturating_add(deleted as u64);
}
pub fn record_workdir_edit(
&self,
path: &workdir::WorkdirPath,
hash: workdir::ContentHash,
replacements: usize,
added_lines_per_replacement: usize,
deleted_lines_per_replacement: usize,
) {
let added = added_lines_per_replacement.saturating_mul(replacements);
let deleted = deleted_lines_per_replacement.saturating_mul(replacements);
let key = PathBuf::from(path.as_str());
let mut inner = self.inner.lock().unwrap_or_else(|e| e.into_inner());
inner.change_stat.added = inner.change_stat.added.saturating_add(added as u64);
inner.change_stat.deleted = inner.change_stat.deleted.saturating_add(deleted as u64);
if let Some(line_count) = inner.line_counts.get_mut(&key) {
*line_count = line_count.saturating_sub(deleted).saturating_add(added);
}
inner.hashes.insert(key.clone(), hash);
inner.recency.retain(|candidate| candidate != &key);
inner.recency.push_front(key);
if inner.recency.len() > RECENCY_CAPACITY {
inner.recency.pop_back();
}
}
pub fn change_stat(&self) -> ChangeStat {
self.inner
.lock()
.unwrap_or_else(|e| e.into_inner())
.change_stat
}
pub fn expected_workdir_hash(
&self,
path: &workdir::WorkdirPath,
@@ -458,6 +531,21 @@ mod tests {
}
}
#[test]
fn change_stat_saturates_and_accumulates_tracked_mutations() {
let tracker = Tracker::new();
tracker.record_change(7, 3);
tracker.record_change(5, 2);
assert_eq!(
tracker.change_stat(),
ChangeStat {
added: 12,
deleted: 5,
}
);
}
#[tokio::test]
async fn mutation_guard_blocks_equivalent_paths_until_drop() {
let dir = tempfile::tempdir().unwrap();
+3
View File
@@ -50,6 +50,7 @@ impl Tool for WriteTool {
Err(error) => return Err(ToolsError::from(error).into()),
};
let old_line_count = self.tracker.observed_workdir_line_count(&path).unwrap_or(0);
let outcome = self
.session
.write(WriteRequest {
@@ -60,6 +61,8 @@ impl Tool for WriteTool {
.await
.map_err(ToolsError::from)?;
self.tracker
.record_change(params.content.lines().count(), old_line_count);
self.tracker
.record_workdir_content(&path, params.content.as_bytes());
+112 -1
View File
@@ -1,4 +1,4 @@
use std::collections::VecDeque;
use std::collections::{HashMap, VecDeque};
use std::path::Path;
use std::time::{Duration, Instant};
@@ -283,6 +283,8 @@ pub struct App {
/// Presentation-only Internal Worker projections keyed by session identity.
/// They are rendered in separate sub-panes and never mixed into `blocks`.
pub internal_workers: Vec<InternalWorkerView>,
/// Terminal child-session fences, reset only by an authoritative snapshot.
removed_internal_workers: HashMap<String, u64>,
pub scroll: Scroll,
pub mode: Mode,
pub cache: FileCache,
@@ -361,6 +363,7 @@ impl App {
blocks: Vec::new(),
run_error_messages: Vec::new(),
internal_workers: Vec::new(),
removed_internal_workers: HashMap::new(),
scroll: Scroll::default(),
mode: Mode::Normal,
cache: FileCache::new(),
@@ -1318,6 +1321,9 @@ impl App {
revision,
event,
} => self.apply_internal_worker_event(worker, revision, *event),
Event::InternalWorkerRemoved { worker, revision } => {
self.remove_internal_worker(worker, revision)
}
Event::Status { status } => {
self.rewind_refresh_fence = false;
self.set_worker_status(status);
@@ -2005,6 +2011,7 @@ impl App {
.into_iter()
.map(Self::internal_worker_view_from_snapshot)
.collect();
self.removed_internal_workers.clear();
}
fn internal_worker_view_from_snapshot(snapshot: InternalWorkerSnapshot) -> InternalWorkerView {
@@ -2032,6 +2039,12 @@ impl App {
revision: u64,
event: Event,
) {
if self
.removed_internal_workers
.contains_key(&worker.session_id)
{
return;
}
let index = self
.internal_workers
.iter()
@@ -2054,6 +2067,26 @@ impl App {
let _ = target.app.handle_worker_event(event);
}
fn remove_internal_worker(&mut self, worker: InternalWorkerRef, revision: u64) {
let Some(index) = self
.internal_workers
.iter()
.position(|candidate| candidate.worker.session_id == worker.session_id)
else {
self.removed_internal_workers
.entry(worker.session_id)
.and_modify(|current| *current = (*current).max(revision))
.or_insert(revision);
return;
};
if revision <= self.internal_workers[index].revision {
return;
}
self.internal_workers.remove(index);
self.removed_internal_workers
.insert(worker.session_id, revision);
}
fn restore_snapshot(
&mut self,
entries: &[serde_json::Value],
@@ -3545,6 +3578,84 @@ mod completion_flow_tests {
);
}
#[test]
fn terminal_internal_worker_removal_drops_descendants_and_fences_late_events() {
let mut app = App::new("parent".into());
let worker = InternalWorkerRef {
session_id: "child-session".into(),
name: "child".into(),
parent_session_id: Some("parent-session".into()),
kind: protocol::InternalWorkerKind::SubWorker,
};
let nested = InternalWorkerRef {
session_id: "grandchild-session".into(),
name: "grandchild".into(),
parent_session_id: Some("child-session".into()),
kind: protocol::InternalWorkerKind::SubWorker,
};
app.handle_worker_event(Event::InternalWorker {
worker: worker.clone(),
revision: 2,
event: Box::new(Event::InternalWorker {
worker: nested,
revision: 1,
event: Box::new(Event::TextDone {
text: "nested".into(),
}),
}),
});
assert_eq!(app.internal_workers.len(), 1);
assert_eq!(app.internal_workers[0].app.internal_workers.len(), 1);
app.handle_worker_event(Event::InternalWorkerRemoved {
worker: worker.clone(),
revision: 3,
});
app.handle_worker_event(Event::InternalWorker {
worker,
revision: 4,
event: Box::new(Event::TextDone {
text: "late".into(),
}),
});
assert!(app.internal_workers.is_empty());
app.handle_worker_event(Event::Snapshot {
greeting: test_greeting(),
entries: Vec::new(),
status: WorkerStatus::Idle,
in_flight: Default::default(),
internal_workers: Vec::new(),
});
assert!(app.internal_workers.is_empty());
assert!(app.removed_internal_workers.is_empty());
}
#[test]
fn stale_internal_worker_removal_keeps_newer_projection() {
let mut app = App::new("parent".into());
let worker = InternalWorkerRef {
session_id: "child-session".into(),
name: "child".into(),
parent_session_id: Some("parent-session".into()),
kind: protocol::InternalWorkerKind::SubWorker,
};
app.handle_worker_event(Event::InternalWorker {
worker: worker.clone(),
revision: 4,
event: Box::new(Event::TextDone {
text: "current".into(),
}),
});
app.handle_worker_event(Event::InternalWorkerRemoved {
worker,
revision: 3,
});
assert_eq!(app.internal_workers.len(), 1);
assert_eq!(app.internal_workers[0].revision, 4);
}
#[test]
fn snapshot_authoritatively_replaces_internal_worker_views() {
let mut app = App::new("parent".into());
@@ -16,7 +16,7 @@ use crate::feature::{
FeatureDescriptor, FeatureInstallContext, FeatureInstallError, FeatureModule,
ServiceDeclaration, ServiceId, ToolContribution, ToolDeclaration,
};
use crate::spawn::registry::SpawnedWorkerRegistry;
use crate::spawn::registry::{SpawnedWorkerRegistry, SubWorkerStopSummary};
use crate::worker::{
WorkspaceClient, WorkspaceClientError, WorkspaceRequest, WorkspaceRequestMethod,
WorkspaceResponse,
@@ -138,14 +138,19 @@ impl WorkerControlService for WorkspaceWorkerControlService {
let registry = self.registry.as_ref().ok_or_else(|| {
WorkspaceClientError::Request("unknown Worker or permission not granted".to_string())
})?;
registry
let summary = registry
.remove_internal(name)
.await
.map_err(|error| WorkspaceClientError::Request(error.to_string()))?;
.map_err(|error| WorkspaceClientError::Request(error.to_string()))?
.ok_or_else(|| {
WorkspaceClientError::Request(
"unknown Worker or permission not granted".to_string(),
)
})?;
Ok(WorkspaceResponse {
status: 200,
body: serde_json::json!({ "subject": { "kind": "sub_worker", "name": name } })
.to_string(),
body: serde_json::to_string(&summary)
.map_err(|error| WorkspaceClientError::Request(error.to_string()))?,
})
}
@@ -839,6 +844,15 @@ fn tool_output(
response.status, response.body
)));
}
if operation == WorkerOperation::Stop
&& let Ok(summary) = serde_json::from_str::<SubWorkerStopSummary>(&response.body)
{
return Ok(ToolOutput {
summary: render_subworker_stop_summary(&summary),
content: Some(response.body),
attachments: Vec::new(),
});
}
Ok(ToolOutput {
summary: format!("{} completed", operation.tool_name()),
content: Some(response.body),
@@ -846,6 +860,37 @@ fn tool_output(
})
}
fn render_subworker_stop_summary(summary: &SubWorkerStopSummary) -> String {
let tools = if summary.tool_counts.is_empty() {
"No tool calls".to_string()
} else {
summary
.tool_counts
.iter()
.map(|tool| format!("{} {}", tool.count, tool.name))
.collect::<Vec<_>>()
.join(", ")
};
let elapsed = format_elapsed(summary.elapsed_ms);
let changes = summary
.change_stat
.as_ref()
.map(|stat| format!("+{}/-{} Changes · ", stat.added, stat.deleted))
.unwrap_or_default();
format!("SubWorkerStop - done\n {tools}\n {changes}{elapsed}",)
}
fn format_elapsed(elapsed_ms: u64) -> String {
let seconds = elapsed_ms / 1_000;
let minutes = seconds / 60;
let seconds = seconds % 60;
if minutes > 0 {
format!("{minutes}m {seconds}s")
} else {
format!("{seconds}s")
}
}
fn definition<I: JsonSchema + 'static>(
operation: WorkerOperation,
control: Arc<dyn WorkerControlService>,
@@ -1252,6 +1297,47 @@ mod tests {
assert!(client.removals.lock().unwrap().is_empty());
}
#[test]
fn subworker_stop_output_is_compact_and_keeps_typed_evidence() {
let summary = SubWorkerStopSummary {
session_id: "session-1".to_string(),
display_name: "research".to_string(),
outcome: crate::spawn::registry::SubWorkerFinalOutcome::Done,
elapsed_ms: 78_000,
tool_counts: vec![
crate::spawn::registry::SubWorkerToolCount {
name: "Read".to_string(),
count: 26,
},
crate::spawn::registry::SubWorkerToolCount {
name: "Grep".to_string(),
count: 5,
},
],
change_stat: Some(crate::spawn::registry::SubWorkerChangeStat {
added: 215,
deleted: 148,
source: "tracked_write_edit_tools".to_string(),
}),
};
let response = WorkspaceResponse {
status: 200,
body: serde_json::to_string(&summary).unwrap(),
};
let output = tool_output(WorkerOperation::Stop, response).unwrap();
assert_eq!(
output.summary,
"SubWorkerStop - done\n 26 Read, 5 Grep\n +215/-148 Changes · 1m 18s"
);
assert_eq!(
serde_json::from_str::<SubWorkerStopSummary>(output.content.as_deref().unwrap())
.unwrap(),
summary
);
}
#[test]
fn worker_inputs_reject_paths_and_parent_traversal() {
assert!(authority_id("https://runtime.example", "runtime_id").is_err());
@@ -1781,6 +1781,7 @@ provider = "github"
assert!(request.contains("\"title\":\"HTTP ticket\""));
let response_body = serde_json::to_string(&TicketRef {
id: "01TEST".to_string(),
human_key: None,
slug: "http-ticket".to_string(),
status: ticket::TicketStatus::Open,
})
+34 -1
View File
@@ -294,6 +294,8 @@ pub(crate) struct InternalWorkerSessionHandle {
last_error: Arc<Mutex<Option<String>>>,
child_registry: Option<Arc<SpawnedWorkerRegistry>>,
sink: SegmentLogSink,
#[cfg(test)]
fail_stop: Arc<std::sync::atomic::AtomicBool>,
}
impl InternalWorkerSessionHandle {
@@ -319,6 +321,9 @@ impl InternalWorkerSessionHandle {
#[cfg(test)]
pub(crate) fn publish_test_entry(&self, entry: LogEntry) {
self.store
.append(self.session_id, self.segment_id, &entry)
.expect("append test Internal Worker entry");
self.sink.publish(entry);
}
@@ -419,7 +424,23 @@ impl InternalWorkerSessionHandle {
}
}
#[cfg(test)]
pub(crate) fn force_status(&self, status: InternalWorkerSessionStatus) {
self.status
.store(status.encode(), std::sync::atomic::Ordering::Release);
}
#[cfg(test)]
pub(crate) fn force_stop_failure(&self) {
self.fail_stop
.store(true, std::sync::atomic::Ordering::Release);
}
pub(crate) async fn stop(&self) -> Result<(), InternalWorkerSessionError> {
#[cfg(test)]
if self.fail_stop.load(std::sync::atomic::Ordering::Acquire) {
return Err(InternalWorkerSessionError::Unavailable);
}
let prior = self.status.swap(
InternalWorkerSessionStatus::Stopping.encode(),
std::sync::atomic::Ordering::AcqRel,
@@ -585,6 +606,8 @@ pub(crate) async fn prepare_internal_worker_session(
last_error: last_error.clone(),
child_registry,
sink,
#[cfg(test)]
fail_stop: Arc::new(std::sync::atomic::AtomicBool::new(false)),
};
tokio::spawn(async move {
@@ -893,8 +916,17 @@ pub(crate) fn test_internal_worker_session(
let session_id = session_store::new_session_id();
let segment_id = session_store::new_segment_id();
let (command_tx, mut command_rx) = tokio::sync::mpsc::channel(1);
tokio::spawn(async move { while command_rx.recv().await.is_some() {} });
let (event_tx, _) = broadcast::channel(256);
let command_event_tx = event_tx.clone();
tokio::spawn(async move {
while let Some(command) = command_rx.recv().await {
if let InternalWorkerSessionCommand::Stop(done_tx) = command {
let _ = command_event_tx.send(Event::Shutdown);
let _ = done_tx.send(());
break;
}
}
});
let sink = SegmentLogSink::new();
spawn_internal_log_event_bridge(sink.clone(), event_tx.clone());
let handle = InternalWorkerSessionHandle {
@@ -912,6 +944,7 @@ pub(crate) fn test_internal_worker_session(
last_error: Arc::new(Mutex::new(None)),
child_registry: None,
sink,
fail_stop: Arc::new(std::sync::atomic::AtomicBool::new(false)),
};
(handle, event_tx)
}
+18 -11
View File
@@ -170,20 +170,27 @@ impl Tool for SubWorkerStopTool {
) -> Result<ToolOutput, ToolError> {
let input: NameInput = serde_json::from_str(input_json)
.map_err(|e| ToolError::InvalidArgument(format!("invalid SubWorkerStop input: {e}")))?;
if let Some(record) = self.registry.get_internal(&input.name) {
record.session.stop().await.map_err(|error| {
ToolError::ExecutionFailed(format!("stop `{}`: {error}", input.name))
})?;
self.registry
.remove_internal(&input.name)
.await
.map_err(|error| ToolError::ExecutionFailed(error.to_string()))?;
if let Some(summary) = self
.registry
.remove_internal(&input.name)
.await
.map_err(|error| ToolError::ExecutionFailed(error.to_string()))?
{
return Ok(ToolOutput {
summary: format!(
"stopped worker `{}` and reclaimed delegated scope",
input.name
"SubWorkerStop - done\n {} tool kind{}\n {}ms",
summary.tool_counts.len(),
if summary.tool_counts.len() == 1 {
""
} else {
"s"
},
summary.elapsed_ms,
),
content: Some(
serde_json::to_string(&summary)
.map_err(|error| ToolError::ExecutionFailed(error.to_string()))?,
),
content: None,
attachments: Vec::new(),
});
}
+296 -8
View File
@@ -7,17 +7,20 @@
//! Parent registry drop closes all session handles and synchronously returns delegated Write deny
//! rules to the parent scope.
use std::collections::HashSet;
use std::collections::{BTreeMap, HashSet};
use std::io;
use std::sync::{
Arc, Mutex,
atomic::{AtomicBool, AtomicU64, Ordering},
};
use std::time::Instant;
use serde::{Deserialize, Serialize};
use manifest::{Permission, ScopeRule, SharedScope};
use protocol::{Event, InternalWorkerKind, InternalWorkerRef, InternalWorkerSnapshot};
use session_store::{
WorkerMetadataStore, WorkerReclaimedChild, WorkerSpawnedChild, WorkerStoreError,
LoggedItem, WorkerMetadataStore, WorkerReclaimedChild, WorkerSpawnedChild, WorkerStoreError,
};
use tokio::sync::broadcast;
use tracing::warn;
@@ -27,6 +30,39 @@ use crate::internal_worker::{InternalWorkerSessionHandle, InternalWorkerVisibili
use crate::runtime::dir::{RuntimeDir, SpawnedWorkerRecord};
use crate::runtime::worker_allocation;
const STOP_SUMMARY_TOOL_LIMIT: usize = 16;
const STOP_SUMMARY_TOOL_NAME_LIMIT: usize = 64;
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub(crate) enum SubWorkerFinalOutcome {
Done,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub(crate) struct SubWorkerToolCount {
pub name: String,
pub count: u64,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub(crate) struct SubWorkerChangeStat {
pub added: u64,
pub deleted: u64,
pub source: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub(crate) struct SubWorkerStopSummary {
pub session_id: String,
pub display_name: String,
pub outcome: SubWorkerFinalOutcome,
pub elapsed_ms: u64,
pub tool_counts: Vec<SubWorkerToolCount>,
#[serde(skip_serializing_if = "Option::is_none")]
pub change_stat: Option<SubWorkerChangeStat>,
}
#[derive(Clone)]
pub(crate) struct InternalSpawnedWorkerRecord {
pub worker_name: String,
@@ -35,8 +71,13 @@ pub(crate) struct InternalSpawnedWorkerRecord {
#[cfg(test)]
pub installed_tools: Arc<[String]>,
pub session: InternalWorkerSessionHandle,
change_tracker: Option<tools::Tracker>,
started_at: Instant,
stop_lock: Arc<tokio::sync::Mutex<()>>,
scope_reclaimed: Arc<AtomicBool>,
protocol_revision: Arc<AtomicU64>,
protocol_emit_lock: Arc<Mutex<()>>,
protocol_terminal: Arc<AtomicBool>,
forwarding_started: Arc<AtomicBool>,
}
@@ -47,6 +88,7 @@ impl InternalSpawnedWorkerRecord {
workdir_delegation: WorkdirDelegation,
#[cfg(test)] installed_tools: Vec<String>,
session: InternalWorkerSessionHandle,
change_tracker: Option<tools::Tracker>,
) -> Self {
Self {
worker_name,
@@ -55,12 +97,64 @@ impl InternalSpawnedWorkerRecord {
#[cfg(test)]
installed_tools: installed_tools.into(),
session,
change_tracker,
started_at: Instant::now(),
stop_lock: Arc::new(tokio::sync::Mutex::new(())),
scope_reclaimed: Arc::new(AtomicBool::new(false)),
protocol_revision: Arc::new(AtomicU64::new(0)),
protocol_emit_lock: Arc::new(Mutex::new(())),
protocol_terminal: Arc::new(AtomicBool::new(false)),
forwarding_started: Arc::new(AtomicBool::new(false)),
}
}
fn stop_summary(&self) -> SubWorkerStopSummary {
let mut counts = BTreeMap::<String, u64>::new();
for entry in self.session.entries() {
if let session_store::LogEntry::AssistantItem {
item: LoggedItem::ToolCall { name, .. },
..
} = entry
{
let count = counts.entry(bounded_tool_name(&name)).or_default();
*count = count.saturating_add(1);
}
}
let mut tool_counts = counts
.into_iter()
.map(|(name, count)| SubWorkerToolCount { name, count })
.collect::<Vec<_>>();
tool_counts.sort_by(|left, right| {
right
.count
.cmp(&left.count)
.then_with(|| left.name.cmp(&right.name))
});
tool_counts.truncate(STOP_SUMMARY_TOOL_LIMIT);
let change_stat = self.change_tracker.as_ref().and_then(|tracker| {
let stat = tracker.change_stat();
(stat.added > 0 || stat.deleted > 0).then(|| SubWorkerChangeStat {
added: stat.added,
deleted: stat.deleted,
source: "tracked_write_edit_tools".to_string(),
})
});
SubWorkerStopSummary {
session_id: self.session.session_id_string(),
display_name: self.worker_name.clone(),
outcome: SubWorkerFinalOutcome::Done,
elapsed_ms: self
.started_at
.elapsed()
.as_millis()
.min(u128::from(u64::MAX)) as u64,
tool_counts,
change_stat,
}
}
fn claim_scope_reclaim(&self) -> bool {
!self.scope_reclaimed.swap(true, Ordering::AcqRel)
}
@@ -277,12 +371,20 @@ impl SpawnedWorkerRegistry {
};
let worker = record.protocol_ref(Some(parent_session_id));
let protocol_revision = record.protocol_revision.clone();
let protocol_emit_lock = record.protocol_emit_lock.clone();
let protocol_terminal = record.protocol_terminal.clone();
let mut child_rx = record.session.subscribe_events();
tokio::spawn(async move {
loop {
match child_rx.recv().await {
Ok(event) => {
let shutdown = matches!(event, Event::Shutdown);
let _emit_guard = protocol_emit_lock
.lock()
.unwrap_or_else(|error| error.into_inner());
if protocol_terminal.load(Ordering::Acquire) {
break;
}
let revision = protocol_revision.fetch_add(1, Ordering::AcqRel) + 1;
let _ = parent_tx.send(Event::InternalWorker {
worker: worker.clone(),
@@ -294,6 +396,12 @@ impl SpawnedWorkerRegistry {
}
}
Err(broadcast::error::RecvError::Lagged(skipped)) => {
let _emit_guard = protocol_emit_lock
.lock()
.unwrap_or_else(|error| error.into_inner());
if protocol_terminal.load(Ordering::Acquire) {
break;
}
let revision = protocol_revision.fetch_add(1, Ordering::AcqRel) + 1;
let _ = parent_tx.send(Event::InternalWorker {
worker: worker.clone(),
@@ -385,13 +493,34 @@ impl SpawnedWorkerRegistry {
result
}
/// Stop one direct Internal SubWorker and discard its registry/scope state.
///
/// The child actor must acknowledge its stop before the registry is removed.
/// After scope reclamation and removal, `InternalWorkerRemoved` is published
/// exactly once as the parent-stream terminal fence. Callers only receive
/// `Done` after all authoritative cleanup succeeds.
pub(crate) async fn remove_internal(
&self,
worker_name: &str,
) -> io::Result<Option<InternalSpawnedWorkerRecord>> {
if let Some(record) = self.get_internal(worker_name) {
self.reclaim_record_scope(&record)?;
) -> io::Result<Option<SubWorkerStopSummary>> {
let Some(record) = self.get_internal(worker_name) else {
return Ok(None);
};
let _stop_guard = record.stop_lock.lock().await;
let still_registered = self.get_internal(worker_name).is_some_and(|current| {
current.session.session_id_string() == record.session.session_id_string()
});
if !still_registered {
return Ok(None);
}
record
.session
.stop()
.await
.map_err(|error| io::Error::other(error.to_string()))?;
let summary = record.stop_summary();
self.reclaim_record_scope(&record)?;
let removed =
{
let mut records = self.internal_records.lock().map_err(|_| {
@@ -402,14 +531,41 @@ impl SpawnedWorkerRegistry {
})?;
let removed = records
.iter()
.position(|record| record.worker_name == worker_name)
.position(|candidate| {
candidate.worker_name == worker_name
&& candidate.session.session_id_string()
== record.session.session_id_string()
})
.map(|index| records.remove(index));
if removed.is_some() {
names.remove(worker_name);
}
removed
};
Ok(removed)
if removed.is_some() {
self.publish_internal_removal(&record);
}
Ok(removed.map(|_| summary))
}
fn publish_internal_removal(&self, record: &InternalSpawnedWorkerRecord) {
if record.session.visibility() != InternalWorkerVisibility::ParentClient {
return;
}
let Some((parent_tx, parent_session_id)) = self.parent_protocol.lock().unwrap().clone()
else {
return;
};
let _emit_guard = record
.protocol_emit_lock
.lock()
.unwrap_or_else(|error| error.into_inner());
record.protocol_terminal.store(true, Ordering::Release);
let revision = record.protocol_revision.fetch_add(1, Ordering::AcqRel) + 1;
let _ = parent_tx.send(Event::InternalWorkerRemoved {
worker: record.protocol_ref(Some(parent_session_id)),
revision,
});
}
}
@@ -508,6 +664,17 @@ fn record_from_worker_state(child: &WorkerSpawnedChild) -> io::Result<SpawnedWor
})
}
fn bounded_tool_name(name: &str) -> String {
let mut bounded = name
.chars()
.take(STOP_SUMMARY_TOOL_NAME_LIMIT)
.collect::<String>();
if name.chars().count() > STOP_SUMMARY_TOOL_NAME_LIMIT {
bounded.push('…');
}
bounded
}
fn store_error_to_io(error: WorkerStoreError) -> io::Error {
io::Error::other(error)
}
@@ -520,7 +687,7 @@ mod tests {
use session_store::LogEntry;
use super::*;
use crate::internal_worker::test_internal_worker_session;
use crate::internal_worker::{InternalWorkerSessionStatus, test_internal_worker_session};
fn registry() -> Arc<SpawnedWorkerRegistry> {
let scope = Scope::from_config(&ScopeConfig {
@@ -577,6 +744,7 @@ mod tests {
delegation,
Vec::new(),
session,
None,
),
sender,
)
@@ -669,4 +837,124 @@ mod tests {
);
assert!(registry.internal_worker_snapshots().is_empty());
}
fn install_record(registry: &SpawnedWorkerRegistry, record: InternalSpawnedWorkerRecord) {
registry
.internal_names
.lock()
.unwrap()
.insert(record.worker_name.clone());
registry.internal_records.lock().unwrap().push(record);
}
#[tokio::test]
async fn stop_removes_internal_worker_and_returns_bounded_summary() {
let registry = registry();
let (parent_tx, mut parent_rx) = broadcast::channel(32);
registry.attach_parent_protocol(parent_tx, "parent-session".into());
let tracker = tools::Tracker::new();
tracker.record_change(12, 4);
let (mut record, _events) = record("child", InternalWorkerVisibility::ParentClient).await;
record.change_tracker = Some(tracker);
for (index, name) in ["Read", "Read", "Grep"].into_iter().enumerate() {
record.session.publish_test_entry(LogEntry::AssistantItem {
ts: index as u64,
item: LoggedItem::ToolCall {
call_id: format!("call-{index}"),
name: name.to_string(),
arguments: "{}".to_string(),
},
});
}
registry.start_protocol_forwarding(record.clone());
install_record(&registry, record);
let summary = registry.remove_internal("child").await.unwrap().unwrap();
assert_eq!(summary.display_name, "child");
assert_eq!(summary.outcome, SubWorkerFinalOutcome::Done);
assert_eq!(
summary.tool_counts,
vec![
SubWorkerToolCount {
name: "Read".to_string(),
count: 2,
},
SubWorkerToolCount {
name: "Grep".to_string(),
count: 1,
},
]
);
assert_eq!(
summary.change_stat,
Some(SubWorkerChangeStat {
added: 12,
deleted: 4,
source: "tracked_write_edit_tools".to_string(),
})
);
assert!(registry.get_internal("child").is_none());
let terminal_revision = loop {
if let Event::InternalWorkerRemoved { worker, revision } =
parent_rx.recv().await.unwrap()
{
assert_eq!(worker.session_id, summary.session_id);
assert!(revision > 0);
break revision;
}
};
assert!(registry.remove_internal("child").await.unwrap().is_none());
while let Ok(Ok(event)) =
tokio::time::timeout(Duration::from_millis(20), parent_rx.recv()).await
{
assert!(!matches!(event, Event::InternalWorkerRemoved { .. }));
if let Event::InternalWorker { revision, .. } = event {
assert!(revision > terminal_revision);
}
}
}
#[tokio::test]
async fn running_worker_is_stopped_before_removal() {
let registry = registry();
let (record, _events) = record("running", InternalWorkerVisibility::ParentClient).await;
record
.session
.force_status(InternalWorkerSessionStatus::Running);
install_record(&registry, record);
let summary = registry.remove_internal("running").await.unwrap().unwrap();
assert_eq!(summary.outcome, SubWorkerFinalOutcome::Done);
assert!(registry.get_internal("running").is_none());
}
#[tokio::test]
async fn stop_failure_keeps_registry_and_emits_no_removal() {
let registry = registry();
let (parent_tx, mut parent_rx) = broadcast::channel(8);
registry.attach_parent_protocol(parent_tx, "parent-session".into());
let (record, _events) = record("child", InternalWorkerVisibility::ParentClient).await;
record.session.force_stop_failure();
install_record(&registry, record);
let error = registry.remove_internal("child").await.unwrap_err();
assert!(error.to_string().contains("unavailable"));
assert!(registry.get_internal("child").is_some());
assert!(matches!(
parent_rx.try_recv(),
Err(broadcast::error::TryRecvError::Empty)
));
}
#[tokio::test]
async fn read_only_summary_omits_unavailable_change_stat() {
let tracker = tools::Tracker::new();
let (mut record, _events) = record("reader", InternalWorkerVisibility::ParentClient).await;
record.change_tracker = Some(tracker);
assert_eq!(record.stop_summary().change_stat, None);
}
}
+2
View File
@@ -481,6 +481,7 @@ impl Tool for SubWorkerSpawnTool {
.map_err(|error| {
ToolError::ExecutionFailed(format!("install Internal Worker features: {error}"))
})?;
let child_change_tracker = child.tracker().cloned();
#[cfg(test)]
let installed_tools = child
.engine()
@@ -587,6 +588,7 @@ impl Tool for SubWorkerSpawnTool {
#[cfg(test)]
installed_tools,
session.clone(),
child_change_tracker,
);
if let Err(error) = name_reservation.commit(record) {
let _ = session.stop().await;
+1
View File
@@ -35,6 +35,7 @@ ticket.workspace = true
memory.workspace = true
merge-request.workspace = true
tokio = { workspace = true, features = ["fs", "macros", "net", "rt-multi-thread", "sync", "time"] }
tower.workspace = true
tokio-tungstenite.workspace = true
worker.workspace = true
workdir = { workspace = true, features = ["http-client"] }
+41 -6
View File
@@ -2843,12 +2843,6 @@ mod tests {
write_ticket(dir.path(), "00000000001J5", "Second ticket", "planning");
write_ticket(dir.path(), "00000000001J6", "Third ticket", "planning");
let db_path = dir.path().join("workspace.db");
SqliteTicketBackend::open(&db_path, "workspace-test")
.unwrap()
.import_from_local_backend(&ticket::LocalTicketBackend::new(
dir.path().join(".yoi/tickets"),
))
.unwrap();
let store = SqliteWorkspaceStore::open(&db_path).unwrap();
store
.upsert_workspace(&WorkspaceRecord {
@@ -2861,6 +2855,27 @@ mod tests {
})
.await
.unwrap();
SqliteTicketBackend::open(&db_path, "workspace-test")
.unwrap()
.import_from_local_backend(&ticket::LocalTicketBackend::new(
dir.path().join(".yoi/tickets"),
))
.unwrap();
rusqlite::Connection::open(&db_path)
.unwrap()
.execute_batch(
r#"
INSERT INTO workspace_resource_human_keys (
workspace_id, resource_kind, resource_id, sequence, human_key, allocated_at
) VALUES
('workspace-test', 'ticket', '00000000001J2', 1, 'T-1', '2026-01-01T00:00:00Z'),
('workspace-test', 'ticket', '00000000001J5', 2, 'T-2', '2026-01-01T00:00:00Z'),
('workspace-test', 'ticket', '00000000001J6', 3, 'T-3', '2026-01-01T00:00:00Z');
INSERT INTO workspace_resource_human_key_counters (workspace_id, resource_kind, next_sequence)
VALUES ('workspace-test', 'ticket', 4);
"#,
)
.unwrap();
store
.upsert_objective(&ObjectiveRecord {
workspace_id: "workspace-test".to_string(),
@@ -3216,6 +3231,26 @@ mod tests {
})
.await
.unwrap();
rusqlite::Connection::open(&db_path)
.unwrap()
.execute_batch(
r#"
INSERT INTO typed_tickets (
workspace_id, ticket_id, slug, title, status, kind, priority, body,
workflow_state, workflow_state_explicit
) VALUES
('workspace-test', '00000000001J2', 'ticket-j2', 'Ticket J2', 'open', 'task', 'normal', '', 'planning', 1),
('workspace-test', '00000000001J3', 'ticket-j3', 'Ticket J3', 'open', 'task', 'normal', '', 'planning', 1);
INSERT INTO workspace_resource_human_keys (
workspace_id, resource_kind, resource_id, sequence, human_key, allocated_at
) VALUES
('workspace-test', 'ticket', '00000000001J2', 1, 'T-1', '2026-01-01T00:00:00Z'),
('workspace-test', 'ticket', '00000000001J3', 2, 'T-2', '2026-01-01T00:00:00Z');
INSERT INTO workspace_resource_human_key_counters (workspace_id, resource_kind, next_sequence)
VALUES ('workspace-test', 'ticket', 3);
"#,
)
.unwrap();
let authority = SqliteWorkspaceAuthority::new(&db_path, "workspace-test").unwrap();
let created = authority
+8
View File
@@ -2447,6 +2447,8 @@ impl WorkspaceWorkerRuntime for EmbeddedWorkerRuntime {
#[derive(Clone)]
pub struct RemoteRuntimeConfig {
pub runtime_id: String,
/// Explicit Workspace assignment granted by Server authority.
pub workspace_id: Option<String>,
pub display_name: String,
pub base_url: String,
pub bearer_token: Option<String>,
@@ -2489,6 +2491,7 @@ impl RemoteRuntimeConfig {
) -> Self {
Self {
runtime_id: runtime_id.into(),
workspace_id: None,
display_name: display_name.into(),
base_url: base_url.into(),
bearer_token,
@@ -2501,6 +2504,11 @@ impl RemoteRuntimeConfig {
}
}
pub fn with_workspace_id(mut self, workspace_id: impl Into<String>) -> Self {
self.workspace_id = Some(workspace_id.into());
self
}
pub fn with_cached_capabilities(mut self, capabilities: RuntimeCapabilitySummary) -> Self {
self.cached_capabilities = capabilities;
self
+9 -1
View File
@@ -27,6 +27,7 @@ pub mod server;
pub mod skills;
pub mod store;
pub mod worker_source;
pub mod workspace_catalog;
mod workspace_subscription;
pub use authority::{
@@ -45,8 +46,15 @@ pub use repositories::{
ConfiguredRepository, GitCommitSummary, GitRemoteSummary, GitRepositorySummary,
RepositoryLogRead, RepositoryRegistryReader, RepositorySummary,
};
pub use server::{AuthConfig, ServerConfig, WorkspaceApi, build_router, serve};
pub use server::{
AuthConfig, ServerConfig, WorkspaceApi, WorkspaceServerApi, build_router,
build_workspace_server_router, serve, serve_workspace_catalog,
};
pub use store::{ControlPlaneStore, SqliteWorkspaceStore, WorkspaceRecord};
pub use workspace_catalog::{
InitialRepositoryIntent, WorkspaceCatalogService, WorkspaceCreateRequest,
WorkspaceCreateResponse,
};
use worker_runtime::identity::RuntimeWorkerRef;
+81 -54
View File
@@ -9,10 +9,11 @@ use serde::{Deserialize, Serialize};
use tokio::net::TcpListener;
use worker_runtime::auth::{RuntimeIdentityMaterial, decode_public_key};
use yoi_workspace_server::hosts::{RemoteRuntimeAuthConfig, RemoteRuntimeConfig};
use yoi_workspace_server::store::{RepositoryRecord, SqliteWorkspaceStore, TrustedRuntimeRecord};
use yoi_workspace_server::store::{SqliteWorkspaceStore, TrustedRuntimeRecord};
use yoi_workspace_server::{
BackendRuntimesConfigFile, ControlPlaneStore, ServerConfig, WORKSPACE_BACKEND_CONFIG_TEMPLATE,
WorkspaceBackendConfigFile, WorkspaceIdentity, WorkspaceRecord, serve,
BackendRuntimesConfigFile, ControlPlaneStore, InitialRepositoryIntent, ServerConfig,
WORKSPACE_BACKEND_CONFIG_TEMPLATE, WorkspaceBackendConfigFile, WorkspaceCatalogService,
WorkspaceCreateRequest, WorkspaceIdentity, WorkspaceRecord, serve_workspace_catalog,
};
#[derive(Debug)]
@@ -152,30 +153,21 @@ async fn run_init_with_database_path(
if let Some(parent) = database_path.parent() {
tokio::fs::create_dir_all(parent).await?;
}
let store = SqliteWorkspaceStore::open(&database_path)?;
store
.upsert_workspace(&WorkspaceRecord {
workspace_id: identity.workspace_id.clone(),
owner_account_id: None,
let store = Arc::new(SqliteWorkspaceStore::open(&database_path)?);
let service = WorkspaceCatalogService::new(store);
service.create_with_workspace_id(
WorkspaceCreateRequest {
operation_key: format!("cli-init:{}", identity.workspace_id),
display_name: identity.display_name.clone(),
state: "active".to_string(),
created_at: identity.created_at.clone(),
updated_at: identity.created_at.clone(),
})
.await?;
store.upsert_repository(&RepositoryRecord {
workspace_id: identity.workspace_id.clone(),
repository_id: "main".to_string(),
name: "Main repository".to_string(),
kind: "git".to_string(),
provider: Some("git".to_string()),
uri: options.workspace.display().to_string(),
default_ref: Some("HEAD".to_string()),
auth_ref_kind: None,
auth_ref_key: None,
created_at: identity.created_at.clone(),
updated_at: identity.created_at.clone(),
})?;
repository: InitialRepositoryIntent {
uri: options.workspace.display().to_string(),
display_name: Some("Main repository".to_string()),
default_ref: Some("HEAD".to_string()),
},
},
None,
Some(identity.workspace_id.clone()),
)?;
eprintln!(
"yoi-server: initialized workspace `{}` ({}) in server DB `{}`",
@@ -358,6 +350,7 @@ fn run_trust_runtime_command(args: Vec<String>) -> Result<(), Box<dyn std::error
match subcommand.as_str() {
"add" => {
let mut runtime_id = None;
let mut workspace_id = None;
let mut base_url = None;
let mut public_key = None;
let mut display_name = None;
@@ -368,6 +361,9 @@ fn run_trust_runtime_command(args: Vec<String>) -> Result<(), Box<dyn std::error
"--runtime-id" => {
runtime_id = Some(take_value(&flag, inline_value, &mut args)?)
}
"--workspace-id" => {
workspace_id = Some(take_value(&flag, inline_value, &mut args)?)
}
"--base-url" | "--endpoint" => {
base_url = Some(take_value(&flag, inline_value, &mut args)?)
}
@@ -390,15 +386,39 @@ fn run_trust_runtime_command(args: Vec<String>) -> Result<(), Box<dyn std::error
}
let runtime_id = runtime_id
.ok_or_else(|| CliError("trust-runtime add requires --runtime-id".to_string()))?;
let workspace_id = workspace_id
.ok_or_else(|| CliError("trust-runtime add requires --workspace-id".to_string()))?;
if !store
.list_workspaces()?
.iter()
.any(|workspace| workspace.workspace_id == workspace_id)
{
return Err(Box::new(CliError(format!(
"Workspace `{workspace_id}` is not registered"
))));
}
let base_url = base_url
.ok_or_else(|| CliError("trust-runtime add requires --base-url".to_string()))?;
let public_key = public_key
.ok_or_else(|| CliError("trust-runtime add requires --public-key".to_string()))?;
decode_public_key(&public_key)?;
ensure_trusted_runtime_replace_allowed(&store, &runtime_id, replace)?;
if let Some(existing) = store
.list_trusted_runtimes(true)?
.into_iter()
.find(|runtime| runtime.runtime_id == runtime_id)
{
if existing.workspace_id.as_deref() != Some(workspace_id.as_str()) {
return Err(Box::new(CliError(format!(
"runtime `{runtime_id}` is already assigned to Workspace `{}` and cannot be reparented",
existing.workspace_id.as_deref().unwrap_or("unassigned")
))));
}
}
let now = Utc::now().to_rfc3339();
store.upsert_trusted_runtime(&TrustedRuntimeRecord {
runtime_id: runtime_id.clone(),
workspace_id: Some(workspace_id.clone()),
display_name: display_name.unwrap_or_else(|| runtime_id.clone()),
base_url,
public_key,
@@ -437,8 +457,9 @@ fn run_trust_runtime_command(args: Vec<String>) -> Result<(), Box<dyn std::error
} else {
for runtime in records {
println!(
"runtime_id={} base_url={} public_key={} revoked_at={}",
"runtime_id={} workspace_id={} base_url={} public_key={} revoked_at={}",
runtime.runtime_id,
runtime.workspace_id.unwrap_or_default(),
runtime.base_url,
runtime.public_key,
runtime.revoked_at.unwrap_or_default()
@@ -582,12 +603,28 @@ async fn run_serve(options: ServeOptions) -> Result<(), Box<dyn std::error::Erro
}
let store = Arc::new(SqliteWorkspaceStore::open(&database_path)?);
let workspace = select_serve_workspace(store.as_ref())?;
let workspace_root = infer_workspace_root_from_repositories(store.as_ref(), &workspace)?;
let identity = WorkspaceIdentity {
workspace_id: workspace.workspace_id.clone(),
created_at: workspace.created_at.clone(),
display_name: workspace.display_name.clone(),
let workspaces = store.list_workspaces()?;
let (identity, workspace_root) = if let Some(workspace) = workspaces.first() {
(
WorkspaceIdentity {
workspace_id: workspace.workspace_id.clone(),
created_at: workspace.created_at.clone(),
display_name: workspace.display_name.clone(),
},
infer_workspace_root_from_repositories(store.as_ref(), workspace)?,
)
} else {
(
WorkspaceIdentity {
workspace_id: "00000000-0000-0000-0000-000000000000".to_string(),
created_at: Utc::now().to_rfc3339(),
display_name: "Server bootstrap".to_string(),
},
database_path
.parent()
.ok_or_else(|| CliError("server database path has no parent".to_string()))?
.to_path_buf(),
)
};
let runtime_config = BackendRuntimesConfigFile::load_default()?;
let mut resolved = WorkspaceBackendConfigFile::default().resolve_with_runtime_config(
@@ -601,6 +638,7 @@ async fn run_serve(options: ServeOptions) -> Result<(), Box<dyn std::error::Erro
if let Some(listen) = options.listen {
resolved = resolved.with_listen(listen);
}
resolved.server.allow_local_workspace_bootstrap = resolved.listen.ip().is_loopback();
let listener = TcpListener::bind(resolved.listen).await?;
let local_addr = listener.local_addr()?;
@@ -608,12 +646,12 @@ async fn run_serve(options: ServeOptions) -> Result<(), Box<dyn std::error::Erro
resolved = resolved.with_backend_base_url(format!("http://{local_addr}"));
}
eprintln!(
"yoi-server: serving workspace `{}` from server DB `{}` on http://{}",
workspace.workspace_id,
"yoi-server: serving {} workspace(s) from server DB `{}` on http://{}",
workspaces.len(),
database_path.display(),
local_addr
);
serve(resolved.server, store, listener).await?;
serve_workspace_catalog(resolved.server, store, listener).await?;
Ok(())
}
@@ -630,6 +668,9 @@ fn append_trusted_runtime_sources(
return Ok(());
};
for runtime in store.list_trusted_runtimes(false)? {
let Some(workspace_id) = runtime.workspace_id.clone() else {
continue;
};
let auth = RemoteRuntimeAuthConfig {
server_id: server_identity.identity.identity_id.clone(),
server_private_key: server_identity.identity.private_key.clone(),
@@ -640,6 +681,7 @@ fn append_trusted_runtime_sources(
runtime.base_url,
None,
)
.with_workspace_id(workspace_id)
.with_auth(auth);
remote_runtime_sources.retain(|existing| existing.runtime_id != runtime.runtime_id);
remote_runtime_sources.push(remote);
@@ -647,23 +689,6 @@ fn append_trusted_runtime_sources(
Ok(())
}
fn select_serve_workspace(store: &SqliteWorkspaceStore) -> Result<WorkspaceRecord, CliError> {
let workspaces = store
.list_workspaces()
.map_err(|error| CliError(format!("failed to list workspaces from server DB: {error}")))?;
match workspaces.as_slice() {
[] => Err(CliError(
"server DB has no workspace records; run `yoi-server init --workspace <PATH>`"
.to_string(),
)),
[workspace] => Ok(workspace.clone()),
_ => Err(CliError(format!(
"server DB contains {} workspaces; serve workspace selection is not implemented yet",
workspaces.len()
))),
}
}
fn infer_workspace_root_from_repositories(
store: &SqliteWorkspaceStore,
workspace: &WorkspaceRecord,
@@ -914,7 +939,7 @@ fn parse_listen(value: &str) -> Result<SocketAddr, CliError> {
fn print_help() {
println!(
"yoi-server\n\nUsage:\n yoi-server init [OPTIONS]\n yoi-server config <COMMAND> [OPTIONS]\n yoi-server identity init --server-id <SERVER_ID> [--replace]\n yoi-server identity show [--json]\n yoi-server trust-runtime add --runtime-id <RUNTIME_ID> --base-url <URL> --public-key <KEY> [--display-name <NAME>] [--replace]\n yoi-server trust-runtime list [--json] [--include-revoked]\n yoi-server trust-runtime revoke --runtime-id <RUNTIME_ID>\n yoi-server skills <COMMAND> [OPTIONS]\n yoi-server migrate --dry-run [--database <PATH>]
"yoi-server\n\nUsage:\n yoi-server init [OPTIONS]\n yoi-server config <COMMAND> [OPTIONS]\n yoi-server identity init --server-id <SERVER_ID> [--replace]\n yoi-server identity show [--json]\n yoi-server trust-runtime add --runtime-id <RUNTIME_ID> --workspace-id <WORKSPACE_ID> --base-url <URL> --public-key <KEY> [--display-name <NAME>] [--replace]\n yoi-server trust-runtime list [--json] [--include-revoked]\n yoi-server trust-runtime revoke --runtime-id <RUNTIME_ID>\n yoi-server skills <COMMAND> [OPTIONS]\n yoi-server migrate --dry-run [--database <PATH>]
yoi-server serve [OPTIONS]\n\nOptions:\n -h, --help Print help"
);
}
@@ -1038,6 +1063,7 @@ mod tests {
store
.upsert_trusted_runtime(&TrustedRuntimeRecord {
runtime_id: "runtime-a".to_string(),
workspace_id: None,
display_name: "Runtime A".to_string(),
base_url: "http://127.0.0.1:18080".to_string(),
public_key,
@@ -1059,6 +1085,7 @@ mod tests {
async fn init_creates_identity_local_config_and_server_records() {
let temp = tempfile::tempdir().unwrap();
let database_path = temp.path().join("data").join("server").join("server.db");
std::fs::create_dir(temp.path().join(".git")).unwrap();
run_init_with_database_path(
InitOptions {
workspace: temp.path().canonicalize().unwrap(),
+9
View File
@@ -1082,6 +1082,11 @@ mod tests {
) VALUES('w',?1,'r','one','builtin:coder','normal','created','rev1')",
[worker_id().to_string()],
)?;
c.execute(
"INSERT INTO typed_tickets (workspace_id, ticket_id, slug, title, status, kind, priority, body, workflow_state, workflow_state_explicit) \
VALUES ('w', 'ticket', 'ticket', 'Ticket', 'open', 'task', 'normal', '', 'planning', 1)",
[],
)?;
Ok(())
})
.unwrap();
@@ -1267,7 +1272,11 @@ mod tests {
fn purge_tombstone_commit_is_idempotent() {
let s = setup();
s.with_conn(|conn| {
conn.execute("INSERT INTO typed_tickets(workspace_id,ticket_id,slug,title,status,kind,priority,body,workflow_state,workflow_state_explicit) VALUES('w','ticket-old','ticket-old','Old Ticket','open','task','normal','','planning',1)", [])?;
conn.execute("INSERT INTO worker_registry(workspace_id,worker_id,runtime_id,display_name,profile,retention_state,created_at,updated_at) VALUES('w','1','r','old worker','builtin:coder','normal','created','rev1')", [])?;
conn.execute("INSERT INTO ticket_worker_assignments(workspace_id,ticket_id,assignment_id,runtime_id,worker_id,assigned_by,assigned_at) VALUES('w','ticket-old','assignment-old','r','1','test','t')", [])?;
conn.execute("DELETE FROM worker_registry WHERE workspace_id='w' AND runtime_id='r' AND worker_id='1'", [])?;
conn.execute("DELETE FROM typed_tickets WHERE workspace_id='w' AND ticket_id='ticket-old'", [])?;
Ok(())
}).unwrap();
let p = s.plan_worker_removal(&req(), &inv()).unwrap();
+584 -70
View File
@@ -32,9 +32,11 @@ use ticket::{
execute_ticket_backend_operation,
};
use tokio::net::TcpListener;
use tokio::sync::Mutex as AsyncMutex;
use tokio_tungstenite::connect_async;
use tokio_tungstenite::tungstenite::Message as TungsteniteMessage;
use tokio_tungstenite::tungstenite::client::IntoClientRequest;
use tower::ServiceExt;
use url::Url;
use uuid::Uuid;
use webauthn_rs::prelude::{
@@ -111,6 +113,7 @@ use crate::store::{
TicketWorkerAssignmentRecord, UserRecord, WorkdirRegistryRecord, WorkerControlGrantRecord,
WorkerRegistryRecord, WorkerWorkdirLinkRecord, WorkspaceRecord, WorkspaceResourceKind,
};
use crate::workspace_catalog::{WorkspaceCatalogService, WorkspaceCreateRequest};
use crate::{Error, Result};
use worker_runtime::catalog::{
ConfigBundleRef, ProfileSelector, RepositorySelector as RuntimeRepositorySelector,
@@ -158,6 +161,9 @@ pub struct ServerConfig {
pub remote_runtime_sources: Vec<RemoteRuntimeConfig>,
pub runtime_config_path: Option<PathBuf>,
pub backend_base_url: Option<String>,
/// Allows the first ownerless Workspace to be created without a session.
/// This must only be enabled for a loopback-bound local Server.
pub allow_local_workspace_bootstrap: bool,
}
impl ServerConfig {
@@ -187,6 +193,7 @@ impl ServerConfig {
remote_runtime_sources: Vec::new(),
runtime_config_path: BackendRuntimesConfigFile::default_path(),
backend_base_url: None,
allow_local_workspace_bootstrap: false,
}
}
@@ -243,10 +250,69 @@ impl ServerConfig {
Self::default_workspace_backend_data_root(workspace_id).join("embedded-runtime")
}
pub fn with_local_workspace_bootstrap(mut self, enabled: bool) -> Self {
self.allow_local_workspace_bootstrap = enabled;
self
}
pub fn with_embedded_runtime_store_root(mut self, root: impl Into<PathBuf>) -> Self {
self.embedded_runtime_store_root = root.into();
self
}
fn for_catalog_workspace(
&self,
workspace: &WorkspaceRecord,
repositories: Vec<RepositoryRecord>,
) -> Result<Self> {
let primary = repositories
.iter()
.find(|repository| repository.repository_id == "main")
.or_else(|| repositories.first())
.ok_or_else(|| {
Error::Config(format!(
"Workspace {} has no registered repository",
workspace.workspace_id
))
})?;
let workspace_root = PathBuf::from(&primary.uri);
if !workspace_root.is_absolute() {
return Err(Error::Config(format!(
"Workspace {} repository uri is not an absolute local path",
workspace.workspace_id
)));
}
let repositories = repositories
.into_iter()
.map(|repository| ConfiguredRepository {
id: repository.repository_id,
provider: repository.provider.unwrap_or(repository.kind),
path: PathBuf::from(&repository.uri),
uri: repository.uri,
display_name: Some(repository.name),
default_selector: repository.default_ref,
})
.collect();
let mut scoped = self.clone();
scoped.workspace_id.clone_from(&workspace.workspace_id);
scoped
.workspace_display_name
.clone_from(&workspace.display_name);
scoped
.workspace_created_at
.clone_from(&workspace.created_at);
scoped.workspace_root = workspace_root;
scoped.embedded_runtime_store_root =
Self::default_embedded_runtime_store_root(&workspace.workspace_id);
scoped.repositories = repositories;
// Runtime trust is server-global. Only explicitly assigned sources enter
// this Workspace's registry and receive Workspace-scoped capabilities.
scoped.remote_runtime_sources.retain(|runtime| {
runtime.workspace_id.as_deref() == Some(workspace.workspace_id.as_str())
});
scoped.runtime_event_sources.clear();
Ok(scoped)
}
}
const ORCHESTRATOR_ATTENTION_TICKET_LIMIT: usize = 20;
@@ -682,6 +748,213 @@ impl crate::worker_source::VerifiedWorkerRemoveExecutor for WorkspaceWorkerRemov
}
}
#[derive(Clone)]
pub struct WorkspaceServerApi {
template: Arc<ServerConfig>,
store: Arc<dyn ControlPlaneStore>,
catalog: WorkspaceCatalogService,
routers: Arc<AsyncMutex<HashMap<String, Router>>>,
}
impl WorkspaceServerApi {
pub fn new(template: ServerConfig, store: Arc<dyn ControlPlaneStore>) -> Self {
Self {
template: Arc::new(template),
catalog: WorkspaceCatalogService::new(store.clone()),
store,
routers: Arc::new(AsyncMutex::new(HashMap::new())),
}
}
async fn router_for_workspace(&self, workspace_id: &str) -> Result<Option<Router>> {
let mut routers = self.routers.lock().await;
if let Some(router) = routers.get(workspace_id) {
return Ok(Some(router.clone()));
}
let Some(workspace) = self.store.get_workspace(workspace_id).await? else {
return Ok(None);
};
let repositories = self.store.list_repositories(workspace_id)?;
let config = self
.template
.for_catalog_workspace(&workspace, repositories)?;
let api = WorkspaceApi::new(config, self.store.clone()).await?;
tokio::spawn(run_orchestrator_turn_end_hook(api.clone()));
let router = build_router(api);
routers.insert(workspace_id.to_string(), router.clone());
Ok(Some(router))
}
async fn preload(&self) -> Result<()> {
for workspace in self.store.list_workspaces()? {
let _ = self
.router_for_workspace(&workspace.workspace_id)
.await?
.ok_or_else(|| {
Error::Config(format!(
"Workspace {} disappeared while loading",
workspace.workspace_id
))
})?;
}
Ok(())
}
}
fn server_error_response(error: Error) -> Response {
ApiError::from(error).into_response()
}
fn forbidden_server_response(message: &str) -> Response {
(
StatusCode::FORBIDDEN,
Json(serde_json::json!({ "error": message })),
)
.into_response()
}
#[derive(Debug, Deserialize)]
struct WorkspaceListQuery {
limit: Option<usize>,
}
async fn list_server_workspaces(
State(api): State<WorkspaceServerApi>,
headers: HeaderMap,
Query(query): Query<WorkspaceListQuery>,
) -> Response {
let owner = match resolve_server_actor(&api, &headers).await {
Ok(Some(actor)) => Some(actor.account_id),
Ok(None) => None,
Err(error) => return server_error_response(error),
};
match api
.catalog
.list(owner.as_deref(), query.limit.unwrap_or(100))
{
Ok(workspaces) => Json(workspaces).into_response(),
Err(error) => server_error_response(error),
}
}
async fn create_server_workspace(
State(api): State<WorkspaceServerApi>,
headers: HeaderMap,
Json(request): Json<WorkspaceCreateRequest>,
) -> Response {
let (owner_account_id, local_bootstrap) = match resolve_server_actor(&api, &headers).await {
Ok(Some(actor)) => (Some(actor.account_id), false),
Ok(None) if api.template.allow_local_workspace_bootstrap => (None, true),
Ok(None) => {
return forbidden_server_response("Workspace creation requires an authenticated owner");
}
Err(error) => return server_error_response(error),
};
let created = match if local_bootstrap {
api.catalog.create_first_ownerless(request)
} else {
api.catalog.create(request, owner_account_id)
} {
Ok(created) => created,
Err(error) => return server_error_response(error),
};
if let Err(error) = api
.router_for_workspace(&created.workspace.workspace_id)
.await
{
return server_error_response(error);
}
let status = if created.replayed {
StatusCode::OK
} else {
StatusCode::CREATED
};
(status, Json(created)).into_response()
}
async fn resolve_server_actor(
api: &WorkspaceServerApi,
headers: &HeaderMap,
) -> std::result::Result<Option<RequestActor>, Error> {
let cookie_name = auth_public_config(api.template.as_ref()).cookie_name;
resolve_request_actor(api.store.as_ref(), headers, &cookie_name).await
}
async fn dispatch_workspace_request(
State(api): State<WorkspaceServerApi>,
request: Request,
) -> Response {
let path = request.uri().path();
let workspace_id = scoped_workspace_id(path);
let router = if let Some(workspace_id) = workspace_id {
match api.router_for_workspace(workspace_id).await {
Ok(Some(router)) => Some(router),
Ok(None) => None,
Err(error) => return server_error_response(error),
}
} else {
let workspaces = match api.store.list_workspaces() {
Ok(workspaces) => workspaces,
Err(error) => return server_error_response(error),
};
if workspaces.len() == 1 || is_server_global_forward(path) {
match workspaces.first() {
Some(workspace) => match api.router_for_workspace(&workspace.workspace_id).await {
Ok(router) => router,
Err(error) => return server_error_response(error),
},
None => None,
}
} else {
None
}
};
let Some(router) = router else {
return StatusCode::NOT_FOUND.into_response();
};
match router.oneshot(request).await {
Ok(response) => response,
Err(error) => match error {},
}
}
fn is_server_global_forward(path: &str) -> bool {
path == "/api/auth"
|| path.starts_with("/api/auth/")
|| path == "/health"
|| path == "/"
|| path.starts_with("/assets/")
}
fn scoped_workspace_id(path: &str) -> Option<&str> {
let mut segments = path.trim_start_matches('/').split('/');
match (segments.next(), segments.next(), segments.next()) {
(Some("api"), Some("w"), Some(workspace_id))
| (Some("internal"), Some("w"), Some(workspace_id))
if !workspace_id.is_empty() =>
{
Some(workspace_id)
}
(Some("w"), Some(workspace_id), _) if !workspace_id.is_empty() => Some(workspace_id),
_ => None,
}
}
pub async fn build_workspace_server_router(
template: ServerConfig,
store: Arc<dyn ControlPlaneStore>,
) -> Result<Router> {
let api = WorkspaceServerApi::new(template, store);
api.preload().await?;
Ok(Router::new()
.route(
"/api/workspaces",
get(list_server_workspaces).post(create_server_workspace),
)
.fallback(dispatch_workspace_request)
.with_state(api))
}
impl WorkspaceApi {
pub fn with_config_schema_provider(
mut self,
@@ -1687,8 +1960,8 @@ pub fn build_router(api: WorkspaceApi) -> Router {
post(scoped_test_remote_runtime_connection),
)
.route(
"/internal/runtime/resources/fetch",
post(post_internal_runtime_resource_fetch),
"/internal/w/{workspace_id}/runtime/resources/fetch",
post(scoped_post_internal_runtime_resource_fetch),
)
.route("/api/companion/status", get(get_companion_status))
.route(
@@ -1872,6 +2145,16 @@ struct ApiFailureLogEvent<'a> {
diagnostics: Option<&'a [RuntimeDiagnostic]>,
}
pub async fn serve_workspace_catalog(
template: ServerConfig,
store: Arc<dyn ControlPlaneStore>,
listener: TcpListener,
) -> Result<()> {
let router = build_workspace_server_router(template, store).await?;
axum::serve(listener, router).await?;
Ok(())
}
pub async fn serve(
config: ServerConfig,
store: Arc<dyn ControlPlaneStore>,
@@ -9506,13 +9789,20 @@ fn browser_worker_response_from_summary(
})
}
async fn post_internal_runtime_resource_fetch(
async fn scoped_post_internal_runtime_resource_fetch(
State(api): State<WorkspaceApi>,
AxumPath(workspace_id): AxumPath<String>,
Json(request): Json<BackendResourceFetchRequest>,
) -> std::result::Result<
Json<worker_runtime::resource::BackendResourceFetchResponse>,
(StatusCode, Json<BackendResourceError>),
> {
if workspace_id != api.workspace_id() {
return Err((
StatusCode::NOT_FOUND,
Json(BackendResourceError::MissingResource),
));
}
api.resource_broker
.fetch_profile_source_archive(request)
.map(Json)
@@ -12889,38 +13179,10 @@ mod tests {
.is_err()
);
let ticket = browser_ticket_backend(&api)
.unwrap()
.create(create_input)
.unwrap();
let flow_ticket_launch = WorkerSpawnRequest {
requested_worker_name: Some("cross-workspace-ticket".to_string()),
intent: WorkerSpawnIntent::TicketRole {
ticket_id: ticket.id,
role: TicketWorkerRole::Coder,
},
acceptance: WorkerSpawnAcceptanceRequirement::RunAccepted {
expected_segments: 2,
},
profile: ProfileSelector::Builtin("builtin:coder".to_string()),
ticket_assignment: None,
initial_submit: vec![
Segment::Flow {
selector: "builtin:coder-review".to_string(),
},
Segment::text("Implement the Ticket"),
],
working_directory_request: None,
resolved_working_directory_request: None,
resolved_working_directory: None,
resolved_config_bundle: None,
resolved_worker_observation_enabled: false,
resolved_worker_observation_grants: Vec::new(),
resolved_control_operation: None,
resolved_workspace_api: None,
};
assert!(
api.validate_worker_spawn_repository_scope(&flow_ticket_launch)
browser_ticket_backend(&api)
.unwrap()
.create(create_input)
.is_err()
);
@@ -14158,6 +14420,194 @@ mod tests {
config
}
#[tokio::test]
async fn server_router_dispatches_two_workspace_contexts_without_state_leakage() {
let dir = tempfile::tempdir().unwrap();
let repository_a = dir.path().join("repository-a");
let repository_b = dir.path().join("repository-b");
std::fs::create_dir_all(repository_a.join(".git")).unwrap();
std::fs::create_dir_all(repository_b.join(".git")).unwrap();
let template = test_server_config(dir.path());
let store = Arc::new(SqliteWorkspaceStore::open(&template.database_path).unwrap());
let catalog = WorkspaceCatalogService::new(store.clone());
let workspace_a = catalog
.create(
WorkspaceCreateRequest {
operation_key: "create-a".to_string(),
display_name: "Workspace A".to_string(),
repository: crate::workspace_catalog::InitialRepositoryIntent {
uri: repository_a.display().to_string(),
display_name: None,
default_ref: None,
},
},
None,
)
.unwrap();
let workspace_b = catalog
.create(
WorkspaceCreateRequest {
operation_key: "create-b".to_string(),
display_name: "Workspace B".to_string(),
repository: crate::workspace_catalog::InitialRepositoryIntent {
uri: repository_b.display().to_string(),
display_name: None,
default_ref: None,
},
},
None,
)
.unwrap();
let app = build_workspace_server_router(template, store)
.await
.unwrap();
let uri_a = format!("/api/w/{}/workspace", workspace_a.workspace.workspace_id);
let uri_b = format!("/api/w/{}/workspace", workspace_b.workspace.workspace_id);
let (a, b) = tokio::join!(get_json(app.clone(), &uri_a), get_json(app.clone(), &uri_b));
assert_eq!(a["workspace_id"], workspace_a.workspace.workspace_id);
assert_eq!(a["display_name"], "Workspace A");
assert_eq!(b["workspace_id"], workspace_b.workspace.workspace_id);
assert_eq!(b["display_name"], "Workspace B");
let handle = missing_resource_handle();
let resource_response = app
.clone()
.oneshot(
Request::post(format!(
"/internal/w/{}/runtime/resources/fetch",
workspace_b.workspace.workspace_id
))
.header(axum::http::header::CONTENT_TYPE, "application/json")
.body(Body::from(
serde_json::to_vec(&BackendResourceFetchRequest {
audit_correlation_id: handle.audit_correlation_id.clone(),
runtime_id: "runtime-test".to_string(),
worker_id: None,
handle,
})
.unwrap(),
))
.unwrap(),
)
.await
.unwrap();
assert_eq!(resource_response.status(), StatusCode::NOT_FOUND);
let resource_error: BackendResourceError = serde_json::from_slice(
&to_bytes(resource_response.into_body(), usize::MAX)
.await
.unwrap(),
)
.unwrap();
assert_eq!(resource_error, BackendResourceError::MissingResource);
let missing = app
.oneshot(
Request::builder()
.uri("/api/w/00000000-0000-0000-0000-000000000001/workspace")
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(missing.status(), StatusCode::NOT_FOUND);
}
#[tokio::test]
async fn local_bootstrap_create_activates_workspace_without_server_restart() {
let dir = tempfile::tempdir().unwrap();
let repository = dir.path().join("repository");
std::fs::create_dir_all(repository.join(".git")).unwrap();
let template = test_server_config(dir.path()).with_local_workspace_bootstrap(true);
let store = Arc::new(SqliteWorkspaceStore::open(&template.database_path).unwrap());
let app = build_workspace_server_router(template, store)
.await
.unwrap();
let payload = json!({
"operation_key": "bootstrap-1",
"display_name": "Created Workspace",
"repository": {
"uri": repository,
"display_name": "Repository",
"default_ref": "HEAD"
}
});
let created = app
.clone()
.oneshot(
Request::builder()
.method(Method::POST)
.uri("/api/workspaces")
.header(axum::http::header::CONTENT_TYPE, "application/json")
.body(Body::from(payload.to_string()))
.unwrap(),
)
.await
.unwrap();
assert_eq!(created.status(), StatusCode::CREATED);
let body = to_bytes(created.into_body(), usize::MAX).await.unwrap();
let body: Value = serde_json::from_slice(&body).unwrap();
let workspace_id = body["workspace"]["workspace_id"].as_str().unwrap();
let workspace = get_json(app.clone(), &format!("/api/w/{workspace_id}/workspace")).await;
assert_eq!(workspace["display_name"], "Created Workspace");
let replayed = app
.clone()
.oneshot(
Request::builder()
.method(Method::POST)
.uri("/api/workspaces")
.header(axum::http::header::CONTENT_TYPE, "application/json")
.body(Body::from(payload.to_string()))
.unwrap(),
)
.await
.unwrap();
assert_eq!(replayed.status(), StatusCode::OK);
let second_payload = json!({
"operation_key": "bootstrap-2",
"display_name": "Second Ownerless Workspace",
"repository": {
"uri": repository,
"display_name": "Repository",
"default_ref": "HEAD"
}
});
let second = app
.oneshot(
Request::builder()
.method(Method::POST)
.uri("/api/workspaces")
.header(axum::http::header::CONTENT_TYPE, "application/json")
.body(Body::from(second_payload.to_string()))
.unwrap(),
)
.await
.unwrap();
assert_eq!(second.status(), StatusCode::CONFLICT);
}
#[test]
fn scoped_workspace_path_requires_an_explicit_workspace_segment() {
assert_eq!(
scoped_workspace_id("/api/w/workspace-a/tickets"),
Some("workspace-a")
);
assert_eq!(
scoped_workspace_id("/w/workspace-b/workers"),
Some("workspace-b")
);
assert_eq!(
scoped_workspace_id("/internal/w/workspace-c/runtime/resources/fetch"),
Some("workspace-c")
);
assert_eq!(scoped_workspace_id("/api/workspaces"), None);
assert_eq!(scoped_workspace_id("/api/workspace"), None);
}
fn memory_staging_record_json(id: &str, claim: &str) -> String {
json!({
"schema_version": 1,
@@ -14392,27 +14842,7 @@ mod tests {
let mut missing = ticket::NewTicket::new("Missing target");
missing.repository_id = Some("unknown".to_owned());
let missing = backend.create(missing).unwrap();
assert!(matches!(
backend.mark_ready(
TicketIdOrSlug::Id(missing.id.clone()),
ticket::TicketMarkReady {
operation_key: "missing-repository".to_owned(),
reason: None,
author: None,
intake_summary: None,
},
),
Err(ticket::TicketError::UnknownTargetRepository(_))
));
assert_eq!(
backend
.show(TicketIdOrSlug::Id(missing.id))
.unwrap()
.meta
.workflow_state,
TicketWorkflowState::Planning
);
assert!(backend.create(missing).is_err());
assert!(matches!(
backend.set_workflow_state(
TicketIdOrSlug::Id(ticket_ref.id),
@@ -14604,6 +15034,21 @@ mod tests {
.create(ticket::NewTicket::new("Assigned Ticket"))
.unwrap();
let ticket_id = created.id;
api.store
.upsert_worker_registry(&WorkerRegistryRecord {
workspace_id: TEST_WORKSPACE_ID.to_string(),
worker: RuntimeWorkerRef::new("embedded", "42"),
display_name: "Worker 42".to_string(),
profile: Some("builtin:coder".to_string()),
retention_state: "normal".to_string(),
transcript_ref: None,
session_ref: None,
summary_ref: None,
diagnostics_ref: None,
created_at: TEST_CREATED_AT.to_string(),
updated_at: TEST_CREATED_AT.to_string(),
})
.unwrap();
let assignment = TicketWorkerAssignmentRecord {
workspace_id: TEST_WORKSPACE_ID.to_string(),
ticket_id: ticket_id.clone(),
@@ -14698,6 +15143,24 @@ mod tests {
.unwrap()
.worker
.unwrap();
api.store
.upsert_worker_registry(&WorkerRegistryRecord {
workspace_id: TEST_WORKSPACE_ID.to_string(),
worker: RuntimeWorkerRef::new(
EMBEDDED_WORKER_RUNTIME_ID,
source_worker.worker.worker_id.clone(),
),
display_name: "Source Worker".to_string(),
profile: Some("builtin:coder".to_string()),
retention_state: "normal".to_string(),
transcript_ref: None,
session_ref: None,
summary_ref: None,
diagnostics_ref: None,
created_at: TEST_CREATED_AT.to_string(),
updated_at: TEST_CREATED_AT.to_string(),
})
.unwrap();
let recipient_worker = api
.runtime
.spawn_worker(
@@ -14708,6 +15171,24 @@ mod tests {
.unwrap()
.worker
.unwrap();
api.store
.upsert_worker_registry(&WorkerRegistryRecord {
workspace_id: TEST_WORKSPACE_ID.to_string(),
worker: RuntimeWorkerRef::new(
EMBEDDED_WORKER_RUNTIME_ID,
recipient_worker.worker.worker_id.clone(),
),
display_name: "Recipient Worker".to_string(),
profile: Some("builtin:coder".to_string()),
retention_state: "normal".to_string(),
transcript_ref: None,
session_ref: None,
summary_ref: None,
diagnostics_ref: None,
created_at: TEST_CREATED_AT.to_string(),
updated_at: TEST_CREATED_AT.to_string(),
})
.unwrap();
let backend = browser_ticket_backend(&api).unwrap();
let ticket_ref = backend
.create(ticket::NewTicket::new("Notify assigned Worker"))
@@ -16019,6 +16500,7 @@ mod tests {
let mut config = test_server_config(temp.path());
config.remote_runtime_sources.push(RemoteRuntimeConfig {
runtime_id: "runtime-remote".to_string(),
workspace_id: Some(TEST_WORKSPACE_ID.to_string()),
display_name: "Remote Runtime".to_string(),
base_url: "https://runtime.invalid".to_string(),
bearer_token: None,
@@ -16046,8 +16528,20 @@ mod tests {
timeout: std::time::Duration::from_secs(1),
});
let store = SqliteWorkspaceStore::open(config.database_path.clone()).unwrap();
store
.upsert_workspace(&WorkspaceRecord {
workspace_id: TEST_WORKSPACE_ID.to_string(),
owner_account_id: None,
display_name: "Test Workspace".to_string(),
state: "active".to_string(),
created_at: "2026-08-11T00:00:00Z".to_string(),
updated_at: "2026-08-11T00:00:00Z".to_string(),
})
.await
.unwrap();
let trust = crate::store::TrustedRuntimeRecord {
runtime_id: "runtime-remote".to_string(),
workspace_id: Some(TEST_WORKSPACE_ID.to_string()),
display_name: "Remote Runtime".to_string(),
base_url: "https://runtime.invalid".to_string(),
public_key: identity.public_key.clone(),
@@ -16711,20 +17205,20 @@ mod tests {
let handle = missing_resource_handle();
let response = app
.oneshot(
Request::post("/internal/runtime/resources/fetch")
.header("content-type", "application/json")
.body(Body::from(
serde_json::to_vec(
&worker_runtime::resource::BackendResourceFetchRequest {
audit_correlation_id: handle.audit_correlation_id.clone(),
runtime_id: "runtime-test".to_string(),
worker_id: None,
handle,
},
)
.unwrap(),
))
Request::post(format!(
"/internal/w/{TEST_WORKSPACE_ID}/runtime/resources/fetch"
))
.header("content-type", "application/json")
.body(Body::from(
serde_json::to_vec(&worker_runtime::resource::BackendResourceFetchRequest {
audit_correlation_id: handle.audit_correlation_id.clone(),
runtime_id: "runtime-test".to_string(),
worker_id: None,
handle,
})
.unwrap(),
))
.unwrap(),
)
.await
.unwrap();
@@ -16747,7 +17241,7 @@ mod tests {
let archive = test_profile_archive();
let runtime_id = "runtime-test";
let handle = broker.issue_profile_source_archive_handle(
"workspace-test",
TEST_WORKSPACE_ID,
crate::resource_broker::BackendResourceTarget::Runtime(runtime_id),
archive,
);
@@ -16756,7 +17250,7 @@ mod tests {
let addr = listener.local_addr().unwrap();
let server = tokio::spawn(async move { axum::serve(listener, app).await.unwrap() });
let client = worker_runtime::resource::HttpBackendResourceClient::new(
format!("http://{addr}/internal/runtime/resources/fetch"),
format!("http://{addr}/internal/w/{TEST_WORKSPACE_ID}/runtime/resources/fetch"),
None,
);
@@ -18947,6 +19441,26 @@ mod tests {
})
.await
.unwrap();
rusqlite::Connection::open(&config.database_path)
.unwrap()
.execute_batch(
r#"
INSERT INTO typed_tickets (
workspace_id, ticket_id, slug, title, status, kind, priority, body,
workflow_state, workflow_state_explicit
) VALUES
('0192f0e8-4d84-7d6e-a000-000000000001', '00000000001J2', 'ticket-j2', 'Ticket J2', 'open', 'task', 'normal', '', 'planning', 1),
('0192f0e8-4d84-7d6e-a000-000000000001', '00000000001J3', 'ticket-j3', 'Ticket J3', 'open', 'task', 'normal', '', 'planning', 1);
INSERT INTO workspace_resource_human_keys (
workspace_id, resource_kind, resource_id, sequence, human_key, allocated_at
) VALUES
('0192f0e8-4d84-7d6e-a000-000000000001', 'ticket', '00000000001J2', 1, 'T-1', '2026-01-01T00:00:00Z'),
('0192f0e8-4d84-7d6e-a000-000000000001', 'ticket', '00000000001J3', 2, 'T-2', '2026-01-01T00:00:00Z');
INSERT INTO workspace_resource_human_key_counters (workspace_id, resource_kind, next_sequence)
VALUES ('0192f0e8-4d84-7d6e-a000-000000000001', 'ticket', 3);
"#,
)
.unwrap();
let api = WorkspaceApi::new_with_execution_backend(
config,
store,
File diff suppressed because it is too large Load Diff
@@ -0,0 +1,394 @@
use std::path::{Path, PathBuf};
use std::sync::Arc;
use chrono::{SecondsFormat, Utc};
use serde::{Deserialize, Serialize};
use sha2::{Digest, Sha256};
use uuid::Uuid;
use crate::store::{
ControlPlaneStore, RepositoryRecord, WorkspaceBootstrapRecord, WorkspaceRecord,
};
use crate::{Error, Result};
const DEFAULT_REPOSITORY_ID: &str = "main";
const MAX_DISPLAY_NAME_BYTES: usize = 200;
const MAX_OPERATION_KEY_BYTES: usize = 200;
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[serde(deny_unknown_fields)]
pub struct InitialRepositoryIntent {
pub uri: String,
#[serde(default)]
pub display_name: Option<String>,
#[serde(default)]
pub default_ref: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[serde(deny_unknown_fields)]
pub struct WorkspaceCreateRequest {
pub operation_key: String,
pub display_name: String,
pub repository: InitialRepositoryIntent,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub struct WorkspaceCreateResponse {
pub workspace: WorkspaceRecord,
pub repository: RepositoryRecord,
pub config_revision: u64,
pub request_fingerprint: String,
pub replayed: bool,
}
#[derive(Clone)]
pub struct WorkspaceCatalogService {
store: Arc<dyn ControlPlaneStore>,
}
impl WorkspaceCatalogService {
pub fn new(store: Arc<dyn ControlPlaneStore>) -> Self {
Self { store }
}
pub fn list(
&self,
owner_account_id: Option<&str>,
limit: usize,
) -> Result<Vec<WorkspaceRecord>> {
let limit = limit.clamp(1, 200);
Ok(self
.store
.list_workspaces()?
.into_iter()
.filter(|workspace| {
workspace.owner_account_id.is_none()
|| owner_account_id
.is_some_and(|owner| workspace.owner_account_id.as_deref() == Some(owner))
})
.take(limit)
.collect())
}
pub fn create(
&self,
request: WorkspaceCreateRequest,
owner_account_id: Option<String>,
) -> Result<WorkspaceCreateResponse> {
self.create_internal(request, owner_account_id, None, false)
}
pub fn create_first_ownerless(
&self,
request: WorkspaceCreateRequest,
) -> Result<WorkspaceCreateResponse> {
self.create_internal(request, None, None, true)
}
pub fn create_with_workspace_id(
&self,
request: WorkspaceCreateRequest,
owner_account_id: Option<String>,
requested_workspace_id: Option<String>,
) -> Result<WorkspaceCreateResponse> {
self.create_internal(request, owner_account_id, requested_workspace_id, false)
}
fn create_internal(
&self,
request: WorkspaceCreateRequest,
owner_account_id: Option<String>,
requested_workspace_id: Option<String>,
require_empty_catalog: bool,
) -> Result<WorkspaceCreateResponse> {
let operation_key = normalize_required(
"operation_key",
request.operation_key,
MAX_OPERATION_KEY_BYTES,
)?;
let display_name =
normalize_required("display_name", request.display_name, MAX_DISPLAY_NAME_BYTES)?;
let repository_path = validate_repository_uri(&request.repository.uri)?;
let repository_uri = repository_path.to_string_lossy().into_owned();
let repository_name = request
.repository
.display_name
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
.unwrap_or("Main repository")
.to_string();
let default_ref = request
.repository
.default_ref
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
.unwrap_or("HEAD")
.to_string();
let requested_workspace_id = requested_workspace_id
.map(|value| {
Uuid::parse_str(value.trim())
.map(|id| id.to_string())
.map_err(|_| Error::InvalidInput("workspace_id must be a UUID".to_string()))
})
.transpose()?;
let workspace_id = requested_workspace_id
.clone()
.unwrap_or_else(|| Uuid::now_v7().to_string());
let fingerprint = workspace_create_fingerprint(
requested_workspace_id.as_deref(),
&display_name,
owner_account_id.as_deref(),
&repository_uri,
&repository_name,
&default_ref,
);
let now = Utc::now().to_rfc3339_opts(SecondsFormat::Millis, true);
let result = self
.store
.create_workspace_bootstrap(&WorkspaceBootstrapRecord {
operation_key,
request_fingerprint: fingerprint.clone(),
require_empty_catalog,
workspace: WorkspaceRecord {
workspace_id: workspace_id.clone(),
owner_account_id,
display_name,
state: "active".to_string(),
created_at: now.clone(),
updated_at: now.clone(),
},
repository: RepositoryRecord {
workspace_id,
repository_id: DEFAULT_REPOSITORY_ID.to_string(),
name: repository_name,
kind: "git".to_string(),
provider: Some("git".to_string()),
uri: repository_uri,
default_ref: Some(default_ref),
auth_ref_kind: None,
auth_ref_key: None,
created_at: now.clone(),
updated_at: now,
},
})?;
Ok(WorkspaceCreateResponse {
workspace: result.workspace,
repository: result.repository,
config_revision: result.config_revision,
request_fingerprint: fingerprint,
replayed: result.replayed,
})
}
}
fn normalize_required(field: &str, value: String, max_bytes: usize) -> Result<String> {
let value = value.trim();
if value.is_empty() || value.len() > max_bytes {
return Err(Error::InvalidInput(format!(
"{field} must be between 1 and {max_bytes} bytes"
)));
}
Ok(value.to_string())
}
fn validate_repository_uri(uri: &str) -> Result<PathBuf> {
let uri = uri.trim();
if uri.is_empty() || uri.contains("://") {
return Err(Error::InvalidInput(
"initial repository uri must be an absolute server-local path".to_string(),
));
}
let path = Path::new(uri);
if !path.is_absolute() {
return Err(Error::InvalidInput(
"initial repository uri must be an absolute server-local path".to_string(),
));
}
let path = path.canonicalize().map_err(|error| {
Error::InvalidInput(format!("initial repository path is unavailable: {error}"))
})?;
if !path.is_dir() {
return Err(Error::InvalidInput(
"initial repository path must be a directory".to_string(),
));
}
let normal_git = path.join(".git").exists();
let bare_git = path.join("HEAD").is_file() && path.join("objects").is_dir();
if !normal_git && !bare_git {
return Err(Error::InvalidInput(
"initial repository path is not a Git repository".to_string(),
));
}
Ok(path)
}
fn workspace_create_fingerprint(
requested_workspace_id: Option<&str>,
display_name: &str,
owner_account_id: Option<&str>,
repository_uri: &str,
repository_name: &str,
default_ref: &str,
) -> String {
let payload = serde_json::json!({
"requested_workspace_id": requested_workspace_id,
"display_name": display_name,
"owner_account_id": owner_account_id,
"repository": {
"repository_id": DEFAULT_REPOSITORY_ID,
"uri": repository_uri,
"display_name": repository_name,
"default_ref": default_ref,
"kind": "git",
}
});
let mut hasher = Sha256::new();
hasher.update(serde_json::to_vec(&payload).expect("workspace fingerprint serializes"));
let digest = hasher.finalize();
let encoded = digest
.iter()
.map(|byte| format!("{byte:02x}"))
.collect::<String>();
format!("sha256:{encoded}")
}
#[cfg(test)]
mod tests {
use super::*;
use crate::store::SqliteWorkspaceStore;
fn git_repository() -> tempfile::TempDir {
let dir = tempfile::tempdir().unwrap();
std::fs::create_dir(dir.path().join(".git")).unwrap();
dir
}
#[tokio::test]
async fn create_is_atomic_and_exact_retries_converge() {
let store = Arc::new(SqliteWorkspaceStore::in_memory().unwrap());
let service = WorkspaceCatalogService::new(store.clone());
let repository = git_repository();
let request = WorkspaceCreateRequest {
operation_key: "request-1".to_string(),
display_name: "Workspace A".to_string(),
repository: InitialRepositoryIntent {
uri: repository.path().display().to_string(),
display_name: None,
default_ref: None,
},
};
let created = service.create(request.clone(), None).unwrap();
let replayed = service.create(request, None).unwrap();
assert!(!created.replayed);
assert!(replayed.replayed);
assert_eq!(
created.workspace.workspace_id,
replayed.workspace.workspace_id
);
assert_eq!(store.list_workspaces().unwrap().len(), 1);
assert_eq!(
store
.list_repositories(&created.workspace.workspace_id)
.unwrap()
.len(),
1
);
assert!(
store
.load_workspace_config(&created.workspace.workspace_id)
.unwrap()
.is_some()
);
}
#[test]
fn concurrent_ownerless_bootstrap_commits_exactly_one_workspace() {
let store = Arc::new(SqliteWorkspaceStore::in_memory().unwrap());
let service = WorkspaceCatalogService::new(store.clone());
let repository_a = git_repository();
let repository_b = git_repository();
let requests = [
WorkspaceCreateRequest {
operation_key: "bootstrap-a".to_string(),
display_name: "Workspace A".to_string(),
repository: InitialRepositoryIntent {
uri: repository_a.path().display().to_string(),
display_name: None,
default_ref: None,
},
},
WorkspaceCreateRequest {
operation_key: "bootstrap-b".to_string(),
display_name: "Workspace B".to_string(),
repository: InitialRepositoryIntent {
uri: repository_b.path().display().to_string(),
display_name: None,
default_ref: None,
},
},
];
let barrier = Arc::new(std::sync::Barrier::new(2));
let results = std::thread::scope(|scope| {
requests
.into_iter()
.map(|request| {
let service = service.clone();
let barrier = barrier.clone();
scope.spawn(move || {
barrier.wait();
service.create_first_ownerless(request)
})
})
.collect::<Vec<_>>()
.into_iter()
.map(|handle| handle.join().unwrap())
.collect::<Vec<_>>()
});
assert_eq!(results.iter().filter(|result| result.is_ok()).count(), 1);
assert_eq!(results.iter().filter(|result| result.is_err()).count(), 1);
assert_eq!(store.list_workspaces().unwrap().len(), 1);
let error = results
.into_iter()
.find_map(Result::err)
.unwrap()
.to_string();
assert!(error.contains("catalog is empty"), "{error}");
}
#[tokio::test]
async fn idempotency_key_reuse_with_different_payload_is_rejected() {
let store = Arc::new(SqliteWorkspaceStore::in_memory().unwrap());
let service = WorkspaceCatalogService::new(store);
let repository = git_repository();
let mut request = WorkspaceCreateRequest {
operation_key: "request-1".to_string(),
display_name: "Workspace A".to_string(),
repository: InitialRepositoryIntent {
uri: repository.path().display().to_string(),
display_name: None,
default_ref: None,
},
};
service.create(request.clone(), None).unwrap();
request.display_name = "Workspace B".to_string();
let error = service.create(request, None).unwrap_err().to_string();
assert!(error.contains("different input"), "{error}");
}
#[test]
fn repository_intent_rejects_remote_and_non_git_paths() {
let remote = validate_repository_uri("https://example.test/repo.git").unwrap_err();
assert!(remote.to_string().contains("server-local path"));
let dir = tempfile::tempdir().unwrap();
let non_git = validate_repository_uri(&dir.path().display().to_string()).unwrap_err();
assert!(non_git.to_string().contains("not a Git repository"));
}
}
+1
View File
@@ -22,6 +22,7 @@ It is not a dumping ground for external research, old plans, API inventories, or
14. [`development/work-items.md`](development/work-items.md) — how project work is recorded and reviewed.
15. [`development/rust-testing-strategy.md`](development/rust-testing-strategy.md) — what Yoi Rust tests should prove, where they belong, and how to name them.
16. [`development/validation.md`](development/validation.md) — how to check changes.
17. [`development/workspace-schema-migrations.md`](development/workspace-schema-migrations.md) — how to preflight, apply, verify, and roll back control-plane SQLite schema changes.
## What belongs here
@@ -0,0 +1,42 @@
# Workspace database schema migration runbook
The Workspace Server owns one control-plane SQLite database. Schema changes are applied by the Server at startup; domain components such as Ticket and Merge Request contribute tables to that same database, but they do not create a second Workspace authority.
## Before deployment
1. Stop writes and shut down every Server process using the database. Do not run two Server generations against one database during migration.
2. Record the current binary revision and database schema version.
3. Take a byte-for-byte backup of the database and its WAL/SHM state using a SQLite-safe backup procedure.
4. Run the read-only plan with the new binary:
```sh
yoi-server migrate --dry-run --database <server.db>
```
The plan runs against an in-memory copy. It reports the current and target schema versions, migration names, Worker identity mappings, and repairs without mutating the source database. Workspace-resource preflight failures name the relation and bounded offending row identities; repair those rows through the owning domain authority before retrying.
## Applying
Start exactly one instance of the new Server binary against the database. Startup applies migration 39 in one SQLite transaction after the Ticket and Merge Request component schemas are available. The migration:
- rebuilds Ticket, Objective, assignment, Artifact, and human-key tables with Workspace-scoped composite identity;
- adds composite foreign keys for repository, Ticket, Objective, Worker, relation-target, and current-assignment references;
- validates new historical assignment/event references with SQLite triggers while allowing those audit rows to survive later Ticket or Worker retention deletion; parent delete/Runtime-move triggers record exact Workspace-scoped tombstones, and startup accepts a missing live parent only when that tombstone exists, so an unrelated same ID in another Workspace cannot change the result; reservation operation ids remain intentionally unconstrained until their resources exist;
- checks the rebuilt schema with `PRAGMA foreign_key_check` before recording the schema version; and
- restores `PRAGMA foreign_keys = ON` whether the transaction commits or rolls back.
After startup, verify:
```sql
SELECT MAX(version) FROM __yoi_schema_migrations;
PRAGMA foreign_key_check;
PRAGMA integrity_check;
```
The expected migration version is `39`, `foreign_key_check` returns no rows, and `integrity_check` returns `ok`.
## Failure and rollback
There is no in-place down migration. A failed migration transaction leaves the prior schema version and data intact. Keep the Server stopped, preserve the failure diagnostics, and either repair the preflight data with the prior generation or restore the complete pre-migration backup before retrying.
Never run an older binary after a newer schema version has committed. Startup fences this case and refuses to serve when the database schema version is newer than the binary supports. Rollback therefore means restoring both the prior binary and its matching pre-migration database backup; it does not mean pointing the old binary at the upgraded database.
+1 -1
View File
@@ -188,4 +188,4 @@ in_flight?: InFlightSnapshot,
* Parent-owned Internal Worker sessions visible to this client.
* Service-private Internal Workers are deliberately excluded.
*/
internal_workers?: Array<InternalWorkerSnapshot>, } } | { "event": "internal_worker", "data": { worker: InternalWorkerRef, revision: number, event: Event, } } | { "event": "segment_rotated", "data": { entry: unknown, } } | { "event": "status", "data": { status: WorkerStatus, } } | { "event": "command", "data": { event: CommandEvent, } } | { "event": "completions", "data": { kind: CompletionKind, entries: Array<CompletionEntry>, } } | { "event": "rewind_targets", "data": { head_entries: number, targets: Array<RewindTarget>, } } | { "event": "rewind_applied", "data": { entries: Array<unknown>, input: Array<Segment>, summary: RewindSummary, } } | { "event": "workers_listed", "data": { workers: unknown, } } | { "event": "worker_restored", "data": { result: unknown, } } | { "event": "peer_registered", "data": { result: unknown, } } | { "event": "alert", "data": Alert } | { "event": "memory_worker", "data": MemoryWorkerEvent } | { "event": "compact_start" } | { "event": "compact_done", "data": { new_segment_id: string, } } | { "event": "compact_failed", "data": { error: string, } } | { "event": "shutdown" };
internal_workers?: Array<InternalWorkerSnapshot>, } } | { "event": "internal_worker", "data": { worker: InternalWorkerRef, revision: number, event: Event, } } | { "event": "internal_worker_removed", "data": { worker: InternalWorkerRef, revision: number, } } | { "event": "segment_rotated", "data": { entry: unknown, } } | { "event": "status", "data": { status: WorkerStatus, } } | { "event": "command", "data": { event: CommandEvent, } } | { "event": "completions", "data": { kind: CompletionKind, entries: Array<CompletionEntry>, } } | { "event": "rewind_targets", "data": { head_entries: number, targets: Array<RewindTarget>, } } | { "event": "rewind_applied", "data": { entries: Array<unknown>, input: Array<Segment>, summary: RewindSummary, } } | { "event": "workers_listed", "data": { workers: unknown, } } | { "event": "worker_restored", "data": { result: unknown, } } | { "event": "peer_registered", "data": { result: unknown, } } | { "event": "alert", "data": Alert } | { "event": "memory_worker", "data": MemoryWorkerEvent } | { "event": "compact_start" } | { "event": "compact_done", "data": { new_segment_id: string, } } | { "event": "compact_failed", "data": { error: string, } } | { "event": "shutdown" };
@@ -1418,7 +1418,10 @@ Deno.test("Internal Worker output stays separate and revision-fenced", () => {
}]);
assertEquals(projection.lines, []);
assertEquals(projection.internalWorkers.length, 1);
assertEquals(projection.internalWorkers[0].console.lines[0].body, "child output");
assertEquals(
projection.internalWorkers[0].console.lines[0].body,
"child output",
);
projection = projector.append([{
eventId: "2",
@@ -1495,15 +1498,112 @@ Deno.test("parent snapshot authoritatively replaces Internal Worker projections"
},
}]);
const projection = projector.append([{ eventId: "snapshot", event }]);
assertEquals(projection.internalWorkers.map((worker) => worker.worker.session_id), [
"replacement",
]);
assertEquals(
projection.internalWorkers.map((worker) => worker.worker.session_id),
[
"replacement",
],
);
const childLines = projection.internalWorkers[0].console.lines;
assertEquals(childLines.length, 1);
assertEquals(new Set(childLines.map((line) => line.id)).size, 1);
assertEquals(childLines[0].kind, "tool");
});
Deno.test("terminal Internal Worker removal drops descendants and fences late events", () => {
const worker = {
session_id: "child-session",
name: "child",
parent_session_id: "parent-session",
kind: "sub_worker" as const,
};
const nestedWorker = {
session_id: "grandchild-session",
name: "grandchild",
parent_session_id: "child-session",
kind: "sub_worker" as const,
};
const projector = createConsoleProjector();
let projection = projector.append([{
eventId: "child",
event: {
event: "internal_worker",
data: {
worker,
revision: 2,
event: {
event: "internal_worker",
data: {
worker: nestedWorker,
revision: 1,
event: { event: "text_done", data: { text: "nested" } },
},
},
},
},
}]);
assertEquals(projection.internalWorkers.length, 1);
assertEquals(
projection.internalWorkers[0].console.internalWorkers.length,
1,
);
projection = projector.append([{
eventId: "removed",
event: {
event: "internal_worker_removed",
data: { worker, revision: 3 },
},
}, {
eventId: "late",
event: {
event: "internal_worker",
data: {
worker,
revision: 4,
event: { event: "text_done", data: { text: "must stay removed" } },
},
},
}]);
assertEquals(projection.internalWorkers, []);
const snapshot = snapshotEvent("/repo");
projection = projector.append([{ eventId: "snapshot", event: snapshot }]);
assertEquals(projection.internalWorkers, []);
assertEquals(projection.removedInternalWorkers, {});
});
Deno.test("stale Internal Worker removal cannot discard a newer projection", () => {
const worker = {
session_id: "child-session",
name: "child",
parent_session_id: "parent-session",
kind: "sub_worker" as const,
};
const projector = createConsoleProjector();
projector.append([{
eventId: "current",
event: {
event: "internal_worker",
data: {
worker,
revision: 4,
event: { event: "text_done", data: { text: "current" } },
},
},
}]);
const projection = projector.append([{
eventId: "stale-removal",
event: {
event: "internal_worker_removed",
data: { worker, revision: 3 },
},
}]);
assertEquals(projection.internalWorkers.length, 1);
assertEquals(projection.internalWorkers[0].revision, 4);
});
Deno.test("snapshot restores TaskStore state from system history", () => {
const taskSnapshot =
`[Session TaskStore snapshot]\n\n\`\`\`json\n{\n "tasks": [{"taskid": 3, "status": "pending", "subject": "Restored", "description": "From compaction"}]\n}\n\`\`\``;
@@ -98,6 +98,8 @@ export type ConsoleProjection = {
cwd: string | null;
lastEventId: string | null;
internalWorkers: InternalWorkerProjection[];
/** Terminal child-session fences, reset only by an authoritative snapshot. */
removedInternalWorkers: Record<string, number>;
};
export type ConsoleTimelineLineSelection = {
@@ -183,6 +185,7 @@ export function emptyConsoleProjection(): ConsoleProjection {
cwd: null,
lastEventId: null,
internalWorkers: [],
removedInternalWorkers: {},
};
}
@@ -424,6 +427,7 @@ export function applyProtocolEvent(
cwd: projection.cwd,
lastEventId: envelope.eventId,
internalWorkers: [...projection.internalWorkers],
removedInternalWorkers: { ...projection.removedInternalWorkers },
};
const event = envelope.event;
@@ -549,9 +553,16 @@ export function applyProtocolEvent(
next.internalWorkers = (event.data.internal_workers ?? []).map((worker) =>
projectInternalWorkerSnapshot(worker, envelope.eventId, next.cwd)
);
next.removedInternalWorkers = {};
break;
}
case "internal_worker": {
if (
Object.hasOwn(
next.removedInternalWorkers,
event.data.worker.session_id,
)
) break;
const existingIndex = next.internalWorkers.findIndex((worker) =>
worker.worker.session_id === event.data.worker.session_id
);
@@ -576,6 +587,19 @@ export function applyProtocolEvent(
else next.internalWorkers.push(updated);
break;
}
case "internal_worker_removed": {
const existingIndex = next.internalWorkers.findIndex((worker) =>
worker.worker.session_id === event.data.worker.session_id
);
const existingRevision = existingIndex >= 0
? next.internalWorkers[existingIndex].revision
: 0;
if (event.data.revision <= existingRevision) break;
next.removedInternalWorkers[event.data.worker.session_id] =
event.data.revision;
if (existingIndex >= 0) next.internalWorkers.splice(existingIndex, 1);
break;
}
case "status":
next.status = event.data.status;
break;
@@ -1492,6 +1516,7 @@ function snapshotProjectionFromEntries(
cwd,
lastEventId: eventId,
internalWorkers: [],
removedInternalWorkers: {},
};
entries.forEach((entry, index) =>
applyLogEntry(projection, `${eventId}-snapshot-${index}`, entry)