Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 16 additions & 8 deletions src/crates/core/src/agentic/tools/implementations/work_tool.rs
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@ impl Tool for WorkTool {
}

async fn description(&self) -> CoreResult<String> {
Ok("Run and manage specialist Work through one control-plane tool. action=start atomically creates and launches an Agent WorkSession and returns its work_id; action=continue sends more instructions to existing Work; action=status reads progress and results; action=control changes lifecycle state. Always target Work by the work_id from start, never by a session id.".to_string())
Ok("Run and manage specialist Work through one control-plane tool. action=start atomically creates and launches an Agent WorkSession and returns its work_id; action=continue sends more instructions to existing Work; action=status reads progress and results; action=control changes lifecycle state; action=reclassify changes kind or topic attachment. Always target Work by the work_id from start, never by a session id. System-managed works are immutable.".to_string())
}

fn input_schema(&self) -> Value {
Expand All @@ -33,17 +33,25 @@ impl Tool for WorkTool {
"properties": {
"action": {
"type": "string",
"enum": ["start", "continue", "status", "control"],
"description": "start: create and launch new Work. continue: add instructions to existing Work. status: read progress and results. control: change lifecycle state."
"enum": ["start", "continue", "status", "control", "reclassify"],
"description": "start: create and launch new Work. continue: add instructions to existing Work. status: read progress and results. control: change lifecycle state. reclassify: change kind or topic attachment."
},
"work_id": {
"type": "string",
"description": "The Work to target. Required for continue and control, and for status of one specific Work. This is the work_id returned by start, not a session id."
"description": "The Work to target. Required for continue, control, reclassify, and for status of one specific Work. This is the work_id returned by start, not a session id."
},
"kind": {
"type": "string",
"enum": ["one_shot", "multi_step", "long_running_session"],
"description": "start only. multi_step (default) for normal multi-step execution; one_shot for a single self-contained task; long_running_session for ongoing work."
"enum": ["one_shot", "multi_step", "long_running_session", "tracking", "topic", "recurring", "app_workflow"],
"description": "start or reclassify. multi_step (default) for normal multi-step execution; one_shot for a single self-contained task; long_running_session/tracking for ongoing user work; topic for a theme container; recurring for user cadence work; app_workflow when an app subject is attached."
},
"topic_work_id": {
"type": "string",
"description": "Optional Topic work_id to attach this Work under. Topic target must have kind=topic."
},
"clear_topic_work_id": {
"type": "boolean",
"description": "reclassify/start only. Clear topic attachment when true."
},
"title": {
"type": "string",
Expand Down Expand Up @@ -293,8 +301,8 @@ mod tests {
let actions = schema["properties"]["action"]["enum"]
.as_array()
.expect("action enum");
assert_eq!(actions.len(), 4);
for action in ["start", "continue", "status", "control"] {
assert_eq!(actions.len(), 5);
for action in ["start", "continue", "status", "control", "reclassify"] {
assert!(
actions.iter().any(|value| value.as_str() == Some(action)),
"missing action {action}"
Expand Down
55 changes: 49 additions & 6 deletions src/crates/core/src/agentic_os/tools/work.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,8 +3,9 @@ use serde_json::{json, Value};

use crate::agentic_os::work::{
AdvanceWorkRequest, ControlWorkAction, ControlWorkRequest, PrimarySurfacePolicy,
StartWorkRequest, WorkAssignmentKind, WorkAssignmentRef, WorkId, WorkKind, WorkOwnerRef,
WorkProjection, WorkRecord, WorkScope, WorkService, WorkStatus, WorkSubject, WorkVisibility,
ReclassifyWorkRequest, StartWorkRequest, UpdateWorkRequest, WorkAssignmentKind,
WorkAssignmentRef, WorkId, WorkKind, WorkOwnerRef, WorkProjection, WorkRecord, WorkScope,
WorkService, WorkStatus, WorkSubject, WorkVisibility,
};
use crate::error::{CoreError, CoreResult};

Expand All @@ -15,6 +16,7 @@ pub enum WorkAction {
Continue,
Status,
Control,
Reclassify,
}

#[derive(Debug, Deserialize)]
Expand Down Expand Up @@ -52,6 +54,10 @@ pub struct WorkInput {
pub control_action: Option<ControlWorkAction>,
#[serde(default)]
pub include_archived: Option<bool>,
#[serde(default)]
pub topic_work_id: Option<WorkId>,
#[serde(default)]
pub clear_topic_work_id: Option<bool>,
#[serde(skip)]
pub owner: Option<WorkOwnerRef>,
}
Expand All @@ -62,6 +68,7 @@ pub async fn handle(service: &WorkService, input: WorkInput) -> CoreResult<Value
WorkAction::Continue => continue_work(service, input).await,
WorkAction::Status => status_work(service, input).await,
WorkAction::Control => control_work(service, input).await,
WorkAction::Reclassify => reclassify_work(service, input).await,
}
}

Expand Down Expand Up @@ -89,18 +96,32 @@ async fn start_work(service: &WorkService, input: WorkInput) -> CoreResult<Value
})
.await?;

let mut work = response.work;
if input.topic_work_id.is_some() || input.clear_topic_work_id.unwrap_or(false) {
work = service
.update(
&work.id,
UpdateWorkRequest {
topic_work_id: input.topic_work_id,
clear_topic_work_id: input.clear_topic_work_id.unwrap_or(false),
..UpdateWorkRequest::default()
},
)
.await?;
}

Ok(json!({
"action": "start",
"work_id": response.work.id,
"status": response.work.status,
"surface": response.work.primary_surface,
"work_id": work.id,
"status": work.status,
"surface": work.primary_surface,
"execution": {
"kind": "agent_session_run",
"execution_binding_id": response.execution_binding_id,
"turn_id": response.turn_id,
"started": response.started,
},
"work": response.work,
"work": work,
}))
}

Expand Down Expand Up @@ -173,6 +194,28 @@ async fn control_work(service: &WorkService, input: WorkInput) -> CoreResult<Val
}))
}

