g1t

syntaqx/g1t

public

Git for AI scale: a forge for thousands of agents working on the same code at once.

g1t/services/work/src/lib.rs

1,293 lines47,178 bytes
//! The work service: issues, pull requests, comments and sessions.
//!
//! Other services reach it over `POST /rpc/<method>`; see
//! `g1t_contracts::work` for the methods and their arguments. It also
//! consumes its queue of events from the bus.

mod checks;
mod rows;

use g1t_contracts::events::{
    CommentCreated, Event, IssueEvent, NewEvent, Publish, PullEvent, SessionAppended,
};
use g1t_contracts::repos::{ForkArgs, GetArgs, HeadArgs, LandArgs, Landed, Repo, RepoPath};
use g1t_contracts::time::rfc3339;
use g1t_contracts::work::*;
use g1t_contracts::{FailureCode, Outcome, User, Viewer, new_id};
use g1t_kit::{args, now_ms, reply, rpc_method};
use serde::Serialize;
use worker::wasm_bindgen::JsValue;
use worker::{
    Context, D1Database, Env, Fetcher, MessageBatch, MessageExt, Request, Response, Result, event,
};

use rows::{CommentRow, IssueRow, MovedRow, NumberRow, PullRow, SessionRow, ValueRow};

const SOURCE: &str = "work";
const MAX_ENTRY_BATCH: usize = 200;
const MAX_ENTRY_CHARS: usize = 64_000;
const MAX_TITLE_CHARS: usize = 200;
const SESSION_PAGE: u32 = 500;
const LIST_PAGE: u32 = 100;
const UNVERIFIED: &str = "Confirm your email address first. Check your inbox, or resend the link from the banner on g1t.sh.";

const ISSUE_COLUMNS: &str = "issues.*,
  (SELECT count(*) FROM pulls WHERE pulls.issue_id = issues.id) AS pull_count,
  (SELECT count(*) FROM comments
   WHERE comments.repo_id = issues.repo_id AND comments.number = issues.number) AS comment_count";

fn no_issue<T>() -> Outcome<T> {
    Outcome::fail(FailureCode::NotFound, "Issue not found.")
}

fn no_pull<T>() -> Outcome<T> {
    Outcome::fail(FailureCode::NotFound, "Pull request not found.")
}

fn optional(value: &Option<String>) -> JsValue {
    value.as_deref().map_or(JsValue::NULL, JsValue::from)
}

fn optional_number(value: Option<u32>) -> JsValue {
    value.map_or(JsValue::NULL, JsValue::from)
}

/// The lowercase name a `State` is stored and sent as.
fn state_name(state: Option<State>) -> Option<&'static str> {
    state.map(|state| match state {
        State::Open => "open",
        State::Closed => "closed",
    })
}

/// A trimmed title, or why it cannot be used.
fn valid_title(title: &str) -> std::result::Result<&str, &'static str> {
    let title = title.trim();
    if title.is_empty() {
        Err("A title is required.")
    } else if title.chars().count() > MAX_TITLE_CHARS {
        Err("That title is too long.")
    } else {
        Ok(title)
    }
}

/// Unwraps an `Outcome`, returning its failure from the enclosing method.
macro_rules! check {
    ($outcome:expr) => {
        match $outcome {
            Outcome::Ok(value) => value,
            Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
        }
    };
}

struct Work {
    db: D1Database,
    repos: Fetcher,
    events: Fetcher,
}

impl Work {
    async fn publish<T: Serialize>(
        &self,
        kind: &'static str,
        repo_id: &str,
        actor: &User,
        data: T,
    ) -> Result<()> {
        self.publish_as(kind, repo_id, Some(actor.id.clone()), data)
            .await
    }

    /// Publishes an event caused by `actor`, or by g1t itself.
    async fn publish_as<T: Serialize>(
        &self,
        kind: &'static str,
        repo_id: &str,
        actor: Option<String>,
        data: T,
    ) -> Result<()> {
        let event = NewEvent {
            kind,
            source: SOURCE,
            repo_id: Some(repo_id.to_owned()),
            actor,
            data,
        };
        g1t_kit::call(
            &self.events,
            "publish",
            &Publish {
                events: vec![event],
            },
        )
        .await
    }

    /// The repository, if the viewer may see it. Whether they may is
    /// decided by the repos service.
    async fn repo(&self, path: &RepoPath, viewer: &Viewer) -> Result<Outcome<Repo>> {
        g1t_kit::call(
            &self.repos,
            "get",
            &GetArgs {
                path: path.clone(),
                viewer: viewer.clone(),
            },
        )
        .await
    }

    /// The next number in the repository's sequence. Taking it is one
    /// statement, so concurrent opens cannot be given the same number.
    async fn next_number(&self, repo_id: &str) -> Result<u32> {
        let row = self
            .db
            .prepare(
                "INSERT INTO counters (repo_id, last) VALUES (?, 1)
                 ON CONFLICT (repo_id) DO UPDATE SET last = last + 1
                 RETURNING last AS n",
            )
            .bind(&[repo_id.into()])?
            .first::<NumberRow>(None)
            .await?;
        row.map(|row| row.n)
            .ok_or_else(|| worker::Error::RustError("no number was assigned".into()))
    }

