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
7 changes: 7 additions & 0 deletions src/apps/cli/AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -165,3 +165,10 @@ guide.
For unattended question lifecycle changes, run `cargo test --locked -p openbitfun-cli --bin openbitfun shared_runtime::` and `cargo test --locked -p openbitfun-agent-runtime-ipc protocol_contract_tests::`.

For Pages account adapters, use the focused Core command in its guide and `cargo check -p openbitfun-cli`.

For `/goal` prompt routing (including the pending-session guard):

```bash
cargo test --locked -p openbitfun-cli --bin openbitfun goal_prompts_
cargo test --locked -p openbitfun-cli --bin openbitfun -- modes::exec::tests:: dispatch::worker::tests::
```
23 changes: 23 additions & 0 deletions src/apps/cli/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,7 @@ equivalent exists:
| `/new` or `/clear` | Start a new session. |
| `/timeline` | Navigate persisted user messages without changing the session. |
| `/fork` | Fork the full session or fork immediately before a selected prompt. |
| `/goal <objective>` | Start a persistent goal through the shared runtime; while working, steer the active turn toward it. |
| `/compact` or `/summarize` | Compact model context without deleting the saved transcript. |
| `/undo` / `/redo` | Move the persisted session timeline backward or forward. |
| `/diff` | Review staged, unstaged, and untracked workspace changes. |
Expand All @@ -71,6 +72,28 @@ selected child Session's active execution subtree.
configure a command that waits until the file is closed; missing commands,
non-zero exits, and empty editor output leave the current draft unchanged.

### Long-running goals

A prompt beginning with `/goal <objective>` activates the goal on the executing
host, including interactive input and `exec`. `exec` and detached dispatch keep
observing the goal's continuation turns; a successful intermediate turn does not
finish the job. Completion finishes successfully; blocked, paused, quota-limited,
or budget-limited goals return an incomplete/error outcome with the saved session
available for inspection and explicit resumption where supported.

Plain `/goal` prompts have no token budget by default. The optional `create_goal`
tool budget is set only on explicit request and accounts for non-cached input plus
output on the main session, not provider-wide billing or child-session usage. It
is a soft budget checked by the runtime, with one final wrap-up turn. An existing
100-continuation safety stop marks an unfinished goal blocked; explicit resume
starts a fresh continuation window without resetting accumulated usage.

The host must contain this behavior; a newer mobile or peer controller cannot add
it to an older target. Goal state survives in session storage, but host shutdown
is not automatic restart/recovery. Review the saved session and explicitly resume
after an interruption. Completion still depends on the model verifying the user's
requirements against real evidence; the runtime does not prove arbitrary tasks.

### Prompt continuity

