use reqwest::StatusCode;
use serde::de::DeserializeOwned;
use serde::{Deserialize, Serialize};
use serde_json::{json, Value};
use tokio::time::sleep;
use tracing::warn;
#[cfg(test)]
use crate::pm::PmKr;
use crate::pm::{
parse_project_content, project_slug, render_project_content, IssueComment, IssueObservation,
PmError, PmItem, PmItemCreate, PmItemUpdate, PmProject, PmResult, PmWave, ProjectContent,
TeamBinding, RATE_LIMIT_RETRIES,
};
const LINEAR_BASE_URL: &str = "https://api.linear.app/graphql";
const LIST_ITEMS_PAGE_SIZE: u32 = 50;
const LIST_PROJECTS_PAGE_SIZE: u32 = 50;
const COMPLETED_STATE_TYPE: &str = "completed";
const REPOSITORY_CLAIM_PREFIX: &str = "<!-- loopflow-repository:";
const LIST_TEAMS_QUERY: &str = r#"query ListTeams {
teams {
nodes {
id
name
key
description
}
}
}"#;
const CREATE_TEAM_MUTATION: &str = r#"mutation CreateTeam($name: String!, $key: String!, $description: String!) {
teamCreate(input: { name: $name, key: $key, description: $description }) {
team {
id
}
}
}"#;
const UPDATE_TEAM_DESCRIPTION_MUTATION: &str = r#"mutation UpdateTeamDescription($id: String!, $description: String!) {
teamUpdate(id: $id, input: { description: $description }) {
team {
id
}
}
}"#;
const CREATE_INITIATIVE_MUTATION: &str = r#"mutation CreateInitiative($name: String!, $description: String!) {
initiativeCreate(input: { name: $name, description: $description }) {
initiative {
id
}
}
}"#;
const UPDATE_INITIATIVE_MUTATION: &str = r#"mutation UpdateInitiative($id: String!, $name: String!) {
initiativeUpdate(id: $id, input: { name: $name }) {
initiative {
id
}
}
}"#;
const LIST_INITIATIVES_QUERY: &str = r#"query ListInitiatives($after: String, $first: Int!) {
initiatives(after: $after, first: $first) {
nodes {
id
name
description
}
pageInfo {
hasNextPage
endCursor
}
}
}"#;
const LIST_INITIATIVE_PROJECTS_QUERY: &str = r#"query ListInitiativeProjects($initiativeId: String!, $after: String, $first: Int!) {
initiative(id: $initiativeId) {
projects(after: $after, first: $first, includeSubInitiatives: false) {
nodes {
id
name
description
content
initiatives(first: 50) {
nodes {
id
}
}
teams(first: 50) {
nodes {
id
}
}
}
pageInfo {
hasNextPage
endCursor
}
}
}
}"#;
const CREATE_PROJECT_MUTATION: &str = r#"mutation CreateProject($name: String!, $description: String!, $content: String!, $teamId: String!) {
projectCreate(input: { name: $name, description: $description, content: $content, teamIds: [$teamId] }) {
project {
id
}
}
}"#;
const UPDATE_PROJECT_MUTATION: &str = r#"mutation UpdateProject($id: String!, $name: String!, $description: String!, $content: String!) {
projectUpdate(id: $id, input: { name: $name, description: $description, content: $content }) {
project {
id
}
}
}"#;
const MOVE_PROJECT_TO_TEAM_MUTATION: &str = r#"mutation MoveProjectToTeam($id: String!, $teamId: String!) {
projectUpdate(id: $id, input: { teamIds: [$teamId] }) {
project {
id
}
}
}"#;
const SET_PROJECT_TEAMS_MUTATION: &str = r#"mutation SetProjectTeams($id: String!, $teamIds: [String!]!) {
projectUpdate(id: $id, input: { teamIds: $teamIds }) {
project {
id
}
}
}"#;
const ARCHIVE_PROJECT_MUTATION: &str = r#"mutation ArchiveProject($id: String!) {
projectArchive(id: $id) {
success
}
}"#;
const ATTACH_PROJECT_MUTATION: &str = r#"mutation AttachProject($initiativeId: String!, $projectId: String!) {
initiativeToProjectCreate(input: { initiativeId: $initiativeId, projectId: $projectId }) {
initiativeToProject {
id
}
}
}"#;
const LIST_ITEMS_QUERY: &str = r#"query ListProjectIssues($projectId: String!, $after: String, $first: Int!) {
project(id: $projectId) {
issues(first: $first, after: $after) {
nodes {
id
identifier
url
title
description
prioritySortOrder
sortOrder
assignee {
id
}
state {
type
}
project {
id
name
}
team {
id
}
}
pageInfo {
hasNextPage
endCursor
}
}
}
}"#;
const ISSUE_OWNERSHIP_QUERY: &str = r#"query IssueOwnership($id: String!) {
issue(id: $id) {
id
identifier
url
title
description
prioritySortOrder
sortOrder
assignee { id }
state { type }
team { id }
project {
id
name
description
content
initiatives(first: 50) { nodes { id } }
teams(first: 50) { nodes { id } }
}
}
}"#;
const PROJECT_OWNERSHIP_QUERY: &str = r#"query ProjectOwnership($id: String!) {
project(id: $id) {
id
name
description
content
initiatives(first: 50) { nodes { id } }
teams(first: 50) { nodes { id } }
}
}"#;
const CREATE_ITEM_MUTATION: &str = r#"mutation CreateIssue($teamId: String!, $projectId: String!, $title: String!, $description: String!, $stateId: String) {
issueCreate(input: { teamId: $teamId, projectId: $projectId, title: $title, description: $description, stateId: $stateId }) {
issue {
id
}
}
}"#;
const UPDATE_ITEM_MUTATION: &str = r#"mutation UpdateIssue($id: String!, $input: IssueUpdateInput!) {
issueUpdate(id: $id, input: $input) {
issue {
id
}
}
}"#;
const MOVE_ITEM_MUTATION: &str = r#"mutation MoveIssueToProject($id: String!, $projectId: String!) {
issueUpdate(id: $id, input: { projectId: $projectId }) {
issue {
id
}
}
}"#;
const SET_ITEM_STATE_MUTATION: &str = r#"mutation SetIssueState($id: String!, $stateId: String!) {
issueUpdate(id: $id, input: { stateId: $stateId }) {
issue {
id
}
}
}"#;
const MOVE_ITEM_TO_TEAM_MUTATION: &str = r#"mutation MoveIssueToTeam($id: String!, $teamId: String!) {
issueUpdate(id: $id, input: { teamId: $teamId }) {
issue {
id
identifier
}
}
}"#;
const LIST_COMPLETED_WORKFLOW_STATES_QUERY: &str = r#"query CompletedWorkflowStates($teamId: ID!) {
workflowStates(filter: { team: { id: { eq: $teamId } }, type: { eq: "completed" } }) {
nodes {
id
}
}
}"#;
const LIST_UNSTARTED_WORKFLOW_STATES_QUERY: &str = r#"query UnstartedWorkflowStates($teamId: ID!) {
workflowStates(filter: { team: { id: { eq: $teamId } }, type: { eq: "unstarted" } }) {
nodes {
id
position
}
}
}"#;
const CREATE_COMMENT_MUTATION: &str = r#"mutation CreateComment($issueId: String!, $body: String!) {
commentCreate(input: { issueId: $issueId, body: $body }) {
comment {
id
}
}
}"#;
const VIEWER_QUERY: &str = r#"query Viewer {
viewer {
id
}
}"#;
const CREATE_WEBHOOK_MUTATION: &str = r#"mutation CreateWebhook($url: String!, $secret: String!, $resourceTypes: [String!]!, $teamId: String!) {
webhookCreate(input: { url: $url, secret: $secret, resourceTypes: $resourceTypes, teamId: $teamId }) {
webhook {
id
}
}
}"#;
const UPDATE_COMMENT_MUTATION: &str = r#"mutation UpdateComment($id: String!, $body: String!) {
commentUpdate(id: $id, input: { body: $body }) {
comment {
id
}
}
}"#;
const LINK_ATTACHMENT_MUTATION: &str = r#"mutation LinkAttachment($issueId: String!, $url: String!, $title: String!) {
attachmentLinkURL(issueId: $issueId, url: $url, title: $title) {
attachment {
id
}
}
}"#;
const ISSUE_OBSERVATION_QUERY: &str = r#"query IssueObservation($id: String!, $comments: Int!) {
issue(id: $id) {
updatedAt
title
description
comments(first: $comments, orderBy: createdAt) {
nodes {
id
body
updatedAt
user {
id
}
}
pageInfo { hasNextPage endCursor }
}
}
}"#;
const ISSUE_COMMENTS_QUERY: &str = r#"query IssueComments($id: String!, $comments: Int!, $after: String) {
issue(id: $id) {
comments(first: $comments, after: $after, orderBy: createdAt) {
nodes {
id
body
updatedAt
user { id }
}
pageInfo {
hasNextPage
endCursor
}
}
}
}"#;
const ISSUE_TEAM_QUERY: &str = r#"query IssueTeam($id: String!) {
issue(id: $id) {
team {
id
}
}
}"#;
const UPDATE_ATTACHMENT_MUTATION: &str = r#"mutation UpdateAttachment($id: String!, $title: String!, $subtitle: String!) {
attachmentUpdate(id: $id, input: { title: $title, subtitle: $subtitle }) {
attachment {
id
}
}
}"#;
const OBSERVATION_COMMENT_PAGE: u32 = 50;
#[derive(Debug, Clone)]
pub struct LinearClient {
client: reqwest::Client,
token: String,
team_id: Option<String>,
base_url: String,
}
impl LinearClient {
pub fn new(token: String, team_id: Option<String>) -> Self {
Self {
client: reqwest::Client::new(),
token,
team_id,
base_url: LINEAR_BASE_URL.to_string(),
}
}
#[cfg(test)]
pub(crate) fn with_base_url(token: String, team_id: Option<String>, base_url: String) -> Self {
Self {
client: reqwest::Client::new(),
token,
team_id,
base_url,
}
}
pub async fn ensure_team(
&self,
name: &str,
key: &str,
repository: &str,
) -> PmResult<TeamBinding> {
let requested_key = key.trim().to_ascii_uppercase();
if requested_key.is_empty() {
return Err(PmError::Message(
"team key cannot be empty; pass --team-key <KEY>".to_string(),
));
}
let response: TeamsData = self.graphql(LIST_TEAMS_QUERY, json!({})).await?;
let teams = response.teams.nodes;
if let Some(team) = teams
.iter()
.find(|team| team.key.eq_ignore_ascii_case(&requested_key))
{
if team.name.eq_ignore_ascii_case(name) {
return self.claim_team(team, repository).await;
}
return Err(PmError::Message(format!(
"Linear team key {requested_key:?} already belongs to team {:?} (id {}). \
Pass a different --team-key or rename that team.",
team.name, team.id
)));
}
if let Some(team) = teams
.iter()
.find(|team| team.name.eq_ignore_ascii_case(name))
{
if team.key.eq_ignore_ascii_case(&requested_key) {
return self.claim_team(team, repository).await;
}
return Err(PmError::Message(format!(
"a Linear team named {name:?} already exists with key {:?} (id {}). \
Pass --team-key {} to adopt it, or rename the team.",
team.key, team.id, team.key
)));
}
let description = repository_claim_marker(repository);
let response: TeamCreateData = self
.graphql(
CREATE_TEAM_MUTATION,
json!({
"name": name,
"key": requested_key,
"description": description,
}),
)
.await?;
let team_id = response.team_create.team.id;
let binding = self.validate_team_claim(&team_id, repository).await?;
Ok(TeamBinding {
created: true,
..binding
})
}
async fn claim_team(&self, team: &TeamNode, repository: &str) -> PmResult<TeamBinding> {
match repository_claim(team.description.as_deref().unwrap_or_default())? {
Some(claimed) if claimed == repository => {}
Some(claimed) => {
return Err(PmError::Message(format!(
"Linear team {} ({}, key {}) is claimed by repository {claimed}; \
repository {repository} cannot use it",
team.name, team.id, team.key
)))
}
None => {
let description = description_with_repository_claim(
team.description.as_deref().unwrap_or_default(),
repository,
);
let _: Value = self
.graphql(
UPDATE_TEAM_DESCRIPTION_MUTATION,
json!({ "id": team.id, "description": description }),
)
.await?;
}
}
let mut binding = self.validate_team_claim(&team.id, repository).await?;
binding.created = false;
Ok(binding)
}
pub async fn claim_configured_team(
&self,
team_id: &str,
repository: &str,
expected_name: Option<&str>,
expected_key: Option<&str>,
) -> PmResult<TeamBinding> {
let response: TeamsData = self.graphql(LIST_TEAMS_QUERY, json!({})).await?;
let team = response
.teams
.nodes
.into_iter()
.find(|team| team.id == team_id)
.ok_or_else(|| PmError::Message(format!("no Linear team with id {team_id}")))?;
if let Some(name) = expected_name.filter(|name| !team.name.eq_ignore_ascii_case(name)) {
return Err(PmError::Message(format!(
"repository {repository} is already bound to Linear Team {} ({}, key {}); \
it cannot rebind through --team-name {name:?}. Run repository-wide `lf pm reteam` instead.",
team.name, team.id, team.key
)));
}
if let Some(key) = expected_key.filter(|key| !team.key.eq_ignore_ascii_case(key.trim())) {
return Err(PmError::Message(format!(
"repository {repository} is already bound to Linear Team {} ({}, key {}); \
it cannot rebind through --team-key {key:?}. Run repository-wide `lf pm reteam` instead.",
team.name, team.id, team.key
)));
}
self.claim_team(&team, repository).await
}
pub async fn validate_team_claim(
&self,
team_id: &str,
repository: &str,
) -> PmResult<TeamBinding> {
let response: TeamsData = self.graphql(LIST_TEAMS_QUERY, json!({})).await?;
let team = response
.teams
.nodes
.into_iter()
.find(|team| team.id == team_id)
.ok_or_else(|| PmError::Message(format!("no Linear team with id {team_id}")))?;
match repository_claim(team.description.as_deref().unwrap_or_default())? {
Some(claimed) if claimed == repository => Ok(TeamBinding {
id: team.id,
key: team.key,
created: false,
}),
Some(claimed) => Err(PmError::Message(format!(
"Linear team {} ({}, key {}) is claimed by repository {claimed}; \
configured repository {repository} must choose another Team",
team.name, team.id, team.key
))),
None => Err(PmError::Message(format!(
"Linear team {} ({}, key {}) has no Loopflow repository claim; \
run `lf pm init --team-key {}` to claim it for {repository}",
team.name, team.id, team.key, team.key
))),
}
}
fn require_team_id(&self) -> PmResult<String> {
self.team_id.clone().ok_or_else(|| {
PmError::Message(
"Linear write requires repository `pm.linear_team` in .lf/config.yaml; \
run `lf pm init --wave <wave> --team-key <KEY>`"
.to_string(),
)
})
}
async fn graphql<T>(&self, query: &str, variables: Value) -> PmResult<T>
where
T: DeserializeOwned,
{
let request = GraphqlRequest { query, variables };
for attempt in 0..=RATE_LIMIT_RETRIES {
let response = self
.client
.post(&self.base_url)
.bearer_auth(&self.token)
.json(&request)
.send()
.await
.map_err(|err| PmError::Message(format!("linear request failed: {err}")))?;
if response.status() == StatusCode::TOO_MANY_REQUESTS && attempt < RATE_LIMIT_RETRIES {
let delay = super::retry_after_delay(response.headers());
warn!(
attempt = attempt + 1,
delay_seconds = delay.as_secs(),
"linear rate limited; retrying"
);
sleep(delay).await;
continue;
}
return parse_graphql_response(response).await;
}
Err(PmError::Message(
"linear request failed after retries".to_string(),
))
}
async fn item_team_id(&self, item_id: &str) -> PmResult<String> {
let response: IssueTeamData = self
.graphql(ISSUE_TEAM_QUERY, json!({ "id": item_id }))
.await?;
response
.issue
.map(|issue| issue.team.id)
.ok_or_else(|| PmError::Message(format!("no Linear issue with id {item_id}")))
}
async fn completed_state_id(&self, team_id: &str) -> PmResult<String> {
let response: WorkflowStatesData = self
.graphql(
LIST_COMPLETED_WORKFLOW_STATES_QUERY,
json!({ "teamId": team_id }),
)
.await?;
response
.workflow_states
.nodes
.into_iter()
.next()
.map(|state| state.id)
.ok_or_else(|| {
PmError::Message(format!(
"no completed Linear workflow state found for team {team_id}"
))
})
}
async fn unstarted_state_id(&self, team_id: &str) -> PmResult<Option<String>> {
let response: WorkflowStatesData = self
.graphql(
LIST_UNSTARTED_WORKFLOW_STATES_QUERY,
json!({ "teamId": team_id }),
)
.await?;
Ok(response
.workflow_states
.nodes
.into_iter()
.min_by(|left, right| left.position.total_cmp(&right.position))
.map(|state| state.id))
}
pub async fn create_wave(&self, name: &str, summary: &str) -> PmResult<String> {
let response: InitiativeCreateData = self
.graphql(
CREATE_INITIATIVE_MUTATION,
json!({
"name": name,
"description": linear_description(summary),
}),
)
.await?;
Ok(response.initiative_create.initiative.id)
}
pub async fn rename_wave(&self, initiative_id: &str, name: &str) -> PmResult<()> {
let _: Value = self
.graphql(
UPDATE_INITIATIVE_MUTATION,
json!({
"id": initiative_id,
"name": name,
}),
)
.await?;
Ok(())
}
pub async fn list_waves(&self) -> PmResult<Vec<PmWave>> {
let mut after = None;
let mut waves = Vec::new();
loop {
let response: InitiativesData = self
.graphql(
LIST_INITIATIVES_QUERY,
json!({
"after": after,
"first": LIST_PROJECTS_PAGE_SIZE,
}),
)
.await?;
let page = response.initiatives;
waves.extend(page.nodes.into_iter().map(|initiative| PmWave {
id: initiative.id,
name: initiative.name,
summary: initiative.description.unwrap_or_default(),
}));
if !page.page_info.has_next_page {
return Ok(waves);
}
after = page.page_info.end_cursor;
}
}
pub async fn create_project(
&self,
initiative_id: &str,
name: &str,
content: &ProjectContent,
) -> PmResult<String> {
let team_id = self.require_team_id()?;
let response: ProjectCreateData = self
.graphql(
CREATE_PROJECT_MUTATION,
json!({
"name": name,
"description": project_description(content),
"content": render_project_content(content),
"teamId": team_id,
}),
)
.await?;
let project_id = response.project_create.project.id;
let _: Value = self
.graphql(
ATTACH_PROJECT_MUTATION,
json!({
"initiativeId": initiative_id,
"projectId": project_id,
}),
)
.await?;
Ok(project_id)
}
pub async fn update_project(
&self,
project_id: &str,
name: &str,
content: &ProjectContent,
) -> PmResult<()> {
let _: Value = self
.graphql(
UPDATE_PROJECT_MUTATION,
json!({
"id": project_id,
"name": name,
"description": project_description(content),
"content": render_project_content(content),
}),
)
.await?;
Ok(())
}
pub async fn archive_project(&self, project_id: &str) -> PmResult<()> {
let response: ProjectArchiveData = self
.graphql(
ARCHIVE_PROJECT_MUTATION,
json!({
"id": project_id,
}),
)
.await?;
if !response.project_archive.success {
return Err(PmError::Message(format!(
"Linear did not archive Project {project_id}"
)));
}
Ok(())
}
pub async fn list_projects(&self, initiative_id: &str) -> PmResult<Vec<PmProject>> {
let mut after = None;
let mut projects = Vec::new();
loop {
let response: InitiativeProjectsData = self
.graphql(
LIST_INITIATIVE_PROJECTS_QUERY,
json!({
"initiativeId": initiative_id,
"after": after,
"first": LIST_PROJECTS_PAGE_SIZE,
}),
)
.await?;
let page = response.initiative.projects;
projects.extend(page.nodes.into_iter().map(ProjectNode::into_pm_project));
if !page.page_info.has_next_page {
return Ok(projects);
}
after = page.page_info.end_cursor;
}
}
pub async fn list_items(&self, project_id: &str) -> PmResult<Vec<PmItem>> {
self.list_issue_nodes(project_id)
.await?
.into_iter()
.enumerate()
.map(|(rank, issue)| issue.into_pm_item(rank as u32))
.collect()
}
async fn list_issue_nodes(&self, project_id: &str) -> PmResult<Vec<IssueNode>> {
let mut after = None;
let mut issues = Vec::new();
loop {
let response: ProjectIssuesData = self
.graphql(
LIST_ITEMS_QUERY,
json!({
"projectId": project_id,
"after": after,
"first": LIST_ITEMS_PAGE_SIZE,
}),
)
.await?;
let page = response.project.issues;
issues.extend(page.nodes);
if !page.page_info.has_next_page {
issues.sort_by(|left, right| {
left.priority_sort_order
.total_cmp(&right.priority_sort_order)
.then_with(|| left.sort_order.total_cmp(&right.sort_order))
});
return Ok(issues);
}
after = page.page_info.end_cursor;
}
}
pub async fn create_item(&self, project_id: &str, item: &PmItemCreate) -> PmResult<String> {
let team_id = self.require_team_id()?;
let state_id = self.unstarted_state_id(&team_id).await?;
let response: IssueCreateData = self
.graphql(
CREATE_ITEM_MUTATION,
json!({
"teamId": team_id,
"projectId": project_id,
"title": item.name,
"description": item.description,
"stateId": state_id,
}),
)
.await?;
Ok(response.issue_create.issue.id)
}
pub async fn update_item(&self, item_id: &str, update: &PmItemUpdate) -> PmResult<()> {
let Some(update) = update.text_update() else {
return Ok(());
};
let mut input = serde_json::Map::new();
if let Some(name) = update.name {
input.insert("title".to_string(), json!(name));
}
if let Some(description) = update.description {
input.insert("description".to_string(), json!(description));
}
let _: Value = self
.graphql(
UPDATE_ITEM_MUTATION,
json!({
"id": item_id,
"input": input,
}),
)
.await?;
Ok(())
}
pub async fn move_item_to_project(&self, item_id: &str, project_id: &str) -> PmResult<()> {
let _: Value = self
.graphql(
MOVE_ITEM_MUTATION,
json!({
"id": item_id,
"projectId": project_id,
}),
)
.await?;
Ok(())
}
pub async fn move_item_to_team(&self, item_id: &str, team_id: &str) -> PmResult<String> {
let response: IssueUpdateIdentifierData = self
.graphql(
MOVE_ITEM_TO_TEAM_MUTATION,
json!({
"id": item_id,
"teamId": team_id,
}),
)
.await?;
Ok(response.issue_update.issue.identifier)
}
pub async fn move_project_to_team(&self, project_id: &str, team_id: &str) -> PmResult<()> {
let _: Value = self
.graphql(
MOVE_PROJECT_TO_TEAM_MUTATION,
json!({
"id": project_id,
"teamId": team_id,
}),
)
.await?;
Ok(())
}
pub async fn set_project_teams(&self, project_id: &str, team_ids: &[String]) -> PmResult<()> {
let _: Value = self
.graphql(
SET_PROJECT_TEAMS_MUTATION,
json!({ "id": project_id, "teamIds": team_ids }),
)
.await?;
Ok(())
}
pub async fn issue_ownership(&self, issue_id: &str) -> PmResult<(PmItem, PmProject)> {
let response: IssueOwnershipData = self
.graphql(ISSUE_OWNERSHIP_QUERY, json!({ "id": issue_id }))
.await?;
let issue = response.issue.ok_or_else(|| {
PmError::Message(format!("no Linear issue with id or identifier {issue_id}"))
})?;
issue.into_ownership()
}
pub async fn project_ownership(&self, project_id: &str) -> PmResult<PmProject> {
let response: ProjectOwnershipData = self
.graphql(PROJECT_OWNERSHIP_QUERY, json!({ "id": project_id }))
.await?;
response
.project
.map(ProjectNode::into_pm_project)
.ok_or_else(|| PmError::Message(format!("no Linear Project with id {project_id}")))
}
pub async fn complete_item(&self, item_id: &str) -> PmResult<()> {
let team_id = self.item_team_id(item_id).await?;
let state_id = self.completed_state_id(&team_id).await?;
let _: Value = self
.graphql(
SET_ITEM_STATE_MUTATION,
json!({
"id": item_id,
"stateId": state_id,
}),
)
.await?;
Ok(())
}
pub async fn reopen_item(&self, item_id: &str) -> PmResult<()> {
let team_id = self.item_team_id(item_id).await?;
let Some(state_id) = self.unstarted_state_id(&team_id).await? else {
return Err(PmError::Message(format!(
"no active Linear workflow state found to reopen issue {item_id}"
)));
};
let _: Value = self
.graphql(
SET_ITEM_STATE_MUTATION,
json!({
"id": item_id,
"stateId": state_id,
}),
)
.await?;
Ok(())
}
pub async fn comment(&self, item_id: &str, body: &str) -> PmResult<String> {
let response: CommentData = self
.graphql(
CREATE_COMMENT_MUTATION,
json!({
"issueId": item_id,
"body": body,
}),
)
.await?;
Ok(response.comment_create.comment.id)
}
pub async fn find_comment_with_marker(
&self,
issue_id: &str,
marker: &str,
) -> PmResult<Option<String>> {
let mut after = None;
loop {
let response: IssueCommentsData = self
.graphql(
ISSUE_COMMENTS_QUERY,
json!({
"id": issue_id,
"comments": OBSERVATION_COMMENT_PAGE,
"after": after,
}),
)
.await?;
let issue = response
.issue
.ok_or_else(|| PmError::Message(format!("linear issue {issue_id} not found")))?;
if let Some(comment) = issue
.comments
.nodes
.into_iter()
.find(|comment| comment.body.contains(marker))
{
return Ok(Some(comment.id));
}
if !issue.comments.page_info.has_next_page {
return Ok(None);
}
after = issue.comments.page_info.end_cursor;
if after.is_none() {
return Err(PmError::Message(format!(
"Linear comments for issue {issue_id} have another page without a cursor"
)));
}
}
}
pub async fn update_comment(&self, comment_id: &str, body: &str) -> PmResult<()> {
let _: Value = self
.graphql(
UPDATE_COMMENT_MUTATION,
json!({
"id": comment_id,
"body": body,
}),
)
.await?;
Ok(())
}
pub async fn link_attachment(
&self,
issue_id: &str,
url: &str,
title: &str,
) -> PmResult<String> {
let response: AttachmentLinkData = self
.graphql(
LINK_ATTACHMENT_MUTATION,
json!({
"issueId": issue_id,
"url": url,
"title": title,
}),
)
.await?;
Ok(response.attachment_link_url.attachment.id)
}
pub async fn update_attachment(
&self,
attachment_id: &str,
title: &str,
subtitle: &str,
) -> PmResult<()> {
let _: Value = self
.graphql(
UPDATE_ATTACHMENT_MUTATION,
json!({
"id": attachment_id,
"title": title,
"subtitle": subtitle,
}),
)
.await?;
Ok(())
}
pub async fn viewer_id(&self) -> PmResult<String> {
let response: ViewerData = self.graphql(VIEWER_QUERY, json!({})).await?;
Ok(response.viewer.id)
}
pub async fn create_webhook(&self, url: &str, secret: &str) -> PmResult<String> {
let team_id = self.require_team_id()?;
let response: WebhookCreateData = self
.graphql(
CREATE_WEBHOOK_MUTATION,
json!({
"url": url,
"secret": secret,
"resourceTypes": ["Issue", "Comment"],
"teamId": team_id,
}),
)
.await?;
Ok(response.webhook_create.webhook.id)
}
pub async fn observe_issue(&self, issue_id: &str) -> PmResult<IssueObservation> {
let response: IssueObservationData = self
.graphql(
ISSUE_OBSERVATION_QUERY,
json!({ "id": issue_id, "comments": OBSERVATION_COMMENT_PAGE }),
)
.await?;
let issue = response
.issue
.ok_or_else(|| PmError::Message(format!("linear issue {issue_id} not found")))?;
let mut page = issue.comments;
let mut comments = Vec::new();
loop {
comments.extend(page.nodes.into_iter().map(|node| IssueComment {
id: node.id,
revision: node.updated_at,
body: node.body,
author_id: node.user.map(|user| user.id),
}));
if !page.page_info.has_next_page {
break;
}
let after = page.page_info.end_cursor.ok_or_else(|| {
PmError::Message("Linear comment page is missing its continuation cursor".into())
})?;
let response: IssueCommentsData = self
.graphql(
ISSUE_COMMENTS_QUERY,
json!({"id": issue_id, "comments": OBSERVATION_COMMENT_PAGE, "after": after}),
)
.await?;
page = response
.issue
.ok_or_else(|| PmError::Message(format!("linear issue {issue_id} not found")))?
.comments;
}
comments.sort_by(|left, right| {
left.revision
.cmp(&right.revision)
.then(left.id.cmp(&right.id))
});
Ok(IssueObservation {
revision: issue.updated_at,
title: issue.title,
description: issue.description.unwrap_or_default(),
comments,
})
}
}
#[derive(Serialize)]
struct GraphqlRequest<'a> {
query: &'a str,
variables: Value,
}
#[derive(Deserialize)]
struct GraphqlResponse {
#[serde(default)]
data: Option<Value>,
#[serde(default)]
errors: Vec<GraphqlError>,
}
#[derive(Debug, Deserialize)]
struct GraphqlError {
message: String,
#[serde(default)]
extensions: Option<GraphqlErrorExtensions>,
}
impl GraphqlError {
fn display_message(&self) -> &str {
self.extensions
.as_ref()
.and_then(|extensions| extensions.user_presentable_message.as_deref())
.filter(|message| !message.trim().is_empty())
.unwrap_or(&self.message)
}
}
#[derive(Debug, Deserialize)]
#[serde(rename_all = "camelCase")]
struct GraphqlErrorExtensions {
#[serde(default)]
user_presentable_message: Option<String>,
}
#[derive(Deserialize)]
struct ProjectCreateData {
#[serde(rename = "projectCreate")]
project_create: ProjectPayload,
}
#[derive(Deserialize)]
struct ProjectArchiveData {
#[serde(rename = "projectArchive")]
project_archive: SuccessPayload,
}
#[derive(Deserialize)]
struct SuccessPayload {
success: bool,
}
#[derive(Deserialize)]
struct InitiativeCreateData {
#[serde(rename = "initiativeCreate")]
initiative_create: InitiativePayload,
}
#[derive(Deserialize)]
struct IssueCreateData {
#[serde(rename = "issueCreate")]
issue_create: IssuePayload,
}
#[derive(Deserialize)]
struct IssueUpdateIdentifierData {
#[serde(rename = "issueUpdate")]
issue_update: IssueIdentifierPayload,
}
#[derive(Deserialize)]
struct IssueIdentifierPayload {
issue: IssueIdentifierNode,
}
#[derive(Deserialize)]
struct IssueIdentifierNode {
identifier: String,
}
#[derive(Deserialize)]
struct ProjectPayload {
project: IdNode,
}
#[derive(Deserialize)]
struct InitiativePayload {
initiative: IdNode,
}
#[derive(Deserialize)]
struct IssuePayload {
issue: IdNode,
}
#[derive(Deserialize)]
struct CommentData {
#[serde(rename = "commentCreate")]
comment_create: CommentPayload,
}
#[derive(Deserialize)]
struct CommentPayload {
comment: IdNode,
}
#[derive(Deserialize)]
struct AttachmentLinkData {
#[serde(rename = "attachmentLinkURL")]
attachment_link_url: AttachmentPayload,
}
#[derive(Deserialize)]
struct AttachmentPayload {
attachment: IdNode,
}
#[derive(Deserialize)]
struct IdNode {
id: String,
}
#[derive(Deserialize)]
struct ViewerData {
viewer: IdNode,
}
#[derive(Deserialize)]
struct WebhookCreateData {
#[serde(rename = "webhookCreate")]
webhook_create: WebhookCreateNode,
}
#[derive(Deserialize)]
struct WebhookCreateNode {
webhook: IdNode,
}
#[derive(Deserialize)]
struct IssueTeamData {
issue: Option<IssueTeamNode>,
}
#[derive(Deserialize)]
struct IssueTeamNode {
team: IdNode,
}
#[derive(Deserialize)]
struct IssueObservationData {
issue: Option<IssueObservationNode>,
}
#[derive(Deserialize)]
struct IssueObservationNode {
#[serde(rename = "updatedAt")]
updated_at: String,
#[serde(default)]
title: String,
#[serde(default)]
description: Option<String>,
comments: PagedCommentConnection,
}
#[derive(Deserialize)]
struct IssueCommentsData {
issue: Option<IssueCommentsNode>,
}
#[derive(Deserialize)]
struct IssueCommentsNode {
comments: PagedCommentConnection,
}
#[derive(Deserialize)]
struct PagedCommentConnection {
nodes: Vec<CommentNode>,
#[serde(rename = "pageInfo")]
page_info: PageInfo,
}
#[derive(Deserialize)]
struct CommentNode {
id: String,
#[serde(rename = "updatedAt")]
updated_at: Option<String>,
#[serde(default)]
body: String,
#[serde(default)]
user: Option<IdNode>,
}
#[derive(Deserialize)]
struct IssueNode {
id: String,
#[serde(default)]
identifier: String,
url: Option<String>,
#[serde(default)]
title: String,
#[serde(default)]
description: Option<String>,
#[serde(rename = "prioritySortOrder", default)]
priority_sort_order: f64,
#[serde(rename = "sortOrder", default)]
sort_order: f64,
#[serde(default)]
assignee: Option<IdNode>,
#[serde(default)]
state: Option<WorkflowStateRef>,
#[serde(default)]
project: Option<ProjectRef>,
#[serde(default)]
team: Option<IdNode>,
}
impl IssueNode {
fn into_pm_item(self, rank: u32) -> PmResult<PmItem> {
let completed = self
.state
.as_ref()
.is_some_and(|state| state.r#type.eq_ignore_ascii_case(COMPLETED_STATE_TYPE));
let identifier = if self.identifier.is_empty() {
self.id.clone()
} else {
self.identifier
};
let project = self
.project
.ok_or_else(|| PmError::Message(format!("Linear issue {identifier} has no Project")))?;
let team = self
.team
.ok_or_else(|| PmError::Message(format!("Linear issue {identifier} has no Team")))?;
Ok(PmItem {
id: self.id,
identifier,
url: self.url,
name: self.title,
description: self.description.unwrap_or_default(),
rank,
completed,
project_id: project.id,
project: project_slug(&project.name),
team_id: team.id,
assignee: self.assignee.map(|assignee| assignee.id),
})
}
}
#[derive(Deserialize)]
struct ProjectRef {
id: String,
name: String,
}
#[derive(Deserialize)]
struct IssueOwnershipData {
issue: Option<OwnedIssueNode>,
}
#[derive(Deserialize)]
struct OwnedIssueNode {
id: String,
#[serde(default)]
identifier: String,
url: Option<String>,
#[serde(default)]
title: String,
#[serde(default)]
description: Option<String>,
#[serde(rename = "prioritySortOrder", default)]
priority_sort_order: f64,
#[serde(rename = "sortOrder", default)]
sort_order: f64,
#[serde(default)]
assignee: Option<IdNode>,
#[serde(default)]
state: Option<WorkflowStateRef>,
#[serde(default)]
team: Option<IdNode>,
#[serde(default)]
project: Option<ProjectNode>,
}
impl OwnedIssueNode {
fn into_ownership(self) -> PmResult<(PmItem, PmProject)> {
let project = self.project.ok_or_else(|| {
PmError::Message(format!("Linear issue {} has no Project", self.identifier))
})?;
let item = IssueNode {
id: self.id,
identifier: self.identifier,
url: self.url,
title: self.title,
description: self.description,
priority_sort_order: self.priority_sort_order,
sort_order: self.sort_order,
assignee: self.assignee,
state: self.state,
project: Some(ProjectRef {
id: project.id.clone(),
name: project.name.clone(),
}),
team: self.team,
}
.into_pm_item(0)?;
Ok((item, project.into_pm_project()))
}
}
#[derive(Deserialize)]
struct WorkflowStateRef {
#[serde(rename = "type")]
r#type: String,
}
#[derive(Deserialize)]
struct ProjectIssuesData {
project: ProjectWithIssues,
}
#[derive(Deserialize)]
struct ProjectWithIssues {
issues: IssuesConnection,
}
#[derive(Deserialize)]
struct IssuesConnection {
nodes: Vec<IssueNode>,
#[serde(rename = "pageInfo")]
page_info: PageInfo,
}
#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
struct PageInfo {
has_next_page: bool,
end_cursor: Option<String>,
}
#[derive(Deserialize)]
struct WorkflowStatesData {
#[serde(rename = "workflowStates")]
workflow_states: WorkflowStatesConnection,
}
#[derive(Deserialize)]
struct WorkflowStatesConnection {
nodes: Vec<WorkflowStateNode>,
}
#[derive(Deserialize)]
struct WorkflowStateNode {
id: String,
#[serde(default)]
position: f64,
}
#[derive(Deserialize)]
struct TeamsData {
teams: TeamsConnection,
}
#[derive(Deserialize)]
struct TeamsConnection {
nodes: Vec<TeamNode>,
}
#[derive(Deserialize)]
struct TeamNode {
id: String,
name: String,
key: String,
#[serde(default)]
description: Option<String>,
}
#[derive(Deserialize)]
struct TeamCreateData {
#[serde(rename = "teamCreate")]
team_create: TeamCreatePayload,
}
#[derive(Deserialize)]
struct TeamCreatePayload {
team: IdNode,
}
#[derive(Deserialize)]
struct ProjectNode {
id: String,
name: String,
#[serde(default)]
description: Option<String>,
#[serde(default)]
content: Option<String>,
#[serde(default)]
initiatives: IdConnection,
#[serde(default)]
teams: IdConnection,
}
impl ProjectNode {
fn into_pm_project(self) -> PmProject {
let content = parse_project_content(self.content.as_deref().unwrap_or_default());
PmProject {
id: self.id,
slug: project_slug(&self.name),
name: self.name,
summary: self.description.unwrap_or_default(),
definition: content.definition,
flows: Some(content.flows),
krs: content.krs,
initiative_ids: self
.initiatives
.nodes
.into_iter()
.map(|initiative| initiative.id)
.collect(),
team_ids: self.teams.nodes.into_iter().map(|team| team.id).collect(),
}
}
}
#[derive(Deserialize)]
struct ProjectOwnershipData {
project: Option<ProjectNode>,
}
#[derive(Default, Deserialize)]
struct IdConnection {
#[serde(default)]
nodes: Vec<IdNode>,
}
#[derive(Deserialize)]
struct InitiativesData {
initiatives: InitiativesConnection,
}
#[derive(Deserialize)]
struct InitiativesConnection {
nodes: Vec<InitiativeNode>,
#[serde(rename = "pageInfo")]
page_info: PageInfo,
}
#[derive(Deserialize)]
struct InitiativeNode {
id: String,
name: String,
#[serde(default)]
description: Option<String>,
}
#[derive(Deserialize)]
struct InitiativeProjectsData {
initiative: InitiativeWithProjects,
}
#[derive(Deserialize)]
struct InitiativeWithProjects {
projects: ProjectsConnection,
}
#[derive(Deserialize)]
struct ProjectsConnection {
nodes: Vec<ProjectNode>,
#[serde(rename = "pageInfo")]
page_info: PageInfo,
}
async fn parse_graphql_response<T: DeserializeOwned>(response: reqwest::Response) -> PmResult<T> {
let status = response.status();
let body = response
.bytes()
.await
.map_err(|err| PmError::Message(format!("failed to read Linear response: {err}")))?;
let parsed = serde_json::from_slice::<GraphqlResponse>(&body)
.map_err(|err| PmError::Message(format!("failed to decode Linear response: {err}")))?;
if let Some(error) = parsed.errors.first() {
if status.is_success() {
return Err(PmError::Message(error.display_message().to_string()));
}
return Err(PmError::Message(format!(
"linear request failed with status {status}: {}",
error.display_message()
)));
}
if !status.is_success() {
let body_text = String::from_utf8_lossy(&body).trim().to_string();
if body_text.is_empty() {
return Err(PmError::Message(format!(
"linear request failed with status {status}"
)));
}
return Err(PmError::Message(format!(
"linear request failed with status {status}: {body_text}"
)));
}
let data = parsed
.data
.ok_or_else(|| PmError::Message("linear response missing data".to_string()))?;
serde_json::from_value(data)
.map_err(|err| PmError::Message(format!("failed to decode Linear response: {err}")))
}
fn repository_claim_marker(repository: &str) -> String {
format!("{REPOSITORY_CLAIM_PREFIX} {repository} -->")
}
fn repository_claim(description: &str) -> PmResult<Option<String>> {
let mut claims = Vec::new();
for line in description.lines() {
if !line.contains(REPOSITORY_CLAIM_PREFIX) {
continue;
}
let trimmed = line.trim();
let Some(value) = trimmed
.strip_prefix(REPOSITORY_CLAIM_PREFIX)
.and_then(|value| value.strip_suffix("-->"))
.map(str::trim)
.filter(|value| !value.is_empty())
else {
return Err(PmError::Message(format!(
"Linear Team has a malformed Loopflow repository marker: {trimmed:?}"
)));
};
claims.push(value.to_string());
}
match claims.as_slice() {
[] => Ok(None),
[claim] => Ok(Some(claim.clone())),
_ => Err(PmError::Message(format!(
"Linear Team has multiple Loopflow repository markers: {}",
claims.join(", ")
))),
}
}
fn description_with_repository_claim(description: &str, repository: &str) -> String {
let human = description.trim_end();
let marker = repository_claim_marker(repository);
if human.is_empty() {
marker
} else {
format!("{human}\n\n{marker}")
}
}
fn linear_description(description: &str) -> String {
let summary = first_meaningful_paragraph(description);
if summary.is_empty() {
return String::new();
}
const MAX_DESCRIPTION_LEN: usize = 255;
summary.chars().take(MAX_DESCRIPTION_LEN).collect()
}
fn project_description(content: &ProjectContent) -> String {
let summary = content
.definition
.split("\n\n")
.map(|paragraph| paragraph.split_whitespace().collect::<Vec<_>>().join(" "))
.find(|paragraph| !paragraph.is_empty())
.unwrap_or_default();
linear_description(&summary)
}
fn first_meaningful_paragraph(description: &str) -> String {
let mut lines = description.lines().peekable();
while let Some(line) = lines.next() {
let trimmed = line.trim();
if trimmed.is_empty() || trimmed.starts_with('#') {
continue;
}
let mut paragraph = vec![trimmed];
while let Some(next_line) = lines.peek() {
let trimmed = next_line.trim();
if trimmed.is_empty() {
break;
}
if trimmed.starts_with('#') {
lines.next();
break;
}
paragraph.push(trimmed);
lines.next();
}
return paragraph.join(" ");
}
String::new()
}
#[cfg(test)]
mod tests {
use super::*;
use crate::pm::test_server::{self, json_response};
use axum::http::StatusCode;
use serde_json::json;
#[test]
fn issue_mutations_use_linear_string_ids() {
for query in [
UPDATE_ITEM_MUTATION,
MOVE_ITEM_MUTATION,
SET_ITEM_STATE_MUTATION,
CREATE_COMMENT_MUTATION,
] {
assert!(!query.contains(": ID!"));
}
assert!(MOVE_ITEM_MUTATION.contains("$projectId: String!"));
assert!(SET_ITEM_STATE_MUTATION.contains("$stateId: String!"));
assert!(CREATE_COMMENT_MUTATION.contains("$issueId: String!"));
}
#[test]
fn attachment_link_omits_subtitle_and_update_keeps_it() {
assert!(
!LINK_ATTACHMENT_MUTATION.contains("subtitle"),
"attachmentLinkURL rejects subtitle; it must not appear in the create mutation"
);
assert!(
UPDATE_ATTACHMENT_MUTATION.contains("subtitle: $subtitle"),
"attachmentUpdate carries PR state as its input subtitle"
);
}
#[test]
fn workflow_state_filters_use_linear_team_id() {
assert!(CREATE_ITEM_MUTATION.contains("$teamId: String!"));
assert!(LIST_COMPLETED_WORKFLOW_STATES_QUERY.contains("$teamId: ID!"));
assert!(LIST_UNSTARTED_WORKFLOW_STATES_QUERY.contains("$teamId: ID!"));
}
#[tokio::test]
async fn create_webhook_registers_issue_and_comment_resources() {
let (base_url, requests) = test_server::spawn(vec![json_response(
StatusCode::OK,
json!({ "data": { "webhookCreate": { "webhook": { "id": "wh-1" } } } }),
)])
.await;
let client = LinearClient::with_base_url(
"linear-secret".to_string(),
Some("team-loo".to_string()),
base_url,
);
let id = client
.create_webhook("https://loopflow.example/linear/webhook", "whsec")
.await
.expect("create webhook");
assert_eq!(id, "wh-1");
let requests = requests.lock().await;
let body: Value = serde_json::from_str(&requests[0].body).expect("body is json");
assert_eq!(
body["variables"]["url"],
"https://loopflow.example/linear/webhook"
);
assert_eq!(
body["variables"]["resourceTypes"],
json!(["Issue", "Comment"])
);
assert_eq!(body["variables"]["teamId"], "team-loo");
assert!(!CREATE_WEBHOOK_MUTATION.contains("allPublicTeams"));
assert!(CREATE_WEBHOOK_MUTATION.contains("$url: String!"));
assert!(!CREATE_WEBHOOK_MUTATION.contains(": ID!"));
}
#[tokio::test]
async fn viewer_id_reads_loopflows_own_user() {
let (base_url, _requests) = test_server::spawn(vec![json_response(
StatusCode::OK,
json!({ "data": { "viewer": { "id": "user-loopflow" } } }),
)])
.await;
let client = LinearClient::with_base_url("linear-secret".to_string(), None, base_url);
assert_eq!(
client.viewer_id().await.expect("viewer id"),
"user-loopflow"
);
}
#[tokio::test]
async fn observe_issue_reads_revision_content_and_comment_authors() {
let (base_url, requests) = test_server::spawn(vec![json_response(
StatusCode::OK,
json!({
"data": {
"issue": {
"updatedAt": "2026-07-15T18:00:00.000Z",
"title": "Stream Linear edits",
"description": "New body",
"comments": {
"pageInfo": {"hasNextPage": false, "endCursor": null},
"nodes": [
{ "id": "c-1", "body": "please prioritize", "user": { "id": "user-human" } },
{ "id": "c-2", "body": "PR: https://x", "user": { "id": "user-loopflow" } },
{ "id": "c-3", "body": "integration note", "user": null }
]
}
}
}
}),
)])
.await;
let client = LinearClient::with_base_url("linear-secret".to_string(), None, base_url);
let observation = client
.observe_issue("issue-1")
.await
.expect("observe issue");
assert_eq!(observation.revision, "2026-07-15T18:00:00.000Z");
assert_eq!(observation.title, "Stream Linear edits");
assert_eq!(observation.description, "New body");
assert_eq!(
observation.comments,
vec![
IssueComment {
id: "c-1".to_string(),
revision: None,
body: "please prioritize".to_string(),
author_id: Some("user-human".to_string()),
},
IssueComment {
id: "c-2".to_string(),
revision: None,
body: "PR: https://x".to_string(),
author_id: Some("user-loopflow".to_string()),
},
IssueComment {
id: "c-3".to_string(),
revision: None,
body: "integration note".to_string(),
author_id: None,
},
]
);
let requests = requests.lock().await;
let body: Value = serde_json::from_str(&requests[0].body).expect("body is json");
assert_eq!(body["variables"]["id"], "issue-1");
assert_eq!(body["variables"]["comments"], OBSERVATION_COMMENT_PAGE);
}
#[tokio::test]
async fn observe_issue_reads_every_comment_page_in_revision_order() {
let comment = |id, revision| json!({"id": id, "body": "advice", "updatedAt": revision, "user": {"id": "me"}});
let (url, _) = test_server::spawn(vec![
json_response(StatusCode::OK, json!({"data": {"issue": {
"updatedAt": "2026-09-23T00:00:00Z", "title": "Task", "description": "",
"comments": {"nodes": [comment("newer", "2026-09-23T00:00:00Z")], "pageInfo": {"hasNextPage": true, "endCursor": "cursor-1"}}
}}})),
json_response(StatusCode::OK, json!({"data": {"issue": {"comments": {
"nodes": [comment("older", "2026-09-22T00:00:00Z")], "pageInfo": {"hasNextPage": false, "endCursor": null}
}}}})),
]).await;
let client = LinearClient::with_base_url("fixture-token".into(), None, url);
let observed = client.observe_issue("issue-1").await.unwrap();
assert_eq!(
observed
.comments
.iter()
.map(|comment| comment.id.as_str())
.collect::<Vec<_>>(),
["older", "newer"]
);
assert_eq!(
observed.comments[1].revision.as_deref(),
Some("2026-09-23T00:00:00Z")
);
}
#[tokio::test]
async fn observe_issue_reports_a_missing_issue() {
let (base_url, _requests) = test_server::spawn(vec![json_response(
StatusCode::OK,
json!({ "data": { "issue": null } }),
)])
.await;
let client = LinearClient::with_base_url("linear-secret".to_string(), None, base_url);
let error = client
.observe_issue("issue-missing")
.await
.expect_err("missing issue errors");
assert!(error.to_string().contains("issue-missing"));
}
#[tokio::test]
async fn list_items_maps_linear_project_issues() {
let (base_url, requests) = test_server::spawn(vec![json_response(
StatusCode::OK,
json!({
"data": {
"project": {
"issues": {
"nodes": [
{
"id": "issue-1",
"identifier": "LOO-1",
"url": "https://linear.app/loopflow/issue/INF-1/first",
"title": "First",
"description": "one",
"prioritySortOrder": 10.0,
"sortOrder": 10.0,
"assignee": { "id": "user-1" },
"state": { "type": "unstarted" },
"project": { "id": "project-123", "name": "Scan" },
"team": { "id": "team-9" }
},
{
"id": "issue-2",
"identifier": "LOO-2",
"title": "Second",
"description": "two",
"prioritySortOrder": 0.0,
"sortOrder": 0.0,
"state": { "type": "completed" },
"project": { "id": "project-123", "name": "Scan" },
"team": { "id": "team-9" }
}
],
"pageInfo": {
"hasNextPage": false,
"endCursor": null
}
}
}
}
}),
)])
.await;
let client = LinearClient::with_base_url(
"linear-secret".to_string(),
Some("team-9".to_string()),
base_url,
);
let items = client
.list_items("project-123")
.await
.expect("list items succeeds");
assert_eq!(items.len(), 2);
assert_eq!(items[0].id, "issue-2");
assert!(items[0].completed);
assert_eq!(items[1].assignee.as_deref(), Some("user-1"));
assert_eq!(
items[1].url.as_deref(),
Some("https://linear.app/loopflow/issue/INF-1/first")
);
let requests = requests.lock().await;
assert_eq!(requests.len(), 1);
assert_eq!(
requests[0].authorization.as_deref(),
Some("Bearer linear-secret")
);
let request: Value = serde_json::from_str(&requests[0].body).expect("request body is json");
assert!(request["query"]
.as_str()
.expect("query string")
.contains("$projectId: String!"));
assert!(request["query"]
.as_str()
.expect("query string")
.contains("\n url\n"));
}
#[tokio::test]
async fn list_projects_resolves_owning_teams() {
let (base_url, requests) = test_server::spawn(vec![json_response(
StatusCode::OK,
json!({ "data": { "initiative": { "projects": {
"nodes": [{
"id": "project-1",
"name": "Unified Practice Targets",
"description": "",
"content": "## Definition\n\nA bet.\n\n## KRs\n",
"initiatives": { "nodes": [{ "id": "initiative-1" }] },
"teams": { "nodes": [{ "id": "team-cadenza" }] }
}],
"pageInfo": { "hasNextPage": false, "endCursor": null }
} } } }),
)])
.await;
let client = LinearClient::with_base_url("linear-secret".to_string(), None, base_url);
let projects = client
.list_projects("initiative-1")
.await
.expect("list projects");
assert_eq!(projects.len(), 1);
assert_eq!(projects[0].team_ids, ["team-cadenza".to_string()]);
let request: Value =
serde_json::from_str(&requests.lock().await[0].body).expect("query json");
assert!(request["query"]
.as_str()
.expect("query string")
.contains("teams(first: 50)"));
}
#[tokio::test]
async fn create_project_writes_content_then_attaches_to_initiative() {
let (base_url, requests) = test_server::spawn(vec![
json_response(
StatusCode::OK,
json!({ "data": { "projectCreate": { "project": { "id": "project-1" } } } }),
),
json_response(
StatusCode::OK,
json!({ "data": { "initiativeToProjectCreate": { "initiativeToProject": { "id": "link-1" } } } }),
),
])
.await;
let client = LinearClient::with_base_url(
"linear-secret".to_string(),
Some("team-9".to_string()),
base_url,
);
let project_id = client
.create_project(
"initiative-1",
"Wave Chat",
&ProjectContent {
definition: "Conversation stays in flow.".to_string(),
flows: crate::pm::ProjectFlowPlan::empty(),
krs: vec![PmKr {
text: "Replies stream".to_string(),
holds: false,
}],
},
)
.await
.expect("create project");
assert_eq!(project_id, "project-1");
let requests = requests.lock().await;
let create: Value = serde_json::from_str(&requests[0].body).expect("create json");
assert_eq!(create["variables"]["name"], "Wave Chat");
assert!(create["variables"]["content"]
.as_str()
.expect("content")
.contains("- [ ] Replies stream"));
let attach: Value = serde_json::from_str(&requests[1].body).expect("attach json");
assert_eq!(attach["variables"]["initiativeId"], "initiative-1");
assert_eq!(attach["variables"]["projectId"], "project-1");
}
#[tokio::test]
async fn update_project_replaces_definition_and_krs() {
let (base_url, requests) = test_server::spawn(vec![json_response(
StatusCode::OK,
json!({ "data": { "projectUpdate": { "project": { "id": "project-1" } } } }),
)])
.await;
let client = LinearClient::with_base_url(
"linear-secret".to_string(),
Some("team-9".to_string()),
base_url,
);
client
.update_project(
"project-1",
"Wave Chat",
&ProjectContent {
definition: "Conversation stays in flow.".to_string(),
flows: crate::pm::ProjectFlowPlan {
recommended: Some("task-design".to_string()),
},
krs: vec![PmKr {
text: "Replies survive every restart boundary".to_string(),
holds: false,
}],
},
)
.await
.expect("update project");
let requests = requests.lock().await;
let update: Value = serde_json::from_str(&requests[0].body).expect("update json");
assert_eq!(update["variables"]["id"], "project-1");
assert!(update["variables"]["content"]
.as_str()
.expect("content")
.contains("Replies survive every restart boundary"));
assert!(update["variables"]["content"]
.as_str()
.expect("content")
.contains("recommended: task-design"));
}
#[tokio::test]
async fn archive_project_uses_linear_archive_mutation() {
let (base_url, requests) = test_server::spawn(vec![json_response(
StatusCode::OK,
json!({ "data": { "projectArchive": { "success": true } } }),
)])
.await;
let client = LinearClient::with_base_url(
"linear-secret".to_string(),
Some("team-9".to_string()),
base_url,
);
client
.archive_project("project-1")
.await
.expect("archive project");
let requests = requests.lock().await;
let archive: Value = serde_json::from_str(&requests[0].body).expect("archive json");
assert!(archive["query"]
.as_str()
.expect("query")
.contains("projectArchive"));
assert_eq!(archive["variables"]["id"], "project-1");
}
#[tokio::test]
async fn move_item_to_team_returns_the_reassigned_identifier() {
let (base_url, requests) = test_server::spawn(vec![json_response(
StatusCode::OK,
json!({ "data": { "issueUpdate": { "issue": { "id": "issue-9", "identifier": "PRD-4" } } } }),
)])
.await;
let client = LinearClient::with_base_url(
"linear-secret".to_string(),
Some("team-old".to_string()),
base_url,
);
let new_identifier = client
.move_item_to_team("issue-9", "team-prd")
.await
.expect("move succeeds");
assert_eq!(new_identifier, "PRD-4");
let requests = requests.lock().await;
let body: Value = serde_json::from_str(&requests[0].body).expect("move json");
assert!(body["query"]
.as_str()
.expect("query")
.contains("issueUpdate"));
assert_eq!(body["variables"]["id"], "issue-9");
assert_eq!(body["variables"]["teamId"], "team-prd");
}
#[tokio::test]
async fn move_project_to_team_sets_the_team_ids() {
let (base_url, requests) = test_server::spawn(vec![json_response(
StatusCode::OK,
json!({ "data": { "projectUpdate": { "project": { "id": "project-1" } } } }),
)])
.await;
let client = LinearClient::with_base_url(
"linear-secret".to_string(),
Some("team-old".to_string()),
base_url,
);
client
.move_project_to_team("project-1", "team-cadenza")
.await
.expect("move project succeeds");
let requests = requests.lock().await;
let body: Value = serde_json::from_str(&requests[0].body).expect("move json");
let query = body["query"].as_str().expect("query");
assert!(query.contains("projectUpdate"));
assert!(query.contains("teamIds: [$teamId]"));
assert_eq!(body["variables"]["id"], "project-1");
assert_eq!(body["variables"]["teamId"], "team-cadenza");
}
#[tokio::test]
async fn create_update_and_comment_map_to_linear_mutations() {
let (base_url, requests) = test_server::spawn(vec![
json_response(
StatusCode::OK,
json!({
"data": {
"workflowStates": {
"nodes": [
{ "id": "state-in-progress", "position": 2.0 },
{ "id": "state-todo", "position": 1.0 }
]
}
}
}),
),
json_response(
StatusCode::OK,
json!({ "data": { "issueCreate": { "issue": { "id": "issue-123" } } } }),
),
json_response(
StatusCode::OK,
json!({ "data": { "issueUpdate": { "issue": { "id": "issue-123" } } } }),
),
json_response(
StatusCode::OK,
json!({ "data": { "commentCreate": { "comment": { "id": "comment-1" } } } }),
),
])
.await;
let client = LinearClient::with_base_url(
"linear-secret".to_string(),
Some("team-9".to_string()),
base_url,
);
let item_id = client
.create_item(
"project-123",
&PmItemCreate {
name: "Implement client".to_string(),
description: "Build the GraphQL adapter".to_string(),
},
)
.await
.expect("create item succeeds");
client
.update_item(
&item_id,
&PmItemUpdate {
name: Some("Implement Linear client".to_string()),
description: Some("Build the GraphQL adapter and tests".to_string()),
},
)
.await
.expect("update item succeeds");
client
.comment(&item_id, "Shipped in v0.9.9")
.await
.expect("comment succeeds");
assert_eq!(item_id, "issue-123");
let requests = requests.lock().await;
assert_eq!(requests.len(), 4);
let states_body: Value =
serde_json::from_str(&requests[0].body).expect("states body is json");
assert!(states_body["query"]
.as_str()
.expect("query present")
.contains("UnstartedWorkflowStates"));
let create_body: Value =
serde_json::from_str(&requests[1].body).expect("create body is json");
assert_eq!(create_body["variables"]["stateId"], json!("state-todo"));
let update_body: Value =
serde_json::from_str(&requests[2].body).expect("update body is json");
assert_eq!(
update_body["variables"]["input"],
json!({
"title": "Implement Linear client",
"description": "Build the GraphQL adapter and tests",
})
);
}
#[tokio::test]
async fn pr_linkage_maps_to_attachment_and_comment_mutations() {
let (base_url, requests) = test_server::spawn(vec![
json_response(
StatusCode::OK,
json!({ "data": { "attachmentLinkURL": { "attachment": { "id": "att-1" } } } }),
),
json_response(
StatusCode::OK,
json!({ "data": { "attachmentUpdate": { "attachment": { "id": "att-1" } } } }),
),
json_response(
StatusCode::OK,
json!({ "data": { "commentUpdate": { "comment": { "id": "comment-1" } } } }),
),
])
.await;
let client = LinearClient::with_base_url("linear-secret".to_string(), None, base_url);
let attachment_id = client
.link_attachment("issue-1", "https://example/pr/7", "GitHub PR #7")
.await
.expect("link attachment succeeds");
assert_eq!(attachment_id, "att-1");
client
.update_attachment("att-1", "GitHub PR #7", "Merged")
.await
.expect("update attachment succeeds");
client
.update_comment("comment-1", "updated body")
.await
.expect("update comment succeeds");
let requests = requests.lock().await;
assert_eq!(requests.len(), 3);
let link: Value = serde_json::from_str(&requests[0].body).expect("link body is json");
assert!(link["query"]
.as_str()
.expect("query present")
.contains("attachmentLinkURL"));
assert_eq!(link["variables"]["issueId"], json!("issue-1"));
assert_eq!(link["variables"]["url"], json!("https://example/pr/7"));
assert!(
link["variables"].get("subtitle").is_none(),
"attachmentLinkURL must not send a subtitle variable"
);
let update: Value =
serde_json::from_str(&requests[1].body).expect("attachment update body is json");
assert!(update["query"]
.as_str()
.expect("query present")
.contains("attachmentUpdate"));
assert_eq!(update["variables"]["subtitle"], json!("Merged"));
let comment: Value =
serde_json::from_str(&requests[2].body).expect("comment update body is json");
assert!(comment["query"]
.as_str()
.expect("query present")
.contains("commentUpdate"));
assert_eq!(comment["variables"]["id"], json!("comment-1"));
}
#[tokio::test]
async fn update_item_omits_absent_text_fields() {
let (base_url, requests) = test_server::spawn(vec![json_response(
StatusCode::OK,
json!({ "data": { "issueUpdate": { "issue": { "id": "issue-123" } } } }),
)])
.await;
let client = LinearClient::with_base_url("linear-secret".to_string(), None, base_url);
client
.update_item(
"issue-123",
&PmItemUpdate {
name: None,
description: Some("Only the description changes".to_string()),
},
)
.await
.expect("description-only update succeeds");
let requests = requests.lock().await;
let update_body: Value =
serde_json::from_str(&requests[0].body).expect("update body is json");
assert_eq!(
update_body["variables"]["input"],
json!({ "description": "Only the description changes" })
);
assert!(update_body["variables"]["input"].get("title").is_none());
}
#[tokio::test]
async fn create_item_omits_state_when_team_has_no_unstarted_state() {
let (base_url, requests) = test_server::spawn(vec![
json_response(
StatusCode::OK,
json!({ "data": { "workflowStates": { "nodes": [] } } }),
),
json_response(
StatusCode::OK,
json!({ "data": { "issueCreate": { "issue": { "id": "issue-123" } } } }),
),
])
.await;
let client = LinearClient::with_base_url(
"linear-secret".to_string(),
Some("team-9".to_string()),
base_url,
);
client
.create_item(
"project-123",
&PmItemCreate {
name: "Implement client".to_string(),
description: "Build the GraphQL adapter".to_string(),
},
)
.await
.expect("create item succeeds");
let requests = requests.lock().await;
let create_body: Value =
serde_json::from_str(&requests[1].body).expect("create body is json");
assert_eq!(create_body["variables"]["stateId"], Value::Null);
}
#[tokio::test]
async fn ensure_team_adopts_matching_key_without_creating() {
let claimed = "<!-- loopflow-repository: loopflowstudio/loopflow -->";
let teams = json_response(
StatusCode::OK,
json!({ "data": { "teams": { "nodes": [
{ "id": "team-prd", "name": "Product", "key": "PRD", "description": claimed },
{ "id": "team-inf", "name": "Infrastructure", "key": "INF", "description": null },
] } } }),
);
let (base_url, requests) = test_server::spawn(vec![teams.clone(), teams]).await;
let client = LinearClient::with_base_url("linear-secret".to_string(), None, base_url);
let binding = client
.ensure_team("Product", "PRD", "loopflowstudio/loopflow")
.await
.expect("adopt existing team");
assert_eq!(
binding,
TeamBinding {
id: "team-prd".to_string(),
key: "PRD".to_string(),
created: false,
}
);
assert_eq!(requests.lock().await.len(), 2);
}
#[tokio::test]
async fn ensure_team_creates_when_absent() {
let (base_url, requests) = test_server::spawn(vec![
json_response(
StatusCode::OK,
json!({ "data": { "teams": { "nodes": [] } } }),
),
json_response(
StatusCode::OK,
json!({ "data": { "teamCreate": { "team": { "id": "team-new" } } } }),
),
json_response(
StatusCode::OK,
json!({ "data": { "teams": { "nodes": [{
"id": "team-new", "name": "Product", "key": "PRD",
"description": "<!-- loopflow-repository: loopflowstudio/loopflow -->"
}] } } }),
),
])
.await;
let client = LinearClient::with_base_url("linear-secret".to_string(), None, base_url);
let binding = client
.ensure_team("Product", "prd", "loopflowstudio/loopflow")
.await
.expect("create team");
assert_eq!(
binding,
TeamBinding {
id: "team-new".to_string(),
key: "PRD".to_string(),
created: true,
}
);
let requests = requests.lock().await;
let create_body: Value =
serde_json::from_str(&requests[1].body).expect("create body is json");
assert_eq!(create_body["variables"]["key"], "PRD");
assert_eq!(create_body["variables"]["name"], "Product");
assert_eq!(
create_body["variables"]["description"],
"<!-- loopflow-repository: loopflowstudio/loopflow -->"
);
}
#[tokio::test]
async fn ensure_team_refuses_key_owned_by_another_team() {
let (base_url, _requests) = test_server::spawn(vec![json_response(
StatusCode::OK,
json!({ "data": { "teams": { "nodes": [
{ "id": "team-x", "name": "Platform", "key": "PRD" },
] } } }),
)])
.await;
let client = LinearClient::with_base_url("linear-secret".to_string(), None, base_url);
let err = client
.ensure_team("Product", "PRD", "loopflowstudio/loopflow")
.await
.expect_err("conflicting key is refused");
let message = err.to_string();
assert!(message.contains("PRD"), "{message}");
assert!(message.contains("Platform"), "{message}");
}
#[tokio::test]
async fn ensure_team_refuses_name_owned_under_a_different_key() {
let (base_url, _requests) = test_server::spawn(vec![json_response(
StatusCode::OK,
json!({ "data": { "teams": { "nodes": [
{ "id": "team-y", "name": "Product", "key": "PROD" },
] } } }),
)])
.await;
let client = LinearClient::with_base_url("linear-secret".to_string(), None, base_url);
let err = client
.ensure_team("Product", "PRD", "loopflowstudio/loopflow")
.await
.expect_err("name reused under a different key is refused");
assert!(err.to_string().contains("PROD"), "{err}");
}
#[tokio::test]
async fn ensure_team_claims_unmarked_team_without_erasing_human_description() {
let (base_url, requests) = test_server::spawn(vec![
json_response(
StatusCode::OK,
json!({ "data": { "teams": { "nodes": [{
"id": "team-loo", "name": "Loopflow", "key": "LOO",
"description": "Human-owned team notes."
}] } } }),
),
json_response(
StatusCode::OK,
json!({ "data": { "teamUpdate": { "team": { "id": "team-loo" } } } }),
),
json_response(
StatusCode::OK,
json!({ "data": { "teams": { "nodes": [{
"id": "team-loo", "name": "Loopflow", "key": "LOO",
"description": "Human-owned team notes.\n\n<!-- loopflow-repository: loopflowstudio/loopflow -->"
}] } } }),
),
])
.await;
let client = LinearClient::with_base_url("linear-secret".to_string(), None, base_url);
let binding = client
.ensure_team("Loopflow", "LOO", "loopflowstudio/loopflow")
.await
.expect("claim unmarked Team");
assert_eq!(binding.id, "team-loo");
let requests = requests.lock().await;
let update: Value = serde_json::from_str(&requests[1].body).unwrap();
let description = update["variables"]["description"].as_str().unwrap();
assert!(description.starts_with("Human-owned team notes."));
assert!(description.contains("loopflowstudio/loopflow"));
}
#[tokio::test]
async fn ensure_team_refuses_a_foreign_repository_claim_before_mutation() {
let (base_url, requests) = test_server::spawn(vec![json_response(
StatusCode::OK,
json!({ "data": { "teams": { "nodes": [{
"id": "team-loo", "name": "Loopflow", "key": "LOO",
"description": "<!-- loopflow-repository: acme/other -->"
}] } } }),
)])
.await;
let client = LinearClient::with_base_url("linear-secret".to_string(), None, base_url);
let error = client
.ensure_team("Loopflow", "LOO", "loopflowstudio/loopflow")
.await
.unwrap_err();
assert!(error.to_string().contains("acme/other"));
assert!(error.to_string().contains("loopflowstudio/loopflow"));
assert_eq!(requests.lock().await.len(), 1);
}
#[tokio::test]
async fn configured_team_rebind_is_refused_before_claiming_another_team() {
let (base_url, requests) = test_server::spawn(vec![json_response(
StatusCode::OK,
json!({ "data": { "teams": { "nodes": [{
"id": "team-loo", "name": "Loopflow", "key": "LOO",
"description": "Human-owned team notes."
}] } } }),
)])
.await;
let client = LinearClient::with_base_url("linear-secret".to_string(), None, base_url);
let error = client
.claim_configured_team("team-loo", "loopflowstudio/loopflow", None, Some("HOO"))
.await
.unwrap_err();
assert!(error.to_string().contains("already bound"));
assert!(error.to_string().contains("lf pm reteam"));
assert_eq!(requests.lock().await.len(), 1);
}
#[test]
fn repository_claim_rejects_ambiguous_markers() {
let description =
"<!-- loopflow-repository: acme/one -->\n<!-- loopflow-repository: acme/two -->";
assert!(repository_claim(description).is_err());
}
#[tokio::test]
async fn issue_ownership_resolves_project_and_team_by_stable_ids() {
let (base_url, _requests) = test_server::spawn(vec![json_response(
StatusCode::OK,
json!({ "data": { "issue": {
"id": "issue-uuid", "identifier": "LOO-42", "url": null,
"title": "Resolve ownership", "description": "",
"prioritySortOrder": 0.0, "sortOrder": 0.0,
"assignee": null, "state": { "type": "unstarted" },
"team": { "id": "team-loo" },
"project": {
"id": "project-api", "name": "Product — Loopflow API",
"description": "", "content": "## Definition\n\nOne model.\n\n## KRs\n",
"initiatives": { "nodes": [{ "id": "initiative-product" }] },
"teams": { "nodes": [{ "id": "team-loo" }] }
}
} } }),
)])
.await;
let client = LinearClient::with_base_url(
"linear-secret".to_string(),
Some("team-loo".to_string()),
base_url,
);
let (item, project) = client.issue_ownership("LOO-42").await.unwrap();
assert_eq!(item.project_id, "project-api");
assert_eq!(item.team_id, "team-loo");
assert_eq!(project.initiative_ids, ["initiative-product"]);
assert_eq!(project.team_ids, ["team-loo"]);
}
#[tokio::test]
async fn complete_item_resolves_state_from_the_issue_team_not_the_wave_team() {
let (base_url, requests) = test_server::spawn(vec![
json_response(
StatusCode::OK,
json!({ "data": { "issue": { "team": { "id": "team-eng" } } } }),
),
json_response(
StatusCode::OK,
json!({ "data": { "workflowStates": { "nodes": [{ "id": "state-done" }] } } }),
),
json_response(
StatusCode::OK,
json!({ "data": { "issueUpdate": { "issue": { "id": "ENG-7" } } } }),
),
])
.await;
let client = LinearClient::with_base_url(
"linear-secret".to_string(),
Some("team-wave".to_string()),
base_url,
);
client.complete_item("ENG-7").await.expect("complete item");
let requests = requests.lock().await;
assert_eq!(requests.len(), 3);
let team_body: Value = serde_json::from_str(&requests[0].body).expect("team body is json");
assert!(team_body["query"]
.as_str()
.expect("query present")
.contains("IssueTeam"));
assert_eq!(team_body["variables"]["id"], json!("ENG-7"));
let states_body: Value =
serde_json::from_str(&requests[1].body).expect("states body is json");
assert_eq!(states_body["variables"]["teamId"], json!("team-eng"));
let set_body: Value = serde_json::from_str(&requests[2].body).expect("set body is json");
assert_eq!(set_body["variables"]["stateId"], json!("state-done"));
assert_eq!(set_body["variables"]["id"], json!("ENG-7"));
}
#[tokio::test]
async fn reopen_item_resolves_state_from_the_issue_team_not_the_wave_team() {
let (base_url, requests) = test_server::spawn(vec![
json_response(
StatusCode::OK,
json!({ "data": { "issue": { "team": { "id": "team-eng" } } } }),
),
json_response(
StatusCode::OK,
json!({ "data": { "workflowStates": { "nodes": [
{ "id": "state-todo", "position": 1.0 }
] } } }),
),
json_response(
StatusCode::OK,
json!({ "data": { "issueUpdate": { "issue": { "id": "ENG-7" } } } }),
),
])
.await;
let client = LinearClient::with_base_url(
"linear-secret".to_string(),
Some("team-wave".to_string()),
base_url,
);
client.reopen_item("ENG-7").await.expect("reopen item");
let requests = requests.lock().await;
let states_body: Value =
serde_json::from_str(&requests[1].body).expect("states body is json");
assert_eq!(states_body["variables"]["teamId"], json!("team-eng"));
}
#[test]
fn linear_description_skips_headings_and_truncates() {
let summary = linear_description(
"## Vision\n\nThis is the first paragraph.\n\n## Strategy\n\nSecond paragraph.",
);
assert_eq!(summary, "This is the first paragraph.");
let long = "a".repeat(300);
assert_eq!(linear_description(&long).len(), 255);
assert_eq!(linear_description(""), "");
}
}