    async fn issue(&self, repo_id: &str, number: u32) -> Result<Option<Issue>> {
        Ok(self
            .db
            .prepare(format!(
                "SELECT {ISSUE_COLUMNS} FROM issues WHERE repo_id = ? AND number = ?"
            ))
            .bind(&[repo_id.into(), number.into()])?
            .first::<IssueRow>(None)
            .await?
            .map(Issue::from))
    }

    async fn pull(&self, repo_id: &str, number: u32) -> Result<Option<Pull>> {
        Ok(self
            .db
            .prepare("SELECT * FROM pulls WHERE repo_id = ? AND number = ?")
            .bind(&[repo_id.into(), number.into()])?
            .first::<PullRow>(None)
            .await?
            .map(Pull::from))
    }

    async fn comments(&self, repo_id: &str, number: u32) -> Result<Vec<Comment>> {
        let rows = self
            .db
            .prepare(
                "SELECT * FROM comments WHERE repo_id = ? AND number = ? ORDER BY id LIMIT 500",
            )
            .bind(&[repo_id.into(), number.into()])?
            .all()
            .await?
            .results::<CommentRow>()?;
        Ok(rows.into_iter().map(Comment::from).collect())
    }

    /// The repository and one of its issues, as seen by `viewer`.
    async fn issue_at(
        &self,
        path: &RepoPath,
        number: u32,
        viewer: &Viewer,
    ) -> Result<Outcome<(Repo, Issue)>> {
        let Outcome::Ok(repo) = self.repo(path, viewer).await? else {
            return Ok(no_issue());
        };
        Ok(match self.issue(&repo.id, number).await? {
            Some(issue) => Outcome::Ok((repo, issue)),
            None => no_issue(),
        })
    }

    /// The repository and one of its pull requests, as seen by `viewer`.
    async fn pull_at(
        &self,
        path: &RepoPath,
        number: u32,
        viewer: &Viewer,
    ) -> Result<Outcome<(Repo, Pull)>> {
        let Outcome::Ok(repo) = self.repo(path, viewer).await? else {
            return Ok(no_pull());
        };
        Ok(match self.pull(&repo.id, number).await? {
            Some(pull) => Outcome::Ok((repo, pull)),
            None => no_pull(),
        })
    }

    fn issue_event(issue: &Issue) -> IssueEvent {
        IssueEvent {
            issue_id: issue.id.clone(),
            repo_id: issue.repo_id.clone(),
            number: issue.number,
            ..IssueEvent::default()
        }
    }

    fn pull_event(pull: &Pull) -> PullEvent {
        PullEvent {
            pull_id: pull.id.clone(),
            repo_id: pull.repo_id.clone(),
            number: pull.number,
            issue: pull.issue,
            ..PullEvent::default()
        }
    }

    // --- Issues ------------------------------------------------------------

    async fn open_issue(&self, a: OpenIssueArgs) -> Result<Outcome<Issue>> {
        if !a.actor.verified {
            return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
        }
        let title = match valid_title(&a.title) {
            Ok(title) => title,
            Err(message) => return Ok(Outcome::fail(FailureCode::Invalid, message)),
        };
        let Some(labels) = normalize_labels(&a.labels) else {
            return Ok(Outcome::fail(
                FailureCode::Invalid,
                "An issue can have up to 10 labels of up to 40 characters each.",
            ));
        };
        let repo = check!(self.repo(&a.repo, &Some(a.actor.clone())).await?);
        let checks: Vec<&str> = a
            .checks
            .iter()
            .map(|check| check.trim())
            .filter(|check| !check.is_empty())
            .collect();

        let now = now_ms();
        let id = new_id("iss", now);
        let number = self.next_number(&repo.id).await?;
        let timestamp = rfc3339(now);
        self.db
            .prepare(
                "INSERT INTO issues
                   (id, repo_id, number, title, body, labels, checks, author_id, author_name,
                    created_at, updated_at)
                 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
            )
            .bind(&[
                id.as_str().into(),
                repo.id.as_str().into(),
                number.into(),
                title.into(),
                a.body.trim().into(),
                serde_json::to_string(&labels)?.into(),
                serde_json::to_string(&checks)?.into(),
                a.actor.id.as_str().into(),
                a.actor.username.as_str().into(),
                timestamp.as_str().into(),
                timestamp.as_str().into(),
            ])?
            .run()
            .await?;
        let Some(issue) = self.issue(&repo.id, number).await? else {
            return Ok(no_issue());
        };
        self.publish(
            "issue.opened",
            &repo.id,
            &a.actor,
            IssueEvent {
                title: Some(issue.title.clone()),
                ..Self::issue_event(&issue)
            },
        )
        .await?;
        Ok(Outcome::Ok(issue))
    }

    async fn list_issues(&self, a: ListIssuesArgs) -> Result<Outcome<Vec<Issue>>> {
        let repo = check!(self.repo(&a.repo, &a.viewer).await?);
        let state = state_name(a.state).map_or(JsValue::NULL, JsValue::from);
        let label = a
            .label
            .map(|label| label.trim().to_lowercase())
            .filter(|label| !label.is_empty());
        let rows = self
            .db
            .prepare(format!(
                "SELECT {ISSUE_COLUMNS} FROM issues
                 WHERE repo_id = ? AND (? IS NULL OR state = ?)
                   AND (? IS NULL OR EXISTS
                     (SELECT 1 FROM json_each(issues.labels) WHERE json_each.value = ?))
                 ORDER BY number DESC LIMIT ?"
            ))
            .bind(&[
                repo.id.into(),
                state.clone(),
                state,
                optional(&label),
                optional(&label),
                LIST_PAGE.into(),
            ])?
            .all()
            .await?
            .results::<IssueRow>()?;
        Ok(Outcome::Ok(rows.into_iter().map(Issue::from).collect()))
    }

