diff --git a/crates/ticket/src/lib.rs b/crates/ticket/src/lib.rs index a7829c56..7dd923a9 100644 --- a/crates/ticket/src/lib.rs +++ b/crates/ticket/src/lib.rs @@ -120,6 +120,26 @@ fn io_err(path: impl Into, source: io::Error) -> TicketError { } } +fn read_ticket_summary_row(row: &rusqlite::Row<'_>) -> rusqlite::Result { + let workflow_state = row.get::<_, String>(7)?; + Ok(TicketSummary { + id: row.get(0)?, + slug: row.get(1)?, + title: row.get(2)?, + status: ExtensibleTicketStatus::from(row.get::<_, String>(3)?.as_str()), + kind: row.get(4)?, + priority: row.get(5)?, + labels: Vec::new(), + readiness: row.get(6)?, + workflow_state: TicketWorkflowState::parse(&workflow_state) + .unwrap_or(TicketWorkflowState::Planning), + workflow_state_explicit: row.get::<_, i64>(8)? != 0, + queued_by: row.get(9)?, + queued_at: row.get(10)?, + updated_at: row.get(11)?, + }) +} + fn sqlite_err(error: impl std::fmt::Display) -> TicketError { TicketError::Sqlite(error.to_string()) } @@ -1572,6 +1592,27 @@ pub struct SqliteTicketListProjection { pub items: Vec, } +#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq, Eq)] +pub struct SqliteTicketListCursor { + pub state_rank: i64, + pub updated_at: Option, + pub ticket_id: String, +} + +#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq, Eq)] +pub struct SqliteTicketListPageQuery { + pub states: Vec, + pub limit: usize, + pub after: Option, +} + +#[derive(Debug, Clone, Default, PartialEq, Eq)] +pub struct SqliteTicketListPage { + pub items: Vec, + pub has_more: bool, + pub next: Option, +} + #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] pub struct TicketInvalidRecord { pub label: String, @@ -2554,6 +2595,112 @@ impl SqliteTicketBackend { }) } + /// Lists one stable keyset-paginated Workspace summary page. + /// + /// Filtering, ordering, and `limit + 1` are applied by SQLite before the bounded blocker + /// hydration query. The returned cursor is storage data; callers must wrap it in their own + /// opaque, query-bound transport cursor. + pub fn list_workspace_projection_page( + &self, + query: SqliteTicketListPageQuery, + ) -> Result { + self.with_read(|conn| { + let states = serde_json::to_string( + &query + .states + .iter() + .map(|state| state.as_str()) + .collect::>(), + ) + .map_err(|error| TicketError::Sqlite(error.to_string()))?; + let cursor_rank = query.after.as_ref().map(|cursor| cursor.state_rank); + let cursor_updated_at = query + .after + .as_ref() + .and_then(|cursor| cursor.updated_at.clone()); + let cursor_id = query + .after + .as_ref() + .map(|cursor| cursor.ticket_id.as_str()); + let fetch_limit = query.limit.saturating_add(1); + let mut statement = conn + .prepare( + "SELECT ticket_id, slug, title, status, kind, priority, readiness, + workflow_state, workflow_state_explicit, queued_by, queued_at, updated_at + FROM typed_tickets AS ticket + WHERE workspace_id = ?1 + AND (json_array_length(?2) = 0 OR EXISTS ( + SELECT 1 FROM json_each(?2) AS state + WHERE state.value = ticket.workflow_state + )) + AND (?3 IS NULL OR + CASE ticket.workflow_state + WHEN 'ready' THEN 0 WHEN 'planning' THEN 1 + WHEN 'inprogress' THEN 2 WHEN 'queued' THEN 3 + WHEN 'done' THEN 4 WHEN 'closed' THEN 5 ELSE 6 END > ?4 + OR (CASE ticket.workflow_state + WHEN 'ready' THEN 0 WHEN 'planning' THEN 1 + WHEN 'inprogress' THEN 2 WHEN 'queued' THEN 3 + WHEN 'done' THEN 4 WHEN 'closed' THEN 5 ELSE 6 END = ?4 + AND (COALESCE(ticket.updated_at, '') < COALESCE(?5, '') + OR (COALESCE(ticket.updated_at, '') = COALESCE(?5, '') + AND ticket.ticket_id > ?3)))) + ORDER BY CASE ticket.workflow_state + WHEN 'ready' THEN 0 WHEN 'planning' THEN 1 + WHEN 'inprogress' THEN 2 WHEN 'queued' THEN 3 + WHEN 'done' THEN 4 WHEN 'closed' THEN 5 ELSE 6 END ASC, + ticket.updated_at DESC, ticket.ticket_id ASC + LIMIT ?6", + ) + .map_err(sqlite_err)?; + let rows = statement + .query_map( + params![ + self.workspace_id, + states, + cursor_id, + cursor_rank, + cursor_updated_at, + i64::try_from(fetch_limit).unwrap_or(i64::MAX) + ], + read_ticket_summary_row, + ) + .map_err(sqlite_err)?; + let mut summaries = rows + .collect::, _>>() + .map_err(sqlite_err)?; + let has_more = summaries.len() > query.limit; + summaries.truncate(query.limit); + let next = has_more.then(|| { + let summary = summaries.last().expect("non-empty page with continuation"); + SqliteTicketListCursor { + state_rank: match summary.workflow_state { + TicketWorkflowState::Ready => 0, + TicketWorkflowState::Planning => 1, + TicketWorkflowState::InProgress => 2, + TicketWorkflowState::Queued => 3, + TicketWorkflowState::Done => 4, + TicketWorkflowState::Closed => 5, + }, + updated_at: summary.updated_at.clone(), + ticket_id: summary.id.clone(), + } + }); + let blockers = self.list_workspace_blockers(conn, &summaries)?; + Ok(SqliteTicketListPage { + items: summaries + .into_iter() + .map(|summary| SqliteTicketListItem { + relation_blockers: blockers.get(&summary.id).cloned().unwrap_or_default(), + summary, + }) + .collect(), + has_more, + next, + }) + }) + } + fn list_workspace_summaries( &self, conn: &Connection, @@ -7234,6 +7381,86 @@ state: planning assert!(!debug.contains("full-artifact-marker")); } + #[test] + fn sqlite_workspace_projection_page_uses_stable_keyset_and_state_filter() { + let tmp = TempDir::new().unwrap(); + let db_path = tmp.path().join("workspace.db"); + let backend = SqliteTicketBackend::open(&db_path, "workspace-test").unwrap(); + let mut ids = Vec::new(); + for (title, state, updated_at) in [ + ("Newest", TicketWorkflowState::Ready, "2026-08-12T03:00:00Z"), + ("Middle", TicketWorkflowState::Ready, "2026-08-12T02:00:00Z"), + ( + "Oldest", + TicketWorkflowState::Planning, + "2026-08-12T01:00:00Z", + ), + ] { + let mut input = NewTicket::new(title); + input.workflow_state = Some(state); + let ticket = backend.create(input).unwrap(); + Connection::open(&db_path) + .unwrap() + .execute( + "UPDATE typed_tickets SET updated_at=?3 WHERE workspace_id=?1 AND ticket_id=?2", + params!["workspace-test", ticket.id, updated_at], + ) + .unwrap(); + ids.push(ticket.id); + } + + Connection::open(&db_path) + .unwrap() + .execute( + "UPDATE typed_tickets SET updated_at='2026-08-12T05:00:00Z' WHERE workspace_id=?1 AND ticket_id=?2", + params!["workspace-test", ids[2]], + ) + .unwrap(); + let combined_lane = backend + .list_workspace_projection_page(SqliteTicketListPageQuery { + states: vec![TicketWorkflowState::Ready, TicketWorkflowState::Planning], + limit: 1, + after: None, + }) + .unwrap(); + assert_eq!( + combined_lane.items[0].summary.workflow_state, + TicketWorkflowState::Ready, + "lane pagination order must match the UI's state-primary order", + ); + + let first = backend + .list_workspace_projection_page(SqliteTicketListPageQuery { + states: vec![TicketWorkflowState::Ready], + limit: 1, + after: None, + }) + .unwrap(); + assert_eq!(first.items[0].summary.id, ids[0]); + assert!(first.has_more); + + let mut inserted = NewTicket::new("Inserted"); + inserted.workflow_state = Some(TicketWorkflowState::Ready); + let inserted = backend.create(inserted).unwrap(); + Connection::open(&db_path) + .unwrap() + .execute( + "UPDATE typed_tickets SET updated_at='2026-08-12T04:00:00Z' WHERE workspace_id=?1 AND ticket_id=?2", + params!["workspace-test", inserted.id], + ) + .unwrap(); + + let second = backend + .list_workspace_projection_page(SqliteTicketListPageQuery { + states: vec![TicketWorkflowState::Ready], + limit: 1, + after: first.next, + }) + .unwrap(); + assert_eq!(second.items[0].summary.id, ids[1]); + assert!(!second.has_more); + } + #[test] fn sqlite_workspace_projection_sql_shape_is_constant_for_ticket_count() { let source = include_str!("lib.rs"); @@ -7245,8 +7472,8 @@ state: planning .map(|offset| start + offset) .expect("following method"); let projection_source = &source[start..end]; - assert_eq!(projection_source.matches("self.with_read(").count(), 1); - assert_eq!(projection_source.matches(".prepare(").count(), 2); + assert_eq!(projection_source.matches("self.with_read(").count(), 2); + assert_eq!(projection_source.matches(".prepare(").count(), 3); let item_loop = projection_source .split("items: summaries") .nth(1) diff --git a/crates/ticket/src/sqlite_schema.rs b/crates/ticket/src/sqlite_schema.rs index 4f876294..ebbbf8f2 100644 --- a/crates/ticket/src/sqlite_schema.rs +++ b/crates/ticket/src/sqlite_schema.rs @@ -7,7 +7,7 @@ use crate::{Result, TicketError, sqlite_err}; const MIGRATION_TABLE: &str = "ticket_schema_migrations"; const MAX_SCHEMA_DIAGNOSTICS: usize = 32; -pub const LATEST_SQLITE_TICKET_SCHEMA_VERSION: i64 = 3; +pub const LATEST_SQLITE_TICKET_SCHEMA_VERSION: i64 = 4; #[derive(Clone, Copy)] struct Migration { @@ -32,6 +32,11 @@ const MIGRATIONS: &[Migration] = &[ name: "convert_legacy_reviews_to_comments", apply: retire_legacy_ticket_review_events, }, + Migration { + version: 4, + name: "add_ticket_query_indexes", + apply: add_ticket_query_indexes, + }, ]; #[derive(Clone, Copy)] @@ -507,6 +512,29 @@ fn retire_legacy_ticket_review_events(connection: &Connection) -> Result<()> { .map_err(sqlite_err) } +fn add_ticket_query_indexes(connection: &Connection) -> Result<()> { + connection + .execute_batch( + r#" + CREATE INDEX IF NOT EXISTS typed_tickets_workspace_state_updated + ON typed_tickets(workspace_id, workflow_state, updated_at DESC, ticket_id); + CREATE INDEX IF NOT EXISTS typed_tickets_workspace_updated + ON typed_tickets(workspace_id, updated_at DESC, ticket_id); + CREATE INDEX IF NOT EXISTS typed_tickets_workspace_created + ON typed_tickets(workspace_id, created_at DESC, ticket_id); + CREATE INDEX IF NOT EXISTS typed_tickets_workspace_title + ON typed_tickets(workspace_id, title COLLATE NOCASE, ticket_id); + CREATE INDEX IF NOT EXISTS typed_ticket_events_workspace_kind_ticket + ON typed_ticket_events(workspace_id, kind, ticket_id, event_index); + CREATE INDEX IF NOT EXISTS typed_ticket_relations_workspace_source_kind + ON typed_ticket_relations(workspace_id, ticket_id, kind, target); + CREATE INDEX IF NOT EXISTS typed_ticket_relations_workspace_target_kind + ON typed_ticket_relations(workspace_id, target, kind, ticket_id); + "#, + ) + .map_err(sqlite_err) +} + fn add_column_if_missing( connection: &Connection, table: &str, @@ -841,10 +869,10 @@ mod tests { verify_sqlite_ticket_schema(&connection).unwrap(); let versions = load_applied_migrations(&connection).unwrap(); - assert_eq!(versions.len(), 3); + assert_eq!(versions.len(), 4); assert_eq!( versions.get(&LATEST_SQLITE_TICKET_SCHEMA_VERSION), - Some(&"convert_legacy_reviews_to_comments".to_string()) + Some(&"add_ticket_query_indexes".to_string()) ); } @@ -989,7 +1017,7 @@ mod tests { .to_string() .contains("unsupported Ticket schema migration version 99") ); - assert_eq!(load_applied_migrations(&connection).unwrap().len(), 4); + assert_eq!(load_applied_migrations(&connection).unwrap().len(), 5); } #[test] @@ -1090,7 +1118,7 @@ mod tests { connection.execute("INSERT INTO typed_ticket_events (workspace_id,ticket_id,event_index,kind,author,at,status,heading,body) VALUES ('workspace-1','ticket-1',0,'review','reviewer','2026-08-11T00:00:00Z','approve','Review','legacy evidence')",[]).unwrap(); connection.execute("INSERT INTO typed_ticket_event_attributes (workspace_id,ticket_id,event_index,key,value) VALUES ('workspace-1','ticket-1',0,'result','approve')",[]).unwrap(); connection - .execute("DELETE FROM ticket_schema_migrations WHERE version=3", []) + .execute("DELETE FROM ticket_schema_migrations WHERE version>=3", []) .unwrap(); migrate_sqlite_ticket_schema(&connection).unwrap(); let (kind,status,heading,body):(String,Option,Option,Option)=connection.query_row("SELECT kind,status,heading,body FROM typed_ticket_events WHERE workspace_id='workspace-1' AND ticket_id='ticket-1' AND event_index=0",[],|row|Ok((row.get(0)?,row.get(1)?,row.get(2)?,row.get(3)?))).unwrap(); @@ -1129,6 +1157,6 @@ mod tests { let connection = Connection::open(database).unwrap(); verify_sqlite_ticket_schema(&connection).unwrap(); - assert_eq!(load_applied_migrations(&connection).unwrap().len(), 3); + assert_eq!(load_applied_migrations(&connection).unwrap().len(), 4); } } diff --git a/crates/workspace-server/src/authority.rs b/crates/workspace-server/src/authority.rs index 9070db63..0c8f0abb 100644 --- a/crates/workspace-server/src/authority.rs +++ b/crates/workspace-server/src/authority.rs @@ -6,9 +6,11 @@ use merge_request::{ ReviewDecision, }; use project_record::{allocate_record_id, unix_epoch_millis_now}; +use rusqlite::{params_from_iter, types::Value as SqlValue}; use ticket::{ - SqliteTicketBackend, TicketBackend, TicketEvent, TicketIdOrSlug, TicketWorkspaceActionPriority, + SqliteTicketBackend, SqliteTicketListCursor, SqliteTicketListItem, SqliteTicketListPageQuery, + TicketBackend, TicketEvent, TicketIdOrSlug, TicketWorkflowState, TicketWorkspaceActionPriority, project_ticket_workspace_item, }; @@ -17,8 +19,9 @@ use crate::records::{ ObjectiveQueryItem, ObjectiveQueryRequest, ObjectiveQueryResponse, ObjectiveResourceSummary, ObjectiveShowRequest, ObjectiveSummary, ProjectRecordList, QueryPage, TicketAssignmentSummary, TicketDetail, TicketEventDetail, TicketEvidenceEvent, TicketEvidenceSummary, - TicketMergeRequestSummary, TicketQueryItem, TicketQueryRequest, TicketQueryResponse, - TicketShowRequest, TicketSummary, summarize_body, truncate_body, validate_project_id, + TicketListPageRequest, TicketMergeRequestSummary, TicketQueryItem, TicketQueryRequest, + TicketQueryResponse, TicketShowRequest, TicketSummary, TicketSummaryPage, summarize_body, + truncate_body, validate_project_id, }; use crate::store::{ ControlPlaneStore, MemoryDocumentRecord, MemoryStagingRecord, MemoryStagingResolutionRecord, @@ -43,6 +46,7 @@ impl WorkspaceAuthority for T where T: ObjectiveAuthority + TicketAuthority + pub trait TicketAuthority { fn list_tickets(&self, limit: usize) -> Result>; + fn list_ticket_page(&self, request: TicketListPageRequest) -> Result; fn query_tickets(&self, query: TicketQueryRequest) -> Result; fn ticket(&self, id: &str) -> Result; fn show_ticket(&self, id: &str, query: TicketShowRequest) -> Result; @@ -317,11 +321,358 @@ impl SqliteWorkspaceAuthority { }) } + fn query_ticket_candidate_ids( + &self, + query: &TicketQueryRequest, + sort: TicketQuerySort, + after: Option<&(String, String)>, + limit: usize, + ) -> Result> { + let mut values = vec![SqlValue::Text(self.workspace_id.clone())]; + let mut predicates = vec!["t.workspace_id=?1".to_string()]; + let mut bind = |value: SqlValue| { + values.push(value); + format!("?{}", values.len()) + }; + if !query.states.is_empty() { + let states = query + .states + .iter() + .map(|state| bind(SqlValue::Text(state.clone()))) + .collect::>(); + predicates.push(format!("t.workflow_state IN ({})", states.join(","))); + } + if let Some(text) = query.query.as_deref().filter(|text| !text.is_empty()) { + let pattern = bind(SqlValue::Text(format!("%{}%", text.to_lowercase()))); + predicates.push(format!( + "(lower(t.title) LIKE {p} OR lower(t.body) LIKE {p} OR EXISTS ( + SELECT 1 FROM typed_ticket_events e + WHERE e.workspace_id=t.workspace_id AND e.ticket_id=t.ticket_id + AND lower(e.body) LIKE {p}))", + p = pattern + )); + } + if let Some(value) = &query.updated_after { + let value = bind(SqlValue::Text(value.clone())); + predicates.push(format!("COALESCE(t.updated_at,'')>{value}")); + } + if let Some(value) = &query.updated_before { + let value = bind(SqlValue::Text(value.clone())); + predicates.push(format!("COALESCE(t.updated_at,'')<{value}")); + } + if let Some(value) = &query.linked_objective_id { + let value = bind(SqlValue::Text(value.clone())); + predicates.push(format!("EXISTS (SELECT 1 FROM objective_ticket_links link WHERE link.workspace_id=t.workspace_id AND link.ticket_id=t.ticket_id AND link.objective_id={value})")); + } + if query.related_ticket_id.is_some() || query.relation_kind.is_some() { + let related = query + .related_ticket_id + .as_ref() + .map(|value| bind(SqlValue::Text(value.clone()))); + let kind = query + .relation_kind + .as_ref() + .map(|value| bind(SqlValue::Text(value.clone()))); + let related = related + .map(|value| format!("AND ((r.ticket_id=t.ticket_id AND r.target={value}) OR (r.target=t.ticket_id AND r.ticket_id={value}))")) + .unwrap_or_else(|| { + "AND (r.ticket_id=t.ticket_id OR r.target=t.ticket_id)".to_string() + }); + let kind = kind + .map(|value| format!("AND r.kind={value}")) + .unwrap_or_default(); + predicates.push(format!("EXISTS (SELECT 1 FROM typed_ticket_relations r WHERE r.workspace_id=t.workspace_id {related} {kind})")); + } + let blocker = "EXISTS (SELECT 1 FROM typed_ticket_relations relation + JOIN typed_tickets blocker ON blocker.workspace_id=relation.workspace_id + AND blocker.ticket_id=CASE WHEN relation.ticket_id=t.ticket_id THEN relation.target ELSE relation.ticket_id END + WHERE relation.workspace_id=t.workspace_id + AND ((relation.ticket_id=t.ticket_id AND relation.kind='depends_on') + OR (relation.target=t.ticket_id AND relation.kind='blocks')) + AND blocker.workflow_state NOT IN ('done','closed'))"; + let active_blocker = "EXISTS (SELECT 1 FROM typed_ticket_relations relation + JOIN typed_tickets blocker ON blocker.workspace_id=relation.workspace_id + AND blocker.ticket_id=CASE WHEN relation.ticket_id=t.ticket_id THEN relation.target ELSE relation.ticket_id END + WHERE relation.workspace_id=t.workspace_id + AND ((relation.ticket_id=t.ticket_id AND relation.kind='depends_on') + OR (relation.target=t.ticket_id AND relation.kind='blocks')) + AND blocker.workflow_state NOT IN ('queued','inprogress','done','closed'))"; + let report_index = "(SELECT max(event.event_index) FROM typed_ticket_events event WHERE event.workspace_id=t.workspace_id AND event.ticket_id=t.ticket_id AND event.kind='implementation_report')"; + let edit_index = "(SELECT max(event.event_index) FROM typed_ticket_events event WHERE event.workspace_id=t.workspace_id AND event.ticket_id=t.ticket_id AND event.kind='item_edit')"; + let current_report = format!( + "({report_index} IS NOT NULL AND ({edit_index} IS NULL OR {report_index}>={edit_index}))" + ); + let merge_request_id = "(SELECT relation.merge_request_id FROM merge_request_ticket_relations relation + JOIN merge_requests request ON request.workspace_id=relation.workspace_id AND request.merge_request_id=relation.merge_request_id + WHERE relation.workspace_id=t.workspace_id AND relation.ticket_id=t.ticket_id + ORDER BY CASE WHEN request.state='open' THEN 0 ELSE 1 END, request.created_at DESC LIMIT 1)"; + let review_subject = format!( + "(SELECT json_extract(requested.payload_json,'$.subject_ref') + FROM merge_request_thread_events requested + WHERE requested.workspace_id=t.workspace_id + AND requested.merge_request_id={merge_request_id} + AND requested.kind='review_requested' + ORDER BY requested.sequence DESC LIMIT 1)" + ); + let review_decision = format!( + "(SELECT json_extract(event.payload_json,'$.decision') + FROM merge_request_thread_events event + WHERE event.workspace_id=t.workspace_id + AND event.merge_request_id={merge_request_id} AND event.kind='review' + AND json_extract(event.payload_json,'$.subject_ref')={review_subject} + AND NOT EXISTS (SELECT 1 FROM merge_request_thread_events revoked + WHERE revoked.workspace_id=event.workspace_id + AND revoked.merge_request_id=event.merge_request_id + AND revoked.kind='review_revoked' + AND json_extract(revoked.payload_json,'$.review_event_id')=event.event_id) + ORDER BY event.sequence DESC LIMIT 1)" + ); + let review_status = format!( + "CASE WHEN {merge_request_id} IS NULL THEN 'none' WHEN {review_decision}='approve' THEN 'approved' WHEN {review_decision}='request_changes' THEN 'request_changes' ELSE 'pending' END" + ); + let has_commit = format!( + "({review_subject} IS NOT NULL OR EXISTS (SELECT 1 FROM typed_ticket_event_references reference WHERE reference.workspace_id=t.workspace_id AND reference.ticket_id=t.ticket_id AND reference.kind='commit'))" + ); + if !query.event_kinds.is_empty() { + let event_kinds = query + .event_kinds + .iter() + .map(|event_kind| bind(SqlValue::Text(event_kind.clone()))) + .collect::>(); + predicates.push(format!("EXISTS (SELECT 1 FROM typed_ticket_events event WHERE event.workspace_id=t.workspace_id AND event.ticket_id=t.ticket_id AND event.kind IN ({}))", event_kinds.join(","))); + } + for evidence in &query.evidence { + predicates.push(match evidence.as_str() { + "implementation_report" => format!("{report_index} IS NOT NULL"), + "implementation_report_after_rescope" => current_report.clone(), + "merge_request" => format!("{merge_request_id} IS NOT NULL"), + "commit" => has_commit.clone(), + "approved_review" => format!("{review_status}='approved'"), + other => { + return Err(Error::InvalidRecordId(format!( + "unsupported evidence filter `{other}`" + ))); + } + }); + } + if let Some(status) = &query.review_status { + let status = if matches!(status.as_str(), "unresolved_changes" | "changes_requested") { + "request_changes" + } else { + status.as_str() + }; + let status = bind(SqlValue::Text(status.to_string())); + predicates.push(format!("{review_status}={status}")); + } + for attention in &query.attention { + predicates.push(match attention.as_str() { + "done_not_closed" => "t.workflow_state='done'".to_string(), + "implementation_report_not_closed" => { + format!("{report_index} IS NOT NULL AND t.workflow_state!='closed'") + } + "report_after_rescope" => current_report.clone(), + "unresolved_review" | "unresolved_changes" => { + format!("{review_status}='request_changes'") + } + "missing_commit" => format!("NOT {has_commit}"), + "blocked" => blocker.to_string(), + "unblocked" => format!("NOT {blocker}"), + "ready" => format!("t.workflow_state='ready' AND NOT {blocker}"), + "awaiting_review" => format!("{review_status}='pending'"), + "stale_after_rescope" => { + format!("{report_index} IS NOT NULL AND NOT {current_report}") + } + "missing_evidence" => format!( + "NOT ({current_report} AND {has_commit} AND {review_status}='approved')" + ), + other => { + return Err(Error::InvalidRecordId(format!( + "unsupported attention filter `{other}`" + ))); + } + }); + } + let rank_expression = match sort { + TicketQuerySort::Priority => format!( + "CASE WHEN t.workflow_state='ready' AND NOT {active_blocker} THEN 0 WHEN t.workflow_state IN ('queued','inprogress') THEN 1 ELSE 2 END" + ), + TicketQuerySort::Relevance => { + if let Some(text) = query.query.as_deref().filter(|text| !text.is_empty()) { + let pattern = bind(SqlValue::Text(format!("%{}%", text.to_lowercase()))); + format!( + "CASE WHEN lower(t.title) LIKE {pattern} THEN 0 WHEN lower(t.body) LIKE {pattern} THEN 1 WHEN EXISTS (SELECT 1 FROM typed_ticket_events event WHERE event.workspace_id=t.workspace_id AND event.ticket_id=t.ticket_id AND lower(event.body) LIKE {pattern}) THEN 2 ELSE 3 END" + ) + } else { + "3".to_string() + } + } + _ => "0".to_string(), + }; + if let Some((key, id)) = after { + match sort { + TicketQuerySort::UpdatedDesc => { + let key = bind(SqlValue::Text(key.clone())); + let id = bind(SqlValue::Text(id.clone())); + predicates.push(format!("(COALESCE(t.updated_at,'')<{key} OR (COALESCE(t.updated_at,'')={key} AND t.ticket_id>{id}))")); + } + TicketQuerySort::CreatedDesc => { + let key = bind(SqlValue::Text(key.clone())); + let id = bind(SqlValue::Text(id.clone())); + predicates.push(format!("(COALESCE(t.created_at,'')<{key} OR (COALESCE(t.created_at,'')={key} AND t.ticket_id>{id}))")); + } + TicketQuerySort::Title => { + let key = bind(SqlValue::Text(key.clone())); + let id = bind(SqlValue::Text(id.clone())); + predicates.push(format!( + "(lower(t.title)>{key} OR (lower(t.title)={key} AND t.ticket_id>{id}))" + )); + } + TicketQuerySort::Priority | TicketQuerySort::Relevance => { + let (rank, updated_at) = key.split_once('|').unwrap_or(("9", "")); + let rank = bind(SqlValue::Integer(rank.parse::().unwrap_or(9))); + let updated_at = bind(SqlValue::Text(updated_at.to_string())); + let id = bind(SqlValue::Text(id.clone())); + predicates.push(format!("({rank_expression}>{rank} OR ({rank_expression}={rank} AND (COALESCE(t.updated_at,'')<{updated_at} OR (COALESCE(t.updated_at,'')={updated_at} AND t.ticket_id>{id}))))")); + } + } + } + let order = match sort { + TicketQuerySort::Title => "t.title COLLATE NOCASE ASC, t.ticket_id ASC".to_string(), + TicketQuerySort::CreatedDesc => "t.created_at DESC, t.ticket_id ASC".to_string(), + TicketQuerySort::UpdatedDesc => "t.updated_at DESC, t.ticket_id ASC".to_string(), + TicketQuerySort::Priority | TicketQuerySort::Relevance => { + format!("{rank_expression} ASC, t.updated_at DESC, t.ticket_id ASC") + } + }; + let limit = bind(SqlValue::Integer(i64::try_from(limit).unwrap_or(i64::MAX))); + let sql = format!( + "SELECT t.ticket_id FROM typed_tickets t WHERE {} ORDER BY {order} LIMIT {limit}", + predicates.join(" AND ") + ); + self.store.with_conn(|connection| { + let mut statement = connection.prepare(&sql)?; + let rows = statement.query_map(params_from_iter(values.iter()), |row| row.get(0))?; + Ok(rows.collect::, _>>()?) + }) + } + + fn query_objective_candidate_ids( + &self, + query: &ObjectiveQueryRequest, + sort: ObjectiveQuerySort, + after: Option<&(String, String)>, + limit: usize, + ) -> Result> { + let mut values = vec![SqlValue::Text(self.workspace_id.clone())]; + let mut predicates = vec!["o.workspace_id=?1".to_string()]; + let mut bind = |value: SqlValue| { + values.push(value); + format!("?{}", values.len()) + }; + if !query.states.is_empty() { + let states = query + .states + .iter() + .map(|state| bind(SqlValue::Text(state.clone()))) + .collect::>(); + predicates.push(format!("o.state IN ({})", states.join(","))); + } + if let Some(text) = query.query.as_deref().filter(|text| !text.is_empty()) { + let pattern = bind(SqlValue::Text(format!("%{}%", text.to_lowercase()))); + predicates.push(format!( + "(lower(o.title) LIKE {pattern} OR lower(o.body_md) LIKE {pattern})" + )); + } + if let Some(value) = &query.updated_after { + let value = bind(SqlValue::Text(value.clone())); + predicates.push(format!("o.updated_at>{value}")); + } + if let Some(value) = &query.updated_before { + let value = bind(SqlValue::Text(value.clone())); + predicates.push(format!("o.updated_at<{value}")); + } + if let Some(value) = &query.linked_ticket_id { + let value = bind(SqlValue::Text(value.clone())); + predicates.push(format!("EXISTS (SELECT 1 FROM objective_ticket_links link WHERE link.workspace_id=o.workspace_id AND link.objective_id=o.objective_id AND link.ticket_id={value})")); + } + let relevance_rank = if let Some(text) = + query.query.as_deref().filter(|text| !text.is_empty()) + { + let pattern = bind(SqlValue::Text(format!("%{}%", text.to_lowercase()))); + format!( + "CASE WHEN lower(o.title) LIKE {pattern} THEN 0 WHEN lower(o.body_md) LIKE {pattern} THEN 1 ELSE 2 END" + ) + } else { + "2".to_string() + }; + if let Some((key, id)) = after { + match sort { + ObjectiveQuerySort::UpdatedDesc => { + let key = bind(SqlValue::Text(key.clone())); + let id = bind(SqlValue::Text(id.clone())); + predicates.push(format!( + "(o.updated_at<{key} OR (o.updated_at={key} AND o.objective_id>{id}))" + )); + } + ObjectiveQuerySort::CreatedDesc => { + let key = bind(SqlValue::Text(key.clone())); + let id = bind(SqlValue::Text(id.clone())); + predicates.push(format!( + "(o.created_at<{key} OR (o.created_at={key} AND o.objective_id>{id}))" + )); + } + ObjectiveQuerySort::Title => { + let key = bind(SqlValue::Text(key.clone())); + let id = bind(SqlValue::Text(id.clone())); + predicates.push(format!( + "(lower(o.title)>{key} OR (lower(o.title)={key} AND o.objective_id>{id}))" + )); + } + ObjectiveQuerySort::Relevance => { + let (rank, updated_at) = key.split_once('|').unwrap_or(("9", "")); + let rank = bind(SqlValue::Integer(rank.parse::().unwrap_or(9))); + let updated_at = bind(SqlValue::Text(updated_at.to_string())); + let id = bind(SqlValue::Text(id.clone())); + predicates.push(format!("({relevance_rank}>{rank} OR ({relevance_rank}={rank} AND (o.updated_at<{updated_at} OR (o.updated_at={updated_at} AND o.objective_id>{id}))))")); + } + } + } + let order = match sort { + ObjectiveQuerySort::Title => { + "o.title COLLATE NOCASE ASC, o.objective_id ASC".to_string() + } + ObjectiveQuerySort::CreatedDesc => "o.created_at DESC, o.objective_id ASC".to_string(), + ObjectiveQuerySort::UpdatedDesc => "o.updated_at DESC, o.objective_id ASC".to_string(), + ObjectiveQuerySort::Relevance => { + format!("{relevance_rank} ASC, o.updated_at DESC, o.objective_id ASC") + } + }; + let limit = bind(SqlValue::Integer(i64::try_from(limit).unwrap_or(i64::MAX))); + let sql = format!( + "SELECT o.objective_id FROM objectives o WHERE {} ORDER BY {order} LIMIT {limit}", + predicates.join(" AND ") + ); + self.store.with_conn(|connection| { + let mut statement = connection.prepare(&sql)?; + let rows = statement.query_map(params_from_iter(values.iter()), |row| row.get(0))?; + Ok(rows.collect::, _>>()?) + }) + } + fn read_ticket_detail(&self, id: &str, request: TicketShowRequest) -> Result { validate_project_id(id)?; let ticket = self .ticket_backend .show(TicketIdOrSlug::Id(id.to_string()))?; + self.ticket_detail_from_ticket(ticket, request) + } + + fn ticket_detail_from_ticket( + &self, + ticket: ticket::Ticket, + request: TicketShowRequest, + ) -> Result { + let id = ticket.meta.id.as_str(); let (body, body_truncated) = truncate_body(ticket.document.body.as_str(), DETAIL_BODY_LIMIT); let event_limit = request @@ -458,42 +809,100 @@ impl TicketAuthority for SqliteWorkspaceAuthority { }) } + fn list_ticket_page(&self, request: TicketListPageRequest) -> Result { + let limit = request.limit.unwrap_or(30).clamp(1, 100); + let mut states = request.states; + states.sort(); + states.dedup(); + let parsed_states = states + .iter() + .map(|state| { + TicketWorkflowState::parse(state).ok_or_else(|| { + Error::InvalidRecordId(format!("unsupported ticket state `{state}`")) + }) + }) + .collect::>>()?; + let fingerprint = format!( + "ticket-summary:v2:sort=priority:states={}", + states.join(",") + ); + let after = request + .cursor + .as_deref() + .map(|cursor| parse_ticket_summary_cursor(cursor, &fingerprint)) + .transpose()?; + let page = + self.ticket_backend + .list_workspace_projection_page(SqliteTicketListPageQuery { + states: parsed_states, + limit, + after, + })?; + let items = page + .items + .into_iter() + .map(ticket_summary_from_sqlite_item) + .collect::>(); + let next_cursor = page + .next + .map(|position| make_ticket_summary_cursor(&fingerprint, position)); + Ok(TicketSummaryPage { + page: QueryPage { + limit, + returned: items.len(), + has_more: page.has_more, + next_cursor, + sort: "priority".to_string(), + source_limit: None, + source_truncated: false, + }, + items, + invalid_records: Vec::new(), + record_authority: RECORD_SOURCE_WORKSPACE_SQLITE.to_string(), + }) + } + fn query_tickets(&self, query: TicketQueryRequest) -> Result { validate_ticket_query(&query)?; let limit = query.limit.unwrap_or(50).clamp(1, 100); let sort = normalize_ticket_sort(query.sort.as_deref(), query.query.is_some())?; + let fingerprint = ticket_query_fingerprint(&query, sort); let cursor = query .cursor .as_deref() - .map(parse_query_cursor) + .map(|cursor| parse_bound_query_cursor(cursor, &fingerprint)) .transpose()?; - let mut summaries = self.list_tickets(1_001)?.items; - let source_truncated = summaries.len() > 1_000; - summaries.truncate(1_000); + let candidate_limit = limit.saturating_add(1); + let candidate_ids = + self.query_ticket_candidate_ids(&query, sort, cursor.as_ref(), candidate_limit)?; + let source_truncated = candidate_ids.len() == candidate_limit; let mut items = Vec::new(); - for summary in summaries { - let detail = self.read_ticket_detail( - &summary.id, + for ticket_id in candidate_ids { + let authoritative = self + .ticket_backend + .show(TicketIdOrSlug::Id(ticket_id.clone()))?; + let summary = ticket_summary_from_ticket(&authoritative); + let authoritative_body = authoritative.document.body.clone(); + let authoritative_events = authoritative.events.clone(); + let detail = self.ticket_detail_from_ticket( + authoritative, TicketShowRequest { event_limit: Some(TICKET_EVENT_LIMIT), event_cursor: None, }, )?; - let authoritative = self - .ticket_backend - .show(TicketIdOrSlug::Id(summary.id.clone()))?; if ticket_matches_query( &summary, &detail, - authoritative.document.body.as_str(), - &authoritative.events, + authoritative_body.as_str(), + &authoritative_events, &query, ) { items.push(ticket_query_item( summary, &detail, - authoritative.document.body.as_str(), - &authoritative.events, + authoritative_body.as_str(), + &authoritative_events, &query, )); } @@ -505,7 +914,11 @@ impl TicketAuthority for SqliteWorkspaceAuthority { let has_more = items.len() > limit; items.truncate(limit); let next_cursor = has_more - .then(|| items.last().map(|item| make_ticket_cursor(item, sort))) + .then(|| { + items + .last() + .map(|item| make_ticket_cursor(item, sort, &fingerprint)) + }) .flatten(); Ok(TicketQueryResponse { page: QueryPage { @@ -514,7 +927,7 @@ impl TicketAuthority for SqliteWorkspaceAuthority { has_more, next_cursor, sort: sort.to_string(), - source_limit: Some(1_000), + source_limit: Some(candidate_limit), source_truncated, }, items, @@ -566,33 +979,36 @@ impl ObjectiveAuthority for SqliteWorkspaceAuthority { )?; let limit = query.limit.unwrap_or(50).clamp(1, 100); let sort = normalize_objective_sort(query.sort.as_deref(), query.query.is_some())?; + let fingerprint = objective_query_fingerprint(&query, sort); let cursor = query .cursor .as_deref() - .map(parse_query_cursor) + .map(|cursor| parse_bound_query_cursor(cursor, &fingerprint)) .transpose()?; - let mut objectives = self.list_objectives(1_001)?.items; - let source_truncated = objectives.len() > 1_000; - objectives.truncate(1_000); + let candidate_limit = limit.saturating_add(1); + let objective_ids = + self.query_objective_candidate_ids(&query, sort, cursor.as_ref(), candidate_limit)?; + let source_truncated = objective_ids.len() == candidate_limit; let mut items = Vec::new(); - for objective in objectives { - let body_md = self.objective_record(&objective.id)?.body_md; - if !objective_matches_query(&objective, &body_md, &query) { - continue; - } + for objective_id in objective_ids { + let record = self.objective_record(&objective_id)?; let linked_tickets = self .store - .list_objective_ticket_links(&self.workspace_id, &objective.id)? + .list_objective_ticket_links(&self.workspace_id, &objective_id)? .into_iter() .map(|link| link.ticket_id) .collect::>(); - if query - .linked_ticket_id - .as_ref() - .is_some_and(|id| !linked_tickets.iter().any(|ticket_id| ticket_id == id)) - { - continue; - } + let body_md = record.body_md.clone(); + let objective = ObjectiveSummary { + id: record.objective_id, + title: record.title, + state: record.state, + created_at: Some(record.created_at), + updated_at: Some(record.updated_at), + summary: summarize_body(&body_md), + linked_tickets: linked_tickets.clone(), + record_source: RECORD_SOURCE_WORKSPACE_SQLITE.to_string(), + }; items.push(objective_query_item( objective, linked_tickets, @@ -607,7 +1023,11 @@ impl ObjectiveAuthority for SqliteWorkspaceAuthority { let has_more = items.len() > limit; items.truncate(limit); let next_cursor = has_more - .then(|| items.last().map(|item| make_objective_cursor(item, sort))) + .then(|| { + items + .last() + .map(|item| make_objective_cursor(item, sort, &fingerprint)) + }) .flatten(); Ok(ObjectiveQueryResponse { page: QueryPage { @@ -616,7 +1036,7 @@ impl ObjectiveAuthority for SqliteWorkspaceAuthority { has_more, next_cursor, sort: sort.to_string(), - source_limit: Some(1_000), + source_limit: Some(candidate_limit), source_truncated, }, items, @@ -1240,6 +1660,35 @@ fn parse_offset_cursor(cursor: &str, field: &str) -> Result { .map_err(|_| Error::InvalidRecordId(format!("invalid {field}"))) } +fn make_bound_query_cursor(fingerprint: &str, key: &str, id: &str) -> String { + make_query_cursor(&format!("{fingerprint}\n{key}"), id) +} + +fn parse_bound_query_cursor(value: &str, fingerprint: &str) -> Result<(String, String)> { + let (key, id) = parse_query_cursor(value)?; + let prefix = format!("{fingerprint}\n"); + let key = key.strip_prefix(&prefix).ok_or_else(|| { + Error::InvalidRecordId("cursor does not match the current filters or sort".to_string()) + })?; + Ok((key.to_string(), id)) +} + +fn ticket_query_fingerprint(query: &TicketQueryRequest, sort: TicketQuerySort) -> String { + let mut query = query.clone(); + query.cursor = None; + query.limit = None; + query.sort = Some(sort.to_string()); + serde_json::to_string(&query).expect("ticket query fingerprint must serialize") +} + +fn objective_query_fingerprint(query: &ObjectiveQueryRequest, sort: ObjectiveQuerySort) -> String { + let mut query = query.clone(); + query.cursor = None; + query.limit = None; + query.sort = Some(sort.to_string()); + serde_json::to_string(&query).expect("objective query fingerprint must serialize") +} + fn make_query_cursor(key: &str, id: &str) -> String { format!("v1:{}:{key}{id}", key.len()) } @@ -1545,8 +1994,8 @@ fn sort_ticket_query_items(items: &mut [TicketQueryItem], sort: TicketQuerySort) }); } -fn make_ticket_cursor(item: &TicketQueryItem, sort: TicketQuerySort) -> String { - make_query_cursor(&ticket_sort_key(item, sort), &item.id) +fn make_ticket_cursor(item: &TicketQueryItem, sort: TicketQuerySort, fingerprint: &str) -> String { + make_bound_query_cursor(fingerprint, &ticket_sort_key(item, sort), &item.id) } fn ticket_item_after_cursor( @@ -1571,31 +2020,6 @@ fn ticket_item_after_cursor( } } -fn objective_matches_query( - objective: &ObjectiveSummary, - body_md: &str, - query: &ObjectiveQueryRequest, -) -> bool { - if !query.states.is_empty() && !query.states.iter().any(|state| state == &objective.state) { - return false; - } - if query - .updated_after - .as_ref() - .is_some_and(|after| objective.updated_at.as_deref().unwrap_or("") <= after.as_str()) - || query - .updated_before - .as_ref() - .is_some_and(|before| objective.updated_at.as_deref().unwrap_or("") >= before.as_str()) - { - return false; - } - query.query.as_ref().is_none_or(|text| { - let needle = text.to_lowercase(); - objective.title.to_lowercase().contains(&needle) || body_md.to_lowercase().contains(&needle) - }) -} - fn objective_query_item( objective: ObjectiveSummary, linked_tickets: Vec, @@ -1673,8 +2097,12 @@ fn sort_objective_query_items(items: &mut [ObjectiveQueryItem], sort: ObjectiveQ }); } -fn make_objective_cursor(item: &ObjectiveQueryItem, sort: ObjectiveQuerySort) -> String { - make_query_cursor(&objective_sort_key(item, sort), &item.id) +fn make_objective_cursor( + item: &ObjectiveQueryItem, + sort: ObjectiveQuerySort, + fingerprint: &str, +) -> String { + make_bound_query_cursor(fingerprint, &objective_sort_key(item, sort), &item.id) } fn objective_item_after_cursor( @@ -1802,6 +2230,80 @@ fn memory_resolution_from_record(record: MemoryStagingResolutionRecord) -> Memor } } +fn ticket_summary_from_ticket(ticket: &ticket::Ticket) -> TicketSummary { + let summary = ticket::TicketSummary { + id: ticket.meta.id.clone(), + slug: ticket.meta.slug.clone(), + title: ticket.meta.title.clone(), + status: ticket.meta.status.clone(), + kind: ticket.meta.kind.clone(), + priority: ticket.meta.priority.clone(), + labels: ticket.meta.labels.clone(), + readiness: ticket.meta.readiness.clone(), + workflow_state: ticket.meta.workflow_state, + workflow_state_explicit: ticket.meta.workflow_state_explicit, + queued_by: ticket.meta.queued_by.clone(), + queued_at: ticket.meta.queued_at.clone(), + updated_at: ticket.meta.updated_at.clone(), + }; + ticket_summary_from_sqlite_item(SqliteTicketListItem { + summary, + relation_blockers: ticket.relations.blockers.clone(), + }) +} + +fn ticket_summary_from_sqlite_item(item: SqliteTicketListItem) -> TicketSummary { + let projection = project_ticket_workspace_item(&item.summary, &item.relation_blockers, None); + TicketSummary { + id: item.summary.id, + title: item.summary.title, + state: item.summary.workflow_state.as_str().to_string(), + priority: item.summary.priority, + updated_at: item.summary.updated_at, + queued_by: item.summary.queued_by, + queued_at: item.summary.queued_at, + workspace_action_priority: workspace_action_priority_name(projection.priority).to_string(), + record_source: "sqlite_yoi_ticket".to_string(), + } +} + +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] +struct TicketSummaryCursorEnvelope { + version: u8, + fingerprint: String, + position: SqliteTicketListCursor, +} + +fn make_ticket_summary_cursor(fingerprint: &str, position: SqliteTicketListCursor) -> String { + make_query_cursor( + &serde_json::to_string(&TicketSummaryCursorEnvelope { + version: 1, + fingerprint: fingerprint.to_string(), + position, + }) + .expect("ticket summary cursor must serialize"), + "", + ) +} + +fn parse_ticket_summary_cursor( + value: &str, + expected_fingerprint: &str, +) -> Result { + let (encoded, trailing) = parse_query_cursor(value)?; + if !trailing.is_empty() { + return Err(Error::InvalidRecordId("cursor is malformed".to_string())); + } + let cursor: TicketSummaryCursorEnvelope = serde_json::from_str(&encoded) + .map_err(|_| Error::InvalidRecordId("cursor is malformed".to_string()))?; + if cursor.version != 1 || cursor.fingerprint != expected_fingerprint { + return Err(Error::InvalidRecordId( + "cursor does not match the current ticket filters or sort".to_string(), + )); + } + Ok(cursor.position) +} + fn workspace_action_priority_name(priority: TicketWorkspaceActionPriority) -> &'static str { match priority { TicketWorkspaceActionPriority::ReadyForQueue => "ready_for_queue", @@ -2083,6 +2585,49 @@ mod tests { .contains(&"body".to_string()) ); assert_eq!(ticket_query.page.limit, 1); + let no_review = authority + .query_tickets(TicketQueryRequest { + review_status: Some("none".to_string()), + sort: Some("updated_desc".to_string()), + limit: Some(10), + ..TicketQueryRequest::default() + }) + .unwrap(); + assert!( + no_review + .items + .iter() + .any(|item| item.id == "00000000001J2"), + "review-status storage predicate must query the authoritative MR thread schema" + ); + for evidence in [ + "implementation_report", + "implementation_report_after_rescope", + ] { + authority + .query_tickets(TicketQueryRequest { + evidence: vec![evidence.to_string()], + limit: Some(10), + ..TicketQueryRequest::default() + }) + .unwrap_or_else(|error| panic!("accepted evidence filter {evidence}: {error}")); + } + for attention in ["implementation_report_not_closed", "report_after_rescope"] { + authority + .query_tickets(TicketQueryRequest { + attention: vec![attention.to_string()], + limit: Some(10), + ..TicketQueryRequest::default() + }) + .unwrap_or_else(|error| panic!("accepted attention filter {attention}: {error}")); + } + authority + .query_tickets(TicketQueryRequest { + review_status: Some("changes_requested".to_string()), + limit: Some(10), + ..TicketQueryRequest::default() + }) + .expect("accepted review-status alias must execute"); let historical_event_query = authority .query_tickets(TicketQueryRequest { query: Some("Historical event marker".to_string()), @@ -2135,6 +2680,35 @@ mod tests { .unwrap(); assert_eq!(exact_relation.items.len(), 1); assert_eq!(exact_relation.items[0].id, "00000000001J2"); + let incoming_relation = authority + .query_tickets(TicketQueryRequest { + related_ticket_id: Some("00000000001J2".to_string()), + relation_kind: Some("related".to_string()), + limit: Some(10), + ..TicketQueryRequest::default() + }) + .unwrap(); + assert_eq!(incoming_relation.items.len(), 1); + assert_eq!(incoming_relation.items[0].id, "00000000001J5"); + let summary_page = authority + .list_ticket_page(TicketListPageRequest { + states: vec!["planning".to_string(), "ready".to_string()], + limit: Some(1), + cursor: None, + }) + .unwrap(); + assert_eq!(summary_page.items.len(), 1); + assert!(summary_page.page.has_more); + let mismatched_summary_cursor = authority.list_ticket_page(TicketListPageRequest { + states: vec!["done".to_string()], + limit: Some(1), + cursor: summary_page.page.next_cursor, + }); + assert!(matches!( + mismatched_summary_cursor, + Err(Error::InvalidRecordId(_)) + )); + let first_page = authority .query_tickets(TicketQueryRequest { sort: Some("title".to_string()), @@ -2154,6 +2728,13 @@ mod tests { assert_eq!(second_page.items.len(), 1); assert_ne!(first_page.items[0].id, second_page.items[0].id); assert!(second_page.page.has_more); + let mismatched_cursor = authority.query_tickets(TicketQueryRequest { + sort: Some("updated_desc".to_string()), + limit: Some(1), + cursor: first_page.page.next_cursor.clone(), + ..TicketQueryRequest::default() + }); + assert!(matches!(mismatched_cursor, Err(Error::InvalidRecordId(_)))); let third_page = authority .query_tickets(TicketQueryRequest { sort: Some("title".to_string()), diff --git a/crates/workspace-server/src/records.rs b/crates/workspace-server/src/records.rs index 35348f60..62f978b7 100644 --- a/crates/workspace-server/src/records.rs +++ b/crates/workspace-server/src/records.rs @@ -12,6 +12,14 @@ pub struct ProjectRecordList { pub record_authority: String, } +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +pub struct TicketSummaryPage { + pub items: Vec, + pub page: QueryPage, + pub invalid_records: Vec, + pub record_authority: String, +} + #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] #[cfg_attr(feature = "typescript", derive(ts_rs::TS))] pub struct InvalidProjectRecord { @@ -33,12 +41,23 @@ pub struct TicketSummary { pub record_source: String, } +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Default)] +pub struct TicketListPageRequest { + #[serde(default)] + pub states: Vec, + #[serde(default)] + pub limit: Option, + #[serde(default)] + pub cursor: Option, +} + #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] #[cfg_attr(feature = "typescript", derive(ts_rs::TS))] pub struct TicketListResponse { pub workspace_id: String, pub limit: usize, pub items: Vec, + pub page: QueryPage, pub invalid_records: Vec, pub record_authority: String, } diff --git a/crates/workspace-server/src/server.rs b/crates/workspace-server/src/server.rs index 5f69c5e8..bac31c20 100644 --- a/crates/workspace-server/src/server.rs +++ b/crates/workspace-server/src/server.rs @@ -2268,6 +2268,9 @@ struct ObjectiveEditRequest { #[derive(Debug, Deserialize)] struct TicketListQuery { limit: Option, + cursor: Option, + /// Comma-separated workflow states. Repeated lane requests normally pass one state group. + states: Option, } #[derive(Debug, Deserialize)] @@ -8476,17 +8479,35 @@ async fn list_tickets( State(api): State, Query(query): Query, ) -> ApiResult> { - let requested_limit = query.limit.unwrap_or(api.config.max_records); - let limit = requested_limit.min(1000); - let ProjectRecordList { + let limit = query.limit.unwrap_or(30).clamp(1, 100); + let states = query + .states + .as_deref() + .map(|states| { + states + .split(',') + .filter(|state| !state.is_empty()) + .map(str::to_string) + .collect::>() + }) + .unwrap_or_default(); + let crate::records::TicketSummaryPage { items, + page, invalid_records, record_authority, - } = api.authority.list_tickets(limit)?; + } = api + .authority + .list_ticket_page(crate::records::TicketListPageRequest { + states, + limit: Some(limit), + cursor: query.cursor, + })?; Ok(Json(crate::records::TicketListResponse { workspace_id: api.config.workspace_id, limit, items, + page, invalid_records, record_authority, })) diff --git a/crates/workspace-server/src/store.rs b/crates/workspace-server/src/store.rs index 06e18e5e..f1d59410 100644 --- a/crates/workspace-server/src/store.rs +++ b/crates/workspace-server/src/store.rs @@ -196,6 +196,11 @@ const MIGRATIONS: &[Migration] = &[ name: "remove Worker control delegation authority", apply: remove_worker_control_delegation_authority, }, + Migration { + version: 36, + name: "add Objective query indexes", + apply: add_objective_query_indexes, + }, ]; struct Migration { @@ -5046,6 +5051,24 @@ fn create_worker_control_delegation_operation_authority(conn: &Connection) -> Re Ok(()) } +fn add_objective_query_indexes(conn: &Connection) -> Result<()> { + conn.execute_batch( + r#" + CREATE INDEX IF NOT EXISTS objectives_workspace_state_updated + ON objectives(workspace_id, state, updated_at DESC, objective_id); + CREATE INDEX IF NOT EXISTS objectives_workspace_updated + ON objectives(workspace_id, updated_at DESC, objective_id); + CREATE INDEX IF NOT EXISTS objectives_workspace_created + ON objectives(workspace_id, created_at DESC, objective_id); + CREATE INDEX IF NOT EXISTS objectives_workspace_title + ON objectives(workspace_id, title COLLATE NOCASE, objective_id); + CREATE INDEX IF NOT EXISTS objective_ticket_links_workspace_ticket_objective + ON objective_ticket_links(workspace_id, ticket_id, objective_id); + "#, + )?; + Ok(()) +} + fn remove_worker_control_delegation_authority(conn: &Connection) -> Result<()> { let mut statement = conn.prepare("SELECT workspace_id, grant_id, permissions_json FROM worker_control_grants")?; @@ -5799,7 +5822,7 @@ INSERT INTO worker_control_grants ( apply_migrations(&conn).unwrap(); - assert_eq!(current_schema_version(&conn).unwrap(), 35); + assert_eq!(current_schema_version(&conn).unwrap(), 36); assert!(!table_exists(&conn, "worker_control_delegation_operations").unwrap()); let (permissions_json, revoked_at): (String, Option) = conn .query_row( @@ -5847,7 +5870,7 @@ INSERT INTO worker_control_grants ( apply_migrations(&conn).unwrap(); - assert_eq!(current_schema_version(&conn).unwrap(), 35); + assert_eq!(current_schema_version(&conn).unwrap(), 36); assert!(table_exists(&conn, "worker_workdir_attachment_reservations").unwrap()); } @@ -5880,7 +5903,7 @@ CREATE TABLE flow_events (event_id TEXT PRIMARY KEY); apply_migrations(&conn).unwrap(); - assert_eq!(current_schema_version(&conn).unwrap(), 35); + assert_eq!(current_schema_version(&conn).unwrap(), 36); assert!(table_exists(&conn, "flow_sources").unwrap()); assert!(table_exists(&conn, "flow_source_revisions").unwrap()); assert!(!table_exists(&conn, "flow_instances").unwrap()); @@ -5947,7 +5970,7 @@ INSERT INTO worker_workdir_attachment_reservations ( apply_migrations(&conn).unwrap(); - assert_eq!(current_schema_version(&conn).unwrap(), 35); + assert_eq!(current_schema_version(&conn).unwrap(), 36); let repositories_sql: String = conn .query_row( "SELECT sql FROM sqlite_master WHERE type = 'table' AND name = 'repositories'", @@ -6127,7 +6150,7 @@ INSERT INTO workdir_registry ( let db = dir.path().join("control-plane.sqlite"); let store = SqliteWorkspaceStore::open(&db).unwrap(); - assert_eq!(store.schema_version().await.unwrap(), 35); + assert_eq!(store.schema_version().await.unwrap(), 36); assert!( !store .with_conn(|conn| table_exists(conn, "worker_workspace_credentials")) @@ -6144,7 +6167,7 @@ INSERT INTO workdir_registry ( store.upsert_workspace(&record).await.unwrap(); let reopened = SqliteWorkspaceStore::open(&db).unwrap(); - assert_eq!(reopened.schema_version().await.unwrap(), 35); + assert_eq!(reopened.schema_version().await.unwrap(), 36); assert_eq!( reopened.get_workspace("local-dev").await.unwrap(), Some(record) @@ -6691,7 +6714,7 @@ INSERT INTO workdir_registry ( .unwrap(); let store = SqliteWorkspaceStore::from_connection(conn).unwrap(); - assert_eq!(store.schema_version().await.unwrap(), 35); + assert_eq!(store.schema_version().await.unwrap(), 36); store .with_conn(|conn| { @@ -6880,7 +6903,7 @@ CREATE TABLE ticket_assignment_operations ( #[tokio::test] async fn repository_records_round_trip() { let store = SqliteWorkspaceStore::in_memory().unwrap(); - assert_eq!(store.schema_version().await.unwrap(), 35); + assert_eq!(store.schema_version().await.unwrap(), 36); let workspace = WorkspaceRecord { workspace_id: "local-dev".to_string(), owner_account_id: None, @@ -6946,7 +6969,7 @@ CREATE TABLE ticket_assignment_operations ( #[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(), 35); + assert_eq!(store.schema_version().await.unwrap(), 36); let workspace = WorkspaceRecord { workspace_id: "local-dev".to_string(), owner_account_id: None, @@ -7337,7 +7360,7 @@ CREATE TABLE ticket_assignment_operations ( #[tokio::test] async fn account_and_login_records_round_trip() { let store = SqliteWorkspaceStore::in_memory().unwrap(); - assert_eq!(store.schema_version().await.unwrap(), 35); + assert_eq!(store.schema_version().await.unwrap(), 36); let now = "2026-07-22T00:00:00Z".to_string(); let account = AccountRecord { account_id: "acct-user-alice".to_string(), diff --git a/web/workspace/src/lib/generated/ticket-api.ts b/web/workspace/src/lib/generated/ticket-api.ts index a24976be..a6815a42 100644 --- a/web/workspace/src/lib/generated/ticket-api.ts +++ b/web/workspace/src/lib/generated/ticket-api.ts @@ -19,6 +19,7 @@ export type TicketListResponse = { workspace_id: string; limit: number; items: Array; + page: QueryPage; invalid_records: Array; record_authority: string; }; diff --git a/web/workspace/src/lib/workspace/styles/tickets.css b/web/workspace/src/lib/workspace/styles/tickets.css index c7098b5f..0fc65d2e 100644 --- a/web/workspace/src/lib/workspace/styles/tickets.css +++ b/web/workspace/src/lib/workspace/styles/tickets.css @@ -169,8 +169,32 @@ display: grid; align-content: start; gap: 0.55rem; + max-height: min(68vh, 48rem); + overflow-y: auto; + overscroll-behavior: contain; padding: 0.6rem; } + .ticket-lane-page-state { + display: flex; + justify-content: center; + gap: 0.5rem; + margin: 0; + padding: 0.45rem; + color: var(--text-muted); + font-size: 0.72rem; + text-align: center; + } + .ticket-lane-page-error { + align-items: center; + color: var(--danger); + } + .ticket-lane-page-error button { + border: 1px solid var(--line); + border-radius: 0.35rem; + background: var(--bg-raised); + color: inherit; + padding: 0.2rem 0.45rem; + } .ticket-card { display: grid; gap: 0.55rem; diff --git a/web/workspace/src/lib/workspace/tickets/ticket-panel.test.ts b/web/workspace/src/lib/workspace/tickets/ticket-panel.test.ts index 842103cf..649eb81d 100644 --- a/web/workspace/src/lib/workspace/tickets/ticket-panel.test.ts +++ b/web/workspace/src/lib/workspace/tickets/ticket-panel.test.ts @@ -105,6 +105,32 @@ Deno.test("ticket worker launch uses the common Worker route and bounded Ticket ); }); +Deno.test("ticket board paginates every lane independently with bounded summary requests", async () => { + const loadSource = await Deno.readTextFile( + "src/routes/w/[workspaceId]/tickets/+page.ts", + ); + const pageSource = await Deno.readTextFile( + "src/routes/w/[workspaceId]/tickets/+page.svelte", + ); + + assertIncludes(loadSource, 'limit: "30"'); + assertIncludes(loadSource, "states: states.join"); + assertIncludes(pageSource, "lane.page.next_cursor"); + assertIncludes(pageSource, "lane.loading"); + assertIncludes(pageSource, "mergeTickets"); + assertIncludes(pageSource, "onscroll"); + assertIncludes(pageSource, "Retry"); + if (loadSource.includes("limit=1000") || pageSource.includes("limit=1000")) { + throw new Error("Ticket board must not fetch the legacy 1000-item list"); + } + if ( + loadSource.includes("/tickets/query") || + pageSource.includes("/tickets/query") + ) { + throw new Error("Ticket board must use the bounded summary endpoint"); + } +}); + Deno.test("ticket panel starts the Orchestrator explicitly and gates orchestration actions", async () => { const panelSource = await Deno.readTextFile( "src/routes/w/[workspaceId]/tickets/+page.svelte", @@ -113,7 +139,10 @@ Deno.test("ticket panel starts the Orchestrator explicitly and gates orchestrati "src/routes/w/[workspaceId]/tickets/[ticketId]/+page.svelte", ); - assertIncludes(panelSource, 'workspaceApiPath(data.workspaceId, "/orchestrator")'); + assertIncludes( + panelSource, + 'workspaceApiPath(data.workspaceId, "/orchestrator")', + ); assertIncludes(panelSource, '{ method: "POST" }'); assertIncludes(panelSource, "Start Orchestrator"); assertIncludes(panelSource, "orchestrator.data?.online"); diff --git a/web/workspace/src/routes/w/[workspaceId]/tickets/+page.svelte b/web/workspace/src/routes/w/[workspaceId]/tickets/+page.svelte index f4c55b04..ff68779f 100644 --- a/web/workspace/src/routes/w/[workspaceId]/tickets/+page.svelte +++ b/web/workspace/src/routes/w/[workspaceId]/tickets/+page.svelte @@ -2,30 +2,94 @@ import { untrack } from "svelte"; import type { ApiResult } from "$lib/workspace/api/http"; import { loadJson, workspaceApiPath } from "$lib/workspace/api/http"; + import type { + QueryPage, + TicketListResponse, + TicketSummary, + } from "$lib/generated/ticket-api"; import { ticketLanes, type WorkspaceOrchestratorStatus, } from "$lib/workspace/tickets/ticket-panel"; - import type { - TicketListResponse, - TicketSummary, - } from "$lib/workspace/sidebar/types"; + import type { PageData } from "./$types"; - const { data } = $props<{ - data: { - workspaceId: string; - tickets: ApiResult; - orchestrator: ApiResult; - }; - }>(); + type LaneState = { + states: string[]; + tickets: TicketSummary[]; + page: QueryPage; + loading: boolean; + error: string | null; + }; - const initialTickets = untrack(() => data.tickets.data?.items ?? []); - let tickets = $state(initialTickets); + let { data }: { data: PageData } = $props(); + // svelte-ignore state_referenced_locally + let laneState = $state>( + Object.fromEntries( + Object.entries(data.ticketLanes).map(([laneId, lane]) => [ + laneId, + { + states: [...lane.states], + tickets: lane.response.items, + page: lane.response.page, + loading: false, + error: null, + }, + ]), + ), + ); let orchestrator = $state>( untrack(() => data.orchestrator), ); let orchestratorStarting = $state(false); - let lanes = $derived(ticketLanes(tickets)); + const tickets = $derived( + Object.values(laneState).flatMap((lane) => lane.tickets), + ); + const lanes = $derived(ticketLanes(tickets)); + + function mergeTickets( + current: TicketSummary[], + incoming: TicketSummary[], + ): TicketSummary[] { + const byId = new Map(current.map((ticket) => [ticket.id, ticket])); + for (const ticket of incoming) byId.set(ticket.id, ticket); + return [...byId.values()]; + } + + async function loadMore(laneId: string): Promise { + const lane = laneState[laneId]; + if (!lane || lane.loading || !lane.page.has_more || !lane.page.next_cursor) { + return; + } + lane.loading = true; + lane.error = null; + try { + const search = new URLSearchParams({ + limit: "30", + states: lane.states.join(","), + cursor: lane.page.next_cursor, + }); + const response = await fetch( + `/api/w/${encodeURIComponent(data.workspaceId)}/tickets?${search}`, + ); + if (!response.ok) { + throw new Error(`追加読み込みに失敗しました (${response.status})`); + } + const page = (await response.json()) as TicketListResponse; + lane.tickets = mergeTickets(lane.tickets, page.items); + lane.page = page.page; + } catch (error) { + lane.error = error instanceof Error ? error.message : String(error); + } finally { + lane.loading = false; + } + } + + function handleLaneScroll(event: Event, laneId: string): void { + const container = event.currentTarget as HTMLElement; + const remaining = + container.scrollHeight - container.scrollTop - container.clientHeight; + if (remaining <= 96) void loadMore(laneId); + } async function startOrchestrator() { if (orchestratorStarting || orchestrator.data?.online) return; @@ -45,7 +109,9 @@ } -Tickets · Yoi + + Tickets · {data.workspaceId} +
@@ -76,7 +142,7 @@
{tickets.length} - tickets + loaded tickets
@@ -93,6 +159,7 @@
{#each lanes as lane (lane.id)} + {@const pagination = laneState[lane.id]}
@@ -101,8 +168,11 @@
{lane.tickets.length}
- -
{/each} diff --git a/web/workspace/src/routes/w/[workspaceId]/tickets/+page.ts b/web/workspace/src/routes/w/[workspaceId]/tickets/+page.ts index e115cab7..04ab28e8 100644 --- a/web/workspace/src/routes/w/[workspaceId]/tickets/+page.ts +++ b/web/workspace/src/routes/w/[workspaceId]/tickets/+page.ts @@ -1,23 +1,56 @@ import { loadJson, workspaceApiPath } from "$lib/workspace/api/http"; +import type { TicketListResponse } from "$lib/generated/ticket-api"; import type { WorkspaceOrchestratorStatus } from "$lib/workspace/tickets/ticket-panel"; -import type { TicketListResponse } from "$lib/workspace/sidebar/types"; import type { PageLoad } from "./$types"; -export const load = (async ({ fetch, params }) => { - const [tickets, orchestrator] = await Promise.all([ - loadJson( - fetch, - `${workspaceApiPath(params.workspaceId, "/tickets")}?limit=1000`, +const LANE_STATES = { + "ready-planning": ["ready", "planning"], + "inprogress-queued": ["inprogress", "queued"], + "done-closed": ["done", "closed"], +} as const; + +export type TicketLaneId = keyof typeof LANE_STATES; + +export type TicketLanePage = { + states: readonly string[]; + response: TicketListResponse; +}; + +export const load: PageLoad = async ({ fetch, params }) => { + const workspaceId = params.workspaceId; + const [entries, orchestrator] = await Promise.all([ + Promise.all( + Object.entries(LANE_STATES).map(async ([laneId, states]) => { + const search = new URLSearchParams({ + limit: "30", + states: states.join(","), + }); + const response = await fetch( + `/api/w/${encodeURIComponent(workspaceId)}/tickets?${search}`, + ); + if (!response.ok) { + throw new Error( + `failed to load ${laneId} Ticket lane (${response.status})`, + ); + } + return [ + laneId, + { states: [...states], response: await response.json() }, + ] as const; + }), ), loadJson( fetch, - workspaceApiPath(params.workspaceId, "/orchestrator"), + workspaceApiPath(workspaceId, "/orchestrator"), ), ]); return { - workspaceId: params.workspaceId, - tickets, + workspaceId, + ticketLanes: Object.fromEntries(entries) as unknown as Record< + TicketLaneId, + TicketLanePage + >, orchestrator, }; -}) satisfies PageLoad; +};