mirror of
https://github.com/akitaonrails/ai-memory.git
synced 2026-10-02 03:24:46 +08:00
fix(handoff): keep parallel OpenCode sessions' batons owned and deliver only quiet ones
With turn checkpoints (#865), several live OpenCode sessions in one directory each publish a baton. Each checkpoint retired the other live sessions' batons, and SessionStart handed a new session the baton of a session still in use, with another conversation's context. A checkpoint now spares other open sessions' batons. Startup delivery (startup_handoff) skips an open source that captured anything in the last ten minutes, and the claim re-checks it in its own transaction. The claim's sweep retires older batons of quiet open sessions but spares a busy one; an explicit accept keeps sparing every open session.
This commit is contained in:
@@ -165,6 +165,15 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
|
||||
- Native `ai-memory upgrade` container refusal now uses the shared
|
||||
`running_in_container` helper (`AI_MEMORY_IN_CONTAINER`, `/.dockerenv`,
|
||||
`/run/.containerenv` / Podman), matching staged-hooks detection. (#802)
|
||||
- Parallel OpenCode 2 sessions in one directory no longer receive each
|
||||
other's context. Each completed turn's checkpoint retired the automatic
|
||||
handoffs of every other live session there, and the next session to start
|
||||
was handed the handoff of a session still in use (mid-turn, or seconds after
|
||||
its last turn). A checkpoint now spares other live sessions' handoffs; a
|
||||
starting session receives a live session's handoff only once that session
|
||||
has captured nothing for ten minutes, re-checked inside the claim; and the
|
||||
claim retires older handoffs of quiet sessions but not of one in use.
|
||||
(#883)
|
||||
- `memory_query`'s vector stream called the generic `Embedder::embed`
|
||||
instead of `embed_query` on the configured embedder, so a
|
||||
query/document-asymmetric embedder (Google's task-typed embeddings, or
|
||||
|
||||
@@ -69,6 +69,14 @@ pub const MAX_HOOK_BATCH_ITEMS: usize = 256;
|
||||
const AUTOMATIC_HANDOFF_ADMISSION_TIMEOUT: std::time::Duration =
|
||||
std::time::Duration::from_millis(750);
|
||||
|
||||
/// How long a still-open session must be quiet before its turn-checkpoint
|
||||
/// baton is handed to a new session. Nothing distinguishes a closed terminal
|
||||
/// from a parallel session in the same directory; a session that captured
|
||||
/// anything this recently is treated as in use. The live incident this
|
||||
/// guards against delivered batons from sessions 74 seconds and 0 seconds
|
||||
/// (mid-turn) away from their last event.
|
||||
const LIVE_BATON_QUIET_PERIOD: jiff::SignedDuration = jiff::SignedDuration::from_mins(10);
|
||||
|
||||
/// Maximum cwd-resolution cache entries kept per server process. The cache is
|
||||
/// an optimization only; evicted entries are re-resolved through the writer.
|
||||
pub const DEFAULT_PROJECT_CACHE_MAX_ENTRIES: usize = 4096;
|
||||
@@ -1249,6 +1257,16 @@ async fn fetch_and_accept_handoff(
|
||||
query: HandoffQuery,
|
||||
actor: Option<IdentityKey>,
|
||||
skip_webhooks: Vec<String>,
|
||||
) -> anyhow::Result<Option<String>> {
|
||||
fetch_and_accept_handoff_at(state, query, actor, skip_webhooks, jiff::Timestamp::now()).await
|
||||
}
|
||||
|
||||
async fn fetch_and_accept_handoff_at(
|
||||
state: &HookState,
|
||||
query: HandoffQuery,
|
||||
actor: Option<IdentityKey>,
|
||||
skip_webhooks: Vec<String>,
|
||||
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
|
||||
@@ -1290,9 +1308,20 @@ async fn fetch_and_accept_handoff(
|
||||
Some(key) => ai_memory_core::OwnerFilter::User(key.storage_key()),
|
||||
None => ai_memory_core::OwnerFilter::Unattributed,
|
||||
};
|
||||
// A duration never fails here; the fallback keeps every live session's
|
||||
// baton, the conservative side.
|
||||
let busy_since = now
|
||||
.saturating_sub(LIVE_BATON_QUIET_PERIOD)
|
||||
.unwrap_or(jiff::Timestamp::MIN);
|
||||
let handoff = state
|
||||
.reader
|
||||
.latest_open_handoff(ws, proj, query.cwd.clone(), owner_filter.clone())
|
||||
.startup_handoff(
|
||||
ws,
|
||||
proj,
|
||||
query.cwd.clone(),
|
||||
owner_filter.clone(),
|
||||
busy_since,
|
||||
)
|
||||
.await?;
|
||||
let handoff_md = handoff.as_ref().map(render_handoff_markdown);
|
||||
// The brief is additive and non-destructive: unlike the handoff (a
|
||||
@@ -1390,6 +1419,7 @@ async fn fetch_and_accept_handoff(
|
||||
}),
|
||||
managed.as_ref().map(|managed| managed.run_id),
|
||||
receiving_session,
|
||||
busy_since,
|
||||
)
|
||||
.await?
|
||||
} else {
|
||||
@@ -8855,6 +8885,9 @@ mod tests {
|
||||
}
|
||||
}
|
||||
}
|
||||
// Both sources are still open, so their batons wait out the quiet period.
|
||||
let quiet =
|
||||
jiff::Timestamp::now() + LIVE_BATON_QUIET_PERIOD + jiff::SignedDuration::from_secs(1);
|
||||
for (project, other) in [("alpha", "beta"), ("beta", "alpha")] {
|
||||
let query = HandoffQuery {
|
||||
workspace: Some("checkpoint-workspace".into()),
|
||||
@@ -8863,13 +8896,13 @@ mod tests {
|
||||
session_id: Some(SessionId::new().to_string()),
|
||||
..Default::default()
|
||||
};
|
||||
let first = fetch_and_accept_handoff(&state, query.clone(), None, Vec::new())
|
||||
let first = fetch_and_accept_handoff_at(&state, query.clone(), None, Vec::new(), quiet)
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
assert!(first.contains(&format!("Latest assistant response: {project}-latest")));
|
||||
assert!(!first.contains(&format!("{other}-")));
|
||||
let second = fetch_and_accept_handoff(&state, query, None, Vec::new())
|
||||
let second = fetch_and_accept_handoff_at(&state, query, None, Vec::new(), quiet)
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(
|
||||
@@ -8879,6 +8912,89 @@ mod tests {
|
||||
}
|
||||
}
|
||||
|
||||
// Live incident: parallel OpenCode sessions in one directory. Every
|
||||
// completed turn retired the other sessions' batons, and the next session
|
||||
// to start was handed the baton of a session still in use (once mid-turn,
|
||||
// once 74 seconds after its last turn), carrying another conversation.
|
||||
#[tokio::test]
|
||||
async fn opencode_parallel_live_batons_are_owned_and_wait_for_quiet() {
|
||||
let tmp = TempDir::new().unwrap();
|
||||
let state = make_state(&tmp).await;
|
||||
let (alpha, beta) = (SessionId::new().to_string(), SessionId::new().to_string());
|
||||
for (session, text) in [(&alpha, "alpha work"), (&beta, "beta work")] {
|
||||
for event in ["user-prompt", "stop"] {
|
||||
let mut env = opencode_turn_event(session, event, text);
|
||||
crate::assistant_capture::apply_assistant_backstop(&mut env, true);
|
||||
process(&state, env, None, Vec::new()).await.unwrap();
|
||||
}
|
||||
}
|
||||
let open_batons = || async {
|
||||
state
|
||||
.reader
|
||||
.list_handoffs(
|
||||
state.workspace_id,
|
||||
state.project_id,
|
||||
None,
|
||||
ai_memory_core::OwnerFilter::Any,
|
||||
10,
|
||||
)
|
||||
.await
|
||||
.unwrap()
|
||||
.into_iter()
|
||||
.filter(|h| h.lifecycle.state == ai_memory_core::HandoffState::Open)
|
||||
.count()
|
||||
};
|
||||
assert_eq!(
|
||||
open_batons().await,
|
||||
2,
|
||||
"a turn must not retire another live session's baton"
|
||||
);
|
||||
|
||||
// Alpha is mid-turn again; beta just finished one. `between` separates
|
||||
// beta's last event from alpha's newest one.
|
||||
tokio::time::sleep(std::time::Duration::from_millis(2)).await;
|
||||
let between = jiff::Timestamp::now();
|
||||
tokio::time::sleep(std::time::Duration::from_millis(2)).await;
|
||||
process(
|
||||
&state,
|
||||
opencode_turn_event(&alpha, "user-prompt", "alpha continues"),
|
||||
None,
|
||||
Vec::new(),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
let receiver = || HandoffQuery {
|
||||
agent: Some("opencode2".into()),
|
||||
session_id: Some(SessionId::new().to_string()),
|
||||
..Default::default()
|
||||
};
|
||||
let busy = fetch_and_accept_handoff_at(
|
||||
&state,
|
||||
receiver(),
|
||||
None,
|
||||
Vec::new(),
|
||||
jiff::Timestamp::now(),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(busy.is_none(), "a session in use keeps its baton: {busy:?}");
|
||||
assert_eq!(open_batons().await, 2);
|
||||
|
||||
// Ten minutes after `between`: beta has been quiet, alpha has not.
|
||||
let later = between + LIVE_BATON_QUIET_PERIOD;
|
||||
let delivered = fetch_and_accept_handoff_at(&state, receiver(), None, Vec::new(), later)
|
||||
.await
|
||||
.unwrap()
|
||||
.expect("a quiet live session's baton is deliverable");
|
||||
assert!(delivered.contains("beta work"), "{delivered}");
|
||||
assert!(!delivered.contains("alpha"), "{delivered}");
|
||||
assert_eq!(
|
||||
open_batons().await,
|
||||
1,
|
||||
"claiming one baton must not sweep the baton of a session in use"
|
||||
);
|
||||
}
|
||||
|
||||
// Live incident: a `.ai-memory.toml` naming a workspace appeared under
|
||||
// running OpenCode sessions, so their later events (same cwd) resolved to
|
||||
// a new `(workspace, project)`. Scope is not identity: those events are
|
||||
@@ -10973,7 +11089,7 @@ mod tests {
|
||||
})
|
||||
.await
|
||||
.unwrap();
|
||||
state
|
||||
let id = state
|
||||
.writer
|
||||
.insert_handoff(NewHandoff {
|
||||
workspace_id: state.workspace_id,
|
||||
@@ -10989,7 +11105,10 @@ mod tests {
|
||||
owner_user: None,
|
||||
})
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
// A SessionEnd baton: its source is over.
|
||||
state.writer.end_session(session_id, None).await.unwrap();
|
||||
id
|
||||
}
|
||||
|
||||
let stale = insert_auto(&state, "/repo/api", "STALE-SPECIFIC").await;
|
||||
|
||||
@@ -5918,6 +5918,7 @@ mod tests {
|
||||
Some(acceptance(first_handoff, None)),
|
||||
Some(run.run_id),
|
||||
None,
|
||||
jiff::Timestamp::now(),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
@@ -5944,6 +5945,7 @@ mod tests {
|
||||
Some(acceptance(second_handoff, None)),
|
||||
Some(run.run_id),
|
||||
None,
|
||||
jiff::Timestamp::now(),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
@@ -6010,6 +6012,7 @@ mod tests {
|
||||
Some(acceptance(selected_auto, Some("/repo/api/src".into()))),
|
||||
Some(run.run_id),
|
||||
None,
|
||||
jiff::Timestamp::now(),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
@@ -2949,6 +2949,10 @@ fn handoff_fields(h: &NewHandoff) -> StoreResult<HandoffFields> {
|
||||
/// boundary: the same directory inside a shared container is the norm, so
|
||||
/// owner equality is the only thing keeping one operator's SessionEnd from
|
||||
/// retiring another's pending baton. `keep` spares the baton being refreshed.
|
||||
///
|
||||
/// Another session that is still open owns its turn-checkpoint baton and
|
||||
/// refreshes it itself, so it is spared too: with parallel sessions in one
|
||||
/// directory, every completed turn would otherwise retire the others' batons.
|
||||
fn expire_same_cwd_auto_handoffs(
|
||||
conn: &Transaction<'_>,
|
||||
h: &NewHandoff,
|
||||
@@ -2961,13 +2965,17 @@ fn expire_same_cwd_auto_handoffs(
|
||||
WHERE workspace_id = ?1 AND project_id = ?2 \
|
||||
AND state = 'open' AND from_session_id IS NOT NULL \
|
||||
AND (cwd = ?3 OR (cwd IS NULL AND ?3 IS NULL)) \
|
||||
AND owner_user IS ?4 AND id IS NOT ?5",
|
||||
AND owner_user IS ?4 AND id IS NOT ?5 \
|
||||
AND (from_session_id IS ?6 OR NOT EXISTS ( \
|
||||
SELECT 1 FROM sessions s \
|
||||
WHERE s.id = handoffs.from_session_id AND s.ended_at IS NULL))",
|
||||
params![
|
||||
h.workspace_id.as_bytes(),
|
||||
h.project_id.as_bytes(),
|
||||
cwd,
|
||||
h.owner_user.as_deref(),
|
||||
keep.map(HandoffId::as_bytes)
|
||||
keep.map(HandoffId::as_bytes),
|
||||
h.from_session_id.as_ref().map(|s| &s.as_bytes()[..])
|
||||
],
|
||||
)?;
|
||||
if expired > 0 {
|
||||
@@ -3082,14 +3090,22 @@ fn insert_handoff_row(conn: &Transaction<'_>, h: &NewHandoff) -> StoreResult<Han
|
||||
/// handoffs.
|
||||
pub fn accept_handoff(conn: &mut Connection, acceptance: &HandoffAcceptance) -> StoreResult<bool> {
|
||||
let tx = conn.transaction()?;
|
||||
let claimed = accept_handoff_in_transaction(&tx, acceptance)?;
|
||||
let claimed = accept_handoff_in_transaction(&tx, acceptance, None)?;
|
||||
tx.commit()?;
|
||||
Ok(claimed)
|
||||
}
|
||||
|
||||
/// Claim one handoff inside `tx`.
|
||||
///
|
||||
/// `busy_since` (microseconds) is the automatic-delivery cutoff: when set, a
|
||||
/// baton whose source session is still open and captured anything after it is
|
||||
/// not claimed, and the post-claim sweep retires older batons of quiet open
|
||||
/// sessions too. `None` (an explicit accept) claims regardless and spares
|
||||
/// every open session's baton from the sweep.
|
||||
pub(crate) fn accept_handoff_in_transaction(
|
||||
tx: &Transaction<'_>,
|
||||
acceptance: &HandoffAcceptance,
|
||||
busy_since: Option<i64>,
|
||||
) -> StoreResult<bool> {
|
||||
let HandoffAcceptance {
|
||||
handoff_id,
|
||||
@@ -3185,13 +3201,18 @@ pub(crate) fn accept_handoff_in_transaction(
|
||||
}
|
||||
let metadata = tx
|
||||
.query_row(
|
||||
"SELECT from_session_id IS NOT NULL, cwd, created_at, owner_user \
|
||||
"SELECT from_session_id IS NOT NULL, cwd, created_at, owner_user, \
|
||||
EXISTS (SELECT 1 FROM sessions s \
|
||||
WHERE s.id = handoffs.from_session_id AND s.ended_at IS NULL \
|
||||
AND EXISTS (SELECT 1 FROM observations o \
|
||||
WHERE o.session_id = s.id AND o.created_at > ?4)) \
|
||||
FROM handoffs \
|
||||
WHERE id = ?1 AND workspace_id = ?2 AND project_id = ?3 AND state = 'open'",
|
||||
params![
|
||||
handoff_id.as_bytes(),
|
||||
workspace_id.as_bytes(),
|
||||
project_id.as_bytes(),
|
||||
busy_since.unwrap_or(i64::MAX),
|
||||
],
|
||||
|row| {
|
||||
Ok((
|
||||
@@ -3199,13 +3220,20 @@ pub(crate) fn accept_handoff_in_transaction(
|
||||
row.get::<_, Option<String>>(1)?,
|
||||
row.get::<_, i64>(2)?,
|
||||
row.get::<_, Option<String>>(3)?,
|
||||
row.get::<_, bool>(4)?,
|
||||
))
|
||||
},
|
||||
)
|
||||
.optional()?;
|
||||
let Some((automatic, cwd, created_at, owner_user)) = metadata else {
|
||||
let Some((automatic, cwd, created_at, owner_user, source_busy)) = metadata else {
|
||||
return Ok(false);
|
||||
};
|
||||
// Startup selection happens on a reader before this claim; the source may
|
||||
// have resumed (or refreshed this very baton) in between. Re-checked here,
|
||||
// in the claim's transaction, so a session in use never loses its baton.
|
||||
if source_busy {
|
||||
return Ok(false);
|
||||
}
|
||||
// The ownership check rides along in the UPDATE's WHERE rather than being a
|
||||
// separate read: the claim stays a single atomic compare-and-set (only one
|
||||
// racing session can flip 'open' -> 'accepted'), and a caller who is not
|
||||
@@ -3274,6 +3302,7 @@ pub(crate) fn accept_handoff_in_transaction(
|
||||
created_at,
|
||||
receiving_cwd.as_deref(),
|
||||
owner_user.as_deref(),
|
||||
busy_since,
|
||||
)?;
|
||||
if expired > 0 {
|
||||
audit(
|
||||
@@ -3307,6 +3336,12 @@ fn validate_identity_storage_key(value: Option<&str>, label: &str) -> StoreResul
|
||||
/// behalf would take it away from Bob too. Equality on `owner_user` keeps the
|
||||
/// unattributed single-operator case (every row NULL) behaving exactly as it
|
||||
/// does without ownership.
|
||||
///
|
||||
/// A session that is still open and in use keeps its baton: it belongs to work
|
||||
/// in progress and is refreshed by it. "In use" means it captured anything
|
||||
/// after `busy_since`; without a cutoff every open session is spared. A quiet
|
||||
/// open session is superseded like an ended one, so abandoned conversations
|
||||
/// that never end cannot pile up and surface one by one to later sessions.
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
fn expire_superseded_auto_handoffs(
|
||||
tx: &Transaction<'_>,
|
||||
@@ -3317,6 +3352,7 @@ fn expire_superseded_auto_handoffs(
|
||||
accepted_created_at: i64,
|
||||
receiving_cwd: Option<&str>,
|
||||
accepted_owner: Option<&str>,
|
||||
busy_since: Option<i64>,
|
||||
) -> StoreResult<usize> {
|
||||
let receiving_cwd = receiving_cwd.or(accepted_cwd);
|
||||
let accepted_key = crate::reader::handoff_selection_key(
|
||||
@@ -3329,15 +3365,22 @@ fn expire_superseded_auto_handoffs(
|
||||
"SELECT id, cwd, created_at FROM handoffs \
|
||||
WHERE workspace_id = ?1 AND project_id = ?2 \
|
||||
AND state = 'open' AND from_session_id IS NOT NULL \
|
||||
AND owner_user IS ?3",
|
||||
AND owner_user IS ?3 \
|
||||
AND NOT EXISTS (SELECT 1 FROM sessions s \
|
||||
WHERE s.id = handoffs.from_session_id AND s.ended_at IS NULL \
|
||||
AND (?4 IS NULL OR EXISTS (SELECT 1 FROM observations o \
|
||||
WHERE o.session_id = s.id AND o.created_at > ?4)))",
|
||||
)?;
|
||||
let rows = stmt.query_map(
|
||||
params![workspace_id, project_id, accepted_owner, busy_since],
|
||||
|row| {
|
||||
Ok((
|
||||
row.get::<_, Vec<u8>>(0)?,
|
||||
row.get::<_, Option<String>>(1)?,
|
||||
row.get::<_, i64>(2)?,
|
||||
))
|
||||
},
|
||||
)?;
|
||||
let rows = stmt.query_map(params![workspace_id, project_id, accepted_owner], |row| {
|
||||
Ok((
|
||||
row.get::<_, Vec<u8>>(0)?,
|
||||
row.get::<_, Option<String>>(1)?,
|
||||
row.get::<_, i64>(2)?,
|
||||
))
|
||||
})?;
|
||||
let mut ids = Vec::new();
|
||||
for row in rows {
|
||||
let (id_bytes, cwd, created_at) = row?;
|
||||
|
||||
@@ -5538,19 +5538,74 @@ impl ReaderPool {
|
||||
project_id: ProjectId,
|
||||
cwd_filter: Option<String>,
|
||||
owner_filter: OwnerFilter,
|
||||
) -> StoreResult<Option<Handoff>> {
|
||||
self.open_handoff_for(workspace_id, project_id, cwd_filter, owner_filter, None)
|
||||
.await
|
||||
}
|
||||
|
||||
/// [`Self::latest_open_handoff`] for automatic delivery to a starting
|
||||
/// session: the baton of a session that is still open and captured
|
||||
/// anything after `busy_since` is not a candidate.
|
||||
///
|
||||
/// A turn checkpoint publishes a live session's baton, and nothing tells
|
||||
/// a closed terminal apart from a parallel session still at work in the
|
||||
/// same directory. Handing the baton of a session in use to a new one
|
||||
/// gives that session another conversation's context and consumes the
|
||||
/// baton its own successor needed. Once the session has been quiet since
|
||||
/// `busy_since`, its baton is deliverable again; ended sessions and
|
||||
/// manual handoffs are unaffected.
|
||||
///
|
||||
/// # Errors
|
||||
/// Propagates any SQL or pool error.
|
||||
pub async fn startup_handoff(
|
||||
&self,
|
||||
workspace_id: WorkspaceId,
|
||||
project_id: ProjectId,
|
||||
cwd_filter: Option<String>,
|
||||
owner_filter: OwnerFilter,
|
||||
busy_since: Timestamp,
|
||||
) -> StoreResult<Option<Handoff>> {
|
||||
self.open_handoff_for(
|
||||
workspace_id,
|
||||
project_id,
|
||||
cwd_filter,
|
||||
owner_filter,
|
||||
Some(busy_since.as_microsecond()),
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
async fn open_handoff_for(
|
||||
&self,
|
||||
workspace_id: WorkspaceId,
|
||||
project_id: ProjectId,
|
||||
cwd_filter: Option<String>,
|
||||
owner_filter: OwnerFilter,
|
||||
busy_since: Option<i64>,
|
||||
) -> StoreResult<Option<Handoff>> {
|
||||
self.with_conn(move |conn| {
|
||||
// Ownership belongs in the query, not just in
|
||||
// `is_handoff_candidate`: prompt-derived fields from another
|
||||
// operator must never be loaded or deserialized for this caller.
|
||||
let (owner_clause, owner_param) = handoff_owner_sql(&owner_filter, 3);
|
||||
let busy_index = if owner_param.is_some() { 4 } else { 3 };
|
||||
let busy_clause = if busy_since.is_some() {
|
||||
format!(
|
||||
" AND NOT EXISTS (SELECT 1 FROM sessions s \
|
||||
WHERE s.id = handoffs.from_session_id AND s.ended_at IS NULL \
|
||||
AND EXISTS (SELECT 1 FROM observations o \
|
||||
WHERE o.session_id = s.id AND o.created_at > ?{busy_index}))"
|
||||
)
|
||||
} else {
|
||||
String::new()
|
||||
};
|
||||
let sql = format!(
|
||||
"SELECT id, workspace_id, project_id, from_session_id, from_agent, to_agent, \
|
||||
cwd, summary, open_questions, next_steps, files_touched, state, \
|
||||
created_at, accepted_by, accepted_at, accepted_by_session, \
|
||||
owner_user, accepted_by_user \
|
||||
FROM handoffs \
|
||||
WHERE workspace_id = ?1 AND project_id = ?2 AND state = 'open'{owner_clause} \
|
||||
WHERE workspace_id = ?1 AND project_id = ?2 AND state = 'open'{owner_clause}{busy_clause} \
|
||||
ORDER BY created_at DESC"
|
||||
);
|
||||
let mut stmt = conn.prepare(&sql)?;
|
||||
@@ -5559,6 +5614,9 @@ impl ReaderPool {
|
||||
if let Some(owner) = owner_param.as_ref() {
|
||||
binds.push(owner);
|
||||
}
|
||||
if let Some(since) = busy_since.as_ref() {
|
||||
binds.push(since);
|
||||
}
|
||||
let rows = stmt.query_map(binds.as_slice(), row_to_handoff)?;
|
||||
let mut selected: Option<Handoff> = None;
|
||||
for r in rows {
|
||||
@@ -10564,7 +10622,7 @@ mod tests {
|
||||
})
|
||||
.await
|
||||
.unwrap();
|
||||
store
|
||||
let id = store
|
||||
.writer
|
||||
.insert_handoff(NewHandoff {
|
||||
workspace_id,
|
||||
@@ -10580,7 +10638,10 @@ mod tests {
|
||||
owner_user: None,
|
||||
})
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
// A SessionEnd baton: its source is over.
|
||||
store.writer.end_session(session_id, None).await.unwrap();
|
||||
id
|
||||
}
|
||||
|
||||
let superseded_same_cwd =
|
||||
|
||||
@@ -644,6 +644,7 @@ pub(crate) enum WriteCmd {
|
||||
handoff: Option<HandoffAcceptance>,
|
||||
managed_run_id: Option<ManagedRunId>,
|
||||
receiving_session: Option<NewSession>,
|
||||
busy_since: jiff::Timestamp,
|
||||
reply: oneshot::Sender<StoreResult<StartupContextAcceptance>>,
|
||||
},
|
||||
FinishWorkstreamRun {
|
||||
@@ -2627,18 +2628,24 @@ impl WriterHandle {
|
||||
/// SessionStart response.
|
||||
///
|
||||
/// When a managed run was requested but is no longer claimable, the
|
||||
/// handoff remains open and both result fields are false.
|
||||
/// handoff remains open and both result fields are false. `busy_since` is
|
||||
/// the same cutoff the selection used
|
||||
/// ([`crate::ReaderPool::startup_handoff`]), re-applied in the claim's
|
||||
/// transaction: a baton whose open source captured anything after it stays
|
||||
/// open.
|
||||
pub async fn accept_startup_context(
|
||||
&self,
|
||||
handoff: Option<HandoffAcceptance>,
|
||||
managed_run_id: Option<ManagedRunId>,
|
||||
receiving_session: Option<NewSession>,
|
||||
busy_since: jiff::Timestamp,
|
||||
) -> StoreResult<StartupContextAcceptance> {
|
||||
let (tx, rx) = oneshot::channel();
|
||||
self.send(WriteCmd::AcceptStartupContext {
|
||||
handoff,
|
||||
managed_run_id,
|
||||
receiving_session,
|
||||
busy_since,
|
||||
reply: tx,
|
||||
})
|
||||
.await?;
|
||||
@@ -3655,6 +3662,7 @@ fn worker_loop(mut conn: Connection, mut rx: mpsc::Receiver<WriteCmd>) {
|
||||
handoff,
|
||||
managed_run_id,
|
||||
receiving_session,
|
||||
busy_since,
|
||||
reply,
|
||||
} => {
|
||||
let result = (|| {
|
||||
@@ -3683,7 +3691,11 @@ fn worker_loop(mut conn: Connection, mut rx: mpsc::Receiver<WriteCmd>) {
|
||||
return Ok(StartupContextAcceptance::default());
|
||||
}
|
||||
let handoff_accepted = match handoff {
|
||||
Some(acceptance) => ops::accept_handoff_in_transaction(&tx, &acceptance)?,
|
||||
Some(acceptance) => ops::accept_handoff_in_transaction(
|
||||
&tx,
|
||||
&acceptance,
|
||||
Some(busy_since.as_microsecond()),
|
||||
)?,
|
||||
None => false,
|
||||
};
|
||||
tx.commit()?;
|
||||
|
||||
@@ -769,7 +769,7 @@ async fn automatic_supersession_does_not_reach_across_operators() {
|
||||
})
|
||||
.await
|
||||
.unwrap();
|
||||
store
|
||||
let id = store
|
||||
.writer
|
||||
.insert_handoff(NewHandoff {
|
||||
workspace_id: ws,
|
||||
@@ -785,7 +785,9 @@ async fn automatic_supersession_does_not_reach_across_operators() {
|
||||
owner_user: stamp,
|
||||
})
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
store.writer.end_session(session_id, None).await.unwrap();
|
||||
id
|
||||
}
|
||||
|
||||
let bob = auto_handoff(&store, ws, proj, "bob", "bob's baton").await;
|
||||
|
||||
@@ -17,8 +17,9 @@
|
||||
//! core capability rather than a detail.
|
||||
|
||||
use ai_memory_core::{
|
||||
ActorContext, AgentKind, HandoffAcceptance, IdentityKey, NewHandoff, NewPage, NewSession,
|
||||
NewUser, OwnerFilter, PagePath, ProjectId, SessionId, Tier, UserRole, WorkspaceId, owner_stamp,
|
||||
ActorContext, AgentKind, HandoffAcceptance, HandoffState, IdentityKey, NewHandoff, NewPage,
|
||||
NewSession, NewUser, OwnerFilter, PagePath, ProjectId, SessionId, Tier, UserRole, WorkspaceId,
|
||||
owner_stamp,
|
||||
};
|
||||
use ai_memory_store::Store;
|
||||
|
||||
@@ -517,3 +518,217 @@ async fn accept_reopens_an_ended_receiver_session_but_keeps_claim_once() {
|
||||
"the resurrection path must not let a second session steal an accepted baton"
|
||||
);
|
||||
}
|
||||
|
||||
/// Parallel live sessions in one directory each own a turn-checkpoint baton.
|
||||
/// One session's checkpoint, and a receiver claiming it, must leave the other
|
||||
/// live session's baton open: before this held, every completed turn retired
|
||||
/// the other sessions' batons and the survivor went to whoever started next.
|
||||
#[tokio::test]
|
||||
async fn parallel_live_sessions_keep_their_own_checkpoint_batons() {
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
let store = Store::open(tmp.path()).unwrap();
|
||||
let (ws, proj) = scope(&store).await;
|
||||
let baton = |session: SessionId, summary: &str| NewHandoff {
|
||||
workspace_id: ws,
|
||||
project_id: proj,
|
||||
from_session_id: Some(session),
|
||||
from_agent: AgentKind::OpenCode,
|
||||
to_agent: None,
|
||||
cwd: Some("/repo".into()),
|
||||
summary: summary.into(),
|
||||
open_questions: Vec::new(),
|
||||
next_steps: Vec::new(),
|
||||
files_touched: Vec::new(),
|
||||
owner_user: None,
|
||||
};
|
||||
let alpha = open_session(&store, ws, proj, AgentKind::OpenCode).await;
|
||||
let beta = open_session(&store, ws, proj, AgentKind::OpenCode).await;
|
||||
let alpha_baton = store
|
||||
.writer
|
||||
.checkpoint_session_handoff(baton(alpha, "alpha"))
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
tokio::time::sleep(std::time::Duration::from_millis(2)).await;
|
||||
let beta_baton = store
|
||||
.writer
|
||||
.checkpoint_session_handoff(baton(beta, "beta"))
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
let state = |id| {
|
||||
let reader = store.reader.clone();
|
||||
async move {
|
||||
reader
|
||||
.handoff_by_id(id)
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap()
|
||||
.lifecycle
|
||||
.state
|
||||
}
|
||||
};
|
||||
assert_eq!(state(alpha_baton).await, HandoffState::Open);
|
||||
|
||||
let receiver = open_session(&store, ws, proj, AgentKind::OpenCode).await;
|
||||
let claimed = store
|
||||
.writer
|
||||
.accept_handoff(HandoffAcceptance {
|
||||
handoff_id: beta_baton,
|
||||
workspace_id: ws,
|
||||
project_id: proj,
|
||||
accepting_agent: AgentKind::OpenCode,
|
||||
accepting_session: Some(receiver),
|
||||
accepting_user: None,
|
||||
owner_filter: OwnerFilter::Any,
|
||||
receiving_cwd: Some("/repo".into()),
|
||||
})
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(claimed);
|
||||
assert_eq!(
|
||||
state(alpha_baton).await,
|
||||
HandoffState::Open,
|
||||
"claiming one live session's baton must not sweep another's"
|
||||
);
|
||||
|
||||
// Once alpha ends, its baton is an ordinary SessionEnd baton again and the
|
||||
// same-cwd supersession applies to it.
|
||||
store.writer.end_session(alpha, None).await.unwrap();
|
||||
let gamma = open_session(&store, ws, proj, AgentKind::OpenCode).await;
|
||||
store
|
||||
.writer
|
||||
.checkpoint_session_handoff(baton(gamma, "gamma"))
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
assert_eq!(state(alpha_baton).await, HandoffState::Expired);
|
||||
}
|
||||
|
||||
/// Automatic delivery re-checks the source inside the claim, and the claim's
|
||||
/// sweep retires batons of quiet (abandoned) open sessions while sparing one
|
||||
/// still in use. The selection runs on a reader before the writer claims, so
|
||||
/// a source can resume in between; and OpenCode sessions that never end would
|
||||
/// otherwise leave one open baton each, surfacing older conversations one by
|
||||
/// one to later sessions.
|
||||
#[tokio::test]
|
||||
async fn startup_claim_rechecks_the_source_and_sweeps_only_quiet_open_batons() {
|
||||
use ai_memory_core::{NewObservation, ObservationKind, Sanitized, Sanitizer};
|
||||
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
let store = Store::open(tmp.path()).unwrap();
|
||||
let (ws, proj) = scope(&store).await;
|
||||
let mut batons = Vec::new();
|
||||
// Checkpoint order, oldest first: busy, abandoned, quiet.
|
||||
for summary in ["busy", "abandoned", "quiet"] {
|
||||
let session = open_session(&store, ws, proj, AgentKind::OpenCode).await;
|
||||
store
|
||||
.writer
|
||||
.insert_observation(Sanitized::new(
|
||||
NewObservation {
|
||||
session_id: session,
|
||||
workspace_id: ws,
|
||||
project_id: proj,
|
||||
kind: ObservationKind::UserPrompt,
|
||||
extension: None,
|
||||
source_event: None,
|
||||
title: "prompt".into(),
|
||||
body: summary.into(),
|
||||
importance: 5,
|
||||
},
|
||||
&Sanitizer::builtin(),
|
||||
))
|
||||
.await
|
||||
.unwrap();
|
||||
let id = store
|
||||
.writer
|
||||
.checkpoint_session_handoff(NewHandoff {
|
||||
workspace_id: ws,
|
||||
project_id: proj,
|
||||
from_session_id: Some(session),
|
||||
from_agent: AgentKind::OpenCode,
|
||||
to_agent: None,
|
||||
cwd: Some("/repo".into()),
|
||||
summary: summary.into(),
|
||||
open_questions: Vec::new(),
|
||||
next_steps: Vec::new(),
|
||||
files_touched: Vec::new(),
|
||||
owner_user: None,
|
||||
})
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
batons.push((session, id));
|
||||
tokio::time::sleep(std::time::Duration::from_millis(2)).await;
|
||||
}
|
||||
let [
|
||||
(_, busy),
|
||||
(abandoned_session, abandoned),
|
||||
(quiet_session, quiet),
|
||||
] = batons[..]
|
||||
else {
|
||||
unreachable!()
|
||||
};
|
||||
// Only "busy" captured anything in the last hour.
|
||||
let hour_ago = (jiff::Timestamp::now() - jiff::SignedDuration::from_hours(1)).as_microsecond();
|
||||
let conn = rusqlite::Connection::open(store.db_path()).unwrap();
|
||||
for session in [abandoned_session, quiet_session] {
|
||||
conn.execute(
|
||||
"UPDATE observations SET created_at = ?1 WHERE session_id = ?2",
|
||||
rusqlite::params![hour_ago, session.as_bytes()],
|
||||
)
|
||||
.unwrap();
|
||||
}
|
||||
drop(conn);
|
||||
let cutoff = jiff::Timestamp::now() - jiff::SignedDuration::from_mins(10);
|
||||
let claim = |handoff_id, receiver| HandoffAcceptance {
|
||||
handoff_id,
|
||||
workspace_id: ws,
|
||||
project_id: proj,
|
||||
accepting_agent: AgentKind::OpenCode,
|
||||
accepting_session: Some(receiver),
|
||||
accepting_user: None,
|
||||
owner_filter: OwnerFilter::Any,
|
||||
receiving_cwd: Some("/repo".into()),
|
||||
};
|
||||
let state = |id| {
|
||||
let reader = store.reader.clone();
|
||||
async move {
|
||||
reader
|
||||
.handoff_by_id(id)
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap()
|
||||
.lifecycle
|
||||
.state
|
||||
}
|
||||
};
|
||||
|
||||
// Selected earlier, but its source is in use by the time of the claim.
|
||||
let receiver = open_session(&store, ws, proj, AgentKind::OpenCode).await;
|
||||
let raced = store
|
||||
.writer
|
||||
.accept_startup_context(Some(claim(busy, receiver)), None, None, cutoff)
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(!raced.handoff_accepted, "a source in use keeps its baton");
|
||||
assert_eq!(state(busy).await, HandoffState::Open);
|
||||
|
||||
let receiver = open_session(&store, ws, proj, AgentKind::OpenCode).await;
|
||||
let delivered = store
|
||||
.writer
|
||||
.accept_startup_context(Some(claim(quiet, receiver)), None, None, cutoff)
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(delivered.handoff_accepted);
|
||||
assert_eq!(
|
||||
state(abandoned).await,
|
||||
HandoffState::Expired,
|
||||
"an older baton of a quiet open session is superseded"
|
||||
);
|
||||
assert_eq!(
|
||||
state(busy).await,
|
||||
HandoffState::Open,
|
||||
"a session in use keeps its baton even when it is older"
|
||||
);
|
||||
}
|
||||
|
||||
+6
-1
@@ -1218,7 +1218,12 @@ The generated plugin targets the OpenCode 2.0.10+ event API (checked against
|
||||
terminal is not a session end: each completed root turn instead writes a
|
||||
deterministic checkpoint (no LLM call) of `sessions/<id>.md` and refreshes the
|
||||
session's automatic handoff, keeping one open baton per live session. The next
|
||||
session claims the latest checkpoint. Startup context is claimed once per root
|
||||
session claims the latest checkpoint of a session that has captured nothing for
|
||||
ten minutes. Nothing tells a closed terminal from a parallel session still at
|
||||
work in the same directory, so this is a heuristic: a session in use keeps its
|
||||
baton, one session's turn never retires another live session's baton, and a
|
||||
claim retires only older batons of quiet sessions. A tool or model call that
|
||||
runs longer than ten minutes without a captured event looks quiet. Startup context is claimed once per root
|
||||
session and retained on every later model request; child sessions never claim
|
||||
it or publish a baton.
|
||||
|
||||
|
||||
@@ -29,6 +29,7 @@ boundary not yet built.
|
||||
| 4b | Handoff owner-scoped recovery (`any_owner` admin gate) | `ai-memory-mcp/src/server.rs` `require_admin_capability` on `memory_handoff_accept`/`_cancel` `any_owner` | `handoff_admission.rs` — cancel gate + **accept gate** (adversarial: non-admin `any_owner` accept refused) | STRONG |
|
||||
| 4c | Turn-checkpoint baton is owner+scope-bound | `ai-memory-store/src/ops.rs` `checkpoint_session_handoff` — reads the session's own `(workspace, project, owner_user)` from `sessions WHERE id=?1 AND ended_at IS NULL`, refuses on scope/owner mismatch, refreshes only that session's own `state='open'` baton (id preserved), returns `None` when the session already ended | `ops.rs` `checkpoint_session_handoff_refreshes_only_the_live_sessions_own_baton` (#865) | STRONG |
|
||||
| 4d | Native session rebind / drifted-end is owner+agent+cwd CAS | `ai-memory-store/src/ops.rs` — `session.moved` rebind gated on `agent==OpenCode` + `ended_at IS NULL` + lexical `normalize_cwd(stored)==normalize_cwd(from)`; `expire_same_cwd_auto_handoffs` AND-gated on `owner_user IS ?4 AND id IS NOT ?5`; drifted `SessionEnd` ends only same owner+agent+cwd | `ops.rs` `native_move_rebinds_live_session_only_from_its_current_cwd`, `native_move_rejects_foreign_owner_and_ignores_completed_replay`, `drifted_session_end_ends_only_the_same_actor_agent_and_cwd` (#865) | STRONG |
|
||||
| 4e | A live session's baton stays its own until the session is quiet | `ai-memory-store/src/ops.rs` - `expire_same_cwd_auto_handoffs` spares batons of other open sessions; `accept_handoff_in_transaction(.., busy_since)` re-checks inside the claim that the source is not an open session with observations after the cutoff, and its sweep `expire_superseded_auto_handoffs(.., busy_since)` spares a busy open session's baton (every open session's when `None`, the explicit MCP accept); `reader.rs` `startup_handoff` selects with the same cutoff; `ai-memory-hooks/src/router.rs` `LIVE_BATON_QUIET_PERIOD` | `multi_session.rs` `parallel_live_sessions_keep_their_own_checkpoint_batons`, `startup_claim_rechecks_the_source_and_sweeps_only_quiet_open_batons`; `router.rs` `opencode_parallel_live_batons_are_owned_and_wait_for_quiet` - each fails with its guard removed | STRONG |
|
||||
| 5a | Pages shared: `author_id` is never a read filter (invariant #16) | `ai-memory-store/src/reader.rs` `search_pages`/`page_body_by_ids` — `author_id` is an attribution JOIN only, never a WHERE term | `multi_session.rs` — a page with a **non-null** `author_id` (operator A) is readable by operator B in the same project | STRONG |
|
||||
| 5b | Page supersession (loser stays reachable) | `ai-memory-store/src/ops.rs` `upsert_page_in_tx` — demote `is_latest=0` (never delete) + `supersedes` chain | `multi_session.rs` `concurrent_writes_to_one_path_supersede_rather_than_destroy`; `retrieval_superseded.rs` | STRONG |
|
||||
| 6 | Active-project pointer (PerActor, no clobber) | `ai-memory-core/src/active_project.rs` `set_for`/`lookup_for` (fail-closed on `Mismatch`) | `active_project.rs` `parallel_harnesses_of_one_user_keep_separate_pointers`, `two_operators_never_read_each_others_pointer`, `a_session_mismatch_fails_closed_once_anything_has_been_keyed` | STRONG |
|
||||
|
||||
Reference in New Issue
Block a user