use actix_web::{web, HttpResponse, Responder};
use super::{ChatRequest, ChatResponse};
use crate::app_state::AppState;
use bamboo_engine::config::GoldConfig;
use bamboo_engine::model_config_helper::{
parse_session_gold_config, resolve_gold_config, GOLD_CONFIG_METADATA_KEY,
};
use bamboo_engine::session_app::chat::{parse_goal_command, GoalCommand};
use bamboo_engine::session_app::metadata::SessionMetadataService;
mod images;
mod request;
fn sync_runtime_workspace(session_id: &str, workspace_path: Option<&str>) {
if let Some(workspace) = workspace_path
.map(str::trim)
.filter(|s| !s.is_empty())
.map(std::path::PathBuf::from)
.and_then(|path| std::fs::canonicalize(&path).ok().or(Some(path)))
{
bamboo_tools::tools::workspace_state::publish_resolved_workspace(session_id, workspace);
}
}
async fn save_and_cache_session_locked(
state: &AppState,
session: &bamboo_agent_core::Session,
) -> Result<(), HttpResponse> {
state
.persistence
.storage()
.save_session(session)
.await
.map_err(|error| {
HttpResponse::InternalServerError().json(serde_json::json!({
"error": crate::error::error_value(format!(
"Failed to persist chat session: {error}"
))
}))
})?;
state.sessions.insert(
session.id.clone(),
std::sync::Arc::new(parking_lot::RwLock::new(session.clone())),
);
Ok(())
}
fn project_context_error_response(
error: bamboo_engine::project_context::ProjectContextError,
) -> HttpResponse {
use bamboo_engine::project_context::ProjectContextError;
match error {
ProjectContextError::WorkspaceConflict {
workspace,
owner_project_id,
session_project_id,
} => HttpResponse::Conflict().json(serde_json::json!({
"error": {
"type": "api_error",
"code": "project_workspace_conflict",
"message": "Workspace belongs to another Project"
},
"workspace": workspace,
"owner_project_id": owner_project_id,
"session_project_id": session_project_id,
})),
ProjectContextError::UnassignedWorkspaceConflict {
workspace,
owner_project_id,
} => HttpResponse::Conflict().json(serde_json::json!({
"error": {
"type": "api_error",
"code": "project_workspace_conflict",
"message": "Workspace belongs to another Project"
},
"workspace": workspace,
"owner_project_id": owner_project_id,
"session_project_id": "unassigned",
})),
ProjectContextError::WorkspaceInvalid { workspace, message } => HttpResponse::BadRequest()
.json(serde_json::json!({
"error": {
"type": "api_error",
"code": "workspace_invalid",
"message": message
},
"workspace": workspace,
})),
ProjectContextError::InvalidProjectIdentity { raw, message } => HttpResponse::BadRequest()
.json(serde_json::json!({
"error": {
"type": "api_error",
"code": "invalid_project_identity",
"message": format!(
"Session carries an invalid Project identity '{raw}': {message}"
)
}
})),
ProjectContextError::ProjectUnavailable { project_id } => {
HttpResponse::Conflict().json(serde_json::json!({
"error": {
"type": "api_error",
"code": "project_unavailable",
"message": "Assigned Project is unavailable"
},
"project_id": project_id,
}))
}
error @ (ProjectContextError::Source(_) | ProjectContextError::IdentityMismatch { .. }) => {
tracing::error!(%error, "failed to resolve Project context");
crate::error::json_error(
actix_web::http::StatusCode::INTERNAL_SERVER_ERROR,
"Failed to resolve Project context",
)
}
}
}
#[cfg(test)]
mod tests;
pub async fn handler(state: web::Data<AppState>, req: web::Json<ChatRequest>) -> impl Responder {
let session_id = request::resolve_session_id(req.session_id.as_deref());
let (existing_session_found, existing_project_id, existing_workspace) = match state
.storage
.load_session(&session_id)
.await
{
Ok(Some(existing)) => {
let project_id = match bamboo_engine::project_context::ProjectContextResolver::session_project_identity(&existing) {
bamboo_engine::project_context::SessionProjectIdentity::Assigned(project_id) => Some(project_id),
bamboo_engine::project_context::SessionProjectIdentity::Unassigned => None,
bamboo_engine::project_context::SessionProjectIdentity::Invalid { raw, message } => {
return HttpResponse::BadRequest().json(serde_json::json!({
"error": {
"type": "api_error",
"code": "invalid_project_identity",
"message": format!(
"Session carries an invalid Project identity '{raw}': {message}"
)
},
"session_id": session_id,
}));
}
};
(true, project_id, existing.workspace_path_meta())
}
Ok(None) => (false, None, None),
Err(error) => {
tracing::error!(%error, "failed to load chat session for Project validation");
return crate::error::json_error(
actix_web::http::StatusCode::INTERNAL_SERVER_ERROR,
"Failed to validate session Project membership",
);
}
};
if let Some(project_id) = req.project_id.as_ref() {
match state.project_store.get(project_id) {
Ok(project) if project.status == bamboo_domain::ProjectStatus::Active => {}
Ok(_) => {
return HttpResponse::Conflict().json(serde_json::json!({
"error": {
"type": "api_error",
"code": "project_archived",
"message": "Sessions can only be created in an active Project"
},
"project_id": project_id,
}));
}
Err(bamboo_projects::ProjectStoreError::NotFound(_)) => {
return crate::error::json_error(
actix_web::http::StatusCode::NOT_FOUND,
"target Project not found",
);
}
Err(error) => {
tracing::error!(%error, "failed to validate chat Project");
return crate::error::json_error(
actix_web::http::StatusCode::INTERNAL_SERVER_ERROR,
"Failed to validate target Project",
);
}
}
if existing_session_found && existing_project_id.as_ref() != Some(project_id) {
return HttpResponse::Conflict().json(serde_json::json!({
"error": {
"type": "api_error",
"code": "session_project_reassignment_required",
"message": "Chat cannot change Project membership; use PATCH /sessions/{id}"
},
"session_id": session_id,
"current_project_id": existing_project_id,
"requested_project_id": project_id,
}));
}
}
let effective_project_id = req.project_id.clone().or(existing_project_id);
let requested_workspace = req
.workspace_path
.as_deref()
.or(existing_workspace.as_deref());
let final_workspace = match crate::project_context::validate_workspace_assignment(
&state.project_store,
effective_project_id.as_ref(),
requested_workspace,
) {
Ok(workspace) => workspace,
Err(error) => {
return match error {
crate::project_context::ProjectWorkspaceValidationError::Invalid {
code,
workspace,
message,
} => HttpResponse::BadRequest().json(serde_json::json!({
"error": {
"type": "api_error",
"code": code,
"message": message
},
"workspace": workspace,
})),
crate::project_context::ProjectWorkspaceValidationError::Conflict {
workspace,
owner_project_id,
session_project_id,
} => HttpResponse::Conflict().json(serde_json::json!({
"error": {
"type": "api_error",
"code": "project_workspace_conflict",
"message": "Workspace belongs to another Project"
},
"workspace": workspace,
"owner_project_id": owner_project_id,
"session_project_id": session_project_id,
})),
crate::project_context::ProjectWorkspaceValidationError::Store(error) => {
tracing::error!(%error, "failed to validate workspace Project ownership");
crate::error::json_error(
actix_web::http::StatusCode::INTERNAL_SERVER_ERROR,
"Failed to validate workspace Project ownership",
)
}
};
}
};
let final_workspace_display = final_workspace
.as_deref()
.map(bamboo_config::paths::path_to_display_string);
tracing::debug!(
"[{}] Chat requested: message_len={}, is_goal_command={}, image_count={}",
session_id,
req.message.len(),
parse_goal_command(&req.message).is_some(),
req.images.as_ref().map(|i| i.len()).unwrap_or(0),
);
let config_snapshot = state.config.read().await.clone();
let default_model = if request::optional_non_empty(req.model.as_deref()).is_some() {
None
} else {
bamboo_engine::resolved_defaults::resolve_default_run_config(
&config_snapshot,
&state.provider_registry,
)
.model_roster
.model
};
let model = match request::resolve_model(req.model.as_deref(), default_model.as_deref()) {
Ok(model) => model,
Err(response) => return response,
};
let global_default_prompt =
bamboo_engine::prompt_defaults::read_global_default_system_prompt_template();
let builtin_fallback_prompt = crate::app_state::DEFAULT_BASE_PROMPT;
let data_dir = Some(state.app_data_dir.clone());
let mut project_preflight = bamboo_agent_core::Session::new(&session_id, &model);
if let Some(project_id) = effective_project_id.as_ref() {
project_preflight.set_project_id_meta(project_id.to_string());
}
if let Some(workspace) = final_workspace_display.as_deref() {
project_preflight.set_workspace_path_meta(workspace);
} else if let Some(workspace) = config_snapshot.get_default_work_area_path() {
project_preflight
.set_workspace_path_meta(bamboo_config::paths::path_to_display_string(&workspace));
}
if let Err(error) = state
.project_context_resolver
.refresh_session_prompt_read_only(&mut project_preflight)
.await
{
return project_context_error_response(error);
}
let workspace_was_explicit = req.workspace_path.is_some();
let mut input = bamboo_engine::session_app::types::ChatTurnInput {
session_id: session_id.clone(),
project_id: effective_project_id,
model: model.clone(),
model_ref: req.model_ref.clone(),
provider: req.provider.clone(),
message: req.message.clone(),
system_prompt: request::optional_non_empty(req.system_prompt.as_deref()).map(String::from),
enhance_prompt: request::optional_non_empty(req.enhance_prompt.as_deref())
.map(String::from),
workspace_path: workspace_was_explicit
.then(|| project_preflight.workspace_path_meta())
.flatten(),
default_workspace_path: config_snapshot
.get_default_work_area_path()
.as_deref()
.map(bamboo_config::paths::path_to_display_string),
selected_skill_ids: req.selected_skill_ids.clone(),
workflow_selection: req.workflow_selection.clone(),
orchestration_opt_in: req.orchestration_opt_in,
copilot_conclusion_with_options_enhancement_enabled: req
.copilot_conclusion_with_options_enhancement_enabled,
data_dir,
};
let persistence_guard = state.persistence.acquire_lock(&session_id).await;
let authoritative_session = match state.persistence.storage().load_session(&session_id).await {
Ok(session) => session,
Err(error) => {
tracing::error!(%error, "failed to load authoritative chat session");
return crate::error::json_error(
actix_web::http::StatusCode::INTERNAL_SERVER_ERROR,
"Failed to load authoritative chat session",
);
}
};
if req.project_id.is_none() {
input.project_id = authoritative_session.as_ref().and_then(|session| {
match bamboo_engine::project_context::ProjectContextResolver::session_project_identity(
session,
) {
bamboo_engine::project_context::SessionProjectIdentity::Assigned(project_id) => {
Some(project_id)
}
bamboo_engine::project_context::SessionProjectIdentity::Unassigned
| bamboo_engine::project_context::SessionProjectIdentity::Invalid { .. } => None,
}
});
}
if let Some(requested_workspace) = req.workspace_path.as_deref() {
input.workspace_path = match crate::project_context::validate_workspace_assignment(
&state.project_store,
input.project_id.as_ref(),
Some(requested_workspace),
) {
Ok(workspace) => workspace
.as_deref()
.map(bamboo_config::paths::path_to_display_string),
Err(error) => {
return match error {
crate::project_context::ProjectWorkspaceValidationError::Invalid {
code,
workspace,
message,
} => HttpResponse::BadRequest().json(serde_json::json!({
"error": {
"type": "api_error",
"code": code,
"message": message
},
"workspace": workspace,
})),
crate::project_context::ProjectWorkspaceValidationError::Conflict {
workspace,
owner_project_id,
session_project_id,
} => HttpResponse::Conflict().json(serde_json::json!({
"error": {
"type": "api_error",
"code": "project_workspace_conflict",
"message": "Workspace belongs to another Project"
},
"workspace": workspace,
"owner_project_id": owner_project_id,
"session_project_id": session_project_id,
})),
crate::project_context::ProjectWorkspaceValidationError::Store(error) => {
tracing::error!(%error, "failed to revalidate chat workspace ownership");
crate::error::json_error(
actix_web::http::StatusCode::INTERNAL_SERVER_ERROR,
"Failed to validate workspace Project ownership",
)
}
};
}
};
} else {
input.workspace_path = None;
}
let mut session =
match bamboo_engine::session_app::chat::prepare_chat_turn_from_authoritative_session(
authoritative_session,
input,
global_default_prompt.as_str(),
builtin_fallback_prompt,
) {
Ok(session) => session,
Err(bamboo_engine::session_app::errors::ChatError::InvalidWorkflowSelection(error)) => {
return HttpResponse::BadRequest().json(serde_json::json!({
"error": crate::error::error_value(error)
}));
}
Err(bamboo_engine::session_app::errors::ChatError::InvalidProjectIdentity {
raw,
message,
}) => {
return HttpResponse::BadRequest().json(serde_json::json!({
"error": {
"type": "api_error",
"code": "invalid_project_identity",
"message": format!(
"Session carries an invalid Project identity '{raw}': {message}"
)
},
"session_id": session_id,
}));
}
Err(bamboo_engine::session_app::errors::ChatError::ProjectIdentityConflict {
expected,
actual,
}) => {
return HttpResponse::Conflict().json(serde_json::json!({
"error": {
"type": "api_error",
"code": "session_project_changed",
"message": "Session Project membership changed while preparing chat"
},
"session_id": session_id,
"expected_project_id": expected,
"actual_project_id": actual,
}));
}
Err(error) => {
tracing::error!("Chat turn preparation failed: {error}");
return HttpResponse::InternalServerError().json(serde_json::json!({
"error": crate::error::error_value(format!("Failed to prepare chat: {error}"))
}));
}
};
let session_was_created = session
.metadata
.get(bamboo_engine::session_app::chat::SESSION_START_SOURCE_METADATA_KEY)
.is_some_and(|source| source == "startup");
if let Err(error) = state
.project_context_resolver
.refresh_session_prompt_read_only(&mut session)
.await
{
return project_context_error_response(error);
}
if let Err(response) = save_and_cache_session_locked(state.as_ref(), &session).await {
return response;
}
sync_runtime_workspace(&session_id, session.workspace_path_meta().as_deref());
if session_was_created {
state.account_sink.record(
Some(&session_id),
&bamboo_agent_core::AgentEvent::SessionCreated {
session_id: session_id.clone(),
project_id: session.project_id_meta(),
title: session.title.clone(),
kind: session.kind,
created_at: session.created_at,
},
);
}
let effective_message = match crate::lifecycle_hooks::apply_user_prompt_submit_hooks(
&config_snapshot.lifecycle_hooks,
Some(state.app_data_dir.clone()),
&mut session,
&req.message,
)
.await
{
Ok(message) => message,
Err(reason) => {
if let Err(response) = save_and_cache_session_locked(state.as_ref(), &session).await {
return response;
}
return HttpResponse::BadRequest().json(serde_json::json!({
"error": crate::error::error_value(reason),
"hook_event": "UserPromptSubmit"
}));
}
};
if let Some(goal_cmd) = parse_goal_command(&req.message) {
tracing::debug!(
"[{}] Chat intercepted as /goal command: {:?}",
session_id,
goal_cmd
);
drop(persistence_guard);
return handle_goal_command(state.as_ref(), &session_id, &goal_cmd).await;
}
if let Err(response) = images::append_user_message(
&state,
&mut session,
&effective_message,
req.images.as_deref(),
)
.await
{
return response;
}
if let Err(response) = save_and_cache_session_locked(state.as_ref(), &session).await {
return response;
}
if let Some(msg) = session.messages.last() {
state.account_sink.record(
Some(&session_id),
&bamboo_agent_core::AgentEvent::MessageAppended {
session_id: session_id.clone(),
message_id: msg.id.clone(),
role: msg.role.clone(),
content: msg.content.clone(),
created_at: msg.created_at,
},
);
}
tracing::debug!(
"[{}] Chat turn persisted: messages={}, last_role={:?} -> client should now POST /execute",
session_id,
session.messages.len(),
session.messages.last().map(|m| format!("{:?}", m.role)),
);
drop(persistence_guard);
HttpResponse::Created().json(ChatResponse {
session_id: session_id.clone(),
stream_url: format!("/api/v1/events/{}", session_id),
status: "streaming".to_string(),
goal_command: None,
})
}
#[derive(Debug, serde::Serialize)]
pub struct GoalCommandResponse {
pub action: String,
pub should_execute: bool,
pub gold_config: Option<GoldConfig>,
}
async fn handle_goal_command(
state: &AppState,
session_id: &str,
cmd: &GoalCommand,
) -> HttpResponse {
let config_snapshot = state.config.read().await.clone();
let session = match state.load_session_merged(session_id).await {
Some(s) => s,
None => {
return HttpResponse::NotFound().json(serde_json::json!({
"error": crate::error::error_value("Session not found"),
"session_id": session_id
}));
}
};
let current_json = session.metadata.get(GOLD_CONFIG_METADATA_KEY).cloned();
let current_effective = resolve_gold_config(&config_snapshot, current_json.as_deref());
let (new_config, should_resume) = match cmd {
GoalCommand::Status => {
let response_config = current_effective.clone();
return HttpResponse::Ok().json(ChatResponse {
session_id: session_id.to_string(),
stream_url: format!("/api/v1/events/{}", session_id),
status: "accepted".to_string(),
goal_command: Some(GoalCommandResponse {
action: "status".to_string(),
should_execute: false,
gold_config: response_config,
}),
});
}
GoalCommand::Off => {
let mut cfg = current_effective.unwrap_or_default();
cfg.enabled = false;
cfg.auto_answer_enabled = false;
cfg.auto_continue_enabled = false;
(cfg, false)
}
GoalCommand::Clear => {
let mut cfg = current_effective.unwrap_or_default();
cfg.enabled = false;
cfg.auto_answer_enabled = false;
cfg.auto_continue_enabled = false;
cfg.goal = None;
cfg.evaluation_prompt = None;
(cfg, false)
}
GoalCommand::On => {
let mut cfg = current_effective.unwrap_or_default();
let has_prompt = cfg.effective_goal().is_some();
if !has_prompt {
return HttpResponse::Ok().json(ChatResponse {
session_id: session_id.to_string(),
stream_url: format!("/api/v1/events/{}", session_id),
status: "accepted".to_string(),
goal_command: Some(GoalCommandResponse {
action: "on_no_prompt".to_string(),
should_execute: false,
gold_config: Some(cfg),
}),
});
}
cfg.enabled = true;
cfg.auto_answer_enabled = true;
cfg.auto_continue_enabled = true;
(cfg, false)
}
GoalCommand::SetPrompt(prompt) => {
let mut cfg = current_effective.unwrap_or_default();
cfg.enabled = true;
cfg.auto_answer_enabled = true;
cfg.auto_continue_enabled = true;
cfg.goal = Some(prompt.clone());
(cfg, true)
}
};
let new_json = serde_json::to_string(&new_config).ok();
match SessionMetadataService::set_gold_config_json(state, session_id, new_json.clone(), None)
.await
{
Ok(_) => {}
Err(e) => {
tracing::error!(session_id = %session_id, "Failed to persist gold_config: {e}");
return HttpResponse::InternalServerError().json(serde_json::json!({
"error": crate::error::error_value(format!("Failed to update goal config: {e}"))
}));
}
}
if let Some(mut session) = state.load_session_merged(session_id).await {
session.metadata.retain(|key, _| !key.starts_with("gold."));
session.metadata.remove("goal.state");
if should_resume {
if let Some(runtime_state) = session.agent_runtime_state.as_mut() {
runtime_state.status = bamboo_domain::AgentStatusState::Idle;
runtime_state.suspension = None;
runtime_state.waiting_for_children = None;
}
session.metadata.remove("runtime.suspend_reason");
let goal_text = parse_session_gold_config(new_json.as_deref())
.as_ref()
.and_then(|cfg| cfg.effective_goal().map(str::to_string))
.unwrap_or_default();
let instruction = format!(
"The user has set a session goal:\n\n{goal_text}\n\nBefore taking any action, briefly confirm your understanding of this goal, surface any ambiguities or assumptions, and outline how you plan to achieve it. If anything is genuinely unclear or you need a decision from the user, ask them now. Otherwise, state your plan and begin working toward the goal."
);
let mut resume_msg = bamboo_domain::Message::user(instruction);
resume_msg.metadata = Some(serde_json::json!({
"hidden_from_ui": true,
"runtime_kind": "gold_goal_resume"
}));
session.add_message(resume_msg);
crate::handlers::agent::events::mark_pending_turn(&mut session);
}
state.save_and_cache_session(&mut session).await;
}
let response_config = parse_session_gold_config(new_json.as_deref());
HttpResponse::Ok().json(ChatResponse {
session_id: session_id.to_string(),
stream_url: format!("/api/v1/events/{}", session_id),
status: "accepted".to_string(),
goal_command: Some(GoalCommandResponse {
action: match cmd {
GoalCommand::Off => "off".to_string(),
GoalCommand::Clear => "clear".to_string(),
GoalCommand::On => "on".to_string(),
GoalCommand::SetPrompt(_) => "set_prompt".to_string(),
GoalCommand::Status => unreachable!(),
},
should_execute: should_resume,
gold_config: response_config,
}),
})
}