    async fn get_issue(&self, a: ViewArgs) -> Result<Outcome<IssueDetail>> {
        let (repo, issue) = check!(self.issue_at(&a.repo, a.number, &a.viewer).await?);
        let pulls = self
            .db
            .prepare("SELECT * FROM pulls WHERE issue_id = ? ORDER BY number")
            .bind(&[issue.id.as_str().into()])?
            .all()
            .await?
            .results::<PullRow>()?;
        Ok(Outcome::Ok(IssueDetail {
            comments: self.comments(&repo.id, issue.number).await?,
            pulls: pulls.into_iter().map(Pull::from).collect(),
            issue,
        }))
    }

    /// The issue, if `actor` wrote it or belongs to the repository's workspace.
    async fn manageable_issue(
        &self,
        actor: &User,
        path: &RepoPath,
        number: u32,
    ) -> Result<Outcome<Issue>> {
        let (repo, issue) = check!(self.issue_at(path, number, &Some(actor.clone())).await?);
        if issue.author.id != actor.id && !actor.is_member(&repo.namespace) {
            return Ok(Outcome::fail(
                FailureCode::Forbidden,
                "Only the author or a member of the workspace can change an issue.",
            ));
        }
        Ok(Outcome::Ok(issue))
    }

    async fn update_issue(&self, a: UpdateIssueArgs) -> Result<Outcome<Issue>> {
        let issue = check!(self.manageable_issue(&a.actor, &a.repo, a.number).await?);
        let title = match a.title.as_deref().map(valid_title) {
            Some(Err(message)) => return Ok(Outcome::fail(FailureCode::Invalid, message)),
            Some(Ok(title)) => Some(title.to_owned()),
            None => None,
        };
        let labels = match a.labels.as_deref().map(normalize_labels) {
            Some(None) => {
                return Ok(Outcome::fail(
                    FailureCode::Invalid,
                    "An issue can have up to 10 labels of up to 40 characters each.",
                ));
            }
            Some(Some(labels)) => Some(serde_json::to_string(&labels)?),
            None => None,
        };
        let body = a.body.map(|body| body.trim().to_owned());
        self.db
            .prepare(
                "UPDATE issues
                 SET title = COALESCE(?, title), body = COALESCE(?, body),
                     labels = COALESCE(?, labels), updated_at = ?
                 WHERE id = ?",
            )
            .bind(&[
                optional(&title),
                optional(&body),
                optional(&labels),
                rfc3339(now_ms()).into(),
                issue.id.as_str().into(),
            ])?
            .run()
            .await?;
        let Some(issue) = self.issue(&issue.repo_id, issue.number).await? else {
            return Ok(no_issue());
        };
        self.publish(
            "issue.updated",
            &issue.repo_id,
            &a.actor,
            Self::issue_event(&issue),
        )
        .await?;
        Ok(Outcome::Ok(issue))
    }

    async fn close_issue(&self, a: IssueActionArgs) -> Result<Outcome<Issue>> {
        let mut issue = check!(self.manageable_issue(&a.actor, &a.repo, a.number).await?);
        if issue.state == State::Closed {
            return Ok(Outcome::fail(
                FailureCode::Conflict,
                "This issue is already closed.",
            ));
        }
        let reason = a.reason.unwrap_or(IssueReason::Completed);
        let now = rfc3339(now_ms());
        self.db
            .prepare(
                "UPDATE issues SET state = 'closed', reason = ?, closed_at = ?, updated_at = ?
                 WHERE id = ?",
            )
            .bind(&[
                reason.as_str().into(),
                now.as_str().into(),
                now.as_str().into(),
                issue.id.as_str().into(),
            ])?
            .run()
            .await?;
        self.publish(
            "issue.closed",
            &issue.repo_id,
            &a.actor,
            IssueEvent {
                reason: Some(reason.as_str()),
                ..Self::issue_event(&issue)
            },
        )
        .await?;
        issue.state = State::Closed;
        issue.reason = Some(reason);
        issue.closed_at = Some(now.clone());
        issue.updated_at = now;
        Ok(Outcome::Ok(issue))
    }

    async fn reopen_issue(&self, a: IssueActionArgs) -> Result<Outcome<Issue>> {
        let mut issue = check!(self.manageable_issue(&a.actor, &a.repo, a.number).await?);
        if issue.state == State::Open {
            return Ok(Outcome::fail(
                FailureCode::Conflict,
                "This issue is already open.",
            ));
        }
        let now = rfc3339(now_ms());
        self.db
            .prepare(
                "UPDATE issues
                 SET state = 'open', reason = NULL, resolved_by = NULL, closed_at = NULL,
                     updated_at = ?
                 WHERE id = ?",
            )
            .bind(&[now.as_str().into(), issue.id.as_str().into()])?
            .run()
            .await?;
        self.publish(
            "issue.reopened",
            &issue.repo_id,
            &a.actor,
            Self::issue_event(&issue),
        )
        .await?;
        issue.state = State::Open;
        issue.reason = None;
        issue.resolved_by = None;
        issue.closed_at = None;
        issue.updated_at = now;
        Ok(Outcome::Ok(issue))
    }

