mirror of
https://github.com/akitaonrails/ai-memory.git
synced 2026-10-02 03:24:46 +08:00
@@ -12,6 +12,12 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
|
||||
provider. (#1026)
|
||||
|
||||
### Fixed
|
||||
- Fixed interrupted launchers blocking an immediate managed-workstream restart
|
||||
by adding explicit `ai-memory run --force-unlock` recovery. The server
|
||||
atomically expires and replaces only a lease attributed to the same
|
||||
authenticated operator (or another unattributed lease in single-user
|
||||
operation), while cross-operator eviction remains refused; the command does
|
||||
not kill a native process. (#795)
|
||||
- Fixed Unix release archives potentially carrying macOS AppleDouble sidecars
|
||||
outside the native upgrader's strict path allowlist. Release packing now
|
||||
disables sidecar generation, and packaging tests lock the staged top-level
|
||||
|
||||
@@ -364,8 +364,16 @@ config home.
|
||||
ai-memory run claude
|
||||
ai-memory run codex --yolo # later: same workstream, different harness
|
||||
ai-memory continue # resume the newest managed checkout
|
||||
# after a dead launcher left its lease behind (same operator only)
|
||||
ai-memory run --force-unlock codex
|
||||
```
|
||||
|
||||
`--force-unlock` immediately expires the selected workstream's active lease;
|
||||
use it only when you know the previous launcher is gone. It does not kill a
|
||||
native process, and it cannot evict another authenticated operator's run. See
|
||||
the [managed-workstream recovery notes](docs/managed-workstreams.md#lease-recovery)
|
||||
for the full safety contract.
|
||||
|
||||
Auto-wiring is on by default; opt out with `ai-memory run --no-autowire` or
|
||||
`AI_MEMORY_RUN_AUTOWIRE=false`. You can still wire agents by hand with
|
||||
`install-hooks` / `install-mcp` (e.g. for a harness you never launch through
|
||||
|
||||
@@ -49,7 +49,7 @@ pub enum Command {
|
||||
RepairBackfillTimestamps(RepairBackfillTimestampsArgs),
|
||||
/// Launch an agent in an opt-in, cross-harness managed workstream.
|
||||
/// Native arguments are forwarded except exact wrapper flags such as
|
||||
/// `--yolo` and `--fresh`.
|
||||
/// `--yolo`, `--fresh`, and `--force-unlock`.
|
||||
Run(RunArgs),
|
||||
/// Pick a local project and installed harness, then launch from that
|
||||
/// checkout. Removes the `cd` step `run` requires.
|
||||
@@ -326,6 +326,12 @@ pub struct RunArgs {
|
||||
/// resuming or adopting an existing harness session.
|
||||
#[arg(long)]
|
||||
pub fresh: bool,
|
||||
/// Force-expire the selected workstream's active lease before launching.
|
||||
/// Use only when the previous launcher is gone: its later heartbeat or
|
||||
/// finish will be refused. The server permits takeover only for the same
|
||||
/// authenticated operator (or the single-user, unattributed owner).
|
||||
#[arg(long)]
|
||||
pub force_unlock: bool,
|
||||
/// Skip the one-time auto-install of this harness's ai-memory hooks + MCP.
|
||||
/// Auto-wire is on by default so a managed launch captures without a manual
|
||||
/// `install-hooks`/`install-mcp` step; pass this (or set
|
||||
|
||||
@@ -76,6 +76,7 @@ pub async fn run(config: &Config, args: ContinueArgs) -> Result<i32> {
|
||||
jail: None,
|
||||
no_jail: false,
|
||||
fresh: args.fresh,
|
||||
force_unlock: false,
|
||||
no_autowire: false,
|
||||
env: Vec::new(),
|
||||
env_file: None,
|
||||
|
||||
@@ -165,6 +165,7 @@ pub async fn run(config: &Config, args: ResumeArgs) -> Result<i32> {
|
||||
jail: None,
|
||||
no_jail: false,
|
||||
fresh: args.fresh,
|
||||
force_unlock: false,
|
||||
no_autowire: false,
|
||||
env: Vec::new(),
|
||||
env_file: None,
|
||||
|
||||
@@ -126,6 +126,7 @@ pub(super) async fn run_from_with_wiring(
|
||||
let trailing_yolo = remove_wrapper_yolo(&mut native_args);
|
||||
let trailing_true_yolo = remove_wrapper_true_yolo(&mut native_args);
|
||||
let trailing_fresh = remove_wrapper_fresh(&mut native_args);
|
||||
let trailing_force_unlock = remove_wrapper_force_unlock(&mut native_args);
|
||||
let trailing_no_autowire = remove_wrapper_no_autowire(&mut native_args);
|
||||
let trailing_jail = remove_wrapper_jail(&mut native_args);
|
||||
let jail_request = jail_request(args.jail, args.no_jail, trailing_jail)?;
|
||||
@@ -136,6 +137,7 @@ pub(super) async fn run_from_with_wiring(
|
||||
);
|
||||
let yolo_requested = yolo_modes.yolo;
|
||||
let force_fresh = args.fresh || trailing_fresh;
|
||||
let force_unlock = args.force_unlock || trailing_force_unlock;
|
||||
let no_autowire = args.no_autowire || trailing_no_autowire;
|
||||
let run_env = resolve_run_env(args.env_file.as_deref(), &args.env)
|
||||
.context("resolving --env/--env-file for the managed run")?;
|
||||
@@ -226,6 +228,7 @@ pub(super) async fn run_from_with_wiring(
|
||||
available_agents: unique_auto_agents(&auto_candidates),
|
||||
workstream: args.workstream,
|
||||
new_workstream: args.new_workstream,
|
||||
force_unlock,
|
||||
lease_owner: lease_owner(),
|
||||
};
|
||||
let interrupted_before_spawn = CancellationToken::new();
|
||||
@@ -1618,6 +1621,12 @@ fn remove_wrapper_fresh(args: &mut Vec<OsString>) -> bool {
|
||||
args.len() != before
|
||||
}
|
||||
|
||||
fn remove_wrapper_force_unlock(args: &mut Vec<OsString>) -> bool {
|
||||
let before = args.len();
|
||||
args.retain(|arg| arg != OsStr::new("--force-unlock"));
|
||||
args.len() != before
|
||||
}
|
||||
|
||||
fn remove_wrapper_no_autowire(args: &mut Vec<OsString>) -> bool {
|
||||
let before = args.len();
|
||||
args.retain(|arg| arg != OsStr::new("--no-autowire"));
|
||||
@@ -2283,6 +2292,15 @@ async fn prepare_managed_run(
|
||||
interactive: bool,
|
||||
interrupted: &CancellationToken,
|
||||
) -> Result<PrepareManagedRunResponse> {
|
||||
if request.force_unlock {
|
||||
return match post_json(endpoint, "/workstream/runs", request).await {
|
||||
Err(error) if is_active_workstream_conflict(&error) => Err(error.context(
|
||||
"--force-unlock was refused: the active lease belongs to another operator, or \
|
||||
the server does not support forced lease recovery",
|
||||
)),
|
||||
other => other,
|
||||
};
|
||||
}
|
||||
let result = prepare_managed_run_with_retry(
|
||||
endpoint,
|
||||
request,
|
||||
@@ -3284,6 +3302,7 @@ mod tests {
|
||||
available_agents: Vec::new(),
|
||||
workstream: None,
|
||||
new_workstream: None,
|
||||
force_unlock: false,
|
||||
lease_owner: "workstation:43".into(),
|
||||
};
|
||||
|
||||
@@ -3402,6 +3421,7 @@ mod tests {
|
||||
available_agents: Vec::new(),
|
||||
workstream: None,
|
||||
new_workstream: None,
|
||||
force_unlock: false,
|
||||
lease_owner: "workstation:43".into(),
|
||||
}
|
||||
}
|
||||
@@ -3475,6 +3495,27 @@ mod tests {
|
||||
server.abort();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn force_unlock_asks_once_instead_of_waiting_out_the_lease() {
|
||||
let (app, attempts) = held_lease_server(usize::MAX, Duration::from_secs(60));
|
||||
let (endpoint, server) = serve(app).await;
|
||||
let mut request = held_lease_request();
|
||||
request.force_unlock = true;
|
||||
let started = std::time::Instant::now();
|
||||
|
||||
let error = prepare_managed_run(&endpoint, &request, true, &CancellationToken::new())
|
||||
.await
|
||||
.expect_err("a server that refuses the takeover must fail immediately");
|
||||
|
||||
assert!(started.elapsed() < Duration::from_secs(5));
|
||||
assert_eq!(attempts.load(Ordering::SeqCst), 1);
|
||||
assert!(
|
||||
format!("{error:#}").contains("--force-unlock was refused"),
|
||||
"{error:#}"
|
||||
);
|
||||
server.abort();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn ctrl_c_aborts_the_held_lease_wait_immediately() {
|
||||
let (app, attempts) = held_lease_server(1, Duration::from_secs(60));
|
||||
@@ -3969,6 +4010,30 @@ mod tests {
|
||||
assert_eq!(trailing, ["--model", "opus"].map(OsString::from));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn force_unlock_is_a_wrapper_flag_before_or_after_the_harness() {
|
||||
let cli = Cli::try_parse_from([
|
||||
"ai-memory",
|
||||
"run",
|
||||
"--force-unlock",
|
||||
"codex",
|
||||
"--model",
|
||||
"gpt-5",
|
||||
])
|
||||
.unwrap();
|
||||
let CliCommand::Run(args) = cli.command else {
|
||||
panic!("expected run command");
|
||||
};
|
||||
assert!(args.force_unlock);
|
||||
assert_eq!(args.native_args, ["--model", "gpt-5"].map(OsString::from));
|
||||
|
||||
let mut trailing = ["--model", "gpt-5", "--force-unlock"]
|
||||
.map(OsString::from)
|
||||
.to_vec();
|
||||
assert!(remove_wrapper_force_unlock(&mut trailing));
|
||||
assert_eq!(trailing, ["--model", "gpt-5"].map(OsString::from));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn wrapper_no_autowire_parses_before_or_after_the_harness() {
|
||||
// Before the harness: clap binds it as the wrapper flag.
|
||||
@@ -4129,6 +4194,7 @@ mod tests {
|
||||
jail: None,
|
||||
no_jail: false,
|
||||
fresh: false,
|
||||
force_unlock: false,
|
||||
no_autowire: true,
|
||||
env: vec![("CLAUDE_CONFIG_DIR".to_string(), "/from/cli".to_string())],
|
||||
env_file: Some(env_file.clone()),
|
||||
@@ -4617,6 +4683,7 @@ mod tests {
|
||||
jail: None,
|
||||
no_jail: false,
|
||||
fresh: false,
|
||||
force_unlock: false,
|
||||
no_autowire: false,
|
||||
env: Vec::new(),
|
||||
env_file: None,
|
||||
@@ -4772,6 +4839,7 @@ mod tests {
|
||||
jail: None,
|
||||
no_jail: false,
|
||||
fresh: false,
|
||||
force_unlock: false,
|
||||
no_autowire: false,
|
||||
env,
|
||||
env_file: None,
|
||||
|
||||
@@ -229,6 +229,7 @@ pub async fn run(config: &Config, args: ShowArgs) -> Result<i32> {
|
||||
jail: None,
|
||||
no_jail: false,
|
||||
fresh: args.fresh,
|
||||
force_unlock: false,
|
||||
no_autowire: false,
|
||||
env: Vec::new(),
|
||||
env_file: None,
|
||||
|
||||
@@ -146,6 +146,12 @@ pub struct PrepareManagedRunRequest {
|
||||
/// Create and select a fresh named workstream.
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub new_workstream: Option<String>,
|
||||
/// Expire an active lease owned by this same operator before opening the
|
||||
/// replacement run. This is an explicit recovery override for a launcher
|
||||
/// that exited without releasing its lease; it never permits cross-owner
|
||||
/// takeover.
|
||||
#[serde(default)]
|
||||
pub force_unlock: bool,
|
||||
/// Diagnostic owner label (host and process id), not an authorization key.
|
||||
pub lease_owner: String,
|
||||
}
|
||||
@@ -398,5 +404,27 @@ mod tests {
|
||||
|
||||
assert!(!request.automatic_harness);
|
||||
assert!(request.available_agents.is_empty());
|
||||
assert!(!request.force_unlock);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn forced_unlock_is_explicit_on_the_wire() {
|
||||
let request: PrepareManagedRunRequest = serde_json::from_value(serde_json::json!({
|
||||
"workspace": "default",
|
||||
"project": "memory",
|
||||
"cwd": "/repo",
|
||||
"repo_fingerprint": "repo",
|
||||
"worktree_fingerprint": "worktree",
|
||||
"agent": "codex",
|
||||
"force_unlock": true,
|
||||
"lease_owner": "host:2"
|
||||
}))
|
||||
.unwrap();
|
||||
|
||||
assert!(request.force_unlock);
|
||||
assert_eq!(
|
||||
serde_json::to_value(request).unwrap()["force_unlock"],
|
||||
serde_json::json!(true)
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -288,9 +288,10 @@ async fn prepare_run(
|
||||
Ok(owner) => owner,
|
||||
Err(failure) => return error(StatusCode::INTERNAL_SERVER_ERROR, failure.to_string()),
|
||||
};
|
||||
let force_unlock = request.force_unlock;
|
||||
let prepared = state
|
||||
.writer
|
||||
.prepare_workstream_run_owned(
|
||||
.prepare_workstream_run_owned_with_unlock(
|
||||
PrepareWorkstreamRun {
|
||||
workspace_id: scope.workspace_id,
|
||||
project_id: scope.project_id,
|
||||
@@ -304,6 +305,7 @@ async fn prepare_run(
|
||||
lease_owner: request.lease_owner,
|
||||
},
|
||||
owner_user,
|
||||
force_unlock,
|
||||
)
|
||||
.await;
|
||||
match prepared {
|
||||
@@ -1439,6 +1441,7 @@ mod tests {
|
||||
available_agents: vec![AgentKind::Codex, AgentKind::ClaudeCode],
|
||||
workstream: None,
|
||||
new_workstream: None,
|
||||
force_unlock: false,
|
||||
lease_owner: "automatic".into(),
|
||||
}),
|
||||
)
|
||||
@@ -1476,6 +1479,7 @@ mod tests {
|
||||
available_agents: Vec::new(),
|
||||
workstream: None,
|
||||
new_workstream: None,
|
||||
force_unlock: false,
|
||||
lease_owner: "explicit".into(),
|
||||
}),
|
||||
)
|
||||
@@ -1515,6 +1519,7 @@ mod tests {
|
||||
available_agents: vec![AgentKind::KimiCode, AgentKind::ClaudeCode],
|
||||
workstream: None,
|
||||
new_workstream: None,
|
||||
force_unlock: false,
|
||||
lease_owner: "automatic".into(),
|
||||
}),
|
||||
)
|
||||
@@ -1561,6 +1566,7 @@ mod tests {
|
||||
available_agents: Vec::new(),
|
||||
workstream: None,
|
||||
new_workstream: None,
|
||||
force_unlock: false,
|
||||
lease_owner: "alice-launcher".into(),
|
||||
}),
|
||||
)
|
||||
@@ -1613,6 +1619,106 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn force_unlock_cannot_evict_another_operator() {
|
||||
let temp = TempDir::new().unwrap();
|
||||
let store = Store::open(temp.path()).unwrap();
|
||||
let state = test_state(&store, temp.path());
|
||||
for name in ["alice", "bob"] {
|
||||
store
|
||||
.writer
|
||||
.create_human_user(
|
||||
ai_memory_core::NewUser {
|
||||
username: name.into(),
|
||||
name: None,
|
||||
email: None,
|
||||
},
|
||||
ai_memory_core::UserRole::User,
|
||||
None,
|
||||
false,
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
}
|
||||
let actor = |name: &str| {
|
||||
Some(Extension(ai_memory_core::ActorContext {
|
||||
user: Some(name.into()),
|
||||
..ai_memory_core::ActorContext::default()
|
||||
}))
|
||||
};
|
||||
let request = |force_unlock: bool, lease_owner: &str| PrepareManagedRunRequest {
|
||||
workspace: "default".into(),
|
||||
project: "managed-force-unlock".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,
|
||||
force_unlock,
|
||||
lease_owner: lease_owner.into(),
|
||||
};
|
||||
|
||||
let first = prepare_run(
|
||||
State(state.clone()),
|
||||
None,
|
||||
actor("alice"),
|
||||
None,
|
||||
Json(request(false, "alice:1")),
|
||||
)
|
||||
.await;
|
||||
assert_eq!(first.status(), StatusCode::OK);
|
||||
let body = to_bytes(first.into_body(), 64 * 1024).await.unwrap();
|
||||
let first: PrepareManagedRunResponse = serde_json::from_slice(&body).unwrap();
|
||||
|
||||
let refused = prepare_run(
|
||||
State(state.clone()),
|
||||
None,
|
||||
actor("bob"),
|
||||
None,
|
||||
Json(request(true, "bob:2")),
|
||||
)
|
||||
.await;
|
||||
assert_eq!(refused.status(), StatusCode::CONFLICT);
|
||||
assert!(
|
||||
store
|
||||
.writer
|
||||
.heartbeat_managed_run(first.run_id)
|
||||
.await
|
||||
.unwrap(),
|
||||
"the refused caller must not disturb Alice's lease"
|
||||
);
|
||||
|
||||
let replacement = prepare_run(
|
||||
State(state),
|
||||
None,
|
||||
actor("alice"),
|
||||
None,
|
||||
Json(request(true, "alice:3")),
|
||||
)
|
||||
.await;
|
||||
assert_eq!(replacement.status(), StatusCode::OK);
|
||||
let body = to_bytes(replacement.into_body(), 64 * 1024).await.unwrap();
|
||||
let replacement: PrepareManagedRunResponse = serde_json::from_slice(&body).unwrap();
|
||||
assert_ne!(replacement.run_id, first.run_id);
|
||||
assert!(
|
||||
!store
|
||||
.writer
|
||||
.heartbeat_managed_run(first.run_id)
|
||||
.await
|
||||
.unwrap()
|
||||
);
|
||||
assert!(
|
||||
store
|
||||
.writer
|
||||
.heartbeat_managed_run(replacement.run_id)
|
||||
.await
|
||||
.unwrap()
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn kiro_is_accepted_as_an_explicit_and_automatic_harness() {
|
||||
let temp = TempDir::new().unwrap();
|
||||
@@ -1635,6 +1741,7 @@ mod tests {
|
||||
available_agents: Vec::new(),
|
||||
workstream: None,
|
||||
new_workstream: None,
|
||||
force_unlock: false,
|
||||
lease_owner: "explicit".into(),
|
||||
}),
|
||||
)
|
||||
@@ -1673,6 +1780,7 @@ mod tests {
|
||||
available_agents: vec![AgentKind::KiroCli],
|
||||
workstream: None,
|
||||
new_workstream: None,
|
||||
force_unlock: false,
|
||||
lease_owner: "automatic".into(),
|
||||
}),
|
||||
)
|
||||
@@ -1702,6 +1810,7 @@ mod tests {
|
||||
available_agents: Vec::new(),
|
||||
workstream: None,
|
||||
new_workstream: None,
|
||||
force_unlock: false,
|
||||
lease_owner: "explicit".into(),
|
||||
}),
|
||||
)
|
||||
@@ -1740,6 +1849,7 @@ mod tests {
|
||||
available_agents: vec![AgentKind::CommandCode, AgentKind::ClaudeCode],
|
||||
workstream: None,
|
||||
new_workstream: None,
|
||||
force_unlock: false,
|
||||
lease_owner: "automatic".into(),
|
||||
}),
|
||||
)
|
||||
@@ -1769,6 +1879,7 @@ mod tests {
|
||||
available_agents: Vec::new(),
|
||||
workstream: None,
|
||||
new_workstream: None,
|
||||
force_unlock: false,
|
||||
lease_owner: "explicit".into(),
|
||||
}),
|
||||
)
|
||||
@@ -1807,6 +1918,7 @@ mod tests {
|
||||
available_agents: vec![AgentKind::Grok],
|
||||
workstream: None,
|
||||
new_workstream: None,
|
||||
force_unlock: false,
|
||||
lease_owner: "automatic".into(),
|
||||
}),
|
||||
)
|
||||
@@ -1839,6 +1951,7 @@ mod tests {
|
||||
available_agents: Vec::new(),
|
||||
workstream: None,
|
||||
new_workstream: None,
|
||||
force_unlock: false,
|
||||
lease_owner: "explicit".into(),
|
||||
}),
|
||||
)
|
||||
@@ -1877,6 +1990,7 @@ mod tests {
|
||||
available_agents: vec![AgentKind::AntigravityCli],
|
||||
workstream: None,
|
||||
new_workstream: None,
|
||||
force_unlock: false,
|
||||
lease_owner: "automatic".into(),
|
||||
}),
|
||||
)
|
||||
@@ -2144,6 +2258,7 @@ mod tests {
|
||||
available_agents: Vec::new(),
|
||||
workstream: None,
|
||||
new_workstream: None,
|
||||
force_unlock: false,
|
||||
lease_owner: "alice".into(),
|
||||
}),
|
||||
)
|
||||
|
||||
@@ -12545,6 +12545,7 @@ pub(crate) mod tests {
|
||||
lease_owner: "test".into(),
|
||||
},
|
||||
None,
|
||||
false,
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
@@ -14457,6 +14458,7 @@ pub(crate) mod tests {
|
||||
lease_owner: "test".to_string(),
|
||||
},
|
||||
None,
|
||||
false,
|
||||
)
|
||||
.expect("opening a managed run should succeed")
|
||||
}
|
||||
|
||||
@@ -276,6 +276,7 @@ pub(crate) fn prepare_run(
|
||||
conn: &mut Connection,
|
||||
input: &PrepareWorkstreamRun,
|
||||
owner_user: Option<&str>,
|
||||
force_unlock: bool,
|
||||
) -> StoreResult<PreparedWorkstreamRun> {
|
||||
if owner_user
|
||||
.is_some_and(|owner| ai_memory_core::IdentityKey::from_storage_key(owner).is_none())
|
||||
@@ -293,21 +294,29 @@ pub(crate) fn prepare_run(
|
||||
)?;
|
||||
|
||||
let (workstream_id, workstream_name) = select_workstream(&tx, input, now)?;
|
||||
let busy: Option<(String, i64)> = tx
|
||||
let busy: Option<(Vec<u8>, String, i64, Option<String>)> = tx
|
||||
.query_row(
|
||||
"SELECT lease_owner, lease_expires_at FROM managed_runs \
|
||||
"SELECT id, lease_owner, lease_expires_at, owner_user FROM managed_runs \
|
||||
WHERE workstream_id = ?1 AND state = 'active'",
|
||||
params![workstream_id.as_bytes()],
|
||||
|row| Ok((row.get(0)?, row.get(1)?)),
|
||||
|row| Ok((row.get(0)?, row.get(1)?, row.get(2)?, row.get(3)?)),
|
||||
)
|
||||
.optional()?;
|
||||
if let Some((owner, expires)) = busy {
|
||||
return Err(StoreError::WorkstreamBusy(format!(
|
||||
"owned by {owner} until {}",
|
||||
Timestamp::from_microsecond(expires)
|
||||
.map(|t| t.to_string())
|
||||
.unwrap_or_else(|_| expires.to_string())
|
||||
)));
|
||||
if let Some((run_id, owner, expires, active_owner_user)) = busy {
|
||||
if force_unlock && active_owner_user.as_deref() == owner_user {
|
||||
tx.execute(
|
||||
"UPDATE managed_runs SET state = 'expired', ended_at = ?1, \
|
||||
lease_expires_at = ?1 WHERE id = ?2 AND state = 'active'",
|
||||
params![now, run_id],
|
||||
)?;
|
||||
} else {
|
||||
return Err(StoreError::WorkstreamBusy(format!(
|
||||
"owned by {owner} until {}",
|
||||
Timestamp::from_microsecond(expires)
|
||||
.map(|t| t.to_string())
|
||||
.unwrap_or_else(|_| expires.to_string())
|
||||
)));
|
||||
}
|
||||
}
|
||||
|
||||
let latest_sequence: i64 = tx.query_row(
|
||||
|
||||
@@ -715,6 +715,7 @@ pub(crate) enum WriteCmd {
|
||||
PrepareWorkstreamRun {
|
||||
input: PrepareWorkstreamRun,
|
||||
owner_user: Option<String>,
|
||||
force_unlock: bool,
|
||||
reply: oneshot::Sender<StoreResult<PreparedWorkstreamRun>>,
|
||||
},
|
||||
HeartbeatManagedRun {
|
||||
@@ -3013,11 +3014,24 @@ impl WriterHandle {
|
||||
&self,
|
||||
input: PrepareWorkstreamRun,
|
||||
owner_user: Option<String>,
|
||||
) -> StoreResult<PreparedWorkstreamRun> {
|
||||
self.prepare_workstream_run_owned_with_unlock(input, owner_user, false)
|
||||
.await
|
||||
}
|
||||
|
||||
/// Select a workstream, optionally replacing this same operator's active
|
||||
/// lease in the same writer transaction.
|
||||
pub async fn prepare_workstream_run_owned_with_unlock(
|
||||
&self,
|
||||
input: PrepareWorkstreamRun,
|
||||
owner_user: Option<String>,
|
||||
force_unlock: bool,
|
||||
) -> StoreResult<PreparedWorkstreamRun> {
|
||||
let (tx, rx) = oneshot::channel();
|
||||
self.send(WriteCmd::PrepareWorkstreamRun {
|
||||
input,
|
||||
owner_user,
|
||||
force_unlock,
|
||||
reply: tx,
|
||||
})
|
||||
.await?;
|
||||
@@ -4269,10 +4283,15 @@ fn worker_loop(mut conn: Connection, mut rx: mpsc::Receiver<WriteCmd>) {
|
||||
WriteCmd::PrepareWorkstreamRun {
|
||||
input,
|
||||
owner_user,
|
||||
force_unlock,
|
||||
reply,
|
||||
} => {
|
||||
let result =
|
||||
crate::workstream::prepare_run(&mut conn, &input, owner_user.as_deref());
|
||||
let result = crate::workstream::prepare_run(
|
||||
&mut conn,
|
||||
&input,
|
||||
owner_user.as_deref(),
|
||||
force_unlock,
|
||||
);
|
||||
send_or_warn(reply, result, "prepare_workstream_run");
|
||||
}
|
||||
WriteCmd::HeartbeatManagedRun { run_id, reply } => {
|
||||
|
||||
@@ -22,7 +22,7 @@ use ai_memory_core::{
|
||||
owner_stamp,
|
||||
};
|
||||
use ai_memory_store::{
|
||||
LinkOrAdoptManagedRunSession, ManagedRunSessionLink, PrepareWorkstreamRun, Store,
|
||||
LinkOrAdoptManagedRunSession, ManagedRunSessionLink, PrepareWorkstreamRun, Store, StoreError,
|
||||
WorkstreamSelection,
|
||||
};
|
||||
|
||||
@@ -917,6 +917,99 @@ async fn stale_codex_run_recovery_cannot_cross_project_owner_or_session_boundari
|
||||
);
|
||||
}
|
||||
|
||||
/// Forced lease recovery is deliberately narrower than project write access:
|
||||
/// the same operator can recover their own abandoned launcher, while another
|
||||
/// operator in the same project cannot evict it. The old run becomes terminal
|
||||
/// in the same transaction that creates the replacement.
|
||||
#[tokio::test]
|
||||
async fn force_unlock_replaces_only_the_same_operators_active_run() {
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
let store = Store::open(tmp.path()).unwrap();
|
||||
let (ws, proj) = scope(&store).await;
|
||||
let alice = operator("alice");
|
||||
let bob = operator("bob");
|
||||
let prepare = |lease_owner: &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::Current,
|
||||
lease_owner: lease_owner.into(),
|
||||
};
|
||||
|
||||
let abandoned = store
|
||||
.writer
|
||||
.prepare_workstream_run_owned(prepare("alice:1"), Some(alice.clone()))
|
||||
.await
|
||||
.unwrap();
|
||||
let refused = store
|
||||
.writer
|
||||
.prepare_workstream_run_owned_with_unlock(prepare("bob:2"), Some(bob), true)
|
||||
.await
|
||||
.unwrap_err();
|
||||
assert!(matches!(refused, StoreError::WorkstreamBusy(_)));
|
||||
assert!(
|
||||
store
|
||||
.writer
|
||||
.heartbeat_managed_run(abandoned.run_id)
|
||||
.await
|
||||
.unwrap(),
|
||||
"a refused cross-owner takeover must leave the original lease active"
|
||||
);
|
||||
|
||||
let replacement = store
|
||||
.writer
|
||||
.prepare_workstream_run_owned_with_unlock(prepare("alice:3"), Some(alice.clone()), true)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_ne!(replacement.run_id, abandoned.run_id);
|
||||
assert!(
|
||||
!store
|
||||
.writer
|
||||
.heartbeat_managed_run(abandoned.run_id)
|
||||
.await
|
||||
.unwrap(),
|
||||
"the replaced run must be terminal"
|
||||
);
|
||||
assert!(
|
||||
store
|
||||
.writer
|
||||
.heartbeat_managed_run(replacement.run_id)
|
||||
.await
|
||||
.unwrap(),
|
||||
"the replacement is the sole live control"
|
||||
);
|
||||
|
||||
let solo = store
|
||||
.writer
|
||||
.prepare_workstream_run_owned(
|
||||
PrepareWorkstreamRun {
|
||||
selection: WorkstreamSelection::New("solo".into()),
|
||||
..prepare("solo:1")
|
||||
},
|
||||
None,
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
let solo_replacement = store
|
||||
.writer
|
||||
.prepare_workstream_run_owned_with_unlock(
|
||||
PrepareWorkstreamRun {
|
||||
selection: WorkstreamSelection::Named("solo".into()),
|
||||
..prepare("solo:2")
|
||||
},
|
||||
None,
|
||||
true,
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_ne!(solo_replacement.run_id, solo.run_id);
|
||||
}
|
||||
|
||||
/// 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.
|
||||
|
||||
@@ -205,8 +205,12 @@ ledger/session state: after any harness establishes the workstream, a newly
|
||||
joining harness starts fresh and receives portable history instead of adopting
|
||||
unrelated old native history. Handled launcher failures cancel their lease;
|
||||
normal reopen retries brief finalization conflicts, while an unclean process
|
||||
death remains bounded by the renewable lease expiry. See [Managed cross-harness
|
||||
workstreams](managed-workstreams.md).
|
||||
death remains bounded by the renewable lease expiry. An explicit
|
||||
`--force-unlock` recovery expires and replaces a selected active lease in the
|
||||
same writer transaction, but only when its durable operator attribution equals
|
||||
the new run's attribution; the informational `host:pid` lease label is never an
|
||||
authorization key. The old run can no longer heartbeat or finish. See [Managed
|
||||
cross-harness workstreams](managed-workstreams.md).
|
||||
|
||||
## Hook event vocabulary
|
||||
|
||||
|
||||
@@ -74,9 +74,9 @@ ai-memory run
|
||||
```
|
||||
|
||||
Everything after the harness name is native argv except the wrapper-owned exact
|
||||
flags `--yolo` and `--fresh`. No `--` separator is needed, and ai-memory does
|
||||
not maintain a second copy of each harness's option schema. Other wrapper
|
||||
options come first:
|
||||
flags `--yolo`, `--fresh`, and `--force-unlock`. No `--` separator is needed,
|
||||
and ai-memory does not maintain a second copy of each harness's option schema.
|
||||
Other wrapper options come first:
|
||||
|
||||
Portable events, handoffs, and project briefs are injected as explicitly
|
||||
delimited, untrusted historical data. Instruction-like text inside stored
|
||||
@@ -88,7 +88,8 @@ file, and the current checkout remain authoritative.
|
||||
```text
|
||||
ai-memory run [--workspace NAME] [--project NAME]
|
||||
[--workstream NAME | --new NAME] [--executable PATH]
|
||||
[--yolo] [--fresh] [--env KEY=VALUE]... [--env-file PATH]
|
||||
[--yolo] [--fresh] [--force-unlock]
|
||||
[--env KEY=VALUE]... [--env-file PATH]
|
||||
[claude|claude*|codex|opencode|opencode2|pi|crush|omp|kimi|command-code|kiro|grok|antigravity]
|
||||
[native arguments...]
|
||||
```
|
||||
@@ -691,6 +692,8 @@ immediately. A new launch retries an active-workstream conflict briefly so a
|
||||
previous launcher can finish; if another harness is genuinely still running,
|
||||
the conflict remains and concurrent writers are still rejected.
|
||||
|
||||
### Lease recovery
|
||||
|
||||
A launcher that dies without releasing its lease — killed, its terminal
|
||||
closed, or a sandbox such as ai-jail torn down — leaves the workstream held
|
||||
until that lease lapses. An interactive relaunch (stdin and stderr are
|
||||
@@ -705,6 +708,25 @@ retry window and fail fast rather than hanging. Terminal
|
||||
interrupts continue to reach the child while the parent stays alive to finish
|
||||
or cancel the run.
|
||||
|
||||
When you know the prior launcher is gone and do not want to wait for the lease,
|
||||
force-expire it explicitly:
|
||||
|
||||
```bash
|
||||
ai-memory run --force-unlock codex
|
||||
# The exact wrapper flag is also accepted after the harness name.
|
||||
ai-memory run codex --force-unlock
|
||||
```
|
||||
|
||||
The replacement is atomic and limited to the same durable authenticated
|
||||
operator; in single-user or otherwise unattributed operation, both runs must be
|
||||
unattributed. A different operator's active run is still refused. The command
|
||||
expires the managed lease only — it does not signal or kill a native process.
|
||||
If the previous launcher is actually alive, its later heartbeats and finish are
|
||||
rejected, and its final transcript tail may not be imported. Use
|
||||
`--force-unlock` only after verifying that launcher has stopped. Older servers
|
||||
do not honor the request and return a refusal, so upgrade the server as well as
|
||||
the client before relying on this recovery path.
|
||||
|
||||
Before the child starts, `Ctrl+C` at the native-session chooser cancels the
|
||||
acquired run and exits without requiring Enter or adopting the selected session.
|
||||
The launcher waits for the server's cancellation response. A request error is
|
||||
|
||||
@@ -61,7 +61,7 @@ boundary not yet built.
|
||||
| 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 |
|
||||
| 13b | Managed-run lease exclusivity and explicit recovery (one active run per workstream, invariant #16) | `ai-memory-store/src/workstream.rs` `prepare_run` expires lapsed leases and otherwise refuses an `active` run inside the writer transaction (`StoreError::WorkstreamBusy`). The explicit `force_unlock` path may expire and replace that row in the same transaction only when its durable `owner_user` equals the incoming actor's attribution; different operators remain isolated. `heartbeat` renews only `active` rows. `ai-memory-cli/src/commands/run.rs` normally uses `wait_out_held_lease`; `--force-unlock` instead asks once and fails immediately when refused. Neither path treats the informational `lease_owner` label as authority | `multi_session.rs` `force_unlock_replaces_only_the_same_operators_active_run` (foreign operator refused with live control, same operator succeeds, old run loses heartbeat, unattributed single-user control); hooks `workstream.rs` `force_unlock_cannot_evict_another_operator` exercises the HTTP actor boundary; CLI `run.rs` `force_unlock_asks_once_instead_of_waiting_out_the_lease`; existing controls `managed_workstream_batches_are_idempotent_and_release_the_lease`, `a_renewed_lease_is_reported_as_a_live_owner_not_taken_over`, `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. A shared unattributed/root context is intentionally one operator bucket; deployments needing person-level refusal must authenticate DB users |
|
||||
| 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 |
|
||||
| 14c | Client-supplied event time is bounded and self-scoped (#919) | `ai-memory-hooks/src/payload.rs` `HookEnvelope::occurred_at_micros` — a `/hook` caller's optional `occurred_at` must be `> 0` and `<= now + 5min`; it only sets the caller's *own* admitted session's `started_at`/`ended_at`/`created_at`, never another session/project/user, and the ingest dedup `seen_at`/TTL stays on `now` | `payload.rs` `occurred_at_micros_rejects_a_far_future_timestamp`, `…_rejects_non_positive_values`, `…_rejects_garbage_strings` | MEDIUM — self-scoped numeric bound; no cross-tenant surface (a client already controls its own content) |
|
||||
|
||||
Reference in New Issue
Block a user