async fn reclassify_work(service: &WorkService, input: WorkInput) -> CoreResult<Value> {
let work = service
.reclassify(ReclassifyWorkRequest {
work_id: required_work_id(input.work_id, "reclassify")?,
kind: input
.kind
.ok_or_else(|| CoreError::validation("kind is required for action=reclassify"))?,
topic_work_id: input.topic_work_id,
clear_topic_work_id: input.clear_topic_work_id.unwrap_or(false),
})
.await?;

Ok(json!({
"action": "reclassify",
"work_id": work.id,
"kind": work.kind,
"topic_work_id": work.topic_work_id,
"status": work.status,
"work": work,
}))
}

fn work_executor_to_assignment(executor: WorkExecutorInput) -> CoreResult<WorkAssignmentRef> {
match executor.kind.unwrap_or(WorkExecutorKind::Agent) {
WorkExecutorKind::Agent => {
Expand Down
24 changes: 24 additions & 0 deletions src/crates/core/src/agentic_os/work/hooks.rs
Original file line number Diff line number Diff line change
Expand Up @@ -289,6 +289,30 @@ impl WorkLifecycleHookBus {
items: reports,
}
}

pub async fn notify_deleted(
&self,
context: &WorkLifecycleHookContext,
report: WorkCleanupReport,
) {
let hook = WorkLifecycleHookKind::Deleted { report };
for handler in self.handlers.iter() {
if !handler
.phases()
.contains(&WorkLifecycleHookPhase::AfterCommit)
{
continue;
}
if let Err(error) = handler.handle(context, &hook).await {
log::warn!(
"Work delete lifecycle hook failed after commit: handler_id={} work_id={} error={}",
handler.id(),
context.work.id,
error
);
}
}
}
}

struct WorkSessionLifecycleHook;
Expand Down
6 changes: 3 additions & 3 deletions src/crates/core/src/agentic_os/work/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -50,9 +50,9 @@ pub use service::{
AdvanceWorkRequest, AdvanceWorkResponse, ControlWorkAction, ControlWorkRequest,
ControlWorkResponse, CreateWorkRequest, DeleteWorkResponse, DispatchNewWorkRequest,
DispatchWorkRequest, DispatchWorkResponse, LinkSessionToWorkRequest, PrimarySurfacePolicy,
ResolveAppWorkRequest, ResolveAppWorkResponse, ResolveComponentWorkRequest,
ResolveComponentWorkResponse, StartWorkRequest, StartWorkResponse, UpdateWorkRequest,
WorkService,
ReclassifyWorkRequest, ResolveAppWorkRequest, ResolveAppWorkResponse,
ResolveComponentWorkRequest, ResolveComponentWorkResponse, StartWorkRequest, StartWorkResponse,
UpdateWorkRequest, WorkService,
};
pub use store::{default_work_store, FileWorkStore, MemoryWorkStore, WorkStore};
pub use subject::{
Expand Down
9 changes: 9 additions & 0 deletions src/crates/core/src/agentic_os/work/projection.rs
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,12 @@ pub struct WorkProjection {
pub scope: WorkScope,
pub primary_surface: WorkSurfaceRef,
pub running: bool,
#[serde(default)]
pub system_managed: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub system_process_kind: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub topic_work_id: Option<WorkId>,
pub updated_at: i64,
}

Expand All @@ -37,6 +43,9 @@ impl From<&WorkRecord> for WorkProjection {
.execution_bindings
.iter()
.any(|binding| binding.is_running()),
system_managed: record.system_managed,
system_process_kind: record.system_process_kind.clone(),
topic_work_id: record.topic_work_id.clone(),
updated_at: record.updated_at,
}
}
Expand Down
9 changes: 9 additions & 0 deletions src/crates/core/src/agentic_os/work/record.rs
Original file line number Diff line number Diff line change
Expand Up @@ -152,6 +152,12 @@ pub struct WorkRecord {
pub builder_issues: Vec<WorkBuilderIssue>,
pub artifact_refs: Vec<ArtifactRef>,
pub memory_refs: Vec<MemoryRef>,
#[serde(default)]
pub system_managed: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub system_process_kind: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub topic_work_id: Option<WorkId>,
pub created_at: i64,
pub updated_at: i64,
}
Expand Down Expand Up @@ -199,6 +205,9 @@ impl WorkRecord {
builder_issues: Vec::new(),
artifact_refs: Vec::new(),
memory_refs: Vec::new(),
system_managed: false,
system_process_kind: None,
topic_work_id: None,
created_at: now,
updated_at: now,
}
Expand Down
Loading
Loading