#![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, StatusMapping, StatusName, 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 ISSUE_NOT_FOUND_MESSAGE: &str = "Entity not found: Issue";
pub const ISSUE_NOT_FOUND_PRESENTABLE: &str = "Could not find referenced Issue.";
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 content 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 content 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){ content 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 RESOLUTION: &str = "query($key:String!){ teams(first:2,filter:{key:{eqIgnoreCase:$key}}){nodes{id states(first:250){nodes{id name type} pageInfo{hasNextPage}}}} projectStatuses(first:250){nodes{id name type position} pageInfo{hasNextPage}} }";
pub const WORKFLOW_STATE_CREATE: &str = "mutation($input:WorkflowStateCreateInput!){ workflowStateCreate(input:$input){success workflowState{id name type}} }";
pub const PROJECT_STATUS_CREATE: &str = "mutation($input:ProjectStatusCreateInput!){ projectStatusCreate(input:$input){success status{id name type position}} }";
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_UPDATE_READ: &str = "mutation($id:String!,$input:IssueUpdateInput!){ issueUpdate(id:$id,input:$input){success issue{ id identifier title description url createdAt updatedAt archivedAt state{name type} priority labels{nodes{id name color}} project{id} }} }";
pub const ISSUE_PRIORITY_UPDATE: &str = "mutation($id:String!,$input:IssueUpdateInput!){ issueUpdate(id:$id,input:$input){success issue{id priority}} }";
pub const ISSUE_REWRITE: &str = "mutation($id:String!,$input:IssueUpdateInput!,$first:Int!){ issueUpdate(id:$id,input:$input){success issue{id relations(first:$first){nodes{id type relatedIssue{id}} pageInfo{hasNextPage endCursor}}}} }";
pub const PROJECT_REWRITE: &str = "mutation($id:String!,$input:ProjectUpdateInput!,$first:Int!){ projectUpdate(id:$id,input:$input){success project{id relations(first:$first){nodes{id type relatedProject{id}} pageInfo{hasNextPage endCursor}}}} }";
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,
status_mapping: StatusMapping,
project: Option<ProjectScope>,
}
#[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: StatusMapping::default(),
project: None,
}
}
}
const fn kind_word(kind: ItemKind) -> &'static str {
match kind {
ItemKind::Task => "task",
ItemKind::Project => "project",
}
}
const fn vocabulary_word(kind: ItemKind) -> &'static str {
match kind {
ItemKind::Task => "workflow state",
ItemKind::Project => "project status",
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum CreatedStateType {
Backlog,
Unstarted,
Started,
Completed,
Canceled,
}
impl CreatedStateType {
const fn as_str(self) -> &'static str {
match self {
Self::Backlog => "backlog",
Self::Unstarted => "unstarted",
Self::Started => "started",
Self::Completed => "completed",
Self::Canceled => "canceled",
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum CreatedStatusType {
Backlog,
Planned,
Started,
Completed,
Canceled,
}
impl CreatedStatusType {
const fn as_str(self) -> &'static str {
match self {
Self::Backlog => "backlog",
Self::Planned => "planned",
Self::Started => "started",
Self::Completed => "completed",
Self::Canceled => "canceled",
}
}
}
const fn created_types(category: StatusCategory) -> (CreatedStateType, CreatedStatusType) {
match category {
StatusCategory::Backlog | StatusCategory::Draft => {
(CreatedStateType::Backlog, CreatedStatusType::Backlog)
}
StatusCategory::Todo | StatusCategory::Queued => {
(CreatedStateType::Unstarted, CreatedStatusType::Planned)
}
StatusCategory::InProgress | StatusCategory::Unknown => {
(CreatedStateType::Started, CreatedStatusType::Started)
}
StatusCategory::Done => (CreatedStateType::Completed, CreatedStatusType::Completed),
StatusCategory::Cancelled => (CreatedStateType::Canceled, CreatedStatusType::Canceled),
}
}
const fn created_type(category: StatusCategory, kind: ItemKind) -> &'static str {
let (state, status) = created_types(category);
match kind {
ItemKind::Task => state.as_str(),
ItemKind::Project => status.as_str(),
}
}
const CREATED_COLOR: &str = "#95a2b3";
#[derive(Debug, Clone)]
struct Held {
id: NativeId,
name: StatusName,
kind: String,
}
impl Held {
fn read(node: &Value) -> Result<Self, SourceError> {
Ok(Self {
id: NativeId(backend_id(node, "id")?.into()),
name: held_name(node)?,
kind: str_at(node, "type")?.to_owned(),
})
}
}
fn position_of(node: &Value) -> Result<f64, SourceError> {
node.get("position")
.and_then(Value::as_f64)
.ok_or_else(|| SourceError::Malformed {
message: "missing number field position".into(),
})
}
#[derive(Debug, Clone)]
enum Created {
State(Held),
Status { held: Held, position: f64 },
}
impl Created {
fn held(&self) -> &Held {
match self {
Self::State(held) | Self::Status { held, .. } => held,
}
}
}
fn held_name(node: &Value) -> Result<StatusName, SourceError> {
StatusName::try_from(str_at(node, "name")?.to_owned()).map_err(|refused| {
SourceError::Malformed {
message: format!("Linear answered a status name this source cannot hold: {refused}"),
}
})
}
#[derive(Debug, Clone)]
struct Vocabulary {
team: NativeId,
states: Vec<Held>,
statuses: Vec<Held>,
last_position: f64,
}
impl Vocabulary {
fn read(data: &Value, source: &SourceName, team: &str) -> Result<Self, SourceError> {
let teams = data
.pointer("/teams/nodes")
.and_then(Value::as_array)
.ok_or_else(|| SourceError::Malformed {
message: "missing teams.nodes".into(),
})?;
let found = match teams.as_slice() {
[found] => found,
other => {
return Err(SourceError::Refused {
message: format!(
"source {source} cannot resolve the configured team {team:?}: found {} \
matches; next: set team to the key of exactly one team this \
credential can see",
other.len()
),
});
}
};
let nodes = |pointer: &str| {
data.pointer(pointer)
.and_then(Value::as_array)
.ok_or_else(|| SourceError::Malformed {
message: format!(
"missing {}",
pointer.trim_start_matches('/').replace('/', ".")
),
})
};
let team = NativeId(backend_id(found, "id")?.into());
for (more, held) in [
(
"/teams/nodes/0/states/pageInfo/hasNextPage",
"workflow states of its team",
),
(
"/projectStatuses/pageInfo/hasNextPage",
"project statuses of its workspace",
),
] {
let more = data.pointer(more).and_then(Value::as_bool).ok_or_else(|| {
SourceError::Malformed {
message: format!(
"missing boolean {}",
more.trim_start_matches('/').replace('/', ".")
),
}
})?;
if more {
return Err(SourceError::Refused {
message: format!(
"source {source} reads the {held} in one page of 250, and Linear holds \
more; next: archive the ones no longer used, so the rest fit one page"
),
});
}
}
let statuses = nodes("/projectStatuses/nodes")?;
Ok(Self {
team,
states: nodes("/teams/nodes/0/states/nodes")?
.iter()
.map(Held::read)
.collect::<Result<_, _>>()?,
statuses: statuses.iter().map(Held::read).collect::<Result<_, _>>()?,
last_position: statuses
.iter()
.map(position_of)
.try_fold(0.0_f64, |last, position| Ok(last.max(position?)))?,
})
}
fn add(&mut self, created: Created) {
match created {
Created::State(held) => self.states.push(held),
Created::Status { held, position } => {
self.statuses.push(held);
self.last_position = self.last_position.max(position);
}
}
}
fn of(&self, kind: ItemKind) -> &[Held] {
match kind {
ItemKind::Task => &self.states,
ItemKind::Project => &self.statuses,
}
}
fn find(
&self,
kind: ItemKind,
name: &str,
source: &SourceName,
) -> Result<Option<&Held>, SourceError> {
let matched = self
.of(kind)
.iter()
.filter(|held| held.name.matches(name))
.collect::<Vec<_>>();
match matched.as_slice() {
[] => Ok(None),
[held] => Ok(Some(held)),
several => Err(SourceError::Refused {
message: format!(
"source {source} cannot resolve {} {name:?}: found {} matches with ids {:?}",
vocabulary_word(kind),
several.len(),
several
.iter()
.map(|held| held.id.0.as_str())
.collect::<Vec<_>>()
),
}),
}
}
}
#[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, serde::Serialize, schemars::JsonSchema)]
pub struct StatusNamesReport {
pub source: SourceName,
team: Team,
pub names: Vec<MappedStatusName>,
#[serde(skip_serializing_if = "Option::is_none")]
pub refused: Option<RefusedCreate>,
}
impl StatusNamesReport {
#[must_use]
pub fn team(&self) -> &str {
&self.team.0
}
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, schemars::JsonSchema)]
pub struct RefusedCreate {
pub kind: ItemKind,
pub name: StatusName,
pub message: String,
}
#[derive(Debug, Clone, PartialEq)]
pub struct MappedStatusName {
kind: ItemKind,
category: StatusCategory,
name: StatusName,
found: Found,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Found {
Present(String),
Missing,
Created(String),
}
impl MappedStatusName {
#[must_use]
pub fn kind(&self) -> ItemKind {
self.kind
}
#[must_use]
pub fn category(&self) -> StatusCategory {
self.category
}
#[must_use]
pub fn name(&self) -> &str {
self.name.as_str()
}
#[must_use]
pub fn found(&self) -> &Found {
&self.found
}
#[must_use]
pub fn created(&self) -> bool {
matches!(self.found, Found::Created(_))
}
#[must_use]
pub fn expected_type(&self) -> &'static str {
created_type(self.category, self.kind)
}
}
#[derive(serde::Serialize, schemars::JsonSchema)]
#[schemars(rename = "MappedStatusName")]
struct MappedStatusNameWire<'a> {
kind: ItemKind,
category: StatusCategory,
name: &'a StatusName,
present: bool,
#[serde(rename = "type", skip_serializing_if = "Option::is_none")]
found_type: Option<&'a str>,
expected_type: &'static str,
#[serde(default, skip_serializing_if = "std::ops::Not::not")]
#[schemars(!skip_serializing_if)]
created: bool,
}
impl<'a> From<&'a MappedStatusName> for MappedStatusNameWire<'a> {
fn from(mapped: &'a MappedStatusName) -> Self {
let found_type = match &mapped.found {
Found::Present(kind) | Found::Created(kind) => Some(kind.as_str()),
Found::Missing => None,
};
Self {
kind: mapped.kind,
category: mapped.category,
name: &mapped.name,
present: found_type.is_some(),
found_type,
expected_type: mapped.expected_type(),
created: mapped.created(),
}
}
}
impl serde::Serialize for MappedStatusName {
fn serialize<S: serde::Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
MappedStatusNameWire::from(self).serialize(serializer)
}
}
impl schemars::JsonSchema for MappedStatusName {
fn schema_name() -> std::borrow::Cow<'static, str> {
MappedStatusNameWire::schema_name()
}
fn json_schema(generator: &mut schemars::SchemaGenerator) -> Schema {
MappedStatusNameWire::json_schema(generator)
}
}
pub async fn status_names(
name: &SourceName,
config: LinearConfig,
secrets: &dyn SecretResolver,
apply: bool,
) -> Result<StatusNamesReport, SourceError> {
LinearSource::new(name, config, secrets)?
.status_names(apply)
.await
}
struct LinearSource {
client: reqwest::Client,
endpoint: Endpoint,
key: SecretString,
team: Option<Team>,
name: SourceName,
statuses: StatusMapping,
project: Option<ProjectScope>,
vocabulary: std::sync::Mutex<Option<std::sync::Arc<Vocabulary>>>,
}
impl LinearSource {
fn new(
name: &SourceName,
config: LinearConfig,
secrets: &dyn SecretResolver,
) -> Result<Self, SourceError> {
for kind in [ItemKind::Task, ItemKind::Project] {
StatusMapping::distinct(
name,
kind,
config
.status_mapping
.names(kind)
.map(|(category, mapped)| (category, mapped.as_str())),
)?;
}
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: config.status_mapping,
project: config.project,
vocabulary: std::sync::Mutex::new(None),
})
}
}
#[derive(Clone, Copy)]
enum WriteKind {
Task,
Project,
}
enum HeldRelations<'a> {
None,
Page(&'a Value),
Unread,
}
enum Lookup<'a> {
IssueLabel(&'a str),
ProjectLabel(&'a str),
}
impl Lookup<'_> {
fn query(&self) -> &'static str {
match self {
Self::IssueLabel(_) => graphql::ISSUE_LABEL,
Self::ProjectLabel(_) => graphql::PROJECT_LABEL,
}
}
fn connection(&self) -> &'static str {
match self {
Self::IssueLabel(_) => "issueLabels",
Self::ProjectLabel(_) => "projectLabels",
}
}
fn diagnostic(&self) -> String {
match self {
Self::IssueLabel(name) | Self::ProjectLabel(name) => format!("label {name:?}"),
}
}
fn variables(&self) -> Value {
match self {
Self::IssueLabel(name) | Self::ProjectLabel(name) => json!({"name":name}),
}
}
}
#[derive(Clone, Copy)]
enum MutationRoot {
IssueCreate,
IssueUpdate,
ProjectCreate,
ProjectUpdate,
IssueRelationCreate,
ProjectRelationCreate,
IssueRelationDelete,
ProjectRelationDelete,
IssueDelete,
ProjectDelete,
DocumentCreate,
DocumentUpdate,
DocumentDelete,
CommentCreate,
CommentUpdate,
CommentDelete,
WorkflowStateCreate,
ProjectStatusCreate,
}
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",
Self::WorkflowStateCreate => "workflowStateCreate",
Self::ProjectStatusCreate => "projectStatusCreate",
}
}
}
#[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,
}
struct Refusal(GqlError);
impl Refusal {
fn entity_missing(&self) -> bool {
let lowered = self.0.message.to_ascii_lowercase();
lowered.trim_end() == graphql::ISSUE_NOT_FOUND_MESSAGE.to_ascii_lowercase()
|| lowered.starts_with(&format!(
"{} ",
graphql::ISSUE_NOT_FOUND_MESSAGE.to_ascii_lowercase()
))
|| self
.0
.extensions
.as_ref()
.and_then(|extensions| extensions.get("userPresentableMessage"))
.and_then(Value::as_str)
.is_some_and(|said| {
said.to_ascii_lowercase().starts_with(
&graphql::ISSUE_NOT_FOUND_PRESENTABLE
.trim_end_matches('.')
.to_ascii_lowercase(),
)
})
}
fn into_error(self) -> SourceError {
SourceError::Refused {
message: self.0.said(),
}
}
}
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> {
self.answer(query, variables)
.await?
.map_err(|refusal| refusal.into_error())
}
async fn send_to_held(
&self,
query: &str,
variables: Value,
) -> Result<Option<Value>, SourceError> {
match self.answer(query, variables).await? {
Ok(data) => Ok(Some(data)),
Err(refusal) if refusal.entity_missing() => Ok(None),
Err(refusal) => Err(refusal.into_error()),
}
}
async fn answer(
&self,
query: &str,
variables: Value,
) -> Result<Result<Value, Refusal>, 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 text = response.text().await.unwrap_or_default();
if let Some(missing) = serde_json::from_str::<Envelope>(&text)
.ok()
.and_then(|body| body.errors.into_iter().next())
.map(Refusal)
.filter(Refusal::entity_missing)
{
return Ok(Err(missing));
}
let said = elided(&text);
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.into_iter().next() {
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 Ok(Err(Refusal(error)));
}
body.data.map(Ok).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.extend(self.status_narrowing(ItemKind::Task, &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, kind: ItemKind, statuses: &[StatusCategory]) -> Option<Value> {
let member = match kind {
ItemKind::Task => "state",
ItemKind::Project => "status",
};
let named = |name: &str, operator: &str| json!({ (member): {"name": {(operator): name}}});
let mut alternatives = Vec::new();
for category in statuses {
if let Ok(name) = self.statuses.name_for(*category, kind) {
alternatives.push(named(name.as_str(), "eqIgnoreCase"));
}
if *category == StatusCategory::Unknown {
let unmapped: Vec<Value> = self
.statuses
.names(kind)
.map(|(_, name)| named(name.as_str(), "neqIgnoreCase"))
.collect();
if unmapped.is_empty() {
return None;
}
alternatives.push(Self::narrowed(unmapped));
}
}
Some(match alternatives.len() {
0 => json!({ (member): {"name": {"in": Vec::<&str>::new()}}}),
1 => alternatives.pop().expect("one alternative"),
_ => json!({ "or": alternatives }),
})
}
fn confirms(query: &TaskQuery, task: &Task) -> bool {
(query.statuses.is_empty() || query.statuses.contains(&task.status.category))
&& 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.extend(self.status_narrowing(ItemKind::Project, statuses));
}
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
),
}),
}
}
fn project_out_of_scope(&self, target: Option<&NativeId>) -> Option<SourceError> {
let scope = self.project.as_ref()?;
if target.is_some_and(|target| target.0 == scope.0) {
return None;
}
Some(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,
target.map_or_else(
|| "a new one".to_owned(),
|target| format!("the project {}", target.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"),
})?;
match nodes.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 vocabulary(
&self,
fresh: bool,
) -> Result<(std::sync::Arc<Vocabulary>, bool), SourceError> {
if !fresh && let Some(held) = self.held_vocabulary() {
return Ok((held, false));
}
let team = self.team.as_ref().ok_or_else(|| SourceError::Refused {
message: format!(
"source {} needs config.team before it can write a Linear item or a status",
self.name
),
})?;
let data = self
.send(graphql::RESOLUTION, json!({"key": team.0}))
.await?;
let read = std::sync::Arc::new(Vocabulary::read(&data, &self.name, &team.0)?);
*self
.vocabulary
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(read.clone());
Ok((read, true))
}
fn held_vocabulary(&self) -> Option<std::sync::Arc<Vocabulary>> {
self.vocabulary
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
}
fn forget_vocabulary(&self) {
*self
.vocabulary
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner) = None;
}
fn remember(&self, created: &Created) {
let mut guard = self
.vocabulary
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if let Some(vocabulary) = guard.as_mut() {
std::sync::Arc::make_mut(vocabulary).add(created.clone());
}
}
async fn team_id(&self) -> Result<NativeId, SourceError> {
Ok(self.vocabulary(false).await?.0.team.clone())
}
async fn status_id(
&self,
category: StatusCategory,
kind: ItemKind,
) -> Result<(NativeId, String), SourceError> {
let name = self
.statuses
.name_for(category, kind)
.map_err(|why| why.refusal(&self.name, category, kind))?;
let (vocabulary, fresh) = self.vocabulary(false).await?;
if let Some(held) = vocabulary.find(kind, name.as_str(), &self.name)? {
return Ok((held.id.clone(), held.name.as_str().to_owned()));
}
if !fresh {
let (vocabulary, _) = self.vocabulary(true).await?;
if let Some(held) = vocabulary.find(kind, name.as_str(), &self.name)? {
return Ok((held.id.clone(), held.name.as_str().to_owned()));
}
}
let category_key = category_word(category);
let kind_key = kind_word(kind);
let held_by = match kind {
ItemKind::Task => format!(
"team {}",
self.team.as_ref().map_or("(none)", |team| team.0.as_str())
),
ItemKind::Project => "this workspace".to_owned(),
};
Err(SourceError::Refused {
message: format!(
"source {} maps the {kind_key} status {category_key} to the {} {:?}, which {held_by} \
does not have; next: run `onetaskgraph sources fields {} --apply` to create it, \
or point status_mapping.{category_key}.{kind_key} of this source at a {} \
{held_by} has",
self.name,
vocabulary_word(kind),
name.as_str(),
self.name,
vocabulary_word(kind),
),
})
}
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.statuses
.name_for(status.category, ItemKind::Task)
.map_err(|why| why.refusal(&self.name, status.category, ItemKind::Task))?;
}
if let Some(status) = update.status.as_ref().filter(|_| {
TaskUpdate {
status: None,
..update.clone()
}
.is_empty()
}) {
let Some(task) = self.status_written(id, status.category).await? else {
return Ok(None);
};
return Ok(Some(TaskUpdateOutcome {
delivers_before: task.delivers.clone(),
written: std::iter::once(UpdatedField::Status).collect(),
task,
}));
}
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.status_id(status.category, ItemKind::Task).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 carries_state = input.contains_key("stateId");
let task = if input.is_empty() {
before.clone()
} else {
let answered = self
.send(
graphql::ISSUE_UPDATE_READ,
json!({"id":before.id.0,"input":Value::Object(input)}),
)
.await
.and_then(|data| self.updated_issue(&data, &before.id));
match answered {
Ok(Some(task)) => task,
Ok(None) => {
return Err(SourceError::Malformed {
message: format!("task {id} was updated and then answered as no issue"),
});
}
Err(error) => {
if carries_state {
self.forget_vocabulary();
}
return Err(error);
}
}
};
if let Some(prepared) = &relations {
self.write_relations(&before.id, prepared, WriteKind::Task, HeldRelations::Unread)
.await?;
}
let mut written = update.changed(&before, &task);
if relations.is_some() {
written.insert(UpdatedField::DependsOn);
}
Ok(Some(TaskUpdateOutcome {
task,
written,
delivers_before: before.delivers,
}))
}
async fn status_written(
&self,
id: &NativeId,
category: StatusCategory,
) -> Result<Option<Task>, SourceError> {
let (state, _) = self.status_id(category, ItemKind::Task).await?;
let target = id.clone();
let answered = self
.send_to_held(
graphql::ISSUE_UPDATE_READ,
json!({"id":target.0,"input":{"stateId":state.0}}),
)
.await;
match answered {
Ok(Some(data)) => self.updated_issue(&data, &target),
Ok(None) => Ok(None),
Err(error) => {
self.forget_vocabulary();
Err(error)
}
}
}
fn updated_issue(&self, data: &Value, asked: &NativeId) -> Result<Option<Task>, SourceError> {
let issue = mutation_payload(data, MutationRoot::IssueUpdate)?
.get("issue")
.filter(|issue| !issue.is_null())
.ok_or_else(|| SourceError::Malformed {
message: "missing issueUpdate.issue".into(),
})?;
if optional_str(issue, "identifier")? != Some(asked.0.as_str()) {
written_is(issue, asked)?;
}
optional(&json!({ "issue": issue }), "issue", |v| {
map_task(v, &self.name, &self.statuses)
})
}
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,
held: HeldRelations<'_>,
) -> Result<(), SourceError> {
let mut cursor: Option<Cursor> = None;
let mut answered = match held {
HeldRelations::None => None,
HeldRelations::Page(page) => Some(page.clone()),
HeldRelations::Unread => Some(Value::Null),
};
while let Some(page) = answered.take() {
let relations = if page.is_null() {
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(),
})?;
root.get("relations")
.cloned()
.ok_or_else(|| SourceError::Malformed {
message: "missing relations".into(),
})?
} else {
page
};
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)?;
}
if let Some(next) = page_next(&relations)? {
cursor = Some(next);
answered = Some(Value::Null);
}
}
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", |v| map_project(v, &self.statuses))?;
(
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,
assets: Support::Unsupported,
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", |v| map_project(v, &self.statuses))?;
Ok(Page {
items: page
.items
.into_iter()
.filter(|project| {
(query.statuses.is_empty() || query.statuses.contains(&project.status.category))
&& 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 })?;
}
let project = self.filed_in(write.item.project.as_ref(), "task")?;
let (state, _) = self
.status_id(write.item.status.category, ItemKind::Task)
.await?;
let edges = self
.prepare_edges(&write.depends_on, WriteKind::Task)
.await?;
let team = self.team_id().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_REWRITE,
json!({"id":id.0,"input":input,"first":MAX_PAGE_SIZE}),
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.inspect_err(|_| {
self.forget_vocabulary();
})?;
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());
let held = held_relations(write.target.as_ref(), issue)?;
self.write_relations(&id, &edges, WriteKind::Task, held)
.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(refused) = self.project_out_of_scope(write.target.as_ref()) {
return Err(refused);
}
let (status, _) = self
.status_id(write.item.status.category, ItemKind::Project)
.await?;
let edges = self
.prepare_edges(&write.depends_on, WriteKind::Project)
.await?;
let team = self.team_id().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,"content":description,"statusId":status,"labelIds":labels});
let (query, variables, root) = match &write.target {
Some(id) => (
graphql::PROJECT_REWRITE,
json!({"id":id.0,"input":input,"first":MAX_PAGE_SIZE}),
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.inspect_err(|_| {
self.forget_vocabulary();
})?;
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());
let held = held_relations(write.target.as_ref(), project)?;
self.write_relations(&id, &edges, WriteKind::Project, held)
.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 check_status_write(
&self,
kind: ItemKind,
category: StatusCategory,
target: Option<&NativeId>,
) -> Result<(), SourceError> {
if kind == ItemKind::Project
&& let Some(refused) = self.project_out_of_scope(target)
{
return Err(refused);
}
self.status_id(category, kind).await.map(|_| ())
}
async fn set_task_status(
&self,
id: &NativeId,
category: StatusCategory,
) -> Result<Option<Status>, SourceError> {
Ok(self
.status_written(id, category)
.await?
.map(|task| task.status))
}
async fn set_task_status_reading(
&self,
id: &NativeId,
category: StatusCategory,
) -> Result<Option<Task>, SourceError> {
self.status_written(id, category).await
}
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_content(&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 set_project_rendering(
&self,
id: &NativeId,
content: &str,
provenance: &Value,
_answers: &std::collections::BTreeMap<String, Value>,
) -> Result<Option<()>, SourceError> {
let Some((project, description)) = self.project_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_project_content(&project.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
}
async fn end_command(&self) -> Result<(), SourceError> {
Ok(())
}
}
fn held_relations<'a>(
target: Option<&NativeId>,
written: &'a Value,
) -> Result<HeldRelations<'a>, SourceError> {
if target.is_none() {
return Ok(HeldRelations::None);
}
written
.get("relations")
.filter(|relations| !relations.is_null())
.map(HeldRelations::Page)
.ok_or_else(|| SourceError::Malformed {
message: "a rewrite answered without the relations it was asked for".into(),
})
}
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 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, &self.statuses)?,
optional_string(v, "content")?,
))
})?
.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_content(
&self,
id: &NativeId,
content: Option<&str>,
) -> Result<(), SourceError> {
let data = self
.send(
graphql::PROJECT_UPDATE,
json!({"id":id.0,"input":{"content":content}}),
)
.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 status_names(&self, apply: bool) -> Result<StatusNamesReport, SourceError> {
let team = self.team.clone().ok_or_else(|| SourceError::Refused {
message: format!(
"source {} needs config.team to report its status names",
self.name
),
})?;
let (mut vocabulary, _) = self.vocabulary(false).await?;
let mut names = Vec::new();
let mut refused = None;
for kind in [ItemKind::Task, ItemKind::Project] {
for (category, name) in self.statuses.names(kind) {
let found = vocabulary
.find(kind, name.as_str(), &self.name)?
.map(|held| held.kind.clone());
let mut mapped = MappedStatusName {
kind,
category,
name: name.clone(),
found: found.map_or(Found::Missing, Found::Present),
};
if apply && refused.is_none() && mapped.found == Found::Missing {
match self
.create_status_name(kind, category, name.as_str(), &vocabulary)
.await
{
Ok(created) => {
mapped.found = Found::Created(created.held().kind.clone());
self.remember(&created);
std::sync::Arc::make_mut(&mut vocabulary).add(created);
}
Err(error) => {
refused = Some(RefusedCreate {
kind,
name: name.clone(),
message: error.to_string(),
});
}
}
}
names.push(mapped);
}
}
Ok(StatusNamesReport {
source: self.name.clone(),
team,
names,
refused,
})
}
async fn create_status_name(
&self,
kind: ItemKind,
category: StatusCategory,
name: &str,
vocabulary: &Vocabulary,
) -> Result<Created, SourceError> {
let kind_of = created_type(category, kind);
let (query, input, root, payload) = match kind {
ItemKind::Task => (
graphql::WORKFLOW_STATE_CREATE,
json!({"teamId": vocabulary.team.0, "name": name, "type": kind_of,
"color": CREATED_COLOR}),
MutationRoot::WorkflowStateCreate,
"workflowState",
),
ItemKind::Project => (
graphql::PROJECT_STATUS_CREATE,
json!({"name": name, "type": kind_of, "color": CREATED_COLOR,
"position": vocabulary.last_position + 1.0}),
MutationRoot::ProjectStatusCreate,
"status",
),
};
let data = self.send(query, json!({ "input": input })).await?;
let created = mutation_payload(&data, root)?
.get(payload)
.filter(|created| !created.is_null())
.ok_or_else(|| SourceError::Malformed {
message: format!("missing {}.{payload}", root.as_str()),
})?;
let held = Held::read(created)?;
if held.name.as_str() != name || held.kind != kind_of {
return Err(SourceError::Malformed {
message: format!(
"Linear answered the create of {} {name:?} of type {kind_of} with {:?} of \
type {}",
vocabulary_word(kind),
held.name.as_str(),
held.kind
),
});
}
Ok(match kind {
ItemKind::Task => Created::State(held),
ItemKind::Project => Created::Status {
position: position_of(created)?,
held,
},
})
}
}
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, root.long_form())?)?;
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 mapped_status(
v: &Value,
statuses: &StatusMapping,
kind: ItemKind,
) -> Result<Status, SourceError> {
let name = str_at(v, "name")?;
Ok(Status {
category: statuses
.category_of(kind, name)
.unwrap_or(StatusCategory::Unknown),
name: name.to_owned(),
})
}
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: mapped_status(
v.get("state").ok_or_else(|| SourceError::Malformed {
message: "missing state".into(),
})?,
statuses,
ItemKind::Task,
)?,
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, statuses: &StatusMapping) -> Result<Project, 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(Project {
id: NativeId(str_at(v, "id")?.into()),
title: str_at(v, "name")?.into(),
content,
status: mapped_status(
v.get("status").ok_or_else(|| SourceError::Malformed {
message: "missing status".into(),
})?,
statuses,
ItemKind::Project,
)?,
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",
}
}
const fn long_form(self) -> &'static str {
match self {
Self::Issue => "description",
Self::Project => "content",
}
}
}
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())))
}