flagon-io/g1t

public

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

g1t/crates/contracts/src/events.rs

295 lines9,574 bytesCodeBlame
//! Events published on the bus, and the events service that carries them.
//! Mirrors `packages/contracts/src/events.ts`.

use serde::Serialize;

/// What a publisher supplies; the bus fills in the id and time.
#[derive(Debug, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct NewEvent<T: Serialize> {
    #[serde(rename = "type")]
    pub kind: &'static str,
    /// The service that published it.
    pub source: &'static str,
    /// The repo the event concerns.
    pub repo_id: Option<String>,
    /// The user or agent that caused it, if any.
    pub actor: Option<String>,
    pub data: T,
}

#[derive(Debug, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct RepoCreated {
    pub repo_id: String,
    pub namespace: String,
    pub name: String,
    pub is_private: bool,
}

#[derive(Debug, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct RepoForked {
    pub repo_id: String,
    pub source_repo_id: String,
    pub pull_id: String,
}

/// One branch or tag moved by a push. `after` is the commit it points to now.
#[derive(Debug, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct GitPush {
    pub repo_id: String,
    /// The full ref, such as `refs/heads/main`.
    #[serde(rename = "ref")]
    pub git_ref: String,
    /// Where it pointed before; absent for a new branch or tag.
    #[serde(skip_serializing_if = "Option::is_none")]
    pub before: Option<String>,
    pub after: String,
    /// Whether the ref is the repository's default branch.
    pub default_branch: bool,
}

/// The payload of `issue.opened`, `issue.updated`, `issue.assigned`,
/// `issue.closed` and `issue.reopened`; each uses the fields that apply to it.
#[derive(Debug, Default, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct IssueEvent {
    pub issue_id: String,
    pub repo_id: String,
    pub number: u32,
    #[serde(skip_serializing_if = "Option::is_none")]
    pub title: Option<String>,
    /// On close: `completed` or `not_planned`.
    #[serde(skip_serializing_if = "Option::is_none")]
    pub reason: Option<&'static str>,
    /// On close: the number of the pull request whose merge closed it.
    #[serde(skip_serializing_if = "Option::is_none")]
    pub resolved_by: Option<u32>,
    /// On `issue.assigned`: the people it is now assigned to.
    #[serde(skip_serializing_if = "Option::is_none")]
    pub assignees: Option<Vec<String>>,
}

/// The payload of `pull.opened`, `pull.ready`, `pull.updated` (its head
/// moved), `pull.closed` and `pull.merged`; each uses the fields that apply to it.
#[derive(Debug, Default, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct PullEvent {
    pub pull_id: String,
    pub repo_id: String,
    pub number: u32,
    /// The number of the issue it is for.
    #[serde(skip_serializing_if = "Option::is_none")]
    pub issue: Option<u32>,
    #[serde(skip_serializing_if = "Option::is_none")]
    pub agent: Option<String>,
    /// On merge: the commit the branch now points to. On update and when
    /// marked ready: the head of the change.
    #[serde(skip_serializing_if = "Option::is_none")]
    pub commit: Option<String>,
    /// On close: the pull request that was merged instead.
    #[serde(skip_serializing_if = "Option::is_none")]
    pub superseded_by: Option<u32>,
}

/// `checks.completed`: a run of an issue's acceptance checks against a pull
/// request finished.
#[derive(Debug, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct ChecksEvent {
    pub pull_id: String,
    pub repo_id: String,
    pub number: u32,
    /// `passed`, `failed` or `errored`.
    pub status: &'static str,
    /// The commit that was checked.
    pub commit: String,
}

/// `workflow.completed`: a GitHub Actions run finished.
#[derive(Debug, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct WorkflowEvent {
    pub run_id: String,
    pub repo_id: String,
    /// The workflow's name, and its file.
    pub workflow: String,
    pub path: String,
    /// The run's number among the workflow's runs.
    pub number: u64,
    /// The GitHub event that started it, such as `push`.
    pub event: String,
    /// `success`, `failure`, `cancelled` or `skipped`.
    pub conclusion: String,
    #[serde(rename = "ref")]
    pub git_ref: String,
    pub sha: String,
    /// The pull request it ran for, if any.
    #[serde(skip_serializing_if = "Option::is_none")]
    pub pull: Option<u32>,
}

/// `review.completed`: a g1t agent finished reviewing a pull request, or
/// could not.
#[derive(Debug, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct ReviewEvent {
    pub pull_id: String,
    pub repo_id: String,
    pub number: u32,
    /// `approve` or `request_changes`; absent when no review was written.
    #[serde(skip_serializing_if = "Option::is_none")]
    pub verdict: Option<&'static str>,
}

