workspace: add typed Ticket API and targets

This commit is contained in:
2026-07-30 22:13:01 +09:00
parent d94ef81b43
commit 1c2d284578
6 changed files with 863 additions and 172 deletions
+33 -1
View File
@@ -10,7 +10,7 @@ use ticket::{
use crate::records::{
ObjectiveDetail, ObjectiveResourceSummary, ObjectiveSummary, ProjectRecordList, TicketDetail,
TicketSummary, summarize_body, truncate_body, validate_project_id,
TicketEventDetail, TicketSummary, summarize_body, truncate_body, validate_project_id,
};
use crate::store::{
ControlPlaneStore, MemoryDocumentRecord, MemoryStagingRecord, MemoryStagingResolutionRecord,
@@ -19,6 +19,8 @@ use crate::store::{
use crate::{Error, Result};
const DETAIL_BODY_LIMIT: usize = 64 * 1024;
const TICKET_EVENT_LIMIT: usize = 100;
const TICKET_EVENT_BODY_LIMIT: usize = 16 * 1024;
const DEFAULT_MEMORY_DOCUMENT_BODY: &str = "# Memory\n\n";
const RECORD_SOURCE_WORKSPACE_SQLITE: &str = "workspace-sqlite";
@@ -244,6 +246,25 @@ impl TicketAuthority for SqliteWorkspaceAuthority {
.show(TicketIdOrSlug::Id(id.to_string()))?;
let (body, body_truncated) =
truncate_body(ticket.document.body.as_str(), DETAIL_BODY_LIMIT);
let event_start = ticket.events.len().saturating_sub(TICKET_EVENT_LIMIT);
let events = ticket.events[event_start..]
.iter()
.enumerate()
.map(|(index, event)| TicketEventDetail {
sequence: event_start + index,
kind: event.kind.as_str().to_owned(),
author: event.author.clone(),
at: event.at.clone(),
status: event.status.clone(),
from: event.from.clone(),
to: event.to.clone(),
reason: event.reason.clone(),
state_field: event.state_field.clone(),
heading: event.heading.clone(),
body: (!event.body.as_str().is_empty())
.then(|| truncate_body(event.body.as_str(), TICKET_EVENT_BODY_LIMIT).0),
})
.collect();
Ok(TicketDetail {
id: ticket.meta.id,
title: ticket.meta.title,
@@ -253,11 +274,22 @@ impl TicketAuthority for SqliteWorkspaceAuthority {
updated_at: ticket.meta.updated_at,
queued_by: ticket.meta.queued_by,
queued_at: ticket.meta.queued_at,
repository_id: ticket.meta.repository_id,
ref_selector: ticket.meta.ref_selector,
risk_flags: ticket.meta.risk_flags,
body,
body_truncated,
event_count: ticket.events.len(),
events,
artifact_count: ticket.artifacts.len(),
artifacts: ticket
.artifacts
.into_iter()
.map(|artifact| artifact.relative_path.display().to_string())
.collect(),
resolution: ticket
.resolution
.map(|resolution| resolution.as_str().to_string()),
record_source: "sqlite_yoi_ticket".to_string(),
})
}
+20
View File
@@ -41,14 +41,34 @@ pub struct TicketDetail {
pub updated_at: Option<String>,
pub queued_by: Option<String>,
pub queued_at: Option<String>,
pub repository_id: Option<String>,
pub ref_selector: Option<String>,
pub risk_flags: Vec<String>,
pub body: String,
pub body_truncated: bool,
pub event_count: usize,
pub events: Vec<TicketEventDetail>,
pub artifact_count: usize,
pub artifacts: Vec<String>,
pub resolution: Option<String>,
pub record_source: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub struct TicketEventDetail {
pub sequence: usize,
pub kind: String,
pub author: Option<String>,
pub at: Option<String>,
pub status: Option<String>,
pub from: Option<String>,
pub to: Option<String>,
pub reason: Option<String>,
pub state_field: Option<String>,
pub heading: Option<String>,
pub body: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub struct ObjectiveSummary {
pub id: String,
+431 -14
View File
@@ -19,6 +19,11 @@ use memory::backend::{
use protocol::stream::{decode_method, encode_event};
use serde::{Deserialize, Serialize};
use sha2::{Digest, Sha256};
use ticket::{
MarkdownText, NewTicketEvent, TicketBackend, TicketBodyReplacement, TicketEventKind,
TicketIdOrSlug, TicketItemEdit, TicketReview, TicketReviewResult, TicketStateChange,
TicketTargetEdit, TicketWorkflowState,
};
use ticket::{
SqliteTicketBackend, TicketBackendHttpResponse, TicketBackendOperation,
execute_ticket_backend_operation,
@@ -511,7 +516,30 @@ pub fn build_router(api: WorkspaceApi) -> Router {
"/api/w/{workspace_id}/skills/{name}/activate",
get(scoped_activate_skill),
)
.route("/api/w/{workspace_id}/tickets/{id}", get(scoped_get_ticket))
.route(
"/api/w/{workspace_id}/tickets/{id}",
get(scoped_get_ticket).patch(scoped_edit_ticket_item),
)
.route(
"/api/w/{workspace_id}/tickets/{id}/state",
post(scoped_transition_ticket_state),
)
.route(
"/api/w/{workspace_id}/tickets/{id}/events",
post(scoped_append_ticket_event),
)
.route(
"/api/w/{workspace_id}/tickets/{id}/reviews",
post(scoped_review_ticket),
)
.route(
"/api/w/{workspace_id}/tickets/{id}/queue",
post(scoped_queue_ticket),
)
.route(
"/api/w/{workspace_id}/tickets/{id}/close",
post(scoped_close_ticket),
)
.route("/api/objectives", get(list_objectives))
.route(
"/api/w/{workspace_id}/objectives",
@@ -1526,6 +1554,220 @@ async fn scoped_get_ticket(
get_ticket(State(api), AxumPath(path.id)).await
}
#[derive(Debug, Deserialize)]
#[serde(deny_unknown_fields)]
struct BrowserEditTicketRequest {
title: Option<String>,
body: Option<String>,
old_string: Option<String>,
new_string: Option<String>,
#[serde(default)]
replace_all: bool,
target: Option<TicketTargetEdit>,
author: Option<String>,
}
#[derive(Debug, Deserialize)]
#[serde(deny_unknown_fields)]
struct BrowserTransitionTicketStateRequest {
state: TicketWorkflowState,
reason: Option<String>,
body: Option<String>,
author: Option<String>,
}
#[derive(Debug, Deserialize)]
#[serde(rename_all = "snake_case")]
enum BrowserTicketThreadRole {
Comment,
Plan,
Decision,
ImplementationReport,
}
impl From<BrowserTicketThreadRole> for TicketEventKind {
fn from(role: BrowserTicketThreadRole) -> Self {
match role {
BrowserTicketThreadRole::Comment => Self::Comment,
BrowserTicketThreadRole::Plan => Self::Plan,
BrowserTicketThreadRole::Decision => Self::Decision,
BrowserTicketThreadRole::ImplementationReport => Self::ImplementationReport,
}
}
}
#[derive(Debug, Deserialize)]
#[serde(deny_unknown_fields)]
struct BrowserAppendTicketEventRequest {
role: BrowserTicketThreadRole,
body: String,
author: Option<String>,
}
#[derive(Debug, Deserialize)]
#[serde(deny_unknown_fields)]
struct BrowserReviewTicketRequest {
result: TicketReviewResult,
body: String,
author: Option<String>,
}
#[derive(Debug, Deserialize)]
#[serde(deny_unknown_fields)]
struct BrowserQueueTicketRequest {
queued_by: Option<String>,
}
#[derive(Debug, Deserialize)]
#[serde(deny_unknown_fields)]
struct BrowserCloseTicketRequest {
resolution: String,
}
fn browser_ticket_backend(api: &WorkspaceApi) -> Result<SqliteTicketBackend> {
let config = ticket::config::TicketConfig::load_workspace(&api.config.workspace_root)
.map_err(|error| Error::Config(format!("load Ticket workspace settings: {error}")))?;
Ok(SqliteTicketBackend::new(
api.config.database_path.clone(),
api.config.workspace_id.clone(),
)
.with_record_language(config.ticket_record_language()))
}
fn browser_ticket_detail(api: &WorkspaceApi, ticket_id: &str) -> ApiResult<Json<TicketDetail>> {
Ok(Json(api.authority.ticket(ticket_id)?))
}
async fn scoped_edit_ticket_item(
State(api): State<WorkspaceApi>,
AxumPath(path): AxumPath<ScopedRecordPath>,
Json(request): Json<BrowserEditTicketRequest>,
) -> ApiResult<Json<TicketDetail>> {
validate_workspace_scope(&api, &path.workspace_id)?;
if let Some(TicketTargetEdit::Set { repository_id, .. }) = request.target.as_ref() {
if !api
.store
.list_repositories(&api.config.workspace_id)?
.iter()
.any(|repository| repository.repository_id == *repository_id)
{
return Err(settings_bad_request(
"unknown_ticket_repository",
"repository_id must identify a repository registered in this Workspace",
));
}
}
browser_ticket_backend(&api)?
.edit_item(
TicketIdOrSlug::Id(path.id.clone()),
TicketItemEdit {
title: request.title,
body: request.body.map(MarkdownText::new),
body_replacement: match (request.old_string, request.new_string) {
(Some(old_string), Some(new_string)) => Some(TicketBodyReplacement {
old_string,
new_string,
replace_all: request.replace_all,
}),
(None, None) => None,
_ => {
return Err(settings_bad_request(
"invalid_ticket_edit_replacement",
"old_string and new_string must be provided together",
));
}
},
target: request.target,
author: request.author,
},
)
.map_err(Error::from)?;
browser_ticket_detail(&api, &path.id)
}
async fn scoped_transition_ticket_state(
State(api): State<WorkspaceApi>,
AxumPath(path): AxumPath<ScopedRecordPath>,
Json(request): Json<BrowserTransitionTicketStateRequest>,
) -> ApiResult<Json<TicketDetail>> {
validate_workspace_scope(&api, &path.workspace_id)?;
let current = api.authority.ticket(&path.id)?;
let mut change = TicketStateChange::new(
current.state,
request.state.as_str(),
request
.reason
.unwrap_or_else(|| "state changed from Web Ticket API".to_owned()),
request.body.unwrap_or_default(),
);
change.author = request.author;
browser_ticket_backend(&api)?
.set_workflow_state(TicketIdOrSlug::Id(path.id.clone()), change)
.map_err(Error::from)?;
browser_ticket_detail(&api, &path.id)
}
async fn scoped_append_ticket_event(
State(api): State<WorkspaceApi>,
AxumPath(path): AxumPath<ScopedRecordPath>,
Json(request): Json<BrowserAppendTicketEventRequest>,
) -> ApiResult<Json<TicketDetail>> {
validate_workspace_scope(&api, &path.workspace_id)?;
let mut event = NewTicketEvent::new(request.role.into(), request.body);
event.author = request.author;
browser_ticket_backend(&api)?
.add_event(TicketIdOrSlug::Id(path.id.clone()), event)
.map_err(Error::from)?;
browser_ticket_detail(&api, &path.id)
}
async fn scoped_review_ticket(
State(api): State<WorkspaceApi>,
AxumPath(path): AxumPath<ScopedRecordPath>,
Json(request): Json<BrowserReviewTicketRequest>,
) -> ApiResult<Json<TicketDetail>> {
validate_workspace_scope(&api, &path.workspace_id)?;
browser_ticket_backend(&api)?
.review(
TicketIdOrSlug::Id(path.id.clone()),
TicketReview {
result: request.result,
body: MarkdownText::new(request.body),
author: request.author,
},
)
.map_err(Error::from)?;
browser_ticket_detail(&api, &path.id)
}
async fn scoped_queue_ticket(
State(api): State<WorkspaceApi>,
AxumPath(path): AxumPath<ScopedRecordPath>,
Json(request): Json<BrowserQueueTicketRequest>,
) -> ApiResult<Json<TicketDetail>> {
validate_workspace_scope(&api, &path.workspace_id)?;
let queued_by = request.queued_by.as_deref().unwrap_or("web");
browser_ticket_backend(&api)?
.queue_ready(TicketIdOrSlug::Id(path.id.clone()), queued_by)
.map_err(Error::from)?;
browser_ticket_detail(&api, &path.id)
}
async fn scoped_close_ticket(
State(api): State<WorkspaceApi>,
AxumPath(path): AxumPath<ScopedRecordPath>,
Json(request): Json<BrowserCloseTicketRequest>,
) -> ApiResult<Json<TicketDetail>> {
validate_workspace_scope(&api, &path.workspace_id)?;
browser_ticket_backend(&api)?
.close(
TicketIdOrSlug::Id(path.id.clone()),
MarkdownText::new(request.resolution),
)
.map_err(Error::from)?;
browser_ticket_detail(&api, &path.id)
}
async fn scoped_ticket_backend_operation(
State(api): State<WorkspaceApi>,
AxumPath(path): AxumPath<ScopedWorkspacePath>,
@@ -4329,19 +4571,35 @@ async fn create_workspace_worker(
},
)
.map_err(|err| err.into_error())?;
Ok(Json(record_browser_worker_spawn(
&api,
request.runtime_id,
display_name,
selected_working_directory_id,
result,
)?))
}
fn record_browser_worker_spawn(
api: &WorkspaceApi,
requested_runtime_id: String,
display_name: String,
selected_working_directory_id: Option<String>,
result: WorkerSpawnResult,
) -> ApiResult<BrowserCreateWorkerResponse> {
if result.state != WorkerOperationState::Accepted {
return Err(worker_create_not_accepted_error(
request.runtime_id.clone(),
requested_runtime_id.clone(),
result.diagnostics,
));
}
let worker = result.worker.ok_or_else(|| Error::RuntimeOperationFailed {
runtime_id: request.runtime_id.clone(),
runtime_id: requested_runtime_id,
code: "workspace_worker_create_missing_summary".to_string(),
message: "Runtime completed worker creation without returning a Worker summary".to_string(),
})?;
let worker_record = record_worker_summary(
&api,
api,
&worker,
display_name.as_str(),
worker.profile.clone(),
@@ -4349,13 +4607,9 @@ async fn create_workspace_worker(
)?;
if let Some(working_directory) = worker.working_directory.as_ref() {
let workdir_record =
workdir_record_from_summary(&api, worker.runtime_id.as_str(), working_directory);
workdir_record_from_summary(api, worker.runtime_id.as_str(), working_directory);
api.store.upsert_workdir_registry(&workdir_record)?;
link_worker_to_workdir(
&api,
&worker_record,
&working_directory.working_directory_id,
)?;
link_worker_to_workdir(api, &worker_record, &working_directory.working_directory_id)?;
}
if let Some(workdir_id) = selected_working_directory_id.as_deref() {
if api
@@ -4370,7 +4624,7 @@ async fn create_workspace_worker(
{
if let Some(status) = result.working_directory {
let record = workdir_record_from_summary(
&api,
api,
worker.runtime_id.as_str(),
&status.summary,
);
@@ -4383,7 +4637,7 @@ async fn create_workspace_worker(
.get_workdir_registry(&api.config.workspace_id, workdir_id)?
.is_some()
{
link_worker_to_workdir(&api, &worker_record, workdir_id)?;
link_worker_to_workdir(api, &worker_record, workdir_id)?;
}
}
let runtime_id = worker.runtime_id.clone();
@@ -4395,14 +4649,14 @@ async fn create_workspace_worker(
encode_path_segment(&runtime_id),
encode_path_segment(&worker_id)
);
Ok(Json(BrowserCreateWorkerResponse {
Ok(BrowserCreateWorkerResponse {
workspace_id,
runtime_id,
worker_id,
console_href,
worker,
diagnostics: result.diagnostics,
}))
})
}
async fn post_internal_runtime_resource_fetch(
@@ -6749,6 +7003,7 @@ fn workspace_id_mismatch_error() -> ApiError {
)
}
#[derive(Debug)]
struct ApiError {
error: Error,
diagnostics: Vec<RuntimeDiagnostic>,
@@ -6762,6 +7017,22 @@ impl From<Error> for ApiError {
severity: DiagnosticSeverity::Error,
message: sanitize_backend_error(message),
}],
Error::Ticket(ticket_error) => vec![RuntimeDiagnostic {
code: match ticket_error {
ticket::TicketError::NotFound(_) => "ticket_not_found",
ticket::TicketError::Ambiguous { .. } => "ticket_ambiguous",
ticket::TicketError::Locked { .. } => "ticket_locked",
ticket::TicketError::Conflict(_) => "ticket_conflict",
ticket::TicketError::InvalidPathComponent(_)
| ticket::TicketError::PathEscapesRoot { .. } => "invalid_ticket_request",
ticket::TicketError::Io { .. }
| ticket::TicketError::Parse { .. }
| ticket::TicketError::Sqlite(_) => "ticket_backend_error",
}
.to_string(),
severity: DiagnosticSeverity::Error,
message: sanitize_backend_error(&ticket_error.to_string()),
}],
_ => Vec::new(),
};
Self { error, diagnostics }
@@ -6778,6 +7049,16 @@ impl IntoResponse for ApiError {
fn into_response(self) -> Response {
let status = match &self.error {
Error::InvalidRuntimeIdentifier { .. } => StatusCode::BAD_REQUEST,
Error::Ticket(ticket::TicketError::NotFound(_)) => StatusCode::NOT_FOUND,
Error::Ticket(
ticket::TicketError::Ambiguous { .. }
| ticket::TicketError::Locked { .. }
| ticket::TicketError::Conflict(_),
) => StatusCode::CONFLICT,
Error::Ticket(
ticket::TicketError::InvalidPathComponent(_)
| ticket::TicketError::PathEscapesRoot { .. },
) => StatusCode::BAD_REQUEST,
Error::InvalidRecordId(_)
| Error::MissingFrontmatter(_)
| Error::UnknownHost(_)
@@ -6896,6 +7177,21 @@ mod tests {
const TEST_REPOSITORY_ID: &str = "main";
const TEST_CREATED_AT: &str = "2026-06-23T06:43:28Z";
#[test]
fn ticket_api_errors_preserve_http_status() {
let not_found = ApiError::from(Error::Ticket(ticket::TicketError::NotFound(
"0000000000000".to_string(),
)))
.into_response();
assert_eq!(not_found.status(), StatusCode::NOT_FOUND);
let conflict = ApiError::from(Error::Ticket(ticket::TicketError::Conflict(
"invalid transition".to_string(),
)))
.into_response();
assert_eq!(conflict.status(), StatusCode::CONFLICT);
}
#[test]
fn backend_worker_projection_preserves_missing_rows_links_and_redacts_paths() {
let worker = WorkerRegistryRecord {
@@ -7669,6 +7965,127 @@ mod tests {
}
}
#[tokio::test]
async fn ticket_browser_endpoints_mutate_typed_backend_and_return_thread() {
let dir = tempfile::tempdir().unwrap();
let api = test_api(dir.path()).await;
let Json(created) = scoped_ticket_backend_operation(
State(api.clone()),
AxumPath(ScopedWorkspacePath {
workspace_id: TEST_WORKSPACE_ID.to_string(),
}),
Json(TicketBackendOperation::Create {
input: ticket::NewTicket::new("Browser Ticket API"),
}),
)
.await
.unwrap();
let ticket_id = match created {
TicketBackendHttpResponse::Ok {
result: ticket::TicketBackendOperationResult::TicketRef(ticket_ref),
} => ticket_ref.id,
other => panic!("unexpected create response: {other:?}"),
};
let path = || ScopedRecordPath {
workspace_id: TEST_WORKSPACE_ID.to_string(),
id: ticket_id.clone(),
};
let Json(edited) = scoped_edit_ticket_item(
State(api.clone()),
AxumPath(path()),
Json(BrowserEditTicketRequest {
title: Some("Browser Ticket API edited".to_string()),
body: Some("Updated from the Browser API.".to_string()),
old_string: None,
new_string: None,
replace_all: false,
target: Some(TicketTargetEdit::Set {
repository_id: "main".to_string(),
ref_selector: Some("feature/api".to_string()),
}),
author: Some("browser-user".to_string()),
}),
)
.await
.unwrap();
assert_eq!(edited.title, "Browser Ticket API edited");
assert_eq!(edited.body, "Updated from the Browser API.");
assert_eq!(edited.repository_id.as_deref(), Some("main"));
assert_eq!(edited.ref_selector.as_deref(), Some("feature/api"));
let Json(commented) = scoped_append_ticket_event(
State(api.clone()),
AxumPath(path()),
Json(BrowserAppendTicketEventRequest {
role: BrowserTicketThreadRole::Comment,
body: "API comment".to_string(),
author: Some("browser-user".to_string()),
}),
)
.await
.unwrap();
assert!(commented.events.iter().any(|event| {
event.kind == "comment" && event.body.as_deref() == Some("API comment")
}));
let Json(ready) = scoped_transition_ticket_state(
State(api.clone()),
AxumPath(path()),
Json(BrowserTransitionTicketStateRequest {
state: TicketWorkflowState::Ready,
reason: Some("intake complete".to_string()),
body: Some("Ready for queue".to_string()),
author: Some("browser-user".to_string()),
}),
)
.await
.unwrap();
assert_eq!(ready.state, "ready");
let Json(queued) = scoped_queue_ticket(
State(api.clone()),
AxumPath(path()),
Json(BrowserQueueTicketRequest {
queued_by: Some("browser-user".to_string()),
}),
)
.await
.unwrap();
assert_eq!(queued.state, "queued");
assert_eq!(queued.queued_by.as_deref(), Some("browser-user"));
let Json(reviewed) = scoped_review_ticket(
State(api.clone()),
AxumPath(path()),
Json(BrowserReviewTicketRequest {
result: TicketReviewResult::Approve,
body: "API review".to_string(),
author: Some("reviewer".to_string()),
}),
)
.await
.unwrap();
assert!(reviewed.events.iter().any(|event| {
event.kind == "review" && event.body.as_deref() == Some("API review")
}));
let Json(closed) = scoped_close_ticket(
State(api),
AxumPath(path()),
Json(BrowserCloseTicketRequest {
resolution: "Closed through the Browser API.".to_string(),
}),
)
.await
.unwrap();
assert_eq!(closed.state, "closed");
assert_eq!(
closed.resolution.as_deref(),
Some("Closed through the Browser API.")
);
}
#[tokio::test]
async fn ticket_backend_endpoint_uses_workspace_sqlite_backend() {
let dir = tempfile::tempdir().unwrap();
+72 -149
View File
@@ -82,6 +82,11 @@ const MIGRATIONS: &[Migration] = &[
name: "objective mutation audit events",
apply: create_objective_event_tables,
},
Migration {
version: 14,
name: "remove unused control-plane Ticket tables",
apply: remove_unused_control_plane_ticket_tables,
},
];
struct Migration {
@@ -2156,6 +2161,20 @@ CREATE TABLE IF NOT EXISTS memory_staging_resolutions (
Ok(())
}
fn remove_unused_control_plane_ticket_tables(conn: &Connection) -> Result<()> {
conn.execute_batch(
r#"
DROP TABLE IF EXISTS ticket_target_paths;
DROP TABLE IF EXISTS ticket_worker_links;
DROP TABLE IF EXISTS ticket_targets;
DROP TABLE IF EXISTS ticket_relations;
DROP TABLE IF EXISTS ticket_events;
DROP TABLE IF EXISTS tickets;
"#,
)?;
Ok(())
}
fn create_objective_event_tables(conn: &Connection) -> Result<()> {
conn.execute_batch(
r#"
@@ -2649,68 +2668,6 @@ CREATE TABLE IF NOT EXISTS workspaces (
updated_at TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS tickets (
workspace_id TEXT NOT NULL REFERENCES workspaces(workspace_id) ON DELETE CASCADE,
ticket_id TEXT PRIMARY KEY,
title TEXT NOT NULL,
state TEXT NOT NULL,
priority TEXT,
assignee_kind TEXT,
assignee_key TEXT,
assignee_display TEXT,
body_md TEXT NOT NULL,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
closed_at TEXT,
resolution_event_id TEXT
);
CREATE TABLE IF NOT EXISTS ticket_events (
workspace_id TEXT NOT NULL REFERENCES workspaces(workspace_id) ON DELETE CASCADE,
event_id TEXT PRIMARY KEY,
ticket_id TEXT NOT NULL REFERENCES tickets(ticket_id) ON DELETE CASCADE,
event_seq INTEGER NOT NULL,
kind TEXT NOT NULL,
activity_id TEXT,
author_kind TEXT NOT NULL,
author_key TEXT NOT NULL,
author_display TEXT NOT NULL,
author_source_kind TEXT,
author_source_key TEXT,
created_at TEXT NOT NULL,
body_md TEXT,
subject_kind TEXT,
subject_id TEXT,
previous_state TEXT,
new_state TEXT,
status TEXT,
artifact_id TEXT,
worker_ref_kind TEXT,
worker_ref_key TEXT,
worker_display TEXT,
host_ref_kind TEXT,
host_ref_key TEXT,
host_display TEXT,
repository_id TEXT,
caused_by_event_id TEXT,
UNIQUE (ticket_id, event_seq)
);
CREATE TABLE IF NOT EXISTS ticket_relations (
workspace_id TEXT NOT NULL REFERENCES workspaces(workspace_id) ON DELETE CASCADE,
source_ticket_id TEXT NOT NULL REFERENCES tickets(ticket_id) ON DELETE CASCADE,
target_ticket_id TEXT NOT NULL REFERENCES tickets(ticket_id) ON DELETE CASCADE,
kind TEXT NOT NULL,
created_at TEXT NOT NULL,
author_kind TEXT NOT NULL,
author_key TEXT NOT NULL,
author_display TEXT NOT NULL,
author_source_kind TEXT,
author_source_key TEXT,
note TEXT,
PRIMARY KEY (source_ticket_id, target_ticket_id, kind)
);
CREATE TABLE IF NOT EXISTS objectives (
workspace_id TEXT NOT NULL REFERENCES workspaces(workspace_id) ON DELETE CASCADE,
objective_id TEXT PRIMARY KEY,
@@ -2793,43 +2750,6 @@ CREATE TABLE IF NOT EXISTS repositories (
updated_at TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS ticket_targets (
workspace_id TEXT NOT NULL REFERENCES workspaces(workspace_id) ON DELETE CASCADE,
ticket_id TEXT NOT NULL REFERENCES tickets(ticket_id) ON DELETE CASCADE,
target_id TEXT NOT NULL,
repository_id TEXT NOT NULL REFERENCES repositories(repository_id) ON DELETE CASCADE,
role TEXT NOT NULL,
intent TEXT NOT NULL,
ref_selector TEXT,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
PRIMARY KEY (ticket_id, target_id)
);
CREATE TABLE IF NOT EXISTS ticket_target_paths (
workspace_id TEXT NOT NULL REFERENCES workspaces(workspace_id) ON DELETE CASCADE,
ticket_id TEXT NOT NULL,
target_id TEXT NOT NULL,
path TEXT NOT NULL,
PRIMARY KEY (ticket_id, target_id, path),
FOREIGN KEY (ticket_id, target_id) REFERENCES ticket_targets(ticket_id, target_id) ON DELETE CASCADE
);
CREATE TABLE IF NOT EXISTS ticket_worker_links (
workspace_id TEXT NOT NULL REFERENCES workspaces(workspace_id) ON DELETE CASCADE,
ticket_id TEXT NOT NULL REFERENCES tickets(ticket_id) ON DELETE CASCADE,
worker_ref_kind TEXT NOT NULL,
worker_ref_key TEXT NOT NULL,
worker_display TEXT,
role TEXT NOT NULL,
status TEXT NOT NULL,
activity_id TEXT,
assigned_at TEXT,
released_at TEXT,
last_event_id TEXT,
PRIMARY KEY (ticket_id, worker_ref_kind, worker_ref_key, role)
);
CREATE TABLE IF NOT EXISTS artifacts (
workspace_id TEXT NOT NULL REFERENCES workspaces(workspace_id) ON DELETE CASCADE,
artifact_id TEXT PRIMARY KEY,
@@ -2904,13 +2824,40 @@ mod tests {
use super::*;
use std::collections::BTreeSet;
#[test]
fn removes_unused_control_plane_ticket_tables() {
let conn = Connection::open_in_memory().unwrap();
conn.execute_batch(
r#"
CREATE TABLE tickets (ticket_id TEXT PRIMARY KEY);
CREATE TABLE ticket_events (event_id TEXT PRIMARY KEY, ticket_id TEXT REFERENCES tickets(ticket_id));
CREATE TABLE ticket_relations (source_ticket_id TEXT, target_ticket_id TEXT);
CREATE TABLE ticket_targets (ticket_id TEXT, target_id TEXT, PRIMARY KEY (ticket_id, target_id));
CREATE TABLE ticket_target_paths (ticket_id TEXT, target_id TEXT, path TEXT);
CREATE TABLE ticket_worker_links (ticket_id TEXT, worker_ref_key TEXT);
"#,
)
.unwrap();
remove_unused_control_plane_ticket_tables(&conn).unwrap();
for table in [
"tickets",
"ticket_events",
"ticket_relations",
"ticket_targets",
"ticket_target_paths",
"ticket_worker_links",
] {
assert!(!table_exists(&conn, table).unwrap(), "{table} still exists");
}
}
#[tokio::test]
async fn migrates_sqlite_and_preserves_workspace_record() {
let dir = tempfile::tempdir().unwrap();
let db = dir.path().join("control-plane.sqlite");
let store = SqliteWorkspaceStore::open(&db).unwrap();
assert_eq!(store.schema_version().await.unwrap(), 13);
assert_eq!(store.schema_version().await.unwrap(), 14);
let record = WorkspaceRecord {
workspace_id: "local-dev".to_string(),
@@ -2923,7 +2870,7 @@ mod tests {
store.upsert_workspace(&record).await.unwrap();
let reopened = SqliteWorkspaceStore::open(&db).unwrap();
assert_eq!(reopened.schema_version().await.unwrap(), 13);
assert_eq!(reopened.schema_version().await.unwrap(), 14);
assert_eq!(
reopened.get_workspace("local-dev").await.unwrap(),
Some(record)
@@ -2939,9 +2886,6 @@ mod tests {
let tables = table_names(&conn);
for expected in [
"workspaces",
"tickets",
"ticket_events",
"ticket_relations",
"objectives",
"objective_ticket_links",
"objective_resources",
@@ -2949,9 +2893,6 @@ mod tests {
"workspace_memory_documents",
"memory_staging_resolutions",
"repositories",
"ticket_targets",
"ticket_target_paths",
"ticket_worker_links",
"artifacts",
"audit_events",
"worker_registry",
@@ -2977,6 +2918,12 @@ mod tests {
"actors",
"validation_results",
"ci_results",
"tickets",
"ticket_events",
"ticket_relations",
"ticket_targets",
"ticket_target_paths",
"ticket_worker_links",
] {
assert!(
!tables.contains(forbidden),
@@ -3017,39 +2964,6 @@ mod tests {
"updated_at",
],
);
assert_columns(
&conn,
"ticket_events",
[
"workspace_id",
"event_id",
"ticket_id",
"event_seq",
"kind",
"activity_id",
"author_kind",
"author_key",
"author_display",
"author_source_kind",
"author_source_key",
"created_at",
"body_md",
"subject_kind",
"subject_id",
"previous_state",
"new_state",
"status",
"artifact_id",
"worker_ref_kind",
"worker_ref_key",
"worker_display",
"host_ref_kind",
"host_ref_key",
"host_display",
"repository_id",
"caused_by_event_id",
],
);
assert_columns(
&conn,
"worker_registry",
@@ -3098,7 +3012,7 @@ mod tests {
],
);
for table in ["workspaces", "repositories", "ticket_events", "artifacts"] {
for table in ["workspaces", "repositories", "artifacts"] {
let columns = table_columns(&conn, table).unwrap();
for forbidden_column in [
"payload",
@@ -3144,7 +3058,7 @@ mod tests {
.unwrap();
let store = SqliteWorkspaceStore::from_connection(conn).unwrap();
assert_eq!(store.schema_version().await.unwrap(), 13);
assert_eq!(store.schema_version().await.unwrap(), 14);
store
.with_conn(|conn| {
@@ -3152,9 +3066,6 @@ mod tests {
for expected in [
"workspaces",
"repositories",
"tickets",
"ticket_events",
"ticket_worker_links",
"artifacts",
"audit_events",
"workspace_memory_documents",
@@ -3172,7 +3083,19 @@ mod tests {
"missing {expected} after upgrade"
);
}
for forbidden in ["runs", "hosts", "workers", "actors", "validation_results"] {
for forbidden in [
"runs",
"hosts",
"workers",
"actors",
"validation_results",
"tickets",
"ticket_events",
"ticket_relations",
"ticket_targets",
"ticket_target_paths",
"ticket_worker_links",
] {
assert!(
!tables.contains(forbidden),
"upgraded schema must not retain forbidden canonical table {forbidden}"
@@ -3238,7 +3161,7 @@ mod tests {
#[tokio::test]
async fn repository_records_round_trip() {
let store = SqliteWorkspaceStore::in_memory().unwrap();
assert_eq!(store.schema_version().await.unwrap(), 13);
assert_eq!(store.schema_version().await.unwrap(), 14);
let workspace = WorkspaceRecord {
workspace_id: "local-dev".to_string(),
owner_account_id: None,
@@ -3276,7 +3199,7 @@ mod tests {
#[tokio::test]
async fn memory_authority_records_round_trip_and_close_staging() {
let store = SqliteWorkspaceStore::in_memory().unwrap();
assert_eq!(store.schema_version().await.unwrap(), 13);
assert_eq!(store.schema_version().await.unwrap(), 14);
let workspace = WorkspaceRecord {
workspace_id: "local-dev".to_string(),
owner_account_id: None,
@@ -3450,7 +3373,7 @@ mod tests {
#[tokio::test]
async fn account_and_login_records_round_trip() {
let store = SqliteWorkspaceStore::in_memory().unwrap();
assert_eq!(store.schema_version().await.unwrap(), 13);
assert_eq!(store.schema_version().await.unwrap(), 14);
let now = "2026-07-22T00:00:00Z".to_string();
let account = AccountRecord {
account_id: "acct-user-alice".to_string(),