mirror of
https://github.com/akitaonrails/ai-memory.git
synced 2026-10-02 03:24:46 +08:00
Merge PR #880 into main
fix(run): import the session linked during a managed run (#820) Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01MDbhmszrjG9s5MrPrTuNtm # Conflicts: # CHANGELOG.md
This commit is contained in:
@@ -28,6 +28,15 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
|
||||
run` after an upgrade, silently disabling `[auto_scope] per_session` for their
|
||||
MCP calls. Auto-wire now detects an existing session-aware bridge and keeps
|
||||
it. (#888)
|
||||
- After a managed launch, `ai-memory run` imported the newest native session
|
||||
in the checkout even when a hook in the launched harness had linked the
|
||||
run's own session, so a concurrent launch in the same checkout could hand it
|
||||
another transcript. The server now records when a session is linked during
|
||||
a run (schema migration V67, adding `managed_runs.native_session_linked_at`)
|
||||
and reports it in the run status, and the launcher imports that session
|
||||
when this checkout's store holds it (a process the child starts inherits
|
||||
the run id; OpenCode is checked by the session's recorded directory). An
|
||||
older server reports no link and keeps the previous behavior. (#820)
|
||||
- The Linux/macOS Docker wrapper now keeps its native host client in
|
||||
`${XDG_DATA_HOME:-~/.local/share}/ai-memory/native-runner` instead of
|
||||
`~/.cache/ai-memory/native-runner`. `ai-memory run` auto-wires hooks whose
|
||||
|
||||
@@ -17,7 +17,7 @@ use ai_memory_workstream::{
|
||||
export_transcript, has_native_session_selector, inspect_repository, kiro_explicit_session_id,
|
||||
kiro_harness_from_source_cursor, kiro_selects_non_default_engine, kiro_selects_v2_engine,
|
||||
kiro_selects_v3_engine, kiro_v3_resume_uses_default_store, list_native_sessions,
|
||||
native_session_exists, wait_for_transcript_flush,
|
||||
native_session_exists, native_session_in_checkout, wait_for_transcript_flush,
|
||||
};
|
||||
use anyhow::{Context as _, Result, anyhow};
|
||||
use tokio::process::Command;
|
||||
@@ -649,10 +649,29 @@ async fn resolve_native_session_after_run(
|
||||
if let Some(native_session_id) = &plan.expected_session_id {
|
||||
return Ok(Some(native_session_id.clone()));
|
||||
}
|
||||
// A session linked under this run's id was reported by this run's child,
|
||||
// which a concurrent launch in the same checkout cannot do; discovery
|
||||
// only sees the newest session there. A descendant process inherits the
|
||||
// id too, so the session must also be this checkout's.
|
||||
let linked = server_status
|
||||
.filter(|status| status.native_session_linked)
|
||||
.and_then(|status| status.native_session_id.as_deref());
|
||||
if let Some(linked) = linked
|
||||
&& native_session_in_checkout(harness, home, cwd, plan.session_dir.as_deref(), linked)
|
||||
.unwrap_or(false)
|
||||
{
|
||||
return Ok(Some(linked.to_string()));
|
||||
}
|
||||
let discovered =
|
||||
discover_native_session(harness, home, cwd, plan.session_dir.as_deref(), started_at)
|
||||
.await?;
|
||||
Ok(discovered.or_else(|| server_status.and_then(|status| status.native_session_id.clone())))
|
||||
// A linked session set aside above belongs to another checkout, so it is
|
||||
// no fallback either.
|
||||
Ok(discovered.or_else(|| {
|
||||
server_status
|
||||
.and_then(|status| status.native_session_id.clone())
|
||||
.filter(|reported| Some(reported.as_str()) != linked)
|
||||
}))
|
||||
}
|
||||
|
||||
async fn list_auto_sessions(home: &Path, cwd: &Path) -> Result<Vec<AutoSessionCandidate>> {
|
||||
@@ -2294,6 +2313,91 @@ mod tests {
|
||||
assert!(error.to_string().contains("--fresh cannot be combined"));
|
||||
}
|
||||
|
||||
/// A session linked during the run was reported by this run's child, so
|
||||
/// it wins over a newer session another launch made in the same checkout,
|
||||
/// even when it repeats the session the run was prepared with. A child's
|
||||
/// own descendants inherit the run id, so a linked session this checkout
|
||||
/// does not hold is set aside. Without a link, discovery still decides.
|
||||
#[tokio::test]
|
||||
async fn a_session_linked_during_the_run_wins_over_discovery() {
|
||||
let temp = tempfile::tempdir().unwrap();
|
||||
let cwd = temp.path().join("repo");
|
||||
let session_root = temp.path().join(".codex/sessions/2026/01/01");
|
||||
std::fs::create_dir_all(&cwd).unwrap();
|
||||
std::fs::create_dir_all(&session_root).unwrap();
|
||||
let started_at = SystemTime::now();
|
||||
let rollout = |name: &str, id: &str, cwd: &Path| {
|
||||
std::fs::write(
|
||||
session_root.join(format!("rollout-{name}.jsonl")),
|
||||
format!(
|
||||
"{}\n",
|
||||
serde_json::json!({
|
||||
"type": "session_meta",
|
||||
"payload": {"id": id, "cwd": cwd}
|
||||
})
|
||||
),
|
||||
)
|
||||
.unwrap();
|
||||
};
|
||||
rollout("prepared", "prepared", &cwd);
|
||||
rollout("nested", "nested", &temp.path().join("other-checkout"));
|
||||
rollout("concurrent", "concurrent-newer", &cwd);
|
||||
let plan = build_launch_plan(ManagedHarness::Codex, None, Vec::new(), None).unwrap();
|
||||
let status = |linked: bool, native: &str| ManagedRunStatus {
|
||||
run_id: ManagedRunId::new(),
|
||||
workstream_id: WorkstreamId::new(),
|
||||
agent: AgentKind::Codex,
|
||||
native_session_id: Some(native.to_string()),
|
||||
native_session_linked: linked,
|
||||
context_delivered: true,
|
||||
state: "active".to_string(),
|
||||
};
|
||||
for (linked, native, expected) in [
|
||||
(true, "prepared", Some("prepared")),
|
||||
(false, "prepared", Some("concurrent-newer")),
|
||||
(true, "nested", Some("concurrent-newer")),
|
||||
] {
|
||||
let status = status(linked, native);
|
||||
assert_eq!(
|
||||
resolve_native_session_after_run(
|
||||
&plan,
|
||||
ManagedHarness::Codex,
|
||||
temp.path(),
|
||||
&cwd,
|
||||
started_at,
|
||||
Some(&status),
|
||||
)
|
||||
.await
|
||||
.unwrap()
|
||||
.as_deref(),
|
||||
expected,
|
||||
"linked={linked} native={native}"
|
||||
);
|
||||
}
|
||||
// With nothing to discover here, the other checkout's session is not
|
||||
// taken as a fallback either; an unlinked report still is.
|
||||
let empty = temp.path().join("empty-checkout");
|
||||
std::fs::create_dir_all(&empty).unwrap();
|
||||
for (linked, expected) in [(true, None), (false, Some("nested"))] {
|
||||
let status = status(linked, "nested");
|
||||
assert_eq!(
|
||||
resolve_native_session_after_run(
|
||||
&plan,
|
||||
ManagedHarness::Codex,
|
||||
temp.path(),
|
||||
&empty,
|
||||
started_at,
|
||||
Some(&status),
|
||||
)
|
||||
.await
|
||||
.unwrap()
|
||||
.as_deref(),
|
||||
expected,
|
||||
"linked={linked}"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn utility_launch_does_not_adopt_a_recent_unrelated_session() {
|
||||
let temp = tempfile::tempdir().unwrap();
|
||||
|
||||
@@ -245,6 +245,12 @@ pub struct ManagedRunStatus {
|
||||
/// Native session linked by SessionStart, if observed.
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub native_session_id: Option<String>,
|
||||
/// Whether `native_session_id` was linked during this run (by a hook in
|
||||
/// the child, or by the launcher before the spawn) rather than carried
|
||||
/// over from the workstream when the run was prepared. An older server
|
||||
/// does not send it and reads as `false`.
|
||||
#[serde(default)]
|
||||
pub native_session_linked: bool,
|
||||
/// Whether the SessionStart context packet was returned successfully.
|
||||
pub context_delivered: bool,
|
||||
/// Current run state (`active`, `finished`, or `expired`).
|
||||
@@ -347,6 +353,21 @@ pub struct WorkstreamEvent {
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn older_run_status_reads_as_nothing_linked() {
|
||||
let status: ManagedRunStatus = serde_json::from_value(serde_json::json!({
|
||||
"run_id": "018f0000-0000-7000-8000-000000000002",
|
||||
"workstream_id": "018f0000-0000-7000-8000-000000000001",
|
||||
"agent": "codex",
|
||||
"native_session_id": "prepared",
|
||||
"context_delivered": false,
|
||||
"state": "active"
|
||||
}))
|
||||
.unwrap();
|
||||
|
||||
assert!(!status.native_session_linked);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn older_prepare_response_defaults_to_no_adoption() {
|
||||
let response: PrepareManagedRunResponse = serde_json::from_value(serde_json::json!({
|
||||
|
||||
@@ -261,6 +261,7 @@ async fn run_status(
|
||||
workstream_id: status.workstream_id,
|
||||
agent: status.agent,
|
||||
native_session_id: status.native_session_id,
|
||||
native_session_linked: status.native_session_linked,
|
||||
context_delivered: status.context_delivered,
|
||||
state: status.state,
|
||||
})
|
||||
@@ -973,6 +974,48 @@ mod tests {
|
||||
(workspace_id, project_id)
|
||||
}
|
||||
|
||||
/// The launcher reads whether the run's child linked a session from the
|
||||
/// run status, so the route must carry it.
|
||||
#[tokio::test]
|
||||
async fn run_status_reports_a_session_linked_during_the_run() {
|
||||
let temp = TempDir::new().unwrap();
|
||||
let store = Store::open(temp.path()).unwrap();
|
||||
let state = test_state(&store, temp.path());
|
||||
let (workspace_id, project_id) = seed_scope(&store).await;
|
||||
let prepared = store
|
||||
.writer
|
||||
.prepare_workstream_run(prepare_input(
|
||||
workspace_id,
|
||||
project_id,
|
||||
AgentKind::Codex,
|
||||
"launcher",
|
||||
))
|
||||
.await
|
||||
.unwrap();
|
||||
let status = async || {
|
||||
let response = run_status(
|
||||
State(state.clone()),
|
||||
None,
|
||||
AxumPath(prepared.run_id.to_string()),
|
||||
)
|
||||
.await;
|
||||
assert_eq!(response.status(), StatusCode::OK);
|
||||
let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
|
||||
serde_json::from_slice::<ManagedRunStatus>(&body).unwrap()
|
||||
};
|
||||
assert!(!status().await.native_session_linked);
|
||||
assert!(
|
||||
store
|
||||
.writer
|
||||
.link_managed_run_session(prepared.run_id, AgentKind::Codex, "native-1")
|
||||
.await
|
||||
.unwrap()
|
||||
);
|
||||
let linked = status().await;
|
||||
assert!(linked.native_session_linked);
|
||||
assert_eq!(linked.native_session_id.as_deref(), Some("native-1"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn rename_endpoint_is_scoped_and_reports_selector_misuse() {
|
||||
let temp = TempDir::new().unwrap();
|
||||
|
||||
@@ -0,0 +1,14 @@
|
||||
-- Record when a session was linked to a managed run.
|
||||
--
|
||||
-- `managed_runs.native_session_id` starts as the workstream's current session
|
||||
-- when the run is prepared, and a hook in the launched child (or the launcher
|
||||
-- itself, before the spawn) replaces it through the link path. A link that
|
||||
-- repeats the prepared session left no trace, so after exit `ai-memory run`
|
||||
-- could not tell "the child confirmed this session" from "nothing linked", and
|
||||
-- fell back to the newest session in the checkout, which a concurrent launch
|
||||
-- may own.
|
||||
--
|
||||
-- NULL means nothing linked during this run; rows written by an older version
|
||||
-- stay NULL, the reading those runs already had.
|
||||
|
||||
ALTER TABLE managed_runs ADD COLUMN native_session_linked_at INTEGER;
|
||||
@@ -411,7 +411,7 @@ mod tests {
|
||||
|row| row.get(0),
|
||||
)
|
||||
.unwrap();
|
||||
assert_eq!(version, 66, "update the pin when adding a migration");
|
||||
assert_eq!(version, 67, "update the pin when adding a migration");
|
||||
let cols: i64 = conn
|
||||
.query_row(
|
||||
"SELECT COUNT(*) FROM pragma_table_info('users') WHERE name = 'token_hash'",
|
||||
|
||||
@@ -4915,6 +4915,150 @@ mod tests {
|
||||
assert!(mismatch.is_empty());
|
||||
}
|
||||
|
||||
/// A run starts with the workstream's current session, which is no
|
||||
/// evidence that its child used it; a link during the run is, even one
|
||||
/// that repeats that session. A link refused after the context packet
|
||||
/// went out marks nothing, and a finish that names another session drops
|
||||
/// the mark, which belonged to the one before.
|
||||
#[tokio::test]
|
||||
async fn managed_run_status_reports_a_link_made_during_the_run() {
|
||||
let tmp = TempDir::new().unwrap();
|
||||
let store = Store::open(tmp.path()).unwrap();
|
||||
let ws = store
|
||||
.writer
|
||||
.get_or_create_workspace("default")
|
||||
.await
|
||||
.unwrap();
|
||||
let project = store
|
||||
.writer
|
||||
.get_or_create_project(ws, "managed", None)
|
||||
.await
|
||||
.unwrap();
|
||||
let prepare = PrepareWorkstreamRun {
|
||||
workspace_id: ws,
|
||||
project_id: project,
|
||||
repo_fingerprint: "repo".into(),
|
||||
worktree_fingerprint: "worktree".into(),
|
||||
cwd: "/repo".into(),
|
||||
agent: AgentKind::Codex,
|
||||
automatic_harness: false,
|
||||
available_agents: Vec::new(),
|
||||
selection: WorkstreamSelection::Current,
|
||||
lease_owner: "test:1".into(),
|
||||
};
|
||||
let status = async |run_id| {
|
||||
let status = store
|
||||
.reader
|
||||
.managed_run_status(run_id)
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
(status.native_session_id, status.native_session_linked)
|
||||
};
|
||||
|
||||
let first = store
|
||||
.writer
|
||||
.prepare_workstream_run(prepare.clone())
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(status(first.run_id).await, (None, false));
|
||||
assert!(
|
||||
store
|
||||
.writer
|
||||
.link_managed_run_session(first.run_id, AgentKind::Codex, "native-1")
|
||||
.await
|
||||
.unwrap()
|
||||
);
|
||||
assert_eq!(status(first.run_id).await, (Some("native-1".into()), true));
|
||||
store
|
||||
.writer
|
||||
.finish_workstream_run(FinishWorkstreamRun {
|
||||
run_id: first.run_id,
|
||||
native_session_id: Some("native-1".into()),
|
||||
source_cursor: None,
|
||||
events: Vec::new(),
|
||||
complete: true,
|
||||
segment_path: None,
|
||||
exit_code: Some(0),
|
||||
})
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let second = store.writer.prepare_workstream_run(prepare).await.unwrap();
|
||||
assert_eq!(second.native_session_id.as_deref(), Some("native-1"));
|
||||
assert_eq!(
|
||||
status(second.run_id).await,
|
||||
(Some("native-1".into()), false)
|
||||
);
|
||||
assert!(
|
||||
store
|
||||
.writer
|
||||
.accept_managed_run_context(second.run_id)
|
||||
.await
|
||||
.unwrap()
|
||||
);
|
||||
assert!(
|
||||
!store
|
||||
.writer
|
||||
.link_managed_run_session(second.run_id, AgentKind::Codex, "native-2")
|
||||
.await
|
||||
.unwrap()
|
||||
);
|
||||
assert_eq!(
|
||||
status(second.run_id).await,
|
||||
(Some("native-1".into()), false)
|
||||
);
|
||||
assert!(
|
||||
store
|
||||
.writer
|
||||
.link_managed_run_session(second.run_id, AgentKind::Codex, "native-1")
|
||||
.await
|
||||
.unwrap()
|
||||
);
|
||||
assert_eq!(status(second.run_id).await, (Some("native-1".into()), true));
|
||||
let finish = |native: &str, complete: bool| FinishWorkstreamRun {
|
||||
run_id: second.run_id,
|
||||
native_session_id: Some(native.into()),
|
||||
source_cursor: None,
|
||||
events: Vec::new(),
|
||||
complete,
|
||||
segment_path: None,
|
||||
exit_code: None,
|
||||
};
|
||||
store
|
||||
.writer
|
||||
.finish_workstream_run(finish("native-1", false))
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(status(second.run_id).await, (Some("native-1".into()), true));
|
||||
store
|
||||
.writer
|
||||
.finish_workstream_run(finish("native-3", false))
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(
|
||||
status(second.run_id).await,
|
||||
(Some("native-3".into()), false)
|
||||
);
|
||||
assert!(
|
||||
store
|
||||
.writer
|
||||
.link_managed_run_session(second.run_id, AgentKind::Codex, "native-3")
|
||||
.await
|
||||
.unwrap()
|
||||
);
|
||||
assert_eq!(status(second.run_id).await, (Some("native-3".into()), true));
|
||||
store
|
||||
.writer
|
||||
.finish_workstream_run(finish("native-4", true))
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(
|
||||
status(second.run_id).await,
|
||||
(Some("native-4".into()), false)
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn managed_workstream_batches_are_idempotent_and_release_the_lease() {
|
||||
let tmp = TempDir::new().unwrap();
|
||||
|
||||
@@ -141,6 +141,9 @@ pub struct StoredManagedRunStatus {
|
||||
pub agent: AgentKind,
|
||||
/// Native session observed by hooks.
|
||||
pub native_session_id: Option<String>,
|
||||
/// Whether a session was linked during this run, rather than carried
|
||||
/// over from the workstream when the run was prepared.
|
||||
pub native_session_linked: bool,
|
||||
/// SessionStart delivery acknowledgement.
|
||||
pub context_delivered: bool,
|
||||
/// State string.
|
||||
@@ -575,8 +578,9 @@ pub(crate) fn link_native_session(
|
||||
],
|
||||
)?;
|
||||
tx.execute(
|
||||
"UPDATE managed_runs SET native_session_id = ?1, sync_after = ?2 WHERE id = ?3",
|
||||
params![native_session_id, initial_delivery, run_id.as_bytes()],
|
||||
"UPDATE managed_runs SET native_session_id = ?1, sync_after = ?2, \
|
||||
native_session_linked_at = ?3 WHERE id = ?4",
|
||||
params![native_session_id, initial_delivery, now, run_id.as_bytes()],
|
||||
)?;
|
||||
tx.commit()?;
|
||||
Ok(true)
|
||||
@@ -793,6 +797,8 @@ pub(crate) fn finish_run(
|
||||
if input.complete {
|
||||
tx.execute(
|
||||
"UPDATE managed_runs SET state = 'finished', native_session_id = COALESCE(?1, native_session_id), \
|
||||
native_session_linked_at = CASE WHEN ?1 IS NULL OR ?1 = native_session_id \
|
||||
THEN native_session_linked_at END, \
|
||||
ended_at = ?2, lease_expires_at = ?2, exit_code = ?3 WHERE id = ?4",
|
||||
params![
|
||||
native_session,
|
||||
@@ -804,6 +810,8 @@ pub(crate) fn finish_run(
|
||||
} else {
|
||||
tx.execute(
|
||||
"UPDATE managed_runs SET native_session_id = COALESCE(?1, native_session_id), \
|
||||
native_session_linked_at = CASE WHEN ?1 IS NULL OR ?1 = native_session_id \
|
||||
THEN native_session_linked_at END, \
|
||||
lease_expires_at = ?2 WHERE id = ?3",
|
||||
params![native_session, now + LEASE_MICROS, input.run_id.as_bytes()],
|
||||
)?;
|
||||
@@ -826,7 +834,8 @@ pub(crate) fn run_status(
|
||||
) -> StoreResult<Option<StoredManagedRunStatus>> {
|
||||
let row = conn
|
||||
.query_row(
|
||||
"SELECT workstream_id, agent_kind, native_session_id, context_delivered, state \
|
||||
"SELECT workstream_id, agent_kind, native_session_id, \
|
||||
native_session_linked_at IS NOT NULL, context_delivered, state \
|
||||
FROM managed_runs WHERE id = ?1",
|
||||
params![run_id.as_bytes()],
|
||||
|row| {
|
||||
@@ -835,18 +844,27 @@ pub(crate) fn run_status(
|
||||
row.get::<_, String>(1)?,
|
||||
row.get::<_, Option<String>>(2)?,
|
||||
row.get::<_, bool>(3)?,
|
||||
row.get::<_, String>(4)?,
|
||||
row.get::<_, bool>(4)?,
|
||||
row.get::<_, String>(5)?,
|
||||
))
|
||||
},
|
||||
)
|
||||
.optional()?;
|
||||
row.map(
|
||||
|(workstream, agent, native_session_id, context_delivered, state)| {
|
||||
|(
|
||||
workstream,
|
||||
agent,
|
||||
native_session_id,
|
||||
native_session_linked,
|
||||
context_delivered,
|
||||
state,
|
||||
)| {
|
||||
Ok(StoredManagedRunStatus {
|
||||
run_id,
|
||||
workstream_id: WorkstreamId::from_slice(&workstream)?,
|
||||
agent: AgentKind::from_wire(&agent),
|
||||
native_session_id,
|
||||
native_session_linked,
|
||||
context_delivered,
|
||||
state,
|
||||
})
|
||||
|
||||
@@ -20,7 +20,7 @@ use ai_memory_core::{
|
||||
ActorContext, AgentKind, HandoffAcceptance, IdentityKey, NewHandoff, NewPage, NewSession,
|
||||
NewUser, OwnerFilter, PagePath, ProjectId, SessionId, Tier, UserRole, WorkspaceId, owner_stamp,
|
||||
};
|
||||
use ai_memory_store::Store;
|
||||
use ai_memory_store::{PrepareWorkstreamRun, Store, WorkstreamSelection};
|
||||
|
||||
fn operator(name: &str) -> String {
|
||||
IdentityKey::User(name.into()).storage_key()
|
||||
@@ -426,3 +426,71 @@ async fn an_owned_handoff_stays_with_its_owner_while_pages_stay_shared() {
|
||||
operator("alice")
|
||||
);
|
||||
}
|
||||
|
||||
/// Two workstreams launched at once in one checkout each get a managed run.
|
||||
/// A session one run's child links marks that run only: the other run's
|
||||
/// status must neither report it nor count as linked, or its launcher would
|
||||
/// import the other launch's transcript.
|
||||
#[tokio::test]
|
||||
async fn a_session_linked_by_one_managed_run_is_not_another_runs() {
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
let store = Store::open(tmp.path()).unwrap();
|
||||
let (ws, proj) = scope(&store).await;
|
||||
let prepare = |name: &str| PrepareWorkstreamRun {
|
||||
workspace_id: ws,
|
||||
project_id: proj,
|
||||
repo_fingerprint: "repo".into(),
|
||||
worktree_fingerprint: "worktree".into(),
|
||||
cwd: "/repo".into(),
|
||||
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(prepare("alpha"))
|
||||
.await
|
||||
.unwrap();
|
||||
let beta = store
|
||||
.writer
|
||||
.prepare_workstream_run(prepare("beta"))
|
||||
.await
|
||||
.unwrap();
|
||||
assert_ne!(alpha.workstream_id, beta.workstream_id);
|
||||
let status = async |run| {
|
||||
let status = store.reader.managed_run_status(run).await.unwrap().unwrap();
|
||||
(status.native_session_id, status.native_session_linked)
|
||||
};
|
||||
|
||||
assert!(
|
||||
store
|
||||
.writer
|
||||
.link_managed_run_session(beta.run_id, AgentKind::Codex, "native-beta")
|
||||
.await
|
||||
.unwrap()
|
||||
);
|
||||
assert_eq!(status(alpha.run_id).await, (None, false));
|
||||
assert_eq!(
|
||||
status(beta.run_id).await,
|
||||
(Some("native-beta".into()), true)
|
||||
);
|
||||
|
||||
// Control: alpha's own link marks alpha, and leaves beta as it was.
|
||||
assert!(
|
||||
store
|
||||
.writer
|
||||
.link_managed_run_session(alpha.run_id, AgentKind::Codex, "native-alpha")
|
||||
.await
|
||||
.unwrap()
|
||||
);
|
||||
assert_eq!(
|
||||
status(alpha.run_id).await,
|
||||
(Some("native-alpha".into()), true)
|
||||
);
|
||||
assert_eq!(
|
||||
status(beta.run_id).await,
|
||||
(Some("native-beta".into()), true)
|
||||
);
|
||||
}
|
||||
|
||||
@@ -13,5 +13,5 @@ pub use repository::{RepositoryIdentity, inspect_repository};
|
||||
pub use transcript::{
|
||||
ExportedTranscript, NativeSessionCandidate, discover_native_session, export_transcript,
|
||||
kiro_harness_from_source_cursor, kiro_v3_resume_uses_default_store, list_native_sessions,
|
||||
native_session_exists, wait_for_transcript_flush,
|
||||
native_session_exists, native_session_in_checkout, wait_for_transcript_flush,
|
||||
};
|
||||
|
||||
@@ -232,6 +232,36 @@ pub async fn list_native_sessions(
|
||||
Ok(sessions)
|
||||
}
|
||||
|
||||
/// Whether the native store holds `native_session_id` for this checkout.
|
||||
/// [`native_session_exists`] answers for the store; OpenCode keeps every
|
||||
/// checkout's sessions in one database, so there the recorded directory must
|
||||
/// match `cwd` too.
|
||||
pub fn native_session_in_checkout(
|
||||
harness: ManagedHarness,
|
||||
home: &Path,
|
||||
cwd: &Path,
|
||||
session_dir: Option<&Path>,
|
||||
native_session_id: &str,
|
||||
) -> Result<bool> {
|
||||
let table = match harness {
|
||||
ManagedHarness::OpenCode => "session",
|
||||
ManagedHarness::OpenCode2 => "session_v2",
|
||||
_ => return native_session_exists(harness, home, cwd, session_dir, native_session_id),
|
||||
};
|
||||
let db = opencode_db(home, session_dir);
|
||||
if !db.is_file() {
|
||||
return Ok(false);
|
||||
}
|
||||
let connection = Connection::open_with_flags(
|
||||
&db,
|
||||
OpenFlags::SQLITE_OPEN_READ_ONLY | OpenFlags::SQLITE_OPEN_NO_MUTEX,
|
||||
)?;
|
||||
let mut statement = connection.prepare(&format!(
|
||||
"SELECT 1 FROM {table} WHERE id = ?1 AND directory = ?2"
|
||||
))?;
|
||||
Ok(statement.exists(params![native_session_id, cwd.to_string_lossy()])?)
|
||||
}
|
||||
|
||||
/// Check whether one exact native session still exists in the harness's
|
||||
/// read-only transcript store. `Ok(false)` means the resume target is
|
||||
/// definitely absent; store access or schema failures remain errors so callers
|
||||
@@ -3745,6 +3775,50 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
/// OpenCode keeps every checkout's sessions in one database, so a session
|
||||
/// is this checkout's only when its recorded directory matches; other
|
||||
/// harnesses answer as `native_session_exists` does.
|
||||
#[test]
|
||||
fn native_session_in_checkout_checks_the_opencode_directory() {
|
||||
let temp = tempfile::tempdir().unwrap();
|
||||
let cwd = temp.path().join("repo");
|
||||
let other = temp.path().join("other");
|
||||
let store = temp.path().join("opencode");
|
||||
fs::create_dir_all(&store).unwrap();
|
||||
let connection = Connection::open(store.join("opencode.db")).unwrap();
|
||||
connection
|
||||
.execute_batch(
|
||||
"CREATE TABLE session(id TEXT PRIMARY KEY, directory TEXT NOT NULL, \
|
||||
time_updated INTEGER NOT NULL); \
|
||||
CREATE TABLE session_v2(id TEXT PRIMARY KEY, directory TEXT NOT NULL, \
|
||||
time_updated INTEGER NOT NULL);",
|
||||
)
|
||||
.unwrap();
|
||||
for table in ["session", "session_v2"] {
|
||||
for (id, dir) in [("here", &cwd), ("elsewhere", &other)] {
|
||||
connection
|
||||
.execute(
|
||||
&format!("INSERT INTO {table} VALUES (?1, ?2, 1)"),
|
||||
params![id, dir.to_string_lossy()],
|
||||
)
|
||||
.unwrap();
|
||||
}
|
||||
}
|
||||
for harness in [ManagedHarness::OpenCode, ManagedHarness::OpenCode2] {
|
||||
let in_checkout = |id: &str| {
|
||||
native_session_in_checkout(harness, temp.path(), &cwd, Some(&store), id).unwrap()
|
||||
};
|
||||
assert!(in_checkout("here"), "{harness:?}");
|
||||
assert!(!in_checkout("elsewhere"), "{harness:?}");
|
||||
assert!(!in_checkout("missing"), "{harness:?}");
|
||||
assert!(
|
||||
native_session_exists(harness, temp.path(), &cwd, Some(&store), "elsewhere")
|
||||
.unwrap(),
|
||||
"{harness:?}: the store holds it"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn crush_candidate_discovery_and_incremental_export_are_read_only() {
|
||||
let temp = tempfile::tempdir().unwrap();
|
||||
|
||||
@@ -331,7 +331,12 @@ resume, continue, session, or fork selector.
|
||||
Direct launches continue to use the same handoff path without a managed
|
||||
packet.
|
||||
4. When the child exits, ai-memory reads the native transcript store without
|
||||
modifying it. Visible user/assistant messages, completed tool calls/results,
|
||||
modifying it. A session named on the command line (or chosen before the
|
||||
spawn) is the one it reads. Otherwise a session linked during the run under
|
||||
its `AI_MEMORY_RUN_ID`, even the workstream's current one, is read when the
|
||||
native store holds it for this checkout; only without such a link does it
|
||||
look for the newest session in the checkout, which a concurrent launch
|
||||
there could own. Visible user/assistant messages, completed tool calls/results,
|
||||
compaction summaries, and a non-mutating Git checkpoint enter an append-only
|
||||
workstream ledger. Hidden reasoning and unsupported/private records are
|
||||
excluded and recorded as extraction-loss annotations. Each delivered
|
||||
|
||||
@@ -42,6 +42,7 @@ boundary not yet built.
|
||||
| 11b | Capture exclusions drop before storage | `ai-memory-hooks` `capture_policy.rs` `inspect`→`Drop` (before semaphore/spawn) | `capture_policy.rs` per-agent `…honors_exclusions` tests | STRONG |
|
||||
| 11c | Capture hook ≤200ms budget (invariant #5) | `hooks/_lib.sh` capture path `curl --max-time 0.2` (context-fetch 1.0s and background drain 2.0s are separate, larger-budget paths) | none (shell-script timeout; hard to unit-test) — watch on any capture-path change | WATCH |
|
||||
| 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), and never falls back to a link it set aside | `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` | PARTIAL: a run whose child links nothing still falls back to discovering the newest session in the checkout |
|
||||
| — | Per-project authorization (#708) | proposal only — `docs/design-per-project-authz.md` (`authorize_project` choke point + unscoped-read/raw-id bypass classes) | none yet — the design's "Verification plan" tests (authz matrix, unscoped-read-leak, raw-id-authz, ship-inert) land WITH the code | FUTURE |
|
||||
|
||||
## Keeping this current (the standing protocol)
|
||||
|
||||
Reference in New Issue
Block a user