mirror of
https://github.com/akitaonrails/ai-memory.git
synced 2026-10-02 03:24:46 +08:00
@@ -12,6 +12,14 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
|
||||
provider. (#1026)
|
||||
|
||||
### Fixed
|
||||
- Fixed managed Codex sessions losing their startup continuity when a shared
|
||||
Codex app-server daemon reports a finished run's stale `AI_MEMORY_RUN_ID`.
|
||||
SessionStart now recovers only the sole live, undelivered Codex run in the
|
||||
already-authorized repository, checkout cwd, and operator bucket, and links
|
||||
it atomically to the new native session. Active boundary mismatches,
|
||||
concurrent matching runs, cross-project or cross-worktree candidates,
|
||||
cross-operator candidates, and a second linker all fail closed instead of
|
||||
guessing or rebinding a run. (#987)
|
||||
- `ai-memory doctor` now warns when the nearest `.ai-memory.toml`'s
|
||||
`[capture]` section is invalid. This fails closed today — every file and
|
||||
shell tool event is reduced to metadata until the marker is fixed, nothing
|
||||
|
||||
@@ -1364,6 +1364,7 @@ pub async fn run(config: &Config, args: ServeArgs) -> Result<()> {
|
||||
reader: store.reader.clone(),
|
||||
sanitizer: sanitizer.clone(),
|
||||
data_dir: config.data_dir.clone(),
|
||||
trusted_proxy_identity: trusted_proxy_identity_enabled(&config.auth),
|
||||
});
|
||||
let admin = admin_router_with_sweep_tuning(
|
||||
AdminState {
|
||||
|
||||
@@ -1361,14 +1361,6 @@ async fn fetch_and_accept_handoff_at(
|
||||
now: jiff::Timestamp,
|
||||
) -> anyhow::Result<Option<String>> {
|
||||
let agent = query.agent.as_deref().map_or(AgentKind::Other, parse_agent);
|
||||
// A managed run's ledger is additive, not a replacement. Returning it here
|
||||
// skipped `latest_open_handoff` below, so a session launched by
|
||||
// `ai-memory run` never consumed the handoff a previous session left for
|
||||
// it — the slot just stayed open, and the next managed session missed it
|
||||
// too. The brief already reaches the managed path (it is recomposed per
|
||||
// session, so resolving it twice was harmless); the handoff is single-use
|
||||
// and had no second chance.
|
||||
let managed = fetch_managed_context(state, &query, agent, viewer).await?;
|
||||
// Keep the active-project key compatible with MCP transports: the native
|
||||
// session id is carried separately below to bind a destructive handoff
|
||||
// claim to its exact receiver.
|
||||
@@ -1398,6 +1390,11 @@ async fn fetch_and_accept_handoff_at(
|
||||
ai_memory_store::ProjectAccess::Read,
|
||||
)
|
||||
.await?;
|
||||
// A managed run's ledger is additive, not a replacement. Resolve and
|
||||
// authorize the request's repository before interpreting the run id: a
|
||||
// shared Codex daemon may carry a terminal id from another checkout, and
|
||||
// recovery must select only inside the current, authorized boundary.
|
||||
let managed = fetch_managed_context(state, &query, agent, ws, proj, actor.as_ref()).await?;
|
||||
// Session-start handoff delivery is a foreground action. Publish it so
|
||||
// static MCP callers resolve to the directory that is opening now. The
|
||||
// query carries no recall preference; the main capture path publishes
|
||||
@@ -1598,7 +1595,9 @@ async fn fetch_managed_context(
|
||||
state: &HookState,
|
||||
query: &HandoffQuery,
|
||||
agent: AgentKind,
|
||||
viewer: Option<ai_memory_core::UserId>,
|
||||
workspace_id: WorkspaceId,
|
||||
project_id: ProjectId,
|
||||
actor: Option<&IdentityKey>,
|
||||
) -> anyhow::Result<Option<PendingManagedContext>> {
|
||||
let Some(raw_run_id) = query.managed_run.as_deref() else {
|
||||
return Ok(None);
|
||||
@@ -1607,30 +1606,70 @@ async fn fetch_managed_context(
|
||||
warn!(managed_run = %raw_run_id, "invalid managed run id on SessionStart");
|
||||
return Ok(None);
|
||||
};
|
||||
// The run id reaches a workstream's event ledger without naming its
|
||||
// repository; the same rule as the `/workstream/runs/*` routes applies.
|
||||
if viewer.is_some() {
|
||||
let scope = state.reader.managed_run_scope(run_id).await?;
|
||||
crate::grants::authorize_resolved(
|
||||
&state.reader,
|
||||
scope,
|
||||
viewer,
|
||||
ai_memory_store::ProjectAccess::Write,
|
||||
)
|
||||
.await?;
|
||||
}
|
||||
if let Some(native_session_id) = query
|
||||
let Some(native_session_id) = query
|
||||
.session_id
|
||||
.as_deref()
|
||||
.filter(|value| !value.trim().is_empty())
|
||||
{
|
||||
let _ = state
|
||||
.writer
|
||||
.link_managed_run_session(run_id, agent, native_session_id)
|
||||
.await?;
|
||||
}
|
||||
let Some(context) = state.reader.managed_run_context(run_id, 256).await? else {
|
||||
warn!(managed_run = %run_id, "managed SessionStart has no active run");
|
||||
else {
|
||||
warn!(managed_run = %run_id, "managed SessionStart has no native session id");
|
||||
return Ok(None);
|
||||
};
|
||||
let Some(cwd) = query
|
||||
.cwd
|
||||
.as_deref()
|
||||
.filter(|value| !value.trim().is_empty())
|
||||
else {
|
||||
warn!(managed_run = %run_id, "managed SessionStart has no checkout cwd");
|
||||
return Ok(None);
|
||||
};
|
||||
let owner_user = crate::workstream::managed_run_owner_stamp(
|
||||
&state.reader,
|
||||
actor,
|
||||
state.trusted_proxy_identity,
|
||||
)
|
||||
.await?;
|
||||
let linked = state
|
||||
.writer
|
||||
.link_or_adopt_managed_run_session(ai_memory_store::LinkOrAdoptManagedRunSession {
|
||||
supplied_run_id: run_id,
|
||||
workspace_id,
|
||||
project_id,
|
||||
cwd: cwd.to_owned(),
|
||||
agent,
|
||||
native_session_id: native_session_id.to_owned(),
|
||||
owner_user,
|
||||
})
|
||||
.await?;
|
||||
let selected_run_id = match linked {
|
||||
ai_memory_store::ManagedRunSessionLink::Exact(run_id) => run_id,
|
||||
ai_memory_store::ManagedRunSessionLink::Adopted(adopted) => {
|
||||
warn!(
|
||||
stale_managed_run = %run_id,
|
||||
managed_run = %adopted,
|
||||
native_session_id,
|
||||
"recovered Codex SessionStart from a stale shared-daemon run id"
|
||||
);
|
||||
adopted
|
||||
}
|
||||
ai_memory_store::ManagedRunSessionLink::Ambiguous => {
|
||||
warn!(managed_run = %run_id, "managed SessionStart recovery is ambiguous");
|
||||
return Ok(None);
|
||||
}
|
||||
ai_memory_store::ManagedRunSessionLink::NoMatch => {
|
||||
warn!(managed_run = %run_id, "managed SessionStart has no live scoped run");
|
||||
return Ok(None);
|
||||
}
|
||||
ai_memory_store::ManagedRunSessionLink::Refused => {
|
||||
warn!(managed_run = %run_id, "managed SessionStart run boundary mismatch");
|
||||
return Ok(None);
|
||||
}
|
||||
};
|
||||
let Some(context) = state
|
||||
.reader
|
||||
.managed_run_context(selected_run_id, 256)
|
||||
.await?
|
||||
else {
|
||||
warn!(managed_run = %selected_run_id, "managed SessionStart has no active run");
|
||||
return Ok(None);
|
||||
};
|
||||
if context.agent != agent {
|
||||
@@ -1652,7 +1691,7 @@ async fn fetch_managed_context(
|
||||
context.sync_after,
|
||||
);
|
||||
Ok(Some(PendingManagedContext {
|
||||
run_id,
|
||||
run_id: selected_run_id,
|
||||
markdown: rendered,
|
||||
}))
|
||||
}
|
||||
@@ -12520,7 +12559,9 @@ mod tests {
|
||||
project_strategy: None,
|
||||
briefing: None,
|
||||
briefing_budget: None,
|
||||
managed_run: Some(run.run_id.to_string()),
|
||||
// Codex's shared app-server may retain the previous finished run
|
||||
// id even though this launcher opened `run` (#987).
|
||||
managed_run: Some(first.run_id.to_string()),
|
||||
session_id: Some("native-2".into()),
|
||||
identity: None,
|
||||
identity_src: None,
|
||||
@@ -12538,6 +12579,13 @@ mod tests {
|
||||
rendered.contains("HANDOFF-MARKER"),
|
||||
"the managed ledger must not swallow the pending handoff: {rendered}"
|
||||
);
|
||||
let linked = state
|
||||
.reader
|
||||
.managed_run_status(run.run_id)
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
assert_eq!(linked.native_session_id.as_deref(), Some("native-2"));
|
||||
assert!(
|
||||
state
|
||||
.reader
|
||||
|
||||
@@ -47,6 +47,22 @@ pub struct WorkstreamState {
|
||||
pub sanitizer: Sanitizer,
|
||||
/// ai-memory data root containing `raw/workstreams`.
|
||||
pub data_dir: PathBuf,
|
||||
/// Whether a trusted proxy can distinguish operators without DB users.
|
||||
pub trusted_proxy_identity: bool,
|
||||
}
|
||||
|
||||
pub(crate) async fn managed_run_owner_stamp(
|
||||
reader: &ReaderPool,
|
||||
identity: Option<&ai_memory_core::IdentityKey>,
|
||||
trusted_proxy_identity: bool,
|
||||
) -> Result<Option<String>, StoreError> {
|
||||
let Some(identity) = identity else {
|
||||
return Ok(None);
|
||||
};
|
||||
let distinguishes = reader
|
||||
.distinguishes_operators(trusted_proxy_identity)
|
||||
.await?;
|
||||
Ok(ai_memory_core::owner_stamp(Some(identity), distinguishes))
|
||||
}
|
||||
|
||||
/// Build the host-wrapper API. It is mounted beside `/hook` and therefore
|
||||
@@ -151,7 +167,8 @@ fn scope_refusal(failure: ScopeResolutionError) -> Response {
|
||||
async fn prepare_run(
|
||||
State(state): State<WorkstreamState>,
|
||||
level: Option<Extension<AuthLevel>>,
|
||||
actor: Option<Extension<ai_memory_core::AuthorizedViewer>>,
|
||||
actor: Option<Extension<ai_memory_core::ActorContext>>,
|
||||
viewer: Option<Extension<ai_memory_core::AuthorizedViewer>>,
|
||||
Json(request): Json<PrepareManagedRunRequest>,
|
||||
) -> Response {
|
||||
if let Err(response) = authorize(level, Capability::NormalWrite) {
|
||||
@@ -239,7 +256,7 @@ async fn prepare_run(
|
||||
&state.writer,
|
||||
request.workspace.trim(),
|
||||
request.project.trim(),
|
||||
actor_user(actor),
|
||||
actor_user(viewer),
|
||||
)
|
||||
.await
|
||||
{
|
||||
@@ -260,20 +277,34 @@ async fn prepare_run(
|
||||
);
|
||||
}
|
||||
};
|
||||
let identity = actor.and_then(|Extension(actor)| actor.identity_key());
|
||||
let owner_user = match managed_run_owner_stamp(
|
||||
&state.reader,
|
||||
identity.as_ref(),
|
||||
state.trusted_proxy_identity,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(owner) => owner,
|
||||
Err(failure) => return error(StatusCode::INTERNAL_SERVER_ERROR, failure.to_string()),
|
||||
};
|
||||
let prepared = state
|
||||
.writer
|
||||
.prepare_workstream_run(PrepareWorkstreamRun {
|
||||
workspace_id: scope.workspace_id,
|
||||
project_id: scope.project_id,
|
||||
repo_fingerprint: request.repo_fingerprint,
|
||||
worktree_fingerprint: request.worktree_fingerprint,
|
||||
cwd: request.cwd,
|
||||
agent: request.agent,
|
||||
automatic_harness: request.automatic_harness,
|
||||
available_agents: request.available_agents,
|
||||
selection,
|
||||
lease_owner: request.lease_owner,
|
||||
})
|
||||
.prepare_workstream_run_owned(
|
||||
PrepareWorkstreamRun {
|
||||
workspace_id: scope.workspace_id,
|
||||
project_id: scope.project_id,
|
||||
repo_fingerprint: request.repo_fingerprint,
|
||||
worktree_fingerprint: request.worktree_fingerprint,
|
||||
cwd: request.cwd,
|
||||
agent: request.agent,
|
||||
automatic_harness: request.automatic_harness,
|
||||
available_agents: request.available_agents,
|
||||
selection,
|
||||
lease_owner: request.lease_owner,
|
||||
},
|
||||
owner_user,
|
||||
)
|
||||
.await;
|
||||
match prepared {
|
||||
Ok(prepared) => Json(PrepareManagedRunResponse {
|
||||
@@ -1093,6 +1124,7 @@ mod tests {
|
||||
reader: store.reader.clone(),
|
||||
sanitizer: Sanitizer::default(),
|
||||
data_dir: data_dir.to_path_buf(),
|
||||
trusted_proxy_identity: false,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1395,6 +1427,7 @@ mod tests {
|
||||
State(state),
|
||||
None,
|
||||
None,
|
||||
None,
|
||||
Json(PrepareManagedRunRequest {
|
||||
workspace: "default".into(),
|
||||
project: "managed".into(),
|
||||
@@ -1431,6 +1464,7 @@ mod tests {
|
||||
State(state.clone()),
|
||||
None,
|
||||
None,
|
||||
None,
|
||||
Json(PrepareManagedRunRequest {
|
||||
workspace: "default".into(),
|
||||
project: "managed".into(),
|
||||
@@ -1469,6 +1503,7 @@ mod tests {
|
||||
State(state),
|
||||
None,
|
||||
None,
|
||||
None,
|
||||
Json(PrepareManagedRunRequest {
|
||||
workspace: "default".into(),
|
||||
project: "managed".into(),
|
||||
@@ -1487,6 +1522,97 @@ mod tests {
|
||||
assert_eq!(automatic.status(), StatusCode::OK);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn prepare_stamps_the_operator_bucket_used_by_session_start_recovery() {
|
||||
let temp = TempDir::new().unwrap();
|
||||
let store = Store::open(temp.path()).unwrap();
|
||||
let state = test_state(&store, temp.path());
|
||||
store
|
||||
.writer
|
||||
.create_human_user(
|
||||
ai_memory_core::NewUser {
|
||||
username: "alice".into(),
|
||||
name: None,
|
||||
email: None,
|
||||
},
|
||||
ai_memory_core::UserRole::User,
|
||||
None,
|
||||
false,
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let response = prepare_run(
|
||||
State(state.clone()),
|
||||
None,
|
||||
Some(Extension(ai_memory_core::ActorContext {
|
||||
user: Some("alice".into()),
|
||||
..ai_memory_core::ActorContext::default()
|
||||
})),
|
||||
None,
|
||||
Json(PrepareManagedRunRequest {
|
||||
workspace: "default".into(),
|
||||
project: "managed-owner".into(),
|
||||
cwd: "/repo".into(),
|
||||
repo_fingerprint: "repo".into(),
|
||||
worktree_fingerprint: "worktree".into(),
|
||||
agent: AgentKind::Codex,
|
||||
automatic_harness: false,
|
||||
available_agents: Vec::new(),
|
||||
workstream: None,
|
||||
new_workstream: None,
|
||||
lease_owner: "alice-launcher".into(),
|
||||
}),
|
||||
)
|
||||
.await;
|
||||
assert_eq!(response.status(), StatusCode::OK);
|
||||
let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
|
||||
let prepared: PrepareManagedRunResponse = serde_json::from_slice(&body).unwrap();
|
||||
let (workspace_id, project_id) = store
|
||||
.reader
|
||||
.managed_run_scope(prepared.run_id)
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(
|
||||
store
|
||||
.writer
|
||||
.link_or_adopt_managed_run_session(ai_memory_store::LinkOrAdoptManagedRunSession {
|
||||
supplied_run_id: prepared.run_id,
|
||||
workspace_id,
|
||||
project_id,
|
||||
cwd: "/repo".into(),
|
||||
agent: AgentKind::Codex,
|
||||
native_session_id: "native-bob".into(),
|
||||
owner_user: Some(
|
||||
ai_memory_core::IdentityKey::User("bob".into()).storage_key(),
|
||||
),
|
||||
})
|
||||
.await
|
||||
.unwrap(),
|
||||
ai_memory_store::ManagedRunSessionLink::Refused
|
||||
);
|
||||
assert_eq!(
|
||||
store
|
||||
.writer
|
||||
.link_or_adopt_managed_run_session(ai_memory_store::LinkOrAdoptManagedRunSession {
|
||||
supplied_run_id: prepared.run_id,
|
||||
workspace_id,
|
||||
project_id,
|
||||
cwd: "/repo".into(),
|
||||
agent: AgentKind::Codex,
|
||||
native_session_id: "native-alice".into(),
|
||||
owner_user: Some(
|
||||
ai_memory_core::IdentityKey::User("alice".into()).storage_key(),
|
||||
),
|
||||
})
|
||||
.await
|
||||
.unwrap(),
|
||||
ai_memory_store::ManagedRunSessionLink::Exact(prepared.run_id)
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn kiro_is_accepted_as_an_explicit_and_automatic_harness() {
|
||||
let temp = TempDir::new().unwrap();
|
||||
@@ -1497,6 +1623,7 @@ mod tests {
|
||||
State(state.clone()),
|
||||
None,
|
||||
None,
|
||||
None,
|
||||
Json(PrepareManagedRunRequest {
|
||||
workspace: "default".into(),
|
||||
project: "managed".into(),
|
||||
@@ -1534,6 +1661,7 @@ mod tests {
|
||||
State(state),
|
||||
None,
|
||||
None,
|
||||
None,
|
||||
Json(PrepareManagedRunRequest {
|
||||
workspace: "default".into(),
|
||||
project: "managed".into(),
|
||||
@@ -1562,6 +1690,7 @@ mod tests {
|
||||
State(state.clone()),
|
||||
None,
|
||||
None,
|
||||
None,
|
||||
Json(PrepareManagedRunRequest {
|
||||
workspace: "default".into(),
|
||||
project: "managed".into(),
|
||||
@@ -1599,6 +1728,7 @@ mod tests {
|
||||
State(state),
|
||||
None,
|
||||
None,
|
||||
None,
|
||||
Json(PrepareManagedRunRequest {
|
||||
workspace: "default".into(),
|
||||
project: "managed".into(),
|
||||
@@ -1627,6 +1757,7 @@ mod tests {
|
||||
State(state.clone()),
|
||||
None,
|
||||
None,
|
||||
None,
|
||||
Json(PrepareManagedRunRequest {
|
||||
workspace: "default".into(),
|
||||
project: "managed".into(),
|
||||
@@ -1664,6 +1795,7 @@ mod tests {
|
||||
State(state),
|
||||
None,
|
||||
None,
|
||||
None,
|
||||
Json(PrepareManagedRunRequest {
|
||||
workspace: "default".into(),
|
||||
project: "managed".into(),
|
||||
@@ -1695,6 +1827,7 @@ mod tests {
|
||||
State(state.clone()),
|
||||
None,
|
||||
None,
|
||||
None,
|
||||
Json(PrepareManagedRunRequest {
|
||||
workspace: "default".into(),
|
||||
project: "managed".into(),
|
||||
@@ -1732,6 +1865,7 @@ mod tests {
|
||||
State(state),
|
||||
None,
|
||||
None,
|
||||
None,
|
||||
Json(PrepareManagedRunRequest {
|
||||
workspace: "default".into(),
|
||||
project: "managed".into(),
|
||||
@@ -1998,6 +2132,7 @@ mod tests {
|
||||
State(state.clone()),
|
||||
None,
|
||||
None,
|
||||
None,
|
||||
Json(PrepareManagedRunRequest {
|
||||
workspace: "default".into(),
|
||||
project: "alice-client-work".into(),
|
||||
|
||||
@@ -0,0 +1,10 @@
|
||||
-- Managed-run recovery may replace a stale harness-exported run id with the
|
||||
-- one live launch in the current repository. Keep that recovery in the same
|
||||
-- operator bucket as the launch so one user cannot bind another user's run.
|
||||
-- NULL retains the historical single-operator/shared-server behavior.
|
||||
|
||||
ALTER TABLE managed_runs ADD COLUMN owner_user TEXT;
|
||||
|
||||
CREATE INDEX idx_managed_runs_adoption
|
||||
ON managed_runs(state, agent_kind, owner_user, context_delivered,
|
||||
native_session_linked_at, lease_expires_at);
|
||||
@@ -98,9 +98,10 @@ pub use users::{
|
||||
};
|
||||
pub use web_sessions::{LiveWebSession, WebSession, hash_session_secret};
|
||||
pub use workstream::{
|
||||
FinishWorkstreamRun, FinishedWorkstreamRun, ManagedRunContext, PrepareWorkstreamRun,
|
||||
PreparedWorkstreamRun, RenameWorkstream, RenamedWorkstream, StoredManagedRunStatus,
|
||||
StoredWorkstreamSummary, WorkstreamSelection, WorkstreamSelector,
|
||||
FinishWorkstreamRun, FinishedWorkstreamRun, LinkOrAdoptManagedRunSession, ManagedRunContext,
|
||||
ManagedRunSessionLink, PrepareWorkstreamRun, PreparedWorkstreamRun, RenameWorkstream,
|
||||
RenamedWorkstream, StoredManagedRunStatus, StoredWorkstreamSummary, WorkstreamSelection,
|
||||
WorkstreamSelector,
|
||||
};
|
||||
pub use writer::{StartupContextAcceptance, WriterHandle};
|
||||
|
||||
|
||||
@@ -12544,6 +12544,7 @@ pub(crate) mod tests {
|
||||
selection: crate::WorkstreamSelection::Current,
|
||||
lease_owner: "test".into(),
|
||||
},
|
||||
None,
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
@@ -14455,6 +14456,7 @@ pub(crate) mod tests {
|
||||
selection: crate::workstream::WorkstreamSelection::Current,
|
||||
lease_owner: "test".to_string(),
|
||||
},
|
||||
None,
|
||||
)
|
||||
.expect("opening a managed run should succeed")
|
||||
}
|
||||
|
||||
@@ -79,6 +79,40 @@ pub struct PreparedWorkstreamRun {
|
||||
pub may_adopt_existing_session: bool,
|
||||
}
|
||||
|
||||
/// Result of binding a SessionStart to an exact or recovered managed run.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
pub enum ManagedRunSessionLink {
|
||||
/// The supplied run was active and accepted the native session.
|
||||
Exact(ManagedRunId),
|
||||
/// A stale supplied id was replaced by the sole safe current candidate.
|
||||
Adopted(ManagedRunId),
|
||||
/// No live run matched the current repository and operator.
|
||||
NoMatch,
|
||||
/// More than one run matched, so choosing one would be unsafe.
|
||||
Ambiguous,
|
||||
/// The supplied active run disagreed with the request boundary.
|
||||
Refused,
|
||||
}
|
||||
|
||||
/// Exact boundary used to link or recover one managed SessionStart.
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct LinkOrAdoptManagedRunSession {
|
||||
/// Run id supplied by the child hook; it may be terminal for Codex.
|
||||
pub supplied_run_id: ManagedRunId,
|
||||
/// Already-resolved workspace receiving the SessionStart.
|
||||
pub workspace_id: WorkspaceId,
|
||||
/// Already-resolved project receiving the SessionStart.
|
||||
pub project_id: ProjectId,
|
||||
/// Canonical checkout cwd recorded by both launcher and hook.
|
||||
pub cwd: String,
|
||||
/// Harness reporting the native session.
|
||||
pub agent: AgentKind,
|
||||
/// Harness-native session id to bind exactly once.
|
||||
pub native_session_id: String,
|
||||
/// Topology-aware qualified operator identity, or shared `None`.
|
||||
pub owner_user: Option<String>,
|
||||
}
|
||||
|
||||
/// Store-level finish input after the raw segment has been made durable.
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct FinishWorkstreamRun {
|
||||
@@ -223,15 +257,33 @@ struct LinkRunRow {
|
||||
workstream: Vec<u8>,
|
||||
agent: String,
|
||||
native_session: Option<String>,
|
||||
native_session_linked: bool,
|
||||
sync_through: i64,
|
||||
context_delivered: bool,
|
||||
}
|
||||
|
||||
struct ManagedRunBoundaryRow {
|
||||
state: String,
|
||||
workspace: Vec<u8>,
|
||||
project: Vec<u8>,
|
||||
cwd: String,
|
||||
agent: String,
|
||||
owner_user: Option<String>,
|
||||
}
|
||||
|
||||
/// Atomically select a workstream, expire stale leases, and open one run.
|
||||
pub(crate) fn prepare_run(
|
||||
conn: &mut Connection,
|
||||
input: &PrepareWorkstreamRun,
|
||||
owner_user: Option<&str>,
|
||||
) -> StoreResult<PreparedWorkstreamRun> {
|
||||
if owner_user
|
||||
.is_some_and(|owner| ai_memory_core::IdentityKey::from_storage_key(owner).is_none())
|
||||
{
|
||||
return Err(StoreError::InvalidState(
|
||||
"invalid managed-run owner identity".into(),
|
||||
));
|
||||
}
|
||||
let now = Timestamp::now().as_microsecond();
|
||||
let tx = conn.transaction()?;
|
||||
tx.execute(
|
||||
@@ -295,15 +347,16 @@ pub(crate) fn prepare_run(
|
||||
let run_id = ManagedRunId::new();
|
||||
tx.execute(
|
||||
"INSERT INTO managed_runs( \
|
||||
id, workstream_id, agent_kind, lease_owner, native_session_id, state, \
|
||||
id, workstream_id, agent_kind, lease_owner, native_session_id, owner_user, state, \
|
||||
sync_after, sync_through, context_delivered, lease_expires_at, started_at \
|
||||
) VALUES (?1, ?2, ?3, ?4, ?5, 'active', ?6, ?7, 0, ?8, ?9)",
|
||||
) VALUES (?1, ?2, ?3, ?4, ?5, ?6, 'active', ?7, ?8, 0, ?9, ?10)",
|
||||
params![
|
||||
run_id.as_bytes(),
|
||||
workstream_id.as_bytes(),
|
||||
agent.as_str(),
|
||||
input.lease_owner,
|
||||
native_session_id,
|
||||
owner_user,
|
||||
sync_after,
|
||||
latest_sequence,
|
||||
now + LEASE_MICROS,
|
||||
@@ -513,9 +566,22 @@ pub(crate) fn link_native_session(
|
||||
}
|
||||
let now = Timestamp::now().as_microsecond();
|
||||
let tx = conn.transaction()?;
|
||||
let linked = link_native_session_in_transaction(&tx, run_id, agent, native_session_id, now)?;
|
||||
tx.commit()?;
|
||||
Ok(linked)
|
||||
}
|
||||
|
||||
fn link_native_session_in_transaction(
|
||||
tx: &Transaction<'_>,
|
||||
run_id: ManagedRunId,
|
||||
agent: AgentKind,
|
||||
native_session_id: &str,
|
||||
now: i64,
|
||||
) -> StoreResult<bool> {
|
||||
let run: Option<LinkRunRow> = tx
|
||||
.query_row(
|
||||
"SELECT workstream_id, agent_kind, native_session_id, sync_through, context_delivered \
|
||||
"SELECT workstream_id, agent_kind, native_session_id, \
|
||||
native_session_linked_at IS NOT NULL, sync_through, context_delivered \
|
||||
FROM managed_runs WHERE id = ?1 AND state = 'active'",
|
||||
params![run_id.as_bytes()],
|
||||
|row| {
|
||||
@@ -523,8 +589,9 @@ pub(crate) fn link_native_session(
|
||||
workstream: row.get(0)?,
|
||||
agent: row.get(1)?,
|
||||
native_session: row.get(2)?,
|
||||
sync_through: row.get(3)?,
|
||||
context_delivered: row.get(4)?,
|
||||
native_session_linked: row.get(3)?,
|
||||
sync_through: row.get(4)?,
|
||||
context_delivered: row.get(5)?,
|
||||
})
|
||||
},
|
||||
)
|
||||
@@ -535,7 +602,7 @@ pub(crate) fn link_native_session(
|
||||
if run.agent != agent.as_str() {
|
||||
return Ok(false);
|
||||
}
|
||||
if run.context_delivered
|
||||
if (run.context_delivered || run.native_session_linked)
|
||||
&& run
|
||||
.native_session
|
||||
.as_deref()
|
||||
@@ -582,10 +649,122 @@ pub(crate) fn link_native_session(
|
||||
native_session_linked_at = ?3 WHERE id = ?4",
|
||||
params![native_session_id, initial_delivery, now, run_id.as_bytes()],
|
||||
)?;
|
||||
tx.commit()?;
|
||||
Ok(true)
|
||||
}
|
||||
|
||||
/// Bind a SessionStart to its exact run, or recover one stale Codex daemon id.
|
||||
///
|
||||
/// Candidate selection and linking share one writer transaction. Recovery is
|
||||
/// deliberately narrow: only a missing/terminal Codex id can be replaced, and
|
||||
/// only by the sole live, undelivered, never-linked run in the same repository
|
||||
/// and operator bucket.
|
||||
pub(crate) fn link_or_adopt_native_session(
|
||||
conn: &mut Connection,
|
||||
input: &LinkOrAdoptManagedRunSession,
|
||||
) -> StoreResult<ManagedRunSessionLink> {
|
||||
if input.native_session_id.trim().is_empty()
|
||||
|| input.cwd.trim().is_empty()
|
||||
|| input
|
||||
.owner_user
|
||||
.as_deref()
|
||||
.is_some_and(|owner| ai_memory_core::IdentityKey::from_storage_key(owner).is_none())
|
||||
{
|
||||
return Ok(ManagedRunSessionLink::Refused);
|
||||
}
|
||||
|
||||
let now = Timestamp::now().as_microsecond();
|
||||
let tx = conn.transaction()?;
|
||||
let supplied: Option<ManagedRunBoundaryRow> = tx
|
||||
.query_row(
|
||||
"SELECT mr.state, w.workspace_id, w.project_id, w.cwd, mr.agent_kind, mr.owner_user \
|
||||
FROM managed_runs mr JOIN workstreams w ON w.id = mr.workstream_id \
|
||||
WHERE mr.id = ?1",
|
||||
params![input.supplied_run_id.as_bytes()],
|
||||
|row| {
|
||||
Ok(ManagedRunBoundaryRow {
|
||||
state: row.get(0)?,
|
||||
workspace: row.get(1)?,
|
||||
project: row.get(2)?,
|
||||
cwd: row.get(3)?,
|
||||
agent: row.get(4)?,
|
||||
owner_user: row.get(5)?,
|
||||
})
|
||||
},
|
||||
)
|
||||
.optional()?;
|
||||
|
||||
let selected = match supplied {
|
||||
Some(run) if run.state == "active" => {
|
||||
let boundary_matches = run.workspace.as_slice() == input.workspace_id.as_bytes()
|
||||
&& run.project.as_slice() == input.project_id.as_bytes()
|
||||
&& run.cwd == input.cwd
|
||||
&& run.agent == input.agent.as_str()
|
||||
&& run.owner_user == input.owner_user;
|
||||
if !boundary_matches {
|
||||
tx.commit()?;
|
||||
return Ok(ManagedRunSessionLink::Refused);
|
||||
}
|
||||
(input.supplied_run_id, false)
|
||||
}
|
||||
_ if input.agent == AgentKind::Codex => {
|
||||
let mut statement = tx.prepare(
|
||||
"SELECT mr.id FROM managed_runs mr \
|
||||
JOIN workstreams w ON w.id = mr.workstream_id \
|
||||
WHERE w.workspace_id = ?1 AND w.project_id = ?2 AND w.cwd = ?3 \
|
||||
AND mr.agent_kind = ?4 AND mr.owner_user IS ?5 \
|
||||
AND mr.state = 'active' AND mr.lease_expires_at > ?6 \
|
||||
AND mr.context_delivered = 0 AND mr.native_session_linked_at IS NULL \
|
||||
ORDER BY mr.started_at DESC, mr.id DESC LIMIT 2",
|
||||
)?;
|
||||
let rows = statement.query_map(
|
||||
params![
|
||||
input.workspace_id.as_bytes(),
|
||||
input.project_id.as_bytes(),
|
||||
input.cwd.as_str(),
|
||||
input.agent.as_str(),
|
||||
input.owner_user.as_deref(),
|
||||
now,
|
||||
],
|
||||
|row| row.get::<_, Vec<u8>>(0),
|
||||
)?;
|
||||
let candidates = rows.collect::<Result<Vec<_>, _>>()?;
|
||||
drop(statement);
|
||||
match candidates.as_slice() {
|
||||
[only] => (ManagedRunId::from_slice(only)?, true),
|
||||
[] => {
|
||||
tx.commit()?;
|
||||
return Ok(ManagedRunSessionLink::NoMatch);
|
||||
}
|
||||
_ => {
|
||||
tx.commit()?;
|
||||
return Ok(ManagedRunSessionLink::Ambiguous);
|
||||
}
|
||||
}
|
||||
}
|
||||
_ => {
|
||||
tx.commit()?;
|
||||
return Ok(ManagedRunSessionLink::NoMatch);
|
||||
}
|
||||
};
|
||||
|
||||
if !link_native_session_in_transaction(
|
||||
&tx,
|
||||
selected.0,
|
||||
input.agent,
|
||||
&input.native_session_id,
|
||||
now,
|
||||
)? {
|
||||
tx.commit()?;
|
||||
return Ok(ManagedRunSessionLink::Refused);
|
||||
}
|
||||
tx.commit()?;
|
||||
Ok(if selected.1 {
|
||||
ManagedRunSessionLink::Adopted(selected.0)
|
||||
} else {
|
||||
ManagedRunSessionLink::Exact(selected.0)
|
||||
})
|
||||
}
|
||||
|
||||
/// Mark the assigned synchronization packet delivered to SessionStart.
|
||||
pub(crate) fn accept_context(conn: &mut Connection, run_id: ManagedRunId) -> StoreResult<bool> {
|
||||
let tx = conn.transaction()?;
|
||||
|
||||
@@ -36,8 +36,9 @@ use crate::session_consolidation::SessionConsolidationJob;
|
||||
use crate::users::{self, TOKEN_HASH_LEN};
|
||||
use crate::web_sessions::{self, WebSession};
|
||||
use crate::workstream::{
|
||||
FinishWorkstreamRun, FinishedWorkstreamRun, PrepareWorkstreamRun, PreparedWorkstreamRun,
|
||||
RenameWorkstream, RenamedWorkstream,
|
||||
FinishWorkstreamRun, FinishedWorkstreamRun, LinkOrAdoptManagedRunSession,
|
||||
ManagedRunSessionLink, PrepareWorkstreamRun, PreparedWorkstreamRun, RenameWorkstream,
|
||||
RenamedWorkstream,
|
||||
};
|
||||
|
||||
/// Result of atomically claiming the startup context assembled for one hook.
|
||||
@@ -713,6 +714,7 @@ pub(crate) enum WriteCmd {
|
||||
},
|
||||
PrepareWorkstreamRun {
|
||||
input: PrepareWorkstreamRun,
|
||||
owner_user: Option<String>,
|
||||
reply: oneshot::Sender<StoreResult<PreparedWorkstreamRun>>,
|
||||
},
|
||||
HeartbeatManagedRun {
|
||||
@@ -729,6 +731,10 @@ pub(crate) enum WriteCmd {
|
||||
native_session_id: String,
|
||||
reply: oneshot::Sender<StoreResult<bool>>,
|
||||
},
|
||||
LinkOrAdoptManagedRunSession {
|
||||
input: LinkOrAdoptManagedRunSession,
|
||||
reply: oneshot::Sender<StoreResult<ManagedRunSessionLink>>,
|
||||
},
|
||||
AcceptManagedRunContext {
|
||||
run_id: ManagedRunId,
|
||||
reply: oneshot::Sender<StoreResult<bool>>,
|
||||
@@ -2998,10 +3004,23 @@ impl WriterHandle {
|
||||
pub async fn prepare_workstream_run(
|
||||
&self,
|
||||
input: PrepareWorkstreamRun,
|
||||
) -> StoreResult<PreparedWorkstreamRun> {
|
||||
self.prepare_workstream_run_owned(input, None).await
|
||||
}
|
||||
|
||||
/// Select a workstream and stamp the operator bucket used by safe recovery.
|
||||
pub async fn prepare_workstream_run_owned(
|
||||
&self,
|
||||
input: PrepareWorkstreamRun,
|
||||
owner_user: Option<String>,
|
||||
) -> StoreResult<PreparedWorkstreamRun> {
|
||||
let (tx, rx) = oneshot::channel();
|
||||
self.send(WriteCmd::PrepareWorkstreamRun { input, reply: tx })
|
||||
.await?;
|
||||
self.send(WriteCmd::PrepareWorkstreamRun {
|
||||
input,
|
||||
owner_user,
|
||||
reply: tx,
|
||||
})
|
||||
.await?;
|
||||
rx.await.map_err(|_| StoreError::WriterClosed)?
|
||||
}
|
||||
|
||||
@@ -3039,6 +3058,17 @@ impl WriterHandle {
|
||||
rx.await.map_err(|_| StoreError::WriterClosed)?
|
||||
}
|
||||
|
||||
/// Link an exact run or safely recover a stale Codex daemon run id.
|
||||
pub async fn link_or_adopt_managed_run_session(
|
||||
&self,
|
||||
input: LinkOrAdoptManagedRunSession,
|
||||
) -> StoreResult<ManagedRunSessionLink> {
|
||||
let (tx, rx) = oneshot::channel();
|
||||
self.send(WriteCmd::LinkOrAdoptManagedRunSession { input, reply: tx })
|
||||
.await?;
|
||||
rx.await.map_err(|_| StoreError::WriterClosed)?
|
||||
}
|
||||
|
||||
/// Acknowledge successful SessionStart delivery for a managed run.
|
||||
pub async fn accept_managed_run_context(&self, run_id: ManagedRunId) -> StoreResult<bool> {
|
||||
let (tx, rx) = oneshot::channel();
|
||||
@@ -4236,8 +4266,13 @@ fn worker_loop(mut conn: Connection, mut rx: mpsc::Receiver<WriteCmd>) {
|
||||
let result = crate::maintenance::record_success(&conn, job);
|
||||
send_or_warn(reply, result, "record_maintenance_job_success");
|
||||
}
|
||||
WriteCmd::PrepareWorkstreamRun { input, reply } => {
|
||||
let result = crate::workstream::prepare_run(&mut conn, &input);
|
||||
WriteCmd::PrepareWorkstreamRun {
|
||||
input,
|
||||
owner_user,
|
||||
reply,
|
||||
} => {
|
||||
let result =
|
||||
crate::workstream::prepare_run(&mut conn, &input, owner_user.as_deref());
|
||||
send_or_warn(reply, result, "prepare_workstream_run");
|
||||
}
|
||||
WriteCmd::HeartbeatManagedRun { run_id, reply } => {
|
||||
@@ -4262,6 +4297,10 @@ fn worker_loop(mut conn: Connection, mut rx: mpsc::Receiver<WriteCmd>) {
|
||||
);
|
||||
send_or_warn(reply, result, "link_managed_run_session");
|
||||
}
|
||||
WriteCmd::LinkOrAdoptManagedRunSession { input, reply } => {
|
||||
let result = crate::workstream::link_or_adopt_native_session(&mut conn, &input);
|
||||
send_or_warn(reply, result, "link_or_adopt_managed_run_session");
|
||||
}
|
||||
WriteCmd::AcceptManagedRunContext { run_id, reply } => {
|
||||
let result = crate::workstream::accept_context(&mut conn, run_id);
|
||||
send_or_warn(reply, result, "accept_managed_run_context");
|
||||
|
||||
@@ -21,7 +21,10 @@ use ai_memory_core::{
|
||||
NewSession, NewUser, OwnerFilter, PagePath, ProjectId, SessionId, Tier, UserRole, WorkspaceId,
|
||||
owner_stamp,
|
||||
};
|
||||
use ai_memory_store::{PrepareWorkstreamRun, Store, WorkstreamSelection};
|
||||
use ai_memory_store::{
|
||||
LinkOrAdoptManagedRunSession, ManagedRunSessionLink, PrepareWorkstreamRun, Store,
|
||||
WorkstreamSelection,
|
||||
};
|
||||
|
||||
fn operator(name: &str) -> String {
|
||||
IdentityKey::User(name.into()).storage_key()
|
||||
@@ -803,3 +806,218 @@ async fn a_session_linked_by_one_managed_run_is_not_another_runs() {
|
||||
(Some("native-beta".into()), true)
|
||||
);
|
||||
}
|
||||
|
||||
/// Stale-run recovery is a scoped, owned, single-candidate operation. The
|
||||
/// controls prove that the same candidate is usable once every boundary
|
||||
/// matches, while foreign owners/projects and a second linker cannot take it.
|
||||
#[tokio::test]
|
||||
async fn stale_codex_run_recovery_cannot_cross_project_owner_or_session_boundaries() {
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
let store = Store::open(tmp.path()).unwrap();
|
||||
let (ws, proj) = scope(&store).await;
|
||||
let other_proj = store
|
||||
.writer
|
||||
.get_or_create_project(ws, "other-app", None)
|
||||
.await
|
||||
.unwrap();
|
||||
let alice = operator("alice");
|
||||
let bob = operator("bob");
|
||||
let prepare = |project_id, name: &str| PrepareWorkstreamRun {
|
||||
workspace_id: ws,
|
||||
project_id,
|
||||
repo_fingerprint: format!("repo-{project_id}"),
|
||||
worktree_fingerprint: format!("worktree-{project_id}"),
|
||||
cwd: format!("/repo/{project_id}"),
|
||||
agent: AgentKind::Codex,
|
||||
automatic_harness: false,
|
||||
available_agents: Vec::new(),
|
||||
selection: WorkstreamSelection::New(name.into()),
|
||||
lease_owner: format!("launcher-{name}"),
|
||||
};
|
||||
|
||||
let stale = store
|
||||
.writer
|
||||
.prepare_workstream_run_owned(prepare(proj, "stale"), Some(alice.clone()))
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(store.writer.cancel_managed_run(stale.run_id).await.unwrap());
|
||||
let current = store
|
||||
.writer
|
||||
.prepare_workstream_run_owned(prepare(proj, "current"), Some(alice.clone()))
|
||||
.await
|
||||
.unwrap();
|
||||
let foreign = store
|
||||
.writer
|
||||
.prepare_workstream_run_owned(prepare(other_proj, "foreign"), Some(alice.clone()))
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let recover = |project_id, owner: Option<String>, native: &'static str| {
|
||||
store
|
||||
.writer
|
||||
.link_or_adopt_managed_run_session(LinkOrAdoptManagedRunSession {
|
||||
supplied_run_id: stale.run_id,
|
||||
workspace_id: ws,
|
||||
project_id,
|
||||
cwd: format!("/repo/{project_id}"),
|
||||
agent: AgentKind::Codex,
|
||||
native_session_id: native.into(),
|
||||
owner_user: owner,
|
||||
})
|
||||
};
|
||||
assert_eq!(
|
||||
recover(proj, Some(bob), "native-bob").await.unwrap(),
|
||||
ManagedRunSessionLink::NoMatch,
|
||||
"a second operator must not see Alice's candidate"
|
||||
);
|
||||
assert_eq!(
|
||||
store
|
||||
.writer
|
||||
.link_or_adopt_managed_run_session(LinkOrAdoptManagedRunSession {
|
||||
supplied_run_id: stale.run_id,
|
||||
workspace_id: ws,
|
||||
project_id: proj,
|
||||
cwd: "/repo/a-different-worktree".into(),
|
||||
agent: AgentKind::Codex,
|
||||
native_session_id: "native-other-worktree".into(),
|
||||
owner_user: Some(alice.clone()),
|
||||
})
|
||||
.await
|
||||
.unwrap(),
|
||||
ManagedRunSessionLink::NoMatch,
|
||||
"a run in another checkout cwd must not be adopted"
|
||||
);
|
||||
assert_eq!(
|
||||
recover(other_proj, Some(alice.clone()), "native-foreign")
|
||||
.await
|
||||
.unwrap(),
|
||||
ManagedRunSessionLink::Adopted(foreign.run_id),
|
||||
"the same owner may recover only the candidate in the named project"
|
||||
);
|
||||
assert_eq!(
|
||||
recover(proj, Some(alice.clone()), "native-alice")
|
||||
.await
|
||||
.unwrap(),
|
||||
ManagedRunSessionLink::Adopted(current.run_id)
|
||||
);
|
||||
assert_eq!(
|
||||
recover(proj, Some(alice), "native-thief").await.unwrap(),
|
||||
ManagedRunSessionLink::NoMatch,
|
||||
"a linked run cannot be rebound by a racing SessionStart"
|
||||
);
|
||||
let current_status = store
|
||||
.reader
|
||||
.managed_run_status(current.run_id)
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
assert_eq!(
|
||||
current_status.native_session_id.as_deref(),
|
||||
Some("native-alice")
|
||||
);
|
||||
}
|
||||
|
||||
/// Two launches in one repository are a real possibility when the operator
|
||||
/// uses separate named workstreams. Recovery must refuse to guess between
|
||||
/// them, and an active foreign id must not be treated as a stale trigger.
|
||||
#[tokio::test]
|
||||
async fn stale_codex_run_recovery_fails_closed_on_ambiguity_and_active_mismatch() {
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
let store = Store::open(tmp.path()).unwrap();
|
||||
let (ws, proj) = scope(&store).await;
|
||||
let other_proj = store
|
||||
.writer
|
||||
.get_or_create_project(ws, "other-app", None)
|
||||
.await
|
||||
.unwrap();
|
||||
let owner = operator("alice");
|
||||
let prepare = |project_id, name: &str| PrepareWorkstreamRun {
|
||||
workspace_id: ws,
|
||||
project_id,
|
||||
repo_fingerprint: format!("repo-{project_id}"),
|
||||
worktree_fingerprint: format!("worktree-{project_id}"),
|
||||
cwd: format!("/repo/{project_id}"),
|
||||
agent: AgentKind::Codex,
|
||||
automatic_harness: false,
|
||||
available_agents: Vec::new(),
|
||||
selection: WorkstreamSelection::New(name.into()),
|
||||
lease_owner: format!("launcher-{name}"),
|
||||
};
|
||||
let alpha = store
|
||||
.writer
|
||||
.prepare_workstream_run_owned(prepare(proj, "alpha"), Some(owner.clone()))
|
||||
.await
|
||||
.unwrap();
|
||||
let beta = store
|
||||
.writer
|
||||
.prepare_workstream_run_owned(prepare(proj, "beta"), Some(owner.clone()))
|
||||
.await
|
||||
.unwrap();
|
||||
let foreign = store
|
||||
.writer
|
||||
.prepare_workstream_run_owned(prepare(other_proj, "foreign"), Some(owner.clone()))
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(
|
||||
store
|
||||
.writer
|
||||
.link_or_adopt_managed_run_session(LinkOrAdoptManagedRunSession {
|
||||
supplied_run_id: ai_memory_core::ManagedRunId::new(),
|
||||
workspace_id: ws,
|
||||
project_id: proj,
|
||||
cwd: format!("/repo/{proj}"),
|
||||
agent: AgentKind::Codex,
|
||||
native_session_id: "native-ambiguous".into(),
|
||||
owner_user: Some(owner.clone()),
|
||||
})
|
||||
.await
|
||||
.unwrap(),
|
||||
ManagedRunSessionLink::Ambiguous
|
||||
);
|
||||
assert_eq!(
|
||||
store
|
||||
.writer
|
||||
.link_or_adopt_managed_run_session(LinkOrAdoptManagedRunSession {
|
||||
supplied_run_id: alpha.run_id,
|
||||
workspace_id: ws,
|
||||
project_id: other_proj,
|
||||
cwd: format!("/repo/{other_proj}"),
|
||||
agent: AgentKind::Codex,
|
||||
native_session_id: "native-cross-project".into(),
|
||||
owner_user: Some(owner.clone()),
|
||||
})
|
||||
.await
|
||||
.unwrap(),
|
||||
ManagedRunSessionLink::Refused,
|
||||
"an active run with a mismatched boundary must not trigger adoption"
|
||||
);
|
||||
assert_eq!(
|
||||
store
|
||||
.writer
|
||||
.link_or_adopt_managed_run_session(LinkOrAdoptManagedRunSession {
|
||||
supplied_run_id: foreign.run_id,
|
||||
workspace_id: ws,
|
||||
project_id: other_proj,
|
||||
cwd: format!("/repo/{other_proj}"),
|
||||
agent: AgentKind::Codex,
|
||||
native_session_id: "native-control".into(),
|
||||
owner_user: Some(owner),
|
||||
})
|
||||
.await
|
||||
.unwrap(),
|
||||
ManagedRunSessionLink::Exact(foreign.run_id)
|
||||
);
|
||||
for run_id in [alpha.run_id, beta.run_id] {
|
||||
assert!(
|
||||
store
|
||||
.reader
|
||||
.managed_run_status(run_id)
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap()
|
||||
.native_session_id
|
||||
.is_none()
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -500,17 +500,23 @@ protocol](managed-harness-contributions.md), including read-only extraction,
|
||||
pre-turn context delivery, migration invariants, deterministic tests, and an
|
||||
opt-in real-harness acceptance pass.
|
||||
|
||||
### Known issue: Codex's shared daemon and stale run ids (#987)
|
||||
### Codex shared-daemon recovery (#987)
|
||||
|
||||
Recent Codex releases run sessions through a shared background app-server
|
||||
daemon (`codex agents` lists it). The daemon keeps the environment it started
|
||||
with, and the lifecycle hooks it launches inherit that environment — including
|
||||
the `AI_MEMORY_RUN_ID` of whichever managed run auto-started it. A later
|
||||
`ai-memory run codex` then reports that finished run, gets no continuity
|
||||
context, and the server logs `managed SessionStart has no active run`.
|
||||
the `AI_MEMORY_RUN_ID` of whichever managed run auto-started it. When a later
|
||||
managed Codex SessionStart reports that finished id, ai-memory ignores it as
|
||||
authority and searches only the request's already-resolved repository, exact
|
||||
checkout cwd, and operator bucket. The server adopts a replacement only when
|
||||
exactly one live, undelivered Codex run is waiting there, and it selects plus
|
||||
links that run in one writer transaction. Multiple candidates, an active id
|
||||
whose scope, checkout, or owner does not match, and a run already linked by
|
||||
another session all fail closed.
|
||||
|
||||
Until the server-side fix lands, launch managed Codex sessions without the
|
||||
daemon. Native arguments after the harness are forwarded to Codex:
|
||||
Ordinary managed Codex launches therefore work with the shared daemon. For
|
||||
diagnosis, or when an intentionally concurrent pair of named workstreams makes
|
||||
recovery ambiguous, native arguments after the harness are still forwarded:
|
||||
|
||||
```bash
|
||||
ai-memory run codex --no-daemon
|
||||
@@ -518,8 +524,8 @@ ai-memory run codex --no-daemon
|
||||
|
||||
`--no-daemon` makes that one session run without the shared background server
|
||||
even if one is already running; it is available on Codex's interactive and
|
||||
`resume` commands (checked on Codex 0.156). Sessions started without it keep
|
||||
the daemon behavior described above.
|
||||
`resume` commands (checked on Codex 0.156). It remains a useful isolation
|
||||
switch, not a requirement for normal workstream continuity.
|
||||
|
||||
## Installation and recovery
|
||||
|
||||
|
||||
@@ -57,6 +57,7 @@ boundary not yet built.
|
||||
| 11d | Hook server-profile routing: a marker-selected server gets only its own capture and its own token (#992) | `ai-memory-cli/src/server_profiles.rs` `resolve` (validated `ProfileName`, strict `servers.toml` parse, `roots` required once two profiles exist, component-wise root match) and `marker.rs` `find_server_selection` (inherited down the tree, any value shape counts); `commands/hook.rs` `resolve_hook_route` drops a `Rejected` route before spool, handoff and backfill, and hands the drainer no live token for a profile route; `commands/hook_spool.rs` `static_retry_token` (a profile entry retries only with its own stored token), the loopback reroot skip, and chunk splitting on `profile`; generated TS `captureServerRouted` drops a routed repository and gates `fetchHandoff`; `hooks/_lib.sh` `ai_memory_server_routed` (flag refused by `ai_memory_post_hook`/`ai_memory_get_handoff`) and `hooks/lib/ai-memory-hook.ps1` `Test-AiMemoryServerRouted` drop it in the script hooks | `hook.rs` `each_repository_spools_to_its_own_profile_with_its_own_token`, `a_selection_that_does_not_resolve_emits_nothing` (unknown / tokenless / outside roots / roots required / invalid name, plus a resolving control), `session_start_handoff_comes_from_the_profile_server_only`, `a_repository_without_a_server_key_keeps_the_install_default`; `hook_spool.rs` `a_profile_entry_is_never_retried_with_the_install_live_token` (server B accepts exactly the install's live token and must still not get it), `a_profile_entry_recovers_with_its_own_rotated_token` (control), `a_profile_entry_on_a_dead_loopback_port_is_not_rerouted_to_the_default`, `profile_and_default_entries_at_one_address_ride_separate_batches`; `server_profiles.rs` roots/registry/name tests; `marker.rs` `nested_markers_without_server_inherit_the_ancestor_selection`; `install_hooks.rs` `generated_integrations_fail_closed_on_a_server_profile_marker`, `openclaw_plugin.rs` `openclaw_plugin_fails_closed_on_a_server_profile_marker`, and the `server-routed-*` checks in `generated_capture_policy_v1_node_runtime_evidence`; `hook.rs` `an_event_without_a_payload_cwd_routes_by_the_process_cwd`, `a_refused_route_prints_nothing_for_kimi_user_prompts`; `hook_spool.rs` `a_profile_entry_is_not_retried_with_a_token_issued_for_a_new_url`; `server_profiles.rs` `changing_the_url_without_a_token_discards_the_old_token`, `omitted_roots_keep_the_registered_ones`; `marker.rs` `outside_home_the_walk_reaches_a_marker_above_the_checkout_root`, `encoding_noise_cannot_hide_a_server_key`, `an_unreadable_marker_is_a_refused_selection`; `backfill.rs` `a_spawned_backfill_authenticates_like_the_hook_that_spawned_it`; `tests/hooks/test_lib.sh` "server profiles (#992)" section; `hook.rs` `a_mixed_spool_drains_each_event_only_to_its_own_server` (two token-gated servers, one spool, one drain: each server receives exactly its own event with exactly its own bearer); `tests/suite/server_profiles.rs` (the built binary: `server add` → `hook` spools to the profile with its token, unknown profile spools nothing, `uninstall` removes the tokens; `two_real_servers_each_receive_only_their_own_repository` runs two real `ai-memory serve` children with different root tokens and checks on each server which repository landed there); `ai-memory-hooks` `powershell_server_routed.rs` `server_routed_guard_mirrors_the_native_walk` (runs `Test-AiMemoryServerRouted` under real PowerShell: inherited, BOM, look-alike keys, the `$HOME` boundary with its control, past a checkout root outside home, current directory); `tests/hooks/test_lib.sh` and `powershell_home.rs` `a_trailing_separator_on_home_keeps_the_walk_boundary` (a `$HOME` ending in `/` or `\` keeps both script walks at home: a `server` marker above it routes nothing, with a root-home control that does) | STRONG for native hooks, generated TS, `.sh` and `.ps1` hooks. Known gap: an older binary draining a shared spool ignores `profile` |
|
||||
| 12 | Network/auth posture | `config.rs` loopback `DEFAULT_BIND`; `serve.rs` `validate_http_exposure`, `require_allowed_host`; `auth.rs` `require_bearer` | `serve.rs` host-guard (missing→400 / forged→403), non-loopback-requires-token; `auth.rs` wrong-token 401 | STRONG |
|
||||
| 13 | Managed-run transcript attribution (concurrent launches in one checkout, invariant #16) | `ai-memory-store/src/workstream.rs` `link_native_session` stamps `native_session_linked_at` on its own run only, both `finish` updates drop the stamp when the session changes, `run_status` reports it; `ai-memory-cli/src/commands/run.rs` `resolve_native_session_after_run` takes a linked session only when `ai-memory-workstream` `native_session_in_checkout` holds it for this checkout (OpenCode by recorded directory), never falls back to a link it set aside, and turns an `AmbiguousNativeSession` into a warning with nothing imported; `ai-memory-workstream/src/transcript.rs` `discover_crush` claims only the one top-level session created (or, with `--continue`, touched) during the run, and in a data directory outside the project only one that edited a file in it | `multi_session.rs` `a_session_linked_by_one_managed_run_is_not_another_runs`; `run.rs` `a_session_linked_during_the_run_wins_over_discovery` (concurrent newer session, another checkout's link refused, no fallback to it, unlinked control); `transcript.rs` `native_session_in_checkout_checks_the_opencode_directory`; `store/src/lib.rs` `managed_run_status_reports_a_link_made_during_the_run`; `run.rs` `ambiguous_crush_discovery_keeps_the_run`; `transcript.rs` `crush_discovery_claims_only_the_session_the_run_created`, `crush_discovery_in_a_shared_store_claims_only_an_edit_here` | PARTIAL: a run whose child links nothing still falls back to discovering the newest session in the checkout; Crush, which has no hooks, relies on that discovery alone and claims nothing when it is ambiguous |
|
||||
| 13a | Managed Codex stale-daemon run recovery (#987; concurrent launches, project/worktree isolation, operator isolation) | `managed_runs.owner_user` (V71) uses the topology-aware `IdentityKey` owner stamp; `router.rs` resolves and authorizes the SessionStart repository before handling the untrusted run id; store `link_or_adopt_native_session` selects and links in one writer transaction and admits only the sole live, undelivered, never-linked Codex run in the exact workspace/project/checkout-cwd/owner bucket. An active supplied id must itself match every boundary; another harness, an ambiguity, or a second native session is refused | `multi_session.rs` `stale_codex_run_recovery_cannot_cross_project_owner_or_session_boundaries` (foreign owner, foreign checkout, and cross-project candidates refused; project-local controls and racing rebind), `stale_codex_run_recovery_fails_closed_on_ambiguity_and_active_mismatch` (two candidates and active foreign id refused, exact control); `workstream.rs` `prepare_stamps_the_operator_bucket_used_by_session_start_recovery` (HTTP actor stamp with foreign and legitimate controls); `router.rs` `managed_session_start_delivers_ledger_and_pending_handoff` supplies a finished stale id and proves the new run is linked and delivered | STRONG for SessionStart recovery. No candidate is selected from a run id's stale scope, another checkout, or another operator bucket; ambiguity deliberately falls back to no managed packet |
|
||||
| 13b | Managed-run lease exclusivity (one active run per workstream, invariant #16) | `ai-memory-store/src/workstream.rs` `prepare_run` expires lapsed leases and refuses any other `active` run on the workstream inside one transaction (`StoreError::WorkstreamBusy`), regardless of the `lease_owner` label; `heartbeat` renews only `active` rows. `ai-memory-cli/src/commands/run.rs` `wait_out_held_lease` (interactive relaunch) only waits for a reported expiry and retries — it never cancels or claims another run, so the server's busy check stays the sole arbiter | `store/src/lib.rs` `managed_workstream_batches_are_idempotent_and_release_the_lease` (second prepare refused while active); `run.rs` `a_renewed_lease_is_reported_as_a_live_owner_not_taken_over` (renewing holder is reported, never displaced), with controls `interactive_launch_waits_out_a_lapsing_lease_then_proceeds` and `ctrl_c_aborts_the_held_lease_wait_immediately` | STRONG for exclusivity. The `lease_owner` label (`host:pid`) is informational only and not unique inside ai-jail (every jailed launcher reports `ai-sandbox:<ns-pid>`), so it must never become an ownership key |
|
||||
| 14 | Per-project authorization (#708) | `ai-memory-store/src/project_authz.rs` `authorize_project` / `ProjectAuthz::authorize` choke point (V68 `project_grants` + `projects.access_mode`, default `open`; V69 `projects.created_by` feeds `is_creator`); `scope.rs` `ScopeResolver::with_project_authz` (reader pool for reads, writer actor for writes) and its free forms `authorize_scope_for` / `*_guarded`, attached for every DB user by `ai-memory-mcp` `scope_resolver_as`, `ai-memory-web` `authorize_read` / `lookup_project`, and `ai-memory-hooks` (`grants.rs` `authorize_resolved` for run/workstream ids, the capture check in `router.rs`); read-shaped mutations resolve at `ProjectAccess::Write` (`resolve_existing_args(.., need)`); unscoped reads filtered before `LIMIT` by `reader.rs` `readable_repository_sql`; `WriterHandle::authorize_project` as defense in depth. Page ids are never taken from a caller (only derived from already-authorized hits), so there is no page-id entry point to guard | `tests/suite/project_authz.rs` (decision matrix, ship-inert, resolver gate); `tests/suite/access_mode.rs` `every_caller_against_both_modes`, `a_restricted_project_admits_the_team_and_refuses_the_outsider`, `new_projects_follow_the_server_default_and_admit_their_creator`; `scope.rs` `the_argument_shape_does_not_decide_the_level`, `a_user_reaches_only_what_they_were_granted`, `creating_authorizes_against_a_project_that_already_exists`, `a_refused_scope_fails_the_search_instead_of_shortening_it`; `grants.rs` `search_finds_only_what_the_viewer_may_read`, `the_limit_counts_only_what_the_viewer_may_see`, `the_workspace_handoff_comes_only_from_readable_repositories`; `ai-memory-mcp` `server.rs` `a_reader_may_read_everything_and_change_nothing`, `bob_cannot_read_alices_page_in_a_restricted_project`, `bob_cannot_find_alices_page_by_searching_in_a_restricted_project`, `bob_cannot_consolidate_a_session_in_alices_repository`, `a_restricted_projects_queues_need_write`; `ai-memory-hooks` `a_capture_needs_writer_on_the_repository_it_lands_in`, `session_start_delivers_nothing_from_a_repository_the_viewer_cannot_read`, `run_and_workstream_ids_only_answer_someone_who_may_reach_the_repository`; `ai-memory-web` `web_reads_honour_grants_in_a_restricted_project`, `metadata_shows_only_what_the_viewer_may_read` — each with a granted or open-project control, and proven to fail with the choke point (18 tests) or the SQL filter (9 tests) neutralized | STRONG — slice 3 closed both bypass classes (unscoped reads, raw-id entry points) and added the root-only management surface. Out of scope by design: per-project administrators (granting/restricting is root-only), and access modes, creators and grants live only in SQLite, so `reindex` resets them |
|
||||
| 14b | Repository identity routing (#708) | `ai-memory-store/src/ops.rs` `resolve_project_by_identity` — one transaction; an identity already on a project is never overwritten; an unclaimed project is claimed only when the capturing user may write it (the choke point's `resolve_project_authz` on the same transaction), otherwise it is returned unclaimed; a different identity under the same name splits into a new project with no shared `repo_path`; `ai-memory-hooks/src/router.rs` `cache_key_for` keys the path cache by identity too; `ai-memory-core/src/repository_identity.rs` strips credentials from remote URLs client-side, and the server accepts only the routing rungs (`explicit`, `git_remote`) from the wire | `tests/suite/identity_resolution.rs` `an_outsider_cannot_take_an_unclaimed_projects_identity` (control: `an_existing_project_is_claimed_in_place`), `two_unrelated_repositories_with_one_folder_name_stay_apart`, `a_created_or_split_project_admits_its_creator`; `router.rs` `one_path_with_two_remotes_is_two_projects`, `two_api_checkouts_with_different_remotes_get_two_projects` (a manifest rung on the wire is ignored); `repository_identity.rs` credential and normalisation tables — the claim guard and the cache key each proven to fail their test when removed | STRONG — first claimant wins: in a restricted workspace a user who creates a project under a remote identity first holds it until root grants others |
|
||||
|
||||
Reference in New Issue
Block a user