    /// The default labels, then every other label in use on the repository.
    async fn list_labels(&self, a: ViewArgs) -> Result<Outcome<Vec<String>>> {
        let repo = check!(self.repo(&a.repo, &a.viewer).await?);
        let used = self
            .db
            .prepare(
                "SELECT DISTINCT json_each.value AS value
                 FROM issues, json_each(issues.labels)
                 WHERE issues.repo_id = ? ORDER BY 1 LIMIT 200",
            )
            .bind(&[repo.id.into()])?
            .all()
            .await?
            .results::<ValueRow>()?;
        let mut labels: Vec<String> = DEFAULT_LABELS.iter().map(|label| (*label).into()).collect();
        for row in used {
            if !labels.contains(&row.value) {
                labels.push(row.value);
            }
        }
        Ok(Outcome::Ok(labels))
    }

    async fn counts(&self, a: ViewArgs) -> Result<Outcome<Counts>> {
        let repo = check!(self.repo(&a.repo, &a.viewer).await?);
        let counts = self
            .db
            .prepare(
                "SELECT
                   (SELECT count(*) FROM issues WHERE repo_id = ? AND state = 'open') AS issues,
                   (SELECT count(*) FROM pulls
                    WHERE repo_id = ? AND status IN ('draft', 'open')) AS pulls",
            )
            .bind(&[repo.id.as_str().into(), repo.id.as_str().into()])?
            .first::<Counts>(None)
            .await?;
        Ok(Outcome::Ok(counts.unwrap_or(Counts {
            issues: 0,
            pulls: 0,
        })))
    }

    // --- Comments ----------------------------------------------------------

    async fn add_comment(&self, a: AddCommentArgs) -> Result<Outcome<Comment>> {
        if !a.actor.verified {
            return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
        }
        let body = a.body.trim();
        // An approval speaks for itself; anything else has to say something.
        if body.is_empty() && a.verdict != Some(Verdict::Approve) {
            return Ok(Outcome::fail(
                FailureCode::Invalid,
                "A comment cannot be empty.",
            ));
        }
        let path = a
            .path
            .as_deref()
            .map(str::trim)
            .filter(|path| !path.is_empty());
        let line = a.line.filter(|line| *line > 0 && path.is_some());
        if body.chars().count() > MAX_ENTRY_CHARS {
            return Ok(Outcome::fail(
                FailureCode::Invalid,
                "That comment is too long.",
            ));
        }
        let repo = check!(self.repo(&a.repo, &Some(a.actor.clone())).await?);
        // The number names an issue or a pull request, never both.
        let table = if self.issue(&repo.id, a.number).await?.is_some() {
            if path.is_some() || a.verdict.is_some() {
                return Ok(Outcome::fail(
                    FailureCode::Invalid,
                    "Only a pull request can be reviewed or commented on by line.",
                ));
            }
            "issues"
        } else if let Some(pull) = self.pull(&repo.id, a.number).await? {
            if a.verdict.is_some() && pull.author.id == a.actor.id {
                return Ok(Outcome::fail(
                    FailureCode::Forbidden,
                    "You cannot approve or request changes on your own pull request.",
                ));
            }
            "pulls"
        } else {
            return Ok(Outcome::fail(
                FailureCode::NotFound,
                "No issue or pull request has that number.",
            ));
        };

        let now = now_ms();
        let comment = Comment {
            id: new_id("cmt", now),
            author: a.actor.clone(),
            body: body.to_owned(),
            path: path.map(str::to_owned),
            line,
            verdict: a.verdict,
            created_at: rfc3339(now),
        };
        self.db
            .batch(vec![
                self.db
                    .prepare(
                        "INSERT INTO comments
                           (id, repo_id, number, author_id, author_name, body, path, line,
                            verdict, created_at)
                         VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
                    )
                    .bind(&[
                        comment.id.as_str().into(),
                        repo.id.as_str().into(),
                        a.number.into(),
                        a.actor.id.as_str().into(),
                        a.actor.username.as_str().into(),
                        body.into(),
                        optional(&comment.path),
                        optional_number(line),
                        a.verdict
                            .map_or(JsValue::NULL, |verdict| verdict.as_str().into()),
                        comment.created_at.as_str().into(),
                    ])?,
                self.db
                    .prepare(format!(
                        "UPDATE {table} SET updated_at = ? WHERE repo_id = ? AND number = ?"
                    ))
                    .bind(&[
                        comment.created_at.as_str().into(),
                        repo.id.as_str().into(),
                        a.number.into(),
                    ])?,
            ])
            .await?;
        self.publish(
            "comment.created",
            &repo.id,
            &a.actor,
            CommentCreated {
                comment_id: comment.id.clone(),
                repo_id: repo.id.clone(),
                number: a.number,
            },
        )
        .await?;
        Ok(Outcome::Ok(comment))
    }

    // --- Pull requests -----------------------------------------------------

