Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
59 changes: 59 additions & 0 deletions docs/architecture/remote-workspace-transport.md
Original file line number Diff line number Diff line change
Expand Up @@ -291,6 +291,65 @@ selection, or cache ownership. UI events retain workspace IDs through every
adapter; session operations recover their workspace from the session's stored ID.
Across devices, IDs are interpreted only by the selected owning host.

Request ownership also includes the device activation epoch. In the Web UI,
`invokePrepared` captures `SurfaceScope` before asynchronous parameter preparation
(including legacy capability negotiation), checks it before dispatch, and preserves
`SurfaceChangedError` through service error translation. `ApiClient` retains the
same activation through middleware, transport, retries and response handling.
Pending reads belong to one epoch through scoped keys or activation-owned caches;
returning to the same device does not revive an earlier activation's pending work.
Settled caches are keyed by the rendered device as well as the workspace ID,
because same-path local workspaces hash to the same ID on every device.
Multi-step preparation must
check the captured scope before starting another host request. Stream listeners
are detached on activation change and must never send cancellation to the newly
selected host for a search started elsewhere. Controller-local commands keep
their authority across activation in `invokePrepared` exactly as in `ApiClient`.
The Web UI lint configuration rejects `await` inside `api.invoke(...)`
arguments, because the activation would be captured after they resolve.

`CoreSessionStorePort` owns session storage resolution. The temporary path adapter
converts a legacy selector to a catalog ID once, then uses the same ID resolver as
current requests. It must not reconstruct identity from local filesystem
existence, choose a worktree's parent when its execution record is missing, or
return an execution path when resolution fails. A registered but unavailable
folder can still own readable persisted history.

After session admission, `SessionManager` retains the committed storage binding.
History restore, persistence, autosave, idle eviction and internal continuations
reuse that binding. The binding survives in-memory session eviction; a process
restart re-admits the session from its persisted ID and storage owner. Internal
queued turns retain the session's workspace ID, while external legacy submissions
still validate their locator at the compatibility boundary. All resolution uses
the persistence owner's `PathManager`. Readers never commit a binding: before
admission they resolve from the session's workspace configuration, and a pending
claim for a different location makes them fail instead of following an
uncommitted index entry.

The supported SSH history layout remains host plus remote root for upgrade
compatibility, so two saved connections to one host and root share a workspace
record and session mirror. Activation never rejects such records: imported or
persisted records for each connection stay listed and activatable, and session
identity verification, not workspace activation, keeps their histories apart.

Reopening an existing remote record with a different connection rebinds the
record only for an allowed reason:

- the two connection IDs are equivalent (the legacy `ssh-user@host:port` form and
the current `ssh-user@host` form);
- the previous owner is no longer a saved SSH connection;
- the user confirmed the rebind, sent as `rebindConnection` on
`open_remote_workspace` by Desktop and the CLI peer host.

Otherwise the open fails with the stable code
`remote_workspace_connection_conflict` as the whole error message, so remote
controllers can match it. The interactive Web UI asks the user and retries with
`rebindConnection`; startup restore defers with a localized notification instead
of rebinding silently. Older hosts ignore the field and keep their previous
behavior. Supporting multiple simultaneous endpoints for one host and root
requires a versioned storage-identity migration covering sessions and mirrors;
changing the directory hash alone is not a safe migration.

Persisted IDs are opaque. Catalog validation checks record/map-key agreement and
reference integrity; it must not recompute IDs from paths or require a working
SSH profile. Keep unavailable records in the catalog. Activation validates the
Expand Down
20 changes: 17 additions & 3 deletions src/apps/cli/src/peer_host/commands/workspace.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2,10 +2,13 @@

use std::path::PathBuf;

use openbitfun_core::service::workspace::{
remote_workspace_connection_conflict_message, RemoteConnectionRebind,
};
use openbitfun_runtime_ports::SessionStoragePathRequest;
use serde_json::{json, Value};

