use actix_web::{http::StatusCode, web, HttpRequest, HttpResponse, Result};
use bamboo_agent_core::AgentEvent;
use bamboo_domain::{
LegacySessionProjectInput, ProjectId, ProjectManifest, ProjectStatus, WorkspaceBinding,
};
use bamboo_projects::{plan_legacy_migration, ProjectStoreError};
use serde::Deserialize;
use crate::app_state::AppState;
#[derive(Debug, Deserialize)]
pub struct CreateProjectRequest {
pub name: String,
#[serde(default)]
pub description: Option<String>,
#[serde(default)]
pub workspace_bindings: Vec<WorkspaceBinding>,
}
fn deserialize_nullable_description<'de, D>(
deserializer: D,
) -> std::result::Result<Option<Option<String>>, D::Error>
where
D: serde::Deserializer<'de>,
{
Option::<String>::deserialize(deserializer).map(Some)
}
#[derive(Debug, Deserialize)]
pub struct PatchProjectRequest {
#[serde(default)]
pub name: Option<String>,
#[serde(default, deserialize_with = "deserialize_nullable_description")]
pub description: Option<Option<String>>,
}
#[derive(Debug, Deserialize)]
pub struct WorkspaceMutationRequest {
pub path: String,
#[serde(default)]
pub label: Option<String>,
#[serde(default)]
pub git_common_dir: Option<String>,
}
#[derive(Debug, Deserialize)]
pub struct LegacyDryRunRequest {
#[serde(default)]
pub sessions: Vec<LegacySessionProjectInput>,
}
#[derive(Debug, Deserialize)]
pub struct LegacyMemoryMigrationRequest {
pub legacy_project_key: String,
}
#[derive(Debug, Deserialize)]
pub struct LegacyMemoryMigrationStatusQuery {
pub legacy_project_key: String,
}
fn parse_if_match(req: &HttpRequest) -> std::result::Result<u64, HttpResponse> {
let Some(raw) = req.headers().get(actix_web::http::header::IF_MATCH) else {
return Err(crate::error::json_error(
StatusCode::PRECONDITION_REQUIRED,
"If-Match with the current Project revision is required",
));
};
let value = raw
.to_str()
.ok()
.map(str::trim)
.and_then(|value| value.strip_prefix("W/").or(Some(value)))
.map(|value| value.trim_matches('"'))
.and_then(|value| value.parse::<u64>().ok());
value.ok_or_else(|| {
crate::error::json_error(
StatusCode::BAD_REQUEST,
"If-Match must be a Project revision integer or quoted ETag",
)
})
}
fn parse_id(raw: &str) -> std::result::Result<ProjectId, HttpResponse> {
raw.parse::<ProjectId>()
.map_err(|_| crate::error::json_error(StatusCode::BAD_REQUEST, "invalid Project id"))
}
fn project_error(error: ProjectStoreError) -> HttpResponse {
let (status, message) = match error {
ProjectStoreError::NotFound(_) => (StatusCode::NOT_FOUND, error.to_string()),
ProjectStoreError::Conflict { .. } => (StatusCode::PRECONDITION_FAILED, error.to_string()),
ProjectStoreError::AlreadyExists(_)
| ProjectStoreError::Validation(_)
| ProjectStoreError::InvalidPathComponent(_) => (StatusCode::CONFLICT, error.to_string()),
ProjectStoreError::Io(_) | ProjectStoreError::Json(_) => {
tracing::error!(%error, "Project registry operation failed");
(
StatusCode::INTERNAL_SERVER_ERROR,
"Project registry operation failed".to_string(),
)
}
};
crate::error::json_error(status, message)
}
fn with_etag(project: &ProjectManifest, status: StatusCode) -> HttpResponse {
HttpResponse::build(status)
.insert_header((
actix_web::http::header::ETAG,
format!("\"{}\"", project.revision),
))
.json(project)
}
pub async fn list_projects(state: web::Data<AppState>) -> Result<HttpResponse> {
match state.project_store.list() {
Ok(projects) => Ok(HttpResponse::Ok().json(serde_json::json!({ "projects": projects }))),
Err(error) => Ok(project_error(error)),
}
}
pub async fn create_project(
state: web::Data<AppState>,
request: web::Json<CreateProjectRequest>,
) -> Result<HttpResponse> {
let project = match state.project_store.create_with_bindings(
request.name.clone(),
request.description.clone(),
request.workspace_bindings.clone(),
) {
Ok(project) => project,
Err(error) => return Ok(project_error(error)),
};
state.account_sink.record(
None,
&AgentEvent::ProjectCreated {
project_id: project.id.to_string(),
revision: project.revision,
},
);
Ok(with_etag(&project, StatusCode::CREATED))
}
pub async fn get_project(
state: web::Data<AppState>,
path: web::Path<String>,
) -> Result<HttpResponse> {
let id = match parse_id(&path) {
Ok(id) => id,
Err(response) => return Ok(response),
};
match state.project_store.get(&id) {
Ok(project) => Ok(with_etag(&project, StatusCode::OK)),
Err(error) => Ok(project_error(error)),
}
}
pub async fn patch_project(
state: web::Data<AppState>,
path: web::Path<String>,
http_request: HttpRequest,
request: web::Json<PatchProjectRequest>,
) -> Result<HttpResponse> {
let id = match parse_id(&path) {
Ok(id) => id,
Err(response) => return Ok(response),
};
let expected = match parse_if_match(&http_request) {
Ok(revision) => revision,
Err(response) => return Ok(response),
};
let project = match state.project_store.update(&id, expected, |project| {
if let Some(name) = request.name.as_ref() {
project.name = name.clone();
}
if let Some(description) = request.description.as_ref() {
project.description = description.clone();
}
Ok(())
}) {
Ok(project) => project,
Err(error) => return Ok(project_error(error)),
};
state.account_sink.record(
None,
&AgentEvent::ProjectUpdated {
project_id: project.id.to_string(),
revision: project.revision,
},
);
Ok(with_etag(&project, StatusCode::OK))
}
pub async fn bind_workspace(
state: web::Data<AppState>,
path: web::Path<String>,
http_request: HttpRequest,
request: web::Json<WorkspaceMutationRequest>,
) -> Result<HttpResponse> {
let id = match parse_id(&path) {
Ok(id) => id,
Err(response) => return Ok(response),
};
let expected = match parse_if_match(&http_request) {
Ok(revision) => revision,
Err(response) => return Ok(response),
};
let project = match state.project_store.bind_workspace(
&id,
expected,
WorkspaceBinding {
path: request.path.clone(),
label: request.label.clone(),
git_common_dir: request.git_common_dir.clone(),
},
) {
Ok(project) => project,
Err(error) => return Ok(project_error(error)),
};
state.account_sink.record(
None,
&AgentEvent::ProjectUpdated {
project_id: project.id.to_string(),
revision: project.revision,
},
);
Ok(with_etag(&project, StatusCode::OK))
}
pub async fn unbind_workspace(
state: web::Data<AppState>,
path: web::Path<String>,
http_request: HttpRequest,
request: web::Json<WorkspaceMutationRequest>,
) -> Result<HttpResponse> {
let id = match parse_id(&path) {
Ok(id) => id,
Err(response) => return Ok(response),
};
let expected = match parse_if_match(&http_request) {
Ok(revision) => revision,
Err(response) => return Ok(response),
};
let project = match state
.project_store
.unbind_workspace(&id, expected, &request.path)
{
Ok(project) => project,
Err(error) => return Ok(project_error(error)),
};
state.account_sink.record(
None,
&AgentEvent::ProjectUpdated {
project_id: project.id.to_string(),
revision: project.revision,
},
);
Ok(with_etag(&project, StatusCode::OK))
}
pub async fn project_resources(
state: web::Data<AppState>,
path: web::Path<String>,
) -> Result<HttpResponse> {
let id = match parse_id(&path) {
Ok(id) => id,
Err(response) => return Ok(response),
};
match state.project_store.resource_summary(&id) {
Ok(summary) => Ok(HttpResponse::Ok().json(summary)),
Err(error) => Ok(project_error(error)),
}
}
pub async fn archive_project(
state: web::Data<AppState>,
path: web::Path<String>,
http_request: HttpRequest,
) -> Result<HttpResponse> {
let id = match parse_id(&path) {
Ok(id) => id,
Err(response) => return Ok(response),
};
let expected = match parse_if_match(&http_request) {
Ok(revision) => revision,
Err(response) => return Ok(response),
};
let project = match state.project_store.archive(&id, expected) {
Ok(project) => project,
Err(error) => return Ok(project_error(error)),
};
state.account_sink.record(
None,
&AgentEvent::ProjectArchived {
project_id: project.id.to_string(),
revision: project.revision,
},
);
Ok(with_etag(&project, StatusCode::OK))
}
pub async fn legacy_dry_run(
state: web::Data<AppState>,
request: web::Json<LegacyDryRunRequest>,
) -> Result<HttpResponse> {
let projects = match state.project_store.list() {
Ok(projects) => projects,
Err(error) => return Ok(project_error(error)),
};
Ok(HttpResponse::Ok().json(plan_legacy_migration(&request.sessions, &projects)))
}
pub async fn migrate_legacy_memory(
state: web::Data<AppState>,
path: web::Path<String>,
http_request: HttpRequest,
request: web::Json<LegacyMemoryMigrationRequest>,
) -> Result<HttpResponse> {
let id = match parse_id(&path) {
Ok(id) => id,
Err(response) => return Ok(response),
};
let expected = match parse_if_match(&http_request) {
Ok(revision) => revision,
Err(response) => return Ok(response),
};
let current = match state.project_store.get(&id) {
Ok(project) => project,
Err(error) => return Ok(project_error(error)),
};
if current.revision != expected {
return Ok(project_error(ProjectStoreError::Conflict {
expected,
actual: current.revision,
}));
}
let project = if current
.legacy_project_keys
.iter()
.any(|key| key == &request.legacy_project_key)
{
current
} else {
match state.project_store.update(&id, expected, |project| {
project
.legacy_project_keys
.push(request.legacy_project_key.clone());
Ok(())
}) {
Ok(project) => {
state.account_sink.record(
None,
&AgentEvent::ProjectUpdated {
project_id: project.id.to_string(),
revision: project.revision,
},
);
project
}
Err(error) => return Ok(project_error(error)),
}
};
match state
.project_store
.migrate_legacy_memory(&id, &request.legacy_project_key)
{
Ok(report) => {
let authoritative = match state.project_store.get(&id) {
Ok(project) => project,
Err(error) => return Ok(project_error(error)),
};
if authoritative.revision != project.revision {
state.account_sink.record(
None,
&AgentEvent::ProjectUpdated {
project_id: authoritative.id.to_string(),
revision: authoritative.revision,
},
);
}
Ok(HttpResponse::Ok()
.insert_header((
actix_web::http::header::ETAG,
format!("\"{}\"", authoritative.revision),
))
.json(serde_json::json!({
"project_id": id,
"project_revision": authoritative.revision,
"migration": report,
})))
}
Err(error) => Ok(project_error(error)),
}
}
pub async fn legacy_memory_migration_status(
state: web::Data<AppState>,
path: web::Path<String>,
query: web::Query<LegacyMemoryMigrationStatusQuery>,
) -> Result<HttpResponse> {
let id = match parse_id(&path) {
Ok(id) => id,
Err(response) => return Ok(response),
};
match state
.project_store
.legacy_memory_migration_status(&id, &query.legacy_project_key)
{
Ok(status) => Ok(HttpResponse::Ok().json(serde_json::json!({
"project_id": id,
"legacy_project_key": query.legacy_project_key,
"migration": status,
}))),
Err(error) => Ok(project_error(error)),
}
}
pub fn is_active(project: &ProjectManifest) -> bool {
project.status == ProjectStatus::Active
}
#[cfg(test)]
mod tests {
use super::*;
use actix_web::{http::header, test, App};
use serde_json::Value;
async fn app_state() -> (tempfile::TempDir, web::Data<AppState>) {
let dir = tempfile::tempdir().unwrap();
let state = AppState::new(dir.path().to_path_buf())
.await
.expect("test AppState");
(dir, web::Data::new(state))
}
macro_rules! project_app {
($state:expr) => {
App::new()
.app_data($state)
.route("/projects", web::get().to(list_projects))
.route("/projects", web::post().to(create_project))
.route("/projects/{id}", web::get().to(get_project))
.route("/projects/{id}", web::patch().to(patch_project))
.route("/projects/{id}/workspaces", web::post().to(bind_workspace))
.route(
"/projects/{id}/workspaces",
web::delete().to(unbind_workspace),
)
.route("/projects/{id}/resources", web::get().to(project_resources))
.route("/projects/{id}/archive", web::post().to(archive_project))
.route(
"/projects/migrations/legacy/dry-run",
web::post().to(legacy_dry_run),
)
.route(
"/projects/{id}/migrations/legacy-memory",
web::post().to(migrate_legacy_memory),
)
.route(
"/projects/{id}/migrations/legacy-memory",
web::get().to(legacy_memory_migration_status),
)
};
}
#[actix_web::test]
async fn project_routes_enforce_etag_cas_and_keep_identity_stable() {
let (dir, state) = app_state().await;
let mut feed = state.account_sink.subscribe();
let app = test::init_service(project_app!(state.clone())).await;
let create = test::call_service(
&app,
test::TestRequest::post()
.uri("/projects")
.set_json(serde_json::json!({
"name": "Zenith",
"description": "first"
}))
.to_request(),
)
.await;
assert_eq!(create.status(), StatusCode::CREATED);
assert_eq!(create.headers().get(header::ETAG).unwrap(), "\"1\"");
let created: ProjectManifest = test::read_body_json(create).await;
let home = state.project_store.paths().project_home(&created.id);
assert!(home.ends_with(created.id.as_str()));
let event = tokio::time::timeout(std::time::Duration::from_secs(1), feed.recv())
.await
.unwrap()
.unwrap();
assert!(matches!(
&event.event,
AgentEvent::ProjectCreated { project_id, .. } if project_id == created.id.as_str()
));
let missing_precondition = test::call_service(
&app,
test::TestRequest::patch()
.uri(&format!("/projects/{}", created.id))
.set_json(serde_json::json!({"name":"Renamed"}))
.to_request(),
)
.await;
assert_eq!(
missing_precondition.status(),
StatusCode::PRECONDITION_REQUIRED
);
let stale = test::call_service(
&app,
test::TestRequest::patch()
.uri(&format!("/projects/{}", created.id))
.insert_header((header::IF_MATCH, "\"9\""))
.set_json(serde_json::json!({"name":"Renamed"}))
.to_request(),
)
.await;
assert_eq!(stale.status(), StatusCode::PRECONDITION_FAILED);
let patched = test::call_service(
&app,
test::TestRequest::patch()
.uri(&format!("/projects/{}", created.id))
.insert_header((header::IF_MATCH, "\"1\""))
.set_json(serde_json::json!({"name":"Renamed"}))
.to_request(),
)
.await;
assert_eq!(patched.status(), StatusCode::OK);
assert_eq!(patched.headers().get(header::ETAG).unwrap(), "\"2\"");
let renamed: ProjectManifest = test::read_body_json(patched).await;
assert_eq!(renamed.id, created.id);
assert_eq!(renamed.name, "Renamed");
assert_eq!(state.project_store.paths().project_home(&renamed.id), home);
let updated_event = tokio::time::timeout(std::time::Duration::from_secs(1), feed.recv())
.await
.unwrap()
.unwrap();
assert!(matches!(
&updated_event.event,
AgentEvent::ProjectUpdated { project_id, .. } if project_id == renamed.id.as_str()
));
let listed =
test::call_service(&app, test::TestRequest::get().uri("/projects").to_request()).await;
assert_eq!(listed.status(), StatusCode::OK);
let listed: Value = test::read_body_json(listed).await;
assert_eq!(listed["projects"][0]["id"], renamed.id.to_string());
let commands = home.join("commands");
std::fs::create_dir_all(&commands).unwrap();
std::fs::write(commands.join("private.md"), "TOP-SECRET-VALUE").unwrap();
let resources = test::call_service(
&app,
test::TestRequest::get()
.uri(&format!("/projects/{}/resources", renamed.id))
.to_request(),
)
.await;
assert_eq!(resources.status(), StatusCode::OK);
let body = test::read_body(resources).await;
let body = String::from_utf8(body.to_vec()).unwrap();
assert!(!body.contains("TOP-SECRET-VALUE"));
assert!(body.contains("resource_revision"));
let legacy_key = "legacy-key";
let legacy_root = dir
.path()
.join("memory/v1/scopes/projects")
.join(legacy_key);
std::fs::create_dir_all(&legacy_root).unwrap();
std::fs::write(legacy_root.join("index.json"), r#"{"legacy":true}"#).unwrap();
let migration_revision = state.project_store.get(&renamed.id).unwrap().revision;
let migration = test::call_service(
&app,
test::TestRequest::post()
.uri(&format!(
"/projects/{}/migrations/legacy-memory",
renamed.id
))
.insert_header((header::IF_MATCH, format!("\"{migration_revision}\"")))
.set_json(serde_json::json!({"legacy_project_key": legacy_key}))
.to_request(),
)
.await;
assert_eq!(migration.status(), StatusCode::OK);
let migration_etag = migration
.headers()
.get(header::ETAG)
.expect("migration ETag")
.to_str()
.unwrap()
.to_string();
let migration_body: Value = test::read_body_json(migration).await;
assert_eq!(migration_body["migration"]["phase"], "committed");
let response_revision = migration_body["project_revision"]
.as_u64()
.expect("project revision");
assert_eq!(migration_etag, format!("\"{response_revision}\""));
assert_eq!(
state.project_store.get(&renamed.id).unwrap().revision,
response_revision,
"endpoint must return the authoritative post-migration CAS revision"
);
let mut observed_final_event = false;
for _ in 0..8 {
let event = tokio::time::timeout(std::time::Duration::from_secs(1), feed.recv())
.await
.unwrap()
.unwrap();
if matches!(
&event.event,
AgentEvent::ProjectUpdated {
project_id,
revision,
} if project_id == renamed.id.as_str() && *revision == response_revision
) {
observed_final_event = true;
break;
}
}
assert!(
observed_final_event,
"post-migration revision must be published to the Project change feed"
);
let before_external_write = state.project_store.get(&renamed.id).unwrap().revision;
std::fs::write(
commands.join("independent.md"),
"independent external change",
)
.unwrap();
let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(3);
loop {
let current = state.project_store.get(&renamed.id).unwrap().revision;
if current > before_external_write {
break;
}
assert!(
tokio::time::Instant::now() < deadline,
"an external resource write interleaved after migration must advance revision"
);
tokio::time::sleep(std::time::Duration::from_millis(25)).await;
}
assert!(legacy_root.join("index.json").is_file());
assert!(home.join("memory/v1/index.json").is_file());
let migration_status = test::call_service(
&app,
test::TestRequest::get()
.uri(&format!(
"/projects/{}/migrations/legacy-memory?legacy_project_key={legacy_key}",
renamed.id
))
.to_request(),
)
.await;
assert_eq!(migration_status.status(), StatusCode::OK);
let status_body: Value = test::read_body_json(migration_status).await;
assert_eq!(status_body["migration"]["phase"], "committed");
let archive_revision = state.project_store.get(&renamed.id).unwrap().revision;
let archive = test::call_service(
&app,
test::TestRequest::post()
.uri(&format!("/projects/{}/archive", renamed.id))
.insert_header((header::IF_MATCH, format!("\"{archive_revision}\"")))
.to_request(),
)
.await;
assert_eq!(archive.status(), StatusCode::OK);
let archived: ProjectManifest = test::read_body_json(archive).await;
assert_eq!(archived.status, ProjectStatus::Archived);
loop {
let event = tokio::time::timeout(std::time::Duration::from_secs(1), feed.recv())
.await
.unwrap()
.unwrap();
if matches!(
&event.event,
AgentEvent::ProjectArchived { project_id, .. }
if project_id == archived.id.as_str()
) {
break;
}
}
let dry_run = test::call_service(
&app,
test::TestRequest::post()
.uri("/projects/migrations/legacy/dry-run")
.set_json(serde_json::json!({"sessions":[]}))
.to_request(),
)
.await;
assert_eq!(dry_run.status(), StatusCode::OK);
drop(dir);
}
#[actix_web::test]
async fn workspace_binding_routes_are_cas_guarded_and_conflict_safe() {
let (dir, state) = app_state().await;
let app = test::init_service(project_app!(state.clone())).await;
let workspace = dir.path().join("workspace");
std::fs::create_dir_all(&workspace).unwrap();
let first = state.project_store.create("First", None).unwrap();
let second = state.project_store.create("Second", None).unwrap();
let bind = test::call_service(
&app,
test::TestRequest::post()
.uri(&format!("/projects/{}/workspaces", first.id))
.insert_header((header::IF_MATCH, format!("\"{}\"", first.revision)))
.set_json(serde_json::json!({"path": workspace}))
.to_request(),
)
.await;
assert_eq!(bind.status(), StatusCode::OK);
let first_bound: ProjectManifest = test::read_body_json(bind).await;
assert_eq!(first_bound.workspace_bindings.len(), 1);
let conflict = test::call_service(
&app,
test::TestRequest::post()
.uri(&format!("/projects/{}/workspaces", second.id))
.insert_header((header::IF_MATCH, format!("\"{}\"", second.revision)))
.set_json(serde_json::json!({"path": workspace}))
.to_request(),
)
.await;
assert_eq!(conflict.status(), StatusCode::CONFLICT);
assert!(state
.project_store
.get(&second.id)
.unwrap()
.workspace_bindings
.is_empty());
let stale_unbind = test::call_service(
&app,
test::TestRequest::delete()
.uri(&format!("/projects/{}/workspaces", first.id))
.insert_header((header::IF_MATCH, "\"1\""))
.set_json(serde_json::json!({"path": workspace}))
.to_request(),
)
.await;
assert_eq!(stale_unbind.status(), StatusCode::PRECONDITION_FAILED);
let unbind = test::call_service(
&app,
test::TestRequest::delete()
.uri(&format!("/projects/{}/workspaces", first.id))
.insert_header((header::IF_MATCH, format!("\"{}\"", first_bound.revision)))
.set_json(serde_json::json!({"path": workspace}))
.to_request(),
)
.await;
assert_eq!(unbind.status(), StatusCode::OK);
}
}