#![deny(missing_docs)]
use chrono::{DateTime, Utc};
use onetaskgraph_plugin_api::{
Capabilities, Comment, CommentBody, Cursor, DependencyEdge, DependencyEndpoint, DependencyKind,
DependencySupport, Direction, Document, DocumentQuery, Health, ItemKind, ItemWrite, Label,
LabelFilter, Location, MetadataKey, NativeId, NewComment, Page, PageRequest, Priority, Project,
ProjectFilter, ProjectQuery, Repository, SecretResolver, SourceError, SourceName, SourcePlugin,
Status, StatusCategory, Support, Task, TaskQuery, TaskRef, TaskSource, TaskUpdate,
TaskUpdateOutcome, TextFields, TextQuery, UpdatedField, WriteSupport,
};
use schemars::{Schema, schema_for};
use secrecy::{ExposeSecret, SecretString};
use serde::Deserialize;
use serde_json::{Value, json};
pub const KIND: &str = "linear";
pub const MAX_PAGE_SIZE: u32 = 100;
const DEFAULT_ENDPOINT: &str = "https://api.linear.app/graphql";
pub mod graphql {
pub const VIEWER: &str = "query { viewer { id } }";
pub const ISSUE: &str = "query($id:String!){ issue(id:$id){ id identifier title description url createdAt updatedAt archivedAt state{name type} priority labels{nodes{id name color}} project{id} } }";
pub const PROJECT: &str = "query($id:String!){ project(id:$id){ id name description url createdAt updatedAt archivedAt status{name type} labels{nodes{id name color}} } }";
pub const ISSUES: &str = "query($first:Int!,$after:String,$filter:IssueFilter){ issues(first:$first,after:$after,filter:$filter){ nodes{id identifier title description url createdAt updatedAt state{name type} priority labels{nodes{id name color}} project{id}} pageInfo{hasNextPage endCursor} } }";
pub const PROJECTS: &str = "query($first:Int!,$after:String,$filter:ProjectFilter){ projects(first:$first,after:$after,filter:$filter){ nodes{id name description url createdAt updatedAt status{name type} labels{nodes{id name color}}} pageInfo{hasNextPage endCursor} } }";
pub const LABELS: &str = "query($first:Int,$after:String){ issueLabels(first:$first,after:$after){ nodes{id name color} pageInfo{hasNextPage endCursor} } }";
pub const ISSUE_RELATIONS: &str = "query($id:String!,$first:Int!,$after:String){ issue(id:$id){ description relations(first:$first,after:$after){nodes{id type relatedIssue{id}} pageInfo{hasNextPage endCursor}} inverseRelations(first:$first,after:$after){nodes{id type issue{id}} pageInfo{hasNextPage endCursor}} } }";
pub const PROJECT_RELATIONS: &str = "query($id:String!,$first:Int!,$after:String){ project(id:$id){ description relations(first:$first,after:$after){nodes{id type relatedProject{id}} pageInfo{hasNextPage endCursor}} inverseRelations(first:$first,after:$after){nodes{id type project{id}} pageInfo{hasNextPage endCursor}} } }";
pub const TEAM: &str =
"query($key:String!){ teams(filter:{key:{eqIgnoreCase:$key}}){nodes{id}} }";
pub const ISSUE_STATE: &str = "query($name:String!,$team:ID!){ workflowStates(filter:{name:{eqIgnoreCase:$name},team:{id:{eq:$team}}}){nodes{id}} }";
pub const ISSUE_STATE_OF_TYPE: &str = "query($type:String!,$team:ID!){ workflowStates(filter:{type:{eq:$type},team:{id:{eq:$team}}}){nodes{id name}} }";
pub const TEAM_WORKFLOW_STATES: &str =
"query($team:ID!){ workflowStates(filter:{team:{id:{eq:$team}}}){nodes{id name type}} }";
pub const PROJECT_STATUS: &str = "query{ projectStatuses{nodes{id name}} }";
pub const ISSUE_LABEL: &str =
"query($name:String!){ issueLabels(filter:{name:{eqIgnoreCase:$name}}){nodes{id}} }";
pub const PROJECT_LABEL: &str =
"query($name:String!){ projectLabels(filter:{name:{eqIgnoreCase:$name}}){nodes{id}} }";
pub const ISSUE_CREATE: &str =
"mutation($input:IssueCreateInput!){ issueCreate(input:$input){success issue{id}} }";
pub const ISSUE_UPDATE: &str = "mutation($id:String!,$input:IssueUpdateInput!){ issueUpdate(id:$id,input:$input){success issue{id}} }";
pub const ISSUE_PRIORITY_UPDATE: &str = "mutation($id:String!,$input:IssueUpdateInput!){ issueUpdate(id:$id,input:$input){success issue{id priority}} }";
pub const PROJECT_CREATE: &str =
"mutation($input:ProjectCreateInput!){ projectCreate(input:$input){success project{id}} }";
pub const PROJECT_UPDATE: &str = "mutation($id:String!,$input:ProjectUpdateInput!){ projectUpdate(id:$id,input:$input){success project{id}} }";
pub const ISSUE_RELATION_CREATE: &str = "mutation($input:IssueRelationCreateInput!){ issueRelationCreate(input:$input){success issueRelation{id}} }";
pub const PROJECT_RELATION_CREATE: &str = "mutation($input:ProjectRelationCreateInput!){ projectRelationCreate(input:$input){success projectRelation{id}} }";
pub const ISSUE_RELATION_DELETE: &str =
"mutation($id:String!){ issueRelationDelete(id:$id){success} }";
pub const PROJECT_RELATION_DELETE: &str =
"mutation($id:String!){ projectRelationDelete(id:$id){success} }";
pub const ISSUE_DELETE: &str = "mutation($id:String!){ issueDelete(id:$id){success} }";
pub const PROJECT_DELETE: &str = "mutation($id:String!){ projectDelete(id:$id){success} }";
pub const DOCUMENT: &str = "query($id:String!){ document(id:$id){ id title content url createdAt updatedAt archivedAt project{id} } }";
pub const DOCUMENTS: &str = "query($first:Int,$after:String,$filter:DocumentFilter){ documents(first:$first,after:$after,filter:$filter){ nodes{id title content url createdAt updatedAt project{id}} pageInfo{hasNextPage endCursor} } }";
pub const DOCUMENT_CREATE: &str = "mutation($input:DocumentCreateInput!){ documentCreate(input:$input){success document{id}} }";
pub const DOCUMENT_UPDATE: &str = "mutation($id:String!,$input:DocumentUpdateInput!){ documentUpdate(id:$id,input:$input){success document{id}} }";
pub const DOCUMENT_DELETE: &str = "mutation($id:String!){ documentDelete(id:$id){success} }";
pub const ISSUE_COMMENTS: &str = "query($id:String!,$last:Int,$before:String){ issue(id:$id){ archivedAt comments(last:$last,before:$before){ nodes{id body url createdAt updatedAt user{displayName}} pageInfo{hasPreviousPage startCursor} } } }";
pub const COMMENT: &str = "query($id:String){ comment(id:$id){ id archivedAt issue{id} } }";
pub const COMMENT_CREATE: &str = "mutation($input:CommentCreateInput!){ commentCreate(input:$input){success comment{id body url createdAt updatedAt user{displayName}}} }";
pub const COMMENT_UPDATE: &str = "mutation($id:String!,$input:CommentUpdateInput!){ commentUpdate(id:$id,input:$input){success comment{id body url createdAt updatedAt user{displayName}}} }";
pub const COMMENT_DELETE: &str = "mutation($id:String!){ commentDelete(id:$id){success} }";
}
use graphql::{
DOCUMENT, DOCUMENTS, ISSUE, ISSUE_RELATIONS, ISSUES, LABELS, PROJECT, PROJECT_RELATIONS,
PROJECTS, VIEWER,
};
#[derive(Debug, Clone, Deserialize, serde::Serialize, schemars::JsonSchema)]
#[serde(default, deny_unknown_fields)]
pub struct LinearConfig {
#[schemars(with = "String")]
api_key_env: EnvName,
team: Option<Team>,
#[schemars(with = "String")]
endpoint: Endpoint,
#[schemars(schema_with = "status_mapping_schema")]
status_mapping: std::collections::HashMap<StatusCategory, Option<StateName>>,
project: Option<ProjectScope>,
}
fn status_mapping_schema(generator: &mut schemars::SchemaGenerator) -> Schema {
let mut schema = <std::collections::HashMap<StatusCategory, Option<StateName>> as schemars::JsonSchema>::json_schema(generator);
schema.insert(
"propertyNames".to_owned(),
generator.subschema_for::<StatusCategory>().to_value(),
);
schema
}
#[derive(Debug, Clone, PartialEq, Eq, Deserialize, serde::Serialize, schemars::JsonSchema)]
#[serde(try_from = "String", into = "String")]
#[schemars(rename = "LinearWorkflowStateName", extend("minLength" = 1))]
struct StateName(String);
impl From<StateName> for String {
fn from(value: StateName) -> Self {
value.0
}
}
impl TryFrom<String> for StateName {
type Error = String;
fn try_from(value: String) -> Result<Self, Self::Error> {
if value.trim().is_empty() {
Err("a status_mapping workflow state name cannot be blank".into())
} else {
Ok(Self(value))
}
}
}
#[derive(Debug, Clone, Deserialize, serde::Serialize, schemars::JsonSchema)]
#[serde(try_from = "String", into = "String")]
#[schemars(rename = "LinearProjectId", extend("minLength" = 1))]
struct ProjectScope(String);
impl From<ProjectScope> for String {
fn from(value: ProjectScope) -> Self {
value.0
}
}
impl TryFrom<String> for ProjectScope {
type Error = String;
fn try_from(value: String) -> Result<Self, Self::Error> {
if value.trim().is_empty() {
Err("a project id cannot be blank".into())
} else {
Ok(Self(value))
}
}
}
#[derive(Debug, Clone, Deserialize, serde::Serialize)]
#[serde(try_from = "String", into = "String")]
struct EnvName(String);
impl From<EnvName> for String {
fn from(value: EnvName) -> Self {
value.0
}
}
impl TryFrom<String> for EnvName {
type Error = String;
fn try_from(value: String) -> Result<Self, Self::Error> {
let mut bytes = value.bytes();
if bytes
.next()
.is_some_and(|byte| byte == b'_' || byte.is_ascii_uppercase())
&& bytes.all(|byte| byte == b'_' || byte.is_ascii_uppercase() || byte.is_ascii_digit())
{
Ok(Self(value))
} else {
Err("must be an uppercase environment-variable name".into())
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Deserialize, serde::Serialize, schemars::JsonSchema)]
#[serde(try_from = "String", into = "String")]
#[schemars(rename = "LinearTeam", extend("minLength" = 1))]
struct Team(String);
impl From<Team> for String {
fn from(value: Team) -> Self {
value.0
}
}
impl TryFrom<String> for Team {
type Error = String;
fn try_from(value: String) -> Result<Self, Self::Error> {
if value.trim().is_empty() {
Err("must not be empty".into())
} else {
Ok(Self(value))
}
}
}
#[derive(Debug, Clone, Deserialize, serde::Serialize)]
#[serde(try_from = "String", into = "String")]
struct Endpoint(String);
impl From<Endpoint> for String {
fn from(value: Endpoint) -> Self {
value.0
}
}
impl TryFrom<String> for Endpoint {
type Error = String;
fn try_from(value: String) -> Result<Self, Self::Error> {
let url = reqwest::Url::parse(&value).map_err(|e| e.to_string())?;
if matches!(url.scheme(), "http" | "https") {
Ok(Self(value))
} else {
Err("must use http or https".into())
}
}
}
impl Default for LinearConfig {
fn default() -> Self {
Self {
api_key_env: EnvName("LINEAR_API_KEY".into()),
team: None,
endpoint: Endpoint(DEFAULT_ENDPOINT.into()),
status_mapping: std::collections::HashMap::new(),
project: None,
}
}
}
const fn category_position(category: StatusCategory) -> usize {
match category {
StatusCategory::Draft => 0,
StatusCategory::Backlog => 1,
StatusCategory::Todo => 2,
StatusCategory::Queued => 3,
StatusCategory::InProgress => 4,
StatusCategory::Done => 5,
StatusCategory::Cancelled => 6,
StatusCategory::Unknown => 7,
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum WorkflowType {
Backlog,
Unstarted,
Started,
Completed,
Canceled,
}
impl WorkflowType {
const ALL: [Self; 5] = [
Self::Backlog,
Self::Unstarted,
Self::Started,
Self::Completed,
Self::Canceled,
];
const fn position(self) -> usize {
match self {
Self::Backlog => 0,
Self::Unstarted => 1,
Self::Started => 2,
Self::Completed => 3,
Self::Canceled => 4,
}
}
const fn of(category: StatusCategory) -> Option<Self> {
match category {
StatusCategory::Backlog => Some(Self::Backlog),
StatusCategory::Todo => Some(Self::Unstarted),
StatusCategory::InProgress => Some(Self::Started),
StatusCategory::Done => Some(Self::Completed),
StatusCategory::Cancelled => Some(Self::Canceled),
StatusCategory::Draft | StatusCategory::Queued | StatusCategory::Unknown => None,
}
}
const fn as_str(self) -> &'static str {
match self {
Self::Backlog => "backlog",
Self::Unstarted => "unstarted",
Self::Started => "started",
Self::Completed => "completed",
Self::Canceled => "canceled",
}
}
}
const _: () = {
let mut index = 0;
while index < WorkflowType::ALL.len() {
assert!(WorkflowType::ALL[index].position() == index);
index += 1;
}
};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum StateTarget<'a> {
Named(&'a StateName),
FirstOfType(WorkflowType),
DisabledByMapping,
NoWorkflowType,
}
#[derive(Debug, Clone, Default)]
struct StatusMapping {
entries: Vec<(StatusCategory, Option<StateName>)>,
}
impl StatusMapping {
fn resolve(
configured: std::collections::HashMap<StatusCategory, Option<StateName>>,
instance: &SourceName,
) -> Result<Self, SourceError> {
let mut configured: Vec<_> = configured.into_iter().collect();
configured.sort_by_key(|(category, _)| category_position(*category));
let mut entries: Vec<(StatusCategory, Option<StateName>)> = Vec::new();
for (category, name) in configured {
if let Some(name) = &name
&& let Some((other, _)) = entries.iter().find(|(_, held)| {
held.as_ref()
.is_some_and(|held| held.0.eq_ignore_ascii_case(&name.0))
})
{
return Err(SourceError::Config {
message: format!(
"source {instance}: status_mapping sends both {} and {} to the workflow \
state {:?}; one state cannot read back as two categories, so map one of \
them to another state or to null",
category_word(*other),
category_word(category),
name.0
),
});
}
entries.push((category, name));
}
Ok(Self { entries })
}
fn is_empty(&self) -> bool {
self.entries.is_empty()
}
fn target(&self, category: StatusCategory) -> StateTarget<'_> {
match self.entries.iter().find(|(held, _)| *held == category) {
Some((_, Some(name))) => StateTarget::Named(name),
Some((_, None)) => StateTarget::DisabledByMapping,
None => WorkflowType::of(category)
.map_or(StateTarget::NoWorkflowType, StateTarget::FirstOfType),
}
}
fn category_of(&self, state: &str) -> Option<StatusCategory> {
self.entries.iter().find_map(|(category, name)| {
name.as_ref()
.is_some_and(|name| name.0.eq_ignore_ascii_case(state))
.then_some(*category)
})
}
fn named(&self) -> impl Iterator<Item = (StatusCategory, &StateName)> {
self.entries
.iter()
.filter_map(|(category, name)| name.as_ref().map(|name| (*category, name)))
}
}
#[derive(Debug, Clone, Copy, Default)]
pub struct Plugin;
impl SourcePlugin for Plugin {
fn kind(&self) -> &'static str {
KIND
}
fn config_schema(&self) -> Schema {
schema_for!(LinearConfig)
}
fn build(
&self,
name: &SourceName,
config: &Value,
secrets: &dyn SecretResolver,
) -> Result<Box<dyn TaskSource>, SourceError> {
let config: LinearConfig =
serde_json::from_value(config.clone()).map_err(|e| SourceError::Config {
message: format!("source {name}: {e}"),
})?;
Ok(Box::new(LinearSource::new(name, config, secrets)?))
}
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, schemars::JsonSchema)]
pub struct WorkflowStatesReport {
pub source: SourceName,
team: Team,
pub states: Vec<MappedWorkflowState>,
}
impl WorkflowStatesReport {
#[must_use]
pub fn team(&self) -> &str {
&self.team.0
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct MappedWorkflowState {
category: StatusCategory,
state: StateName,
found: Found,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Found {
Present(String),
Missing,
}
impl MappedWorkflowState {
#[must_use]
pub fn category(&self) -> StatusCategory {
self.category
}
#[must_use]
pub fn state(&self) -> &str {
&self.state.0
}
#[must_use]
pub fn found(&self) -> &Found {
&self.found
}
}
#[derive(serde::Serialize, schemars::JsonSchema)]
#[schemars(rename = "MappedWorkflowState")]
struct MappedWorkflowStateWire<'a> {
category: StatusCategory,
state: &'a StateName,
present: bool,
#[serde(rename = "type", skip_serializing_if = "Option::is_none")]
state_type: Option<&'a str>,
}
impl<'a> From<&'a MappedWorkflowState> for MappedWorkflowStateWire<'a> {
fn from(mapped: &'a MappedWorkflowState) -> Self {
let state_type = match &mapped.found {
Found::Present(kind) => Some(kind.as_str()),
Found::Missing => None,
};
Self {
category: mapped.category,
state: &mapped.state,
present: state_type.is_some(),
state_type,
}
}
}
impl serde::Serialize for MappedWorkflowState {
fn serialize<S: serde::Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
MappedWorkflowStateWire::from(self).serialize(serializer)
}
}
impl schemars::JsonSchema for MappedWorkflowState {
fn schema_name() -> std::borrow::Cow<'static, str> {
MappedWorkflowStateWire::schema_name()
}
fn json_schema(generator: &mut schemars::SchemaGenerator) -> Schema {
MappedWorkflowStateWire::json_schema(generator)
}
}
pub async fn workflow_states(
name: &SourceName,
config: LinearConfig,
secrets: &dyn SecretResolver,
) -> Result<WorkflowStatesReport, SourceError> {
LinearSource::new(name, config, secrets)?
.workflow_states()
.await
}
struct LinearSource {
client: reqwest::Client,
endpoint: Endpoint,
key: SecretString,
team: Option<Team>,
name: SourceName,
statuses: StatusMapping,
project: Option<ProjectScope>,
}
impl LinearSource {
fn new(
name: &SourceName,
config: LinearConfig,
secrets: &dyn SecretResolver,
) -> Result<Self, SourceError> {
let statuses = StatusMapping::resolve(config.status_mapping, name)?;
let key = secrets
.get(&config.api_key_env.0)
.filter(|v| !v.expose_secret().trim().is_empty())
.ok_or_else(|| SourceError::Auth {
message: format!("set environment variable {}", config.api_key_env.0),
})?;
Ok(Self {
client: reqwest::Client::new(),
endpoint: config.endpoint,
key,
team: config.team,
name: name.clone(),
statuses,
project: config.project,
})
}
}
#[derive(Clone, Copy)]
enum WriteKind {
Task,
Project,
}
enum Lookup<'a> {
Team(&'a str),
IssueState { name: &'a str, team: &'a NativeId },
ProjectStatus(&'a str),
IssueLabel(&'a str),
ProjectLabel(&'a str),
}
impl Lookup<'_> {
fn query(&self) -> &'static str {
match self {
Self::Team(_) => graphql::TEAM,
Self::IssueState { .. } => graphql::ISSUE_STATE,
Self::ProjectStatus(_) => graphql::PROJECT_STATUS,
Self::IssueLabel(_) => graphql::ISSUE_LABEL,
Self::ProjectLabel(_) => graphql::PROJECT_LABEL,
}
}
fn connection(&self) -> &'static str {
match self {
Self::Team(_) => "teams",
Self::IssueState { .. } => "workflowStates",
Self::ProjectStatus(_) => "projectStatuses",
Self::IssueLabel(_) => "issueLabels",
Self::ProjectLabel(_) => "projectLabels",
}
}
fn diagnostic(&self) -> String {
match self {
Self::Team(_) => "configured team".into(),
Self::IssueState { name, .. } => format!("workflow state {name:?}"),
Self::ProjectStatus(name) => format!("project status {name:?}"),
Self::IssueLabel(name) | Self::ProjectLabel(name) => format!("label {name:?}"),
}
}
fn variables(&self) -> Value {
match self {
Self::Team(key) => json!({"key":key}),
Self::IssueState { name, team } => json!({"name":name,"team":team.0}),
Self::IssueLabel(name) | Self::ProjectLabel(name) => json!({"name":name}),
Self::ProjectStatus(_) => json!({}),
}
}
fn local_name(&self) -> Option<&str> {
match self {
Self::ProjectStatus(name) => Some(name),
_ => None,
}
}
}
#[derive(Clone, Copy)]
enum MutationRoot {
IssueCreate,
IssueUpdate,
ProjectCreate,
ProjectUpdate,
IssueRelationCreate,
ProjectRelationCreate,
IssueRelationDelete,
ProjectRelationDelete,
IssueDelete,
ProjectDelete,
DocumentCreate,
DocumentUpdate,
DocumentDelete,
CommentCreate,
CommentUpdate,
CommentDelete,
}
impl MutationRoot {
fn as_str(self) -> &'static str {
match self {
Self::IssueCreate => "issueCreate",
Self::IssueUpdate => "issueUpdate",
Self::ProjectCreate => "projectCreate",
Self::ProjectUpdate => "projectUpdate",
Self::IssueRelationCreate => "issueRelationCreate",
Self::ProjectRelationCreate => "projectRelationCreate",
Self::IssueRelationDelete => "issueRelationDelete",
Self::ProjectRelationDelete => "projectRelationDelete",
Self::IssueDelete => "issueDelete",
Self::ProjectDelete => "projectDelete",
Self::DocumentCreate => "documentCreate",
Self::DocumentUpdate => "documentUpdate",
Self::DocumentDelete => "documentDelete",
Self::CommentCreate => "commentCreate",
Self::CommentUpdate => "commentUpdate",
Self::CommentDelete => "commentDelete",
}
}
}
#[derive(Deserialize)]
struct Envelope {
data: Option<Value>,
#[serde(default)]
errors: Vec<GqlError>,
}
#[derive(Deserialize)]
struct GqlError {
message: String,
extensions: Option<Value>,
}
#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
struct GqlExtensions {
code: GqlErrorCode,
retry_after: Option<u64>,
}
impl GqlError {
fn coded(&self) -> Option<GqlExtensions> {
self.extensions
.as_ref()
.and_then(|value| serde_json::from_value(value.clone()).ok())
}
fn said(&self) -> String {
let Some(extensions) = &self.extensions else {
return elided(&self.message);
};
match extensions
.get("userPresentableMessage")
.and_then(Value::as_str)
.filter(|sentence| !sentence.is_empty())
{
Some(sentence) => elided(&format!("{}: {sentence} {extensions}", self.message)),
None => elided(&format!("{}: {extensions}", self.message)),
}
}
}
#[derive(Deserialize)]
enum GqlErrorCode {
#[serde(rename = "RATELIMITED", alias = "RATE_LIMITED")]
RateLimited,
#[serde(other)]
Other,
}
const SAID_LIMIT: usize = 400;
fn elided(said: &str) -> String {
let mut printable = String::new();
let mut spaced = true;
for character in said.chars() {
if character.is_control() || character.is_whitespace() {
if !spaced {
printable.push(' ');
spaced = true;
}
continue;
}
printable.push(character);
spaced = false;
}
let printable = printable.trim_end();
if printable.chars().count() <= SAID_LIMIT {
return printable.to_owned();
}
let kept: String = printable.chars().take(SAID_LIMIT).collect();
format!("{kept}…")
}
impl LinearSource {
async fn send(&self, query: &str, variables: Value) -> Result<Value, SourceError> {
let response = self
.client
.post(&self.endpoint.0)
.header("Authorization", self.key.expose_secret())
.json(&json!({"query": query, "variables": variables}))
.send()
.await
.map_err(|e| SourceError::Unavailable {
message: e.to_string(),
})?;
let status = response.status();
let retry = response
.headers()
.get("retry-after")
.and_then(|v| v.to_str().ok())
.and_then(|v| v.parse().ok());
if status.as_u16() == 429 {
return Err(SourceError::RateLimited {
retry_after_seconds: retry,
message: None,
});
}
if status.as_u16() == 401 || status.as_u16() == 403 {
return Err(SourceError::Auth {
message: "Linear rejected the configured credential".into(),
});
}
if !status.is_success() {
let said = elided(&response.text().await.unwrap_or_default());
return Err(SourceError::Unavailable {
message: if said.is_empty() {
format!("Linear returned HTTP {status}")
} else {
format!("Linear returned HTTP {status}: {said}")
},
});
}
let body: Envelope = response.json().await.map_err(|e| SourceError::Malformed {
message: e.to_string(),
})?;
if let Some(error) = body.errors.first() {
if let Some(extensions) = error
.coded()
.filter(|extensions| matches!(extensions.code, GqlErrorCode::RateLimited))
{
return Err(SourceError::RateLimited {
retry_after_seconds: extensions.retry_after.or(retry),
message: None,
});
}
return Err(SourceError::Refused {
message: error.said(),
});
}
body.data.ok_or_else(|| SourceError::Malformed {
message: "GraphQL response has no data".into(),
})
}
fn label_parts(labels: &onetaskgraph_plugin_api::LabelFilter) -> Vec<Value> {
let mut parts = Vec::new();
if !labels.any_of.is_empty() {
parts.push(json!({"or": labels
.any_of
.iter()
.map(|name| json!({"labels": {"some": {"name": {"eqIgnoreCase": name}}}}))
.collect::<Vec<_>>()}));
}
for name in &labels.all_of {
parts.push(json!({"labels": {"some": {"name": {"eqIgnoreCase": name}}}}));
}
for name in &labels.none_of {
parts.push(json!({"labels": {"every": {"name": {"neqIgnoreCase": name}}}}));
}
parts
}
fn narrowed(mut parts: Vec<Value>) -> Value {
if parts.len() == 1 {
parts.pop().unwrap()
} else {
json!({"and": parts})
}
}
fn issue_filter(&self, query: &TaskQuery) -> Result<Value, SourceError> {
let mut parts = self.issue_scope();
parts.extend(Self::label_parts(&query.labels));
if !query.statuses.is_empty() {
parts.push(self.status_narrowing(&query.statuses));
}
match &query.project {
ProjectFilter::Orphans => parts.push(json!({"project": {"null": true}})),
ProjectFilter::Is(id) => parts.push(json!({"project": {"id": {"eq": id.0}}})),
ProjectFilter::Any => {}
}
if !query.priorities.is_empty() {
parts.push(json!({"priority": {"in": query
.priorities
.iter()
.map(|priority| linear_priority(*priority))
.collect::<Vec<_>>()}}));
}
if let Some(since) = query.commented_since {
let since = since.to_rfc3339_opts(chrono::SecondsFormat::AutoSi, true);
parts.push(json!({"comments": {"some": {"or": [
{"createdAt": {"gte": since}},
{"updatedAt": {"gte": since}},
]}}}));
}
for wanted in &query.metadata {
parts.extend(slot_phrase(wanted.value()));
}
if let Some(origin) = &query.origin {
parts.extend(slot_phrase(origin));
}
if let Some(text) = &query.text {
let title = json!({"title": {"containsIgnoreCase": text.terms}});
let content = json!({"description": {"containsIgnoreCase": text.terms}});
parts.push(match text.fields {
TextFields::Title => title,
TextFields::Content => content,
TextFields::TitleOrContent => json!({"or": [title, content]}),
});
}
Ok(Self::narrowed(parts))
}
fn issue_scope(&self) -> Vec<Value> {
let mut parts = Vec::new();
if let Some(team) = &self.team {
parts.push(json!({"team": {"key": {"eqIgnoreCase": team.0}}}));
}
if let Some(project) = &self.project {
parts.push(json!({"project": {"id": {"eq": project.0}}}));
}
parts
}
fn status_narrowing(&self, statuses: &[StatusCategory]) -> Value {
if self.statuses.is_empty() {
return json!({"state": {"type": {"in": statuses
.iter()
.flat_map(workflow_state_types)
.collect::<Vec<_>>()}}});
}
let mut alternatives = Vec::new();
for category in statuses {
let target = self.statuses.target(*category);
if let StateTarget::Named(name) = target {
alternatives.push(json!({"state": {"name": {"eqIgnoreCase": name.0}}}));
}
let types = workflow_state_types(category);
let by_type = if *category == StatusCategory::Unknown {
Some(
json!({"state": {"type": {"nin": WorkflowType::ALL.map(WorkflowType::as_str)}}}),
)
} else {
(!types.is_empty()).then(|| json!({"state": {"type": {"in": types}}}))
};
if let Some(by_type) = by_type {
let mut parts = vec![by_type];
parts.extend(
self.statuses
.named()
.map(|(_, name)| json!({"state": {"name": {"neqIgnoreCase": name.0}}})),
);
alternatives.push(Self::narrowed(parts));
}
}
match alternatives.len() {
0 => json!({"state": {"type": {"in": Vec::<&str>::new()}}}),
1 => alternatives.pop().expect("one alternative"),
_ => json!({ "or": alternatives }),
}
}
fn confirms(query: &TaskQuery, task: &Task) -> bool {
query.metadata_matches(&task.metadata)
&& query.origin_matches(&task.metadata)
&& (query.priorities.is_empty() || query.priorities.contains(&task.priority))
&& query
.text
.as_ref()
.is_none_or(|text| text_holds(&task.title, task.content.as_deref(), text))
}
fn project_filter(
&self,
labels: &onetaskgraph_plugin_api::LabelFilter,
statuses: &[StatusCategory],
) -> Value {
let mut parts = self.project_scope();
parts.extend(Self::label_parts(labels));
if !statuses.is_empty() {
parts.push(json!({"status": {"type": {"in": statuses.iter().flat_map(project_status_types).collect::<Vec<_>>()}}}));
}
Self::narrowed(parts)
}
fn project_scope(&self) -> Vec<Value> {
let mut parts = Vec::new();
if let Some(team) = &self.team {
parts.push(json!({"accessibleTeams": {"some": {"key": {"eqIgnoreCase": team.0}}}}));
}
if let Some(project) = &self.project {
parts.push(json!({"id": {"eq": project.0}}));
}
parts
}
fn in_scope(&self, project: Option<&NativeId>) -> bool {
self.project
.as_ref()
.is_none_or(|scope| project.is_some_and(|project| project.0 == scope.0))
}
fn filed_in(
&self,
project: Option<&NativeId>,
what: &str,
) -> Result<Option<String>, SourceError> {
match (&self.project, project) {
(None, project) => Ok(project.map(|id| id.0.clone())),
(Some(scope), None) => Ok(Some(scope.0.clone())),
(Some(scope), Some(project)) if project.0 == scope.0 => Ok(Some(scope.0.clone())),
(Some(scope), Some(project)) => Err(SourceError::Refused {
message: format!(
"source {} is scoped to the Linear project {} and cannot hold a {what} in the project {}; next: write it with no project, or with {}, or to a source scoped to {}",
self.name, scope.0, project.0, scope.0, project.0
),
}),
}
}
async fn one_id(&self, lookup: Lookup<'_>) -> Result<NativeId, SourceError> {
let data = self.send(lookup.query(), lookup.variables()).await?;
let connection = lookup.connection();
let nodes = data
.get(connection)
.and_then(|v| v.get("nodes"))
.and_then(Value::as_array)
.ok_or_else(|| SourceError::Malformed {
message: format!("missing {connection}.nodes"),
})?;
let matched = match lookup.local_name() {
Some(name) => {
let mut matched = Vec::new();
for node in nodes {
if str_at(node, "name")?.eq_ignore_ascii_case(name) {
matched.push(node);
}
}
matched
}
None => nodes.iter().collect::<Vec<_>>(),
};
match matched.as_slice() {
[] => Err(SourceError::Refused {
message: format!(
"source {} cannot resolve {}: found 0 matches",
self.name,
lookup.diagnostic()
),
}),
[node] => Ok(NativeId(backend_id(node, "id")?.to_owned())),
nodes => {
let ids = nodes
.iter()
.map(|node| backend_id(node, "id"))
.collect::<Result<Vec<_>, _>>()?;
Err(SourceError::Refused {
message: format!(
"source {} cannot resolve {}: found {} matches with ids {ids:?}",
self.name,
lookup.diagnostic(),
nodes.len()
),
})
}
}
}
async fn team_id(&self) -> Result<NativeId, SourceError> {
let team = self.team.as_ref().ok_or_else(|| SourceError::Refused {
message: format!(
"source {} needs config.team before it can create Linear items",
self.name
),
})?;
self.one_id(Lookup::Team(&team.0)).await
}
async fn label_ids(
&self,
labels: &[Label],
kind: WriteKind,
) -> Result<Vec<NativeId>, SourceError> {
let mut ids = Vec::with_capacity(labels.len());
for label in labels {
ids.push(
self.one_id(if matches!(kind, WriteKind::Project) {
Lookup::ProjectLabel(&label.name)
} else {
Lookup::IssueLabel(&label.name)
})
.await?,
);
}
Ok(ids)
}
fn write_description(
&self,
content: Option<&str>,
metadata: &std::collections::BTreeMap<String, Value>,
repositories: &[Repository],
edges: &[DependencyEdge],
kind: WriteKind,
) -> Result<Option<String>, SourceError> {
Self::long_form(
content,
metadata,
repositories,
self.recorded_ends(edges, kind),
)
}
fn recorded_ends(&self, edges: &[DependencyEdge], kind: WriteKind) -> Vec<Value> {
edges
.iter()
.filter(|edge| {
edge.to.kind
!= match kind {
WriteKind::Task => ItemKind::Task,
WriteKind::Project => ItemKind::Project,
}
|| edge
.to
.id()
.split_once(':')
.is_some_and(|(source, _)| source != self.name.as_str())
})
.map(|edge| json!({"id":edge.to.id(),"kind":edge.to.kind}))
.collect()
}
async fn forward_edges(&self, id: &NativeId) -> Result<Vec<DependencyEdge>, SourceError> {
let mut edges = Vec::new();
let mut cursor = None;
loop {
let page = self
.dependencies(
ISSUE_RELATIONS,
DependencyRoot::Issue,
id,
Direction::DependsOn,
&PageRequest {
cursor,
limit: MAX_PAGE_SIZE,
},
)
.await?;
edges.extend(page.items);
match page.next {
Some(next) => cursor = Some(next),
None => return Ok(edges),
}
}
}
async fn targeted_update(
&self,
id: &NativeId,
update: &TaskUpdate,
) -> Result<Option<TaskUpdateOutcome>, SourceError> {
update.consistent()?;
if let Some(delivers) = &update.delivers {
TaskRef::listed(
TaskRef::DELIVERS_KEY,
id,
Some(&self.name),
delivers.clone(),
)
.map_err(|message| SourceError::Refused { message })?;
}
if let Some(status) = &update.status {
self.writable(status.category)?;
}
let Some((before, description)) = self.issue_held(id).await? else {
return Ok(None);
};
let (visible, held) = metadata_description(description)?;
let mut slot = held.clone();
for (key, value) in &update.metadata_set {
slot.insert(key.as_str().to_owned(), value.clone());
}
for key in &update.metadata_remove {
slot.remove(key.as_str());
}
if let Some(delivers) = &update.delivers {
set_task_list(&mut slot, TaskRef::DELIVERS_KEY, delivers);
}
let mut relations = None;
if let Some(wanted) = &update.depends_on {
let current = self.forward_edges(&before.id).await?;
let ends = |edges: &[DependencyEdge]| {
let mut ends: Vec<(String, String)> = edges
.iter()
.map(|edge| {
(
edge.to.id().to_owned(),
format!("{:?}{:?}", edge.to.kind, edge.kind),
)
})
.collect();
ends.sort();
ends
};
if ends(¤t) != ends(wanted) {
let prepared = self.prepare_edges(wanted, WriteKind::Task).await?;
let recorded = self.recorded_ends(&prepared, WriteKind::Task);
if recorded.is_empty() {
slot.remove(DependencyEdge::RECORDED_KEY);
} else {
slot.insert(DependencyEdge::RECORDED_KEY.into(), Value::Array(recorded));
}
relations = Some(prepared);
}
}
let mut input = serde_json::Map::new();
if let Some(title) = update
.title
.as_ref()
.filter(|title| **title != before.title)
{
input.insert("title".into(), json!(title));
}
let content = update.content.as_deref().or(visible.as_deref());
if content != visible.as_deref() || slot != held {
let written = Self::described(content, &slot)?;
let (reads, read) = metadata_description(written.clone())?;
if reads.as_deref().unwrap_or_default() != content.unwrap_or_default() || read != slot {
return Err(SourceError::Refused {
message: format!(
"this content would read back from source {} as something other than \
itself, or ends in what it reads as its own metadata slot; next: change \
how the content ends",
self.name
),
});
}
input.insert("description".into(), json!(written));
}
if let Some(status) = update
.status
.as_ref()
.filter(|status| status.category != before.status.category)
{
let (state, _) = self.resolve_state(status.category, &before.id).await?;
input.insert("stateId".into(), json!(state.0));
}
if let Some(priority) = update
.priority
.filter(|priority| *priority != before.priority)
{
input.insert("priority".into(), json!(linear_priority(priority)));
}
let sent = !input.is_empty();
if sent {
let data = self
.send(
graphql::ISSUE_UPDATE,
json!({"id":before.id.0,"input":Value::Object(input)}),
)
.await?;
let issue = mutation_payload(&data, MutationRoot::IssueUpdate)?
.get("issue")
.filter(|issue| !issue.is_null())
.ok_or_else(|| SourceError::Malformed {
message: "missing issueUpdate.issue".into(),
})?;
written_is(issue, &before.id)?;
}
if let Some(prepared) = &relations {
self.write_relations(&before.id, prepared, WriteKind::Task)
.await?;
}
let task = if sent || relations.is_some() {
self.get_task(&before.id)
.await?
.ok_or_else(|| SourceError::Malformed {
message: format!("task {id} was updated and then could not be read back"),
})?
} else {
before.clone()
};
let mut written = update.changed(&before, &task);
if relations.is_some() {
written.insert(UpdatedField::DependsOn);
}
Ok(Some(TaskUpdateOutcome {
task,
written,
delivers_before: before.delivers,
}))
}
fn long_form(
content: Option<&str>,
metadata: &std::collections::BTreeMap<String, Value>,
repositories: &[Repository],
recorded: Vec<Value>,
) -> Result<Option<String>, SourceError> {
let mut metadata = metadata.clone();
if repositories.is_empty() {
metadata.remove(Repository::METADATA_KEY);
} else {
metadata.insert(Repository::METADATA_KEY.into(), json!(repositories));
}
if recorded.is_empty() {
metadata.remove(DependencyEdge::RECORDED_KEY);
} else {
metadata.insert(DependencyEdge::RECORDED_KEY.into(), Value::Array(recorded));
}
Self::described(content, &metadata)
}
fn described(
content: Option<&str>,
metadata: &std::collections::BTreeMap<String, Value>,
) -> Result<Option<String>, SourceError> {
let visible = content.unwrap_or_default();
if metadata.is_empty() {
return Ok((!visible.is_empty()).then(|| visible.to_owned()));
}
let slot = slot_text(metadata)?;
Ok(Some(if visible.is_empty() {
slot
} else {
format!("{visible}\n\n{slot}")
}))
}
fn unordered_project_relation(&self, near: &NativeId, far: &str) -> SourceError {
SourceError::Refused {
message: format!(
"source {} cannot carry an unordered dependency between projects, because \
Linear types every project relation `dependency` and that is an ordering; \
record {near} to {far} as a dependency, or between tasks",
self.name,
near = near.0,
),
}
}
fn unordered_project_edge(edges: &[DependencyEdge]) -> Option<&DependencyEdge> {
edges
.iter()
.find(|edge| edge.to.kind == ItemKind::Project && edge.kind == DependencyKind::Related)
}
async fn write_relations(
&self,
near: &NativeId,
edges: &[DependencyEdge],
kind: WriteKind,
) -> Result<(), SourceError> {
let mut cursor: Option<Cursor> = None;
loop {
let data = self
.send(
if matches!(kind, WriteKind::Project) {
PROJECT_RELATIONS
} else {
ISSUE_RELATIONS
},
json!({"id":near.0,"first":MAX_PAGE_SIZE,"after":cursor.as_ref().map(|cursor|&cursor.0)}),
)
.await?;
let root = data
.get(if matches!(kind, WriteKind::Project) {
"project"
} else {
"issue"
})
.ok_or_else(|| SourceError::Malformed {
message: "missing relation item".into(),
})?;
let relations = root
.get("relations")
.ok_or_else(|| SourceError::Malformed {
message: "missing relations".into(),
})?;
for relation in relations
.get("nodes")
.and_then(Value::as_array)
.ok_or_else(|| SourceError::Malformed {
message: "missing relations.nodes".into(),
})?
{
let id = backend_id(relation, "id")?;
let (query, mutation) = if matches!(kind, WriteKind::Project) {
(
graphql::PROJECT_RELATION_DELETE,
MutationRoot::ProjectRelationDelete,
)
} else {
(
graphql::ISSUE_RELATION_DELETE,
MutationRoot::IssueRelationDelete,
)
};
let deleted = self.send(query, json!({"id":id})).await?;
mutation_payload(&deleted, mutation)?;
}
let Some(next) = page_next(relations)? else {
break;
};
cursor = Some(next);
}
const NEAR_ANCHOR: &str = "start";
const FAR_ANCHOR: &str = "end";
for edge in edges {
if edge.to.kind
!= match kind {
WriteKind::Task => ItemKind::Task,
WriteKind::Project => ItemKind::Project,
}
{
continue;
}
let far = match edge.to.id().split_once(':') {
Some((source, native)) if source == self.name.as_str() => native,
Some(_) => continue,
None => edge.to.id(),
};
let relation_type = match (kind, edge.kind) {
(WriteKind::Project, DependencyKind::Blocks) => "dependency",
(WriteKind::Task, DependencyKind::Blocks) => "blocks",
(WriteKind::Task, DependencyKind::Related) => "related",
(WriteKind::Project, DependencyKind::Related) => {
return Err(self.unordered_project_relation(near, edge.to.id()));
}
};
let (query, input) = if matches!(kind, WriteKind::Project) {
(
graphql::PROJECT_RELATION_CREATE,
json!({"projectId":near.0,"relatedProjectId":far,"type":relation_type,"anchorType":NEAR_ANCHOR,"relatedAnchorType":FAR_ANCHOR}),
)
} else {
(
graphql::ISSUE_RELATION_CREATE,
json!({"issueId":near.0,"relatedIssueId":far,"type":relation_type}),
)
};
let data = self.send(query, json!({"input":input})).await?;
let mutation = if matches!(kind, WriteKind::Project) {
MutationRoot::ProjectRelationCreate
} else {
MutationRoot::IssueRelationCreate
};
let payload = mutation_payload(&data, mutation)?;
let relation = payload
.get(if matches!(kind, WriteKind::Project) {
"projectRelation"
} else {
"issueRelation"
})
.ok_or_else(|| SourceError::Malformed {
message: format!("missing {} relation", mutation.as_str()),
})?;
backend_id(relation, "id")?;
}
Ok(())
}
async fn prepare_edges(
&self,
edges: &[DependencyEdge],
kind: WriteKind,
) -> Result<Vec<DependencyEdge>, SourceError> {
let mut prepared = Vec::with_capacity(edges.len());
for edge in edges {
let mut edge = edge.clone();
if edge.to.kind
== match kind {
WriteKind::Task => ItemKind::Task,
WriteKind::Project => ItemKind::Project,
}
&& edge
.to
.id()
.split_once(':')
.is_some_and(|(source, _)| source != self.name.as_str())
{
let filter = match kind {
WriteKind::Project => Self::narrowed(self.project_scope()),
WriteKind::Task => {
let mut parts = self.issue_scope();
parts.extend(slot_phrase(edge.to.id()));
Self::narrowed(parts)
}
};
let mut cursor: Option<Cursor> = None;
loop {
let data = self.send(if matches!(kind, WriteKind::Project) { PROJECTS } else { ISSUES }, json!({"first":MAX_PAGE_SIZE,"after":cursor.as_ref().map(|cursor|&cursor.0),"filter":filter})).await?;
let (items, next) = if matches!(kind, WriteKind::Project) {
let page = connection(&data, "projects", map_project)?;
(
page.items
.into_iter()
.map(|item| (item.id, item.metadata))
.collect::<Vec<_>>(),
page.next,
)
} else {
let page = connection(&data, "issues", |v| {
map_task(v, &self.name, &self.statuses)
})?;
(
page.items
.into_iter()
.map(|item| (item.id, item.metadata))
.collect::<Vec<_>>(),
page.next,
)
};
if let Some((id, _)) = items.into_iter().find(|(_, metadata)| {
metadata.get("onetaskgraph.origin").and_then(Value::as_str)
== Some(edge.to.id())
}) {
edge.to = DependencyEndpoint::from_native(id, edge.to.kind);
break;
}
let Some(next) = next else { break };
cursor = Some(next);
}
}
prepared.push(edge);
}
Ok(prepared)
}
}
#[async_trait::async_trait]
impl TaskSource for LinearSource {
fn kind(&self) -> &'static str {
KIND
}
fn capabilities(&self) -> Capabilities {
Capabilities {
projects: Support::Native,
documents: Support::Native,
comments: Support::Native,
priority: Support::Native,
filter_by_priority: Support::Native,
filter_by_comment_activity: Support::Native,
filter_by_metadata: Support::Native,
filter_by_origin: Support::Native,
orphan_tasks: Support::Native,
filter_by_label: Support::Native,
filter_by_status: Support::Native,
search_title: Support::Native,
search_content: Support::Native,
task_dependencies: DependencySupport::BothDirections,
project_dependencies: DependencySupport::BothDirections,
max_page_size: MAX_PAGE_SIZE,
}
}
fn writes(&self) -> WriteSupport {
WriteSupport::Supported
}
async fn health(&self) -> Result<Health, SourceError> {
let data = self.send(VIEWER, json!({})).await?;
str_at(
data.get("viewer").ok_or_else(|| SourceError::Malformed {
message: "missing viewer".into(),
})?,
"id",
)?;
Ok(Health {
reachable: true,
detail: None,
})
}
async fn get_task(&self, id: &NativeId) -> Result<Option<Task>, SourceError> {
Ok(self.issue_held(id).await?.map(|(task, _)| task))
}
async fn get_project(&self, id: &NativeId) -> Result<Option<Project>, SourceError> {
Ok(self.project_held(id).await?.map(|(project, _)| project))
}
async fn query_tasks(
&self,
query: &TaskQuery,
page: &PageRequest,
) -> Result<Page<Task>, SourceError> {
let d=self.send(ISSUES,json!({"first":page.limit.min(MAX_PAGE_SIZE),"after":page.cursor.as_ref().map(|c|&c.0),"filter":self.issue_filter(query)?})).await?;
let page = connection(&d, "issues", |v| map_task(v, &self.name, &self.statuses))?;
Ok(Page {
items: page
.items
.into_iter()
.filter(|task| Self::confirms(query, task))
.collect(),
next: page.next,
})
}
async fn query_projects(
&self,
query: &ProjectQuery,
page: &PageRequest,
) -> Result<Page<Project>, SourceError> {
let d=self.send(PROJECTS,json!({"first":page.limit.min(MAX_PAGE_SIZE),"after":page.cursor.as_ref().map(|c|&c.0),"filter":self.project_filter(&query.labels,&query.statuses)})).await?;
let page = connection(&d, "projects", map_project)?;
Ok(Page {
items: page
.items
.into_iter()
.filter(|project| {
query.text.as_ref().is_none_or(|text| {
text_holds(&project.title, project.content.as_deref(), text)
})
})
.collect(),
next: page.next,
})
}
async fn labels(&self, page: &PageRequest) -> Result<Page<Label>, SourceError> {
let d = self
.send(
LABELS,
json!({"first":page.limit.min(MAX_PAGE_SIZE),"after":page.cursor.as_ref().map(|c|&c.0)}),
)
.await?;
connection(&d, "issueLabels", map_label)
}
async fn task_dependencies(
&self,
id: &NativeId,
direction: Direction,
page: &PageRequest,
) -> Result<Page<DependencyEdge>, SourceError> {
self.dependencies(ISSUE_RELATIONS, DependencyRoot::Issue, id, direction, page)
.await
}
async fn project_dependencies(
&self,
id: &NativeId,
direction: Direction,
page: &PageRequest,
) -> Result<Page<DependencyEdge>, SourceError> {
self.dependencies(
PROJECT_RELATIONS,
DependencyRoot::Project,
id,
direction,
page,
)
.await
}
async fn write_task(&self, write: &ItemWrite<Task>) -> Result<NativeId, SourceError> {
let near = write.target.as_ref().unwrap_or(&write.item.id);
for (key, entries) in [
(TaskRef::DELIVERS_KEY, &write.item.delivers),
(TaskRef::DELIVERED_BY_KEY, &write.item.delivered_by),
] {
TaskRef::listed(key, near, Some(&self.name), entries.clone())
.map_err(|message| SourceError::Refused { message })?;
}
if self.statuses.target(write.item.status.category) == StateTarget::DisabledByMapping {
return Err(self.disabled(write.item.status.category, true));
}
let project = self.filed_in(write.item.project.as_ref(), "task")?;
let edges = self
.prepare_edges(&write.depends_on, WriteKind::Task)
.await?;
let team = self.team_id().await?;
let state = match self.statuses.target(write.item.status.category) {
StateTarget::Named(name) => {
self.named_state(&name.0, write.item.status.category)
.await?
.0
}
StateTarget::FirstOfType(_)
| StateTarget::DisabledByMapping
| StateTarget::NoWorkflowType => {
self.one_id(Lookup::IssueState {
name: &write.item.status.name,
team: &team,
})
.await?
}
};
let labels = self.label_ids(&write.item.labels, WriteKind::Task).await?;
let mut metadata = write.item.metadata.clone();
set_task_list(&mut metadata, TaskRef::DELIVERS_KEY, &write.item.delivers);
set_task_list(
&mut metadata,
TaskRef::DELIVERED_BY_KEY,
&write.item.delivered_by,
);
let description = self.write_description(
write.item.content.as_deref(),
&metadata,
&write.item.repositories,
&edges,
WriteKind::Task,
)?;
let input = json!({"title":write.item.title,"description":description,"stateId":state,"priority":linear_priority(write.item.priority),"labelIds":labels,"projectId":project});
let (query, variables, root) = match &write.target {
Some(id) => (
graphql::ISSUE_UPDATE,
json!({"id":id.0,"input":input}),
MutationRoot::IssueUpdate,
),
None => (
graphql::ISSUE_CREATE,
{
let mut input = input;
input["teamId"] = Value::String(team.0);
json!({"input":input})
},
MutationRoot::IssueCreate,
),
};
let data = self.send(query, variables).await?;
let issue =
mutation_payload(&data, root)?
.get("issue")
.ok_or_else(|| SourceError::Malformed {
message: format!("missing {}.issue", root.as_str()),
})?;
let id = NativeId(backend_id(issue, "id")?.into());
self.write_relations(&id, &edges, WriteKind::Task).await?;
Ok(id)
}
async fn write_project(&self, write: &ItemWrite<Project>) -> Result<NativeId, SourceError> {
if let Some(edge) = Self::unordered_project_edge(&write.depends_on) {
return Err(self.unordered_project_relation(&write.item.id, edge.to.id()));
}
if let Some(key) = delivery_key_in(&write.item.metadata) {
return Err(self.undeliverable(key, "project"));
}
if let Some(scope) = &self.project
&& write
.target
.as_ref()
.is_none_or(|target| target.0 != scope.0)
{
return Err(SourceError::Refused {
message: format!(
"source {} is scoped to the Linear project {} and holds no other project, \
so it cannot write {}; next: copy the project to a source with no \
`project`, or copy its tasks here",
self.name,
scope.0,
write.target.as_ref().map_or_else(
|| "a new one".to_owned(),
|target| format!("the project {}", target.0)
),
),
});
}
let edges = self
.prepare_edges(&write.depends_on, WriteKind::Project)
.await?;
let team = self.team_id().await?;
let status = self
.one_id(Lookup::ProjectStatus(&write.item.status.name))
.await?;
let labels = self
.label_ids(&write.item.labels, WriteKind::Project)
.await?;
let description = self.write_description(
write.item.content.as_deref(),
&write.item.metadata,
&write.item.repositories,
&edges,
WriteKind::Project,
)?;
let input = json!({"name":write.item.title,"description":description,"statusId":status,"labelIds":labels});
let (query, variables, root) = match &write.target {
Some(id) => (
graphql::PROJECT_UPDATE,
json!({"id":id.0,"input":input}),
MutationRoot::ProjectUpdate,
),
None => (
graphql::PROJECT_CREATE,
{
let mut input = input;
input["teamIds"] = json!([team]);
json!({"input":input})
},
MutationRoot::ProjectCreate,
),
};
let data = self.send(query, variables).await?;
let project = mutation_payload(&data, root)?
.get("project")
.ok_or_else(|| SourceError::Malformed {
message: format!("missing {}.project", root.as_str()),
})?;
let id = NativeId(backend_id(project, "id")?.into());
self.write_relations(&id, &edges, WriteKind::Project)
.await?;
Ok(id)
}
async fn delete_task(&self, id: &NativeId) -> Result<(), SourceError> {
if self.get_task(id).await?.is_none() {
return Ok(());
}
let data = self.send(graphql::ISSUE_DELETE, json!({"id":id.0})).await?;
mutation_payload(&data, MutationRoot::IssueDelete)?;
Ok(())
}
async fn delete_project(&self, id: &NativeId) -> Result<(), SourceError> {
if self.get_project(id).await?.is_none() {
return Ok(());
}
let data = self
.send(graphql::PROJECT_DELETE, json!({"id":id.0}))
.await?;
mutation_payload(&data, MutationRoot::ProjectDelete)?;
Ok(())
}
async fn get_document(&self, id: &NativeId) -> Result<Option<Document>, SourceError> {
Ok(self.document_held(id).await?.map(|(document, _)| document))
}
async fn query_documents(
&self,
query: &DocumentQuery,
page: &PageRequest,
) -> Result<Page<Document>, SourceError> {
let want = page.limit.min(MAX_PAGE_SIZE) as usize;
let mut filter = serde_json::Map::new();
if let ProjectFilter::Is(id) = &query.project {
if !self.in_scope(Some(id)) {
return Ok(Page::last(Vec::new()));
}
filter.insert("project".into(), json!({"id": {"eq": id.0}}));
} else if let Some(scope) = &self.project {
filter.insert("project".into(), json!({"id": {"eq": scope.0}}));
}
let filter = Value::Object(filter);
let mut items = Vec::new();
let mut cursor = page.cursor.clone();
loop {
let first = want.saturating_sub(items.len()).max(1);
let d = self
.send(
DOCUMENTS,
json!({"first":first,"after":cursor.as_ref().map(|cursor|&cursor.0),"filter":filter}),
)
.await?;
let fetched = connection(&d, "documents", map_document)?;
items.extend(fetched.items.into_iter().filter(|document| {
document_matches(document, &query.project, &query.labels)
&& self.in_scope(document.project.as_ref())
&& query.text.as_ref().is_none_or(|text| {
text_holds(&document.title, document.content.as_deref(), text)
})
}));
cursor = fetched.next;
if cursor.is_none() || items.len() >= want {
return Ok(Page {
items,
next: cursor,
});
}
}
}
async fn write_document(&self, write: &ItemWrite<Document>) -> Result<NativeId, SourceError> {
if !write.item.labels.is_empty() {
let named = write
.item
.labels
.iter()
.map(|label| label.name.as_str())
.collect::<Vec<_>>()
.join(", ");
return Err(SourceError::Refused {
message: format!(
"source {} cannot carry a document's labels, because Linear's own \
document type has none: {named}",
self.name
),
});
}
if !write.depends_on.is_empty()
|| write
.item
.metadata
.contains_key(DependencyEdge::RECORDED_KEY)
{
return Err(SourceError::Refused {
message: format!(
"source {} cannot carry {} on a document, because a document is not \
work and nothing may depend on one",
self.name,
DependencyEdge::RECORDED_KEY
),
});
}
if let Some(key) = delivery_key_in(&write.item.metadata) {
return Err(self.undeliverable(key, "document"));
}
let content = Self::long_form(
write.item.content.as_deref(),
&write.item.metadata,
&write.item.repositories,
Vec::new(),
)?;
let project = self.filed_in(write.item.project.as_ref(), "document")?;
let (query, variables, root) = match &write.target {
Some(id) => {
if self.get_document(id).await?.is_none() {
return Err(SourceError::Refused {
message: format!("source {} holds no document {}", self.name, id.0),
});
}
(
graphql::DOCUMENT_UPDATE,
json!({"id":id.0,"input":{"title":write.item.title,"content":content,"projectId":project}}),
MutationRoot::DocumentUpdate,
)
}
None => {
let mut input = json!({"title":write.item.title,"content":content});
match &project {
Some(project) => input["projectId"] = Value::String(project.clone()),
None => input["teamId"] = Value::String(self.team_id().await?.0),
}
(
graphql::DOCUMENT_CREATE,
json!({ "input": input }),
MutationRoot::DocumentCreate,
)
}
};
let data = self.send(query, variables).await?;
let document = mutation_payload(&data, root)?
.get("document")
.ok_or_else(|| SourceError::Malformed {
message: format!("missing {}.document", root.as_str()),
})?;
Ok(NativeId(backend_id(document, "id")?.into()))
}
async fn delete_document(&self, id: &NativeId) -> Result<(), SourceError> {
if self.get_document(id).await?.is_none() {
return Ok(());
}
let data = self
.send(graphql::DOCUMENT_DELETE, json!({"id":id.0}))
.await?;
mutation_payload(&data, MutationRoot::DocumentDelete)?;
Ok(())
}
async fn task_comments(
&self,
task: &NativeId,
page: &PageRequest,
) -> Result<Option<Page<Comment>>, SourceError> {
if page.limit == 0 {
return Err(SourceError::Config {
message: "a page limit of 0 is not a page; ask for at least 1 comment".to_owned(),
});
}
let d = self
.send(
graphql::ISSUE_COMMENTS,
json!({"id":task.0,"last":page.limit.min(MAX_PAGE_SIZE),"before":page.cursor.as_ref().map(|c|&c.0)}),
)
.await?;
optional(&d, "issue", comment_page)
}
async fn add_comment(
&self,
task: &NativeId,
comment: &NewComment,
) -> Result<Option<Comment>, SourceError> {
if let Some(author) = &comment.author {
return Err(SourceError::Refused {
message: format!(
"source {} cannot post a comment as {author:?}, because Linear records the \
user whose API key makes the request as the author of every comment; \
leave --author out to post as that user",
self.name
),
});
}
let Some(issue) = self.commented_issue(task).await? else {
return Ok(None);
};
let data = self
.send(
graphql::COMMENT_CREATE,
json!({"input":{"issueId":issue.0,"body":comment.body.as_str()}}),
)
.await?;
written_comment(&data, MutationRoot::CommentCreate).map(Some)
}
async fn edit_comment(
&self,
task: &NativeId,
comment: &NativeId,
body: &CommentBody,
) -> Result<Option<Comment>, SourceError> {
if !self.comment_is_on(task, comment).await? {
return Ok(None);
}
let data = self
.send(
graphql::COMMENT_UPDATE,
json!({"id":comment.0,"input":{"body":body.as_str()}}),
)
.await?;
written_comment(&data, MutationRoot::CommentUpdate).map(Some)
}
async fn delete_comment(
&self,
task: &NativeId,
comment: &NativeId,
) -> Result<Option<NativeId>, SourceError> {
if !self.comment_is_on(task, comment).await? {
return Ok(None);
}
let data = self
.send(graphql::COMMENT_DELETE, json!({"id":comment.0}))
.await?;
mutation_payload(&data, MutationRoot::CommentDelete)?;
Ok(Some(comment.clone()))
}
async fn set_task_status(
&self,
id: &NativeId,
category: StatusCategory,
) -> Result<Option<Status>, SourceError> {
self.writable(category)?;
let Some(task) = self.get_task(id).await? else {
return Ok(None);
};
if task.status.category == category {
return Ok(Some(task.status));
}
let (state_id, name) = self.resolve_state(category, id).await?;
let data = self
.send(
graphql::ISSUE_UPDATE,
json!({"id":task.id.0,"input":{"stateId":state_id.0}}),
)
.await?;
let issue = mutation_payload(&data, MutationRoot::IssueUpdate)?
.get("issue")
.ok_or_else(|| SourceError::Malformed {
message: "missing issueUpdate.issue".into(),
})?;
backend_id(issue, "id")?;
Ok(Some(Status { category, name }))
}
async fn set_task_priority(
&self,
id: &NativeId,
priority: Priority,
) -> Result<Option<Priority>, SourceError> {
let Some(task) = self.get_task(id).await? else {
return Ok(None);
};
let data = self
.send(
graphql::ISSUE_PRIORITY_UPDATE,
json!({"id":task.id.0,"input":{"priority":linear_priority(priority)}}),
)
.await?;
let issue = mutation_payload(&data, MutationRoot::IssueUpdate)?
.get("issue")
.filter(|issue| !issue.is_null())
.ok_or_else(|| SourceError::Malformed {
message: "missing issueUpdate.issue".into(),
})?;
written_is(issue, &task.id)?;
issue_priority(issue).map(Some)
}
async fn set_task_content(
&self,
id: &NativeId,
content: &str,
) -> Result<Option<()>, SourceError> {
let Some((task, description)) = self.issue_held(id).await? else {
return Ok(None);
};
let issue = task.id;
let slot = match description.as_deref() {
Some(description) => metadata_slot(description)?.map(str::to_owned),
None => None,
};
let description = match &slot {
None => content.to_owned(),
Some(slot) if content.is_empty() => slot.clone(),
Some(slot) => format!("{content}\n\n{slot}"),
};
if metadata_slot(&description)? != slot.as_deref() {
return Err(SourceError::Refused {
message: format!(
"this content ends in what source {} reads as its own metadata slot, so part \
of it would read back as metadata rather than as content; next: remove that \
trailing block from the content",
self.name
),
});
}
let (reads, _) = metadata_description(Some(description.clone()))?;
if reads.as_deref().unwrap_or_default() != content {
return Err(SourceError::Refused {
message: format!(
"this content would read back from source {} as {:?} rather than as itself; \
next: change how the content ends",
self.name,
reads.as_deref().unwrap_or_default()
),
});
}
let data = self
.send(
graphql::ISSUE_UPDATE,
json!({"id":issue.0,"input":{"description":description}}),
)
.await?;
let written = mutation_payload(&data, MutationRoot::IssueUpdate)?
.get("issue")
.filter(|issue| !issue.is_null())
.ok_or_else(|| SourceError::Malformed {
message: "missing issueUpdate.issue".into(),
})?;
written_is(written, &issue)?;
Ok(Some(()))
}
async fn set_delivered_by(
&self,
id: &NativeId,
delivered_by: &[TaskRef],
) -> Result<Option<()>, SourceError> {
let entries = TaskRef::listed(
TaskRef::DELIVERED_BY_KEY,
id,
Some(&self.name),
delivered_by.to_vec(),
)
.map_err(|message| SourceError::Refused { message })?;
let Some((task, description)) = self.issue_held(id).await? else {
return Ok(None);
};
let (_, held) = metadata_description(description.clone())?;
let mut slot = held.clone();
set_task_list(&mut slot, TaskRef::DELIVERED_BY_KEY, &entries);
if slot != held {
let rewritten = reslotted(description.as_deref(), &slot)?;
self.write_description_alone(&task.id, rewritten.as_deref())
.await?;
}
Ok(Some(()))
}
async fn set_task_metadata(
&self,
id: &NativeId,
key: &MetadataKey,
value: &Value,
) -> Result<Option<Task>, SourceError> {
let Some((task, description)) = self.issue_held(id).await? else {
return Ok(None);
};
let (_, mut slot) = metadata_description(description.clone())?;
if slot.get(key.as_str()) == Some(value) {
return Ok(Some(task));
}
slot.insert(key.as_str().to_owned(), value.clone());
let rewritten = reslotted(description.as_deref(), &slot)?;
self.write_description_alone(&task.id, rewritten.as_deref())
.await?;
self.get_task(&task.id)
.await?
.map(Some)
.ok_or_else(|| SourceError::Malformed {
message: format!(
"task {} was written and then could not be read back",
task.id
),
})
}
async fn set_project_metadata(
&self,
id: &NativeId,
key: &MetadataKey,
value: &Value,
) -> Result<Option<Project>, SourceError> {
let Some((project, description)) = self.project_held(id).await? else {
return Ok(None);
};
let (_, mut slot) = metadata_description(description.clone())?;
if slot.get(key.as_str()) == Some(value) {
return Ok(Some(project));
}
slot.insert(key.as_str().to_owned(), value.clone());
let rewritten = reslotted(description.as_deref(), &slot)?;
self.write_project_description(&project.id, rewritten.as_deref())
.await?;
self.get_project(&project.id)
.await?
.map(Some)
.ok_or_else(|| SourceError::Malformed {
message: format!(
"project {} was written and then could not be read back",
project.id
),
})
}
async fn set_document_metadata(
&self,
id: &NativeId,
key: &MetadataKey,
value: &Value,
) -> Result<Option<Document>, SourceError> {
let Some((document, content)) = self.document_held(id).await? else {
return Ok(None);
};
let (_, mut slot) = metadata_description(content.clone())?;
if slot.get(key.as_str()) == Some(value) {
return Ok(Some(document));
}
slot.insert(key.as_str().to_owned(), value.clone());
let rewritten = reslotted(content.as_deref(), &slot)?;
self.write_document_content(&document.id, rewritten.as_deref())
.await?;
self.get_document(&document.id)
.await?
.map(Some)
.ok_or_else(|| SourceError::Malformed {
message: format!(
"document {} was written and then could not be read back",
document.id
),
})
}
async fn set_task_rendering(
&self,
id: &NativeId,
content: &str,
provenance: &Value,
_answers: &std::collections::BTreeMap<String, Value>,
) -> Result<Option<()>, SourceError> {
let Some((task, description)) = self.issue_held(id).await? else {
return Ok(None);
};
let (_, mut slot) = metadata_description(description.clone())?;
slot.insert(MetadataKey::TEMPLATE_KEY.to_owned(), provenance.clone());
let rewritten = Self::described(Some(content), &slot)?;
if rewritten != description {
self.write_description_alone(&task.id, rewritten.as_deref())
.await?;
}
Ok(Some(()))
}
async fn set_document_rendering(
&self,
id: &NativeId,
content: &str,
provenance: &Value,
_answers: &std::collections::BTreeMap<String, Value>,
) -> Result<Option<()>, SourceError> {
let Some((document, held)) = self.document_held(id).await? else {
return Ok(None);
};
let (_, mut slot) = metadata_description(held.clone())?;
slot.insert(MetadataKey::TEMPLATE_KEY.to_owned(), provenance.clone());
let rewritten = Self::described(Some(content), &slot)?;
if rewritten != held {
self.write_document_content(&document.id, rewritten.as_deref())
.await?;
}
Ok(Some(()))
}
async fn update_task(
&self,
id: &NativeId,
update: &TaskUpdate,
) -> Result<Option<TaskUpdateOutcome>, SourceError> {
self.targeted_update(id, update).await
}
}
fn acknowledged(item: &Value, asked: &NativeId, root: MutationRoot) -> Result<(), SourceError> {
let written = backend_id(item, "id")?;
if written == asked.0 {
return Ok(());
}
Err(SourceError::Malformed {
message: format!(
"{} for {asked} answered with the item {written}",
root.as_str()
),
})
}
fn written_is(issue: &Value, asked: &NativeId) -> Result<(), SourceError> {
let written = backend_id(issue, "id")?;
if written == asked.0 {
return Ok(());
}
Err(SourceError::Malformed {
message: format!("issueUpdate for {asked} answered with the issue {written}"),
})
}
const NO_DELIVERY: &str = "only a task delivers or is delivered, so neither list has a place \
on anything else";
fn delivery_key_in(metadata: &std::collections::BTreeMap<String, Value>) -> Option<&'static str> {
[TaskRef::DELIVERS_KEY, TaskRef::DELIVERED_BY_KEY]
.into_iter()
.find(|key| metadata.contains_key(*key))
}
fn category_word(category: StatusCategory) -> String {
serde_json::to_value(category)
.ok()
.and_then(|value| value.as_str().map(str::to_owned))
.unwrap_or_else(|| format!("{category:?}"))
}
impl LinearSource {
fn writable(&self, category: StatusCategory) -> Result<(), SourceError> {
match self.statuses.target(category) {
StateTarget::DisabledByMapping => Err(self.disabled(category, true)),
StateTarget::NoWorkflowType => Err(self.disabled(category, false)),
StateTarget::Named(_) | StateTarget::FirstOfType(_) => Ok(()),
}
}
fn disabled(&self, category: StatusCategory, configured: bool) -> SourceError {
let word = category_word(category);
SourceError::Refused {
message: if configured {
format!(
"source {} cannot set a task's status to {word}: that category is disabled \
for this source, because its status_mapping sets {word} to null; map \
status_mapping.{word} to one of the team's workflow states to write it",
self.name
)
} else {
format!(
"source {} cannot set a task's status to {word}: that category is disabled \
for this source, because Linear has no workflow state of that kind — its \
workflow states are triage, backlog, unstarted, started, completed and \
canceled; choose backlog, todo, in-progress, done or cancelled, or map \
status_mapping.{word} to one of the team's workflow states",
self.name
)
},
}
}
async fn state_of(
&self,
category: StatusCategory,
state_type: &str,
task: &NativeId,
) -> Result<(NativeId, String), SourceError> {
let team = self.team_id().await?;
let data = self
.send(
graphql::ISSUE_STATE_OF_TYPE,
json!({"type":state_type,"team":team.0}),
)
.await?;
let nodes = data
.get("workflowStates")
.and_then(|v| v.get("nodes"))
.and_then(Value::as_array)
.ok_or_else(|| SourceError::Malformed {
message: "missing workflowStates.nodes".into(),
})?;
let Some(state) = nodes.first() else {
return Err(SourceError::Refused {
message: format!(
"source {} cannot set task {} to {}: its configured team has no workflow \
state of type {state_type}; add one to the team in Linear",
self.name,
task.0,
category_word(category)
),
});
};
Ok((
NativeId(backend_id(state, "id")?.into()),
str_at(state, "name")?.to_owned(),
))
}
fn undeliverable(&self, named: &str, what: &str) -> SourceError {
SourceError::Refused {
message: format!(
"source {} cannot carry {named} on a {what}: {NO_DELIVERY}; write the {what} \
without it",
self.name
),
}
}
async fn commented_issue(&self, task: &NativeId) -> Result<Option<NativeId>, SourceError> {
Ok(self.get_task(task).await?.map(|task| task.id))
}
async fn comment_is_on(
&self,
task: &NativeId,
comment: &NativeId,
) -> Result<bool, SourceError> {
let Some(issue) = self.commented_issue(task).await? else {
return Ok(false);
};
let data = self.send(graphql::COMMENT, json!({"id":comment.0})).await?;
Ok(optional(&data, "comment", comment_issue)?.flatten() == Some(issue))
}
}
impl LinearSource {
async fn issue_held(
&self,
id: &NativeId,
) -> Result<Option<(Task, Option<String>)>, SourceError> {
let data = self.send(ISSUE, json!({"id":id.0})).await?;
Ok(optional(&data, "issue", |v| {
Ok((
map_task(v, &self.name, &self.statuses)?,
optional_string(v, "description")?,
))
})?
.filter(|(task, _)| self.in_scope(task.project.as_ref())))
}
async fn project_held(
&self,
id: &NativeId,
) -> Result<Option<(Project, Option<String>)>, SourceError> {
let data = self.send(PROJECT, json!({"id":id.0})).await?;
Ok(optional(&data, "project", |v| {
Ok((map_project(v)?, optional_string(v, "description")?))
})?
.filter(|(project, _)| self.in_scope(Some(&project.id))))
}
async fn document_held(
&self,
id: &NativeId,
) -> Result<Option<(Document, Option<String>)>, SourceError> {
let data = self.send(DOCUMENT, json!({"id":id.0})).await?;
Ok(optional(&data, "document", |v| {
Ok((map_document(v)?, optional_string(v, "content")?))
})?
.filter(|(document, _)| self.in_scope(document.project.as_ref())))
}
async fn write_description_alone(
&self,
id: &NativeId,
description: Option<&str>,
) -> Result<(), SourceError> {
let data = self
.send(
graphql::ISSUE_UPDATE,
json!({"id":id.0,"input":{"description":description}}),
)
.await?;
let issue = mutation_payload(&data, MutationRoot::IssueUpdate)?
.get("issue")
.filter(|issue| !issue.is_null())
.ok_or_else(|| SourceError::Malformed {
message: "missing issueUpdate.issue".into(),
})?;
written_is(issue, id)
}
async fn write_project_description(
&self,
id: &NativeId,
description: Option<&str>,
) -> Result<(), SourceError> {
let data = self
.send(
graphql::PROJECT_UPDATE,
json!({"id":id.0,"input":{"description":description}}),
)
.await?;
let project = mutation_payload(&data, MutationRoot::ProjectUpdate)?
.get("project")
.filter(|project| !project.is_null())
.ok_or_else(|| SourceError::Malformed {
message: "missing projectUpdate.project".into(),
})?;
acknowledged(project, id, MutationRoot::ProjectUpdate)
}
async fn write_document_content(
&self,
id: &NativeId,
content: Option<&str>,
) -> Result<(), SourceError> {
let data = self
.send(
graphql::DOCUMENT_UPDATE,
json!({"id":id.0,"input":{"content":content}}),
)
.await?;
let document = mutation_payload(&data, MutationRoot::DocumentUpdate)?
.get("document")
.filter(|document| !document.is_null())
.ok_or_else(|| SourceError::Malformed {
message: "missing documentUpdate.document".into(),
})?;
acknowledged(document, id, MutationRoot::DocumentUpdate)
}
async fn team_states(&self) -> Result<Vec<TeamState>, SourceError> {
let team = self.team_id().await?;
let data = self
.send(graphql::TEAM_WORKFLOW_STATES, json!({"team":team.0}))
.await?;
data.get("workflowStates")
.and_then(|v| v.get("nodes"))
.and_then(Value::as_array)
.ok_or_else(|| SourceError::Malformed {
message: "missing workflowStates.nodes".into(),
})?
.iter()
.map(|node| {
Ok(TeamState {
id: NativeId(backend_id(node, "id")?.into()),
name: StateName::try_from(str_at(node, "name")?.to_owned())
.map_err(|message| SourceError::Malformed { message })?,
kind: str_at(node, "type")?.to_owned(),
})
})
.collect()
}
async fn named_state(
&self,
name: &str,
category: StatusCategory,
) -> Result<(NativeId, String), SourceError> {
self.team_states()
.await?
.into_iter()
.find(|state| state.name.0.eq_ignore_ascii_case(name))
.map(|state| (state.id, state.name.0))
.ok_or_else(|| SourceError::Refused {
message: format!(
"source {} maps status {} to the workflow state {name:?}, which team {} \
does not have; next: add {name:?} to that team in Linear's team settings, \
or point status_mapping.{} of this source at a state the team has",
self.name,
category_word(category),
self.team.as_ref().map_or("(none)", |team| team.0.as_str()),
category_word(category)
),
})
}
async fn resolve_state(
&self,
category: StatusCategory,
task: &NativeId,
) -> Result<(NativeId, String), SourceError> {
match self.statuses.target(category) {
StateTarget::Named(name) => self.named_state(&name.0, category).await,
StateTarget::FirstOfType(kind) => self.state_of(category, kind.as_str(), task).await,
StateTarget::DisabledByMapping => Err(self.disabled(category, true)),
StateTarget::NoWorkflowType => Err(self.disabled(category, false)),
}
}
async fn workflow_states(&self) -> Result<WorkflowStatesReport, SourceError> {
let held = self.team_states().await?;
let states = self
.statuses
.named()
.map(|(category, name)| {
let found = held
.iter()
.find(|state| state.name.0.eq_ignore_ascii_case(&name.0));
MappedWorkflowState {
category,
state: name.clone(),
found: found.map_or(Found::Missing, |state| Found::Present(state.kind.clone())),
}
})
.collect();
let team = self.team.clone().ok_or_else(|| SourceError::Refused {
message: format!(
"source {} needs config.team to report its states",
self.name
),
})?;
Ok(WorkflowStatesReport {
source: self.name.clone(),
team,
states,
})
}
}
struct TeamState {
id: NativeId,
name: StateName,
kind: String,
}
fn reslotted(
held: Option<&str>,
slot: &std::collections::BTreeMap<String, Value>,
) -> Result<Option<String>, SourceError> {
let held = held.unwrap_or_default();
let Some((start, _, _)) = slot_bounds(held)? else {
return LinearSource::described((!held.is_empty()).then_some(held), slot);
};
let above = &held[..start];
if slot.is_empty() {
let visible = above
.strip_suffix("\n\n")
.or_else(|| above.strip_suffix('\n'))
.unwrap_or(above);
return Ok((!visible.is_empty()).then(|| visible.to_owned()));
}
Ok(Some(format!("{above}{}", slot_text(slot)?)))
}
fn set_task_list(
metadata: &mut std::collections::BTreeMap<String, Value>,
key: &str,
entries: &[TaskRef],
) {
if entries.is_empty() {
metadata.remove(key);
} else {
metadata.insert(
key.to_owned(),
Value::Array(
entries
.iter()
.map(|entry| Value::String(entry.as_str().to_owned()))
.collect(),
),
);
}
}
fn comment_page(v: &Value) -> Result<Page<Comment>, SourceError> {
let c = v.get("comments").ok_or_else(|| SourceError::Malformed {
message: "missing comments connection".into(),
})?;
let mut items = c
.get("nodes")
.and_then(Value::as_array)
.ok_or_else(|| SourceError::Malformed {
message: "missing comment nodes".into(),
})?
.iter()
.map(map_comment)
.collect::<Result<Vec<_>, _>>()?;
items.reverse();
let info = c.get("pageInfo").ok_or_else(|| SourceError::Malformed {
message: "missing pageInfo".into(),
})?;
let older = info
.get("hasPreviousPage")
.and_then(Value::as_bool)
.ok_or_else(|| SourceError::Malformed {
message: "missing boolean pageInfo.hasPreviousPage".into(),
})?;
let next = if older {
Some(Cursor(str_at(info, "startCursor")?.into()))
} else {
None
};
Ok(Page { items, next })
}
fn map_comment(v: &Value) -> Result<Comment, SourceError> {
let author = match v.get("user") {
None => {
return Err(SourceError::Malformed {
message: "missing comment user field".into(),
});
}
Some(Value::Null) => None,
Some(user) => Some(str_at(user, "displayName")?.to_owned()),
};
Ok(Comment {
id: NativeId(backend_id(v, "id")?.into()),
author,
created_at: time(v, "createdAt")?,
updated_at: time(v, "updatedAt")?,
body: str_at(v, "body")?.into(),
url: optional_string(v, "url")?,
})
}
fn written_comment(data: &Value, root: MutationRoot) -> Result<Comment, SourceError> {
let comment = mutation_payload(data, root)?
.get("comment")
.ok_or_else(|| SourceError::Malformed {
message: format!("missing {}.comment", root.as_str()),
})?;
map_comment(comment)
}
fn comment_issue(v: &Value) -> Result<Option<NativeId>, SourceError> {
match v.get("issue") {
None => Err(SourceError::Malformed {
message: "missing comment issue field".into(),
}),
Some(Value::Null) => Ok(None),
Some(issue) => Ok(Some(NativeId(backend_id(issue, "id")?.into()))),
}
}
const RECORDED_CURSOR: &str = "onetaskgraph.depends_on:";
impl LinearSource {
async fn dependencies(
&self,
query: &str,
root: DependencyRoot,
id: &NativeId,
direction: Direction,
page: &PageRequest,
) -> Result<Page<DependencyEdge>, SourceError> {
let limit = page.limit.min(MAX_PAGE_SIZE);
let cursor = page.cursor.as_ref().map(|c| c.0.as_str());
if let Some(offset) = cursor.and_then(|c| c.strip_prefix(RECORDED_CURSOR)) {
if direction != Direction::DependsOn {
return Err(SourceError::Malformed {
message: format!(
"{RECORDED_CURSOR}{offset} resumes recorded forward edges, which a reverse dependency read never issues; resume it in the direction that reported it"
),
});
}
let offset: usize = offset.parse().map_err(|_| SourceError::Malformed {
message: format!("{RECORDED_CURSOR}{offset} is not a recorded-edge cursor"),
})?;
let d = self
.send(query, json!({"id":id.0,"first":1,"after":null}))
.await?;
return Ok(recorded_page(
recorded(&d, root, id, &self.name)?,
offset,
limit as usize,
));
}
let d = self
.send(query, json!({"id":id.0,"first":limit,"after":cursor}))
.await?;
let mut answered = relation_page(&d, root, id, direction)?;
if answered.next.is_none()
&& direction == Direction::DependsOn
&& !recorded(&d, root, id, &self.name)?.is_empty()
{
answered.next = Some(Cursor(format!("{RECORDED_CURSOR}0")));
}
Ok(answered)
}
}
fn recorded(
d: &Value,
root: DependencyRoot,
id: &NativeId,
name: &SourceName,
) -> Result<Vec<DependencyEdge>, SourceError> {
let item = d.get(root.as_str()).ok_or_else(|| SourceError::Malformed {
message: format!("missing {}", root.as_str()),
})?;
let (_, metadata) = metadata_description(optional_string(item, "description")?)?;
DependencyEdge::recorded(
&metadata,
id,
root.item_kind(),
name,
Some(root.item_kind()),
)
.map_err(|message| SourceError::Malformed { message })
}
fn recorded_page(edges: Vec<DependencyEdge>, offset: usize, limit: usize) -> Page<DependencyEdge> {
let total = edges.len();
let items: Vec<DependencyEdge> = edges.into_iter().skip(offset).take(limit.max(1)).collect();
let end = offset.saturating_add(items.len());
Page {
items,
next: (end < total).then(|| Cursor(format!("{RECORDED_CURSOR}{end}"))),
}
}
fn workflow_state_types(s: &StatusCategory) -> Vec<&'static str> {
WorkflowType::of(*s)
.map(WorkflowType::as_str)
.into_iter()
.collect()
}
fn project_status_types(s: &StatusCategory) -> Vec<&'static str> {
match s {
StatusCategory::Draft => vec![],
StatusCategory::Backlog => vec!["backlog"],
StatusCategory::Todo => vec!["planned"],
StatusCategory::Queued => vec![],
StatusCategory::InProgress => vec!["started", "paused"],
StatusCategory::Done => vec!["completed"],
StatusCategory::Cancelled => vec!["canceled"],
StatusCategory::Unknown => vec![],
}
}
fn status(v: &Value) -> Result<Status, SourceError> {
let name = str_at(v, "name")?.into();
let category = match str_at(v, "type")? {
"backlog" => StatusCategory::Backlog,
"unstarted" | "planned" => StatusCategory::Todo,
"started" | "paused" => StatusCategory::InProgress,
"completed" => StatusCategory::Done,
"canceled" => StatusCategory::Cancelled,
_ => StatusCategory::Unknown,
};
Ok(Status { category, name })
}
fn issue_status(v: &Value, statuses: &StatusMapping) -> Result<Status, SourceError> {
let name = str_at(v, "name")?;
match statuses.category_of(name) {
Some(category) => Ok(Status {
category,
name: name.to_owned(),
}),
None => status(v),
}
}
fn str_at<'a>(v: &'a Value, k: &str) -> Result<&'a str, SourceError> {
v.get(k)
.and_then(Value::as_str)
.ok_or_else(|| SourceError::Malformed {
message: format!("missing string field {k}"),
})
}
fn map_label(v: &Value) -> Result<Label, SourceError> {
Ok(Label {
id: NativeId(str_at(v, "id")?.into()),
name: str_at(v, "name")?.into(),
color: optional_string(v, "color")?,
})
}
fn labels_of(v: &Value) -> Result<Vec<Label>, SourceError> {
v.get("nodes")
.and_then(Value::as_array)
.ok_or_else(|| SourceError::Malformed {
message: "missing label nodes".into(),
})?
.iter()
.map(map_label)
.collect()
}
fn time(v: &Value, k: &str) -> Result<Option<DateTime<Utc>>, SourceError> {
optional_str(v, k)?
.map(|s| {
s.parse().map_err(|e| SourceError::Malformed {
message: format!("invalid {k}: {e}"),
})
})
.transpose()
}
fn map_task(v: &Value, source: &SourceName, statuses: &StatusMapping) -> Result<Task, SourceError> {
let (content, mut metadata) = metadata_description(optional_string(v, "description")?)?;
let repositories = Repository::from_metadata(&metadata)
.map_err(|message| SourceError::Malformed { message })?;
let url = optional_string(v, "url")?;
let id = NativeId(str_at(v, "id")?.into());
let delivers = delivery_list(&mut metadata, TaskRef::DELIVERS_KEY, &id, source)?;
let delivered_by = delivery_list(&mut metadata, TaskRef::DELIVERED_BY_KEY, &id, source)?;
Ok(Task {
id,
key: Some(str_at(v, "identifier")?.into()),
title: str_at(v, "title")?.into(),
content,
status: issue_status(
v.get("state").ok_or_else(|| SourceError::Malformed {
message: "missing state".into(),
})?,
statuses,
)?,
priority: issue_priority(v)?,
labels: labels_of(v.get("labels").ok_or_else(|| SourceError::Malformed {
message: "missing labels".into(),
})?)?,
project: filed_under(v)?,
location: web_address(url.as_deref()),
url,
created_at: time(v, "createdAt")?,
updated_at: time(v, "updatedAt")?,
metadata,
repositories,
delivers,
delivered_by,
})
}
const fn linear_priority(priority: Priority) -> u8 {
match priority {
Priority::None => 0,
Priority::Urgent => 1,
Priority::High => 2,
Priority::Medium => 3,
Priority::Low => 4,
}
}
fn issue_priority(v: &Value) -> Result<Priority, SourceError> {
let raw = v.get("priority").ok_or_else(|| SourceError::Malformed {
message: "missing number field priority".into(),
})?;
let level = raw.as_f64().filter(|level| level.fract() == 0.0);
Priority::ALL
.into_iter()
.find(|priority| level == Some(f64::from(linear_priority(*priority))))
.ok_or_else(|| SourceError::Malformed {
message: format!(
"field priority is {raw}, which is none of Linear's priorities 0 (none), \
1 (urgent), 2 (high), 3 (normal) and 4 (low)"
),
})
}
fn delivery_list(
metadata: &mut std::collections::BTreeMap<String, Value>,
key: &str,
task: &NativeId,
source: &SourceName,
) -> Result<Vec<TaskRef>, SourceError> {
let held = metadata.remove(key);
TaskRef::from_value(key, task, Some(source), held.as_ref())
.map_err(|message| SourceError::Malformed { message })
}
fn strip_delivery_keys(metadata: &mut std::collections::BTreeMap<String, Value>) {
metadata.remove(TaskRef::DELIVERS_KEY);
metadata.remove(TaskRef::DELIVERED_BY_KEY);
}
fn map_project(v: &Value) -> Result<Project, SourceError> {
let (content, mut metadata) = metadata_description(optional_string(v, "description")?)?;
strip_delivery_keys(&mut metadata);
let repositories = Repository::from_metadata(&metadata)
.map_err(|message| SourceError::Malformed { message })?;
let url = optional_string(v, "url")?;
Ok(Project {
id: NativeId(str_at(v, "id")?.into()),
title: str_at(v, "name")?.into(),
content,
status: status(v.get("status").ok_or_else(|| SourceError::Malformed {
message: "missing status".into(),
})?)?,
labels: labels_of(v.get("labels").ok_or_else(|| SourceError::Malformed {
message: "missing project labels".into(),
})?)?,
location: web_address(url.as_deref()),
url,
created_at: time(v, "createdAt")?,
updated_at: time(v, "updatedAt")?,
metadata,
repositories,
})
}
fn web_address(url: Option<&str>) -> Option<Location> {
url.map(|url| Location::Url(url.to_owned()))
}
fn filed_under(v: &Value) -> Result<Option<NativeId>, SourceError> {
match v.get("project") {
None => Err(SourceError::Malformed {
message: "missing project field".into(),
}),
Some(Value::Null) => Ok(None),
Some(project) => Ok(Some(NativeId(str_at(project, "id")?.into()))),
}
}
fn map_document(v: &Value) -> Result<Document, SourceError> {
let (content, mut metadata) = metadata_description(optional_string(v, "content")?)?;
strip_delivery_keys(&mut metadata);
let repositories = Repository::from_metadata(&metadata)
.map_err(|message| SourceError::Malformed { message })?;
let url = optional_string(v, "url")?;
Ok(Document {
id: NativeId(str_at(v, "id")?.into()),
title: str_at(v, "title")?.into(),
content,
project: filed_under(v)?,
labels: Vec::new(),
location: web_address(url.as_deref()),
url,
created_at: time(v, "createdAt")?,
updated_at: time(v, "updatedAt")?,
metadata,
repositories,
})
}
fn document_matches(document: &Document, project: &ProjectFilter, labels: &LabelFilter) -> bool {
let carries = |name: &String| {
document
.labels
.iter()
.any(|label| label.name.eq_ignore_ascii_case(name))
};
let filed = match project {
ProjectFilter::Any => true,
ProjectFilter::Orphans => document.project.is_none(),
ProjectFilter::Is(id) => document.project.as_ref() == Some(id),
};
filed
&& (labels.any_of.is_empty() || labels.any_of.iter().any(&carries))
&& labels.all_of.iter().all(&carries)
&& !labels.none_of.iter().any(&carries)
}
fn optional<T>(
d: &Value,
k: &str,
f: impl Fn(&Value) -> Result<T, SourceError>,
) -> Result<Option<T>, SourceError> {
match d.get(k) {
None => Err(SourceError::Malformed {
message: format!("missing {k}"),
}),
Some(Value::Null) => Ok(None),
Some(value) if !matches!(value.get("archivedAt"), None | Some(Value::Null)) => Ok(None),
Some(value) => f(value).map(Some),
}
}
fn connection<T>(
d: &Value,
k: &str,
f: impl Fn(&Value) -> Result<T, SourceError>,
) -> Result<Page<T>, SourceError> {
let c = d.get(k).ok_or_else(|| SourceError::Malformed {
message: format!("missing {k} connection"),
})?;
let items = c
.get("nodes")
.and_then(Value::as_array)
.ok_or_else(|| SourceError::Malformed {
message: "missing nodes".into(),
})?
.iter()
.map(f)
.collect::<Result<_, _>>()?;
let next = page_next(c)?;
Ok(Page { items, next })
}
#[derive(Clone, Copy)]
enum DependencyRoot {
Issue,
Project,
}
impl DependencyRoot {
const fn item_kind(self) -> ItemKind {
match self {
Self::Issue => ItemKind::Task,
Self::Project => ItemKind::Project,
}
}
const fn as_str(self) -> &'static str {
match self {
Self::Issue => "issue",
Self::Project => "project",
}
}
}
fn relation_page(
d: &Value,
root: DependencyRoot,
id: &NativeId,
direction: Direction,
) -> Result<Page<DependencyEdge>, SourceError> {
let key = if direction == Direction::DependsOn {
"relations"
} else {
"inverseRelations"
};
let c = d
.get(root.as_str())
.and_then(|v| v.get(key))
.ok_or_else(|| SourceError::Malformed {
message: format!("missing {key}"),
})?;
let nodes = c
.get("nodes")
.and_then(Value::as_array)
.ok_or_else(|| SourceError::Malformed {
message: "missing relation nodes".into(),
})?;
let mut items = Vec::new();
for n in nodes {
let other = n
.get(if direction == Direction::DependsOn {
"relatedIssue"
} else {
"issue"
})
.or_else(|| {
n.get(if direction == Direction::DependsOn {
"relatedProject"
} else {
"project"
})
})
.and_then(|v| v.get("id"))
.and_then(Value::as_str)
.ok_or_else(|| SourceError::Malformed {
message: "missing related id".into(),
})?;
let (from, to) = if direction == Direction::DependsOn {
(id.clone(), NativeId(other.into()))
} else {
(NativeId(other.into()), id.clone())
};
let relation_type =
n.get("type")
.and_then(Value::as_str)
.ok_or_else(|| SourceError::Malformed {
message: "missing relation type".into(),
})?;
let kind = match (root, relation_type) {
(DependencyRoot::Issue, "blocks") | (DependencyRoot::Project, "dependency") => {
DependencyKind::Blocks
}
(DependencyRoot::Issue, "related") => DependencyKind::Related,
_ => {
return Err(SourceError::Malformed {
message: format!(
"invalid relation type: {relation_type} on a {} relation",
root.as_str()
),
});
}
};
let item_kind = root.item_kind();
items.push(DependencyEdge {
from: DependencyEndpoint::from_native(from, item_kind),
to: DependencyEndpoint::from_native(to, item_kind),
kind,
});
}
let next = page_next(c)?;
Ok(Page { items, next })
}
fn optional_str<'a>(v: &'a Value, k: &str) -> Result<Option<&'a str>, SourceError> {
match v.get(k) {
None => Err(SourceError::Malformed {
message: format!("missing field {k}"),
}),
Some(Value::Null) => Ok(None),
Some(value) => value
.as_str()
.map(Some)
.ok_or_else(|| SourceError::Malformed {
message: format!("field {k} is not a string"),
}),
}
}
const METADATA_PREFIX: &str = "<!-- onetaskgraph.metadata";
const METADATA_OPEN_SPAN: &str = "<!-- onetaskgraph.metadata `";
const METADATA_CLOSE_SPAN: &str = "` -->";
const METADATA_OPEN: &str = "<!-- onetaskgraph.metadata\n";
const METADATA_CLOSE: &str = "\n-->";
const METADATA_CLOSE_ESCAPED: &str = "\n\\-->";
fn slot_json(value: &impl serde::Serialize) -> Result<String, SourceError> {
let encoded = serde_json::to_string(value).map_err(|error| SourceError::Malformed {
message: error.to_string(),
})?;
Ok(encoded
.replace('<', "\\u003c")
.replace('>', "\\u003e")
.replace('`', "\\u0060"))
}
fn slot_text(metadata: &std::collections::BTreeMap<String, Value>) -> Result<String, SourceError> {
Ok(format!(
"{METADATA_OPEN_SPAN}{}{METADATA_CLOSE_SPAN}",
slot_json(metadata)?
))
}
fn slot_bounds(description: &str) -> Result<Option<(usize, usize, usize)>, SourceError> {
let Some(start) = description.rfind(METADATA_PREFIX) else {
return Ok(None);
};
let rest = &description[start..];
let (encoded_start, close) = if rest.starts_with(METADATA_OPEN_SPAN) {
let encoded_start = start + METADATA_OPEN_SPAN.len();
(
encoded_start,
description[encoded_start..]
.rfind(METADATA_CLOSE_SPAN)
.map(|at| (at, METADATA_CLOSE_SPAN.len())),
)
} else if rest.starts_with(METADATA_OPEN) {
let encoded_start = start + METADATA_OPEN.len();
(
encoded_start,
[METADATA_CLOSE, METADATA_CLOSE_ESCAPED]
.into_iter()
.filter_map(|close| {
description[encoded_start..]
.find(close)
.map(|at| (at, close.len()))
})
.min(),
)
} else {
return Ok(None);
};
let Some((relative_end, close_len)) = close else {
return Err(SourceError::Malformed {
message: "unterminated onetaskgraph metadata slot in Linear description".into(),
});
};
let encoded_end = encoded_start + relative_end;
if !description[encoded_end + close_len..].trim().is_empty() {
return Ok(None);
}
Ok(Some((start, encoded_start, encoded_end)))
}
fn metadata_slot(description: &str) -> Result<Option<&str>, SourceError> {
Ok(slot_bounds(description)?.map(|(start, _, _)| &description[start..]))
}
fn metadata_description(
description: Option<String>,
) -> Result<(Option<String>, std::collections::BTreeMap<String, Value>), SourceError> {
let Some(description) = description else {
return Ok((None, Default::default()));
};
let Some((start, encoded_start, encoded_end)) = slot_bounds(&description)? else {
return Ok((Some(description), Default::default()));
};
let metadata =
serde_json::from_str(&description[encoded_start..encoded_end]).map_err(|error| {
SourceError::Malformed {
message: format!(
"invalid canonical JSON in Linear onetaskgraph metadata slot: {error}"
),
}
})?;
let above = &description[..start];
let visible = above
.strip_suffix("\n\n")
.or_else(|| above.strip_suffix('\n'))
.unwrap_or(above);
Ok(((!visible.is_empty()).then(|| visible.to_owned()), metadata))
}
fn slot_phrase(value: &str) -> Option<Value> {
let verbatim = value.chars().all(|character| {
(character.is_ascii_graphic() && !"\"\\/<>&'`".contains(character)) || character == ' '
});
verbatim.then(|| json!({"description": {"contains": format!("\"{value}\"")}}))
}
fn text_holds(title: &str, content: Option<&str>, query: &TextQuery) -> bool {
let terms = query.terms.to_lowercase();
let in_title = title.to_lowercase().contains(&terms);
let in_content = content.is_some_and(|body| body.to_lowercase().contains(&terms));
match query.fields {
TextFields::Title => in_title,
TextFields::Content => in_content,
TextFields::TitleOrContent => in_title || in_content,
}
}
fn optional_string(v: &Value, k: &str) -> Result<Option<String>, SourceError> {
Ok(optional_str(v, k)?.map(Into::into))
}
fn backend_id<'a>(value: &'a Value, field: &str) -> Result<&'a str, SourceError> {
let id = str_at(value, field)?;
(!id.is_empty())
.then_some(id)
.ok_or_else(|| SourceError::Malformed {
message: format!("field {field} is an empty backend id"),
})
}
fn mutation_payload(data: &Value, root: MutationRoot) -> Result<&Value, SourceError> {
let root = root.as_str();
let payload = data.get(root).ok_or_else(|| SourceError::Malformed {
message: format!("missing {root}"),
})?;
match payload.get("success").and_then(Value::as_bool) {
Some(true) => Ok(payload),
Some(false) => Err(SourceError::Refused {
message: format!("Linear reported {root} was unsuccessful"),
}),
None => Err(SourceError::Malformed {
message: format!("missing boolean {root}.success"),
}),
}
}
fn page_next(c: &Value) -> Result<Option<Cursor>, SourceError> {
let info = c.get("pageInfo").ok_or_else(|| SourceError::Malformed {
message: "missing pageInfo".into(),
})?;
let more = info
.get("hasNextPage")
.and_then(Value::as_bool)
.ok_or_else(|| SourceError::Malformed {
message: "missing boolean pageInfo.hasNextPage".into(),
})?;
if !more {
return Ok(None);
}
let cursor = str_at(info, "endCursor")?;
Ok(Some(Cursor(cursor.into())))
}