Skip to main content

aven_core/operations/
dependencies.rs

1use crate::ids::WorkspaceId;
2use std::collections::HashSet;
3
4use anyhow::{Result, bail};
5use sqlx::SqliteConnection;
6
7use crate::change_log::{ChangeEntity, ChangePayload, append_change, op_type};
8use crate::db::{Database, begin_immediate};
9use crate::ids::{TaskId, now};
10use crate::refs::get_task_in_workspace;
11use crate::workspaces::Workspace;
12
13impl Database {
14    pub async fn add_task_dependency(
15        &self,
16        workspace: &Workspace,
17        task_id: &TaskId,
18        depends_on_id: &TaskId,
19    ) -> Result<DependencyOutcome> {
20        let mut conn = self.acquire().await?;
21        add_task_dependency(&mut conn, workspace, task_id, depends_on_id).await
22    }
23
24    pub async fn remove_task_dependency(
25        &self,
26        workspace: &Workspace,
27        task_id: &TaskId,
28        depends_on_id: &TaskId,
29    ) -> Result<DependencyOutcome> {
30        let mut conn = self.acquire().await?;
31        remove_task_dependency(&mut conn, workspace, task_id, depends_on_id).await
32    }
33}
34
35pub struct DependencyOutcome {
36    pub task: crate::types::Task,
37    pub depends_on: crate::types::Task,
38    pub changed: bool,
39}
40
41struct DependencyPair {
42    task: crate::types::Task,
43    depends_on: crate::types::Task,
44}
45
46async fn load_dependency_pair(
47    conn: &mut SqliteConnection,
48    workspace: &Workspace,
49    task_id: &crate::ids::TaskId,
50    depends_on_id: &crate::ids::TaskId,
51) -> Result<DependencyPair> {
52    if task_id == depends_on_id {
53        bail!("error dependency-self task_id={task_id}");
54    }
55
56    let task = get_task_in_workspace(conn, workspace, task_id).await?;
57    let depends_on = get_task_in_workspace(conn, workspace, depends_on_id).await?;
58
59    Ok(DependencyPair { task, depends_on })
60}
61
62async fn record_dependency_change(
63    conn: &mut SqliteConnection,
64    workspace: &Workspace,
65    pair: &DependencyPair,
66    op_type: &'static str,
67) -> Result<()> {
68    append_change(
69        conn,
70        ChangeEntity::Task,
71        &pair.task.id,
72        Some("dependencies"),
73        op_type,
74        ChangePayload::workspace(workspace).set("depends_on_task_id", pair.depends_on.id.clone()),
75    )
76    .await?;
77    Ok(())
78}
79
80pub async fn add_task_dependency(
81    conn: &mut SqliteConnection,
82    workspace: &Workspace,
83    task_id: &crate::ids::TaskId,
84    depends_on_id: &crate::ids::TaskId,
85) -> Result<DependencyOutcome> {
86    let mut tx = begin_immediate(conn).await?;
87    let pair = load_dependency_pair(&mut tx, workspace, task_id, depends_on_id).await?;
88
89    if dependency_path_exists(
90        &mut tx,
91        &pair.task.workspace_id,
92        &pair.depends_on.id,
93        &pair.task.id,
94    )
95    .await?
96    {
97        bail!("error dependency-cycle task_id={task_id} depends_on_task_id={depends_on_id}");
98    }
99
100    let created_at = now();
101    let changed = sqlx::query(
102        "INSERT OR IGNORE INTO task_dependencies(workspace_id, task_id, depends_on_task_id, created_at)
103         VALUES (?, ?, ?, ?)",
104    )
105    .bind(&pair.task.workspace_id)
106    .bind(&pair.task.id)
107    .bind(&pair.depends_on.id)
108    .bind(&created_at)
109    .execute(&mut *tx)
110    .await?
111    .rows_affected()
112        > 0;
113
114    if changed {
115        record_dependency_change(&mut tx, workspace, &pair, op_type::DEPENDENCY_ADD).await?;
116    }
117
118    tx.commit().await?;
119    Ok(DependencyOutcome {
120        task: pair.task,
121        depends_on: pair.depends_on,
122        changed,
123    })
124}
125
126pub async fn remove_task_dependency(
127    conn: &mut SqliteConnection,
128    workspace: &Workspace,
129    task_id: &crate::ids::TaskId,
130    depends_on_id: &crate::ids::TaskId,
131) -> Result<DependencyOutcome> {
132    let mut tx = begin_immediate(conn).await?;
133    let pair = load_dependency_pair(&mut tx, workspace, task_id, depends_on_id).await?;
134
135    let changed = sqlx::query(
136        "DELETE FROM task_dependencies
137         WHERE workspace_id = ? AND task_id = ? AND depends_on_task_id = ?",
138    )
139    .bind(&pair.task.workspace_id)
140    .bind(&pair.task.id)
141    .bind(&pair.depends_on.id)
142    .execute(&mut *tx)
143    .await?
144    .rows_affected()
145        > 0;
146
147    if changed {
148        record_dependency_change(&mut tx, workspace, &pair, op_type::DEPENDENCY_REMOVE).await?;
149    }
150
151    tx.commit().await?;
152    Ok(DependencyOutcome {
153        task: pair.task,
154        depends_on: pair.depends_on,
155        changed,
156    })
157}
158
159pub async fn dependency_path_exists(
160    conn: &mut SqliteConnection,
161    workspace_id: &WorkspaceId,
162    from_task_id: &crate::ids::TaskId,
163    to_task_id: &crate::ids::TaskId,
164) -> Result<bool> {
165    let mut visited = HashSet::new();
166    let mut stack = vec![from_task_id.clone()];
167    while let Some(current) = stack.pop() {
168        if !visited.insert(current.clone()) {
169            continue;
170        }
171        if &current == to_task_id {
172            return Ok(true);
173        }
174        let next = sqlx::query_scalar::<_, crate::ids::TaskId>(
175            "SELECT depends_on_task_id
176             FROM task_dependencies
177             WHERE workspace_id = ? AND task_id = ?",
178        )
179        .bind(workspace_id)
180        .bind(&current)
181        .fetch_all(&mut *conn)
182        .await?;
183        stack.extend(next);
184    }
185    Ok(false)
186}