Unsent drafts stay with their session while the TUI remains open. Switching,
Expand Down
17 changes: 17 additions & 0 deletions src/apps/cli/src/actions.rs
Original file line number Diff line number Diff line change
Expand Up @@ -100,6 +100,7 @@ pub(crate) enum ActionHandler {
Status,
WorkspaceDiff,
CompactSession,
GoalPrompt,
Usage,
Editor,
PromptStash,
Expand Down Expand Up @@ -168,6 +169,7 @@ impl ActionHandler {
| Self::Status
| Self::WorkspaceDiff
| Self::CompactSession
| Self::GoalPrompt
| Self::Editor
| Self::PromptStash
| Self::PromptStashPop
Expand Down Expand Up @@ -689,6 +691,21 @@ static ACTION_SPECS: &[ActionSpec] = &[
shortcut_label: None,
slash_on_startup: false,
},
ActionSpec {
id: "goal",
name: "Set a goal",
aliases: &["/goal"],
description: "Send /goal <objective> to start goal mode",
contexts: CHAT,
availability: ActionAvailability::Always,
handler: ActionHandler::GoalPrompt,
default_bindings: &[],
fallback_bindings: &[],
shortcut_field: None,
palette: palette("Session", false),
shortcut_label: None,
slash_on_startup: false,
},
ActionSpec {
id: "compact_session",
name: "Compact context",
Expand Down
7 changes: 7 additions & 0 deletions src/apps/cli/src/agent/runtime_client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1979,6 +1979,13 @@ impl CliAgentRuntimeClient {
}
}

/// Keep cancellation bound to the runtime-owned continuation accepted by a job observer.
pub(crate) async fn observe_active_turn(&self, session_id: &str, turn_id: &str) {
if self.session_id.lock().await.as_deref() == Some(session_id) {
*self.current_turn_id.lock().await = Some(turn_id.to_string());
}
}

pub(crate) async fn cancel_current_turn(&self) -> Result<()> {
let session_id = self.session_id.lock().await.clone();
let turn_id = self.current_turn_id.lock().await.clone();
Expand Down
52 changes: 52 additions & 0 deletions src/apps/cli/src/dispatch/store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -436,6 +436,23 @@ impl DispatchStore {
Ok((current, true))
}

/// Persist the runtime-owned next turn without re-claiming or finishing the job.
pub(crate) fn advance_goal_turn(
&self,
job_id: &str,
previous_turn: &str,
next_turn: &str,
) -> Result<()> {
let job_dir = self.existing_job_dir(job_id)?;
let _lock = JobLock::exclusive(&job_dir.join(".lock"))?;
let mut current = self.load_state_unlocked(&job_dir)?;
if current.state.is_terminal() || current.turn_id.as_deref() != Some(previous_turn) {
bail!("cannot advance goal turn: dispatch job owner changed");
}
current.turn_id = Some(next_turn.to_string());
atomic_write_json(&job_dir.join(STATE_FILE), &current)
}

pub(crate) fn request_cancel(&self, job_id: &str) -> Result<DispatchStateRecord> {
let job_dir = self.existing_job_dir(job_id)?;
let _lock = JobLock::exclusive(&job_dir.join(".lock"))?;
Expand Down Expand Up @@ -2681,6 +2698,41 @@ mod tests {
assert!(error.to_string().contains("still running"));
}

#[test]
fn goal_continuation_persists_its_turn_and_rejects_stale_or_terminal_owners() {
let (_dir, store) = store();
store
.create_job(request("goal-turn"), "Goal".into())
.unwrap();
store
.mark_state("goal-turn", DispatchJobState::Running, Some("turn-1"), None)
.unwrap();
store
.advance_goal_turn("goal-turn", "turn-1", "turn-2")
.unwrap();
let restored = store.load_state("goal-turn").unwrap();
assert_eq!(restored.state, DispatchJobState::Running);
assert_eq!(restored.turn_id.as_deref(), Some("turn-2"));
assert!(store
.advance_goal_turn("goal-turn", "turn-1", "stale")
.is_err());
store
.mark_state(
"goal-turn",
DispatchJobState::Succeeded,
Some("turn-2"),
None,
)
.unwrap();
assert!(store
.advance_goal_turn("goal-turn", "turn-2", "late")
.is_err());
assert_eq!(
store.load_state("goal-turn").unwrap().turn_id.as_deref(),
Some("turn-2")
);
}

#[test]
fn terminal_state_is_idempotent() {
let (_dir, store) = store();
Expand Down
41 changes: 40 additions & 1 deletion src/apps/cli/src/dispatch/worker.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ use openbitfun_agent_runtime::sdk::{
AgentSessionRestoreRequest, AgentTurnCancellationRequest, AgentTurnSettlementRequest,
PermissionReply, PermissionReplySource, PermissionRequest, PermissionRequestEvent,
};
use openbitfun_agent_runtime::thread_goal::{ThreadGoalRunDisposition, ThreadGoalRunTracker};
use openbitfun_events::{project_agentic_frontend_event, AgenticEvent};
use openbitfun_runtime_ports::{
AgentSubmissionSource, DialogSubmissionPolicy, SessionExecutionTarget,
Expand Down Expand Up @@ -226,7 +227,18 @@ async fn run_inner(store: &DispatchStore, job_id: &str) -> Result<()> {
.context("apply dispatch turn model and reasoning preset")?;
}

let turn_id = uuid::Uuid::new_v4().to_string();
let mut goal_run = ThreadGoalRunTracker::new(
agent_runtime
.get_thread_goal(openbitfun_agent_runtime::sdk::AgentThreadGoalGetRequest {
session_id: job.request.session_id.clone(),
workspace_path: workspace_path.clone(),
remote_connection_id: None,
remote_ssh_host: None,
})
.await
.map_err(|error| anyhow!(error.into_message()))?,
);
let mut turn_id = uuid::Uuid::new_v4().to_string();
// Claim the queued follow-up and persist the turn id in one step. A crash
// after the Runtime accepts the turn must never make a replacement worker
// submit the prompt a second time.
Expand Down Expand Up @@ -328,6 +340,19 @@ async fn run_inner(store: &DispatchStore, job_id: &str) -> Result<()> {
);
}
};
let mut goal_disposition = None;
if let AgenticEvent::ThreadGoalUpdated { session_id: owner, goal } = &envelope.event {
if owner == &job.request.session_id {
goal_disposition = goal_run.observe_goal(goal.clone().map(serde_json::from_value).transpose()?);
}
}
if let AgenticEvent::DialogTurnStarted { session_id: owner, turn_id: next_turn, user_message_metadata, .. } = &envelope.event {
if owner == &job.request.session_id && goal_run.accept_continuation(user_message_metadata.as_ref()) {
store.advance_goal_turn(job_id, &turn_id, next_turn)?;
turn_id = next_turn.clone();
event_scope.turn_id = turn_id.clone();
}
}
if !event_scope.admit(&envelope.event) {
continue;
}
Expand All @@ -339,7 +364,21 @@ async fn run_inner(store: &DispatchStore, job_id: &str) -> Result<()> {
job_id,
&DispatchEvent::agent_event(raw, projection),
)?;
if let Some(disposition) = goal_disposition {
match disposition {
ThreadGoalRunDisposition::Complete => break (DispatchJobState::Succeeded, None),
ThreadGoalRunDisposition::Stopped(status) => break (DispatchJobState::Failed, Some(format!("Thread goal stopped with status {status:?}; the objective is not complete"))),
ThreadGoalRunDisposition::Continue => {}
}
}
if let Some(outcome) = terminal_outcome(&envelope.event, &turn_id) {
if outcome.0 == DispatchJobState::Succeeded && turn_kind == DispatchTurnKind::Prompt {
match goal_run.after_successful_turn() {
ThreadGoalRunDisposition::Continue => continue,
ThreadGoalRunDisposition::Complete => {}
ThreadGoalRunDisposition::Stopped(status) => break (DispatchJobState::Failed, Some(format!("Thread goal stopped with status {status:?}; the objective is not complete"))),
}
}
break outcome;
}
}
Expand Down
23 changes: 18 additions & 5 deletions src/apps/cli/src/modes/chat/commands.rs
Original file line number Diff line number Diff line change
@@ -1,5 +1,13 @@
fn is_runtime_goal_prompt(input: &str) -> bool {
openbitfun_core::agentic::goal_mode::goal_objective_from_prompt(input).is_some()
}

fn is_local_slash_command(input: &str) -> bool {
input.trim().starts_with('/') && !is_runtime_goal_prompt(input)
}

fn session_update_blocks_typed_submission(pending_for_current_session: bool, input: &str) -> bool {
pending_for_current_session && !input.trim().starts_with('/')
pending_for_current_session && !is_local_slash_command(input)
}

fn steering_unsupported_reason(draft: &crate::ui::composer::ComposerDraft) -> Option<&'static str> {
Expand Down Expand Up @@ -155,6 +163,7 @@ fn builtin_arguments_error(
fn selected_command_prefill(handler: ActionHandler) -> Option<&'static str> {
match handler {
ActionHandler::RenameSession => Some("/rename "),
ActionHandler::GoalPrompt => Some("/goal "),
_ => None,
}
}
Expand Down Expand Up @@ -1159,6 +1168,10 @@ impl ChatMode {
});
self.pending_workspace_diff = Some(PendingWorkspaceDiff { handle });
}
ActionHandler::GoalPrompt => {
chat_view.set_input("/goal ");
chat_view.set_status(Some("Usage: /goal <objective>. Goal controls are available in the desktop goal menu.".to_string()));
}
ActionHandler::CompactSession => {
self.start_session_compaction(chat_view, chat_state, rt_handle);
}
Expand Down Expand Up @@ -1611,11 +1624,11 @@ impl ChatMode {
chat_view.set_status(Some("Images are unavailable in Shell mode".to_string()));
return Ok(None);
}
if !shell_mode && draft_has_images && trimmed.starts_with('/') {
if !shell_mode && draft_has_images && is_local_slash_command(trimmed) {
chat_view.set_status(Some(IMAGE_ATTACHMENTS_REQUIRE_MESSAGE.to_string()));
return Ok(None);
}
if shell_mode || !trimmed.starts_with('/') {
if shell_mode || !is_local_slash_command(trimmed) {
self.selected_native_command_once = None;
}
let pending_for_current_session = self
Expand All @@ -1632,7 +1645,7 @@ impl ChatMode {
}

if chat_state.is_processing {
if !shell_mode && trimmed.starts_with('/') {
if !shell_mode && is_local_slash_command(trimmed) {
if let Some(input) = chat_view.send_input() {
return self.handle_command(&input.text, chat_view, chat_state, rt_handle);
}
Expand All @@ -1658,7 +1671,7 @@ impl ChatMode {
return Ok(None);
}
tracing::info!("User input: {}", input.text);
if input.text.starts_with('/') {
if is_local_slash_command(&input.text) {
return self.handle_command(&input.text, chat_view, chat_state, rt_handle);
}
self.send_draft_to_agent(input, chat_view, chat_state, rt_handle);
Expand Down
2 changes: 1 addition & 1 deletion src/apps/cli/src/modes/chat/run.rs
Original file line number Diff line number Diff line change
Expand Up @@ -591,7 +591,7 @@ impl ChatMode {
"The restored session uses fallback settings. Review them, then send the preserved input explicitly."
.to_string(),
));
} else if draft.text.starts_with('/') {
} else if is_local_slash_command(&draft.text) {
// Slash commands will be handled in the main loop
chat_view.set_draft(draft);
} else {
Expand Down
31 changes: 31 additions & 0 deletions src/apps/cli/src/modes/chat/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -73,6 +73,37 @@ mod tests {
snapshot.into()
}

#[test]
fn goal_prompts_use_normal_runtime_submission_and_pending_guards() {
use crate::actions::ActionContext;
for input in [
"/goal finish tests",
" /GOAL\nship feature ",
"/goal pause\nthen verify",
] {
assert!(super::is_runtime_goal_prompt(input));
assert!(!super::is_local_slash_command(input));
assert!(session_update_blocks_typed_submission(true, input));
}
for input in [
"/goal",
"/goal pause",
"/goal resume",
"/goal clear",
"/goalie work",
"/compact",
] {
assert!(super::is_local_slash_command(input));
}
assert_eq!(
selected_command_prefill(ActionHandler::GoalPrompt),
Some("/goal ")
);
let action = crate::actions::action_for_alias("/goal", ActionContext::Chat).unwrap();
assert_eq!(action.handler, ActionHandler::GoalPrompt);
assert!(action.handler.available_in_shared_tui(ActionContext::Chat));
}

#[test]
fn explicit_same_id_agent_selection_rebinds_through_the_runtime_owner() {
let source = include_str!("selection.rs").replace("\r\n", "\n");
Expand Down
Loading
Loading