/// `comment.created`. `number` is the issue or pull request commented on.
#[derive(Debug, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct CommentCreated {
    pub comment_id: String,
    pub repo_id: String,
    pub number: u32,
    /// Set when the comment is on a pull request.
    #[serde(skip_serializing_if = "Option::is_none")]
    pub pull_id: Option<String>,
    /// Set when the comment is a review: approve or request changes.
    #[serde(skip_serializing_if = "Option::is_none")]
    pub verdict: Option<crate::work::Verdict>,
}

#[derive(Debug, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct SessionAppended {
    pub pull_id: String,
    pub repo_id: String,
    pub number: u32,
    pub count: u32,
}

/// An event as stored in the log and delivered to subscribers. `data` is
/// left as JSON; each reader decodes the types it cares about.
#[derive(Clone, Debug, Serialize, serde::Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct Event {
    /// Sorts by the time the event was published.
    pub id: String,
    #[serde(rename = "type")]
    pub kind: String,
    /// The service that published it.
    pub source: String,
    /// RFC 3339.
    pub time: String,
    /// The repo the event concerns.
    pub repo_id: Option<String>,
    /// The user or agent that caused it, if any.
    pub actor: Option<String>,
    pub data: serde_json::Value,
}

/// `publish`, as a publisher sends it. Returns nothing.
#[derive(Debug, Serialize)]
pub struct Publish<T: Serialize> {
    pub events: Vec<NewEvent<T>>,
}

/// `publish`, as the events service reads it.
#[derive(Debug, serde::Deserialize)]
pub struct PublishArgs {
    pub events: Vec<Published>,
}

/// A [`NewEvent`] of any type, as received.
#[derive(Debug, serde::Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct Published {
    #[serde(rename = "type")]
    pub kind: String,
    pub source: String,
    #[serde(default)]
    pub repo_id: Option<String>,
    #[serde(default)]
    pub actor: Option<String>,
    pub data: serde_json::Value,
}

/// `list`: events from the log, newest first. Returns `Vec<Event>`.
#[derive(Debug, Default, Serialize, serde::Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct ListArgs {
    #[serde(default)]
    pub repo_id: Option<String>,
    /// Only these types; all types when empty.
    #[serde(default)]
    pub types: Vec<String>,
    /// Only events older than this event id.
    #[serde(default)]
    pub before: Option<String>,
    #[serde(default)]
    pub limit: Option<u32>,
}

/// `workspace.renamed`: a workspace's slug changed from `from` to `to`.
/// Every service that stores a slug moves its rows to the workspace's
/// *current* slug (ask identity by `workspace_id`), so that a repeated or
/// late delivery after a second rename still lands in the right place.
#[derive(Clone, Debug, Serialize, serde::Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct WorkspaceRenamed {
    pub workspace_id: String,
    pub from: String,
    pub to: String,
}

impl WorkspaceRenamed {
    /// The slugs whose rows move to `current`: the two this rename names,
    /// minus `current` itself. Moving rows keyed by either converges on the
    /// current slug whatever order renames are delivered in.
    pub fn stale_slugs(&self, current: &str) -> Vec<String> {
        let mut slugs: Vec<String> = Vec::new();
        for slug in [&self.from, &self.to] {
            if slug != current && !slugs.contains(slug) {
                slugs.push(slug.clone());
            }
        }
        slugs
    }
}

/// `queue.changed`: a repository's merge queue gained, lost or settled an
/// entry, so the next batch may be ready to test.
#[derive(Debug, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct QueueChanged {
    pub repo_id: String,
}

#[cfg(test)]
mod tests {
    use super::*;

    fn renamed(from: &str, to: &str) -> WorkspaceRenamed {
        WorkspaceRenamed {
            workspace_id: "wsp_1".into(),
            from: from.into(),
            to: to.into(),
        }
    }

    #[test]
    fn stale_slugs_leave_out_the_current_one() {
        assert_eq!(renamed("a", "b").stale_slugs("b"), vec!["a"]);
        // Delivered after a second rename, b → c: both move to c.
        assert_eq!(renamed("a", "b").stale_slugs("c"), vec!["a", "b"]);
        // Renamed back: a → b → a.
        assert_eq!(renamed("a", "b").stale_slugs("a"), vec!["b"]);
    }

    #[test]
    fn reads_the_published_payload() {
        let data = serde_json::json!({ "workspaceId": "wsp_1", "from": "a", "to": "b" });
        let event: WorkspaceRenamed = serde_json::from_value(data).unwrap();
        assert_eq!((event.from.as_str(), event.to.as_str()), ("a", "b"));
    }
}