mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
feat(acp): thread-scoped concurrent sessions — scope-keyed dispatch, affinity, control resolution
Scheduling identity becomes ConversationSessionKey (channel_id, root_event_id):
- queue.rs: all in-flight control state (deadlines, batch sizes, retry,
cancel sets, withheld native steer) keyed by scope; flush_next selects
the oldest non-in-flight scope so distinct roots flush concurrently;
restore_unclaimed is the exact lossless undo of flush_next.
- pool.rs: TaskMeta carries the scope; try_claim returns
ClaimOutcome{Claimed,BusyOwner,Exhausted} with strict retained-root
slot affinity and lazy owner pruning; send_steer scope-keyed.
- lib.rs: dispatch_pending dispatches multiple scopes per pass, holding
BusyOwner/Exhausted batches until after the loop (prevents
flush→restore→reflush livelock) then restoring losslessly; control
resolution never fans out across roots — explicit root targets
exactly, channel-level control resolves only when exactly one scope
is in flight (Err(n) = ambiguous, touch nothing); steer fork is
scope-exact. Channel mode remains the degenerate root=None path.
Deterministic coverage: cross-root concurrent dispatch, busy-owner skip
without root migration, pool-exhaustion lossless restore, three
restore_unclaimed exactness tests, control-scope resolution (channel
exactly-one rule, explicit control-frame root, threaded-event root
targeting, channel-mode fallback). cargo test -p buzz-acp: 557 passed.
Co-authored-by: Tyler Longwell <tlongwell@block.xyz>
Signed-off-by: Tyler Longwell <tlongwell@block.xyz>
This commit is contained in:
co-authored by
Tyler Longwell
parent
2444b17b37
commit
eca05aa257
+666
-161
File diff suppressed because it is too large
Load Diff
+64
-47
@@ -40,13 +40,18 @@ use crate::queue::{
|
|||||||
};
|
};
|
||||||
use crate::relay::{ChannelInfo, RestClient};
|
use crate::relay::{ChannelInfo, RestClient};
|
||||||
|
|
||||||
|
// Scheduling/session identity lives in queue.rs next to the scope-keyed
|
||||||
|
// queue state machine; re-exported here so pool consumers keep one name.
|
||||||
|
pub use crate::queue::ConversationSessionKey;
|
||||||
|
|
||||||
// FlushBatch and BatchEvent derive Clone (added in queue.rs) so we can store
|
// FlushBatch and BatchEvent derive Clone (added in queue.rs) so we can store
|
||||||
// a recoverable copy in TaskMeta for panic recovery in Queue mode.
|
// a recoverable copy in TaskMeta for panic recovery in Queue mode.
|
||||||
|
|
||||||
/// Metadata stored per in-flight task for panic recovery.
|
/// Metadata stored per in-flight task for panic recovery.
|
||||||
pub struct TaskMeta {
|
pub struct TaskMeta {
|
||||||
pub agent_index: usize,
|
pub agent_index: usize,
|
||||||
pub channel_id: Option<Uuid>,
|
/// Conversation scope of the in-flight prompt; `None` for heartbeats.
|
||||||
|
pub scope: Option<ConversationSessionKey>,
|
||||||
/// Identifies terminal events when the task panics before returning a result.
|
/// Identifies terminal events when the task panics before returning a result.
|
||||||
pub turn_id: String,
|
pub turn_id: String,
|
||||||
/// Clone of batch for Queue mode panic recovery.
|
/// Clone of batch for Queue mode panic recovery.
|
||||||
@@ -74,22 +79,6 @@ pub struct AgentModelCapabilities {
|
|||||||
pub available_models_raw: Option<serde_json::Value>,
|
pub available_models_raw: Option<serde_json::Value>,
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
|
|
||||||
pub struct ConversationSessionKey {
|
|
||||||
pub channel_id: Uuid,
|
|
||||||
pub root_event_id: Option<String>,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl ConversationSessionKey {
|
|
||||||
#[cfg(test)]
|
|
||||||
pub fn channel(channel_id: Uuid) -> Self {
|
|
||||||
Self {
|
|
||||||
channel_id,
|
|
||||||
root_event_id: None,
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Per-channel session IDs and turn counters.
|
/// Per-channel session IDs and turn counters.
|
||||||
///
|
///
|
||||||
/// Separated from `OwnedAgent` so the state machine is testable without
|
/// Separated from `OwnedAgent` so the state machine is testable without
|
||||||
@@ -579,25 +568,32 @@ impl AgentPool {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Try to claim an idle agent for the given channel (or heartbeat if `None`).
|
/// Try to claim an idle agent for the given scope (or heartbeat if `None`).
|
||||||
///
|
///
|
||||||
/// Pass 1: prefer an agent that already has a session for `channel_id`.
|
/// Pass 1: prefer an agent that already has a session for the scope.
|
||||||
/// Pass 2: any idle agent.
|
/// Pass 2: any idle agent.
|
||||||
///
|
///
|
||||||
/// Returns `None` if all agents are checked out.
|
/// Retained roots have strict slot affinity: if the root's owning slot is
|
||||||
pub fn try_claim(
|
/// busy, the claim returns [`ClaimOutcome::BusyOwner`] so the dispatcher
|
||||||
&mut self,
|
/// can skip this scope and keep dispatching other scopes — a busy owner
|
||||||
session_key: Option<&ConversationSessionKey>,
|
/// must not stall unrelated roots, and the root must never migrate to
|
||||||
) -> Option<OwnedAgent> {
|
/// another slot (that would fork the provider session).
|
||||||
|
/// [`ClaimOutcome::Exhausted`] means no idle agent exists at all, so the
|
||||||
|
/// dispatch loop should stop.
|
||||||
|
pub fn try_claim(&mut self, session_key: Option<&ConversationSessionKey>) -> ClaimOutcome {
|
||||||
// Pass 1: retained roots have strict ownership. If their slot is busy,
|
// Pass 1: retained roots have strict ownership. If their slot is busy,
|
||||||
// wait rather than creating a duplicate provider session in another slot.
|
// wait rather than creating a duplicate provider session in another slot.
|
||||||
if let Some(key) = session_key {
|
if let Some(key) = session_key {
|
||||||
if key.root_event_id.is_some() {
|
if key.root_event_id.is_some() {
|
||||||
self.prune_session_owners();
|
self.prune_session_owners();
|
||||||
if let Some(&owner) = self.session_owners.get(key) {
|
if let Some(&owner) = self.session_owners.get(key) {
|
||||||
// Surviving the prune means the owner is busy (wait by
|
// Surviving the prune means the owner is busy (skip this
|
||||||
// returning None) or idle and retaining the session.
|
// scope, keep dispatching others) or idle and retaining
|
||||||
return self.agents.get_mut(owner).and_then(Option::take);
|
// the session.
|
||||||
|
return match self.agents.get_mut(owner).and_then(Option::take) {
|
||||||
|
Some(agent) => ClaimOutcome::Claimed(agent),
|
||||||
|
None => ClaimOutcome::BusyOwner,
|
||||||
|
};
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
let idx = self.agents.iter().position(|slot| {
|
let idx = self.agents.iter().position(|slot| {
|
||||||
@@ -607,17 +603,21 @@ impl AgentPool {
|
|||||||
});
|
});
|
||||||
if let Some(i) = idx {
|
if let Some(i) = idx {
|
||||||
self.session_owners.insert(key.clone(), i);
|
self.session_owners.insert(key.clone(), i);
|
||||||
return self.agents[i].take();
|
return ClaimOutcome::Claimed(
|
||||||
|
self.agents[i].take().expect("position matched idle slot"),
|
||||||
|
);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// Pass 2: first idle agent. Reserve root affinity immediately so another
|
// Pass 2: first idle agent. Reserve root affinity immediately so another
|
||||||
// turn for the same root cannot fall through while this slot is checked out.
|
// turn for the same root cannot fall through while this slot is checked out.
|
||||||
let idx = self.agents.iter().position(|slot| slot.is_some())?;
|
let Some(idx) = self.agents.iter().position(|slot| slot.is_some()) else {
|
||||||
|
return ClaimOutcome::Exhausted;
|
||||||
|
};
|
||||||
if let Some(key) = session_key.filter(|key| key.root_event_id.is_some()) {
|
if let Some(key) = session_key.filter(|key| key.root_event_id.is_some()) {
|
||||||
self.session_owners.insert(key.clone(), idx);
|
self.session_owners.insert(key.clone(), idx);
|
||||||
}
|
}
|
||||||
self.agents[idx].take()
|
ClaimOutcome::Claimed(self.agents[idx].take().expect("position matched idle slot"))
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Drop root reservations that no longer bind: the owning slot is neither
|
/// Drop root reservations that no longer bind: the owning slot is neither
|
||||||
@@ -703,19 +703,19 @@ impl AgentPool {
|
|||||||
/// watcher, to close the result-vs-ack race.
|
/// watcher, to close the result-vs-ack race.
|
||||||
///
|
///
|
||||||
/// Returns `Err(SteerError::PromptCompleted)` if no task is in flight
|
/// Returns `Err(SteerError::PromptCompleted)` if no task is in flight
|
||||||
/// for `channel_id` (the prompt completed between the mode-gate check
|
/// for `scope` (the prompt completed between the mode-gate check
|
||||||
/// and this call, or the channel was never in flight). This is
|
/// and this call, or the scope was never in flight). This is
|
||||||
/// semantically a soft no-op — the caller should release any withheld
|
/// semantically a soft no-op — the caller should release any withheld
|
||||||
/// event and let normal dispatch handle delivery.
|
/// event and let normal dispatch handle delivery.
|
||||||
pub fn send_steer(
|
pub fn send_steer(
|
||||||
&mut self,
|
&mut self,
|
||||||
channel_id: Uuid,
|
scope: &ConversationSessionKey,
|
||||||
request: SteerRequest,
|
request: SteerRequest,
|
||||||
) -> Result<(), SteerError> {
|
) -> Result<(), SteerError> {
|
||||||
let meta = self
|
let meta = self
|
||||||
.task_map
|
.task_map
|
||||||
.values_mut()
|
.values_mut()
|
||||||
.find(|m| m.channel_id == Some(channel_id))
|
.find(|m| m.scope.as_ref() == Some(scope))
|
||||||
.ok_or(SteerError::PromptCompleted)?;
|
.ok_or(SteerError::PromptCompleted)?;
|
||||||
let tx = meta
|
let tx = meta
|
||||||
.steer_tx
|
.steer_tx
|
||||||
@@ -855,6 +855,33 @@ impl AgentPool {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Outcome of [`AgentPool::try_claim`].
|
||||||
|
// The variant size skew is fine: a ClaimOutcome is matched and consumed
|
||||||
|
// immediately at the claim site, never stored — boxing would add churn
|
||||||
|
// on the hot dispatch path for no benefit.
|
||||||
|
#[allow(clippy::large_enum_variant)]
|
||||||
|
pub enum ClaimOutcome {
|
||||||
|
/// An idle agent was checked out for the scope.
|
||||||
|
Claimed(OwnedAgent),
|
||||||
|
/// The scope is a retained root whose owning slot is busy on another
|
||||||
|
/// turn. Skip this scope — do NOT stop the dispatch loop, and do NOT
|
||||||
|
/// run the root on a different slot (that would fork its session).
|
||||||
|
BusyOwner,
|
||||||
|
/// No idle agent in the pool. Stop the dispatch loop.
|
||||||
|
Exhausted,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl ClaimOutcome {
|
||||||
|
/// The claimed agent, if any. Test-only convenience.
|
||||||
|
#[cfg(test)]
|
||||||
|
pub fn claimed(self) -> Option<OwnedAgent> {
|
||||||
|
match self {
|
||||||
|
ClaimOutcome::Claimed(agent) => Some(agent),
|
||||||
|
ClaimOutcome::BusyOwner | ClaimOutcome::Exhausted => None,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/// Outcome of [`AgentPool::switch_idle_agent_model`].
|
/// Outcome of [`AgentPool::switch_idle_agent_model`].
|
||||||
#[derive(Debug, PartialEq, Eq)]
|
#[derive(Debug, PartialEq, Eq)]
|
||||||
pub enum IdleSwitchResult {
|
pub enum IdleSwitchResult {
|
||||||
@@ -1340,16 +1367,6 @@ fn send_prompt_result(
|
|||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Session identity for a batch. `FlushBatch::conversation_root` is populated
|
|
||||||
/// at queue-push time only when the top-level-sessions experiment is enabled,
|
|
||||||
/// so this is purely data-driven: no root means the legacy channel-scoped key.
|
|
||||||
pub(crate) fn conversation_session_key(batch: &FlushBatch) -> ConversationSessionKey {
|
|
||||||
ConversationSessionKey {
|
|
||||||
channel_id: batch.channel_id,
|
|
||||||
root_event_id: batch.conversation_root.clone(),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Core async function spawned for each prompt.
|
/// Core async function spawned for each prompt.
|
||||||
///
|
///
|
||||||
/// Lifecycle:
|
/// Lifecycle:
|
||||||
@@ -1373,7 +1390,7 @@ pub async fn run_prompt_task(
|
|||||||
) {
|
) {
|
||||||
// Is this a channel prompt or a heartbeat?
|
// Is this a channel prompt or a heartbeat?
|
||||||
let source = match &batch {
|
let source = match &batch {
|
||||||
Some(b) => PromptSource::Channel(conversation_session_key(b)),
|
Some(b) => PromptSource::Channel(b.scope_key()),
|
||||||
None => PromptSource::Heartbeat,
|
None => PromptSource::Heartbeat,
|
||||||
};
|
};
|
||||||
let observer_channel_id = match &source {
|
let observer_channel_id = match &source {
|
||||||
@@ -4557,7 +4574,7 @@ mod tests {
|
|||||||
#[test]
|
#[test]
|
||||||
fn test_conversation_session_key_is_channel_scoped_without_root() {
|
fn test_conversation_session_key_is_channel_scoped_without_root() {
|
||||||
let batch = one_event_batch(Uuid::new_v4());
|
let batch = one_event_batch(Uuid::new_v4());
|
||||||
let key = conversation_session_key(&batch);
|
let key = batch.scope_key();
|
||||||
assert_eq!(key.channel_id, batch.channel_id);
|
assert_eq!(key.channel_id, batch.channel_id);
|
||||||
assert_eq!(key.root_event_id, None);
|
assert_eq!(key.root_event_id, None);
|
||||||
}
|
}
|
||||||
@@ -4566,7 +4583,7 @@ mod tests {
|
|||||||
fn test_conversation_session_key_uses_batch_root_when_present() {
|
fn test_conversation_session_key_uses_batch_root_when_present() {
|
||||||
let mut batch = one_event_batch(Uuid::new_v4());
|
let mut batch = one_event_batch(Uuid::new_v4());
|
||||||
batch.conversation_root = Some("root-a".into());
|
batch.conversation_root = Some("root-a".into());
|
||||||
let key = conversation_session_key(&batch);
|
let key = batch.scope_key();
|
||||||
assert_eq!(key.channel_id, batch.channel_id);
|
assert_eq!(key.channel_id, batch.channel_id);
|
||||||
assert_eq!(key.root_event_id.as_deref(), Some("root-a"));
|
assert_eq!(key.root_event_id.as_deref(), Some("root-a"));
|
||||||
}
|
}
|
||||||
|
|||||||
+579
-415
File diff suppressed because it is too large
Load Diff
Reference in New Issue
Block a user