    async fn open_pull(&self, a: OpenPullArgs) -> Result<Outcome<Pull>> {
        if !a.actor.verified {
            return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
        }
        let repo = check!(self.repo(&a.repo, &Some(a.actor.clone())).await?);
        let issue = match a.issue {
            Some(number) => match self.issue(&repo.id, number).await? {
                Some(issue) if issue.state == State::Open => Some(issue),
                Some(_) => {
                    return Ok(Outcome::fail(
                        FailureCode::Conflict,
                        "This issue is closed.",
                    ));
                }
                None => return Ok(no_issue()),
            },
            None => None,
        };
        // A pull request for an issue takes the issue's title unless given one.
        let title = match (a.title.trim(), &issue) {
            ("", Some(issue)) => issue.title.clone(),
            (title, _) => match valid_title(title) {
                Ok(title) => title.to_owned(),
                Err(message) => return Ok(Outcome::fail(FailureCode::Invalid, message)),
            },
        };
        let agent = match a.agent.trim() {
            "" => "agent",
            agent => agent,
        };
        let runtime = match a.runtime {
            Runtime::Hosted => "hosted",
            Runtime::External => "external",
        };

        let now = now_ms();
        let id = new_id("pr", now);
        let branch = a
            .branch
            .as_deref()
            .map(str::trim)
            .filter(|branch| !branch.is_empty());
        // The change is on a branch already pushed to the repository, or
        // will be made in a fork created for this pull request.
        let (fork, head) = match branch {
            Some(branch) => {
                if branch == repo.default_branch {
                    return Ok(Outcome::fail(
                        FailureCode::Invalid,
                        format!("Choose a branch other than {branch}."),
                    ));
                }
                let head: Option<String> = g1t_kit::call(
                    &self.repos,
                    "head",
                    &HeadArgs {
                        repo_id: repo.id.clone(),
                        branch: branch.to_owned(),
                    },
                )
                .await?;
                let Some(head) = head else {
                    return Ok(Outcome::fail(
                        FailureCode::NotFound,
                        format!("There is no branch named {branch}. Push it first."),
                    ));
                };
                let existing = self
                    .db
                    .prepare(
                        "SELECT number AS n FROM pulls
                         WHERE repo_id = ? AND source_branch = ? AND status IN ('draft', 'open')",
                    )
                    .bind(&[repo.id.as_str().into(), branch.into()])?
                    .first::<NumberRow>(None)
                    .await?;
                if let Some(existing) = existing {
                    return Ok(Outcome::fail(
                        FailureCode::Conflict,
                        format!("Pull request #{} is already open for {branch}.", existing.n),
                    ));
                }
                (None, Some(head))
            }
            None => {
                let fork: Outcome<Repo> = g1t_kit::call(
                    &self.repos,
                    "fork_for_pull",
                    &ForkArgs {
                        source_id: repo.id.clone(),
                        pull_id: id.clone(),
                        actor: a.actor.clone(),
                    },
                )
                .await?;
                (Some(check!(fork)), None)
            }
        };
        // A branch already holds the work, so its pull request is ready for
        // review from the start; one with a fork starts as a draft.
        let status = if branch.is_some() { "open" } else { "draft" };
        let body = Some(a.body.trim().to_owned()).filter(|body| !body.is_empty());

        let number = self.next_number(&repo.id).await?;
        let timestamp = rfc3339(now);
        self.db
            .prepare(
                "INSERT INTO pulls
                   (id, repo_id, number, issue_id, issue_number, title, body, agent, runtime,
                    status, fork_repo_id, fork_namespace, fork_name, source_branch, head_commit,
                    author_id, author_name, created_at, updated_at)
                 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
            )
            .bind(&[
                id.as_str().into(),
                repo.id.as_str().into(),
                number.into(),
                optional(&issue.as_ref().map(|issue| issue.id.clone())),
                optional_number(issue.as_ref().map(|issue| issue.number)),
                title.into(),
                optional(&body),
                agent.into(),
                runtime.into(),
                status.into(),
                optional(&fork.as_ref().map(|fork| fork.id.clone())),
                optional(&fork.as_ref().map(|fork| fork.namespace.clone())),
                optional(&fork.as_ref().map(|fork| fork.name.clone())),
                optional(&branch.map(str::to_owned)),
                optional(&head),
                a.actor.id.as_str().into(),
                a.actor.username.as_str().into(),
                timestamp.as_str().into(),
                timestamp.as_str().into(),
            ])?
            .run()
            .await?;
        let Some(pull) = self.pull(&repo.id, number).await? else {
            return Ok(no_pull());
        };
        self.publish(
            "pull.opened",
            &repo.id,
            &a.actor,
            PullEvent {
                agent: Some(pull.agent.clone()),
                ..Self::pull_event(&pull)
            },
        )
        .await?;
        Ok(Outcome::Ok(pull))
    }

    async fn list_pulls(&self, a: ListPullsArgs) -> Result<Outcome<Vec<Pull>>> {
        let repo = check!(self.repo(&a.repo, &a.viewer).await?);
        let filter = match a.state {
            Some(State::Open) => "AND status IN ('draft', 'open')",
            Some(State::Closed) => "AND status IN ('merged', 'closed')",
            None => "",
        };
        let rows = self
            .db
            .prepare(format!(
                "SELECT * FROM pulls WHERE repo_id = ? {filter} ORDER BY number DESC LIMIT ?"
            ))
            .bind(&[repo.id.into(), LIST_PAGE.into()])?
            .all()
            .await?
            .results::<PullRow>()?;
        Ok(Outcome::Ok(rows.into_iter().map(Pull::from).collect()))
    }