use crate::peer_host::args::{get_string, request_value};
use crate::peer_host::args::{get_string, optional_bool, optional_string, request_value};
use crate::peer_host::state::PeerHostState;
use crate::peer_host::workspace_dto::{workspace_info_to_json, workspace_list_to_json};

Expand Down Expand Up @@ -113,7 +116,12 @@ pub(crate) async fn open_remote_workspace(
let request = request_value(args);
let path = get_string(request, "remotePath")?;
let connection_id = get_string(request, "connectionId")?;
let host = crate::peer_host::args::optional_string(request, "sshHost");
let host = optional_string(request, "sshHost");
let rebind = if optional_bool(request, "rebindConnection").unwrap_or(false) {
RemoteConnectionRebind::UserConfirmed
} else {
RemoteConnectionRebind::Reject
};
let coordinator = openbitfun_core::agentic::coordination::get_global_coordinator()
.ok_or("Conversation coordinator is unavailable")?;
let info = coordinator
Expand All @@ -122,9 +130,15 @@ pub(crate) async fn open_remote_workspace(
&path,
&connection_id,
host.as_deref(),
rebind,
)
.await
.map_err(|e| e.to_string())?;
.map_err(|e| {
let message = e.to_string();
remote_workspace_connection_conflict_message(&message)
.map(str::to_string)
.unwrap_or(message)
})?;
Ok(workspace_info_to_json(&info))
}

Expand Down
19 changes: 18 additions & 1 deletion src/apps/desktop/src/api/commands.rs
Original file line number Diff line number Diff line change
Expand Up @@ -389,6 +389,10 @@ pub struct OpenRemoteWorkspaceRequest {
/// SSH config `host` (DNS or alias). When set, used for session mirror paths even if not connected.
#[serde(default)]
pub ssh_host: Option<String>,
/// The user confirmed moving an existing record owned by another saved
/// connection to `connection_id`.
#[serde(default)]
pub rebind_connection: bool,
}

