hive-runtime: read the ACP provider key from bao
An opencode ACP agent got its provider API key only from the hand-placed backendEnvironmentFile. It now also reads it from the swarm secret store at swarm/agents/<agent>/acp-provider, field api_key, under its own certificate, and sets it in the spawned ACP agent's environment only. Nothing is written to disk. Precedence: a value already in the process environment (the env file) wins and the store is not asked. Otherwise the stored key is used when present. With no store, nothing stored, or a failed read, the agent is spawned without the key as before, and one line is logged without the value. The variable name comes from the existing per-agent option acp.opencode.provider.apiKeyEnv, exported as HIVE_ACP_API_KEY_ENV on the harness only for the opencode preset. Other ACP commands are unchanged. The read lives in hive-runtime, where the ACP child is spawned, so both hive-agent and hive-subagent-daemon use it. The subagent daemon unit gets the key name and, when the agent has a store, the agent's store identity (the same credentials queue-identity.nix gives the harness). No new option or setting. Closes #4841.
This commit is contained in:
parent
c2bdf30e05
commit
c5b21403a6
13 changed files with 486 additions and 9 deletions
|
|
@ -11,6 +11,9 @@ workspace = true
|
|||
hive-claude.workspace = true
|
||||
serde.workspace = true
|
||||
serde_json.workspace = true
|
||||
# The ACP agent's provider key, read under the agent's own certificate
|
||||
# (`acp::provider_key`).
|
||||
swarm-secret-client.workspace = true
|
||||
thiserror.workspace = true
|
||||
tokio.workspace = true
|
||||
tracing.workspace = true
|
||||
|
|
|
|||
|
|
@ -5,6 +5,7 @@
|
|||
//! reporting no usage fails with [`AcpError::EmptyEndTurn`]. An agent that
|
||||
//! never reports usage therefore surfaces every such turn as that error.
|
||||
|
||||
mod provider_key;
|
||||
mod rpc;
|
||||
mod stream;
|
||||
|
||||
|
|
@ -418,7 +419,17 @@ impl<P: CompactionPolicy> AcpRuntime<P> {
|
|||
async fn start(&self, config: &Config) -> Result<Live> {
|
||||
let (_, servers) = mcp_servers(config)?;
|
||||
let cwd = session_cwd(config);
|
||||
let conn = Connection::spawn(&self.command, &cwd, self.permit.clone(), servers)?;
|
||||
let key = match &self.command.api_key_env {
|
||||
Some(name) => provider_key::resolve(name).await,
|
||||
None => None,
|
||||
};
|
||||
let conn = Connection::spawn(
|
||||
&self.command,
|
||||
key.as_ref(),
|
||||
&cwd,
|
||||
self.permit.clone(),
|
||||
servers,
|
||||
)?;
|
||||
let init = conn
|
||||
.request(
|
||||
"initialize",
|
||||
|
|
@ -1164,6 +1175,7 @@ done
|
|||
.iter()
|
||||
.map(|(k, v)| ((*k).to_owned(), (*v).to_owned()))
|
||||
.collect::<BTreeMap<_, _>>(),
|
||||
api_key_env: None,
|
||||
};
|
||||
AcpRuntime::new(
|
||||
command,
|
||||
|
|
|
|||
281
hive-runtime/src/acp/provider_key.rs
Normal file
281
hive-runtime/src/acp/provider_key.rs
Normal file
|
|
@ -0,0 +1,281 @@
|
|||
//! The ACP agent's provider API key, read from the swarm secret store at
|
||||
//! `swarm/agents/<agent>/acp-provider` (field `api_key`) under the agent's own
|
||||
//! certificate, and handed to the spawned agent's environment only.
|
||||
//!
|
||||
//! A value the process already inherited (`backendEnvironmentFile`) wins and
|
||||
//! the store is not asked. Without either, the agent is spawned without it.
|
||||
|
||||
use std::future::Future;
|
||||
use std::time::Duration;
|
||||
|
||||
use swarm_secret_client::{
|
||||
SecretStore,
|
||||
acp::{ProviderKey, provider_key_path},
|
||||
client::{DEFAULT_CERT_MOUNT, ENV_ADDR, ENV_CACERT, Settings},
|
||||
policy,
|
||||
};
|
||||
|
||||
/// Names the agent whose key is read; the cert-auth role and the store path
|
||||
/// are both built from it.
|
||||
const ENV_AGENT_NAME: &str = "HIVE_AGENT_NAME";
|
||||
|
||||
/// How long one read of the store may take before the agent is spawned
|
||||
/// without the key.
|
||||
const STORE_READ_TIMEOUT: Duration = Duration::from_secs(10);
|
||||
|
||||
/// A variable to add to the spawned agent's environment.
|
||||
///
|
||||
/// No `Debug`: `value` is the key.
|
||||
pub(super) struct KeyVar {
|
||||
pub(super) name: String,
|
||||
pub(super) value: String,
|
||||
}
|
||||
|
||||
/// Where [`resolve_with`] reads the key from. A trait so the precedence can
|
||||
/// be exercised without a store.
|
||||
trait KeySource {
|
||||
/// The store path, for log lines. Never the value.
|
||||
fn origin(&self) -> &str;
|
||||
|
||||
/// The key stored now, or `None` when nothing is stored.
|
||||
fn fetch(&self) -> impl Future<Output = Result<Option<String>, swarm_secret_client::Error>>;
|
||||
}
|
||||
|
||||
struct StoreSource {
|
||||
settings: Settings,
|
||||
role: String,
|
||||
path: String,
|
||||
}
|
||||
|
||||
impl KeySource for StoreSource {
|
||||
fn origin(&self) -> &str {
|
||||
&self.path
|
||||
}
|
||||
|
||||
async fn fetch(&self) -> Result<Option<String>, swarm_secret_client::Error> {
|
||||
let store = SecretStore::connect(&self.settings, &self.role, DEFAULT_CERT_MOUNT).await?;
|
||||
let stored: Option<ProviderKey> = store.read_optional(&self.path).await?;
|
||||
Ok(stored.map(|k| k.api_key))
|
||||
}
|
||||
}
|
||||
|
||||
/// Where this agent's key would be read from, or `None` when this process
|
||||
/// was given no store. Same variables as `hive-agent`'s `swarm_queue`.
|
||||
fn store_source(
|
||||
get: impl Fn(&str) -> Option<String>,
|
||||
) -> Result<Option<StoreSource>, swarm_secret_client::Error> {
|
||||
if get(ENV_ADDR).is_none_or(|v| v.is_empty()) {
|
||||
return Ok(None);
|
||||
}
|
||||
let agent = get(ENV_AGENT_NAME)
|
||||
.filter(|v| !v.is_empty())
|
||||
.ok_or(swarm_secret_client::Error::MissingEnv(ENV_AGENT_NAME))?;
|
||||
let settings =
|
||||
Settings::from_lookup(|k| get(k).filter(|v| k != ENV_CACERT || ca_is_usable(v)))?;
|
||||
Ok(Some(StoreSource {
|
||||
settings,
|
||||
role: policy::agent_object_name(&agent)?,
|
||||
path: provider_key_path(&agent)?,
|
||||
}))
|
||||
}
|
||||
|
||||
/// Whether the CA bundle at `path` is a file with bytes in it. Absent means
|
||||
/// the container's own trust store.
|
||||
fn ca_is_usable(path: &str) -> bool {
|
||||
std::fs::metadata(path).is_ok_and(|m| m.len() > 0)
|
||||
}
|
||||
|
||||
/// The variable to add for `name` to the agent's environment, if any.
|
||||
pub(super) async fn resolve(name: &str) -> Option<KeyVar> {
|
||||
let inherited = std::env::var(name).ok();
|
||||
resolve_with(
|
||||
name,
|
||||
inherited,
|
||||
|| {
|
||||
store_source(|k| std::env::var(k).ok()).unwrap_or_else(|e| {
|
||||
tracing::warn!(
|
||||
error = %e,
|
||||
"this agent's secret store coordinates are incomplete, so its ACP \
|
||||
provider key cannot be read from the store"
|
||||
);
|
||||
None
|
||||
})
|
||||
},
|
||||
STORE_READ_TIMEOUT,
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
/// [`resolve`] over an inherited value and a store, each logged once without
|
||||
/// the key.
|
||||
async fn resolve_with<S: KeySource>(
|
||||
name: &str,
|
||||
inherited: Option<String>,
|
||||
source: impl FnOnce() -> Option<S>,
|
||||
timeout: Duration,
|
||||
) -> Option<KeyVar> {
|
||||
if inherited.is_some_and(|v| !v.trim().is_empty()) {
|
||||
tracing::info!(var = name, "ACP provider key is set in the environment");
|
||||
return None;
|
||||
}
|
||||
let Some(source) = source() else {
|
||||
tracing::info!(
|
||||
var = name,
|
||||
"ACP provider key is not set and this agent has no secret store; spawning without it"
|
||||
);
|
||||
return None;
|
||||
};
|
||||
let value = match tokio::time::timeout(timeout, source.fetch()).await {
|
||||
Ok(Ok(value)) => value.map(|v| v.trim().to_owned()).filter(|v| !v.is_empty()),
|
||||
Ok(Err(e)) => {
|
||||
// `%e`, not the source chain: a decode error's source can quote the stored value.
|
||||
tracing::warn!(
|
||||
var = name,
|
||||
path = source.origin(),
|
||||
error = %e,
|
||||
"reading the ACP provider key from the store failed; spawning without it"
|
||||
);
|
||||
return None;
|
||||
}
|
||||
Err(_) => {
|
||||
tracing::warn!(
|
||||
var = name,
|
||||
path = source.origin(),
|
||||
"the store did not answer within {timeout:?}; spawning the ACP agent without \
|
||||
its provider key"
|
||||
);
|
||||
return None;
|
||||
}
|
||||
};
|
||||
let Some(value) = value else {
|
||||
tracing::info!(
|
||||
var = name,
|
||||
path = source.origin(),
|
||||
"no ACP provider key is stored; spawning without it"
|
||||
);
|
||||
return None;
|
||||
};
|
||||
tracing::info!(
|
||||
var = name,
|
||||
path = source.origin(),
|
||||
"ACP provider key read from the store"
|
||||
);
|
||||
Some(KeyVar {
|
||||
name: name.to_owned(),
|
||||
value,
|
||||
})
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use std::cell::Cell;
|
||||
|
||||
use super::*;
|
||||
|
||||
const VAR: &str = "ACP_PROVIDER_API_KEY";
|
||||
|
||||
/// A store answering `answer`, counting how often it is asked.
|
||||
struct Fake {
|
||||
answer: fn() -> Result<Option<String>, swarm_secret_client::Error>,
|
||||
asked: Cell<u32>,
|
||||
}
|
||||
|
||||
impl Fake {
|
||||
fn new(answer: fn() -> Result<Option<String>, swarm_secret_client::Error>) -> Self {
|
||||
Self {
|
||||
answer,
|
||||
asked: Cell::new(0),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl KeySource for &Fake {
|
||||
fn origin(&self) -> &'static str {
|
||||
"swarm/agents/a1/acp-provider"
|
||||
}
|
||||
|
||||
async fn fetch(&self) -> Result<Option<String>, swarm_secret_client::Error> {
|
||||
self.asked.set(self.asked.get() + 1);
|
||||
(self.answer)()
|
||||
}
|
||||
}
|
||||
|
||||
async fn resolve(inherited: Option<&str>, store: Option<&Fake>) -> Option<(String, String)> {
|
||||
resolve_with(
|
||||
VAR,
|
||||
inherited.map(str::to_owned),
|
||||
|| store,
|
||||
Duration::from_secs(1),
|
||||
)
|
||||
.await
|
||||
.map(|k| (k.name, k.value))
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn an_inherited_key_wins_and_the_store_is_not_asked() {
|
||||
let store = Fake::new(|| Ok(Some("from-store".into())));
|
||||
assert!(resolve(Some("from-file"), Some(&store)).await.is_none());
|
||||
assert_eq!(store.asked.get(), 0);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn without_an_inherited_key_the_stored_one_is_added() {
|
||||
let store = Fake::new(|| Ok(Some("from-store\n".into())));
|
||||
assert_eq!(
|
||||
resolve(None, Some(&store)).await,
|
||||
Some((VAR.to_owned(), "from-store".to_owned()))
|
||||
);
|
||||
// An empty inherited value is no value.
|
||||
assert!(resolve(Some(" "), Some(&store)).await.is_some());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn neither_adds_nothing() {
|
||||
assert!(resolve(None, None).await.is_none());
|
||||
let empty = Fake::new(|| Ok(None));
|
||||
assert!(resolve(None, Some(&empty)).await.is_none());
|
||||
let blank = Fake::new(|| Ok(Some(" ".into())));
|
||||
assert!(resolve(None, Some(&blank)).await.is_none());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn an_unreachable_store_adds_nothing() {
|
||||
let failing = Fake::new(|| Err(swarm_secret_client::Error::MissingEnv("BAO_ADDR")));
|
||||
assert!(resolve(None, Some(&failing)).await.is_none());
|
||||
assert_eq!(failing.asked.get(), 1);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn no_store_address_means_no_store() {
|
||||
let source = store_source(|_| None).expect("an absent store is legal");
|
||||
assert!(source.is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_store_without_the_agent_name_is_refused() {
|
||||
let Err(e) = store_source(|k| (k == ENV_ADDR).then(|| "https://bao:8200".to_owned()))
|
||||
else {
|
||||
panic!("a store address without an agent name is a half-delivered container");
|
||||
};
|
||||
assert!(e.to_string().contains(ENV_AGENT_NAME), "{e}");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn the_store_source_reads_the_agents_own_path() {
|
||||
let env = [
|
||||
(ENV_ADDR, "https://bao:8200"),
|
||||
(ENV_AGENT_NAME, "a1"),
|
||||
("BAO_CLIENT_CERT", "/run/credentials/c"),
|
||||
("BAO_CLIENT_KEY", "/run/credentials/k"),
|
||||
];
|
||||
let source = store_source(|k| {
|
||||
env.iter()
|
||||
.find(|(n, _)| *n == k)
|
||||
.map(|(_, v)| (*v).to_owned())
|
||||
})
|
||||
.expect("a complete environment")
|
||||
.expect("a store address was given");
|
||||
assert_eq!(source.path, "swarm/agents/a1/acp-provider");
|
||||
assert_eq!(source.role, policy::agent_object_name("a1").unwrap());
|
||||
}
|
||||
}
|
||||
|
|
@ -12,6 +12,7 @@ use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
|
|||
use tokio::process::{Child, ChildStdin, Command};
|
||||
use tokio::sync::{mpsc, oneshot};
|
||||
|
||||
use super::provider_key::KeyVar;
|
||||
use super::stream::mcp_server_of;
|
||||
use super::{AcpError, PermissionAsk, PermissionPolicy};
|
||||
use crate::spec::AcpCommand;
|
||||
|
|
@ -38,16 +39,21 @@ pub(super) struct Connection {
|
|||
}
|
||||
|
||||
impl Connection {
|
||||
/// Spawn the agent in `cwd` and start reading its output.
|
||||
/// Spawn the agent in `cwd`, with `key` added to its environment, and
|
||||
/// start reading its output.
|
||||
pub(super) fn spawn(
|
||||
command: &AcpCommand,
|
||||
key: Option<&KeyVar>,
|
||||
cwd: &Path,
|
||||
permit: PermissionPolicy,
|
||||
servers: Vec<String>,
|
||||
) -> Result<Self, AcpError> {
|
||||
let mut child = Command::new(&command.command)
|
||||
.args(&command.args)
|
||||
.envs(&command.env)
|
||||
let mut cmd = Command::new(&command.command);
|
||||
cmd.args(&command.args).envs(&command.env);
|
||||
if let Some(key) = key {
|
||||
cmd.env(&key.name, &key.value);
|
||||
}
|
||||
let mut child = cmd
|
||||
.current_dir(cwd)
|
||||
.stdin(Stdio::piped())
|
||||
.stdout(Stdio::piped())
|
||||
|
|
|
|||
|
|
@ -11,6 +11,10 @@ pub const ACP_ARGS_ENV: &str = "HIVE_ACP_ARGS";
|
|||
/// Extra environment for the ACP agent only, as a JSON object of strings.
|
||||
/// Optional. The agent also inherits the harness's own environment.
|
||||
pub const ACP_ENV_ENV: &str = "HIVE_ACP_ENV";
|
||||
/// The variable the ACP agent reads its provider API key from. Optional. When
|
||||
/// set and the process environment leaves that variable unset, the agent is
|
||||
/// spawned with it read from the swarm secret store.
|
||||
pub const ACP_API_KEY_ENV_ENV: &str = "HIVE_ACP_API_KEY_ENV";
|
||||
|
||||
/// The runtime an agent is configured with.
|
||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||
|
|
@ -25,6 +29,8 @@ pub struct AcpCommand {
|
|||
pub command: String,
|
||||
pub args: Vec<String>,
|
||||
pub env: BTreeMap<String, String>,
|
||||
/// See [`ACP_API_KEY_ENV_ENV`].
|
||||
pub api_key_env: Option<String>,
|
||||
}
|
||||
|
||||
/// A runtime configuration that cannot be acted on.
|
||||
|
|
@ -63,6 +69,7 @@ impl RuntimeSpec {
|
|||
command,
|
||||
args: json_or_default(&lookup, ACP_ARGS_ENV)?,
|
||||
env: json_or_default(&lookup, ACP_ENV_ENV)?,
|
||||
api_key_env: lookup(ACP_API_KEY_ENV_ENV).filter(|v| !v.trim().is_empty()),
|
||||
}))
|
||||
}
|
||||
other => Err(SpecError::UnknownRuntime(other.to_owned())),
|
||||
|
|
@ -119,6 +126,7 @@ mod tests {
|
|||
"HIVE_ACP_ENV",
|
||||
r#"{"AGENT_CONFIG":"/nix/store/y/config.json"}"#,
|
||||
),
|
||||
("HIVE_ACP_API_KEY_ENV", "ACP_PROVIDER_API_KEY"),
|
||||
])
|
||||
.unwrap();
|
||||
assert_eq!(
|
||||
|
|
@ -127,6 +135,7 @@ mod tests {
|
|||
command: "/nix/store/x/bin/agent".into(),
|
||||
args: vec!["acp".into(), "--flag".into()],
|
||||
env: [("AGENT_CONFIG".into(), "/nix/store/y/config.json".into())].into(),
|
||||
api_key_env: Some("ACP_PROVIDER_API_KEY".into()),
|
||||
})
|
||||
);
|
||||
}
|
||||
|
|
@ -139,6 +148,7 @@ mod tests {
|
|||
};
|
||||
assert!(cmd.args.is_empty());
|
||||
assert!(cmd.env.is_empty());
|
||||
assert_eq!(cmd.api_key_env, None);
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
|
|
|||
Loading…
Reference in a new issue