    async fn get_pull(&self, a: ViewArgs) -> Result<Outcome<PullDetail>> {
        let (repo, pull) = check!(self.pull_at(&a.repo, a.number, &a.viewer).await?);
        let issue = match pull.issue {
            Some(number) => self.issue(&repo.id, number).await?,
            None => None,
        };
        Ok(Outcome::Ok(PullDetail {
            comments: self.comments(&repo.id, pull.number).await?,
            checks: self.latest_checks(&pull.id).await?,
            issue,
            pull,
        }))
    }

    /// The pull request, if it is still active and `actor` opened it or
    /// belongs to the repository's workspace.
    async fn manageable_pull(
        &self,
        actor: &User,
        path: &RepoPath,
        number: u32,
    ) -> Result<Outcome<Pull>> {
        let (repo, pull) = check!(self.pull_at(path, number, &Some(actor.clone())).await?);
        if pull.author.id != actor.id && !actor.is_member(&repo.namespace) {
            return Ok(Outcome::fail(
                FailureCode::Forbidden,
                "Only whoever opened a pull request, or a member of the workspace, can change it.",
            ));
        }
        if !pull.status.is_active() {
            return Ok(Outcome::fail(
                FailureCode::Conflict,
                format!("This pull request is already {}.", pull.status.as_str()),
            ));
        }
        Ok(Outcome::Ok(pull))
    }

    /// Marks a draft ready for review, or updates the description of one
    /// that already is.
    async fn ready_pull(&self, a: PullActionArgs) -> Result<Outcome<Pull>> {
        let mut pull = check!(self.manageable_pull(&a.actor, &a.repo, a.number).await?);
        let summary = Some(a.summary.trim().to_owned()).filter(|summary| !summary.is_empty());
        let now = rfc3339(now_ms());
        self.db
            .prepare(
                "UPDATE pulls SET status = 'open', body = COALESCE(?, body), updated_at = ?
                 WHERE id = ?",
            )
            .bind(&[
                optional(&summary),
                now.as_str().into(),
                pull.id.as_str().into(),
            ])?
            .run()
            .await?;
        if pull.status == PullStatus::Draft {
            self.publish(
                "pull.ready",
                &pull.repo_id,
                &a.actor,
                Self::pull_event(&pull),
            )
            .await?;
        }
        pull.status = PullStatus::Open;
        pull.body = summary.or(pull.body);
        pull.updated_at = now;
        Ok(Outcome::Ok(pull))
    }

    async fn close_pull(&self, a: PullActionArgs) -> Result<Outcome<Pull>> {
        let mut pull = check!(self.manageable_pull(&a.actor, &a.repo, a.number).await?);
        let now = rfc3339(now_ms());
        self.db
            .prepare("UPDATE pulls SET status = 'closed', updated_at = ? WHERE id = ?")
            .bind(&[now.as_str().into(), pull.id.as_str().into()])?
            .run()
            .await?;
        self.publish(
            "pull.closed",
            &pull.repo_id,
            &a.actor,
            Self::pull_event(&pull),
        )
        .await?;
        pull.status = PullStatus::Closed;
        pull.updated_at = now;
        Ok(Outcome::Ok(pull))
    }

    /// Lands the pull request on the repository's default branch. Unless
    /// told to keep it open, that resolves the issue it was for: the issue
    /// closes naming this pull request, and the others still in progress
    /// for it close as superseded.
    async fn merge_pull(&self, a: PullActionArgs) -> Result<Outcome<Pull>> {
        let viewer = Some(a.actor.clone());
        let (repo, mut pull) = check!(self.pull_at(&a.repo, a.number, &viewer).await?);
        match pull.status {
            PullStatus::Open => {}
            PullStatus::Draft => {
                return Ok(Outcome::fail(
                    FailureCode::Conflict,
                    "This pull request is still a draft. Mark it ready for review first.",
                ));
            }
            status => {
                return Ok(Outcome::fail(
                    FailureCode::Conflict,
                    format!("This pull request is already {}.", status.as_str()),
                ));
            }
        }
        if !a.ignore_checks {
            let waiting = match pull.check_status {
                Some(CheckStatus::Queued | CheckStatus::Running) => {
                    Some("The acceptance checks are still running.")
                }
                Some(CheckStatus::Failed) => Some("The acceptance checks did not pass."),
                Some(CheckStatus::Errored) => Some("The acceptance checks could not be run."),
                Some(CheckStatus::Passed) | None => None,
            };
            if let Some(reason) = waiting {
                return Ok(Outcome::fail(
                    FailureCode::Conflict,
                    format!("{reason} Wait or fix them, or merge anyway by ignoring the checks."),
                ));
            }
        }
        let issue = match pull.issue {
            Some(number) if !a.keep_issue_open => self
                .issue(&repo.id, number)
                .await?
                .filter(|issue| issue.state == State::Open),
            _ => None,
        };

        // Whether the actor may write to the repository is decided by repos.
        let landed: Outcome<Landed> = g1t_kit::call(
            &self.repos,
            "land",
            &LandArgs {
                // A pull request from a branch lands from the repository itself.
                source_id: pull.fork_repo_id.clone().unwrap_or_else(|| repo.id.clone()),
                branch: pull.branch.clone(),
                actor: a.actor.clone(),
            },
        )
        .await?;
        let landed = check!(landed);

        let now = rfc3339(now_ms());
        let mut statements = vec![
            self.db
                .prepare(
                    "UPDATE pulls
                     SET status = 'merged', head_commit = ?, merge_base = ?, merged_by = ?,
                         merged_at = ?, updated_at = ?
                     WHERE id = ?",
                )
                .bind(&[
                    landed.commit.as_str().into(),
                    optional(&landed.previous),
                    a.actor.username.as_str().into(),
                    now.as_str().into(),
                    now.as_str().into(),
                    pull.id.as_str().into(),
                ])?,
        ];
        if let Some(issue) = &issue {
            statements.push(
                self.db
                    .prepare(
                        "UPDATE issues
                         SET state = 'closed', reason = 'completed', resolved_by = ?,
                             closed_at = ?, updated_at = ?
                         WHERE id = ?",
                    )
                    .bind(&[
                        pull.number.into(),
                        now.as_str().into(),
                        now.as_str().into(),
                        issue.id.as_str().into(),
                    ])?,
            );
            statements.push(
                self.db
                    .prepare(
                        "UPDATE pulls SET status = 'closed', superseded_by = ?, updated_at = ?
                         WHERE issue_id = ? AND id != ? AND status IN ('draft', 'open')",
                    )
                    .bind(&[
                        pull.number.into(),
                        now.as_str().into(),
                        issue.id.as_str().into(),
                        pull.id.as_str().into(),
                    ])?,
            );
        }
        self.db.batch(statements).await?;

        self.publish(
            "pull.merged",
            &repo.id,
            &a.actor,
            PullEvent {
                commit: Some(landed.commit.clone()),
                ..Self::pull_event(&pull)
            },
        )
        .await?;
        if let Some(issue) = &issue {
            self.publish(
                "issue.closed",
                &repo.id,
                &a.actor,
                IssueEvent {
                    reason: Some(IssueReason::Completed.as_str()),
                    resolved_by: Some(pull.number),
                    ..Self::issue_event(issue)
                },
            )
            .await?;
        }

        pull.status = PullStatus::Merged;
        pull.head_commit = Some(landed.commit);
        pull.merge_base = landed.previous;
        pull.merged_by = Some(a.actor.username);
        pull.merged_at = Some(now.clone());
        pull.updated_at = now;
        Ok(Outcome::Ok(pull))
    }

