Skip to content
3 changes: 3 additions & 0 deletions migrations/20260305000002_worker_directory.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
-- Persist the working directory for opencode workers so that idle workers
-- can be resumed into the correct directory after a restart.
ALTER TABLE worker_runs ADD COLUMN directory TEXT;
1 change: 1 addition & 0 deletions src/agent/channel.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2250,6 +2250,7 @@ impl Channel {
worker_type,
&self.deps.agent_id,
*interactive,
None,
);
}
ProcessEvent::WorkerStatus {
Expand Down
270 changes: 265 additions & 5 deletions src/agent/channel_dispatch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -466,6 +466,12 @@ pub async fn spawn_opencode_worker_from_state(
directory: &str,
interactive: bool,
) -> std::result::Result<crate::WorkerId, AgentError> {
if !interactive {
return Err(AgentError::Other(anyhow::anyhow!(
"OpenCode workers must be interactive"
)));
}

check_worker_limit(state).await?;
ensure_dispatch_readiness(state, "opencode_worker");
let task = task.into();
Expand All @@ -487,6 +493,17 @@ pub async fn spawn_opencode_worker_from_state(

let server_pool = rc.opencode_server_pool.load().clone();

// Prevent multiple opencode workers on the same directory.
server_pool
.claim_directory(&directory)
.await
.map_err(AgentError::Other)?;

// Clone for the release call in the async worker task.
let release_pool = server_pool.clone();
let release_directory = directory.clone();
let persist_directory = directory.clone();

let oc_secrets_store = state.deps.runtime_config.secrets.load().as_ref().clone();

let worker = if interactive {
Expand All @@ -504,10 +521,11 @@ pub async fn spawn_opencode_worker_from_state(
.write()
.await
.insert(worker_id, input_tx);
match &oc_secrets_store {
let worker = match &oc_secrets_store {
Some(store) => worker.with_secrets_store(store.clone()),
None => worker,
}
};
worker.with_sqlite_pool(state.deps.sqlite_pool.clone())
} else {
let worker = crate::opencode::OpenCodeWorker::new(
Some(state.channel_id.clone()),
Expand All @@ -517,10 +535,11 @@ pub async fn spawn_opencode_worker_from_state(
server_pool,
state.deps.event_tx.clone(),
);
match &oc_secrets_store {
let worker = match &oc_secrets_store {
Some(store) => worker.with_secrets_store(store.clone()),
None => worker,
}
};
worker.with_sqlite_pool(state.deps.sqlite_pool.clone())
};

let worker_id = worker.id;
Expand All @@ -540,7 +559,12 @@ pub async fn spawn_opencode_worker_from_state(
Some(state.channel_id.clone()),
oc_secrets_store,
async move {
let result = worker.run().await.map_err(SpacebotError::from)?;
let result = worker.run().await.map_err(SpacebotError::from);

// Release the directory claim regardless of success or failure.
release_pool.release_directory(&release_directory).await;

let result = result?;

// Persist the transcript built from SSE events so the worker detail
// view can show the full conversation (text + tool calls + results).
Expand Down Expand Up @@ -592,6 +616,12 @@ pub async fn spawn_opencode_worker_from_state(
})
.ok();

// Persist the directory so idle workers can be resumed into the correct
// directory after a restart.
state
.process_run_logger
.log_worker_directory(worker_id, &persist_directory);

tracing::info!(worker_id = %worker_id, task = %task, interactive, "OpenCode worker spawned");

Ok(worker_id)
Expand Down Expand Up @@ -704,6 +734,236 @@ where
})
}

/// Resume an idle interactive worker into a channel's state after restart.
///
/// Loads the prior transcript, creates a resumed worker (builtin or opencode),
/// registers it into the channel's worker_inputs/worker_handles/status_block,
/// and spawns the follow-up loop. Returns `Ok(worker_id)` on success, or
/// an error string if the worker couldn't be resumed.
pub async fn resume_idle_worker_into_state(
state: &ChannelState,
idle_worker: &crate::conversation::history::IdleWorkerRow,
) -> std::result::Result<WorkerId, String> {
let worker_id: WorkerId = idle_worker
.id
.parse::<uuid::Uuid>()
.map_err(|error| format!("invalid worker ID '{}': {error}", idle_worker.id))?;

match idle_worker.worker_type.as_str() {
"opencode" => {
let session_id = idle_worker
.opencode_session_id
.as_deref()
.ok_or("opencode worker has no session_id, cannot resume")?;

let rc = &state.deps.runtime_config;
let opencode_config = rc.opencode.load();
if !opencode_config.enabled {
return Err("OpenCode workers are not enabled".into());
}

let directory = idle_worker
.directory
.as_deref()
.map(std::path::PathBuf::from)
.unwrap_or_else(|| rc.workspace_dir.clone());
let server_pool = rc.opencode_server_pool.load().clone();

let result = crate::opencode::OpenCodeWorker::resume_interactive(
worker_id,
Some(state.channel_id.clone()),
state.deps.agent_id.clone(),
&idle_worker.task,
directory,
server_pool,
state.deps.event_tx.clone(),
session_id.to_string(),
idle_worker.transcript.clone(),
)
.await;
Comment on lines +772 to +783

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⚠️ Potential issue | 🟠 Major

Missing claim_directory call for resumed OpenCode workers.

spawn_opencode_worker_from_state calls server_pool.claim_directory(&directory) to prevent concurrent workers on the same directory, but resume_idle_worker_into_state does not. This creates a race condition at startup:

  1. An idle OpenCode worker for /path/to/project is resumed
  2. Before it completes, a channel spawns a new OpenCode worker for the same directory
  3. The new spawn succeeds because no claim exists for the resumed worker

Additionally, there's no corresponding release_directory call in the resume path's completion handler.

🐛 Proposed fix: Add directory claim/release to resume path
             let directory = rc.workspace_dir.clone();
             let server_pool = rc.opencode_server_pool.load().clone();

+            // Claim directory to prevent concurrent workers (same as spawn path).
+            server_pool
+                .claim_directory(&directory)
+                .await
+                .map_err(|e| e.to_string())?;
+
+            let release_pool = server_pool.clone();
+            let release_directory = directory.clone();
+
             let result = crate::opencode::OpenCodeWorker::resume_interactive(

And in the async worker task (around line 807):

                 async move {
-                    let result = worker.run().await.map_err(SpacebotError::from)?;
+                    let result = worker.run().await.map_err(SpacebotError::from);
+
+                    // Release directory claim regardless of success/failure.
+                    release_pool.release_directory(&release_directory).await;
+
+                    let result = result?;
                     // Persist final transcript.
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@src/agent/channel_dispatch.rs` around lines 761 - 772,
resume_idle_worker_into_state is missing the directory claim/release that
spawn_opencode_worker_from_state uses, creating a race when resuming OpenCode
workers; before calling OpenCodeWorker::resume_interactive in
resume_idle_worker_into_state, call server_pool.claim_directory(&directory) and
ensure you release it (server_pool.release_directory(&directory)) in the resumed
worker's completion/error handler or finally block to mirror
spawn_opencode_worker_from_state behavior and avoid leaving directories claimed
on early exits.


let (mut worker, input_tx) = result.ok_or_else(|| {
"failed to reconnect to OpenCode session (server dead or session expired)"
.to_string()
})?;

// Apply builder chain (same as spawn_opencode_worker_from_state).
let oc_secrets_store = state.deps.runtime_config.secrets.load().as_ref().clone();
if let Some(store) = &oc_secrets_store {
worker = worker.with_secrets_store(store.clone());
}
worker = worker.with_sqlite_pool(state.deps.sqlite_pool.clone());

state
.worker_inputs
.write()
.await
.insert(worker_id, input_tx);

let worker_span = tracing::info_span!(
"worker.resume",
worker_id = %worker_id,
channel_id = %state.channel_id,
task = %idle_worker.task,
worker_type = "opencode",
);
let sqlite_pool = state.deps.sqlite_pool.clone();
let handle = spawn_worker_task(
worker_id,
state.deps.event_tx.clone(),
state.deps.agent_id.clone(),
Some(state.channel_id.clone()),
oc_secrets_store,
async move {
let result = worker.run().await.map_err(SpacebotError::from)?;
// Persist final transcript.
if !result.transcript.is_empty() {
let blob = crate::conversation::worker_transcript::serialize_steps(
&result.transcript,
);
let tool_calls = result.tool_calls;
let wid = worker_id.to_string();
let pool = sqlite_pool.clone();
tokio::spawn(async move {
if let Err(error) = sqlx::query(
"UPDATE worker_runs SET transcript = ?, tool_calls = ? WHERE id = ?",
)
.bind(&blob)
.bind(tool_calls)
.bind(&wid)
.execute(&pool)
.await
{
tracing::warn!(%error, worker_id = wid, "failed to persist OpenCode transcript");
}
});
}
Ok::<String, SpacebotError>(result.result_text)
}
.instrument(worker_span),
);

state.worker_handles.write().await.insert(worker_id, handle);

let opencode_task = format!("[opencode] {}", idle_worker.task);
{
let mut status = state.status_block.write().await;
status.add_worker(worker_id, &opencode_task, false, true);
}

state
.deps
.event_tx
.send(ProcessEvent::WorkerStarted {
agent_id: state.deps.agent_id.clone(),
worker_id,
channel_id: Some(state.channel_id.clone()),
task: opencode_task,
worker_type: "opencode".into(),
interactive: true,
})
.ok();

tracing::info!(worker_id = %worker_id, task = %idle_worker.task, "OpenCode worker resumed");
Ok(worker_id)
}
_ => {
// Builtin worker resume: deserialize transcript blob back into
// Rig message history so the LLM can continue the conversation.
let prior_history = if let Some(blob) = &idle_worker.transcript {
let steps = crate::conversation::worker_transcript::deserialize_transcript(blob)
.map_err(|error| format!("failed to deserialize transcript: {error}"))?;
crate::conversation::worker_transcript::transcript_to_history(&steps)
} else {
return Err("no transcript blob to restore history from".into());
};

let rc = &state.deps.runtime_config;
let prompt_engine = rc.prompts.load();
let sandbox_enabled = state.deps.sandbox.mode_enabled();
let sandbox_containment_active = state.deps.sandbox.containment_active();
let sandbox_read_allowlist = state.deps.sandbox.prompt_read_allowlist();
let sandbox_write_allowlist = state.deps.sandbox.prompt_write_allowlist();
let secrets_guard = rc.secrets.load();
let tool_secret_names = match (*secrets_guard).as_ref() {
Some(store) => store.tool_secret_names(),
None => Vec::new(),
};
let system_prompt = prompt_engine
.render_worker_prompt(
&rc.instance_dir.display().to_string(),
&rc.workspace_dir.display().to_string(),
sandbox_enabled,
sandbox_containment_active,
sandbox_read_allowlist,
sandbox_write_allowlist,
&tool_secret_names,
)
.map_err(|error| format!("failed to render worker prompt: {error}"))?;
let browser_config = (**rc.browser_config.load()).clone();
let brave_search_key = (**rc.brave_search_key.load()).clone();

let (worker, input_tx) = Worker::resume_interactive(
worker_id,
Some(state.channel_id.clone()),
&idle_worker.task,
&system_prompt,
state.deps.clone(),
browser_config,
state.screenshot_dir.clone(),
brave_search_key,
state.logs_dir.clone(),
prior_history,
);

state
.worker_inputs
.write()
.await
.insert(worker_id, input_tx);

let worker_span = tracing::info_span!(
"worker.resume",
worker_id = %worker_id,
channel_id = %state.channel_id,
task = %idle_worker.task,
);
let secrets_store = state.deps.runtime_config.secrets.load().as_ref().clone();
let handle = spawn_worker_task(
worker_id,
state.deps.event_tx.clone(),
state.deps.agent_id.clone(),
Some(state.channel_id.clone()),
secrets_store,
worker.run().instrument(worker_span),
);

state.worker_handles.write().await.insert(worker_id, handle);

{
let mut status = state.status_block.write().await;
status.add_worker(worker_id, &idle_worker.task, false, true);
}

state
.deps
.event_tx
.send(ProcessEvent::WorkerStarted {
agent_id: state.deps.agent_id.clone(),
worker_id,
channel_id: Some(state.channel_id.clone()),
task: idle_worker.task.clone(),
worker_type: "builtin".into(),
interactive: true,
})
.ok();

tracing::info!(worker_id = %worker_id, task = %idle_worker.task, "builtin worker resumed");
Ok(worker_id)
}
}
}

/// Expand a leading `~` or `~/` in a path to the user's home directory.
///
/// LLMs consistently produce tilde-prefixed paths because that's what appears
Expand Down
7 changes: 5 additions & 2 deletions src/agent/channel_history.rs
Original file line number Diff line number Diff line change
Expand Up @@ -416,8 +416,11 @@ pub(crate) fn event_is_for_channel(event: &ProcessEvent, channel_id: &ChannelId)
channel_id: event_channel,
..
} => event_channel.as_ref() == Some(channel_id),
ProcessEvent::OpenCodeSessionCreated { .. }
| ProcessEvent::OpenCodePartUpdated { .. }
ProcessEvent::OpenCodeSessionCreated {
channel_id: event_channel,
..
} => event_channel.as_ref() == Some(channel_id),
ProcessEvent::OpenCodePartUpdated { .. }
| ProcessEvent::StatusUpdate { .. }
| ProcessEvent::TaskUpdated { .. } => false,
}
Expand Down
2 changes: 2 additions & 0 deletions src/agent/cortex.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2491,6 +2491,7 @@ async fn pickup_one_ready_task(deps: &AgentDeps, logger: &CortexLogger) -> anyho
"task",
&deps.agent_id,
false,
None,
);

let task_store = deps.task_store.clone();
Expand Down Expand Up @@ -3613,6 +3614,7 @@ mod tests {
ProcessEvent::OpenCodeSessionCreated {
agent_id: Arc::from("agent"),
worker_id,
channel_id: Some(channel_id.clone()),
session_id: "session-1".to_string(),
port: 19898,
},
Expand Down
Loading