Skip to main content

aven_core/
mutation.rs

1use crate::ids::{ProjectId, WorkspaceId};
2use anyhow::{Result, bail, ensure};
3use sqlx::SqliteConnection;
4use tracing::{debug, info};
5
6use crate::change_log::op_type;
7use crate::choices::TaskPriority;
8use crate::db::{
9    Database, conflict_exists, field_version, insert_change, set_field_version, task_from_row,
10};
11use crate::ids::now;
12use crate::projects::resolve_or_create_project_in_workspace;
13use crate::refs::get_task_in_workspace;
14use crate::task_fields::TaskField;
15use crate::types::{Project, Task};
16use crate::workspaces::Workspace;
17
18#[derive(Debug)]
19pub(crate) struct OpenConflictError {
20    task_id: crate::ids::TaskId,
21    field: &'static str,
22}
23
24impl std::fmt::Display for OpenConflictError {
25    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
26        write!(
27            formatter,
28            "error conflicted-field ref={} field={} hint=\"use conflict resolve\"",
29            self.task_id, self.field
30        )
31    }
32}
33
34impl std::error::Error for OpenConflictError {}
35
36impl Database {
37    pub async fn set_task_status(
38        &self,
39        workspace: &Workspace,
40        task: &Task,
41        status: &str,
42    ) -> Result<Task> {
43        let mut conn = self.acquire().await?;
44        set_status(&mut conn, workspace, task, status).await
45    }
46
47    pub async fn set_task_priority(
48        &self,
49        workspace: &Workspace,
50        task: &Task,
51        priority: &str,
52    ) -> Result<Task> {
53        let mut conn = self.acquire().await?;
54        set_priority(&mut conn, workspace, task, priority).await
55    }
56
57    pub async fn cycle_task_priority(
58        &self,
59        workspace: &Workspace,
60        task: &Task,
61        reverse: bool,
62    ) -> Result<Task> {
63        let mut conn = self.acquire().await?;
64        cycle_priority(&mut conn, workspace, task, reverse).await
65    }
66
67    pub async fn set_task_deleted_state(
68        &self,
69        workspace: &Workspace,
70        task: &Task,
71        deleted: bool,
72    ) -> Result<Task> {
73        let mut conn = self.acquire().await?;
74        set_deleted(&mut conn, workspace, task, deleted).await
75    }
76
77    pub async fn set_task_field(
78        &self,
79        workspace: &Workspace,
80        task_id: &crate::ids::TaskId,
81        field: &str,
82        value: &str,
83    ) -> Result<bool> {
84        let mut conn = self.acquire().await?;
85        set_task_field(&mut conn, workspace, task_id, field, value).await
86    }
87
88    pub async fn set_task_project(
89        &self,
90        workspace: &Workspace,
91        task_id: &crate::ids::TaskId,
92        project: &Project,
93    ) -> Result<bool> {
94        let mut conn = self.acquire().await?;
95        set_task_project(&mut conn, workspace, task_id, project).await
96    }
97
98    pub async fn set_task_fields(
99        &self,
100        workspace: &Workspace,
101        updates: &[(crate::ids::TaskId, String, String)],
102    ) -> Result<Vec<bool>> {
103        let mut conn = self.acquire().await?;
104        let mut outcomes = Vec::with_capacity(updates.len());
105        for (task_id, field, value) in updates {
106            outcomes.push(set_task_field(&mut conn, workspace, task_id, field, value).await?);
107        }
108        Ok(outcomes)
109    }
110
111    pub async fn cycle_task_priorities(
112        &self,
113        workspace: &Workspace,
114        tasks: &[Task],
115        reverse: bool,
116    ) -> Result<Vec<Task>> {
117        let mut conn = self.acquire().await?;
118        let mut outcomes = Vec::with_capacity(tasks.len());
119        for task in tasks {
120            outcomes.push(cycle_priority(&mut conn, workspace, task, reverse).await?);
121        }
122        Ok(outcomes)
123    }
124}
125
126pub async fn set_status(
127    conn: &mut SqliteConnection,
128    workspace: &Workspace,
129    task: &Task,
130    status: &str,
131) -> Result<Task> {
132    set_task_field(conn, workspace, &task.id, "status", status).await?;
133    get_task_in_workspace(conn, workspace, &task.id).await
134}
135
136pub async fn set_priority(
137    conn: &mut SqliteConnection,
138    workspace: &Workspace,
139    task: &Task,
140    priority: &str,
141) -> Result<Task> {
142    set_task_field(conn, workspace, &task.id, "priority", priority).await?;
143    get_task_in_workspace(conn, workspace, &task.id).await
144}
145
146pub async fn cycle_priority(
147    conn: &mut SqliteConnection,
148    workspace: &Workspace,
149    task: &Task,
150    reverse: bool,
151) -> Result<Task> {
152    let index = TaskPriority::ALL
153        .iter()
154        .position(|priority| *priority == task.priority)
155        .unwrap_or(0);
156    let next = if reverse {
157        (index + TaskPriority::ALL.len() - 1) % TaskPriority::ALL.len()
158    } else {
159        (index + 1) % TaskPriority::ALL.len()
160    };
161    set_priority(conn, workspace, task, TaskPriority::ALL[next].as_str()).await
162}
163
164pub async fn set_deleted(
165    conn: &mut SqliteConnection,
166    workspace: &Workspace,
167    task: &Task,
168    deleted: bool,
169) -> Result<Task> {
170    set_task_field(
171        conn,
172        workspace,
173        &task.id,
174        "deleted",
175        if deleted { "1" } else { "0" },
176    )
177    .await?;
178    get_task_in_workspace(conn, workspace, &task.id).await
179}
180
181pub async fn set_task_field(
182    conn: &mut SqliteConnection,
183    workspace: &Workspace,
184    task_id: &crate::ids::TaskId,
185    field: &str,
186    value: &str,
187) -> Result<bool> {
188    let task_field = TaskField::parse_or_unknown(field)?;
189    if task_field.is_project() {
190        let project = resolve_or_create_project_in_workspace(conn, &workspace.id, value).await?;
191        set_task_project(conn, workspace, task_id, &project).await
192    } else {
193        set_task_scalar_field(conn, workspace, task_id, task_field, value).await
194    }
195}
196
197pub async fn set_task_project(
198    conn: &mut SqliteConnection,
199    workspace: &Workspace,
200    task_id: &crate::ids::TaskId,
201    project: &Project,
202) -> Result<bool> {
203    let field = TaskField::Project.as_str();
204    let current = current_task(conn, &workspace.id, task_id).await?;
205    if current.project_id == project.id {
206        return Ok(false);
207    }
208    if conflict_exists(conn, &workspace.id, task_id, field).await? {
209        return Err(anyhow::Error::new(OpenConflictError {
210            task_id: task_id.clone(),
211            field,
212        }));
213    }
214    debug!(task_id = %task_id, field = %field, "task field mutation started");
215    let base = field_version(conn, task_id, field).await?;
216    apply_project_id_in_workspace(conn, &workspace.id, task_id, &project.id).await?;
217    let payload = TaskField::project_payload(&workspace.id, &workspace.key, project);
218    finish_task_field_change(conn, task_id, field, payload, base.as_deref()).await?;
219    Ok(true)
220}
221
222async fn set_task_scalar_field(
223    conn: &mut SqliteConnection,
224    workspace: &Workspace,
225    task_id: &crate::ids::TaskId,
226    task_field: TaskField,
227    value: &str,
228) -> Result<bool> {
229    task_field.validate_value(value)?;
230
231    let field = task_field.as_str();
232    let current = current_task(conn, &workspace.id, task_id).await?;
233    if task_field.current_value(&current) == value {
234        return Ok(false);
235    }
236    if conflict_exists(conn, &workspace.id, task_id, field).await? {
237        return Err(anyhow::Error::new(OpenConflictError {
238            task_id: task_id.clone(),
239            field,
240        }));
241    }
242    debug!(task_id = %task_id, field = %field, "task field mutation started");
243    let base = field_version(conn, task_id, field).await?;
244    apply_scalar_field_value_in_workspace(conn, &workspace.id, task_id, task_field, value).await?;
245    let payload = task_field.scalar_payload(&workspace.id, &workspace.key, value)?;
246    finish_task_field_change(conn, task_id, field, payload, base.as_deref()).await?;
247    Ok(true)
248}
249
250async fn finish_task_field_change(
251    conn: &mut SqliteConnection,
252    task_id: &crate::ids::TaskId,
253    field: &str,
254    payload: serde_json::Value,
255    base: Option<&str>,
256) -> Result<()> {
257    let change_id = insert_change(
258        conn,
259        "task",
260        task_id,
261        Some(field),
262        op_type::SET_FIELD,
263        payload,
264        base,
265    )
266    .await?;
267    set_field_version(conn, task_id, field, &change_id).await?;
268    info!(
269        task_id = %task_id,
270        field = %field,
271        change_id = %change_id,
272        "task field mutated"
273    );
274    Ok(())
275}
276
277async fn current_task(
278    conn: &mut SqliteConnection,
279    workspace_id: &WorkspaceId,
280    task_id: &crate::ids::TaskId,
281) -> Result<Task> {
282    let row = sqlx::query(
283        "SELECT t.id, t.workspace_id, t.title, t.description, t.project_id,
284                p.key AS project_key, p.prefix AS project_prefix, t.status,
285                t.priority, t.created_at, t.updated_at, t.queue_activity_at,
286                t.available_at, t.due_on, t.deleted, t.is_epic
287         FROM tasks t
288         JOIN projects p ON p.workspace_id = t.workspace_id AND p.id = t.project_id
289         WHERE t.workspace_id = ? AND t.id = ?",
290    )
291    .bind(workspace_id)
292    .bind(task_id)
293    .fetch_optional(&mut *conn)
294    .await?
295    .ok_or_else(|| {
296        anyhow::anyhow!(
297            "error task-not-found task_id={} workspace_id={}",
298            task_id,
299            workspace_id
300        )
301    })?;
302    task_from_row(&row)
303}
304
305#[allow(dead_code)]
306pub async fn apply_field_value(
307    conn: &mut SqliteConnection,
308    workspace_id: &WorkspaceId,
309    task_id: &crate::ids::TaskId,
310    field: &str,
311    value: &str,
312) -> Result<()> {
313    apply_field_value_in_workspace(conn, workspace_id, task_id, field, value).await
314}
315
316pub async fn apply_project_id_in_workspace(
317    conn: &mut SqliteConnection,
318    workspace_id: &WorkspaceId,
319    task_id: &crate::ids::TaskId,
320    project_id: &ProjectId,
321) -> Result<()> {
322    let project_exists = sqlx::query_scalar::<_, i64>(
323        "SELECT count(*) FROM projects WHERE workspace_id = ? AND id = ? AND deleted = 0",
324    )
325    .bind(workspace_id)
326    .bind(project_id)
327    .fetch_one(&mut *conn)
328    .await?
329        > 0;
330    if !project_exists {
331        bail!("error unknown-project-id id={project_id}");
332    }
333    let ts = now();
334    let rows_affected = sqlx::query(
335        "UPDATE tasks SET project_id = ?, updated_at = ? WHERE workspace_id = ? AND id = ?",
336    )
337    .bind(project_id)
338    .bind(&ts)
339    .bind(workspace_id)
340    .bind(task_id)
341    .execute(&mut *conn)
342    .await?
343    .rows_affected();
344    ensure!(
345        rows_affected == 1,
346        "error task-not-found task_id={} workspace_id={}",
347        task_id,
348        workspace_id
349    );
350    Ok(())
351}
352
353pub async fn apply_field_value_in_workspace(
354    conn: &mut SqliteConnection,
355    workspace_id: &WorkspaceId,
356    task_id: &crate::ids::TaskId,
357    field: &str,
358    value: &str,
359) -> Result<()> {
360    let task_field = TaskField::parse_or_unknown(field)?;
361    apply_scalar_field_value_in_workspace(conn, workspace_id, task_id, task_field, value).await
362}
363
364async fn apply_scalar_field_value_in_workspace(
365    conn: &mut SqliteConnection,
366    workspace_id: &WorkspaceId,
367    task_id: &crate::ids::TaskId,
368    task_field: TaskField,
369    value: &str,
370) -> Result<()> {
371    task_field.validate_value(value)?;
372
373    let ts = now();
374    let activity_at = if task_field.updates_queue_activity() {
375        ts.as_str()
376    } else {
377        ""
378    };
379    let deleted_value = i64::from(value == "1");
380    let epic_value = i64::from(value == "1");
381    let rows_affected = match task_field {
382        TaskField::Title => sqlx::query(
383            "UPDATE tasks SET title = ?, updated_at = ?, queue_activity_at = COALESCE(NULLIF(?, ''), queue_activity_at) WHERE workspace_id = ? AND id = ?",
384        )
385        .bind(value)
386        .bind(&ts)
387        .bind(activity_at)
388        .bind(workspace_id)
389        .bind(task_id)
390        .execute(&mut *conn)
391        .await?
392        .rows_affected(),
393        TaskField::Description => sqlx::query(
394            "UPDATE tasks SET description = ?, updated_at = ?, queue_activity_at = COALESCE(NULLIF(?, ''), queue_activity_at) WHERE workspace_id = ? AND id = ?",
395        )
396        .bind(value)
397        .bind(&ts)
398        .bind(activity_at)
399        .bind(workspace_id)
400        .bind(task_id)
401        .execute(&mut *conn)
402        .await?
403        .rows_affected(),
404        TaskField::Project => bail!("error project-update-requires-project-id"),
405        TaskField::Status => sqlx::query(
406            "UPDATE tasks SET status = ?, updated_at = ?, queue_activity_at = COALESCE(NULLIF(?, ''), queue_activity_at) WHERE workspace_id = ? AND id = ?",
407        )
408        .bind(value)
409        .bind(&ts)
410        .bind(activity_at)
411        .bind(workspace_id)
412        .bind(task_id)
413        .execute(&mut *conn)
414        .await?
415        .rows_affected(),
416        TaskField::Priority => sqlx::query(
417            "UPDATE tasks SET priority = ?, updated_at = ?, queue_activity_at = COALESCE(NULLIF(?, ''), queue_activity_at) WHERE workspace_id = ? AND id = ?",
418        )
419        .bind(value)
420        .bind(&ts)
421        .bind(activity_at)
422        .bind(workspace_id)
423        .bind(task_id)
424        .execute(&mut *conn)
425        .await?
426        .rows_affected(),
427        TaskField::AvailableAt => sqlx::query(
428            "UPDATE tasks SET available_at = ?, updated_at = ?, queue_activity_at = COALESCE(NULLIF(?, ''), queue_activity_at) WHERE workspace_id = ? AND id = ?",
429        )
430        .bind(value)
431        .bind(&ts)
432        .bind(activity_at)
433        .bind(workspace_id)
434        .bind(task_id)
435        .execute(&mut *conn)
436        .await?
437        .rows_affected(),
438        TaskField::DueOn => sqlx::query(
439            "UPDATE tasks SET due_on = ?, updated_at = ? WHERE workspace_id = ? AND id = ?",
440        )
441        .bind(value)
442        .bind(&ts)
443        .bind(workspace_id)
444        .bind(task_id)
445        .execute(&mut *conn)
446        .await?
447        .rows_affected(),
448        TaskField::Deleted => sqlx::query(
449            "UPDATE tasks SET deleted = ?, updated_at = ?, queue_activity_at = COALESCE(NULLIF(?, ''), queue_activity_at) WHERE workspace_id = ? AND id = ?",
450        )
451        .bind(deleted_value)
452        .bind(&ts)
453        .bind(activity_at)
454        .bind(workspace_id)
455        .bind(task_id)
456        .execute(&mut *conn)
457        .await?
458        .rows_affected(),
459        TaskField::IsEpic => sqlx::query(
460            "UPDATE tasks SET is_epic = ?, updated_at = ? WHERE workspace_id = ? AND id = ?",
461        )
462        .bind(epic_value)
463        .bind(&ts)
464        .bind(workspace_id)
465        .bind(task_id)
466        .execute(&mut *conn)
467        .await?
468        .rows_affected(),
469    };
470    ensure!(
471        rows_affected == 1,
472        "error task-not-found task_id={} workspace_id={}",
473        task_id,
474        workspace_id
475    );
476    Ok(())
477}