    async fn list_active_pulls(&self, a: ViewerArgs) -> Result<Vec<ActivePull>> {
        let Some(viewer) = a.viewer else {
            return Ok(Vec::new());
        };
        let rows = self
            .db
            .prepare(
                "SELECT * FROM pulls
                 WHERE author_id = ? AND status IN ('draft', 'open')
                 ORDER BY updated_at DESC LIMIT 50",
            )
            .bind(&[viewer.id.into()])?
            .all()
            .await?
            .results::<PullRow>()?;
        let mut active = Vec::with_capacity(rows.len());
        for pull in rows.into_iter().map(Pull::from) {
            let issue = match pull.issue {
                Some(number) => self.issue(&pull.repo_id, number).await?,
                None => None,
            };
            active.push(ActivePull { pull, issue });
        }
        Ok(active)
    }

    // --- Sessions ----------------------------------------------------------

    async fn append_session(&self, a: AppendSessionArgs) -> Result<Outcome<Appended>> {
        if a.entries.is_empty() {
            return Ok(Outcome::Ok(Appended { count: 0 }));
        }
        if a.entries.len() > MAX_ENTRY_BATCH {
            return Ok(Outcome::fail(
                FailureCode::Invalid,
                format!("Send at most {MAX_ENTRY_BATCH} entries at a time."),
            ));
        }
        let viewer = Some(a.actor.clone());
        let (_, pull) = check!(self.pull_at(&a.repo, a.number, &viewer).await?);
        if pull.author.id != a.actor.id {
            return Ok(Outcome::fail(
                FailureCode::Forbidden,
                "Only whoever opened a pull request can record its session.",
            ));
        }

        let now = rfc3339(now_ms());
        let count = a.entries.len() as u32;
        let mut statements = Vec::with_capacity(a.entries.len() + 1);
        for entry in a.entries {
            let kind = serde_json::to_value(entry.kind)?;
            let text: String = entry.text.chars().take(MAX_ENTRY_CHARS).collect();
            // Each insert takes the next sequence number itself, so two
            // writers appending at once cannot collide.
            statements.push(
                self.db
                    .prepare(
                        "INSERT INTO session_entries (pull_id, seq, kind, text, tool, \"commit\", at)
                         SELECT ?, COALESCE(MAX(seq), 0) + 1, ?, ?, ?, ?, ?
                         FROM session_entries WHERE pull_id = ?",
                    )
                    .bind(&[
                        pull.id.as_str().into(),
                        kind.as_str().unwrap_or("note").into(),
                        text.into(),
                        optional(&entry.tool),
                        optional(&entry.commit.or_else(|| pull.head_commit.clone())),
                        now.as_str().into(),
                        pull.id.as_str().into(),
                    ])?,
            );
        }
        statements.push(
            self.db
                .prepare("UPDATE pulls SET updated_at = ? WHERE id = ?")
                .bind(&[now.as_str().into(), pull.id.as_str().into()])?,
        );
        self.db.batch(statements).await?;
        self.publish(
            "session.appended",
            &pull.repo_id,
            &a.actor,
            SessionAppended {
                pull_id: pull.id.clone(),
                repo_id: pull.repo_id.clone(),
                number: pull.number,
                count,
            },
        )
        .await?;
        Ok(Outcome::Ok(Appended { count }))
    }

