merge: integrate orchestration

# Conflicts:
#	crates/flow/src/builtin.rs
#	crates/manifest/src/profile.rs
#	crates/worker/src/prompt/catalog.rs
#	crates/worker/src/prompt/system.rs
#	resources/flows/coder-review.dcdl
#	resources/prompts/role/coder.md
#	resources/prompts/role/orchestrator.md
#	web/workspace/src/lib/workspace/console/worker-console.ui.test.ts
#	web/workspace/src/lib/workspace/styles/tickets.css
#	web/workspace/src/routes/w/[workspaceId]/tickets/+page.svelte
#	web/workspace/src/routes/w/[workspaceId]/tickets/+page.ts
This commit is contained in:
2026-08-18 09:55:27 +09:00
37 changed files with 2964 additions and 727 deletions
+627 -67
View File
@@ -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<T> WorkspaceAuthority for T where T: ObjectiveAuthority + TicketAuthority +
pub trait TicketAuthority {
fn list_tickets(&self, limit: usize) -> Result<ProjectRecordList<TicketSummary>>;
fn list_ticket_page(&self, request: TicketListPageRequest) -> Result<TicketSummaryPage>;
fn query_tickets(&self, query: TicketQueryRequest) -> Result<TicketQueryResponse>;
fn ticket(&self, id: &str) -> Result<TicketDetail>;
fn show_ticket(&self, id: &str, query: TicketShowRequest) -> Result<TicketDetail>;
@@ -339,11 +343,358 @@ impl SqliteWorkspaceAuthority {
})
}
fn query_ticket_candidate_ids(
&self,
query: &TicketQueryRequest,
sort: TicketQuerySort,
after: Option<&(String, String)>,
limit: usize,
) -> Result<Vec<String>> {
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::<Vec<_>>();
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::<Vec<_>>();
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::<i64>().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::<std::result::Result<Vec<_>, _>>()?)
})
}
fn query_objective_candidate_ids(
&self,
query: &ObjectiveQueryRequest,
sort: ObjectiveQuerySort,
after: Option<&(String, String)>,
limit: usize,
) -> Result<Vec<String>> {
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::<Vec<_>>();
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::<i64>().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::<std::result::Result<Vec<_>, _>>()?)
})
}
fn read_ticket_detail(&self, id: &str, request: TicketShowRequest) -> Result<TicketDetail> {
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<TicketDetail> {
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
@@ -490,42 +841,100 @@ impl TicketAuthority for SqliteWorkspaceAuthority {
})
}
fn list_ticket_page(&self, request: TicketListPageRequest) -> Result<TicketSummaryPage> {
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::<Result<Vec<_>>>()?;
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::<Vec<_>>();
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<TicketQueryResponse> {
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,
));
}
@@ -537,7 +946,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 {
@@ -546,7 +959,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,
@@ -598,33 +1011,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::<Vec<_>>();
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,
@@ -639,7 +1055,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 {
@@ -648,7 +1068,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,
@@ -1358,6 +1778,35 @@ fn parse_offset_cursor(cursor: &str, field: &str) -> Result<usize> {
.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())
}
@@ -1677,8 +2126,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(
@@ -1703,31 +2152,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<String>,
@@ -1805,8 +2229,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(
@@ -1934,6 +2362,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<SqliteTicketListCursor> {
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",
@@ -2352,6 +2854,28 @@ 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"
);
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()),
@@ -2404,6 +2928,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()),
@@ -2423,6 +2976,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()),
+19
View File
@@ -12,6 +12,14 @@ pub struct ProjectRecordList<T> {
pub record_authority: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub struct TicketSummaryPage {
pub items: Vec<TicketSummary>,
pub page: QueryPage,
pub invalid_records: Vec<InvalidProjectRecord>,
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<String>,
#[serde(default)]
pub limit: Option<usize>,
#[serde(default)]
pub cursor: Option<String>,
}
#[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<TicketSummary>,
pub page: QueryPage,
pub invalid_records: Vec<InvalidProjectRecord>,
pub record_authority: String,
}
+8 -79
View File
@@ -97,38 +97,13 @@ pub struct CommitObservation {
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum RepositoryLookupError {
UnknownRepository {
id: RepositoryId,
},
UnsupportedProvider {
id: RepositoryId,
provider: String,
},
MissingDefaultSelector {
id: RepositoryId,
},
InvalidSelector {
id: RepositoryId,
selector: String,
},
CommitNotFound {
id: RepositoryId,
commit: String,
},
InvalidCommitRelation {
id: RepositoryId,
detail: String,
},
TargetMoved {
id: RepositoryId,
selector: String,
expected: String,
observed: Option<String>,
},
ProviderFailure {
id: RepositoryId,
operation: String,
},
UnknownRepository { id: RepositoryId },
UnsupportedProvider { id: RepositoryId, provider: String },
MissingDefaultSelector { id: RepositoryId },
InvalidSelector { id: RepositoryId, selector: String },
CommitNotFound { id: RepositoryId, commit: String },
InvalidCommitRelation { id: RepositoryId, detail: String },
ProviderFailure { id: RepositoryId, operation: String },
}
#[derive(Debug, Clone)]
@@ -306,45 +281,6 @@ impl RepositoryRegistryReader {
}
}
pub fn update_merge_target(
&self,
id: &str,
selector: &str,
expected_target: &str,
result_commit: &str,
) -> Result<(), RepositoryLookupError> {
let repository = self.merge_repository(id)?;
let target_ref = normalize_target_branch_selector(id, selector)?;
self.observe_commit(id, result_commit)?;
let status = Command::new("git")
.arg("-C")
.arg(&repository.path)
.args([
"update-ref",
target_ref.as_str(),
result_commit,
expected_target,
])
.status()
.map_err(|_| RepositoryLookupError::ProviderFailure {
id: id.into(),
operation: "guarded target update".into(),
})?;
if status.success() {
return Ok(());
}
let observed = self
.observe_merge_target(id, Some(selector))
.ok()
.map(|target| target.commit);
Err(RepositoryLookupError::TargetMoved {
id: id.into(),
selector: selector.into(),
expected: expected_target.into(),
observed,
})
}
fn merge_repository(&self, id: &str) -> Result<&ConfiguredRepository, RepositoryLookupError> {
let repository = self
.find(id)
@@ -755,20 +691,13 @@ mod tests {
vec![base.clone()]
);
reader.ensure_ancestor("main", &base, &source).unwrap();
reader
.update_merge_target("main", "main", &base, &source)
.unwrap();
assert_eq!(
reader
.observe_merge_target("main", Some("refs/heads/main"))
.unwrap()
.commit,
source
base
);
assert!(matches!(
reader.update_merge_target("main", "refs/heads/main", &base, &base),
Err(RepositoryLookupError::TargetMoved { .. })
));
assert!(matches!(
reader.ensure_ancestor("main", &source, &base),
Err(RepositoryLookupError::InvalidCommitRelation { .. })
+430 -123
View File
@@ -1064,11 +1064,18 @@ impl WorkspaceApi {
&self,
request: &WorkerSpawnRequest,
) -> ApiResult<()> {
let selected_repository_id =
let (selected_repository_id, selected_ref_selector) =
if let Some(working_directory) = request.resolved_working_directory_request.as_ref() {
let repository_id = working_directory.repository.id.as_str();
self.require_workspace_repository(repository_id)?;
Some(repository_id.to_string())
(
Some(repository_id.to_string()),
working_directory
.repository
.selector
.as_deref()
.map(str::to_owned),
)
} else if let Some(claim) = request.resolved_working_directory.as_ref() {
let workdir = self
.store
@@ -1080,20 +1087,48 @@ impl WorkspaceApi {
)))
})?;
self.require_workspace_repository(&workdir.repository_id)?;
Some(workdir.repository_id)
(Some(workdir.repository_id), workdir.creation_selector)
} else {
None
(None, None)
};
if let WorkerSpawnIntent::TicketRole { ticket_id, .. } = &request.intent {
let ticket = self.authority.ticket(ticket_id)?;
if let Some(repository_id) = ticket.repository_id.as_deref() {
self.require_workspace_repository(repository_id)?;
if selected_repository_id.as_deref() != Some(repository_id) {
return Err(ApiError::from(Error::Config(format!(
"Ticket `{ticket_id}` targets repository `{repository_id}`, but the Worker launch does not resolve that repository in this Workspace"
))));
}
// Workdir-less Ticket Workers cannot execute repository implementation.
// Preserve that control-plane launch while still validating any persisted
// target (including its Workspace ownership) when one exists.
if selected_repository_id.is_none() && ticket.repository_id.is_none() {
return Ok(());
}
let repository_id = ticket.repository_id.as_deref().ok_or_else(|| {
ApiError::from(Error::Config(
"Ticket implementation target must be validated and persisted before spawning a Ticket Worker".to_owned(),
))
})?;
let ref_selector = ticket.ref_selector.as_deref().ok_or_else(|| {
ApiError::from(Error::Config(
"Ticket implementation target selector must be validated and persisted before spawning a Ticket Worker".to_owned(),
))
})?;
self.require_workspace_repository(repository_id)?;
self.repository_reader()
.observe_merge_target(repository_id, Some(ref_selector))
.map_err(|error| {
ApiError::from(Error::Config(format!(
"Ticket implementation target is no longer resolvable: {error:?}"
)))
})?;
if selected_repository_id.as_deref() != Some(repository_id) {
return Err(ApiError::from(Error::Config(format!(
"Ticket `{ticket_id}` targets repository `{repository_id}`, but the Worker launch resolves `{}`",
selected_repository_id.as_deref().unwrap_or("none")
))));
}
if selected_ref_selector.as_deref() != Some(ref_selector) {
return Err(ApiError::from(Error::Config(format!(
"Ticket `{ticket_id}` targets selector `{ref_selector}`, but the Worker launch resolves `{}`",
selected_ref_selector.as_deref().unwrap_or("none")
))));
}
}
Ok(())
@@ -1319,8 +1354,8 @@ pub fn build_router(api: WorkspaceApi) -> Router {
post(scoped_set_ticket_workflow_state),
)
.route(
"/api/w/{workspace_id}/tickets/{id}/intake-ready",
post(scoped_prepare_ticket_intake_ready),
"/api/w/{workspace_id}/tickets/{id}/workflow/mark-ready",
post(scoped_mark_ticket_ready),
)
.route(
"/api/w/{workspace_id}/tickets/{id}/workflow/queue",
@@ -1401,6 +1436,10 @@ pub fn build_router(api: WorkspaceApi) -> Router {
"/api/w/{workspace_id}/tickets/{id}/state",
post(scoped_transition_ticket_state),
)
.route(
"/api/w/{workspace_id}/tickets/{id}/ready",
post(scoped_mark_ticket_ready_from_browser),
)
.route(
"/api/w/{workspace_id}/tickets/{id}/events",
post(scoped_append_ticket_event),
@@ -2233,6 +2272,9 @@ struct ObjectiveEditRequest {
#[derive(Debug, Deserialize)]
struct TicketListQuery {
limit: Option<usize>,
cursor: Option<String>,
/// Comma-separated workflow states. Repeated lane requests normally pass one state group.
states: Option<String>,
}
#[derive(Debug, Deserialize)]
@@ -3111,6 +3153,55 @@ struct BrowserCloseTicketRequest {
resolution: String,
}
#[derive(Clone)]
struct WorkspaceTicketTargetAuthority {
api: WorkspaceApi,
}
impl ticket::TicketTargetAuthority for WorkspaceTicketTargetAuthority {
fn resolve_target(
&self,
workspace_id: &str,
repository_id: Option<&str>,
ref_selector: Option<&str>,
) -> ticket::Result<ticket::ResolvedTicketTarget> {
if workspace_id != self.api.config.workspace_id {
return Err(ticket::TicketError::UnknownTargetRepository(
repository_id.unwrap_or_default().to_owned(),
));
}
let repository_id = repository_id
.map(str::trim)
.filter(|value| !value.is_empty())
.ok_or(ticket::TicketError::MissingTargetRepository)?;
let repository = self
.api
.store
.get_repository(workspace_id, repository_id)
.map_err(|error| ticket::TicketError::Conflict(error.to_string()))?
.ok_or_else(|| {
ticket::TicketError::UnknownTargetRepository(repository_id.to_owned())
})?;
let selector = ref_selector
.map(str::trim)
.filter(|value| !value.is_empty())
.or(repository.default_ref.as_deref())
.ok_or_else(|| ticket::TicketError::MissingTargetSelector(repository_id.to_owned()))?;
self.api
.repository_reader()
.observe_merge_target(repository_id, Some(selector))
.map_err(|error| ticket::TicketError::InvalidTargetSelector {
repository_id: repository_id.to_owned(),
selector: selector.to_owned(),
reason: format!("{error:?}"),
})?;
Ok(ticket::ResolvedTicketTarget {
repository_id: repository_id.to_owned(),
ref_selector: selector.to_owned(),
})
}
}
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}")))?;
@@ -3118,7 +3209,10 @@ fn browser_ticket_backend(api: &WorkspaceApi) -> Result<SqliteTicketBackend> {
api.config.database_path.clone(),
api.config.workspace_id.clone(),
)?
.with_record_language(config.ticket_record_language()))
.with_record_language(config.ticket_record_language())
.with_target_authority(Arc::new(WorkspaceTicketTargetAuthority {
api: api.clone(),
})))
}
fn browser_ticket_detail(api: &WorkspaceApi, ticket_id: &str) -> ApiResult<Json<TicketDetail>> {
@@ -3212,6 +3306,26 @@ async fn scoped_append_ticket_event(
browser_ticket_detail(&api, &path.id)
}
async fn scoped_mark_ticket_ready_from_browser(
State(api): State<WorkspaceApi>,
AxumPath(path): AxumPath<ScopedRecordPath>,
Json(request): Json<TicketMarkReadyRequest>,
) -> ApiResult<Json<TicketDetail>> {
validate_workspace_scope(&api, &path.workspace_id)?;
browser_ticket_backend(&api)?
.mark_ready(
TicketIdOrSlug::Id(path.id.clone()),
ticket::TicketMarkReady {
operation_key: request.operation_key,
reason: request.reason,
author: Some("web".to_owned()),
intake_summary: None,
},
)
.map_err(Error::from)?;
browser_ticket_detail(&api, &path.id)
}
async fn scoped_queue_ticket(
State(api): State<WorkspaceApi>,
AxumPath(path): AxumPath<ScopedRecordPath>,
@@ -3273,7 +3387,10 @@ async fn execute_worker_ticket_rest_operation(
api.config.workspace_id.clone(),
)
.map_err(Error::from)?
.with_record_language(config.ticket_record_language());
.with_record_language(config.ticket_record_language())
.with_target_authority(Arc::new(WorkspaceTicketTargetAuthority {
api: api.clone(),
}));
let operation_kind = ticket_mutation_operation_kind(&operation);
let is_mutation = operation_kind != "read";
let target = ticket_mutation_target(&operation).cloned();
@@ -3425,6 +3542,15 @@ async fn scoped_create_ticket_record(
headers: HeaderMap,
Json(input): Json<ticket::NewTicket>,
) -> ApiResult<Json<ticket::TicketRef>> {
if input
.workflow_state
.is_some_and(|state| state != TicketWorkflowState::Planning)
{
return Err(settings_bad_request(
"ticket_create_state_bypass",
"Ticket creation must start in planning; use guarded workflow operations for later states",
));
}
let result = execute_worker_ticket_rest_operation(
&api,
&path.workspace_id,
@@ -3538,9 +3664,12 @@ async fn scoped_add_ticket_intake_summary(
}
#[derive(Debug, Deserialize)]
struct TicketIntakeReadyRequest {
summary: ticket::TicketIntakeSummary,
change: TicketStateChange,
struct TicketMarkReadyRequest {
operation_key: String,
#[serde(default)]
reason: Option<String>,
#[serde(default)]
intake_summary: Option<ticket::TicketIntakeSummary>,
}
async fn scoped_set_ticket_state_field(
@@ -3582,24 +3711,31 @@ async fn scoped_set_ticket_workflow_state(
ticket_rest_unit(result)
}
async fn scoped_prepare_ticket_intake_ready(
async fn scoped_mark_ticket_ready(
State(api): State<WorkspaceApi>,
AxumPath((workspace_id, id)): AxumPath<(String, String)>,
headers: HeaderMap,
Json(request): Json<TicketIntakeReadyRequest>,
) -> ApiResult<StatusCode> {
Json(request): Json<TicketMarkReadyRequest>,
) -> ApiResult<Json<ticket::Ticket>> {
let result = execute_worker_ticket_rest_operation(
&api,
&workspace_id,
headers,
TicketBackendOperation::MarkIntakeReady {
TicketBackendOperation::MarkReady {
id: TicketIdOrSlug::Query(id),
summary: request.summary,
change: request.change,
request: ticket::TicketMarkReady {
operation_key: request.operation_key,
reason: request.reason,
author: None,
intake_summary: request.intake_summary,
},
},
)
.await?;
ticket_rest_unit(result)
ticket_rest_result(result, |result| match result {
TicketBackendOperationResult::Ticket(ticket) => Some(ticket),
_ => None,
})
}
async fn scoped_queue_ticket_record(
@@ -3773,6 +3909,37 @@ fn repository_merge_evidence_error(error: RepositoryLookupError) -> ApiError {
.into()
}
fn recorded_merge_completion<'a>(
thread: &'a [merge_request::MergeRequestThreadEvent],
operation_id: &str,
) -> Option<&'a merge_request::MergeEvent> {
thread.iter().find_map(|event| match event {
merge_request::MergeRequestThreadEvent::Merge(event)
if event.operation_id == operation_id =>
{
Some(event)
}
_ => None,
})
}
fn require_completed_target_observation(
observed: &str,
target_ref_before: &str,
target_ref_after: &str,
) -> ApiResult<()> {
if observed == target_ref_after {
return Ok(());
}
if observed == target_ref_before {
return Err(Error::InvalidInput(
"target selector is still at target_ref_before; push the verified result from the Orchestrator Workdir before MergeRequestComplete".into(),
)
.into());
}
Err(Error::InvalidInput("target selector moved outside completion evidence".into()).into())
}
async fn scoped_show_merge_request(
State(api): State<WorkspaceApi>,
AxumPath((workspace_id, ticket_id)): AxumPath<(String, String)>,
@@ -4106,19 +4273,40 @@ async fn scoped_complete_merge_request(
require_workspace_access(&workspace_id, &api)?;
let source = authenticate_worker_mutation_source(&api, &workspace_id, &headers)?;
require_online_workspace_orchestrator_source(&api, &source)?;
let store = merge_request_store(&api, &workspace_id)?;
let mr = store.get(&workspace_id, &ticket_id)?;
let repositories = api.repository_reader();
if let Some(existing) = recorded_merge_completion(&mr.thread, &input.operation_id) {
let replay = merge_request::CompleteMergeRequest {
ticket_id,
operation_id: input.operation_id,
approval_event_id: input.approval_event_id,
current_subject_ref: existing.approved_source_ref.clone(),
target_ref_before: input.target_ref_before,
target_ref_after: input.target_ref_after,
strategy: input.strategy,
resolution: input.resolution,
auth: merge_request::MergeRequestAuth {
workspace_id,
repository_id: mr.repository_id.clone(),
runtime_id: source.runtime_id,
worker_id: source.worker_id,
assignment_id: String::new(),
},
now: Utc::now(),
};
return store.complete(replay).map(Json).map_err(Into::into);
}
let assignment = api
.store
.get_current_ticket_worker_assignment(&workspace_id, &ticket_id)?
.ok_or_else(|| {
Error::TicketAssignmentConflict("Ticket has no current assigned Coder".into())
})?;
let store = merge_request_store(&api, &workspace_id)?;
let mr = store.get(&workspace_id, &ticket_id)?;
let selector = mr
.selector_from
.as_deref()
.ok_or_else(|| Error::InvalidInput("selector_from requires repair".into()))?;
let repositories = api.repository_reader();
let current_source_ref = repositories
.observe_merge_target(&mr.repository_id, Some(selector))
.map_err(repository_merge_evidence_error)?
@@ -4126,12 +4314,11 @@ async fn scoped_complete_merge_request(
let observed = repositories
.observe_merge_target(&mr.repository_id, Some(&mr.selector_to))
.map_err(repository_merge_evidence_error)?;
if observed.commit != input.target_ref_before && observed.commit != input.target_ref_after {
return Err(Error::InvalidInput(
"target selector moved outside completion evidence".into(),
)
.into());
}
require_completed_target_observation(
&observed.commit,
&input.target_ref_before,
&input.target_ref_after,
)?;
let completion = merge_request::CompleteMergeRequest {
ticket_id,
operation_id: input.operation_id,
@@ -4151,31 +4338,7 @@ async fn scoped_complete_merge_request(
now: Utc::now(),
};
store.validate_completion(&completion)?;
let already = observed.commit == input.target_ref_after;
if !already {
repositories
.update_merge_target(
&mr.repository_id,
&mr.selector_to,
&input.target_ref_before,
&input.target_ref_after,
)
.map_err(repository_merge_evidence_error)?
}
match store.complete(completion) {
Ok(v) => Ok(Json(v)),
Err(e) => {
if !already {
let _ = repositories.update_merge_target(
&mr.repository_id,
&mr.selector_to,
&input.target_ref_after,
&input.target_ref_before,
);
}
Err(e.into())
}
}
store.complete(completion).map(Json).map_err(Into::into)
}
fn reject_non_browser_reopen_auth(headers: &HeaderMap) -> Result<()> {
@@ -4381,7 +4544,7 @@ fn ticket_mutation_target(operation: &TicketBackendOperation) -> Option<&TicketI
| TicketBackendOperation::AddIntakeSummary { id, .. }
| TicketBackendOperation::SetStateField { id, .. }
| TicketBackendOperation::SetWorkflowState { id, .. }
| TicketBackendOperation::MarkIntakeReady { id, .. }
| TicketBackendOperation::MarkReady { id, .. }
| TicketBackendOperation::QueueReady { id, .. }
| TicketBackendOperation::Close { id, .. }
| TicketBackendOperation::AddTicketRelation { id, .. }
@@ -4422,11 +4585,11 @@ fn bind_worker_ticket_operation_source(
| TicketBackendOperation::SetStateField { change, .. }
| TicketBackendOperation::SetWorkflowState { change, .. } => change.author = Some(author),
TicketBackendOperation::AddIntakeSummary { summary, .. } => summary.author = Some(author),
TicketBackendOperation::MarkIntakeReady {
summary, change, ..
} => {
summary.author = Some(author.clone());
change.author = Some(author);
TicketBackendOperation::MarkReady { request, .. } => {
request.author = Some(author.clone());
if let Some(summary) = request.intake_summary.as_mut() {
summary.author = Some(author);
}
}
TicketBackendOperation::QueueReady { queued_by, .. } => *queued_by = author,
TicketBackendOperation::AddTicketRelation { relation, .. } => {
@@ -4448,7 +4611,7 @@ fn ticket_mutation_operation_kind(operation: &TicketBackendOperation) -> &'stati
TicketBackendOperation::AddIntakeSummary { .. } => "add_intake_summary",
TicketBackendOperation::SetStateField { .. } => "set_state_field",
TicketBackendOperation::SetWorkflowState { .. } => "set_workflow_state",
TicketBackendOperation::MarkIntakeReady { .. } => "mark_intake_ready",
TicketBackendOperation::MarkReady { .. } => "mark_ready",
TicketBackendOperation::QueueReady { .. } => "queue_ready",
TicketBackendOperation::Close { .. } => "close",
TicketBackendOperation::AddTicketRelation { .. } => "add_relation",
@@ -8329,17 +8492,35 @@ async fn list_tickets(
State(api): State<WorkspaceApi>,
Query(query): Query<TicketListQuery>,
) -> ApiResult<Json<crate::records::TicketListResponse>> {
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::<Vec<_>>()
})
.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,
}))
@@ -11896,7 +12077,26 @@ impl From<Error> for ApiError {
ticket::TicketError::NotFound(_) => "ticket_not_found",
ticket::TicketError::Ambiguous { .. } => "ticket_ambiguous",
ticket::TicketError::Locked { .. } => "ticket_locked",
ticket::TicketError::Conflict(_) => "ticket_conflict",
ticket::TicketError::Conflict(_)
| ticket::TicketError::StaleWorkflowState { .. }
| ticket::TicketError::InvalidWorkflowTransition { .. }
| ticket::TicketError::BlockingRelations(_)
| ticket::TicketError::OperationFingerprintMismatch { .. } => "ticket_conflict",
ticket::TicketError::MissingTargetRepository => {
"ticket_target_repository_missing"
}
ticket::TicketError::UnknownTargetRepository(_) => {
"ticket_target_repository_unknown"
}
ticket::TicketError::MissingTargetSelector(_) => {
"ticket_target_selector_missing"
}
ticket::TicketError::InvalidTargetSelector { .. } => {
"ticket_target_selector_invalid"
}
ticket::TicketError::TargetAuthorityUnavailable => {
"ticket_target_authority_unavailable"
}
ticket::TicketError::InvalidPathComponent(_)
| ticket::TicketError::PathEscapesRoot { .. } => "invalid_ticket_request",
ticket::TicketError::Io { .. }
@@ -12465,6 +12665,54 @@ mod tests {
);
}
#[test]
fn recorded_completion_replay_is_identified_before_later_target_observation() {
let event = merge_request::MergeEvent {
event_id: "merge-event".into(),
sequence: 1,
operation_id: "operation".into(),
approval_event_id: "approval".into(),
approved_source_ref: "source".into(),
target_ref_before: "before".into(),
target_ref_after: "after".into(),
strategy: merge_request::MergeStrategy::FastForward,
resolution: merge_request::ConflictResolution::None,
merged_by: merge_request::WorkerIdentity {
runtime_id: "runtime".into(),
worker_id: "orchestrator".into(),
},
created_at: Utc::now(),
};
let thread = vec![merge_request::MergeRequestThreadEvent::Merge(event.clone())];
assert_eq!(
recorded_merge_completion(&thread, "operation"),
Some(&event)
);
assert!(recorded_merge_completion(&thread, "different").is_none());
assert!(require_completed_target_observation("later", "before", "after").is_err());
}
#[test]
fn merge_request_completion_records_only_an_observed_remote_target_update() {
require_completed_target_observation("after", "before", "after").unwrap();
let not_pushed =
require_completed_target_observation("before", "before", "after").unwrap_err();
assert!(matches!(
not_pushed.error,
Error::InvalidInput(ref message)
if message.contains("push the verified result from the Orchestrator Workdir")
));
let moved = require_completed_target_observation("other", "before", "after").unwrap_err();
assert!(matches!(
moved.error,
Error::InvalidInput(ref message)
if message.contains("moved outside completion evidence")
));
}
#[test]
fn worker_ticket_assignment_projects_coder_intent_and_run_acceptance() {
let initial_submit = vec![
@@ -13562,7 +13810,7 @@ mod tests {
config.repositories = vec![ConfiguredRepository {
id: TEST_REPOSITORY_ID.to_string(),
provider: "git".to_string(),
uri: ".".to_string(),
uri: workspace_root.display().to_string(),
path: workspace_root,
display_name: Some("Test Repository".to_string()),
default_selector: Some("HEAD".to_string()),
@@ -13719,7 +13967,7 @@ mod tests {
fn init_clean_git_workspace(path: &std::path::Path) {
for args in [
vec!["init"],
vec!["init", "--initial-branch=develop"],
vec!["config", "user.email", "test@example.invalid"],
vec!["config", "user.name", "Yoi Test"],
] {
@@ -13751,6 +13999,88 @@ mod tests {
}
}
#[tokio::test]
async fn mark_ready_resolves_workspace_target_and_closes_lifecycle_bypasses() {
let dir = tempfile::tempdir().unwrap();
init_clean_git_workspace(dir.path());
let api = test_api(dir.path()).await;
let backend = browser_ticket_backend(&api).unwrap();
let mut input = ticket::NewTicket::new("Validated target");
input.repository_id = Some(TEST_REPOSITORY_ID.to_owned());
input.ref_selector = Some("develop".to_owned());
let ticket_ref = backend.create(input).unwrap();
let request = ticket::TicketMarkReady {
operation_key: "ready-server-test".to_owned(),
reason: Some("target accepted".to_owned()),
author: Some("test".to_owned()),
intake_summary: None,
};
let ready = backend
.mark_ready(TicketIdOrSlug::Id(ticket_ref.id.clone()), request.clone())
.unwrap();
assert_eq!(ready.meta.workflow_state, TicketWorkflowState::Ready);
assert_eq!(
ready.meta.repository_id.as_deref(),
Some(TEST_REPOSITORY_ID)
);
assert_eq!(ready.meta.ref_selector.as_deref(), Some("develop"));
assert_eq!(
backend
.mark_ready(TicketIdOrSlug::Id(ticket_ref.id.clone()), request)
.unwrap()
.events
.iter()
.filter(|event| event.attributes.contains_key("operation_key"))
.count(),
1
);
assert!(matches!(
backend.edit_item(
TicketIdOrSlug::Id(ticket_ref.id.clone()),
ticket::TicketItemEdit {
target: Some(ticket::TicketTargetEdit::Set {
repository_id: TEST_REPOSITORY_ID.to_owned(),
ref_selector: Some("other".to_owned()),
}),
..Default::default()
},
),
Err(ticket::TicketError::Conflict(_))
));
let mut missing = ticket::NewTicket::new("Missing target");
missing.repository_id = Some("unknown".to_owned());
let missing = backend.create(missing).unwrap();
assert!(matches!(
backend.mark_ready(
TicketIdOrSlug::Id(missing.id.clone()),
ticket::TicketMarkReady {
operation_key: "missing-repository".to_owned(),
reason: None,
author: None,
intake_summary: None,
},
),
Err(ticket::TicketError::UnknownTargetRepository(_))
));
assert_eq!(
backend
.show(TicketIdOrSlug::Id(missing.id))
.unwrap()
.meta
.workflow_state,
TicketWorkflowState::Planning
);
assert!(matches!(
backend.set_workflow_state(
TicketIdOrSlug::Id(ticket_ref.id),
TicketStateChange::new("ready", "queued", "bypass", "must use TicketQueue",),
),
Err(ticket::TicketError::InvalidWorkflowTransition { .. })
));
}
#[test]
fn worker_source_actor_roles_use_canonical_vocabulary() {
assert_eq!(worker_source_actor_role(true, false), "coder");
@@ -13792,6 +14122,7 @@ mod tests {
#[tokio::test]
async fn orchestrator_ticket_notifications_project_authoritative_post_mutation_state() {
let dir = tempfile::tempdir().unwrap();
init_clean_git_workspace(dir.path());
let (api, execution) = test_api_with_recording_backend(dir.path()).await;
let source_worker = api
.runtime
@@ -13849,29 +14180,24 @@ mod tests {
let orchestrator = started.worker.unwrap().worker;
execution.take_inputs();
let ticket = browser_ticket_backend(&api)
.unwrap()
.create(ticket::NewTicket::new("Bounded notification"))
.unwrap();
let mut input = ticket::NewTicket::new("Bounded notification");
input.repository_id = Some(TEST_REPOSITORY_ID.to_owned());
input.ref_selector = Some("develop".to_owned());
let ticket = browser_ticket_backend(&api).unwrap().create(input).unwrap();
let ticket_id = TicketIdOrSlug::Id(ticket.id.clone());
let operations = [
TicketBackendOperation::SetWorkflowState {
TicketBackendOperation::MarkReady {
id: ticket_id.clone(),
change: TicketStateChange::new(
"planning",
"ready",
"ready for implementation",
"test transition",
),
request: ticket::TicketMarkReady {
operation_key: "notification-ready".to_owned(),
reason: Some("ready for implementation".to_owned()),
author: None,
intake_summary: None,
},
},
TicketBackendOperation::SetWorkflowState {
TicketBackendOperation::QueueReady {
id: ticket_id.clone(),
change: TicketStateChange::new(
"ready",
"queued",
"queued for implementation",
"test transition",
),
queued_by: "spoofed".to_owned(),
},
TicketBackendOperation::SetWorkflowState {
id: ticket_id.clone(),
@@ -14276,30 +14602,11 @@ mod tests {
let dir = tempfile::tempdir().unwrap();
let api = test_api(dir.path()).await;
let backend = browser_ticket_backend(&api).unwrap();
let ticket_ref = backend
.create(ticket::NewTicket::new("Recover queued work"))
.unwrap();
backend
.mark_intake_ready(
TicketIdOrSlug::Id(ticket_ref.id.clone()),
ticket::TicketIntakeSummary {
author: Some("intake".to_string()),
body: MarkdownText::new("Ready"),
references: Vec::new(),
},
ticket::TicketStateChange {
from: "planning".to_string(),
to: "ready".to_string(),
reason: "ready".to_string(),
author: Some("intake".to_string()),
body: MarkdownText::new("Ready"),
references: Vec::new(),
},
)
.unwrap();
backend
.queue_ready(TicketIdOrSlug::Id(ticket_ref.id.clone()), "browser-user")
.unwrap();
let mut input = ticket::NewTicket::new("Recover queued work");
input.workflow_state = Some(TicketWorkflowState::Queued);
input.repository_id = Some(TEST_REPOSITORY_ID.to_owned());
input.ref_selector = Some("HEAD".to_owned());
let ticket_ref = backend.create(input).unwrap();
*api.orchestrator_attention_fingerprint.lock().unwrap() = Some(ticket_ref.id.clone());
let Json(started) = scoped_start_workspace_orchestrator(
@@ -14725,6 +15032,7 @@ mod tests {
#[tokio::test]
async fn ticket_browser_endpoints_mutate_typed_backend_and_return_thread() {
let dir = tempfile::tempdir().unwrap();
init_clean_git_workspace(dir.path());
let api = test_api(dir.path()).await;
let ticket_ref = browser_ticket_backend(&api)
.unwrap()
@@ -14764,7 +15072,7 @@ mod tests {
replace_all: false,
target: Some(TicketTargetEdit::Set {
repository_id: "main".to_string(),
ref_selector: Some("feature/api".to_string()),
ref_selector: Some("develop".to_string()),
}),
author: Some("browser-user".to_string()),
}),
@@ -14774,7 +15082,7 @@ mod tests {
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"));
assert_eq!(edited.ref_selector.as_deref(), Some("develop"));
assert_eq!(edited.assignee, None);
assert_eq!(edited.relations.outgoing.len(), 1);
assert_eq!(edited.relations.outgoing[0].target, related_ticket_id);
@@ -14795,14 +15103,13 @@ mod tests {
event.kind == "comment" && event.body.as_deref() == Some("API comment")
}));
let Json(ready) = scoped_transition_ticket_state(
let Json(ready) = scoped_mark_ticket_ready_from_browser(
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()),
Json(TicketMarkReadyRequest {
operation_key: "browser-ready".to_owned(),
reason: Some("intake complete".to_owned()),
intake_summary: None,
}),
)
.await
+33 -10
View File
@@ -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<String>) = 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(),