server: route Ticket mutation notifications

This commit is contained in:
2026-07-31 20:15:11 +09:00
parent 32d3f8f75e
commit 005f6cb498
16 changed files with 2385 additions and 367 deletions
+43 -11
View File
@@ -329,6 +329,8 @@ pub struct WorkerSpawnRequest {
pub resolved_working_directory: Option<WorkingDirectoryClaim>,
#[serde(skip, default)]
pub resolved_config_bundle: Option<ConfigBundle>,
#[serde(skip, default)]
pub resolved_workspace_api: Option<WorkspaceApiRef>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
@@ -1703,13 +1705,16 @@ impl WorkspaceWorkerRuntime for EmbeddedWorkerRuntime {
initial_input: request.initial_input.clone(),
working_directory_request: request.resolved_working_directory_request.clone(),
working_directory: request.resolved_working_directory.clone(),
workspace_api: self
.backend_base_url
.as_ref()
.map(|base_url| WorkspaceApiRef {
workspace_id: self.workspace_id.clone(),
base_url: base_url.clone(),
}),
workspace_api: request.resolved_workspace_api.clone().or_else(|| {
self.backend_base_url
.as_ref()
.map(|base_url| WorkspaceApiRef {
workspace_id: self.workspace_id.clone(),
base_url: base_url.clone(),
runtime_id: Some(self.runtime_id.clone()),
access_token: None,
})
}),
};
match self.runtime.create_worker(create_request) {
Ok(detail) => WorkerSpawnResult {
@@ -2677,9 +2682,13 @@ impl WorkspaceWorkerRuntime for RemoteWorkerRuntime {
initial_input: request.initial_input.clone(),
working_directory_request: request.resolved_working_directory_request.clone(),
working_directory: request.resolved_working_directory.clone(),
workspace_api: Some(WorkspaceApiRef {
workspace_id: self.workspace_id.clone(),
base_url: self.backend_base_url.clone(),
workspace_api: request.resolved_workspace_api.clone().or_else(|| {
Some(WorkspaceApiRef {
workspace_id: self.workspace_id.clone(),
base_url: self.backend_base_url.clone(),
runtime_id: Some(self.runtime_id.clone()),
access_token: None,
})
}),
};
match self.post_json::<_, RuntimeHttpWorkerResponse>("/v1/workers", &create) {
@@ -3121,8 +3130,11 @@ fn embedded_profile_path(profile: &ProfileSelector) -> Result<String, String> {
fn embedded_profile_label(profile: &ProfileSelector) -> Option<String> {
Some(match profile {
ProfileSelector::Builtin(name) | ProfileSelector::Named(name) => {
if name.strip_prefix("builtin:").unwrap_or(name) == MEMORY_CONSOLIDATION_PROFILE {
let builtin_name = name.strip_prefix("builtin:").unwrap_or(name);
if builtin_name == MEMORY_CONSOLIDATION_PROFILE {
MEMORY_CONSOLIDATION_PROFILE.to_string()
} else if builtin_name == WORKSPACE_ORCHESTRATOR_PROFILE {
WORKSPACE_ORCHESTRATOR_PROFILE.to_string()
} else {
safe_display_hint(name)
}
@@ -3132,6 +3144,8 @@ fn embedded_profile_label(profile: &ProfileSelector) -> Option<String> {
const MEMORY_CONSOLIDATION_PROFILE: &str = "memory-consolidation";
const MEMORY_CONSOLIDATION_SINGLETON_KEY: &str = "workspace-memory-consolidation";
const WORKSPACE_ORCHESTRATOR_PROFILE: &str = "orchestrator";
pub(crate) const WORKSPACE_ORCHESTRATOR_SINGLETON_KEY: &str = "workspace-orchestrator";
struct WorkerDisplayMetadata {
display_name: String,
@@ -3160,6 +3174,20 @@ fn worker_display_metadata(
tags,
};
}
if profile_label == Some(WORKSPACE_ORCHESTRATOR_PROFILE) {
let mut tags = vec!["orchestrator".to_string(), "singleton".to_string()];
if internal {
tags.insert(0, "internal".to_string());
}
return WorkerDisplayMetadata {
display_name: requested_display_name
.filter(|value| !value.trim().is_empty())
.map(safe_display_hint)
.unwrap_or_else(|| "Workspace Orchestrator".to_string()),
singleton_key: Some(WORKSPACE_ORCHESTRATOR_SINGLETON_KEY.to_string()),
tags,
};
}
let display_name = requested_display_name
.filter(|value| !value.trim().is_empty())
.map(safe_display_hint)
@@ -4154,6 +4182,7 @@ mod tests {
resolved_working_directory_request: None,
resolved_working_directory: None,
resolved_config_bundle: None,
resolved_workspace_api: None,
}
}
@@ -4280,6 +4309,7 @@ mod tests {
resolved_working_directory_request: None,
resolved_working_directory: None,
resolved_config_bundle: None,
resolved_workspace_api: None,
},
)
.unwrap();
@@ -4376,6 +4406,7 @@ mod tests {
resolved_working_directory_request: None,
resolved_working_directory: None,
resolved_config_bundle: None,
resolved_workspace_api: None,
},
)
.unwrap();
@@ -4408,6 +4439,7 @@ mod tests {
resolved_working_directory_request: None,
resolved_working_directory: None,
resolved_config_bundle: None,
resolved_workspace_api: None,
},
)
.unwrap();
+4
View File
@@ -85,6 +85,10 @@ pub enum Error {
UnknownRepository(String),
#[error("workspace id does not match this Workspace backend")]
WorkspaceIdMismatch,
#[error("Ticket assignment conflict: {0}")]
TicketAssignmentConflict(String),
#[error("Worker Workspace authentication failed: {0}")]
WorkerWorkspaceAuthentication(String),
#[error("workspace identity error: {0}")]
WorkspaceIdentity(String),
#[error("store error: {0}")]
File diff suppressed because it is too large Load Diff
+862 -6
View File
@@ -87,6 +87,16 @@ const MIGRATIONS: &[Migration] = &[
name: "remove unused control-plane Ticket tables",
apply: remove_unused_control_plane_ticket_tables,
},
Migration {
version: 15,
name: "ticket worker current assignment authority",
apply: create_ticket_worker_assignment_tables,
},
Migration {
version: 16,
name: "worker workspace credentials and Ticket notification outbox",
apply: create_ticket_notification_tables,
},
];
struct Migration {
@@ -232,6 +242,66 @@ pub struct WorkerRegistryRecord {
pub updated_at: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub struct TicketWorkerAssignmentRecord {
pub workspace_id: String,
pub ticket_id: String,
pub assignment_id: String,
pub runtime_id: String,
pub worker_id: String,
pub assigned_by: String,
pub assigned_at: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub struct TicketWorkerAssignmentEventRecord {
pub workspace_id: String,
pub ticket_id: String,
pub event_id: String,
pub action: String,
pub assignment_id: Option<String>,
pub previous_assignment_id: Option<String>,
pub actor: String,
pub created_at: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub struct TicketWorkerAssignmentUpdate {
pub current: TicketWorkerAssignmentRecord,
pub previous: Option<TicketWorkerAssignmentRecord>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub struct WorkerWorkspaceCredentialRecord {
pub credential_id: String,
pub token: String,
pub workspace_id: String,
pub runtime_id: String,
pub worker_id: Option<String>,
pub created_at: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub struct TicketNotificationRecipient {
pub runtime_id: String,
pub worker_id: String,
pub recipient_kind: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub struct TicketNotificationDeliveryRecord {
pub notification_id: String,
pub workspace_id: String,
pub ticket_id: String,
pub event_sequence: i64,
pub source_runtime_id: String,
pub source_worker_id: String,
pub recipient_runtime_id: String,
pub recipient_worker_id: String,
pub recipient_kind: String,
pub attempts: i64,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub struct WorkdirRegistryRecord {
pub workspace_id: String,
@@ -483,6 +553,83 @@ pub trait ControlPlaneStore: Send + Sync {
runtime_worker_id: u64,
) -> Result<bool>;
fn get_current_ticket_worker_assignment(
&self,
workspace_id: &str,
ticket_id: &str,
) -> Result<Option<TicketWorkerAssignmentRecord>>;
fn set_current_ticket_worker_assignment(
&self,
record: &TicketWorkerAssignmentRecord,
expected_assignment_id: Option<&str>,
event_id: &str,
) -> Result<TicketWorkerAssignmentUpdate>;
fn clear_current_ticket_worker_assignment(
&self,
workspace_id: &str,
ticket_id: &str,
expected_assignment_id: Option<&str>,
event_id: &str,
actor: &str,
created_at: &str,
) -> Result<Option<TicketWorkerAssignmentRecord>>;
fn list_ticket_worker_assignment_events(
&self,
workspace_id: &str,
ticket_id: &str,
limit: usize,
) -> Result<Vec<TicketWorkerAssignmentEventRecord>>;
fn upsert_worker_workspace_credential(
&self,
record: &WorkerWorkspaceCredentialRecord,
) -> Result<()>;
fn authenticate_worker_workspace_credential(
&self,
token: &str,
workspace_id: &str,
worker_id: &str,
) -> Result<Option<WorkerWorkspaceCredentialRecord>>;
fn enqueue_ticket_notification(
&self,
notification_id: &str,
workspace_id: &str,
ticket_id: &str,
event_sequence: i64,
source_runtime_id: &str,
source_worker_id: &str,
previous_state: &str,
current_state: &str,
created_at: &str,
recipients: &[TicketNotificationRecipient],
) -> Result<()>;
fn list_pending_ticket_notification_deliveries(
&self,
workspace_id: &str,
limit: usize,
) -> Result<Vec<TicketNotificationDeliveryRecord>>;
fn count_ticket_notification_deliveries_for_recipient(
&self,
workspace_id: &str,
ticket_id: &str,
runtime_id: &str,
worker_id: &str,
) -> Result<usize>;
fn mark_ticket_notification_delivered(
&self,
notification_id: &str,
recipient_runtime_id: &str,
recipient_worker_id: &str,
delivered_at: &str,
) -> Result<()>;
fn mark_ticket_notification_failed(
&self,
notification_id: &str,
recipient_runtime_id: &str,
recipient_worker_id: &str,
error: &str,
) -> Result<()>;
fn upsert_workdir_registry(&self, record: &WorkdirRegistryRecord) -> Result<()>;
fn get_workdir_registry(
&self,
@@ -1566,6 +1713,388 @@ impl ControlPlaneStore for SqliteWorkspaceStore {
})
}
fn get_current_ticket_worker_assignment(
&self,
workspace_id: &str,
ticket_id: &str,
) -> Result<Option<TicketWorkerAssignmentRecord>> {
self.with_conn(|conn| {
conn.query_row(
current_ticket_worker_assignment_select_sql().as_str(),
params![workspace_id, ticket_id],
read_ticket_worker_assignment_record,
)
.optional()
.map_err(Error::from)
})
}
fn set_current_ticket_worker_assignment(
&self,
record: &TicketWorkerAssignmentRecord,
expected_assignment_id: Option<&str>,
event_id: &str,
) -> Result<TicketWorkerAssignmentUpdate> {
self.with_conn(|conn| {
let tx = conn.unchecked_transaction()?;
let previous = tx
.query_row(
current_ticket_worker_assignment_select_sql().as_str(),
params![record.workspace_id, record.ticket_id],
read_ticket_worker_assignment_record,
)
.optional()?;
require_expected_ticket_assignment(
record.ticket_id.as_str(),
previous.as_ref(),
expected_assignment_id,
)?;
tx.execute(
r#"INSERT INTO ticket_worker_assignments (
workspace_id, ticket_id, assignment_id, runtime_id, worker_id,
assigned_by, assigned_at
) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)"#,
params![
record.workspace_id,
record.ticket_id,
record.assignment_id,
record.runtime_id,
record.worker_id,
record.assigned_by,
record.assigned_at,
],
)?;
tx.execute(
r#"INSERT INTO ticket_current_worker_assignments (
workspace_id, ticket_id, assignment_id, updated_at
) VALUES (?1, ?2, ?3, ?4)
ON CONFLICT(workspace_id, ticket_id) DO UPDATE SET
assignment_id = excluded.assignment_id,
updated_at = excluded.updated_at"#,
params![
record.workspace_id,
record.ticket_id,
record.assignment_id,
record.assigned_at,
],
)?;
tx.execute(
r#"INSERT INTO ticket_worker_assignment_events (
workspace_id, ticket_id, event_id, action, assignment_id,
previous_assignment_id, actor, created_at
) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)"#,
params![
record.workspace_id,
record.ticket_id,
event_id,
if previous.is_some() {
"reassigned"
} else {
"assigned"
},
record.assignment_id,
previous
.as_ref()
.map(|assignment| assignment.assignment_id.as_str()),
record.assigned_by,
record.assigned_at,
],
)?;
tx.commit()?;
Ok(TicketWorkerAssignmentUpdate {
current: record.clone(),
previous,
})
})
}
fn clear_current_ticket_worker_assignment(
&self,
workspace_id: &str,
ticket_id: &str,
expected_assignment_id: Option<&str>,
event_id: &str,
actor: &str,
created_at: &str,
) -> Result<Option<TicketWorkerAssignmentRecord>> {
self.with_conn(|conn| {
let tx = conn.unchecked_transaction()?;
let previous = tx
.query_row(
current_ticket_worker_assignment_select_sql().as_str(),
params![workspace_id, ticket_id],
read_ticket_worker_assignment_record,
)
.optional()?;
require_expected_ticket_assignment(ticket_id, previous.as_ref(), expected_assignment_id)?;
let Some(previous) = previous else {
tx.commit()?;
return Ok(None);
};
tx.execute(
"DELETE FROM ticket_current_worker_assignments WHERE workspace_id = ?1 AND ticket_id = ?2",
params![workspace_id, ticket_id],
)?;
tx.execute(
r#"INSERT INTO ticket_worker_assignment_events (
workspace_id, ticket_id, event_id, action, assignment_id,
previous_assignment_id, actor, created_at
) VALUES (?1, ?2, ?3, 'unassigned', NULL, ?4, ?5, ?6)"#,
params![
workspace_id,
ticket_id,
event_id,
previous.assignment_id,
actor,
created_at,
],
)?;
tx.commit()?;
Ok(Some(previous))
})
}
fn list_ticket_worker_assignment_events(
&self,
workspace_id: &str,
ticket_id: &str,
limit: usize,
) -> Result<Vec<TicketWorkerAssignmentEventRecord>> {
self.with_conn(|conn| {
let mut stmt = conn.prepare(
r#"SELECT workspace_id, ticket_id, event_id, action, assignment_id,
previous_assignment_id, actor, created_at
FROM ticket_worker_assignment_events
WHERE workspace_id = ?1 AND ticket_id = ?2
ORDER BY created_at DESC, event_id DESC
LIMIT ?3"#,
)?;
let rows = stmt.query_map(
params![workspace_id, ticket_id, limit as i64],
read_ticket_worker_assignment_event_record,
)?;
rows.collect::<std::result::Result<Vec<_>, _>>()
.map_err(Error::from)
})
}
fn upsert_worker_workspace_credential(
&self,
record: &WorkerWorkspaceCredentialRecord,
) -> Result<()> {
self.with_conn(|conn| {
conn.execute(
r#"INSERT INTO worker_workspace_credentials (
credential_id, token, workspace_id, runtime_id, worker_id, created_at
) VALUES (?1, ?2, ?3, ?4, ?5, ?6)
ON CONFLICT(credential_id) DO UPDATE SET
token = excluded.token,
workspace_id = excluded.workspace_id,
runtime_id = excluded.runtime_id,
worker_id = excluded.worker_id,
created_at = excluded.created_at"#,
params![
record.credential_id,
record.token,
record.workspace_id,
record.runtime_id,
record.worker_id,
record.created_at,
],
)?;
Ok(())
})
}
fn authenticate_worker_workspace_credential(
&self,
token: &str,
workspace_id: &str,
worker_id: &str,
) -> Result<Option<WorkerWorkspaceCredentialRecord>> {
self.with_conn(|conn| {
let tx = conn.unchecked_transaction()?;
let record = tx
.query_row(
r#"SELECT credential_id, token, workspace_id, runtime_id, worker_id, created_at
FROM worker_workspace_credentials
WHERE token = ?1 AND workspace_id = ?2"#,
params![token, workspace_id],
|row| {
Ok(WorkerWorkspaceCredentialRecord {
credential_id: row.get(0)?,
token: row.get(1)?,
workspace_id: row.get(2)?,
runtime_id: row.get(3)?,
worker_id: row.get(4)?,
created_at: row.get(5)?,
})
},
)
.optional()?;
let Some(mut record) = record else {
tx.commit()?;
return Ok(None);
};
if record.worker_id.as_deref().is_some_and(|bound| bound != worker_id) {
tx.commit()?;
return Ok(None);
}
if record.worker_id.is_none() {
tx.execute(
"UPDATE worker_workspace_credentials SET worker_id = ?1 WHERE credential_id = ?2 AND worker_id IS NULL",
params![worker_id, record.credential_id],
)?;
record.worker_id = Some(worker_id.to_string());
}
tx.commit()?;
Ok(Some(record))
})
}
fn enqueue_ticket_notification(
&self,
notification_id: &str,
workspace_id: &str,
ticket_id: &str,
event_sequence: i64,
source_runtime_id: &str,
source_worker_id: &str,
previous_state: &str,
current_state: &str,
created_at: &str,
recipients: &[TicketNotificationRecipient],
) -> Result<()> {
self.with_conn(|conn| {
let tx = conn.unchecked_transaction()?;
tx.execute(
r#"INSERT INTO ticket_notification_outbox (
notification_id, workspace_id, ticket_id, event_sequence,
source_runtime_id, source_worker_id, previous_state, current_state, created_at
) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)"#,
params![
notification_id,
workspace_id,
ticket_id,
event_sequence,
source_runtime_id,
source_worker_id,
previous_state,
current_state,
created_at,
],
)?;
for recipient in recipients {
tx.execute(
r#"INSERT OR IGNORE INTO ticket_notification_deliveries (
notification_id, recipient_runtime_id, recipient_worker_id,
recipient_kind, attempts
) VALUES (?1, ?2, ?3, ?4, 0)"#,
params![
notification_id,
recipient.runtime_id,
recipient.worker_id,
recipient.recipient_kind,
],
)?;
}
tx.commit()?;
Ok(())
})
}
fn list_pending_ticket_notification_deliveries(
&self,
workspace_id: &str,
limit: usize,
) -> Result<Vec<TicketNotificationDeliveryRecord>> {
self.with_conn(|conn| {
let mut stmt = conn.prepare(
r#"SELECT o.notification_id, o.workspace_id, o.ticket_id, o.event_sequence,
o.source_runtime_id, o.source_worker_id,
d.recipient_runtime_id, d.recipient_worker_id, d.recipient_kind, d.attempts
FROM ticket_notification_deliveries AS d
JOIN ticket_notification_outbox AS o ON o.notification_id = d.notification_id
WHERE o.workspace_id = ?1 AND d.delivered_at IS NULL
ORDER BY o.created_at ASC, o.notification_id ASC
LIMIT ?2"#,
)?;
let rows = stmt.query_map(params![workspace_id, limit as i64], |row| {
Ok(TicketNotificationDeliveryRecord {
notification_id: row.get(0)?,
workspace_id: row.get(1)?,
ticket_id: row.get(2)?,
event_sequence: row.get(3)?,
source_runtime_id: row.get(4)?,
source_worker_id: row.get(5)?,
recipient_runtime_id: row.get(6)?,
recipient_worker_id: row.get(7)?,
recipient_kind: row.get(8)?,
attempts: row.get(9)?,
})
})?;
rows.collect::<std::result::Result<Vec<_>, _>>()
.map_err(Error::from)
})
}
fn count_ticket_notification_deliveries_for_recipient(
&self,
workspace_id: &str,
ticket_id: &str,
runtime_id: &str,
worker_id: &str,
) -> Result<usize> {
self.with_conn(|conn| {
let count = conn.query_row(
r#"SELECT COUNT(*)
FROM ticket_notification_deliveries AS d
JOIN ticket_notification_outbox AS o ON o.notification_id = d.notification_id
WHERE o.workspace_id = ?1 AND o.ticket_id = ?2
AND d.recipient_runtime_id = ?3 AND d.recipient_worker_id = ?4"#,
params![workspace_id, ticket_id, runtime_id, worker_id],
|row| row.get::<_, i64>(0),
)?;
Ok(count as usize)
})
}
fn mark_ticket_notification_delivered(
&self,
notification_id: &str,
recipient_runtime_id: &str,
recipient_worker_id: &str,
delivered_at: &str,
) -> Result<()> {
self.with_conn(|conn| {
conn.execute(
r#"UPDATE ticket_notification_deliveries
SET delivered_at = ?4, last_error = NULL, attempts = attempts + 1
WHERE notification_id = ?1 AND recipient_runtime_id = ?2 AND recipient_worker_id = ?3"#,
params![notification_id, recipient_runtime_id, recipient_worker_id, delivered_at],
)?;
Ok(())
})
}
fn mark_ticket_notification_failed(
&self,
notification_id: &str,
recipient_runtime_id: &str,
recipient_worker_id: &str,
error: &str,
) -> Result<()> {
self.with_conn(|conn| {
conn.execute(
r#"UPDATE ticket_notification_deliveries
SET last_error = ?4, attempts = attempts + 1
WHERE notification_id = ?1 AND recipient_runtime_id = ?2 AND recipient_worker_id = ?3"#,
params![notification_id, recipient_runtime_id, recipient_worker_id, error],
)?;
Ok(())
})
}
fn upsert_workdir_registry(&self, record: &WorkdirRegistryRecord) -> Result<()> {
self.with_conn(|conn| {
conn.execute(
@@ -1996,6 +2525,63 @@ fn read_worker_registry_record(row: &rusqlite::Row<'_>) -> rusqlite::Result<Work
})
}
fn current_ticket_worker_assignment_select_sql() -> String {
"SELECT a.workspace_id, a.ticket_id, a.assignment_id, a.runtime_id, a.worker_id, \
a.assigned_by, a.assigned_at \
FROM ticket_current_worker_assignments AS current \
JOIN ticket_worker_assignments AS a \
ON a.workspace_id = current.workspace_id \
AND a.ticket_id = current.ticket_id \
AND a.assignment_id = current.assignment_id \
WHERE current.workspace_id = ?1 AND current.ticket_id = ?2"
.to_owned()
}
fn read_ticket_worker_assignment_record(
row: &rusqlite::Row<'_>,
) -> rusqlite::Result<TicketWorkerAssignmentRecord> {
Ok(TicketWorkerAssignmentRecord {
workspace_id: row.get(0)?,
ticket_id: row.get(1)?,
assignment_id: row.get(2)?,
runtime_id: row.get(3)?,
worker_id: row.get(4)?,
assigned_by: row.get(5)?,
assigned_at: row.get(6)?,
})
}
fn read_ticket_worker_assignment_event_record(
row: &rusqlite::Row<'_>,
) -> rusqlite::Result<TicketWorkerAssignmentEventRecord> {
Ok(TicketWorkerAssignmentEventRecord {
workspace_id: row.get(0)?,
ticket_id: row.get(1)?,
event_id: row.get(2)?,
action: row.get(3)?,
assignment_id: row.get(4)?,
previous_assignment_id: row.get(5)?,
actor: row.get(6)?,
created_at: row.get(7)?,
})
}
fn require_expected_ticket_assignment(
ticket_id: &str,
current: Option<&TicketWorkerAssignmentRecord>,
expected_assignment_id: Option<&str>,
) -> Result<()> {
let Some(expected_assignment_id) = expected_assignment_id else {
return Ok(());
};
if current.map(|assignment| assignment.assignment_id.as_str()) == Some(expected_assignment_id) {
return Ok(());
}
Err(Error::TicketAssignmentConflict(format!(
"Ticket {ticket_id} is no longer assigned to {expected_assignment_id}"
)))
}
fn workdir_registry_select_sql(where_clause: &str) -> String {
format!(
"SELECT workspace_id, workdir_id, runtime_id, repository_id, selector, resolved_commit, \
@@ -2175,6 +2761,95 @@ DROP TABLE IF EXISTS tickets;
Ok(())
}
fn create_ticket_worker_assignment_tables(conn: &Connection) -> Result<()> {
conn.execute_batch(
r#"
CREATE TABLE IF NOT EXISTS ticket_worker_assignments (
workspace_id TEXT NOT NULL REFERENCES workspaces(workspace_id) ON DELETE CASCADE,
ticket_id TEXT NOT NULL,
assignment_id TEXT NOT NULL,
runtime_id TEXT NOT NULL,
worker_id TEXT NOT NULL,
assigned_by TEXT NOT NULL,
assigned_at TEXT NOT NULL,
PRIMARY KEY (workspace_id, assignment_id),
UNIQUE (workspace_id, ticket_id, assignment_id)
);
CREATE TABLE IF NOT EXISTS ticket_current_worker_assignments (
workspace_id TEXT NOT NULL,
ticket_id TEXT NOT NULL,
assignment_id TEXT NOT NULL,
updated_at TEXT NOT NULL,
PRIMARY KEY (workspace_id, ticket_id),
FOREIGN KEY (workspace_id, ticket_id, assignment_id)
REFERENCES ticket_worker_assignments(workspace_id, ticket_id, assignment_id)
ON DELETE CASCADE
);
CREATE TABLE IF NOT EXISTS ticket_worker_assignment_events (
workspace_id TEXT NOT NULL REFERENCES workspaces(workspace_id) ON DELETE CASCADE,
ticket_id TEXT NOT NULL,
event_id TEXT NOT NULL,
action TEXT NOT NULL CHECK (action IN ('assigned', 'reassigned', 'unassigned')),
assignment_id TEXT,
previous_assignment_id TEXT,
actor TEXT NOT NULL,
created_at TEXT NOT NULL,
PRIMARY KEY (workspace_id, event_id)
);
CREATE INDEX IF NOT EXISTS idx_ticket_assignments_worker
ON ticket_worker_assignments(workspace_id, runtime_id, worker_id, assigned_at DESC);
CREATE INDEX IF NOT EXISTS idx_ticket_assignment_events_ticket
ON ticket_worker_assignment_events(workspace_id, ticket_id, created_at DESC);
"#,
)?;
Ok(())
}
fn create_ticket_notification_tables(conn: &Connection) -> Result<()> {
conn.execute_batch(
r#"
CREATE TABLE IF NOT EXISTS worker_workspace_credentials (
credential_id TEXT PRIMARY KEY,
token TEXT NOT NULL UNIQUE,
workspace_id TEXT NOT NULL REFERENCES workspaces(workspace_id) ON DELETE CASCADE,
runtime_id TEXT NOT NULL,
worker_id TEXT,
created_at TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS ticket_notification_outbox (
notification_id TEXT PRIMARY KEY,
workspace_id TEXT NOT NULL REFERENCES workspaces(workspace_id) ON DELETE CASCADE,
ticket_id TEXT NOT NULL,
event_sequence INTEGER NOT NULL,
source_runtime_id TEXT NOT NULL,
source_worker_id TEXT NOT NULL,
previous_state TEXT NOT NULL,
current_state TEXT NOT NULL,
created_at TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS ticket_notification_deliveries (
notification_id TEXT NOT NULL REFERENCES ticket_notification_outbox(notification_id) ON DELETE CASCADE,
recipient_runtime_id TEXT NOT NULL,
recipient_worker_id TEXT NOT NULL,
recipient_kind TEXT NOT NULL CHECK (recipient_kind IN ('assigned', 'orchestrator')),
attempts INTEGER NOT NULL DEFAULT 0,
delivered_at TEXT,
last_error TEXT,
PRIMARY KEY (notification_id, recipient_runtime_id, recipient_worker_id)
);
CREATE INDEX IF NOT EXISTS idx_ticket_notification_pending
ON ticket_notification_deliveries(delivered_at, attempts);
"#,
)?;
Ok(())
}
fn create_objective_event_tables(conn: &Connection) -> Result<()> {
conn.execute_batch(
r#"
@@ -2857,7 +3532,7 @@ CREATE TABLE ticket_worker_links (ticket_id TEXT, worker_ref_key TEXT);
let db = dir.path().join("control-plane.sqlite");
let store = SqliteWorkspaceStore::open(&db).unwrap();
assert_eq!(store.schema_version().await.unwrap(), 14);
assert_eq!(store.schema_version().await.unwrap(), 16);
let record = WorkspaceRecord {
workspace_id: "local-dev".to_string(),
@@ -2870,13 +3545,191 @@ CREATE TABLE ticket_worker_links (ticket_id TEXT, worker_ref_key TEXT);
store.upsert_workspace(&record).await.unwrap();
let reopened = SqliteWorkspaceStore::open(&db).unwrap();
assert_eq!(reopened.schema_version().await.unwrap(), 14);
assert_eq!(reopened.schema_version().await.unwrap(), 16);
assert_eq!(
reopened.get_workspace("local-dev").await.unwrap(),
Some(record)
);
}
#[tokio::test]
async fn ticket_worker_assignment_replaces_current_and_preserves_audit_history() {
let dir = tempfile::tempdir().unwrap();
let store = SqliteWorkspaceStore::open(dir.path().join("server.db")).unwrap();
store
.upsert_workspace(&WorkspaceRecord {
workspace_id: "workspace-a".to_string(),
owner_account_id: None,
display_name: "Workspace A".to_string(),
state: "active".to_string(),
created_at: "2026-07-31T00:00:00Z".to_string(),
updated_at: "2026-07-31T00:00:00Z".to_string(),
})
.await
.unwrap();
let first = TicketWorkerAssignmentRecord {
workspace_id: "workspace-a".to_string(),
ticket_id: "ticket-1".to_string(),
assignment_id: "assignment-1".to_string(),
runtime_id: "runtime-1".to_string(),
worker_id: "worker-1".to_string(),
assigned_by: "user-1".to_string(),
assigned_at: "2026-07-31T00:00:01Z".to_string(),
};
let created = store
.set_current_ticket_worker_assignment(&first, None, "event-1")
.unwrap();
assert_eq!(created.current, first);
assert_eq!(created.previous, None);
let second = TicketWorkerAssignmentRecord {
assignment_id: "assignment-2".to_string(),
runtime_id: "runtime-2".to_string(),
worker_id: "worker-2".to_string(),
assigned_by: "user-2".to_string(),
assigned_at: "2026-07-31T00:00:02Z".to_string(),
..first.clone()
};
let replaced = store
.set_current_ticket_worker_assignment(&second, Some("assignment-1"), "event-2")
.unwrap();
assert_eq!(replaced.current, second);
assert_eq!(replaced.previous, Some(first.clone()));
assert_eq!(
store
.get_current_ticket_worker_assignment("workspace-a", "ticket-1")
.unwrap(),
Some(second.clone())
);
let stale = store
.clear_current_ticket_worker_assignment(
"workspace-a",
"ticket-1",
Some("assignment-1"),
"event-stale",
"user-1",
"2026-07-31T00:00:03Z",
)
.unwrap_err();
assert!(matches!(stale, Error::TicketAssignmentConflict(_)));
let cleared = store
.clear_current_ticket_worker_assignment(
"workspace-a",
"ticket-1",
Some("assignment-2"),
"event-3",
"user-2",
"2026-07-31T00:00:03Z",
)
.unwrap();
assert_eq!(cleared, Some(second));
assert_eq!(
store
.get_current_ticket_worker_assignment("workspace-a", "ticket-1")
.unwrap(),
None
);
let events = store
.list_ticket_worker_assignment_events("workspace-a", "ticket-1", 10)
.unwrap();
assert_eq!(
events
.iter()
.map(|event| event.action.as_str())
.collect::<Vec<_>>(),
vec!["unassigned", "reassigned", "assigned"]
);
assert_eq!(events[1].assignment_id.as_deref(), Some("assignment-2"));
assert_eq!(
events[1].previous_assignment_id.as_deref(),
Some("assignment-1")
);
}
#[tokio::test]
async fn worker_credential_binds_once_and_notification_outbox_is_durable() {
let dir = tempfile::tempdir().unwrap();
let store = SqliteWorkspaceStore::open(dir.path().join("server.db")).unwrap();
store
.upsert_workspace(&WorkspaceRecord {
workspace_id: "workspace-a".to_string(),
owner_account_id: None,
display_name: "Workspace A".to_string(),
state: "active".to_string(),
created_at: "2026-07-31T00:00:00Z".to_string(),
updated_at: "2026-07-31T00:00:00Z".to_string(),
})
.await
.unwrap();
store
.upsert_worker_workspace_credential(&WorkerWorkspaceCredentialRecord {
credential_id: "credential-1".to_string(),
token: "secret-token".to_string(),
workspace_id: "workspace-a".to_string(),
runtime_id: "runtime-1".to_string(),
worker_id: None,
created_at: "2026-07-31T00:00:01Z".to_string(),
})
.unwrap();
let bound = store
.authenticate_worker_workspace_credential("secret-token", "workspace-a", "worker-1")
.unwrap()
.unwrap();
assert_eq!(bound.worker_id.as_deref(), Some("worker-1"));
assert!(
store
.authenticate_worker_workspace_credential(
"secret-token",
"workspace-a",
"worker-2",
)
.unwrap()
.is_none()
);
store
.enqueue_ticket_notification(
"notification-1",
"workspace-a",
"ticket-1",
4,
"runtime-1",
"worker-1",
"queued",
"inprogress",
"2026-07-31T00:00:02Z",
&[TicketNotificationRecipient {
runtime_id: "runtime-1".to_string(),
worker_id: "worker-2".to_string(),
recipient_kind: "assigned".to_string(),
}],
)
.unwrap();
let pending = store
.list_pending_ticket_notification_deliveries("workspace-a", 10)
.unwrap();
assert_eq!(pending.len(), 1);
assert_eq!(pending[0].event_sequence, 4);
store
.mark_ticket_notification_delivered(
"notification-1",
"runtime-1",
"worker-2",
"2026-07-31T00:00:03Z",
)
.unwrap();
assert!(
store
.list_pending_ticket_notification_deliveries("workspace-a", 10)
.unwrap()
.is_empty()
);
}
#[test]
fn fresh_schema_matches_workspace_db_v0_boundaries() {
let conn = Connection::open_in_memory().unwrap();
@@ -2896,6 +3749,9 @@ CREATE TABLE ticket_worker_links (ticket_id TEXT, worker_ref_key TEXT);
"artifacts",
"audit_events",
"worker_registry",
"ticket_worker_assignments",
"ticket_current_worker_assignments",
"ticket_worker_assignment_events",
"workdir_registry",
"worker_workdir_links",
"accounts",
@@ -3058,7 +3914,7 @@ CREATE TABLE ticket_worker_links (ticket_id TEXT, worker_ref_key TEXT);
.unwrap();
let store = SqliteWorkspaceStore::from_connection(conn).unwrap();
assert_eq!(store.schema_version().await.unwrap(), 14);
assert_eq!(store.schema_version().await.unwrap(), 16);
store
.with_conn(|conn| {
@@ -3161,7 +4017,7 @@ CREATE TABLE ticket_worker_links (ticket_id TEXT, worker_ref_key TEXT);
#[tokio::test]
async fn repository_records_round_trip() {
let store = SqliteWorkspaceStore::in_memory().unwrap();
assert_eq!(store.schema_version().await.unwrap(), 14);
assert_eq!(store.schema_version().await.unwrap(), 16);
let workspace = WorkspaceRecord {
workspace_id: "local-dev".to_string(),
owner_account_id: None,
@@ -3199,7 +4055,7 @@ CREATE TABLE ticket_worker_links (ticket_id TEXT, worker_ref_key TEXT);
#[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(), 14);
assert_eq!(store.schema_version().await.unwrap(), 16);
let workspace = WorkspaceRecord {
workspace_id: "local-dev".to_string(),
owner_account_id: None,
@@ -3373,7 +4229,7 @@ CREATE TABLE ticket_worker_links (ticket_id TEXT, worker_ref_key TEXT);
#[tokio::test]
async fn account_and_login_records_round_trip() {
let store = SqliteWorkspaceStore::in_memory().unwrap();
assert_eq!(store.schema_version().await.unwrap(), 14);
assert_eq!(store.schema_version().await.unwrap(), 16);
let now = "2026-07-22T00:00:00Z".to_string();
let account = AccountRecord {
account_id: "acct-user-alice".to_string(),