#[derive(Debug, Deserialize, Default)]
Expand Down Expand Up @@ -1343,7 +1347,10 @@ pub async fn open_remote_workspace(
) -> Result<WorkspaceInfoDto, String> {
use openbitfun_core::service::remote_ssh::normalize_remote_workspace_path;
use openbitfun_core::service::remote_ssh::workspace_state::remote_workspace_stable_id;
use openbitfun_core::service::workspace::WorkspaceCreateOptions;
use openbitfun_core::service::workspace::{
remote_workspace_connection_conflict_message, RemoteConnectionRebind,
WorkspaceCreateOptions,
};

let ssh = state.get_ssh_manager_async().await?;
let saved = ssh
Expand Down Expand Up @@ -1431,6 +1438,11 @@ pub async fn open_remote_workspace(
remote_connection_id: Some(request.connection_id.clone()),
remote_ssh_host: Some(ssh_host.clone()),
stable_workspace_id: Some(stable_workspace_id),
remote_connection_rebind: if request.rebind_connection {
RemoteConnectionRebind::UserConfirmed
} else {
RemoteConnectionRebind::Reject
},
};

match state
Expand Down Expand Up @@ -1489,6 +1501,11 @@ pub async fn open_remote_workspace(
Ok(WorkspaceInfoDto::from_workspace_info(&workspace_info))
}
Err(e) => {
let message = e.to_string();
if let Some(conflict) = remote_workspace_connection_conflict_message(&message) {
warn!("Remote workspace open needs a connection rebind decision: {conflict}");
return Err(conflict.to_string());
}
error!("Failed to open remote workspace: {}", e);
Err(format!("Failed to open remote workspace: {}", e))
}
Expand Down
16 changes: 16 additions & 0 deletions src/crates/assembly/core/AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -221,6 +221,22 @@ are not the default Core precheck. For documentation-only changes, run
For assistant discovery, opened-state persistence, and reopening by workspace ID:
`cargo test --locked -p openbitfun-core --no-default-features --features agent-runtime,git --lib service::workspace::service::tests::assistant_`.

For local/SSH workspace identity collisions, SSH connection rebinding, committed
session storage, legacy path compatibility, and queue admission, use the matching
owner filter:

```bash
cargo test --locked -p openbitfun-core --no-default-features --features agent-runtime,git,remote-workspace --lib service::workspace::
cargo test --locked -p openbitfun-core --no-default-features --features agent-runtime,git,remote-workspace --lib agentic::session::
cargo test --locked -p openbitfun-core --no-default-features --features agent-runtime,git,remote-workspace --lib agentic::coordination::
```

The coordination filter includes the scheduler admission tests. These catalog and
storage fixtures do not require or validate a live SSH connection. CI also runs
them under `product-full`. `PathManager::with_user_root_for_tests` derives product
home from the user root's parent, so a fixture that relocates managed paths must
place its user root below the directory it means to exercise.

For disk-backed history paging and legacy sessions without a catalog:
`cargo test --locked -p openbitfun-core --no-default-features --features remote-connect,git --lib history_page_`.
Also run the `staged_revert_catalog_projection` and `load_relay_session_turns_`
Expand Down
90 changes: 23 additions & 67 deletions src/crates/assembly/core/src/agentic/coordination/coordinator.rs
Original file line number Diff line number Diff line change
Expand Up @@ -82,8 +82,8 @@ use crate::service::session::{
ToolItemIdentityExt, TurnStatus,
};
use crate::service::workspace::{
get_global_workspace_service, WorkspaceActivityMode, WorkspaceInfo, WorkspaceKind,
WorkspaceService,
get_global_workspace_service, RemoteConnectionRebind, WorkspaceActivityMode, WorkspaceInfo,
WorkspaceKind, WorkspaceService,
};
use crate::service_agent_runtime::CoreServiceAgentRuntime;
use crate::util::errors::{OpenBitFunError, OpenBitFunResult};
Expand Down Expand Up @@ -1885,31 +1885,9 @@ impl ConversationCoordinator {
&self,
session_id: &str,
) -> OpenBitFunResult<PathBuf> {
if let Some(binding) = self
.session_manager
.resolve_session_workspace_binding(session_id)
self.session_manager
.require_session_storage_path(session_id)
.await
{
return Ok(binding.session_storage_dir());
}

let session = self
.session_manager
.get_session(session_id)
.ok_or_else(|| {
OpenBitFunError::NotFound(format!("Session not found: {}", session_id))
})?;
session
.config
.workspace_path
.as_deref()
.map(PathBuf::from)
.ok_or_else(|| {
OpenBitFunError::Validation(format!(
"workspace_path is required when restoring session: {}",
session_id
))
})
}

async fn is_chinese_locale() -> bool {
Expand Down Expand Up @@ -2386,6 +2364,7 @@ Update the persona files and delete BOOTSTRAP.md as soon as bootstrap is complet
path: &str,
connection_id: &str,
ssh_host: Option<&str>,
remote_connection_rebind: RemoteConnectionRebind,
) -> OpenBitFunResult<WorkspaceInfo> {
let workspace = workspace_service
.prepare_remote_workspace(path, connection_id, ssh_host)
Expand All @@ -2399,7 +2378,7 @@ Update the persona files and delete BOOTSTRAP.md as soon as bootstrap is complet
.and_then(|host| host.as_str()),
)?;
workspace_service
.open_known_remote_workspace(&workspace)
.open_known_remote_workspace_with_rebind(&workspace, remote_connection_rebind)
.await
}

Expand Down Expand Up @@ -4178,8 +4157,7 @@ Update the persona files and delete BOOTSTRAP.md as soon as bootstrap is complet
|| (context_messages.len() == 1 && !session.dialog_turn_ids.is_empty()))
&& !session.dialog_turn_ids.is_empty()
{
let restore_path =
Self::resolve_session_restore_path(&project_workspace_path, None, None).await?;
let restore_path = self.restore_path_for_existing_session(&session_id).await?;
self.restore_session_from_storage_path(&restore_path, &session_id)
.await?;
session = self
Expand Down Expand Up @@ -4816,6 +4794,7 @@ Update the persona files and delete BOOTSTRAP.md as soon as bootstrap is complet
}

