chore: merge develop into hare/develop

This commit is contained in:
2026-08-29 13:25:32 +09:00
88 changed files with 8005 additions and 1472 deletions
+2 -1
View File
@@ -66,11 +66,12 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
WorkerRunResult::Finished => println!("(finished)"),
WorkerRunResult::Paused => println!("(paused)"),
WorkerRunResult::LimitReached => println!("(turn limit reached)"),
WorkerRunResult::Interrupted { message, .. } => println!("(interrupted: {message})"),
WorkerRunResult::RolledBack => println!("(empty turn rolled back)"),
}
// 5. Extract the assistant's reply from history
let history = worker.engine().history();
let history = worker.history();
if let Some(text) = history
.iter()
.rev()
+1 -1
View File
@@ -22,7 +22,7 @@ use crate::compact::token_counter::{
EstimateSource, savings_for_prune_impl, token_estimates_for_prune_impl,
};
impl<C: LlmClient, St: Store> Worker<C, St> {
impl<C: LlmClient + 'static, St: Store> Worker<C, St> {
/// Enable prune projection on the underlying Engine.
///
/// Registers the config and token/savings-estimator closures on the Engine.
+4 -4
View File
@@ -242,13 +242,13 @@ pub(crate) fn savings_for_prune_impl(
// ── Worker に生やす公開 API ───────────────────────────────────────────────
impl<C: LlmClient, St: Store> Worker<C, St> {
impl<C: LlmClient + 'static, St: Store> Worker<C, St> {
/// 現在の history 全体の推定トークン数。
///
/// 最後の measurement と、その後に追加された未測定分の byte/4 外挿。
pub fn total_tokens(&self) -> TokenEstimate {
let usage = self.usage_history();
agen::token_counter::total_tokens(self.history(), &usage)
agen::token_counter::total_tokens(&self.history(), &usage)
}
/// 任意の history index 時点でのプロンプト全長推定。
@@ -259,7 +259,7 @@ impl<C: LlmClient, St: Store> Worker<C, St> {
/// pointer 以降に増えたプロンプト長を測るのに使う。
pub fn total_tokens_at(&self, history_len: usize) -> TokenEstimate {
let usage = self.usage_history();
agen::token_counter::total_tokens_at(self.history(), &usage, history_len)
agen::token_counter::total_tokens_at(&self.history(), &usage, history_len)
}
/// 末尾から `retained` トークン以上を残すための分割位置。
@@ -267,7 +267,7 @@ impl<C: LlmClient, St: Store> Worker<C, St> {
/// `history[..cut.index]` が要約/破棄される側、`history[cut.index..]` が残る側。
pub fn split_for_retained(&self, retained: u64) -> SplitPoint {
let usage = self.usage_history();
split_for_retained_impl(self.history(), &usage, retained)
split_for_retained_impl(&self.history(), &usage, retained)
}
}
+178 -19
View File
@@ -485,6 +485,7 @@ impl WorkerController {
// into the controller task so the in-flight turn can be reached
// via these handles while worker itself is borrowed by drive_turn.
let cancel_tx = worker.engine_mut().cancel_sender();
let pause_tx = worker.engine_mut().pause_sender();
let notify_buffer = worker.notify_buffer_handle();
tokio::spawn(controller_loop(
@@ -494,6 +495,7 @@ impl WorkerController {
shared_state,
runtime_dir,
cancel_tx,
pause_tx,
notify_buffer,
self_parent_socket,
spawner_name,
@@ -763,6 +765,19 @@ pub(crate) fn wire_event_bridges_on_engine<C, St>(
id: result.tool_use_id.clone(),
summary: result.summary.clone(),
output: result.content.clone(),
disposition: Some(match result.disposition {
agen::ToolResultDisposition::Success => protocol::ToolResultDisposition::Success,
agen::ToolResultDisposition::Error => protocol::ToolResultDisposition::Error,
agen::ToolResultDisposition::Interrupted => {
protocol::ToolResultDisposition::Interrupted
}
agen::ToolResultDisposition::Cancelled => {
protocol::ToolResultDisposition::Cancelled
}
agen::ToolResultDisposition::OutcomeUnknown => {
protocol::ToolResultDisposition::OutcomeUnknown
}
}),
is_error: result.is_error,
});
});
@@ -1115,6 +1130,7 @@ async fn controller_loop<C, St>(
shared_state: Arc<WorkerSharedState>,
runtime_dir: Arc<RuntimeDir>,
cancel_tx: mpsc::Sender<()>,
pause_tx: mpsc::Sender<()>,
notify_buffer: NotifyBuffer,
self_parent_socket: Option<PathBuf>,
spawner_name: String,
@@ -1161,22 +1177,35 @@ async fn controller_loop<C, St>(
// clear at run start prevents stale partial output left by an older
// interrupted/error turn from being carried into the next snapshot.
worker.clear_in_flight_events();
set_controller_status(
&shared_state,
&runtime_dir,
&event_tx,
WorkerStatus::Running,
)
.await;
let parent_originated = run.is_parent_originated();
let user_input_run = matches!(&run, PendingRun::Run(_) | PendingRun::RunTracked { .. });
if !user_input_run {
set_controller_status(
&shared_state,
&runtime_dir,
&event_tx,
WorkerStatus::Running,
)
.await;
}
let (mut new_status, shutdown) = match run {
PendingRun::Run(input) => {
let (input_commit_tx, input_commit_rx) = oneshot::channel();
drive_turn(
worker.run(input),
worker.run_with_input_extensions_and_commit_hook(
input,
Vec::new(),
move || {
let _ = input_commit_tx.send(());
},
),
&mut method_rx,
&event_tx,
&cancel_tx,
&pause_tx,
&shared_state,
&runtime_dir,
Some(input_commit_rx),
&notify_buffer,
self_parent_socket.as_ref(),
&spawner_name,
@@ -1186,12 +1215,22 @@ async fn controller_loop<C, St>(
.await
}
PendingRun::RunTracked { input, extension } => {
let (input_commit_tx, input_commit_rx) = oneshot::channel();
drive_turn(
worker.run_with_input_extensions(input, vec![extension]),
worker.run_with_input_extensions_and_commit_hook(
input,
vec![extension],
move || {
let _ = input_commit_tx.send(());
},
),
&mut method_rx,
&event_tx,
&cancel_tx,
&pause_tx,
&shared_state,
&runtime_dir,
Some(input_commit_rx),
&notify_buffer,
self_parent_socket.as_ref(),
&spawner_name,
@@ -1206,7 +1245,10 @@ async fn controller_loop<C, St>(
&mut method_rx,
&event_tx,
&cancel_tx,
&pause_tx,
&shared_state,
&runtime_dir,
None,
&notify_buffer,
self_parent_socket.as_ref(),
&spawner_name,
@@ -1221,7 +1263,10 @@ async fn controller_loop<C, St>(
&mut method_rx,
&event_tx,
&cancel_tx,
&pause_tx,
&shared_state,
&runtime_dir,
None,
&notify_buffer,
self_parent_socket.as_ref(),
&spawner_name,
@@ -1346,7 +1391,7 @@ async fn controller_loop<C, St>(
});
}
},
WorkerStatus::Idle => {
WorkerStatus::Idle | WorkerStatus::Stopped => {
let _ = event_tx.send(Event::Error {
code: ErrorCode::NotRunning,
message: "Worker is not running".into(),
@@ -1387,7 +1432,7 @@ async fn controller_loop<C, St>(
.into(),
});
}
WorkerStatus::Running => {
WorkerStatus::Running | WorkerStatus::Stopped => {
let _ = event_tx.send(Event::Error {
code: ErrorCode::AlreadyRunning,
message:
@@ -1401,7 +1446,7 @@ async fn controller_loop<C, St>(
WorkerStatus::Idle | WorkerStatus::Paused => {
emit_rewind_targets(&worker, &event_tx)
}
WorkerStatus::Running => {
WorkerStatus::Running | WorkerStatus::Stopped => {
let _ = event_tx.send(Event::Error {
code: ErrorCode::AlreadyRunning,
message: "Worker is already executing a turn; rewind can only run while idle or paused"
@@ -1430,7 +1475,7 @@ async fn controller_loop<C, St>(
.into(),
});
}
WorkerStatus::Running => {
WorkerStatus::Running | WorkerStatus::Stopped => {
let _ = event_tx.send(Event::Error {
code: ErrorCode::AlreadyRunning,
message: "Worker is already executing a turn; rewind can only run while idle or paused"
@@ -1618,7 +1663,10 @@ async fn drive_turn<F>(
method_rx: &mut mpsc::Receiver<Method>,
event_tx: &broadcast::Sender<Event>,
cancel_tx: &mpsc::Sender<()>,
pause_tx: &mpsc::Sender<()>,
shared_state: &Arc<WorkerSharedState>,
runtime_dir: &RuntimeDir,
mut input_commit_rx: Option<oneshot::Receiver<()>>,
notify_buffer: &NotifyBuffer,
parent_socket: Option<&PathBuf>,
self_name: &str,
@@ -1634,14 +1682,58 @@ where
loop {
tokio::select! {
// If input commit and provider completion become ready together, expose
// Running only after processing the commit fence. This makes the
// Running snapshot contract deterministic even for immediate clients.
biased;
committed = async {
input_commit_rx
.as_mut()
.expect("input commit receiver guarded by select condition")
.await
}, if input_commit_rx.is_some() => {
input_commit_rx = None;
if committed.is_ok() {
set_controller_status(
shared_state,
runtime_dir,
event_tx,
WorkerStatus::Running,
)
.await;
}
}
result = &mut worker_future => {
return match result {
Ok(r) => {
let (status, run_result) = match r {
WorkerRunResult::Finished if pause_requested => {
(WorkerStatus::Paused, RunResult::Paused)
}
WorkerRunResult::Finished => (WorkerStatus::Idle, RunResult::Finished),
WorkerRunResult::Paused => (WorkerStatus::Paused, RunResult::Paused),
WorkerRunResult::LimitReached => (WorkerStatus::Idle, RunResult::LimitReached),
WorkerRunResult::RolledBack => (WorkerStatus::Idle, RunResult::RolledBack),
WorkerRunResult::Interrupted { .. } if pause_requested => {
let _ = event_tx.send(Event::RunEnd { result: RunResult::Paused });
return (WorkerStatus::Paused, shutdown_requested);
}
WorkerRunResult::Interrupted { code, message } => {
let _ = event_tx.send(Event::Error {
code,
message: message.clone(),
});
if parent_originated {
crate::ipc::event::fire_and_forget(
parent_socket.cloned(),
protocol::WorkerEvent::Errored {
worker_name: self_name.to_string(),
message,
},
);
}
return (WorkerStatus::Idle, shutdown_requested);
}
};
let _ = event_tx.send(Event::RunEnd { result: run_result });
if parent_originated && matches!(run_result, RunResult::Finished) {
@@ -1690,7 +1782,7 @@ where
}
Some(Method::Pause) => {
pause_requested = true;
let _ = cancel_tx.try_send(());
let _ = pause_tx.try_send(());
}
Some(Method::Shutdown) => {
shutdown_requested = true;
@@ -1752,7 +1844,7 @@ where
fn emit_rewind_targets<C, St>(worker: &Worker<C, St>, event_tx: &broadcast::Sender<Event>)
where
C: LlmClient,
C: LlmClient + 'static,
St: Store,
{
match worker.list_rewind_targets() {
@@ -1778,7 +1870,7 @@ fn apply_rewind<C, St>(
expected_head_entries: usize,
) -> bool
where
C: LlmClient,
C: LlmClient + 'static,
St: Store,
{
match worker.rewind_to(target, expected_head_entries) {
@@ -1826,7 +1918,7 @@ fn model_supports_image_attachments(model: &manifest::ModelManifest) -> bool {
fn build_greeting<C, St>(worker: &Worker<C, St>) -> protocol::Greeting
where
C: LlmClient,
C: LlmClient + 'static,
St: Store,
{
let manifest = worker.manifest();
@@ -1942,11 +2034,13 @@ mod tests {
event_tx: broadcast::Sender<Event>,
cancel_tx: mpsc::Sender<()>,
_cancel_rx: mpsc::Receiver<()>,
pause_tx: mpsc::Sender<()>,
_pause_rx: mpsc::Receiver<()>,
shared_state: Arc<WorkerSharedState>,
notify_buffer: NotifyBuffer,
spawned_registry: Arc<SpawnedWorkerRegistry>,
parent_socket_path: PathBuf,
_runtime_dir: Arc<RuntimeDir>,
runtime_dir: Arc<RuntimeDir>,
_temp: TempDir,
}
@@ -1960,6 +2054,7 @@ mod tests {
let (method_tx, method_rx) = mpsc::channel::<Method>(16);
let (event_tx, _) = broadcast::channel::<Event>(16);
let (cancel_tx, cancel_rx) = mpsc::channel::<()>(1);
let (pause_tx, pause_rx) = mpsc::channel::<()>(1);
let shared_state = Arc::new(WorkerSharedState::new(
"child-worker".to_string(),
session_store::new_segment_id(),
@@ -1985,11 +2080,13 @@ mod tests {
event_tx,
cancel_tx,
_cancel_rx: cancel_rx,
pause_tx,
_pause_rx: pause_rx,
shared_state,
notify_buffer,
spawned_registry,
parent_socket_path,
_runtime_dir: runtime_dir,
runtime_dir,
_temp: temp,
}
}
@@ -2042,7 +2139,10 @@ mod tests {
&mut env.method_rx,
&env.event_tx,
&env.cancel_tx,
&env.pause_tx,
&env.shared_state,
&env.runtime_dir,
None,
&env.notify_buffer,
Some(&env.parent_socket_path),
"child-worker",
@@ -2063,6 +2163,44 @@ mod tests {
}
}
#[tokio::test]
async fn pause_waits_for_run_boundary_and_uses_safe_pause_channel() {
let mut env = make_env().await;
let method_tx = env._method_tx.clone();
tokio::spawn(async move {
tokio::time::sleep(Duration::from_millis(10)).await;
method_tx.send(Method::Pause).await.expect("send pause");
});
let worker_future = async {
tokio::time::sleep(Duration::from_millis(100)).await;
Ok::<_, WorkerError>(WorkerRunResult::Finished)
};
let started_at = std::time::Instant::now();
let (status, shutdown) = drive_turn(
worker_future,
&mut env.method_rx,
&env.event_tx,
&env.cancel_tx,
&env.pause_tx,
&env.shared_state,
&env.runtime_dir,
None,
&env.notify_buffer,
None,
"child-worker",
&env.spawned_registry,
true,
)
.await;
assert_eq!(status, WorkerStatus::Paused);
assert!(!shutdown);
assert!(started_at.elapsed() >= Duration::from_millis(100));
assert!(env._pause_rx.try_recv().is_ok());
assert!(env._cancel_rx.try_recv().is_err());
}
#[tokio::test]
async fn non_parent_originated_finished_stays_silent() {
let mut env = make_env().await;
@@ -2074,7 +2212,10 @@ mod tests {
&mut env.method_rx,
&env.event_tx,
&env.cancel_tx,
&env.pause_tx,
&env.shared_state,
&env.runtime_dir,
None,
&env.notify_buffer,
Some(&env.parent_socket_path),
"child-worker",
@@ -2109,7 +2250,10 @@ mod tests {
&mut env.method_rx,
&env.event_tx,
&env.cancel_tx,
&env.pause_tx,
&env.shared_state,
&env.runtime_dir,
None,
&env.notify_buffer,
Some(&env.parent_socket_path),
"child-worker",
@@ -2150,7 +2294,10 @@ mod tests {
&mut env.method_rx,
&env.event_tx,
&env.cancel_tx,
&env.pause_tx,
&env.shared_state,
&env.runtime_dir,
None,
&env.notify_buffer,
Some(&env.parent_socket_path),
"child-worker",
@@ -2189,7 +2336,10 @@ mod tests {
&mut env.method_rx,
&env.event_tx,
&env.cancel_tx,
&env.pause_tx,
&env.shared_state,
&env.runtime_dir,
None,
&env.notify_buffer,
Some(&env.parent_socket_path),
"parent",
@@ -2225,7 +2375,10 @@ mod tests {
&mut env.method_rx,
&env.event_tx,
&env.cancel_tx,
&env.pause_tx,
&env.shared_state,
&env.runtime_dir,
None,
&env.notify_buffer,
Some(&env.parent_socket_path),
"parent",
@@ -2259,7 +2412,10 @@ mod tests {
&mut env.method_rx,
&env.event_tx,
&env.cancel_tx,
&env.pause_tx,
&env.shared_state,
&env.runtime_dir,
None,
&env.notify_buffer,
Some(&env.parent_socket_path),
"parent",
@@ -2292,7 +2448,10 @@ mod tests {
&mut env.method_rx,
&env.event_tx,
&env.cancel_tx,
&env.pause_tx,
&env.shared_state,
&env.runtime_dir,
None,
&env.notify_buffer,
Some(&env.parent_socket_path),
"child-worker",
+2 -2
View File
@@ -1795,9 +1795,9 @@ impl FeatureRegistryBuilder {
}
/// Install modules into the existing Engine tool path and hook builder.
pub(crate) fn install_into_engine<C: LlmClient>(
pub(crate) fn install_into_engine<C: LlmClient, A>(
self,
worker: &mut Engine<C, Mutable>,
worker: &mut Engine<C, Mutable, A>,
hook_builder: &mut HookRegistryBuilder,
) -> FeatureRegistryInstallReport {
let mut pending_tools = Vec::new();
+1
View File
@@ -12,6 +12,7 @@ pub mod memory_extract;
pub mod merge_request;
pub mod objective;
pub mod orchestration;
mod resource_projection;
pub mod session_explore;
pub mod task;
pub mod ticket;
@@ -6,7 +6,9 @@ use memory::backend::{
MemoryBackendOperation, MemoryBackendOperationResult, MemoryStageCandidateOperation,
};
use memory::extract::{CandidateKind, ExtractedCandidate, StagingEvidence};
use memory::schema::{EvidenceKind, SourceEvidenceRef, SourceRef};
use memory::schema::{
EvidenceKind, EvidenceOrigin, EvidenceOriginKind, SourceEvidenceRef, SourceRef,
};
use schemars::JsonSchema;
use serde::Deserialize;
@@ -174,17 +176,29 @@ impl Tool for StageMemoryCandidateTool {
"StageMemoryCandidate requires at least one entry_ref".to_string(),
));
}
let mut evidence = Vec::with_capacity(params.entry_refs.len());
let mut source_refs = Vec::with_capacity(params.entry_refs.len());
let mut entries = Vec::with_capacity(params.entry_refs.len());
for entry_ref in &params.entry_refs {
let projection = self.state.view.evidence_for(entry_ref).ok_or_else(|| {
entries.push(self.state.view.evidence_for(entry_ref).ok_or_else(|| {
ToolError::InvalidArgument(format!(
"unknown SessionEntryRef {entry_ref:?} for this extraction capture"
))
})?;
evidence.push(staging_evidence(&projection));
source_refs.push(source_evidence_ref(&projection));
})?);
}
if matches!(params.kind, CandidateKind::Preference)
&& entries.iter().any(|entry| {
!matches!(
entry.origin,
crate::WorkerHistoryProvenance::HumanInput { .. }
)
})
{
return Err(ToolError::InvalidArgument(
"preference candidates require exclusively HumanInput evidence; model, Worker, Flow, backend, derived, and legacy-unknown origins are not preference authority"
.to_string(),
));
}
let evidence = entries.iter().map(staging_evidence).collect();
let source_refs = entries.iter().map(source_evidence_ref).collect();
let candidate = ExtractedCandidate {
kind: params.kind,
claim: params.claim,
@@ -310,11 +324,65 @@ fn evidence_kind(entry: &SessionEntryEvidence) -> EvidenceKind {
}
}
fn evidence_origin(origin: &crate::WorkerHistoryProvenance) -> EvidenceOrigin {
use crate::WorkerHistoryProvenance as Origin;
let mut evidence = EvidenceOrigin {
kind: EvidenceOriginKind::LegacyUnknown,
account_id: None,
workspace_id: None,
runtime_id: None,
worker_id: None,
flow_selector: None,
flow_definition_id: None,
flow_definition_revision: None,
};
match origin {
Origin::HumanInput { account_id } => {
evidence.kind = EvidenceOriginKind::HumanInput;
evidence.account_id = Some(account_id.clone());
}
Origin::WorkerInput { actor } => {
evidence.kind = EvidenceOriginKind::WorkerInput;
evidence.workspace_id = actor.workspace_id.clone();
evidence.runtime_id = actor.runtime_id.clone();
evidence.worker_id = Some(actor.worker_id.clone());
}
Origin::FlowInstruction {
selector,
definition_id,
definition_revision,
..
} => {
evidence.kind = EvidenceOriginKind::FlowInstruction;
evidence.flow_selector = Some(selector.clone());
evidence.flow_definition_id = Some(definition_id.clone());
evidence.flow_definition_revision = Some(*definition_revision);
}
Origin::BackendInstruction { .. } => evidence.kind = EvidenceOriginKind::BackendInstruction,
Origin::ModelOutput { worker } => {
evidence.kind = EvidenceOriginKind::ModelOutput;
evidence.workspace_id = worker.workspace_id.clone();
evidence.runtime_id = worker.runtime_id.clone();
evidence.worker_id = Some(worker.worker_id.clone());
}
Origin::ToolOutput { worker } => {
evidence.kind = EvidenceOriginKind::ToolOutput;
evidence.workspace_id = worker.workspace_id.clone();
evidence.runtime_id = worker.runtime_id.clone();
evidence.worker_id = Some(worker.worker_id.clone());
}
Origin::DerivedSummary => evidence.kind = EvidenceOriginKind::DerivedSummary,
Origin::LegacyUnknown => evidence.kind = EvidenceOriginKind::LegacyUnknown,
}
evidence
}
fn staging_evidence(entry: &SessionEntryEvidence) -> StagingEvidence {
StagingEvidence {
id: entry.entry_ref.to_string(),
kind: evidence_kind(entry),
entry_range: Some(entry.entry_range),
origin: Some(evidence_origin(&entry.origin)),
excerpt: Some(entry.excerpt.clone()),
summary: Some(entry.summary.clone()),
}
@@ -325,6 +393,7 @@ fn source_evidence_ref(entry: &SessionEntryEvidence) -> SourceEvidenceRef {
segment_id: Some(entry.segment_id.clone()),
entry_range: Some(entry.entry_range),
evidence_id: Some(entry.entry_ref.to_string()),
origin: Some(evidence_origin(&entry.origin)),
evidence_kind: Some(evidence_kind(entry)),
label: Some(entry.label.clone()),
summary: Some(entry.summary.clone()),
@@ -432,6 +501,15 @@ mod tests {
assert!(input.contains("StageMemoryCandidate.entry_refs"));
}
#[test]
fn human_origin_projects_account_authority_into_evidence() {
let origin = evidence_origin(&crate::WorkerHistoryProvenance::HumanInput {
account_id: "account-1".into(),
});
assert_eq!(origin.kind, EvidenceOriginKind::HumanInput);
assert_eq!(origin.account_id.as_deref(), Some("account-1"));
}
#[test]
fn backend_input_failures_remain_invalid_argument_tool_errors() {
let backend = map_memory_stage_error(WorkspaceMemoryBackendError::Backend(
@@ -445,6 +523,19 @@ mod tests {
assert!(matches!(http, ToolError::InvalidArgument(_)));
}
#[tokio::test]
async fn preference_rejects_legacy_unknown_before_backend_mutation() {
let tool = StageMemoryCandidateTool { state: state() };
let error = tool
.execute(
r#"{"kind":"preference","claim":"claim","why_useful":"useful","entry_refs":["E00000000"]}"#,
agen::tool::ToolExecutionContext::direct(),
)
.await
.unwrap_err();
assert!(format!("{error:?}").contains("exclusively HumanInput evidence"));
}
#[tokio::test]
async fn stage_rejects_entry_ref_outside_capture_before_backend_mutation() {
let tool = StageMemoryCandidateTool { state: state() };
+210 -25
View File
@@ -14,6 +14,8 @@ use serde_json::json;
use crate::worker::{WorkspaceClient, WorkspaceRequest, WorkspaceRequestMethod};
use super::resource_projection::{project_objective_detail, project_objective_query};
#[derive(Clone, Debug)]
pub struct WorkspaceHttpObjectiveBackend {
client: Arc<dyn WorkspaceClient>,
@@ -37,6 +39,7 @@ impl WorkspaceHttpObjectiveBackend {
)
.await
.map_err(backend_error)?;
let response = project_objective_query(response).map_err(ToolError::ExecutionFailed)?;
Ok(ToolOutput {
summary: "Queried Objectives".to_string(),
content: Some(serde_json::to_string_pretty(&response).map_err(decode_error)?),
@@ -58,8 +61,10 @@ impl WorkspaceHttpObjectiveBackend {
)
.await
.map_err(backend_error)?;
let response = project_objective_detail(response).map_err(ToolError::ExecutionFailed)?;
let objective_ref = response.objective_ref().to_string();
Ok(ToolOutput {
summary: format!("Read objective {id}"),
summary: format!("Read objective {objective_ref}"),
content: Some(serde_json::to_string_pretty(&response).map_err(decode_error)?),
attachments: Vec::new(),
})
@@ -84,7 +89,7 @@ impl WorkspaceHttpObjectiveBackend {
.await
.map_err(backend_error)?;
Ok(objective_output(
format!("Created objective {}", response.id),
format!("Created objective {}", &response.resource_key),
response,
)?)
}
@@ -112,7 +117,7 @@ impl WorkspaceHttpObjectiveBackend {
.await
.map_err(backend_error)?;
Ok(objective_output(
format!("Edited objective {}", response.id),
format!("Edited objective {}", &response.resource_key),
response,
)?)
}
@@ -134,7 +139,7 @@ impl WorkspaceHttpObjectiveBackend {
.await
.map_err(backend_error)?;
Ok(objective_output(
format!("Updated objective {} state", response.id),
format!("Updated objective {} state", &response.resource_key),
response,
)?)
}
@@ -142,6 +147,7 @@ impl WorkspaceHttpObjectiveBackend {
async fn link_ticket(&self, input: ObjectiveLinkTicketInput) -> Result<ToolOutput, ToolError> {
let id = validate_id(&input.id, "ObjectiveLinkTicket")?;
let ticket_id = validate_id(&input.ticket_id, "ObjectiveLinkTicket")?;
let ticket_resource_key = self.ticket_resource_key(ticket_id).await?;
let url = format!("{}/ticket-links", self.objective_url(id));
let response = send_json::<ObjectiveLinkTicketRequest, ObjectiveDetail>(
self.client.as_ref(),
@@ -154,7 +160,10 @@ impl WorkspaceHttpObjectiveBackend {
.await
.map_err(backend_error)?;
Ok(objective_output(
format!("Linked ticket {ticket_id} to objective {}", response.id),
format!(
"Linked ticket {ticket_resource_key} to objective {}",
&response.resource_key
),
response,
)?)
}
@@ -165,16 +174,46 @@ impl WorkspaceHttpObjectiveBackend {
) -> Result<ToolOutput, ToolError> {
let id = validate_id(&input.id, "ObjectiveUnlinkTicket")?;
let ticket_id = validate_id(&input.ticket_id, "ObjectiveUnlinkTicket")?;
let ticket_resource_key = self.ticket_resource_key(ticket_id).await?;
let url = format!("{}/ticket-links/{}", self.objective_url(id), ticket_id);
let response = delete_json::<ObjectiveDetail>(self.client.as_ref(), &url)
.await
.map_err(backend_error)?;
Ok(objective_output(
format!("Unlinked ticket {ticket_id} from objective {}", response.id),
format!(
"Unlinked ticket {ticket_resource_key} from objective {}",
&response.resource_key
),
response,
)?)
}
async fn ticket_resource_key(&self, ticket_reference: &str) -> Result<String, ToolError> {
let workspace_id = self.client.workspace_id().unwrap_or_default();
let response: serde_json::Value = decode_response(
self.client
.execute(WorkspaceRequest::get(format!(
"/api/w/{workspace_id}/tickets/{ticket_reference}"
)))
.map_err(WorkspaceObjectiveBackendError::from)
.map_err(backend_error)?,
)
.map_err(backend_error)?;
response
.get("resource_key")
.or_else(|| {
response
.get("meta")
.and_then(|meta| meta.get("resource_key"))
})
.and_then(serde_json::Value::as_str)
.filter(|key| is_canonical_resource_key(key, "T-"))
.map(ToOwned::to_owned)
.ok_or_else(|| {
ToolError::ExecutionFailed("required T- human key is unavailable".to_string())
})
}
fn objective_url(&self, id: &str) -> String {
let workspace_id = self.client.workspace_id().unwrap_or_default();
format!("/api/w/{workspace_id}/objectives/{id}")
@@ -185,7 +224,7 @@ impl WorkspaceHttpObjectiveBackend {
pub enum WorkspaceObjectiveBackendError {
#[error("workspace objective backend request failed: {0}")]
Request(#[from] crate::worker::WorkspaceClientError),
#[error("workspace objective backend returned HTTP {status}: {body}")]
#[error("workspace objective backend returned HTTP {status}")]
Http {
status: reqwest::StatusCode,
body: String,
@@ -247,10 +286,26 @@ fn decode_response<T: for<'de> Deserialize<'de>>(
serde_json::from_str(&response.body).map_err(Into::into)
}
fn is_canonical_resource_key(resource_key: &str, prefix: &str) -> bool {
resource_key.strip_prefix(prefix).is_some_and(|sequence| {
!sequence.is_empty() && sequence.bytes().all(|byte| byte.is_ascii_digit())
})
}
fn objective_output(summary: String, response: ObjectiveDetail) -> Result<ToolOutput, ToolError> {
if !is_canonical_resource_key(&response.resource_key, "O-") {
return Err(ToolError::ExecutionFailed(
"required O- human key is unavailable".to_string(),
));
}
let projected = serde_json::json!({
"objective": &response.resource_key,
"title": response.title,
"state": response.state,
});
Ok(ToolOutput {
summary,
content: Some(serde_json::to_string_pretty(&response).map_err(decode_error)?),
content: Some(serde_json::to_string_pretty(&projected).map_err(decode_error)?),
attachments: Vec::new(),
})
@@ -260,7 +315,7 @@ fn validate_id<'a>(id: &'a str, tool_name: &str) -> Result<&'a str, ToolError> {
let id = id.trim();
if id.is_empty() || id.contains('/') {
return Err(ToolError::InvalidArgument(format!(
"{tool_name} requires non-empty canonical id without '/'"
"{tool_name} requires a non-empty Objective reference without '/'"
)));
}
Ok(id)
@@ -411,9 +466,9 @@ const EDIT_DESCRIPTION: &str =
const SET_STATE_DESCRIPTION: &str =
"Set an Objective state through Backend Workspace API authority.";
const LINK_TICKET_DESCRIPTION: &str =
"Link a Ticket id to an Objective through Backend Workspace API authority.";
"Link a Ticket reference to an Objective through Backend Workspace API authority.";
const UNLINK_TICKET_DESCRIPTION: &str =
"Unlink a Ticket id from an Objective through Backend Workspace API authority.";
"Unlink a Ticket reference from an Objective through Backend Workspace API authority.";
fn list_schema() -> serde_json::Value {
json!({
@@ -422,7 +477,7 @@ fn list_schema() -> serde_json::Value {
"properties":{
"query":{"type":["string","null"]},
"states":{"type":"array","items":{"type":"string"},"default":[]},
"linked_ticket_id":{"type":["string","null"]},
"linked_ticket_id":{"type":["string","null"],"description":"Linked Ticket reference. Prefer T-*; canonical internal ids remain accepted for compatibility."},
"updated_after":{"type":["string","null"]},
"updated_before":{"type":["string","null"]},
"sort":{"type":["string","null"],"enum":["relevance","updated_desc","created_desc","title",null]},
@@ -438,7 +493,7 @@ fn show_schema() -> serde_json::Value {
"additionalProperties": false,
"required":["id"],
"properties":{
"id":{"type":"string"},
"id":{"type":"string","description":"Objective reference. Prefer O-*; canonical internal ids remain accepted for compatibility."},
"event_limit":{"type":["integer","null"],"minimum":1,"maximum":50},
"event_cursor":{"type":["string","null"]}
}
@@ -454,7 +509,7 @@ fn create_schema() -> serde_json::Value {
"title":{"type":"string","minLength":1},
"body_md":{"type":"string"},
"state":{"type":"string","default":"active"},
"linked_tickets":{"type":"array","items":{"type":"string"}}
"linked_tickets":{"type":"array","items":{"type":"string"},"description":"Linked Ticket references. Prefer T-*; canonical internal ids remain accepted for compatibility."}
}
})
}
@@ -465,7 +520,7 @@ fn edit_schema() -> serde_json::Value {
"additionalProperties": false,
"required":["id"],
"properties":{
"id":{"type":"string"},
"id":{"type":"string","description":"Objective reference. Prefer O-*; canonical internal ids remain accepted for compatibility."},
"title":{"type":["string","null"]},
"old_string":{"type":["string","null"]},
"new_string":{"type":["string","null"]},
@@ -480,7 +535,7 @@ fn set_state_schema() -> serde_json::Value {
"additionalProperties": false,
"required":["id","state"],
"properties":{
"id":{"type":"string"},
"id":{"type":"string","description":"Objective reference. Prefer O-*; canonical internal ids remain accepted for compatibility."},
"state":{"type":"string","minLength":1}
}
})
@@ -500,8 +555,8 @@ fn id_ticket_schema(required: &[&str]) -> serde_json::Value {
"additionalProperties": false,
"required": required,
"properties":{
"id":{"type":"string"},
"ticket_id":{"type":"string"}
"id":{"type":"string","description":"Objective reference. Prefer O-*; canonical internal ids remain accepted for compatibility."},
"ticket_id":{"type":"string","description":"Ticket reference. Prefer T-*; canonical internal ids remain accepted for compatibility."}
}
})
}
@@ -595,21 +650,20 @@ fn default_state() -> String {
#[derive(Debug, Serialize, Deserialize, PartialEq, Eq)]
struct ObjectiveDetail {
id: String,
resource_key: String,
title: String,
state: String,
created_at: Option<String>,
updated_at: Option<String>,
linked_tickets: Vec<String>,
body: String,
body_truncated: bool,
record_source: String,
}
#[cfg(test)]
mod tests {
use super::*;
use agen::tool::ToolDefinition;
use std::{
io::{Read, Write},
net::TcpListener,
thread,
};
fn tool_names(definitions: Vec<ToolDefinition>) -> Vec<String> {
let mut names = definitions
@@ -656,4 +710,135 @@ mod tests {
let link = link_ticket_schema();
assert_eq!(link["required"], json!(["id", "ticket_id"]));
}
#[tokio::test(flavor = "multi_thread")]
async fn objective_show_summary_uses_projected_human_key() {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let base_url = format!("http://{}", listener.local_addr().unwrap());
let server = thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
let mut buffer = [0_u8; 8192];
let len = stream.read(&mut buffer).unwrap();
let request = String::from_utf8_lossy(&buffer[..len]);
assert!(
request.starts_with("POST /api/w/workspace/objectives/00001INTERNAL/show HTTP/1.1")
);
let body = serde_json::json!({
"id": "00001INTERNAL",
"resource_key": "O-3",
"title": "Objective",
"body": "Body",
"state": "active",
"created_at": null,
"updated_at": null,
"linked_ticket_summaries": [],
"events": [],
"event_page": {"next_cursor": null, "has_more": false}
})
.to_string();
write!(
stream,
"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\n\r\n{}",
body.len(),
body
)
.unwrap();
});
let backend = WorkspaceHttpObjectiveBackend::new(Arc::new(
crate::worker::TestWorkspaceHttpClient::new("workspace", base_url),
));
let output = backend
.show(ShowObjectiveInput {
id: "00001INTERNAL".to_string(),
event_limit: None,
event_cursor: None,
})
.await
.unwrap();
server.join().unwrap();
assert_eq!(output.summary, "Read objective O-3");
assert!(!output.summary.contains("00001INTERNAL"));
assert!(!output.content.unwrap().contains("00001INTERNAL"));
}
#[tokio::test(flavor = "multi_thread")]
async fn objective_link_summaries_resolve_internal_ticket_ids_to_human_keys() {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let base_url = format!("http://{}", listener.local_addr().unwrap());
let server = thread::spawn(move || {
for mutation in ["POST", "DELETE"] {
let (mut stream, _) = listener.accept().unwrap();
let mut buffer = [0_u8; 8192];
let len = stream.read(&mut buffer).unwrap();
let request = String::from_utf8_lossy(&buffer[..len]);
assert!(request.starts_with("GET /api/w/workspace/tickets/00001INTERNAL HTTP/1.1"));
let response_body = serde_json::json!({"resource_key": "T-7"}).to_string();
write!(
stream,
"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\n\r\n{}",
response_body.len(),
response_body
)
.unwrap();
let (mut stream, _) = listener.accept().unwrap();
let mut buffer = [0_u8; 8192];
let len = stream.read(&mut buffer).unwrap();
let request = String::from_utf8_lossy(&buffer[..len]);
assert!(request.starts_with(&format!(
"{mutation} /api/w/workspace/objectives/O-3/ticket-links"
)));
let response_body = serde_json::json!({
"resource_key": "O-3",
"title": "Objective",
"state": "active"
})
.to_string();
write!(
stream,
"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\n\r\n{}",
response_body.len(),
response_body
)
.unwrap();
}
});
let backend = WorkspaceHttpObjectiveBackend::new(Arc::new(
crate::worker::TestWorkspaceHttpClient::new("workspace", base_url),
));
let linked = backend
.link_ticket(ObjectiveLinkTicketInput {
id: "O-3".to_string(),
ticket_id: "00001INTERNAL".to_string(),
})
.await
.unwrap();
let unlinked = backend
.unlink_ticket(ObjectiveUnlinkTicketInput {
id: "O-3".to_string(),
ticket_id: "00001INTERNAL".to_string(),
})
.await
.unwrap();
server.join().unwrap();
for output in [linked, unlinked] {
assert!(output.summary.contains("T-7"));
assert!(!output.summary.contains("00001INTERNAL"));
assert!(!output.content.unwrap().contains("00001INTERNAL"));
}
}
#[test]
fn objective_output_rejects_noncanonical_human_keys() {
let response = ObjectiveDetail {
resource_key: "O-internal".to_string(),
title: "Objective".to_string(),
state: "active".to_string(),
};
assert!(objective_output("created".to_string(), response).is_err());
}
}
@@ -0,0 +1,794 @@
use serde::Serialize;
use serde_json::{Map, Value};
#[derive(Debug, Serialize)]
pub(super) struct ModelTicketQueryResponse {
tickets: Vec<ModelTicketQueryItem>,
next_cursor: Option<String>,
has_more: bool,
}
#[derive(Debug, Serialize)]
struct ModelTicketQueryItem {
ticket: String,
title: String,
state: String,
readiness: Option<String>,
priority: Option<String>,
created_at: Option<String>,
updated_at: Option<String>,
workspace_action_priority: Option<String>,
matched_fields: Vec<String>,
snippet: Option<String>,
current_coder: Option<ModelWorkerSummary>,
linked_objectives: Vec<String>,
relation_count: usize,
blocker_count: usize,
unresolved_blocker_count: usize,
unresolved_review_count: usize,
evidence: Option<ModelTicketEvidence>,
merge_request: Option<ModelMergeRequest>,
}
#[derive(Debug, Serialize)]
pub(super) struct ModelTicketDetail {
ticket: String,
title: String,
body: String,
state: String,
readiness: Option<String>,
priority: Option<String>,
created_at: Option<String>,
updated_at: Option<String>,
thread: Vec<ModelTicketEvent>,
relations: ModelTicketRelations,
linked_objectives: Vec<ModelObjectiveSummary>,
assignments: Vec<ModelAssignment>,
current_coder: Option<ModelWorkerSummary>,
implementation_reports: Vec<ModelEvidenceEvent>,
merge_request: Option<ModelMergeRequest>,
evidence: Option<ModelTicketEvidence>,
actions: Option<ModelTicketActions>,
event_page: Option<ModelEventPage>,
}
#[derive(Debug, Serialize)]
pub(super) struct ModelObjectiveQueryResponse {
objectives: Vec<ModelObjectiveQueryItem>,
next_cursor: Option<String>,
has_more: bool,
}
#[derive(Debug, Serialize)]
struct ModelObjectiveQueryItem {
objective: String,
title: String,
summary: Option<String>,
state: String,
created_at: Option<String>,
updated_at: Option<String>,
linked_tickets: Vec<String>,
linked_ticket_count: usize,
}
#[derive(Debug, Serialize)]
pub(super) struct ModelObjectiveDetail {
objective: String,
title: String,
body: String,
state: String,
created_at: Option<String>,
updated_at: Option<String>,
linked_tickets: Vec<ModelTicketSummary>,
events: Vec<ModelObjectiveEvent>,
event_page: ModelObjectiveEventPage,
}
impl ModelObjectiveDetail {
pub(super) fn objective_ref(&self) -> &str {
&self.objective
}
}
#[derive(Debug, Serialize)]
struct ModelWorkerSummary {
worker: String,
}
#[derive(Debug, Serialize)]
struct ModelTicketEvent {
sequence: usize,
kind: String,
body: Option<String>,
created_at: Option<String>,
}
#[derive(Debug, Serialize, Default)]
struct ModelTicketRelations {
outgoing: Vec<ModelRelation>,
incoming: Vec<ModelRelation>,
blockers: Vec<ModelBlocker>,
notices: Vec<ModelNotice>,
}
#[derive(Debug, Serialize)]
struct ModelRelation {
ticket: String,
kind: String,
note: Option<String>,
created_at: Option<String>,
}
#[derive(Debug, Serialize)]
struct ModelBlocker {
ticket: String,
kind: String,
state: Option<String>,
resolved: bool,
}
#[derive(Debug, Serialize)]
struct ModelNotice {
kind: String,
}
#[derive(Debug, Serialize)]
struct ModelObjectiveSummary {
objective: String,
title: String,
state: String,
}
#[derive(Debug, Serialize)]
struct ModelTicketSummary {
ticket: String,
title: String,
state: String,
}
#[derive(Debug, Serialize)]
struct ModelAssignment {
role: String,
principal: String,
assigned_at: String,
}
#[derive(Debug, Serialize)]
struct ModelEvidenceEvent {
sequence: usize,
kind: String,
created_at: Option<String>,
excerpt: String,
}
#[derive(Debug, Serialize)]
struct ModelMergeRequest {
state: String,
selector_from: Option<String>,
selector_to: String,
review_status: String,
subject_ref: Option<String>,
review_excerpt: Option<String>,
}
#[derive(Debug, Serialize)]
struct ModelTicketEvidence {
has_merge_request: bool,
has_current_subject_ref: bool,
has_review_request: bool,
has_commit: bool,
review_status: Option<String>,
approved_current_subject: bool,
unresolved_request_changes: bool,
complete_for_integration: bool,
missing: Vec<String>,
}
#[derive(Debug, Serialize)]
struct ModelTicketActions {
can_assign_orchestrator: bool,
can_unassign_orchestrator: bool,
can_queue: bool,
can_start_manual_coder: bool,
}
#[derive(Debug, Serialize)]
struct ModelEventPage {
next_cursor: Option<String>,
has_more: bool,
}
#[derive(Debug, Serialize)]
struct ModelObjectiveEvent {
kind: String,
created_at: String,
body: Option<String>,
}
#[derive(Debug, Serialize)]
struct ModelObjectiveEventPage {
next_cursor: Option<String>,
has_more: bool,
}
pub(super) fn project_ticket_query(value: Value) -> Result<ModelTicketQueryResponse, String> {
let root = object(&value, "Ticket query response")?;
let page = object_field(root, "page")?;
let tickets = array_field(root, "items")?
.iter()
.map(project_ticket_query_item)
.collect::<Result<Vec<_>, _>>()?;
Ok(ModelTicketQueryResponse {
tickets,
next_cursor: optional_string(page, "next_cursor")?,
has_more: bool_field(page, "has_more")?,
})
}
fn project_ticket_query_item(value: &Value) -> Result<ModelTicketQueryItem, String> {
let item = object(value, "Ticket query item")?;
Ok(ModelTicketQueryItem {
ticket: human_ref(item, "resource_key", "T-")?,
title: string_field(item, "title")?,
state: string_field(item, "state")?,
readiness: optional_string(item, "readiness")?,
priority: optional_string(item, "priority")?,
created_at: optional_string(item, "created_at")?,
updated_at: optional_string(item, "updated_at")?,
workspace_action_priority: optional_string(item, "workspace_action_priority")?,
matched_fields: string_array(item, "matched_fields")?,
snippet: optional_string(item, "snippet")?,
current_coder: item
.get("current_coder")
.filter(|value| !value.is_null())
.map(project_worker)
.transpose()?,
linked_objectives: string_array(item, "linked_objective_keys")?
.into_iter()
.map(|key| validate_human_ref(key, "O-"))
.collect::<Result<Vec<_>, _>>()?,
relation_count: usize_field(item, "relation_count")?,
blocker_count: usize_field(item, "blocker_count")?,
unresolved_blocker_count: usize_field(item, "unresolved_blocker_count")?,
unresolved_review_count: usize_field(item, "unresolved_review_count")?,
evidence: item.get("evidence").map(project_evidence).transpose()?,
merge_request: item
.get("merge_request")
.filter(|value| !value.is_null())
.map(project_merge_request)
.transpose()?,
})
}
pub(super) fn project_ticket_detail(value: Value) -> Result<ModelTicketDetail, String> {
let root = object(&value, "Ticket detail response")?;
let current_coder = root
.get("current_coder")
.filter(|value| !value.is_null())
.map(project_worker)
.transpose()?;
let assignments = array_field(root, "assignments")?
.iter()
.map(|assignment| project_assignment(assignment, current_coder.as_ref()))
.collect::<Result<Vec<_>, _>>()?;
Ok(ModelTicketDetail {
ticket: human_ref(root, "resource_key", "T-")?,
title: string_field(root, "title")?,
body: string_field(root, "body")?,
state: string_field(root, "state")?,
readiness: optional_string(root, "readiness")?,
priority: optional_string(root, "priority")?,
created_at: optional_string(root, "created_at")?,
updated_at: optional_string(root, "updated_at")?,
thread: array_field(root, "events")?
.iter()
.map(project_ticket_event)
.collect::<Result<Vec<_>, _>>()?,
relations: project_relations(root.get("relations"))?,
linked_objectives: array_field(root, "linked_objectives")?
.iter()
.map(project_objective_summary)
.collect::<Result<Vec<_>, _>>()?,
assignments,
current_coder,
implementation_reports: array_field(root, "implementation_reports")?
.iter()
.map(project_evidence_event)
.collect::<Result<Vec<_>, _>>()?,
merge_request: root
.get("merge_request")
.filter(|value| !value.is_null())
.map(project_merge_request)
.transpose()?,
evidence: root.get("evidence").map(project_evidence).transpose()?,
actions: root
.get("action_eligibility")
.filter(|value| !value.is_null())
.map(project_actions)
.transpose()?,
event_page: root
.get("event_page")
.filter(|value| !value.is_null())
.map(project_event_page)
.transpose()?,
})
}
pub(super) fn project_objective_query(value: Value) -> Result<ModelObjectiveQueryResponse, String> {
let root = object(&value, "Objective query response")?;
let page = object_field(root, "page")?;
Ok(ModelObjectiveQueryResponse {
objectives: array_field(root, "items")?
.iter()
.map(project_objective_query_item)
.collect::<Result<Vec<_>, _>>()?,
next_cursor: optional_string(page, "next_cursor")?,
has_more: bool_field(page, "has_more")?,
})
}
fn project_objective_query_item(value: &Value) -> Result<ModelObjectiveQueryItem, String> {
let item = object(value, "Objective query item")?;
let linked_tickets = string_array(item, "linked_ticket_keys")?
.into_iter()
.map(|key| validate_human_ref(key, "T-"))
.collect::<Result<Vec<_>, _>>()?;
Ok(ModelObjectiveQueryItem {
objective: human_ref(item, "resource_key", "O-")?,
title: string_field(item, "title")?,
summary: optional_string(item, "snippet")?,
state: string_field(item, "state")?,
created_at: optional_string(item, "created_at")?,
updated_at: optional_string(item, "updated_at")?,
linked_ticket_count: linked_tickets.len(),
linked_tickets,
})
}
pub(super) fn project_objective_detail(value: Value) -> Result<ModelObjectiveDetail, String> {
let root = object(&value, "Objective detail response")?;
Ok(ModelObjectiveDetail {
objective: human_ref(root, "resource_key", "O-")?,
title: string_field(root, "title")?,
body: string_field(root, "body")?,
state: string_field(root, "state")?,
created_at: optional_string(root, "created_at")?,
updated_at: optional_string(root, "updated_at")?,
linked_tickets: array_field(root, "linked_ticket_summaries")?
.iter()
.map(project_ticket_summary)
.collect::<Result<Vec<_>, _>>()?,
events: array_field(root, "events")?
.iter()
.map(project_objective_event)
.collect::<Result<Vec<_>, _>>()?,
event_page: project_objective_event_page(
root.get("event_page")
.ok_or_else(|| "Objective detail response is missing event_page".to_string())?,
)?,
})
}
fn project_worker(value: &Value) -> Result<ModelWorkerSummary, String> {
let worker = object(value, "Worker summary")?;
Ok(ModelWorkerSummary {
worker: human_ref(worker, "worker_resource_key", "W-")?,
})
}
fn project_ticket_event(value: &Value) -> Result<ModelTicketEvent, String> {
let event = object(value, "Ticket event")?;
Ok(ModelTicketEvent {
sequence: usize_field(event, "sequence")?,
kind: string_field(event, "kind")?,
body: match event.get("body") {
None | Some(Value::Null) => None,
Some(Value::String(body)) => Some(body.clone()),
Some(_) => return Err("invalid Ticket event body".to_string()),
},
created_at: optional_string(event, "at")?,
})
}
fn project_relations(value: Option<&Value>) -> Result<ModelTicketRelations, String> {
let Some(value) = value else {
return Ok(ModelTicketRelations::default());
};
let relations = object(value, "Ticket relations")?;
Ok(ModelTicketRelations {
outgoing: array_field(relations, "outgoing")?
.iter()
.map(|value| project_relation(value, "target_resource_key", "kind"))
.collect::<Result<Vec<_>, _>>()?,
incoming: array_field(relations, "incoming")?
.iter()
.map(|value| project_relation(value, "source_resource_key", "forward_kind"))
.collect::<Result<Vec<_>, _>>()?,
blockers: array_field(relations, "blockers")?
.iter()
.map(project_blocker)
.collect::<Result<Vec<_>, _>>()?,
notices: array_field(relations, "notices")?
.iter()
.map(project_notice)
.collect::<Result<Vec<_>, _>>()?,
})
}
fn project_relation(
value: &Value,
ticket_key: &str,
kind_key: &str,
) -> Result<ModelRelation, String> {
let relation = object(value, "Ticket relation")?;
let relation_data = relation.get("relation").and_then(Value::as_object);
let kind = if kind_key == "kind" {
relation_data
.ok_or_else(|| "Ticket relation is missing relation data".to_string())
.and_then(|data| string_field(data, "kind"))?
} else {
string_field(relation, kind_key)?
};
let note = match relation_data {
Some(data) => optional_string(data, "note")?,
None => optional_string(relation, "note")?,
};
let created_at = match relation_data {
Some(data) => optional_string(data, "at")?,
None => optional_string(relation, "at")?,
};
Ok(ModelRelation {
ticket: human_ref(relation, ticket_key, "T-")?,
kind,
note,
created_at,
})
}
fn project_blocker(value: &Value) -> Result<ModelBlocker, String> {
let blocker = object(value, "Ticket blocker")?;
Ok(ModelBlocker {
ticket: human_ref(blocker, "blocking_resource_key", "T-")?,
kind: string_field(blocker, "relation_kind")?,
state: optional_string(blocker, "blocking_state")?,
resolved: bool_field(blocker, "resolved")?,
})
}
fn project_notice(value: &Value) -> Result<ModelNotice, String> {
let notice = object(value, "Ticket notice")?;
Ok(ModelNotice {
kind: string_field(notice, "kind")?,
})
}
fn project_objective_summary(value: &Value) -> Result<ModelObjectiveSummary, String> {
let summary = object(value, "Objective summary")?;
Ok(ModelObjectiveSummary {
objective: human_ref(summary, "resource_key", "O-")?,
title: string_field(summary, "title")?,
state: string_field(summary, "state")?,
})
}
fn project_ticket_summary(value: &Value) -> Result<ModelTicketSummary, String> {
let summary = object(value, "Ticket summary")?;
Ok(ModelTicketSummary {
ticket: human_ref(summary, "resource_key", "T-")?,
title: string_field(summary, "title")?,
state: string_field(summary, "state")?,
})
}
fn project_assignment(
value: &Value,
current_coder: Option<&ModelWorkerSummary>,
) -> Result<ModelAssignment, String> {
let assignment = object(value, "Ticket assignment")?;
let principal = object_field(assignment, "principal")?;
let kind = string_field(principal, "kind")?;
let principal = match kind.as_str() {
"worker" => current_coder
.map(|coder| coder.worker.clone())
.ok_or_else(|| {
"Worker assignment is missing a Workspace human key projection".to_string()
})?,
"workspace_agent" => format!("workspace-agent:{}", string_field(principal, "agent_key")?),
"user" => "user".to_string(),
other => format!("source:{other}"),
};
Ok(ModelAssignment {
role: string_field(assignment, "role")?,
principal,
assigned_at: string_field(assignment, "assigned_at")?,
})
}
fn project_evidence_event(value: &Value) -> Result<ModelEvidenceEvent, String> {
let event = object(value, "Ticket evidence event")?;
Ok(ModelEvidenceEvent {
sequence: usize_field(event, "sequence")?,
kind: string_field(event, "kind")?,
created_at: optional_string(event, "at")?,
excerpt: string_field(event, "excerpt")?,
})
}
fn project_merge_request(value: &Value) -> Result<ModelMergeRequest, String> {
let merge = object(value, "Merge Request summary")?;
Ok(ModelMergeRequest {
state: string_field(merge, "state")?,
selector_from: optional_string(merge, "selector_from")?,
selector_to: string_field(merge, "selector_to")?,
review_status: string_field(merge, "review_status")?,
subject_ref: optional_string(merge, "subject_ref")?,
review_excerpt: optional_string(merge, "review_excerpt")?,
})
}
fn project_evidence(value: &Value) -> Result<ModelTicketEvidence, String> {
let evidence = object(value, "Ticket evidence")?;
Ok(ModelTicketEvidence {
has_merge_request: bool_field(evidence, "has_merge_request")?,
has_current_subject_ref: bool_field(evidence, "has_current_subject_ref")?,
has_review_request: bool_field(evidence, "has_review_request")?,
has_commit: bool_field(evidence, "has_commit")?,
review_status: optional_string(evidence, "review_status")?,
approved_current_subject: bool_field(evidence, "approved_current_subject")?,
unresolved_request_changes: bool_field(evidence, "unresolved_request_changes")?,
complete_for_integration: bool_field(evidence, "complete_for_integration")?,
missing: string_array(evidence, "missing")?,
})
}
fn project_actions(value: &Value) -> Result<ModelTicketActions, String> {
let actions = object(value, "Ticket actions")?;
Ok(ModelTicketActions {
can_assign_orchestrator: bool_field(actions, "can_assign_orchestrator")?,
can_unassign_orchestrator: bool_field(actions, "can_unassign_orchestrator")?,
can_queue: bool_field(actions, "can_queue")?,
can_start_manual_coder: bool_field(actions, "can_start_manual_coder")?,
})
}
fn project_event_page(value: &Value) -> Result<ModelEventPage, String> {
let page = object(value, "Ticket event page")?;
Ok(ModelEventPage {
next_cursor: optional_string(page, "next_cursor")?,
has_more: bool_field(page, "has_more")?,
})
}
fn project_objective_event(value: &Value) -> Result<ModelObjectiveEvent, String> {
let event = object(value, "Objective event")?;
let body = optional_string(event, "body")?;
Ok(ModelObjectiveEvent {
kind: string_field(event, "kind")?,
created_at: string_field(event, "created_at")?,
body,
})
}
fn project_objective_event_page(value: &Value) -> Result<ModelObjectiveEventPage, String> {
let page = object(value, "Objective event page")?;
Ok(ModelObjectiveEventPage {
next_cursor: optional_string(page, "next_cursor")?,
has_more: bool_field(page, "has_more")?,
})
}
fn object<'a>(value: &'a Value, context: &str) -> Result<&'a Map<String, Value>, String> {
value
.as_object()
.ok_or_else(|| format!("{context} must be an object"))
}
fn object_field<'a>(
object: &'a Map<String, Value>,
key: &str,
) -> Result<&'a Map<String, Value>, String> {
object
.get(key)
.and_then(Value::as_object)
.ok_or_else(|| format!("missing or invalid {key}"))
}
fn array_field<'a>(object: &'a Map<String, Value>, key: &str) -> Result<&'a [Value], String> {
object
.get(key)
.and_then(Value::as_array)
.map(Vec::as_slice)
.ok_or_else(|| format!("missing or invalid {key}"))
}
fn string_field(object: &Map<String, Value>, key: &str) -> Result<String, String> {
object
.get(key)
.and_then(Value::as_str)
.map(ToOwned::to_owned)
.ok_or_else(|| format!("missing or invalid {key}"))
}
fn optional_string(object: &Map<String, Value>, key: &str) -> Result<Option<String>, String> {
match object.get(key) {
None | Some(Value::Null) => Ok(None),
Some(Value::String(value)) => Ok(Some(value.clone())),
Some(_) => Err(format!("invalid {key}")),
}
}
fn bool_field(object: &Map<String, Value>, key: &str) -> Result<bool, String> {
object
.get(key)
.and_then(Value::as_bool)
.ok_or_else(|| format!("missing or invalid {key}"))
}
fn usize_field(object: &Map<String, Value>, key: &str) -> Result<usize, String> {
object
.get(key)
.and_then(Value::as_u64)
.and_then(|value| usize::try_from(value).ok())
.ok_or_else(|| format!("missing or invalid {key}"))
}
fn string_array(object: &Map<String, Value>, key: &str) -> Result<Vec<String>, String> {
array_field(object, key)?
.iter()
.map(|value| {
value
.as_str()
.map(ToOwned::to_owned)
.ok_or_else(|| format!("invalid {key}"))
})
.collect()
}
fn human_ref(object: &Map<String, Value>, key: &str, prefix: &str) -> Result<String, String> {
let value = object
.get(key)
.and_then(Value::as_str)
.map(ToOwned::to_owned)
.ok_or_else(|| format!("required {prefix} human key is unavailable"))?;
validate_human_ref(value, prefix)
}
fn validate_human_ref(value: String, prefix: &str) -> Result<String, String> {
let valid = value.strip_prefix(prefix).is_some_and(|sequence| {
!sequence.is_empty() && sequence.bytes().all(|byte| byte.is_ascii_digit())
});
if valid {
Ok(value)
} else {
Err(format!("required {prefix} human key is unavailable"))
}
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
#[test]
fn objective_projection_exposes_only_human_resource_references() {
let projected = project_objective_detail(json!({
"id": "00001M10HW6BV",
"resource_key": "O-543",
"title": "Objective",
"body": "Body",
"state": "active",
"created_at": "2026-01-01T00:00:00Z",
"updated_at": "2026-01-02T00:00:00Z",
"linked_tickets": ["00001M0E82D1V"],
"linked_ticket_summaries": [{
"id": "00001M0E82D1V",
"resource_key": "T-496",
"title": "Ticket",
"state": "done",
"updated_at": "2026-01-02T00:00:00Z"
}],
"events": [{
"sequence": 3,
"event_ref": "objective-event-3",
"kind": "linked_ticket",
"created_at": "2026-01-02T00:00:00Z",
"body": "linked"
}],
"event_page": {"next_cursor": null, "has_more": false, "window_start_sequence": 3, "window_end_sequence": 3}
})).expect("projection");
let json = serde_json::to_value(projected).expect("serialize");
let text = json.to_string();
assert!(text.contains("O-543"));
assert!(text.contains("T-496"));
assert!(!text.contains("00001M10HW6BV"));
assert!(!text.contains("00001M0E82D1V"));
assert!(!text.contains("event_ref"));
}
#[test]
fn query_projections_accept_workspace_api_shapes_and_scrub_internal_ids() {
let ticket = project_ticket_query(json!({
"page": {"next_cursor": null, "has_more": false},
"record_authority": "workspace_sqlite",
"items": [{
"id": "00001TICKETINTERNAL",
"resource_key": "T-543",
"title": "Ticket",
"state": "inprogress",
"readiness": null,
"priority": "high",
"created_at": null,
"updated_at": "2026-01-01T00:00:00Z",
"workspace_action_priority": "active_work",
"matched_fields": ["title"],
"snippet": "Ticket",
"current_coder": {"runtime_id": "runtime-internal", "worker_id": "worker-internal", "worker_resource_key": "W-12"},
"linked_objective_ids": ["00001OBJECTIVEINTERNAL"],
"linked_objective_keys": ["O-6"],
"relation_count": 0,
"blocker_count": 0,
"unresolved_blocker_count": 0,
"unresolved_review_count": 0,
"evidence": {
"has_merge_request": false,
"has_current_subject_ref": false,
"has_review_request": false,
"has_commit": false,
"review_status": null,
"approved_current_subject": false,
"unresolved_request_changes": false,
"complete_for_integration": false,
"missing": ["merge_request"]
},
"merge_request": null
}]
})).expect("Ticket query projection");
let ticket_json = serde_json::to_string(&ticket).expect("serialize Ticket query");
assert!(ticket_json.contains("T-543"));
assert!(ticket_json.contains("O-6"));
assert!(ticket_json.contains("W-12"));
assert!(!ticket_json.contains("00001TICKETINTERNAL"));
assert!(!ticket_json.contains("runtime-internal"));
assert!(!ticket_json.contains("worker-internal"));
let objective = project_objective_query(json!({
"page": {"next_cursor": null, "has_more": false},
"record_authority": "workspace_sqlite",
"items": [{
"id": "00001OBJECTIVEINTERNAL",
"resource_key": "O-6",
"title": "Objective",
"state": "active",
"created_at": null,
"updated_at": null,
"matched_fields": [],
"snippet": null,
"linked_ticket_count": 1,
"linked_tickets": ["00001TICKETINTERNAL"],
"linked_ticket_keys": ["T-543"]
}]
}))
.expect("Objective query projection");
let objective_json = serde_json::to_string(&objective).expect("serialize Objective query");
assert!(objective_json.contains("O-6"));
assert!(objective_json.contains("T-543"));
assert!(objective_json.contains("\"summary\":null"));
assert!(!objective_json.contains("00001OBJECTIVEINTERNAL"));
assert!(!objective_json.contains("00001TICKETINTERNAL"));
}
#[test]
fn human_resource_projection_rejects_noncanonical_keys() {
for (key, prefix) in [("T-key", "T-"), ("O-", "O-"), ("W-1x", "W-")] {
assert!(validate_human_ref(key.to_string(), prefix).is_err());
}
}
#[test]
fn ticket_projection_fails_closed_without_worker_resource_key() {
let error = project_worker(&json!({"worker_resource_key": null}))
.expect_err("missing W-key must fail");
assert!(error.contains("W-"));
}
}
@@ -193,6 +193,7 @@ impl Tool for ShowOverviewTool {
.map(|entry| {
serde_json::json!({
"entry_ref": entry.id,
"origin": entry.origin,
"entry_range": entry.entry_range,
"kind": entry.kind.as_str(),
"label": entry.label,
@@ -234,15 +235,16 @@ impl Tool for SearchEntriesTool {
.transpose()?;
let from = params.from.as_deref().map(parse_entry_ref).transpose()?;
let through = params.through.as_deref().map(parse_entry_ref).transpose()?;
let view = self.state.view();
if let (Some(from), Some(through)) = (&from, &through) {
if from.source_index() > through.source_index() {
if view.source_index_for_ref(from) > view.source_index_for_ref(through) {
return Err(ToolError::InvalidArgument(
"SearchEntries from must not be after through".to_string(),
));
}
}
let limit = bounded_limit(params.limit, DEFAULT_PAGE_LIMIT, MAX_PAGE_LIMIT);
let hits = self.state.view().search(&SearchOptions {
let hits = view.search(&SearchOptions {
query: params.query,
kind,
tool_part,
@@ -318,6 +320,7 @@ impl Tool for ReadEntryTool {
.map(|entry| {
serde_json::json!({
"entry_ref": entry.id,
"origin": entry.origin,
"entry_range": entry.entry_range,
"kind": entry.kind.as_str(),
"tool_part": entry.tool_part.map(|part| format!("{part:?}").to_lowercase()),
+222 -13
View File
@@ -33,6 +33,8 @@ use crate::feature::{
use crate::worker::{WorkspaceClient, WorkspaceRequest, WorkspaceRequestMethod};
use agen::tool::{Tool, ToolError, ToolExecutionContext, ToolMeta, ToolOutput};
use super::resource_projection::{project_ticket_detail, project_ticket_query};
#[derive(Clone, Copy)]
enum WorkspaceTicketReadKind {
Query,
@@ -153,8 +155,10 @@ struct WorkspaceQueryTicketInput {
/// stale_after_rescope, and missing_evidence.
#[serde(default)]
attention: Vec<WorkspaceTicketAttentionFilter>,
/// Related Ticket reference. Prefer `T-*`; canonical internal ids remain accepted for compatibility.
related_ticket_id: Option<String>,
relation_kind: Option<WorkspaceTicketRelationFilter>,
/// Linked Objective reference. Prefer `O-*`; canonical internal ids remain accepted for compatibility.
linked_objective_id: Option<String>,
updated_after: Option<String>,
updated_before: Option<String>,
@@ -169,6 +173,7 @@ struct WorkspaceQueryTicketInput {
#[derive(Debug, Deserialize, Serialize, JsonSchema)]
struct WorkspaceShowTicketInput {
/// Ticket reference. Prefer `T-*`; canonical internal ids remain accepted for compatibility.
id: String,
/// Most-recent thread entries to return, bounded by the Backend to 1..=50.
event_limit: Option<usize>,
@@ -229,13 +234,27 @@ impl Tool for WorkspaceTicketReadTool {
.map_err(|error| ToolError::ExecutionFailed(error.to_string()))?;
if !response.is_success() {
return Err(ToolError::ExecutionFailed(format!(
"Workspace Ticket API returned HTTP {}: {}",
response.status, response.body
"Workspace Ticket API request failed with HTTP status {}",
response.status
)));
}
let response_value: Value = serde_json::from_str(&response.body).map_err(|error| {
ToolError::ExecutionFailed(format!(
"Workspace Ticket API returned invalid JSON: {error}"
))
})?;
let content = match self.kind {
WorkspaceTicketReadKind::Query => serde_json::to_string(
&project_ticket_query(response_value).map_err(ToolError::ExecutionFailed)?,
),
WorkspaceTicketReadKind::Show => serde_json::to_string(
&project_ticket_detail(response_value).map_err(ToolError::ExecutionFailed)?,
),
}
.map_err(|error| ToolError::Internal(error.to_string()))?;
Ok(ToolOutput {
summary: self.kind.name().to_string(),
content: Some(response.body),
content: Some(content),
attachments: Vec::new(),
})
}
@@ -733,14 +752,69 @@ impl WorkspaceHttpTicketBackend {
})?;
if !response.is_success() {
return Err(TicketError::Conflict(format!(
"ticket REST API returned HTTP {}: {}",
response.status, response.body
"ticket REST API request failed with HTTP status {}",
response.status
)));
}
serde_json::from_str(&response.body)
let mut value: Value = serde_json::from_str(&response.body).map_err(|error| {
TicketError::Conflict(format!("decode ticket REST response: {error}"))
})?;
Self::canonicalize_ticket_references(&mut value);
serde_json::from_value(value)
.map_err(|error| TicketError::Conflict(format!("decode ticket REST response: {error}")))
}
fn canonicalize_ticket_references(value: &mut Value) {
match value {
Value::Array(values) => {
for value in values {
Self::canonicalize_ticket_references(value);
}
}
Value::Object(object) => {
for value in object.values_mut() {
Self::canonicalize_ticket_references(value);
}
if let Some(resource_key) = object
.get("resource_key")
.and_then(Value::as_str)
.filter(|key| is_canonical_ticket_resource_key(key))
.map(ToOwned::to_owned)
&& object.contains_key("id")
{
object.insert("id".to_string(), Value::String(resource_key));
}
}
_ => {}
}
}
fn resolve_ticket_resource_key(
client: Arc<dyn WorkspaceClient>,
base: &str,
reference: &TicketIdOrSlug,
) -> TicketResult<String> {
let response: Value = Self::request(
client,
WorkspaceRequestMethod::Get,
format!("{base}/{}", Self::ticket_path(reference)),
None,
)?;
response
.get("resource_key")
.or_else(|| {
response
.get("meta")
.and_then(|meta| meta.get("resource_key"))
})
.and_then(Value::as_str)
.filter(|key| is_canonical_ticket_resource_key(key))
.map(ToOwned::to_owned)
.ok_or_else(|| {
TicketError::Conflict("required Ticket human key is unavailable".to_string())
})
}
fn request_unit(
client: Arc<dyn WorkspaceClient>,
method: WorkspaceRequestMethod,
@@ -760,8 +834,8 @@ impl WorkspaceHttpTicketBackend {
})?;
if !response.is_success() {
return Err(TicketError::Conflict(format!(
"ticket REST API returned HTTP {}: {}",
response.status, response.body
"ticket REST API request failed with HTTP status {}",
response.status
)));
}
Ok(TicketBackendOperationResult::Unit)
@@ -802,12 +876,22 @@ impl WorkspaceHttpTicketBackend {
Ok(TicketBackendOperationResult::Tickets(tickets))
}
TicketBackendOperation::Show { id } => {
let ticket = Self::request(
let ticket: Ticket = Self::request(
client,
WorkspaceRequestMethod::Get,
format!("{base}/{}/record", Self::ticket_path(&id)),
None,
)?;
if !ticket
.meta
.resource_key
.as_deref()
.is_some_and(is_canonical_ticket_resource_key)
{
return Err(TicketError::Conflict(
"required Ticket human key is unavailable".to_string(),
));
}
Ok(TicketBackendOperationResult::Ticket(ticket))
}
TicketBackendOperation::Create { input } => {
@@ -910,7 +994,14 @@ impl WorkspaceHttpTicketBackend {
})?),
),
TicketBackendOperation::AddTicketRelation { id, relation } => {
let relation = Self::request(
let source_resource_key =
Self::resolve_ticket_resource_key(client.clone(), &base, &id)?;
let target_resource_key = Self::resolve_ticket_resource_key(
client.clone(),
&base,
&TicketIdOrSlug::Id(relation.target.clone()),
)?;
let mut relation: TicketRelation = Self::request(
client,
WorkspaceRequestMethod::Post,
format!("{base}/{}/relations", Self::ticket_path(&id)),
@@ -918,20 +1009,30 @@ impl WorkspaceHttpTicketBackend {
TicketError::Conflict(format!("serialize Ticket relation: {error}"))
})?),
)?;
relation.ticket_id = source_resource_key;
relation.target = target_resource_key;
relation.author = "workspace".to_string();
Ok(TicketBackendOperationResult::Relation(relation))
}
TicketBackendOperation::RemoveTicketRelation { id, kind, target } => {
let source_resource_key =
Self::resolve_ticket_resource_key(client.clone(), &base, &id)?;
let target_resource_key =
Self::resolve_ticket_resource_key(client.clone(), &base, &target)?;
let target = match target {
TicketIdOrSlug::Id(value)
| TicketIdOrSlug::Slug(value)
| TicketIdOrSlug::Query(value) => value,
};
let relation = Self::request(
let mut relation: TicketRelation = Self::request(
client,
WorkspaceRequestMethod::Delete,
format!("{base}/{}/relations", Self::ticket_path(&id)),
Some(serde_json::json!({ "kind": kind, "target": target })),
)?;
relation.ticket_id = source_resource_key;
relation.target = target_resource_key;
relation.author = "workspace".to_string();
Ok(TicketBackendOperationResult::Relation(relation))
}
TicketBackendOperation::QueryTicketRelations { ticket, kind } => {
@@ -1266,6 +1367,23 @@ mod tests {
.expect("tool exists")
}
#[test]
fn workspace_ticket_backend_canonicalizes_model_facing_ticket_ids() {
let mut value = serde_json::json!({
"id": "00001INTERNAL",
"resource_key": "T-42",
"nested": {
"id": "00002INTERNAL",
"resource_key": "T-43"
},
"body": "user-authored 00003BODY stays unchanged"
});
WorkspaceHttpTicketBackend::canonicalize_ticket_references(&mut value);
assert_eq!(value["id"], "T-42");
assert_eq!(value["nested"]["id"], "T-43");
assert_eq!(value["body"], "user-authored 00003BODY stays unchanged");
}
#[test]
fn workspace_ticket_reads_expose_bounded_query_and_show_contracts_without_legacy_aliases() {
let client: Arc<dyn WorkspaceClient> = Arc::new(
@@ -1742,11 +1860,102 @@ provider = "github"
server.join().unwrap();
}
#[test]
fn workspace_http_backend_records_relation_with_authoritative_human_keys() {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let addr = listener.local_addr().unwrap();
let server = thread::spawn(move || {
for (expected_path, resource_key) in [
("GET /api/w/workspace-a/tickets/01SOURCE HTTP/1.1", "T-1"),
("GET /api/w/workspace-a/tickets/01TARGET HTTP/1.1", "T-2"),
] {
let (mut stream, _) = listener.accept().unwrap();
let mut buffer = [0_u8; 8192];
let len = stream.read(&mut buffer).unwrap();
let request = String::from_utf8_lossy(&buffer[..len]);
assert!(request.starts_with(expected_path));
let body = serde_json::json!({"meta": {"resource_key": resource_key}}).to_string();
write!(
stream,
"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\n\r\n{}",
body.len(), body
)
.unwrap();
}
let (mut stream, _) = listener.accept().unwrap();
let mut buffer = [0_u8; 8192];
let len = stream.read(&mut buffer).unwrap();
let request = String::from_utf8_lossy(&buffer[..len]);
assert!(
request.starts_with("POST /api/w/workspace-a/tickets/01SOURCE/relations HTTP/1.1")
);
let body = serde_json::to_string(&TicketRelation {
ticket_id: "01SOURCE".to_string(),
kind: TicketRelationKind::DependsOn,
target: "01TARGET".to_string(),
note: None,
author: "worker-internal".to_string(),
at: "2026-08-06T00:00:00Z".to_string(),
})
.unwrap();
write!(
stream,
"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\n\r\n{}",
body.len(),
body
)
.unwrap();
});
let backend = WorkspaceHttpTicketBackend::new(Arc::new(
crate::worker::TestWorkspaceHttpClient::new("workspace-a", format!("http://{addr}")),
));
let relation = backend
.add_ticket_relation(
TicketIdOrSlug::Id("01SOURCE".to_string()),
NewTicketRelation {
kind: TicketRelationKind::DependsOn,
target: "01TARGET".to_string(),
note: None,
author: None,
},
)
.unwrap();
server.join().unwrap();
assert_eq!(relation.ticket_id, "T-1");
assert_eq!(relation.target, "T-2");
assert_eq!(relation.author, "workspace");
}
#[test]
fn workspace_http_backend_deletes_exact_ticket_relation() {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let base_url = format!("http://{}", listener.local_addr().unwrap());
let server = thread::spawn(move || {
for (expected_path, resource_key) in [
("GET /api/w/workspace-a/tickets/01SOURCE HTTP/1.1", "T-1"),
("GET /api/w/workspace-a/tickets/01TARGET HTTP/1.1", "T-2"),
] {
let (mut stream, _) = listener.accept().unwrap();
let mut buffer = [0_u8; 8192];
let len = stream.read(&mut buffer).unwrap();
let request = String::from_utf8_lossy(&buffer[..len]);
assert!(request.starts_with(expected_path));
let response_body = serde_json::json!({
"meta": {"resource_key": resource_key}
})
.to_string();
write!(
stream,
"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\n\r\n{}",
response_body.len(),
response_body
)
.unwrap();
}
let (mut stream, _) = listener.accept().unwrap();
let mut buffer = [0_u8; 8192];
let len = stream.read(&mut buffer).unwrap();
@@ -1787,8 +1996,8 @@ provider = "github"
.unwrap();
server.join().unwrap();
assert_eq!(removed.ticket_id, "01SOURCE");
assert_eq!(removed.target, "01TARGET");
assert_eq!(removed.ticket_id, "T-1");
assert_eq!(removed.target, "T-2");
}
#[test]
@@ -1,11 +1,12 @@
use std::sync::Arc;
#[cfg(test)]
use agen::Item;
use agen::tool::{Tool, ToolDefinition, ToolError, ToolMeta, ToolOutput};
use async_trait::async_trait;
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
use session_store::collect_state;
use session_store::{LogEntry, collect_state};
use super::manage_worker::{WORKER_CONTROL_SERVICE_ID, WorkerControlService};
use crate::feature::{
@@ -60,7 +61,27 @@ pub struct WorkerObservationSubject {
#[derive(Debug, Clone)]
pub struct WorkerSessionCapture {
pub segment_id: String,
pub items: Vec<Item>,
pub entries: Vec<agen::HistoryEntry<crate::SessionHistoryMetadata>>,
}
impl WorkerSessionCapture {
pub fn from_log_entries(
segment_id: impl Into<String>,
log_entries: &[LogEntry],
) -> Result<Self, String> {
let segment_id = segment_id.into();
let state = collect_state(log_entries);
let parsed_segment_id = segment_id.parse().unwrap_or_default();
let entries = crate::session_history::restore_history_entries(
state.session_id.unwrap_or_default(),
parsed_segment_id,
log_entries,
)?;
Ok(Self {
segment_id,
entries,
})
}
}
#[derive(Debug, thiserror::Error)]
@@ -161,9 +182,17 @@ impl WorkerObservationProvider for WorkspaceClientWorkerObservationProvider {
})
.collect::<Result<Vec<session_store::LogEntry>, _>>()?;
let state = collect_state(&entries);
let segment_id = response.segment_id;
let parsed_segment_id = segment_id.parse().unwrap_or_default();
let typed_entries = crate::session_history::restore_history_entries(
state.session_id.unwrap_or_default(),
parsed_segment_id,
&entries,
)
.map_err(WorkerObservationError::Unavailable)?;
Ok(WorkerSessionCapture {
segment_id: response.segment_id,
items: state.history,
segment_id,
entries: typed_entries,
})
}
}
@@ -392,9 +421,15 @@ impl WorkerObservationProvider for SpawnedSubWorkerObservationProvider {
.ok_or(WorkerObservationError::NotFound)?;
let entries = record.session.entries();
let state = collect_state(&entries);
let typed_entries = crate::session_history::restore_history_entries(
state.session_id.unwrap_or_default(),
Default::default(),
&entries,
)
.map_err(WorkerObservationError::Unavailable)?;
Ok(WorkerSessionCapture {
segment_id: format!("subworker:{name}"),
items: state.history,
entries: typed_entries,
})
}
}
@@ -508,6 +543,7 @@ impl Tool for ViewSessionOverviewTool {
.map(|entry| {
serde_json::json!({
"entry_ref": entry.id,
"origin": entry.origin,
"entry_range": entry.entry_range,
"kind": entry.kind.as_str(),
"label": entry.label,
@@ -547,7 +583,7 @@ impl Tool for SearchSessionEntriesTool {
let from = params.from.as_deref().map(parse_entry_ref).transpose()?;
let through = params.through.as_deref().map(parse_entry_ref).transpose()?;
if let (Some(from), Some(through)) = (&from, &through) {
if from.source_index() > through.source_index() {
if view.source_index_for_ref(from) > view.source_index_for_ref(through) {
return Err(ToolError::InvalidArgument(
"SearchSessionEntries from must not be after through".to_string(),
));
@@ -573,6 +609,7 @@ impl Tool for SearchSessionEntriesTool {
.map(|entry| {
serde_json::json!({
"entry_ref": entry.id,
"origin": entry.origin,
"entry_range": entry.entry_range,
"kind": entry.kind.as_str(),
"tool_part": entry.tool_part.map(|part| format!("{part:?}").to_lowercase()),
@@ -628,6 +665,7 @@ impl Tool for ReadSessionEntryTool {
.map(|entry| {
serde_json::json!({
"entry_ref": entry.id,
"origin": entry.origin,
"entry_range": entry.entry_range,
"kind": entry.kind.as_str(),
"tool_part": entry.tool_part.map(|part| format!("{part:?}").to_lowercase()),
@@ -661,7 +699,10 @@ async fn latest_view(
.capture_worker_session(subject)
.await
.map_err(tool_error)?;
Ok(SessionCapture::new(capture.segment_id, capture.items))
Ok(SessionCapture::from_history_entries(
capture.segment_id,
capture.entries,
))
}
fn parse_input<T: serde::de::DeserializeOwned>(
@@ -751,9 +792,23 @@ mod tests {
if subject != &granted_subject() {
return Err(WorkerObservationError::NotFound);
}
let entries = self
.captures
.lock()
.unwrap()
.clone()
.into_iter()
.enumerate()
.map(|(index, item)| {
let mut metadata = crate::SessionHistoryMetadata::legacy_unknown();
metadata.entry_id =
session_store::LoggedSessionHistoryEntryId(format!("fake-{index:08}"));
agen::HistoryEntry::new(item, metadata)
})
.collect();
Ok(WorkerSessionCapture {
segment_id: "segment".to_string(),
items: self.captures.lock().unwrap().clone(),
entries,
})
}
}
@@ -796,7 +851,7 @@ mod tests {
let read = read_definition(provider.clone())().1;
let hidden = read
.execute(
r#"{"subject":{"kind":"runtime_worker","runtime_id":"runtime-1","worker_id":"unauthorized"},"entry_ref":"E00000000"}"#,
r#"{"subject":{"kind":"runtime_worker","runtime_id":"runtime-1","worker_id":"unauthorized"},"entry_ref":"Efake-00000000"}"#,
agen::tool::ToolExecutionContext::direct(),
)
.await
@@ -810,7 +865,7 @@ mod tests {
.push(message("a1", Role::Assistant, "second"));
let output = read
.execute(
r#"{"subject":{"kind":"runtime_worker","runtime_id":"runtime-1","worker_id":"granted"},"entry_ref":"E00000000"}"#,
r#"{"subject":{"kind":"runtime_worker","runtime_id":"runtime-1","worker_id":"granted"},"entry_ref":"Efake-00000000"}"#,
agen::tool::ToolExecutionContext::direct(),
)
.await
@@ -819,7 +874,7 @@ mod tests {
let output = read
.execute(
r#"{"subject":{"kind":"runtime_worker","runtime_id":"runtime-1","worker_id":"granted"},"entry_ref":"E00000001"}"#,
r#"{"subject":{"kind":"runtime_worker","runtime_id":"runtime-1","worker_id":"granted"},"entry_ref":"Efake-00000001"}"#,
agen::tool::ToolExecutionContext::direct(),
)
.await
+177 -25
View File
@@ -10,7 +10,7 @@ use std::collections::HashMap;
use std::sync::{Arc, Mutex};
use agen::timeline::event::UsageEvent;
use agen::{Engine, llm_client::LlmClient};
use agen::{Engine, EngineError, llm_client::LlmClient};
use manifest::{Scope, WorkerManifest};
use protocol::{Event, InFlightSnapshot, WorkerStatus};
use session_store::{LogEntry, SegmentId, SessionId, Store, StoreError, TraceEntry};
@@ -55,7 +55,17 @@ pub(crate) struct InternalWorkerSpec {
pub input: String,
pub cache_key: Option<String>,
pub max_turns: Option<u32>,
pub engine_configurator: Option<Box<dyn FnOnce(&mut Engine<Box<dyn LlmClient>>) + Send>>,
pub engine_configurator: Option<
Box<
dyn FnOnce(
&mut Engine<
Box<dyn LlmClient>,
agen::state::Mutable,
crate::SessionHistoryMetadata,
>,
) + Send,
>,
>,
pub features: FeatureRegistryBuilder,
pub required_tools: &'static [&'static str],
pub authority: InternalWorkerAuthority,
@@ -124,7 +134,9 @@ where
let last_usage = Arc::new(Mutex::new(None::<UsageEvent>));
let usage_slot = last_usage.clone();
let mut engine = Engine::new(client).system_prompt(system_prompt);
let mut engine =
Engine::<_, agen::state::Mutable, crate::SessionHistoryMetadata>::new_annotated(client)
.system_prompt(system_prompt);
engine.on_usage(move |usage| {
if let Ok(mut slot) = usage_slot.lock() {
*slot = Some(usage.clone());
@@ -199,12 +211,28 @@ where
on_cancel_sender(worker.engine_mut().cancel_sender());
match worker.run_text(&input).await {
Ok(lifecycle) => Ok(InternalWorkerResult {
Ok(lifecycle @ WorkerRunResult::Finished)
| Ok(lifecycle @ WorkerRunResult::Paused)
| Ok(lifecycle @ WorkerRunResult::RolledBack) => Ok(InternalWorkerResult {
usage: last_usage.lock().ok().and_then(|slot| slot.clone()),
identity,
lifecycle,
history_entries: store.entries_count(session_id, segment_id),
}),
Ok(WorkerRunResult::LimitReached) => Err(InternalWorkerError {
source: WorkerError::Engine(EngineError::Aborted(
"internal Worker reached its turn limit".to_string(),
)),
usage: last_usage.lock().ok().and_then(|slot| slot.clone()),
identity,
history_entries: store.entries_count(session_id, segment_id),
}),
Ok(WorkerRunResult::Interrupted { message, .. }) => Err(InternalWorkerError {
source: WorkerError::Engine(EngineError::Aborted(message)),
usage: last_usage.lock().ok().and_then(|slot| slot.clone()),
identity,
history_entries: store.entries_count(session_id, segment_id),
}),
Err(source) => Err(InternalWorkerError {
source,
usage: last_usage.lock().ok().and_then(|slot| slot.clone()),
@@ -232,6 +260,7 @@ impl Default for InternalWorkerVisibility {
pub(crate) enum InternalWorkerSessionStatus {
Idle,
Running,
Paused,
Stopping,
Stopped,
Failed,
@@ -242,9 +271,10 @@ impl InternalWorkerSessionStatus {
match self {
Self::Idle => 0,
Self::Running => 1,
Self::Stopping => 2,
Self::Stopped => 3,
Self::Failed => 4,
Self::Paused => 2,
Self::Stopping => 3,
Self::Stopped => 4,
Self::Failed => 5,
}
}
@@ -252,13 +282,35 @@ impl InternalWorkerSessionStatus {
match value {
0 => Self::Idle,
1 => Self::Running,
2 => Self::Stopping,
3 => Self::Stopped,
2 => Self::Paused,
3 => Self::Stopping,
4 => Self::Stopped,
_ => Self::Failed,
}
}
}
fn classify_internal_turn_result(
result: Result<WorkerRunResult, WorkerError>,
) -> (InternalWorkerSessionStatus, Option<String>) {
match result {
Ok(WorkerRunResult::Finished) => (InternalWorkerSessionStatus::Idle, None),
Ok(WorkerRunResult::Paused) => (InternalWorkerSessionStatus::Paused, None),
Ok(WorkerRunResult::LimitReached) => (
InternalWorkerSessionStatus::Stopped,
Some("internal Worker reached its turn limit".to_string()),
),
Ok(WorkerRunResult::Interrupted { message, .. }) => {
(InternalWorkerSessionStatus::Stopped, Some(message))
}
Ok(WorkerRunResult::RolledBack) => (
InternalWorkerSessionStatus::Stopped,
Some("internal Worker run was cancelled before AI output".to_string()),
),
Err(error) => (InternalWorkerSessionStatus::Failed, Some(error.to_string())),
}
}
#[derive(Debug, thiserror::Error)]
pub(crate) enum InternalWorkerSessionError {
#[error("failed to build internal Worker session: {message}")]
@@ -353,10 +405,11 @@ impl InternalWorkerSessionHandle {
entries,
status: match self.status() {
InternalWorkerSessionStatus::Running => WorkerStatus::Running,
InternalWorkerSessionStatus::Paused => WorkerStatus::Paused,
InternalWorkerSessionStatus::Idle => WorkerStatus::Idle,
InternalWorkerSessionStatus::Stopping
| InternalWorkerSessionStatus::Stopped
| InternalWorkerSessionStatus::Failed => WorkerStatus::Paused,
| InternalWorkerSessionStatus::Failed => WorkerStatus::Stopped,
},
error: self.last_error.lock().unwrap().clone(),
in_flight,
@@ -388,6 +441,7 @@ impl InternalWorkerSessionHandle {
.map_err(
|current| match InternalWorkerSessionStatus::decode(current) {
InternalWorkerSessionStatus::Running
| InternalWorkerSessionStatus::Paused
| InternalWorkerSessionStatus::Stopping => InternalWorkerSessionError::Busy,
InternalWorkerSessionStatus::Stopped | InternalWorkerSessionStatus::Failed => {
InternalWorkerSessionError::Stopped
@@ -494,7 +548,9 @@ pub(crate) async fn spawn_internal_worker_session(
let last_usage = Arc::new(Mutex::new(None::<UsageEvent>));
let usage_slot = last_usage.clone();
let mut engine = Engine::new(client).system_prompt(system_prompt);
let mut engine =
Engine::<_, agen::state::Mutable, crate::SessionHistoryMetadata>::new_annotated(client)
.system_prompt(system_prompt);
engine.on_usage(move |usage| {
if let Ok(mut slot) = usage_slot.lock() {
*slot = Some(usage.clone());
@@ -591,7 +647,9 @@ pub(crate) fn prepare_internal_worker_from_spec(
manifest.compaction = None;
manifest.memory = None;
let mut engine = Engine::new(client).system_prompt(system_prompt);
let mut engine =
Engine::<_, agen::state::Mutable, crate::SessionHistoryMetadata>::new_annotated(client)
.system_prompt(system_prompt);
engine.set_cache_key(cache_key);
engine.set_max_turns(max_turns);
if let Some(configure) = engine_configurator {
@@ -733,13 +791,7 @@ pub(crate) async fn prepare_internal_worker_session(
loop {
tokio::select! {
result = &mut run => {
let (turn_status, error) = match result {
Ok(_) => (InternalWorkerSessionStatus::Idle, None),
Err(error) => (
InternalWorkerSessionStatus::Failed,
Some(error.to_string()),
),
};
let (turn_status, error) = classify_internal_turn_result(result);
actor_in_flight.clear();
status.store(turn_status.encode(), std::sync::atomic::Ordering::Release);
if let Some(message) = error {
@@ -748,11 +800,20 @@ pub(crate) async fn prepare_internal_worker_session(
code: protocol::ErrorCode::Internal,
message,
});
} else {
let _ = event_tx.send(Event::Status {
status: WorkerStatus::Idle,
});
}
let protocol_status = match turn_status {
InternalWorkerSessionStatus::Idle => WorkerStatus::Idle,
InternalWorkerSessionStatus::Paused => WorkerStatus::Paused,
InternalWorkerSessionStatus::Stopped
| InternalWorkerSessionStatus::Failed => WorkerStatus::Stopped,
InternalWorkerSessionStatus::Running
| InternalWorkerSessionStatus::Stopping => {
unreachable!("run completion cannot remain active")
}
};
let _ = event_tx.send(Event::Status {
status: protocol_status,
});
if let Some(callback) = &on_turn_end {
callback(turn_status);
}
@@ -766,7 +827,7 @@ pub(crate) async fn prepare_internal_worker_session(
let _ = (&mut run).await;
actor_in_flight.clear();
status.store(InternalWorkerSessionStatus::Stopped.encode(), std::sync::atomic::Ordering::Release);
let _ = event_tx.send(Event::Status { status: WorkerStatus::Paused });
let _ = event_tx.send(Event::Status { status: WorkerStatus::Stopped });
let _ = event_tx.send(Event::Shutdown);
state_changed.notify_waiters();
let _ = done.send(());
@@ -792,7 +853,7 @@ pub(crate) async fn prepare_internal_worker_session(
std::sync::atomic::Ordering::Release,
);
let _ = event_tx.send(Event::Status {
status: WorkerStatus::Paused,
status: WorkerStatus::Stopped,
});
let _ = event_tx.send(Event::Shutdown);
state_changed.notify_waiters();
@@ -1102,6 +1163,26 @@ mod tests {
}
}
#[derive(Clone)]
struct FailingClient;
#[async_trait]
impl LlmClient for FailingClient {
fn clone_boxed(&self) -> Box<dyn LlmClient> {
Box::new(self.clone())
}
async fn stream(
&self,
_request: Request,
) -> Result<Pin<Box<dyn Stream<Item = Result<LlmEvent, ClientError>> + Send>>, ClientError>
{
Err(ClientError::Config(
"intentional internal failure".to_string(),
))
}
}
#[derive(Clone)]
struct CancelBeforeAiClient {
calls: Arc<AtomicUsize>,
@@ -1215,6 +1296,77 @@ permission = "write"
assert_eq!(result.identity.kind, "test");
}
#[test]
fn internal_turn_result_mapping_is_exhaustive() {
let cases = [
(
WorkerRunResult::Finished,
InternalWorkerSessionStatus::Idle,
false,
),
(
WorkerRunResult::Paused,
InternalWorkerSessionStatus::Paused,
false,
),
(
WorkerRunResult::LimitReached,
InternalWorkerSessionStatus::Stopped,
true,
),
(
WorkerRunResult::Interrupted {
code: protocol::ErrorCode::Internal,
message: "cancelled".to_string(),
},
InternalWorkerSessionStatus::Stopped,
true,
),
(
WorkerRunResult::RolledBack,
InternalWorkerSessionStatus::Stopped,
true,
),
];
for (result, expected_status, expects_error) in cases {
let (status, error) = classify_internal_turn_result(Ok(result));
assert_eq!(status, expected_status);
assert_eq!(error.is_some(), expects_error);
}
let (status, error) = classify_internal_turn_result(Err(WorkerError::Engine(
EngineError::Aborted("fatal".to_string()),
)));
assert_eq!(status, InternalWorkerSessionStatus::Failed);
assert!(error.is_some_and(|message| message.contains("fatal")));
}
#[tokio::test]
async fn fatal_internal_run_transitions_to_stopped_protocol_status() {
let calls = Arc::new(AtomicUsize::new(0));
let mut internal_spec = spec(calls, &[]);
internal_spec.client = Box::new(FailingClient);
let handle = spawn_internal_worker_session(internal_spec)
.await
.expect("spawn failing Internal Worker session");
assert_eq!(
handle.wait_until_idle().await,
InternalWorkerSessionStatus::Stopped
);
assert_eq!(handle.status(), InternalWorkerSessionStatus::Stopped);
assert_eq!(handle.protocol_snapshot().status, WorkerStatus::Stopped);
assert!(
handle
.last_error
.lock()
.unwrap()
.as_ref()
.is_some_and(|message| message.contains("intentional internal failure"))
);
}
#[tokio::test]
async fn session_accepts_follow_up_turns_and_stops_without_runtime_registration() {
let calls = Arc::new(AtomicUsize::new(0));
+13 -2
View File
@@ -13,7 +13,7 @@
#[cfg(test)]
use crate::prompt::catalog::PromptCatalog;
use agen::Item;
use agen::{Item, ToolResultDisposition};
/// Build synthetic `Item::ToolResult` items for every unanswered
/// `Item::ToolCall` in `history`, preserving order.
@@ -28,7 +28,16 @@ pub(crate) fn orphan_tool_result_closures(history: &[Item], summary: &str) -> Ve
for item in history {
if let Item::ToolCall { call_id, .. } = item {
if !answered.contains(call_id.as_str()) {
out.push(Item::tool_result(call_id.clone(), summary));
out.push(Item::tool_result_item_with_disposition_and_attachments(
call_id.clone(),
summary,
Some(
"Execution ended before completion could be confirmed. Completion and side effects are unknown."
.to_string(),
),
ToolResultDisposition::OutcomeUnknown,
Vec::new(),
));
}
}
}
@@ -77,10 +86,12 @@ mod tests {
Item::ToolResult {
call_id,
summary: got,
disposition,
..
} => {
assert_eq!(call_id, "c1");
assert_eq!(got, &summary);
assert_eq!(*disposition, ToolResultDisposition::OutcomeUnknown);
}
other => panic!("expected ToolResult, got {other:?}"),
}
+39 -2
View File
@@ -8,6 +8,7 @@
//! decisions (continue / skip / abort / pause).
use std::borrow::Cow;
use std::collections::VecDeque;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
@@ -33,7 +34,9 @@ use crate::hook::{
};
use crate::ipc::notify_buffer::{NotifyBuffer, build_system_item_with_provenance};
use crate::prompt::catalog::PromptCatalog;
use crate::session_history::SessionHistoryMetadata;
use crate::worker::SystemItemCommitter;
use agen::HistoryEntry;
use agen::token_counter::total_tokens;
/// Maximum number of bytes copied into `TurnEndInfo::final_text_preview`.
@@ -73,6 +76,7 @@ pub(crate) struct WorkerInterceptor {
/// worker. `None` in tests / `Worker::new` paths where no writer is
/// attached.
log_writer: Option<Arc<dyn SystemItemCommitter>>,
pending_committed_history: Arc<Mutex<VecDeque<HistoryEntry<SessionHistoryMetadata>>>>,
/// Next turn index assigned by `on_prompt_submit`.
next_turn_index: AtomicUsize,
/// Tool calls observed in the current turn (reset on each new prompt).
@@ -80,6 +84,7 @@ pub(crate) struct WorkerInterceptor {
}
impl WorkerInterceptor {
#[cfg(test)]
pub(crate) fn new(
registry: Arc<HookRegistry>,
compact_state: Option<Arc<CompactState>>,
@@ -88,6 +93,28 @@ impl WorkerInterceptor {
pending_attachments: Arc<Mutex<Vec<SystemItem>>>,
prompts: Arc<ArcSwap<PromptCatalog>>,
log_writer: Option<Arc<dyn SystemItemCommitter>>,
) -> Self {
Self::new_with_history_queue(
registry,
compact_state,
usage_history,
pending_notifies,
pending_attachments,
prompts,
log_writer,
Arc::new(Mutex::new(VecDeque::new())),
)
}
pub(crate) fn new_with_history_queue(
registry: Arc<HookRegistry>,
compact_state: Option<Arc<CompactState>>,
usage_history: Option<Arc<Mutex<Vec<UsageRecord>>>>,
pending_notifies: NotifyBuffer,
pending_attachments: Arc<Mutex<Vec<SystemItem>>>,
prompts: Arc<ArcSwap<PromptCatalog>>,
log_writer: Option<Arc<dyn SystemItemCommitter>>,
pending_committed_history: Arc<Mutex<VecDeque<HistoryEntry<SessionHistoryMetadata>>>>,
) -> Self {
Self {
registry,
@@ -99,6 +126,7 @@ impl WorkerInterceptor {
prompts,
prompt_workspace_id: None,
log_writer,
pending_committed_history,
next_turn_index: AtomicUsize::new(0),
tool_calls_this_turn: AtomicUsize::new(0),
}
@@ -125,7 +153,11 @@ impl WorkerInterceptor {
return Ok(());
};
for item in items {
writer.commit_system_item(item.clone())?;
let entry = writer.commit_system_item(item.clone())?;
self.pending_committed_history
.lock()
.expect("pending committed history poisoned")
.push_back(entry);
}
Ok(())
}
@@ -507,7 +539,12 @@ mod tests {
&self,
entry: session_store::LogEntry,
) -> Result<(), session_store::StoreError> {
if let session_store::LogEntry::SystemItem { item, .. } = entry {
let item = match entry {
session_store::LogEntry::SystemItem { item, .. } => Some(item),
session_store::LogEntry::AnnotatedSystemItem { entry, .. } => Some(entry.item),
_ => None,
};
if let Some(item) = item {
self.committed
.lock()
.expect("committed system-item list poisoned")
+8 -2
View File
@@ -29,15 +29,21 @@ pub fn subscribe_worker_protocol_session(handle: &WorkerHandle) -> WorkerProtoco
pub fn live_log_entry_event(entry: LogEntry) -> Option<Event> {
match entry {
LogEntry::SegmentStart { .. } => {
entry @ (LogEntry::SegmentStart { .. } | LogEntry::AnnotatedSegmentStart { .. }) => {
let value = serde_json::to_value(&entry).expect("LogEntry is Serialize");
Some(Event::SegmentRotated { entry: value })
}
LogEntry::UserInput { segments, .. } => Some(Event::UserMessage { segments }),
LogEntry::UserInput { segments, .. } | LogEntry::AnnotatedUserInput { segments, .. } => {
Some(Event::UserMessage { segments })
}
LogEntry::SystemItem { item, .. } => {
let value = serde_json::to_value(&item).expect("SystemItem is Serialize");
Some(Event::SystemItem { item: value })
}
LogEntry::AnnotatedSystemItem { entry, .. } => {
let value = serde_json::to_value(&entry.item).expect("SystemItem is Serialize");
Some(Event::SystemItem { item: value })
}
LogEntry::Invoke { trigger, .. } => Some(Event::InvokeStart { kind: trigger }),
other => {
// `SegmentLogSink::is_live_relevant` keeps non-live-relevant
+8 -2
View File
@@ -12,6 +12,7 @@ pub mod prompt;
pub mod runtime;
pub mod segment_log_sink;
mod session_capture;
mod session_history;
pub mod shared_state;
mod shutdown_after_idle;
pub mod skill;
@@ -33,14 +34,19 @@ pub use manifest::{
};
pub use model_client::{ProviderError, build_client};
pub use prompt::catalog::{
CatalogError, EffectivePromptCatalog, PromptCatalog, WorkerPrompt, WorkspacePromptProjection,
prompt_schema_source,
CatalogError, EffectivePromptCatalog, OrchestratorQueueAttentionContext,
OrchestratorQueueAttentionPrompt, OrchestratorQueueAttentionTicket, PromptCatalog,
WorkerPrompt, WorkspacePromptProjection, prompt_schema_source,
};
pub use prompt::source::PromptCatalogSource;
pub use prompt::system::{SystemPromptContext, SystemPromptError, SystemPromptTemplate};
pub use protocol::{ErrorCode, Event, Method, TurnResult, WorkerStatus};
pub use runtime::dir::RuntimeDir;
pub use segment_log_sink::SegmentLogSink;
pub use session_history::{
SessionHistoryDerivation, SessionHistoryEntryId, SessionHistoryMetadata,
WorkerHistoryProvenance, WorkerSubjectSnapshot,
};
pub use shared_state::WorkerSharedState;
pub use worker::{
LocalWorkingDirectory, WORKER_INPUT_SUBMISSION_EXTENSION_DOMAIN, Worker, WorkerError,
+1 -1
View File
@@ -34,7 +34,7 @@ impl PermissionHook {
}
}
impl<C: LlmClient, St: Store> Worker<C, St> {
impl<C: LlmClient + 'static, St: Store> Worker<C, St> {
pub(crate) fn apply_permissions_from_manifest(&mut self) {
let Some(permissions) = self.manifest().permissions.clone() else {
return;
+149
View File
@@ -141,8 +141,93 @@ impl WorkerPrompt {
];
}
/// Model-visible queued Ticket projection shared by Server and TUI backlog attention paths.
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct OrchestratorQueueAttentionTicket {
resource_key: String,
title: String,
}
impl OrchestratorQueueAttentionTicket {
pub fn new(
resource_key: impl Into<String>,
title: impl Into<String>,
) -> Result<Self, CatalogError> {
let resource_key = resource_key.into();
if !is_ticket_resource_key(&resource_key) {
return Err(CatalogError::InvalidQueueAttentionResourceKey);
}
Ok(Self {
resource_key,
title: bounded_queue_attention_text(&title.into(), 240),
})
}
}
/// Shared model-visible context for every Orchestrator backlog attention renderer.
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct OrchestratorQueueAttentionContext {
tickets: Vec<OrchestratorQueueAttentionTicket>,
separator: &'static str,
omitted_ticket_count: usize,
}
impl OrchestratorQueueAttentionContext {
pub const MAX_TICKETS: usize = 20;
pub fn new(tickets: Vec<OrchestratorQueueAttentionTicket>) -> Self {
let omitted_ticket_count = tickets.len().saturating_sub(Self::MAX_TICKETS);
Self {
tickets: tickets.into_iter().take(Self::MAX_TICKETS).collect(),
separator: "",
omitted_ticket_count,
}
}
}
/// Prompt-catalog entries that must share the same backlog-attention body contract.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum OrchestratorQueueAttentionPrompt {
Server,
Tui,
}
impl OrchestratorQueueAttentionPrompt {
fn key(self) -> &'static str {
match self {
Self::Server => "internal.workspace_orchestrator_queue_attention",
Self::Tui => "panel.orchestrator_idle_queue_notice",
}
}
}
fn is_ticket_resource_key(input: &str) -> bool {
input.len() <= 32
&& input.strip_prefix("T-").is_some_and(|suffix| {
!suffix.is_empty() && suffix.bytes().all(|byte| byte.is_ascii_digit())
})
}
fn bounded_queue_attention_text(input: &str, max_chars: usize) -> String {
let mut output = String::new();
for (index, character) in input.chars().enumerate() {
if index == max_chars {
output.push('…');
break;
}
output.push(if character.is_control() {
' '
} else {
character
});
}
output
}
#[derive(Debug, Error)]
pub enum CatalogError {
#[error("queued Ticket resource key is missing or invalid")]
InvalidQueueAttentionResourceKey,
#[error("failed to build builtin Prompt source tree: {0}")]
BuiltinTree(String),
#[error("failed to evaluate builtin Prompt source tree: {0}")]
@@ -319,6 +404,14 @@ impl PromptCatalog {
self.render_name(key, Value::from_serialize(context))
}
pub fn orchestrator_queue_attention(
&self,
prompt: OrchestratorQueueAttentionPrompt,
context: &OrchestratorQueueAttentionContext,
) -> Result<String, CatalogError> {
self.render_serializable(prompt.key(), context)
}
pub fn render_name(&self, key: &str, ctx: Value) -> Result<String, CatalogError> {
let template = self
.env
@@ -653,6 +746,62 @@ mod tests {
assert!(reviewer.contains("target-only movement does not invalidate approval"));
}
#[test]
fn queue_attention_prompts_share_sanitized_contract_and_true_truncation() {
let catalog = PromptCatalog::builtins_only().unwrap();
let tickets = (1..=OrchestratorQueueAttentionContext::MAX_TICKETS + 1)
.map(|index| {
OrchestratorQueueAttentionTicket::new(
format!("T-{index}"),
format!("Ticket {index}\nwith control\u{7}"),
)
.unwrap()
})
.collect();
let context = OrchestratorQueueAttentionContext::new(tickets);
let server = catalog
.orchestrator_queue_attention(OrchestratorQueueAttentionPrompt::Server, &context)
.unwrap();
let tui = catalog
.orchestrator_queue_attention(OrchestratorQueueAttentionPrompt::Tui, &context)
.unwrap();
assert_eq!(server, tui);
assert!(server.starts_with("Queued Tickets require attention:"));
assert!(server.contains("- T-1 — Ticket 1 with control "));
assert!(!server.contains("T-21"));
assert!(server.contains("were omitted from this notice: 1"));
assert!(server.contains("Re-query current Ticket authority"));
assert!(server.contains("Reread the current Ticket state before acting"));
for secret in [
"workspace_id",
"Workspace:",
"runtime_id",
"worker_id",
"bounded",
] {
assert!(!server.contains(secret), "leaked {secret}: {server}");
}
}
#[test]
fn queue_attention_prompt_omits_truncation_text_for_complete_list() {
let catalog = PromptCatalog::builtins_only().unwrap();
let context = OrchestratorQueueAttentionContext::new(vec![
OrchestratorQueueAttentionTicket::new("T-541", "Attention contract").unwrap(),
]);
let rendered = catalog
.orchestrator_queue_attention(OrchestratorQueueAttentionPrompt::Server, &context)
.unwrap();
assert!(rendered.contains("- T-541 — Attention contract"));
assert!(!rendered.contains("omitted"));
assert!(matches!(
OrchestratorQueueAttentionTicket::new("opaque-id", "must fail"),
Err(CatalogError::InvalidQueueAttentionResourceKey)
));
}
#[test]
fn graph_rejects_dynamic_legacy_missing_and_cycles() {
let invalid = BTreeMap::from([
+3
View File
@@ -121,8 +121,11 @@ impl SegmentLogSink {
matches!(
entry,
LogEntry::SegmentStart { .. }
| LogEntry::AnnotatedSegmentStart { .. }
| LogEntry::UserInput { .. }
| LogEntry::AnnotatedUserInput { .. }
| LogEntry::SystemItem { .. }
| LogEntry::AnnotatedSystemItem { .. }
| LogEntry::Invoke { .. }
)
}
+144 -20
View File
@@ -6,7 +6,8 @@
use std::sync::Arc;
use agen::{Item, Role};
use crate::session_history::{SessionHistoryMetadata, WorkerHistoryProvenance};
use agen::{HistoryEntry, Item, Role};
use serde::{Deserialize, Serialize};
const DEFAULT_SEARCH_LIMIT: usize = 20;
@@ -21,14 +22,21 @@ const OVERVIEW_ANCHOR_STRIDE: usize = 8;
pub(crate) struct SessionEntryRef(String);
impl SessionEntryRef {
pub(crate) fn new(source_index: usize) -> Self {
Self(format!("E{source_index:08}"))
pub(crate) fn from_history_entry_id(entry_id: &crate::SessionHistoryEntryId) -> Self {
Self(format!("E{}", entry_id.0))
}
pub(crate) fn parse(value: &str) -> Option<Self> {
let reference = Self(value.to_string());
reference.source_index()?;
Some(reference)
let suffix = value.strip_prefix('E')?;
if suffix.is_empty()
|| suffix.len() > 64
|| !suffix
.bytes()
.all(|byte| byte.is_ascii_alphanumeric() || byte == b'-')
{
return None;
}
Some(Self(value.to_string()))
}
pub(crate) fn as_str(&self) -> &str {
@@ -97,6 +105,7 @@ impl ToolPart {
#[derive(Debug, Clone)]
pub(crate) struct OverviewItem {
pub id: SessionEntryRef,
pub origin: WorkerHistoryProvenance,
pub entry_range: [u64; 2],
pub kind: ReferenceKind,
pub label: String,
@@ -107,6 +116,7 @@ pub(crate) struct OverviewItem {
#[derive(Debug, Clone)]
pub(crate) struct ReferenceEntry {
pub id: SessionEntryRef,
pub origin: WorkerHistoryProvenance,
pub entry_range: [u64; 2],
pub kind: ReferenceKind,
pub tool_part: Option<ToolPart>,
@@ -132,6 +142,7 @@ pub(crate) struct SearchOptions {
#[derive(Debug, Clone)]
pub(crate) struct SearchHit {
pub id: SessionEntryRef,
pub origin: WorkerHistoryProvenance,
pub kind: ReferenceKind,
pub tool_part: Option<ToolPart>,
pub tool_name: Option<String>,
@@ -177,6 +188,7 @@ impl Default for ReadOptions {
#[derive(Debug, Clone)]
pub(crate) struct ReadEntry {
pub id: SessionEntryRef,
pub origin: WorkerHistoryProvenance,
pub kind: ReferenceKind,
pub tool_part: Option<ToolPart>,
pub tool_name: Option<String>,
@@ -195,6 +207,7 @@ pub(crate) struct ReadResult {
pub(crate) struct SessionEntryEvidence {
pub segment_id: String,
pub entry_ref: SessionEntryRef,
pub origin: WorkerHistoryProvenance,
pub entry_range: [u64; 2],
pub kind: ReferenceKind,
pub tool_part: Option<ToolPart>,
@@ -206,26 +219,42 @@ pub(crate) struct SessionEntryEvidence {
#[derive(Debug, Clone)]
pub(crate) struct SessionCapture {
segment_id: String,
items: Arc<Vec<Item>>,
entries: Arc<Vec<HistoryEntry<SessionHistoryMetadata>>>,
overview: Vec<OverviewItem>,
index: Vec<ReferenceEntry>,
}
impl SessionCapture {
pub(crate) fn new(segment_id: impl Into<String>, items: Vec<Item>) -> Self {
let entries = items
.into_iter()
.enumerate()
.map(|(index, item)| {
let mut metadata = SessionHistoryMetadata::legacy_unknown();
metadata.entry_id =
session_store::LoggedSessionHistoryEntryId(format!("{index:08}"));
HistoryEntry::new(item, metadata)
})
.collect();
Self::from_history_entries(segment_id, entries)
}
pub(crate) fn from_history_entries(
segment_id: impl Into<String>,
entries: Vec<HistoryEntry<SessionHistoryMetadata>>,
) -> Self {
let segment_id = segment_id.into();
let items = Arc::new(items);
let entries = Arc::new(entries);
let mut overview = Vec::new();
let mut index = Vec::new();
for (idx, item) in items.iter().enumerate() {
for (idx, entry) in entries.iter().enumerate() {
let item = &entry.item;
let entry_range = [idx as u64, idx as u64];
match item {
Item::Message { role, content, .. } => {
let kind = match role {
Role::User => ReferenceKind::User,
Role::Assistant => ReferenceKind::Assistant,
Role::System => continue,
let Some(kind) = message_reference_kind(&entry.annotation.origin, role) else {
continue;
};
let text = content
.iter()
@@ -234,9 +263,10 @@ impl SessionCapture {
.join("");
let label = format!("{} message", kind.as_str());
let summary = truncate_chars(&text, 240);
let id = SessionEntryRef::new(idx);
let id = SessionEntryRef::from_history_entry_id(&entry.annotation.entry_id);
index.push(ReferenceEntry {
id: id.clone(),
origin: entry.annotation.origin.clone(),
entry_range,
kind,
tool_part: None,
@@ -248,6 +278,7 @@ impl SessionCapture {
if matches!(kind, ReferenceKind::User | ReferenceKind::Assistant) {
overview.push(OverviewItem {
id: id.clone(),
origin: entry.annotation.origin.clone(),
entry_range,
kind,
label,
@@ -261,7 +292,8 @@ impl SessionCapture {
} => {
let text = format!("{name}\n{arguments}");
index.push(ReferenceEntry {
id: SessionEntryRef::new(idx),
id: SessionEntryRef::from_history_entry_id(&entry.annotation.entry_id),
origin: entry.annotation.origin.clone(),
entry_range,
kind: ReferenceKind::Tool,
tool_part: Some(ToolPart::Input),
@@ -287,7 +319,8 @@ impl SessionCapture {
content.as_deref().unwrap_or_default(),
);
index.push(ReferenceEntry {
id: SessionEntryRef::new(idx),
id: SessionEntryRef::from_history_entry_id(&entry.annotation.entry_id),
origin: entry.annotation.origin.clone(),
entry_range,
kind: ReferenceKind::Tool,
tool_part: Some(ToolPart::Output),
@@ -327,7 +360,7 @@ impl SessionCapture {
Self {
segment_id,
items,
entries,
overview,
index,
}
@@ -337,6 +370,14 @@ impl SessionCapture {
&self.overview
}
pub(crate) fn source_index_for_ref(&self, reference: &SessionEntryRef) -> Option<u64> {
self.index
.iter()
.find(|entry| entry.id == *reference)
.map(|entry| entry.entry_range[0])
.or_else(|| reference.source_index())
}
pub(crate) fn search(&self, options: &SearchOptions) -> Vec<SearchHit> {
let query = options.query.trim().to_lowercase();
let limit = options
@@ -347,12 +388,12 @@ impl SessionCapture {
let min_entry_index = options
.from
.as_ref()
.and_then(SessionEntryRef::source_index)
.and_then(|reference| self.source_index_for_ref(reference))
.unwrap_or_else(|| options.min_entry_index.unwrap_or(0));
let max_entry_index = options
.through
.as_ref()
.and_then(SessionEntryRef::source_index)
.and_then(|reference| self.source_index_for_ref(reference))
.unwrap_or(u64::MAX);
let mut skipped = 0usize;
let mut hits = Vec::new();
@@ -391,6 +432,7 @@ impl SessionCapture {
}
hits.push(SearchHit {
id: entry.id.clone(),
origin: entry.origin.clone(),
kind: entry.kind,
tool_part: entry.tool_part,
tool_name: entry.tool_name.clone(),
@@ -442,13 +484,18 @@ impl SessionCapture {
}
}
}
let Some(item) = self.items.get(entry.entry_range[0] as usize) else {
let Some(item) = self
.entries
.get(entry.entry_range[0] as usize)
.map(|entry| &entry.item)
else {
continue;
};
let text = render_item(item, entry, options.detail, max_bytes.saturating_sub(bytes));
bytes = bytes.saturating_add(text.len());
entries.push(ReadEntry {
id: entry.id.clone(),
origin: entry.origin.clone(),
kind: entry.kind,
tool_part: entry.tool_part,
tool_name: entry.tool_name.clone(),
@@ -485,6 +532,7 @@ impl SessionCapture {
Some(SessionEntryEvidence {
segment_id: self.segment_id.clone(),
entry_ref: entry.id.clone(),
origin: entry.origin.clone(),
entry_range: entry.entry_range,
kind: entry.kind,
tool_part: entry.tool_part,
@@ -495,6 +543,28 @@ impl SessionCapture {
}
}
fn message_reference_kind(
origin: &WorkerHistoryProvenance,
provider_role: &Role,
) -> Option<ReferenceKind> {
match origin {
WorkerHistoryProvenance::HumanInput { .. }
| WorkerHistoryProvenance::WorkerInput { .. } => Some(ReferenceKind::User),
WorkerHistoryProvenance::ModelOutput { .. } => Some(ReferenceKind::Assistant),
WorkerHistoryProvenance::ToolOutput { .. } => Some(ReferenceKind::Tool),
WorkerHistoryProvenance::LegacyUnknown => match provider_role {
Role::User => Some(ReferenceKind::User),
Role::Assistant => Some(ReferenceKind::Assistant),
Role::System => None,
},
// Flow/backend/system content remains out of the observation surface
// even when represented with a provider user/system role.
WorkerHistoryProvenance::FlowInstruction { .. }
| WorkerHistoryProvenance::BackendInstruction { .. }
| WorkerHistoryProvenance::DerivedSummary => None,
}
}
fn render_item(
item: &Item,
entry: &ReferenceEntry,
@@ -563,6 +633,60 @@ fn truncate_chars(text: &str, max_chars: usize) -> String {
mod tests {
use super::*;
#[test]
fn flow_user_role_is_excluded_while_explicit_human_origin_remains_evidence() {
let entries = vec![
crate::session_history::history_entry(
Item::user_message("trusted flow instruction"),
WorkerHistoryProvenance::FlowInstruction {
selector: "builtin:coder-review".into(),
definition_id: "coder-review".into(),
definition_revision: 3,
instance_id: "instance".into(),
state_id: "implement".into(),
},
),
crate::session_history::history_entry(
Item::user_message("remember my preference"),
WorkerHistoryProvenance::HumanInput {
account_id: "account-1".into(),
},
),
];
let capture = SessionCapture::from_history_entries("segment", entries);
let overview = capture.overview();
assert_eq!(overview.len(), 1);
assert!(matches!(
overview[0].origin,
WorkerHistoryProvenance::HumanInput { .. }
));
let evidence = capture.evidence_for(overview[0].id.as_str()).unwrap();
assert!(evidence.excerpt.ends_with("remember my preference"));
assert!(matches!(
evidence.origin,
WorkerHistoryProvenance::HumanInput { .. }
));
}
#[test]
fn stable_logical_ref_survives_retention_and_restore_projection() {
let retained = crate::session_history::history_entry(
Item::assistant_message("retained"),
WorkerHistoryProvenance::ModelOutput {
worker: crate::session_history::worker_subject(Default::default()),
},
);
let expected_ref = SessionEntryRef::from_history_entry_id(&retained.annotation.entry_id);
let before = SessionCapture::from_history_entries("old", vec![retained.clone()]);
let after = SessionCapture::from_history_entries("new", vec![retained]);
assert_eq!(before.overview()[0].id, expected_ref);
assert_eq!(after.overview()[0].id, expected_ref);
assert_eq!(
after.evidence_for(expected_ref.as_str()).unwrap().entry_ref,
expected_ref
);
}
#[test]
fn overview_contains_user_and_assistant_only() {
let view = SessionCapture::new(
+219
View File
@@ -0,0 +1,219 @@
//! Restore-authoritative metadata for model-visible Worker history.
//!
//! Agen transports this annotation without interpreting it. Session Log v2
//! stores each item and metadata in one typed record; legacy records are
//! retained only as explicit `LegacyUnknown` entries.
use agen::{HistoryEntry, Item};
use protocol::Segment;
use session_store::{
LogEntry, LoggedHistoryDerivation, LoggedHistoryEntry, LoggedSessionHistoryEntryId,
LoggedSessionHistoryMetadata, LoggedSessionHistoryOrigin, LoggedWorkerSubject, SegmentId,
SessionId,
};
pub type SessionHistoryEntryId = LoggedSessionHistoryEntryId;
pub type SessionHistoryMetadata = LoggedSessionHistoryMetadata;
pub type WorkerHistoryProvenance = LoggedSessionHistoryOrigin;
pub type SessionHistoryDerivation = LoggedHistoryDerivation;
pub type WorkerSubjectSnapshot = LoggedWorkerSubject;
pub(crate) fn worker_subject(session_id: SessionId) -> WorkerSubjectSnapshot {
WorkerSubjectSnapshot {
workspace_id: None,
runtime_id: None,
worker_id: session_id.to_string(),
}
}
pub(crate) fn metadata(
origin: WorkerHistoryProvenance,
derivation: Option<SessionHistoryDerivation>,
) -> SessionHistoryMetadata {
SessionHistoryMetadata {
entry_id: SessionHistoryEntryId::new(),
origin,
derivation,
}
}
pub(crate) fn history_entry(
item: Item,
origin: WorkerHistoryProvenance,
) -> HistoryEntry<SessionHistoryMetadata> {
HistoryEntry::new(item, metadata(origin, None))
}
pub(crate) fn to_logged_history_entry(
entry: &HistoryEntry<SessionHistoryMetadata>,
) -> LoggedHistoryEntry {
LoggedHistoryEntry {
item: entry.item.clone().into(),
metadata: entry.annotation.clone(),
}
}
fn legacy_entry(item: Item) -> HistoryEntry<SessionHistoryMetadata> {
HistoryEntry::new(item, SessionHistoryMetadata::legacy_unknown())
}
fn from_logged(entry: &LoggedHistoryEntry) -> HistoryEntry<SessionHistoryMetadata> {
HistoryEntry::new(Item::from(entry.item.clone()), entry.metadata.clone())
}
/// Rebuild typed Worker history directly from the append-only Session Log.
/// Missing legacy metadata is never inferred from role or plaintext.
pub(crate) fn restore_history_entries(
_session_id: SessionId,
_segment_id: SegmentId,
entries: &[LogEntry],
) -> Result<Vec<HistoryEntry<SessionHistoryMetadata>>, String> {
let mut history = Vec::new();
for entry in entries {
match entry {
LogEntry::AnnotatedSegmentStart { history: seed, .. } => {
history = seed.iter().map(from_logged).collect();
}
LogEntry::SegmentStart { history: seed, .. } => {
history = seed
.iter()
.cloned()
.map(Item::from)
.map(legacy_entry)
.collect();
}
LogEntry::AnnotatedUserInput { history: input, .. } => {
history.extend(input.iter().map(from_logged))
}
LogEntry::UserInput { segments, .. } => history.push(legacy_entry(Item::user_message(
Segment::flatten_to_text(segments),
))),
LogEntry::AnnotatedAssistantItem { entry, .. }
| LogEntry::AnnotatedToolResult { entry, .. } => history.push(from_logged(entry)),
LogEntry::AssistantItem { item, .. } | LogEntry::ToolResult { item, .. } => {
history.push(legacy_entry(Item::from(item.clone())));
}
LogEntry::AnnotatedSystemItem { entry, .. } => history.push(HistoryEntry::new(
entry.item.to_history_item(),
entry.metadata.clone(),
)),
LogEntry::SystemItem { item, .. } => {
history.push(legacy_entry(item.to_history_item()));
}
_ => {}
}
}
Ok(history)
}
#[cfg(test)]
mod tests {
use super::*;
use agen::llm_client::RequestConfig;
use session_store::LogEntry;
#[test]
fn legacy_user_role_is_not_inferred_as_human_authority() {
let entries = vec![LogEntry::UserInput {
ts: 1,
segments: vec![Segment::text("legacy")],
extensions: Vec::new(),
}];
let restored =
restore_history_entries(SessionId::now_v7(), SegmentId::now_v7(), &entries).unwrap();
assert!(matches!(
restored[0].annotation.origin,
WorkerHistoryProvenance::LegacyUnknown
));
}
#[test]
fn typed_flow_and_unknown_caller_input_round_trip_without_role_inference() {
let session_id = SessionId::now_v7();
let projected = vec![
history_entry(
Item::user_message("flow instructions"),
WorkerHistoryProvenance::FlowInstruction {
selector: "builtin:coder-review".to_string(),
definition_id: "coder-review".to_string(),
definition_revision: 7,
instance_id: "flow-instance".to_string(),
state_id: "implement".to_string(),
},
),
history_entry(
Item::user_message("implement"),
WorkerHistoryProvenance::LegacyUnknown,
),
];
let entries = vec![
LogEntry::AnnotatedSegmentStart {
ts: 0,
session_id,
system_prompt: None,
config: RequestConfig::default(),
history: Vec::new(),
forked_from: None,
compacted_from: None,
},
LogEntry::AnnotatedUserInput {
ts: 1,
segments: vec![
Segment::Flow {
selector: "builtin:coder-review".to_string(),
},
Segment::text("implement"),
],
extensions: Vec::new(),
history: projected.iter().map(to_logged_history_entry).collect(),
},
];
let restored = restore_history_entries(session_id, SegmentId::now_v7(), &entries).unwrap();
assert_eq!(restored, projected);
}
#[test]
fn annotated_restore_preserves_logical_ids_across_reboot() {
let session_id = SessionId::now_v7();
let entry = history_entry(
Item::assistant_message("persisted"),
WorkerHistoryProvenance::ModelOutput {
worker: worker_subject(session_id),
},
);
let log = vec![LogEntry::AnnotatedSegmentStart {
ts: 0,
session_id,
system_prompt: None,
config: RequestConfig::default(),
history: vec![to_logged_history_entry(&entry)],
forked_from: None,
compacted_from: None,
}];
let first = restore_history_entries(session_id, SegmentId::now_v7(), &log).unwrap();
let second = restore_history_entries(session_id, SegmentId::now_v7(), &log).unwrap();
assert_eq!(first[0].annotation.entry_id, entry.annotation.entry_id);
assert_eq!(second[0].annotation.entry_id, entry.annotation.entry_id);
}
#[test]
fn compacted_derivation_uses_stable_logical_entry_ids() {
let source = history_entry(
Item::user_message("source"),
WorkerHistoryProvenance::LegacyUnknown,
);
let summary = HistoryEntry::new(
Item::system_message("summary"),
metadata(
WorkerHistoryProvenance::DerivedSummary,
Some(SessionHistoryDerivation {
sources: vec![source.annotation.entry_id.clone()],
}),
),
);
assert_eq!(
summary.annotation.derivation.unwrap().sources,
vec![source.annotation.entry_id]
);
}
}
+9 -6
View File
@@ -499,7 +499,10 @@ impl Tool for SubWorkerSpawnTool {
InternalWorkerVisibility::ParentClient,
Some(child_registry.clone()),
Some(Arc::new(move |status| {
if status == InternalWorkerSessionStatus::Failed {
if matches!(
status,
InternalWorkerSessionStatus::Failed | InternalWorkerSessionStatus::Stopped
) {
if let Some(registry) = registry.upgrade() {
if let Err(error) = registry.reclaim_internal_scope(&child_name) {
tracing::warn!(
@@ -1249,7 +1252,7 @@ extract_threshold = 4000
)
.await
.unwrap();
assert!(first_capture.items.iter().any(|item| {
assert!(first_capture.entries.iter().map(|entry| &entry.item).any(|item| {
matches!(item, Item::Message { role: Role::Assistant, content, .. } if content.iter().any(|part| matches!(part, ContentPart::Text { text } if text.contains("reviewed"))))
}));
@@ -1271,7 +1274,7 @@ extract_threshold = 4000
)
.await
.unwrap();
assert!(latest_capture.items.len() > first_capture.items.len());
assert!(latest_capture.entries.len() > first_capture.entries.len());
fail_requests.store(true, Ordering::SeqCst);
send.execute(
@@ -1282,16 +1285,16 @@ extract_threshold = 4000
.unwrap();
assert_eq!(
record.session.wait_until_idle().await,
InternalWorkerSessionStatus::Failed
InternalWorkerSessionStatus::Stopped
);
assert_eq!(calls.load(Ordering::SeqCst), 3);
assert!(
spawner_scope.snapshot().is_writable(&workspace_root),
"Failed terminal child must release its delegated Workdir session"
"Stopped terminal child must release its delegated Workdir session"
);
assert!(
!record.workdir_delegation.is_active(),
"failed child must revoke cloned scoped sessions"
"stopped child must revoke cloned scoped sessions"
);
assert!(registry.get_internal("reviewer-child").is_some());
+1219 -273
View File
File diff suppressed because it is too large Load Diff
+35 -23
View File
@@ -163,7 +163,8 @@ async fn make_worker_with_manifest(
let scope = worker::Scope::writable(&pwd).unwrap();
std::mem::forget(pwd_tmp);
let worker = Engine::new(client);
let worker =
Engine::<_, agen::state::Mutable, worker::SessionHistoryMetadata>::new_annotated(client);
let mut worker = Worker::new(
manifest,
worker,
@@ -204,28 +205,34 @@ fn system_texts_in_sink_session_start(
) -> Vec<String> {
let (entries, _rx) = worker.sink().subscribe_with_snapshot();
for entry in entries.into_iter().rev() {
if let session_store::LogEntry::SegmentStart { history, .. } = entry {
return history
let history = match entry {
session_store::LogEntry::AnnotatedSegmentStart { history, .. } => history
.into_iter()
.filter_map(|logged| {
let item: Item = logged.into();
match item {
Item::Message {
role: agen::Role::System,
content,
..
} => Some(
content
.iter()
.map(|p| p.as_text().to_owned())
.collect::<Vec<_>>()
.join(""),
),
_ => None,
}
})
.collect();
}
.map(|entry| entry.item)
.collect::<Vec<_>>(),
session_store::LogEntry::SegmentStart { history, .. } => history,
_ => continue,
};
return history
.into_iter()
.filter_map(|logged| {
let item: Item = logged.into();
match item {
Item::Message {
role: agen::Role::System,
content,
..
} => Some(
content
.iter()
.map(|p| p.as_text().to_owned())
.collect::<Vec<_>>()
.join(""),
),
_ => None,
}
})
.collect();
}
Vec::new()
}
@@ -337,7 +344,12 @@ permission = "write"
// New segment records forked_from pointing at the source.
let new_entries = store.read_all(session_id, new_segment_id).unwrap();
match &new_entries[0] {
LogEntry::SegmentStart {
LogEntry::AnnotatedSegmentStart {
session_id: seg_session,
forked_from: Some(origin),
..
}
| LogEntry::SegmentStart {
session_id: seg_session,
forked_from: Some(origin),
..
+74 -31
View File
@@ -32,16 +32,29 @@ fn history_from_sink(handle: &WorkerHandle) -> Vec<Item> {
let mut items = Vec::new();
for entry in entries {
match entry {
LogEntry::AnnotatedSegmentStart { history, .. } => {
items.extend(history.into_iter().map(|entry| Item::from(entry.item)));
}
LogEntry::SegmentStart { history, .. } => {
items.extend(history.into_iter().map(Item::from));
}
LogEntry::AnnotatedUserInput { history, .. } => {
items.extend(history.into_iter().map(|entry| Item::from(entry.item)));
}
LogEntry::UserInput { segments, .. } => {
let text = protocol::Segment::flatten_to_text(&segments);
items.push(Item::user_message(text));
}
LogEntry::AnnotatedAssistantItem { entry, .. }
| LogEntry::AnnotatedToolResult { entry, .. } => {
items.push(Item::from(entry.item));
}
LogEntry::AssistantItem { item, .. } | LogEntry::ToolResult { item, .. } => {
items.push(Item::from(item));
}
LogEntry::AnnotatedSystemItem { entry, .. } => {
items.push(entry.item.to_history_item());
}
LogEntry::SystemItem { item, .. } => {
items.push(item.to_history_item());
}
@@ -51,6 +64,14 @@ fn history_from_sink(handle: &WorkerHandle) -> Vec<Item> {
items
}
fn system_item(entry: &LogEntry) -> Option<&session_store::SystemItem> {
match entry {
LogEntry::AnnotatedSystemItem { entry, .. } => Some(&entry.item),
LogEntry::SystemItem { item, .. } => Some(item),
_ => None,
}
}
// ---------------------------------------------------------------------------
// Mock LLM Client
// ---------------------------------------------------------------------------
@@ -192,7 +213,8 @@ async fn make_worker_with_pwd_and_manifest(
let scope = manifest::Scope::writable(&pwd).unwrap();
std::mem::forget(pwd_tmp);
let worker = Engine::new(client);
let worker =
Engine::<_, agen::state::Mutable, worker::SessionHistoryMetadata>::new_annotated(client);
let authority = WorkerFilesystemAuthority::local(pwd.clone(), pwd.clone());
let worker = Worker::new(
manifest,
@@ -784,13 +806,30 @@ async fn snapshot_includes_user_input_for_in_flight_turn() {
let client = MockClient::sequential(vec![MockResponse::Hang(simple_text_events())]);
let worker = make_worker(client).await;
let handle = spawn_controller(worker).await;
let mut events = handle.subscribe();
handle
.send(Method::run_text("hello in-flight"))
.await
.unwrap();
wait_for_status(&handle, WorkerStatus::Running).await;
tokio::time::timeout(std::time::Duration::from_secs(2), async {
loop {
if matches!(
events.recv().await,
Ok(Event::Status {
status: WorkerStatus::Running,
})
) {
break;
}
}
})
.await
.expect("running status event");
// The Running event is the in-flight visibility fence: the committed
// annotated input must already be available to an immediately attaching
// subscriber rather than racing behind this status transition.
let stream = tokio::net::UnixStream::connect(handle.runtime_dir.socket_path())
.await
.unwrap();
@@ -804,10 +843,12 @@ async fn snapshot_includes_user_input_for_in_flight_turn() {
// Walk the entries, find a `LogEntry::UserInput` and
// confirm its segments flatten to our submitted text.
let mut found = false;
for value in entries {
for value in &entries {
let entry: session_store::LogEntry =
serde_json::from_value(value).expect("LogEntry deserialise");
if let session_store::LogEntry::UserInput { segments, .. } = entry {
serde_json::from_value(value.clone()).expect("LogEntry deserialise");
if let session_store::LogEntry::UserInput { segments, .. }
| session_store::LogEntry::AnnotatedUserInput { segments, .. } = entry
{
let text = protocol::Segment::flatten_to_text(&segments);
if text == "hello in-flight" {
found = true;
@@ -815,7 +856,10 @@ async fn snapshot_includes_user_input_for_in_flight_turn() {
}
}
}
assert!(found, "snapshot must carry the in-flight UserInput entry");
assert!(
found,
"snapshot must carry the in-flight UserInput entry: {entries:?}"
);
return;
}
Event::Alert(_) => continue,
@@ -1086,7 +1130,7 @@ async fn run_with_paste_segment_inlines_content_and_emits_typed_user_message() {
_ => {}
},
entry = entry_rx.recv() => match entry {
Ok(session_store::LogEntry::UserInput { segments, .. }) => {
Ok(session_store::LogEntry::UserInput { segments, .. } | session_store::LogEntry::AnnotatedUserInput { segments, .. }) => {
user_input_segments = Some(segments);
if saw_turn_end {
break;
@@ -1317,11 +1361,8 @@ async fn notify_while_idle_auto_starts_turn_and_injects_system_message() {
let (entries, _) = handle.sink.subscribe_with_snapshot();
let saw_notify_in_mirror = entries.iter().any(|e| {
matches!(
e,
session_store::LogEntry::SystemItem {
item: session_store::SystemItem::Notification { message, .. },
..
} if message == "turn finished"
system_item(e),
Some(session_store::SystemItem::Notification { message, .. }) if message == "turn finished"
)
});
assert!(
@@ -1463,14 +1504,11 @@ async fn worker_event_turn_ended_while_idle_auto_starts_turn_and_injects_system_
let (entries, _) = handle.sink.subscribe_with_snapshot();
let saw_worker_event_in_mirror = entries.iter().any(|e| {
matches!(
e,
session_store::LogEntry::SystemItem {
item: session_store::SystemItem::WorkerEvent {
event: protocol::WorkerEvent::TurnEnded { worker_name },
..
},
system_item(e),
Some(session_store::SystemItem::WorkerEvent {
event: protocol::WorkerEvent::TurnEnded { worker_name },
..
} if worker_name == "child"
}) if worker_name == "child"
)
});
assert!(
@@ -1552,14 +1590,11 @@ async fn worker_event_scope_sub_delegated_while_idle_stays_control_plane_only()
let (entries, _) = handle.sink.subscribe_with_snapshot();
let saw_scope_event_in_mirror = entries.iter().any(|entry| {
matches!(
entry,
session_store::LogEntry::SystemItem {
item: session_store::SystemItem::WorkerEvent {
event: protocol::WorkerEvent::ScopeSubDelegated { .. },
..
},
system_item(entry),
Some(session_store::SystemItem::WorkerEvent {
event: protocol::WorkerEvent::ScopeSubDelegated { .. },
..
}
})
)
});
assert!(
@@ -2134,9 +2169,13 @@ async fn paused_then_run_closes_orphan_tool_use_for_next_request() {
for item in items {
match item {
agen::Item::ToolResult {
call_id, summary, ..
call_id,
summary,
disposition,
..
} if call_id == "call_orphan" => {
assert_eq!(summary, "[Interrupted by user]");
assert_eq!(summary, "Tool execution outcome unknown");
assert_eq!(*disposition, agen::ToolResultDisposition::OutcomeUnknown);
saw_synthetic_tool_result = true;
}
agen::Item::Message { role, content, .. } if *role == agen::Role::System => {
@@ -2327,8 +2366,11 @@ async fn paused_cancel_abandons_resume_and_next_input_is_fresh_run() {
assert!(
items.iter().any(|item| matches!(
item,
agen::Item::ToolResult { call_id, summary, .. }
if call_id == "call_cancelled" && summary == "[Interrupted by user]"
agen::Item::ToolResult {
call_id,
disposition: agen::ToolResultDisposition::OutcomeUnknown,
..
} if call_id == "call_cancelled"
)),
"paused cancel should close orphan tool_use before future requests: {items:?}"
);
@@ -2373,7 +2415,8 @@ async fn snapshot_contains_user_input(handle: &WorkerHandle, needle: &str) -> bo
let entry: session_store::LogEntry =
serde_json::from_value(value).expect("LogEntry deserialise");
match entry {
session_store::LogEntry::UserInput { segments, .. } => {
session_store::LogEntry::UserInput { segments, .. }
| session_store::LogEntry::AnnotatedUserInput { segments, .. } => {
protocol::Segment::flatten_to_text(&segments).contains(needle)
}
_ => false,
+6 -3
View File
@@ -188,7 +188,8 @@ async fn make_worker(
let pwd = pwd_tmp.path().to_path_buf();
let scope = worker::Scope::writable(&pwd).unwrap();
let mut worker = Engine::new(client);
let mut worker =
Engine::<_, agen::state::Mutable, worker::SessionHistoryMetadata>::new_annotated(client);
worker.register_tool(big_content_tool_definition(tool_name));
let worker = Worker::new(
@@ -460,7 +461,8 @@ async fn metric_write_failure_emits_warn_alert_and_does_not_abort_run() {
// protected token budget covers the only user message). That is enough to drive
// the failure path: at least one metric attempts to write.
let client = MockClient::new(vec![text_response_with_cache("hi", 0, 0)]);
let worker = Engine::new(client);
let worker =
Engine::<_, agen::state::Mutable, worker::SessionHistoryMetadata>::new_annotated(client);
let mut worker = Worker::new(
manifest,
worker,
@@ -536,7 +538,8 @@ permission = "write"
let pwd_tmp = tempfile::tempdir().unwrap();
let pwd = pwd_tmp.path().to_path_buf();
let scope = worker::Scope::writable(&pwd).unwrap();
let worker = Engine::new(client);
let worker =
Engine::<_, agen::state::Mutable, worker::SessionHistoryMetadata>::new_annotated(client);
let mut worker = Worker::new(
manifest,
worker,
@@ -130,7 +130,8 @@ async fn make_worker_with_body(
EffectivePromptCatalog::new(templates, 1, "test-schema", "test-toolchain").unwrap();
let loader = PromptCatalogSource::builtins_only().with_effective_catalog(projection);
let worker = Engine::new(client);
let worker =
Engine::<_, agen::state::Mutable, worker::SessionHistoryMetadata>::new_annotated(client);
let mut worker = Worker::new(
manifest,
worker,