    async fn read_session(&self, a: ViewArgs) -> Result<Outcome<Vec<SessionEntry>>> {
        let (_, pull) = check!(self.pull_at(&a.repo, a.number, &a.viewer).await?);
        let rows = self
            .db
            .prepare(
                "SELECT seq, kind, text, tool, \"commit\", at FROM session_entries
                 WHERE pull_id = ? AND seq > ? ORDER BY seq LIMIT ?",
            )
            .bind(&[pull.id.into(), a.after_seq.into(), SESSION_PAGE.into()])?
            .all()
            .await?
            .results::<SessionRow>()?;
        Ok(Outcome::Ok(
            rows.into_iter().map(SessionEntry::from).collect(),
        ))
    }

    /// A push moves the head of the pull request it concerns: the one whose
    /// fork was pushed to, or the one opened from the branch that moved.
    async fn on_event(&self, event: &Event) -> Result<()> {
        if event.kind != "git.push" {
            return Ok(());
        }
        let (Some(repo_id), Some(after), Some(git_ref)) = (
            event.repo_id.as_deref(),
            event.data["after"].as_str(),
            event.data["ref"].as_str(),
        ) else {
            return Ok(());
        };
        let now = rfc3339(now_ms());
        // The head moved, so whatever the checks said no longer applies.
        let moved = "UPDATE pulls
             SET head_commit = ?, updated_at = ?, check_status = NULL, check_run_id = NULL";
        let active = "status IN ('draft', 'open') AND head_commit IS NOT ?";
        let returning = "RETURNING id, repo_id, number, issue_number, status";
        let mut pulls: Vec<MovedRow> = Vec::new();
        // A fork carries its pull request on its default branch.
        if event.data["defaultBranch"].as_bool() == Some(true) {
            pulls.extend(
                self.db
                    .prepare(format!(
                        "{moved} WHERE fork_repo_id = ? AND {active} {returning}"
                    ))
                    .bind(&[
                        after.into(),
                        now.as_str().into(),
                        repo_id.into(),
                        after.into(),
                    ])?
                    .all()
                    .await?
                    .results::<MovedRow>()?,
            );
        }
        if let Some(branch) = git_ref.strip_prefix("refs/heads/") {
            pulls.extend(
                self.db
                    .prepare(format!(
                        "{moved} WHERE repo_id = ? AND source_branch = ? AND {active} {returning}"
                    ))
                    .bind(&[
                        after.into(),
                        now.as_str().into(),
                        repo_id.into(),
                        branch.into(),
                        after.into(),
                    ])?
                    .all()
                    .await?
                    .results::<MovedRow>()?,
            );
        }
        // A draft is announced when it is marked ready instead.
        for pull in pulls
            .into_iter()
            .filter(|pull| pull.status == PullStatus::Open)
        {
            self.publish_as(
                "pull.updated",
                &pull.repo_id,
                event.actor.clone(),
                PullEvent {
                    pull_id: pull.id,
                    repo_id: pull.repo_id.clone(),
                    number: pull.number,
                    issue: pull.issue_number,
                    commit: Some(after.to_owned()),
                    ..PullEvent::default()
                },
            )
            .await?;
        }
        Ok(())
    }
}

fn service(env: &Env) -> Result<Work> {
    Ok(Work {
        db: env.d1("DB")?,
        repos: env.service("REPOS")?,
        events: env.service("EVENTS")?,
    })
}

#[event(fetch)]
async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> {
    let Some(method) = rpc_method(&request) else {
        return Response::error("Not found", 404);
    };
    let body: serde_json::Value = request.json().await?;
    let work = service(&env)?;

    match method.as_str() {
        "open_issue" => reply(&work.open_issue(args(body)?).await?),
        "list_issues" => reply(&work.list_issues(args(body)?).await?),
        "get_issue" => reply(&work.get_issue(args(body)?).await?),
        "update_issue" => reply(&work.update_issue(args(body)?).await?),
        "close_issue" => reply(&work.close_issue(args(body)?).await?),
        "reopen_issue" => reply(&work.reopen_issue(args(body)?).await?),
        "list_labels" => reply(&work.list_labels(args(body)?).await?),
        "counts" => reply(&work.counts(args(body)?).await?),
        "add_comment" => reply(&work.add_comment(args(body)?).await?),
        "start_checks" => reply(&work.start_checks(args(body)?).await?),
        "report_checks" => reply(&work.report_checks(args(body)?).await?),
        "open_pull" => reply(&work.open_pull(args(body)?).await?),
        "list_pulls" => reply(&work.list_pulls(args(body)?).await?),
        "get_pull" => reply(&work.get_pull(args(body)?).await?),
        "ready_pull" => reply(&work.ready_pull(args(body)?).await?),
        "close_pull" => reply(&work.close_pull(args(body)?).await?),
        "merge_pull" => reply(&work.merge_pull(args(body)?).await?),
        "list_active_pulls" => reply(&work.list_active_pulls(args(body)?).await?),
        "append_session" => reply(&work.append_session(args(body)?).await?),
        "read_session" => reply(&work.read_session(args(body)?).await?),
        _ => Response::error("Unknown method", 404),
    }
}

/// Events from the bus, delivered on this service's own queue.
#[event(queue)]
async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: Context) -> Result<()> {
    let work = service(&env)?;
    for message in batch.messages()? {
        work.on_event(message.body()).await?;
        message.ack();
    }
    Ok(())
}