async fn resolve_session_restore_scope(
&self,
workspace_path: &str,
remote_connection_id: Option<&str>,
remote_ssh_host: Option<&str>,
Expand All @@ -4826,18 +4805,19 @@ Update the persona files and delete BOOTSTRAP.md as soon as bootstrap is complet
remote_ssh_host: remote_ssh_host.map(ToOwned::to_owned),
};

CoreSessionStorePort::default()
CoreSessionStorePort::with_path_manager(self.session_manager.path_manager())
.resolve_session_storage_path(request)
.await
.map_err(|error| OpenBitFunError::Session(error.to_string()))
}

async fn resolve_session_restore_path(
&self,
workspace_path: &str,
remote_connection_id: Option<&str>,
remote_ssh_host: Option<&str>,
) -> OpenBitFunResult<PathBuf> {
Self::resolve_session_restore_scope(workspace_path, remote_connection_id, remote_ssh_host)
self.resolve_session_restore_scope(workspace_path, remote_connection_id, remote_ssh_host)
.await
.map(|resolution| resolution.effective_storage_path)
}
Expand Down Expand Up @@ -6162,7 +6142,7 @@ Update the persona files and delete BOOTSTRAP.md as soon as bootstrap is complet
);
let requested_restore = match storage_workspace_path.as_deref() {
Some(workspace_path) => Some(
Self::resolve_session_restore_scope(
self.resolve_session_restore_scope(
workspace_path,
remote_connection_id.as_deref(),
remote_ssh_host.as_deref(),
Expand Down Expand Up @@ -6408,32 +6388,7 @@ Update the persona files and delete BOOTSTRAP.md as soon as bootstrap is complet
"Starting session history restore: session_id={}",
session_id
);
let restore_workspace_path = session
.config
.project_workspace_path
.as_deref()
.or(session.config.workspace_path.as_deref())
.or(storage_workspace_path.as_deref())
.ok_or_else(|| {
OpenBitFunError::Validation(format!(
"workspace_path is required when restoring session: {}",
session_id
))
})?;
let restore_path = Self::resolve_session_restore_path(
restore_workspace_path,
session
.config
.remote_connection_id
.as_deref()
.or(remote_connection_id.as_deref()),
session
.config
.remote_ssh_host
.as_deref()
.or(remote_ssh_host.as_deref()),
)
.await?;
let restore_path = self.restore_path_for_existing_session(&session_id).await?;
match self
.restore_session_from_storage_path(&restore_path, &session_id)
.await
Expand Down Expand Up @@ -8354,7 +8309,7 @@ Update the persona files and delete BOOTSTRAP.md as soon as bootstrap is complet
let session_storage_path = self
.session_manager
.resolve_storage_path_for_workspace_path(workspace_path)
.await;
.await?;
let has_revert_state = self
.session_manager
.persistence_manager()
Expand Down Expand Up @@ -15040,7 +14995,7 @@ impl openbitfun_runtime_ports::AgentThreadGoalManagementPort for ConversationCoo
.await
.map_err(runtime_port_error_preserving_message)?
} else {
Self::resolve_session_restore_path(
self.resolve_session_restore_path(
&request.workspace_path,
request.remote_connection_id.as_deref(),
request.remote_ssh_host.as_deref(),
Expand Down Expand Up @@ -20300,13 +20255,14 @@ mod tests {
ssh_host,
)
.await;
let storage_path = ConversationCoordinator::resolve_session_restore_path(
logical_workspace_path,
Some(connection_id),
Some(ssh_host),
)
.await
.expect("remote storage path should resolve");
let storage_path = coordinator
.resolve_session_restore_path(
logical_workspace_path,
Some(connection_id),
Some(ssh_host),
)
.await
.expect("remote storage path should resolve");
let goal = ThreadGoal {
goal_id: format!("goal-{index}"),
session_id: session_id.clone(),
Expand Down
Loading
Loading