Skip to main content

onetaskgraph_core/subprocess/
source.rs

1//! The engine's half of the protocol: a [`TaskSource`] that is another process.
2//!
3//! Every method here is one line out and one line back. What it deliberately does *not*
4//! do is decide anything: a `forward-only` plugin is never asked for
5//! [`Direction::DependedOnBy`] because the layer above reads that off the capabilities
6//! this handshake returned and emulates the reverse scan itself, and a predicate a plugin
7//! declared unsupported is removed from the query before it ever reaches here. Putting
8//! either decision in this file would give the product a second compensation layer that
9//! only subprocess-hosted sources went through.
10
11use std::collections::BTreeMap;
12use std::num::NonZeroU64;
13use std::path::Path;
14use std::time::Duration;
15
16use async_trait::async_trait;
17use onetaskgraph_plugin_api::{
18    Capabilities, Comment, CommentBody, DependencyEdge, Direction, Document, DocumentQuery, Health,
19    ItemWrite, Label, MetadataKey, MetadataRecord, Metering, NativeId, NewComment, Page,
20    PageRequest, Priority, Project, ProjectQuery, SourceError, SourceName, Status, StatusCategory,
21    Task, TaskQuery, TaskRef, TaskSource, WriteSupport, unwritable_field, unwritable_metadata,
22};
23use serde::Deserialize;
24use serde_json::{Value, json};
25
26use super::connection::{Connection, Peer};
27use super::wire::{
28    AddCommentParams, CommentResult, CommentsParams, CommentsResult, ContentParams, ContentResult,
29    DeleteCommentParams, DeleteParams, DeletedCommentResult, DeliveredByParams, DeliveredByResult,
30    DependencyParams, DocumentDir, DocumentQueryParams, DocumentResult, DocumentWriteParams,
31    EditCommentParams, EngineIdentity, IdParams, InitializeParams, InitializeResult, LabelParams,
32    MetadataParams, MeteringResult, PROTOCOL_VERSION, PriorityParams, PriorityResult,
33    ProjectQueryParams, ProjectResult, ProjectWriteParams, Request, StatusParams, StatusResult,
34    TaskQueryParams, TaskResult, TaskWriteParams, WriteResult, after_the_first_vocabulary,
35    knows_every_category, spelled, vocabulary,
36};
37
38/// The id the handshake is sent under. §3 makes it the first request on a connection, so
39/// nothing else can have been sent under it, and an answer addressed elsewhere is a
40/// violation rather than an ordering the engine could accommodate.
41const HANDSHAKE_ID: &str = "0";
42
43/// A positive per-request deadline, measured in milliseconds at the configuration edge.
44#[derive(Debug, Clone, Copy, PartialEq, Eq)]
45pub struct RequestDeadline(NonZeroU64);
46
47impl RequestDeadline {
48    /// The protocol's default deadline.
49    pub const DEFAULT: Self = Self(NonZeroU64::new(30_000).expect("non-zero default"));
50
51    /// Validate a millisecond value from a configuration or another public boundary.
52    #[must_use]
53    pub const fn from_millis(milliseconds: NonZeroU64) -> Self {
54        Self(milliseconds)
55    }
56
57    /// The positive millisecond count used by configuration and diagnostics.
58    #[must_use]
59    pub const fn milliseconds(self) -> NonZeroU64 {
60        self.0
61    }
62
63    fn duration(self) -> Duration {
64        Duration::from_millis(self.0.get())
65    }
66}
67
68/// The two bounds a spawned plugin's exchanges are held to, carried together so that the
69/// one constructor every other reaches takes a pair rather than two arguments of one type
70/// a caller could transpose.
71#[derive(Debug, Clone, Copy)]
72struct Deadlines {
73    /// Bounds the `initialize` exchange, and so the child's own start-up with it.
74    handshake: RequestDeadline,
75    /// Bounds each exchange after the handshake, against an already running child.
76    requests: RequestDeadline,
77}
78
79/// A source served by a spawned program speaking `docs/plugin-protocol.md`.
80pub struct SubprocessSource {
81    /// What the plugin called itself in the handshake.
82    ///
83    /// Leaked once per connection because [`TaskSource::kind`] returns `&'static str` for
84    /// the compiled-in plugins, whose kinds really are static, and a subprocess-hosted
85    /// plugin's kind is not known until it answers. One small allocation per configured
86    /// source, for the life of a process that was going to hold that source anyway, is
87    /// the cheapest way to keep the trait honest for both.
88    kind: &'static str,
89    /// Read once at the handshake; §3 says the engine does not ask again.
90    capabilities: Capabilities,
91    /// Whether the plugin said it can be written through, read at the same handshake.
92    ///
93    /// A plugin that said nothing is read as read-only, which is what §3.3 makes an
94    /// absent member mean and what every version-1 plugin written before there was a
95    /// write side is.
96    writes: WriteSupport,
97    /// Whether the plugin said it answers `metering`, read at the same handshake.
98    ///
99    /// A plugin that said nothing is never sent the method and is reported as not metering,
100    /// which is what §3.4 makes an absent member mean.
101    meters: bool,
102    /// Whether the plugin's handshake listed every status category this build knows (§3.5).
103    ///
104    /// A plugin that listed none was written against the first vocabulary, and is never
105    /// handed a category added after it: not in a query, not in a write.
106    knows_every_category: bool,
107    /// Whether the plugin said it answers the two narrow task writes and holds a task's two
108    /// lists (§3.6), read at the same handshake.
109    task_updates: bool,
110    /// Whether the plugin said it answers the three narrow metadata writes (§3.7), read at the
111    /// same handshake.
112    metadata_updates: bool,
113    /// Whether the plugin said it answers the narrow content write (§3.9), read at the same
114    /// handshake.
115    content_updates: bool,
116    /// The live process.
117    connection: Connection,
118}
119
120impl std::fmt::Debug for SubprocessSource {
121    /// Named without its connection, which holds a live child and a credential the
122    /// handshake forwarded — neither belongs in a diagnostic.
123    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
124        f.debug_struct("SubprocessSource")
125            .field("kind", &self.kind)
126            .finish_non_exhaustive()
127    }
128}
129
130impl SubprocessSource {
131    /// Spawn `program`, complete the handshake, and adopt the connection.
132    ///
133    /// # Errors
134    ///
135    /// Returns [`SourceError::Unavailable`] when the program cannot be run or stops
136    /// answering, the plugin's own error when it refuses the handshake, and
137    /// [`SourceError::Config`] when the two sides do not speak the same protocol version
138    /// — refused by name, never guessed at (§6.1).
139    pub fn connect(
140        program: &str,
141        args: &[String],
142        name: &SourceName,
143        config: &Value,
144        secrets: BTreeMap<String, String>,
145    ) -> Result<Self, SourceError> {
146        Self::connect_with_deadline(
147            program,
148            args,
149            name,
150            config,
151            secrets,
152            RequestDeadline::DEFAULT,
153        )
154    }
155
156    /// Spawn a plugin with a deadline applying independently to every exchange.
157    ///
158    /// # Errors
159    ///
160    /// What [`connect`](Self::connect) returns.
161    pub fn connect_with_deadline(
162        program: &str,
163        args: &[String],
164        name: &SourceName,
165        config: &Value,
166        secrets: BTreeMap<String, String>,
167        deadline: RequestDeadline,
168    ) -> Result<Self, SourceError> {
169        Self::connect_with_deadlines(program, args, name, config, secrets, deadline, deadline)
170    }
171
172    /// Spawn a plugin, bounding the initialization handshake and every later request
173    /// separately.
174    ///
175    /// `handshake` bounds the one `initialize` exchange of §3, and a newly spawned child
176    /// does its own starting up inside that exchange: whatever a runtime loads before it
177    /// can read its first line is counted against this bound. `requests` bounds every
178    /// exchange after the handshake succeeded, each on its own, against a child that is by
179    /// then already running. They are the same bound in
180    /// [`connect_with_deadline`](Self::connect_with_deadline) and in every configured
181    /// source, because one `deadline_ms` is what `docs/plugin-protocol.md` §1 gives a user
182    /// to set. Separating them is for a caller that means to hold a *request* to a span
183    /// shorter than a program takes to start — the engine's own probe of what a silent
184    /// child's expired request reports — where one bound would fail the handshake on a
185    /// loaded host instead of reaching the behaviour it was after.
186    ///
187    /// # Errors
188    ///
189    /// What [`connect`](Self::connect) returns.
190    pub fn connect_with_deadlines(
191        program: &str,
192        args: &[String],
193        name: &SourceName,
194        config: &Value,
195        secrets: BTreeMap<String, String>,
196        handshake: RequestDeadline,
197        requests: RequestDeadline,
198    ) -> Result<Self, SourceError> {
199        Self::connect_bounded(
200            program,
201            args,
202            name,
203            config,
204            secrets,
205            Deadlines {
206                handshake,
207                requests,
208            },
209            None,
210        )
211    }
212
213    /// Spawn a plugin, telling it which document's directory its settings are measured
214    /// from — `document_dir` in `docs/plugin-protocol.md` §3, sent only when there is one.
215    ///
216    /// # Errors
217    ///
218    /// What [`connect`](Self::connect) returns, and [`SourceError::Config`] when
219    /// `document_dir` is not an absolute directory whose name is valid UTF-8, and so
220    /// cannot be written into the handshake — refused rather than dropped, because dropping
221    /// it would silently measure the child's paths from its working directory instead.
222    pub fn connect_from_document(
223        program: &str,
224        args: &[String],
225        name: &SourceName,
226        config: &Value,
227        secrets: BTreeMap<String, String>,
228        deadline: RequestDeadline,
229        document_dir: Option<&Path>,
230    ) -> Result<Self, SourceError> {
231        Self::connect_bounded(
232            program,
233            args,
234            name,
235            config,
236            secrets,
237            Deadlines {
238                handshake: deadline,
239                requests: deadline,
240            },
241            document_dir,
242        )
243    }
244
245    fn connect_bounded(
246        program: &str,
247        args: &[String],
248        name: &SourceName,
249        config: &Value,
250        secrets: BTreeMap<String, String>,
251        deadlines: Deadlines,
252        document_dir: Option<&Path>,
253    ) -> Result<Self, SourceError> {
254        let document_dir = document_dir
255            .map(|directory| {
256                DocumentDir::new(directory).map_err(|problem| SourceError::Config {
257                    message: format!(
258                        "source {name}: its settings are measured from the directory holding \
259                         the configuration document that set them, and {problem}; give the \
260                         settings absolute paths, or move the document under a directory \
261                         whose name is valid UTF-8"
262                    ),
263                })
264            })
265            .transpose()?;
266        Self::adopt(
267            Peer::spawn(
268                program,
269                args,
270                deadlines.handshake.duration(),
271                deadlines.requests.duration(),
272            )?,
273            name,
274            config,
275            secrets,
276            document_dir,
277        )
278    }
279
280    /// Connect to a plugin that is already running, over streams somebody else owns.
281    ///
282    /// The handshake, the framing and every refusal are the same as [`connect`]'s, because
283    /// they are the protocol's rather than the process's. What this constructor adds is
284    /// the ability to hold the *other* end: it is how the engine's own tests drive this
285    /// half against [`serve`](super::serve) over a real pipe, including the answers a
286    /// well-behaved program would never give.
287    ///
288    /// # Errors
289    ///
290    /// Returns what [`connect`](Self::connect) returns, minus the failures that belong to
291    /// spawning a program.
292    ///
293    /// [`connect`]: Self::connect
294    pub fn over(
295        to_plugin: impl std::io::Write + Send + 'static,
296        from_plugin: impl std::io::Read + Send + 'static,
297        name: &SourceName,
298        config: &Value,
299        secrets: BTreeMap<String, String>,
300    ) -> Result<Self, SourceError> {
301        Self::over_with_request_deadline(
302            to_plugin,
303            from_plugin,
304            name,
305            config,
306            secrets,
307            RequestDeadline::DEFAULT,
308        )
309    }
310
311    /// Connect over existing streams with a deadline for requests after initialization.
312    ///
313    /// Unlike [`connect_with_deadline`](Self::connect_with_deadline), this engine does
314    /// not own a process it can interrupt while the synchronous handshake is blocked.
315    /// The supplied deadline therefore begins only after initialization succeeds.
316    pub fn over_with_request_deadline(
317        to_plugin: impl std::io::Write + Send + 'static,
318        from_plugin: impl std::io::Read + Send + 'static,
319        name: &SourceName,
320        config: &Value,
321        secrets: BTreeMap<String, String>,
322        deadline: RequestDeadline,
323    ) -> Result<Self, SourceError> {
324        Self::adopt(
325            Peer::over(to_plugin, from_plugin, deadline.duration()),
326            name,
327            config,
328            secrets,
329            None,
330        )
331    }
332
333    /// Shake hands with `peer` and take the connection over.
334    fn adopt(
335        mut peer: Peer,
336        name: &SourceName,
337        config: &Value,
338        secrets: BTreeMap<String, String>,
339        document_dir: Option<DocumentDir>,
340    ) -> Result<Self, SourceError> {
341        let result = Self::handshake(&mut peer, name, config, secrets, document_dir);
342        let InitializeResult {
343            protocol_version,
344            kind,
345            capabilities,
346            writes,
347            meters,
348            statuses,
349            task_updates,
350            metadata_updates,
351            content_updates,
352        } = match result {
353            Ok(result) => result,
354            Err(error) => return Err(with_diagnostics(error, &mut peer)),
355        };
356        let kind = kind.into_string();
357        if protocol_version != Some(PROTOCOL_VERSION) {
358            return Err(SourceError::Config {
359                message: match protocol_version {
360                    Some(spoken) => format!(
361                        "the {kind:?} plugin was asked for protocol version \
362                         {PROTOCOL_VERSION} and answered in version {spoken}; the two are \
363                         incompatible and this engine does not guess between them"
364                    ),
365                    None => format!(
366                        "the {kind:?} plugin did not say which protocol version it \
367                         answered in; this engine speaks version {PROTOCOL_VERSION} and \
368                         does not guess"
369                    ),
370                },
371            });
372        }
373        Ok(Self {
374            kind: String::leak(kind),
375            capabilities,
376            writes: writes.unwrap_or(WriteSupport::Unsupported),
377            meters,
378            knows_every_category: knows_every_category(statuses.as_deref()),
379            task_updates,
380            metadata_updates,
381            content_updates,
382            connection: Connection::adopt(peer),
383        })
384    }
385
386    /// Send `initialize` and read what came back (§3).
387    fn handshake(
388        peer: &mut Peer,
389        name: &SourceName,
390        config: &Value,
391        secrets: BTreeMap<String, String>,
392        document_dir: Option<DocumentDir>,
393    ) -> Result<InitializeResult, SourceError> {
394        let params = InitializeParams {
395            protocol_version: PROTOCOL_VERSION,
396            engine: EngineIdentity {
397                name: "onetaskgraph".to_owned(),
398                version: env!("CARGO_PKG_VERSION").to_owned(),
399            },
400            source_name: name.as_str().to_owned(),
401            config: config.clone(),
402            secrets,
403            statuses: Some(vocabulary()),
404            document_dir,
405        };
406        let request = Request {
407            id: HANDSHAKE_ID.to_owned(),
408            method: "initialize".to_owned(),
409            // Plain data throughout: a `BTreeMap<String, String>` and a `Value` the
410            // configuration layer already parsed.
411            params: serde_json::to_value(&params).expect("a handshake is plain data"),
412        };
413        let line = peer.exchange(
414            &serde_json::to_string(&request).expect("a handshake request is plain data"),
415        )?;
416        let response: super::wire::Response =
417            serde_json::from_str(&line).map_err(|error| SourceError::Malformed {
418                message: format!(
419                    "the plugin's handshake answer is not a response envelope: {error}"
420                ),
421            })?;
422        // §6.3: an envelope addressed to an id this side never sent is a violation, and it
423        // is one here for the same reason it is later — a plugin whose first line answers
424        // something else has not answered the handshake, and reading it as one would build
425        // a source out of a message that was about something different.
426        if response.id != HANDSHAKE_ID {
427            return Err(SourceError::Malformed {
428                message: format!(
429                    "the plugin answered the handshake with an envelope addressed to {:?} \
430                     rather than to {HANDSHAKE_ID:?}",
431                    response.id
432                ),
433            });
434        }
435        let outcome = response.outcome().ok_or_else(|| SourceError::Malformed {
436            message: "the plugin's handshake answer carried both a result and an error, or \
437                      neither"
438                .to_owned(),
439        })?;
440        let result = outcome?;
441        serde_json::from_value(result).map_err(|error| SourceError::Malformed {
442            message: format!("the plugin's handshake answer is not an initialize result: {error}"),
443        })
444    }
445
446    /// The statuses a query may hand this plugin, or `None` when every one asked for is a
447    /// category it was not written against — which no row it holds can be in.
448    ///
449    /// Dropping those categories narrows nothing: a plugin that does not know a category
450    /// cannot report a row in it, so the rows the rest of the list matches are every row the
451    /// whole list matches.
452    fn statuses_for(&self, statuses: &[StatusCategory]) -> Option<Vec<StatusCategory>> {
453        if self.knows_every_category {
454            return Some(statuses.to_vec());
455        }
456        let known: Vec<StatusCategory> = statuses
457            .iter()
458            .copied()
459            .filter(|category| !after_the_first_vocabulary(*category))
460            .collect();
461        (known.len() == statuses.len() || !known.is_empty()).then_some(known)
462    }
463
464    /// Refuse a status this plugin's handshake says it was not written against, before it is
465    /// sent (§3.5).
466    fn knows(&self, category: StatusCategory) -> Result<(), SourceError> {
467        if self.knows_every_category || !after_the_first_vocabulary(category) {
468            return Ok(());
469        }
470        Err(SourceError::Refused {
471            message: format!(
472                "the {:?} plugin's handshake does not list the status category {}, so this \
473                 engine does not hand it one (docs/plugin-protocol.md §3.5); next: upgrade the \
474                 plugin to one whose handshake lists it, or use a category it knows",
475                self.kind,
476                spelled(category)
477            ),
478        })
479    }
480
481    /// Refuse what only a plugin declaring `task_updates` is handed, before it is sent (§3.6).
482    fn updates(&self, what: &str) -> Result<(), SourceError> {
483        if self.task_updates {
484            return Ok(());
485        }
486        Err(SourceError::Refused {
487            message: format!(
488                "the {:?} plugin's handshake does not say it answers the narrow task writes, so \
489                 this engine does not send it {what} (docs/plugin-protocol.md §3.6); next: \
490                 upgrade the plugin to one whose handshake sets task_updates",
491                self.kind
492            ),
493        })
494    }
495
496    /// Refuse a narrow metadata write of `record` to a plugin whose handshake did not declare
497    /// `metadata_updates`, before it is sent (§3.7), in the contract's own words.
498    fn metadata_updates(&self, record: MetadataRecord) -> Result<(), SourceError> {
499        if self.metadata_updates {
500            return Ok(());
501        }
502        Err(unwritable_metadata(self.kind, record))
503    }
504
505    /// Refuse one of the template operations — a write from a rendering, a read of stored
506    /// answers, a regenerate in place — before anything is sent.
507    ///
508    /// §4 carries none of them: a plugin answering the protocol has no method to keep answers
509    /// beside an item or replace a rendering with, so a `write_task` standing in for a create
510    /// from a template would land the content and the provenance without the answers a
511    /// regenerate needs. Refusing is what keeps an item from being half of what it claims.
512    fn unrendered(&self, operation: &str) -> SourceError {
513        SourceError::Refused {
514            message: format!(
515                "the {:?} plugin is hosted over the stdio plugin protocol, which does not carry \
516                 {operation} (docs/plugin-protocol.md §4), so this engine does not send it; \
517                 nothing was written; next: configure the source in-process under its own \
518                 plugin name rather than through a command",
519                self.kind
520            ),
521        }
522    }
523
524    /// Refuse a task write this plugin could only drop part of in silence.
525    fn writable_task(&self, task: &Task) -> Result<(), SourceError> {
526        self.knows(task.status.category)?;
527        // A plugin whose handshake declares no priority is one written before there were
528        // any, and it would drop one in silence (§6). The engine refuses such a write before
529        // it reaches here; this keeps the seam itself from ever sending one.
530        if task.priority != Priority::None && !self.capabilities.priority.is_native() {
531            return Err(unwritable_field(self.kind, "priority"));
532        }
533        if !task.delivers.is_empty() || !task.delivered_by.is_empty() {
534            self.updates("a task carrying delivers or delivered_by")?;
535        }
536        Ok(())
537    }
538
539    /// One call, with its result parsed into the shape the method promises.
540    async fn ask<T: for<'de> Deserialize<'de>>(
541        &self,
542        method: &str,
543        params: Value,
544    ) -> Result<T, SourceError> {
545        let result = self.connection.call(method, params).await?;
546        serde_json::from_value(result).map_err(|error| SourceError::Malformed {
547            message: format!(
548                "the plugin's answer to {method} is not the shape it promises: {error}"
549            ),
550        })
551    }
552}
553
554/// Append whatever the plugin said on standard error to a handshake failure.
555///
556/// A plugin that refuses the handshake and exits has usually said why there and nowhere
557/// else, and a bare "could not read the plugin's answer" would throw that away.
558fn with_diagnostics(error: SourceError, peer: &mut Peer) -> SourceError {
559    let said = peer.said();
560    if said.is_empty() {
561        return error;
562    }
563    let message = format!("{error}; the plugin wrote: {said}");
564    match error {
565        // The wait it asked for is preserved: what the plugin wrote is extra reason, not a
566        // replacement for the one piece of this refusal the engine acts on.
567        SourceError::RateLimited {
568            retry_after_seconds,
569            ..
570        } => SourceError::RateLimited {
571            retry_after_seconds,
572            message: Some(message),
573        },
574        SourceError::Config { .. } => SourceError::Config { message },
575        SourceError::Auth { .. } => SourceError::Auth { message },
576        SourceError::Refused { .. } => SourceError::Refused { message },
577        SourceError::Malformed { .. } => SourceError::Malformed { message },
578        SourceError::Unavailable { .. } => SourceError::Unavailable { message },
579    }
580}
581
582#[async_trait]
583impl TaskSource for SubprocessSource {
584    fn kind(&self) -> &'static str {
585        self.kind
586    }
587
588    fn capabilities(&self) -> Capabilities {
589        self.capabilities.clone()
590    }
591
592    async fn health(&self) -> Result<Health, SourceError> {
593        self.ask("health", json!({})).await
594    }
595
596    async fn get_task(&self, id: &NativeId) -> Result<Option<Task>, SourceError> {
597        let result: TaskResult = self
598            .ask("get_task", params(&IdParams { id: id.clone() }))
599            .await?;
600        Ok(result.task)
601    }
602
603    async fn get_project(&self, id: &NativeId) -> Result<Option<Project>, SourceError> {
604        let result: ProjectResult = self
605            .ask("get_project", params(&IdParams { id: id.clone() }))
606            .await?;
607        Ok(result.project)
608    }
609
610    async fn query_tasks(
611        &self,
612        query: &TaskQuery,
613        page: &PageRequest,
614    ) -> Result<Page<Task>, SourceError> {
615        let Some(statuses) = self.statuses_for(&query.statuses) else {
616            return Ok(Page::last(Vec::new()));
617        };
618        self.ask(
619            "query_tasks",
620            params(&TaskQueryParams {
621                query: TaskQuery {
622                    statuses,
623                    ..query.clone()
624                },
625                page: page.clone(),
626            }),
627        )
628        .await
629    }
630
631    async fn query_projects(
632        &self,
633        query: &ProjectQuery,
634        page: &PageRequest,
635    ) -> Result<Page<Project>, SourceError> {
636        let Some(statuses) = self.statuses_for(&query.statuses) else {
637            return Ok(Page::last(Vec::new()));
638        };
639        self.ask(
640            "query_projects",
641            params(&ProjectQueryParams {
642                query: ProjectQuery {
643                    statuses,
644                    ..query.clone()
645                },
646                page: page.clone(),
647            }),
648        )
649        .await
650    }
651
652    async fn labels(&self, page: &PageRequest) -> Result<Page<Label>, SourceError> {
653        self.ask("labels", params(&LabelParams { page: page.clone() }))
654            .await
655    }
656
657    async fn task_dependencies(
658        &self,
659        id: &NativeId,
660        direction: Direction,
661        page: &PageRequest,
662    ) -> Result<Page<DependencyEdge>, SourceError> {
663        self.ask(
664            "task_dependencies",
665            params(&DependencyParams {
666                id: id.clone(),
667                direction,
668                page: page.clone(),
669            }),
670        )
671        .await
672    }
673
674    async fn project_dependencies(
675        &self,
676        id: &NativeId,
677        direction: Direction,
678        page: &PageRequest,
679    ) -> Result<Page<DependencyEdge>, SourceError> {
680        self.ask(
681            "project_dependencies",
682            params(&DependencyParams {
683                id: id.clone(),
684                direction,
685                page: page.clone(),
686            }),
687        )
688        .await
689    }
690
691    fn writes(&self) -> WriteSupport {
692        self.writes
693    }
694
695    async fn write_task(&self, write: &ItemWrite<Task>) -> Result<NativeId, SourceError> {
696        self.writable_task(&write.item)?;
697        let result: WriteResult = self
698            .ask(
699                "write_task",
700                params(&TaskWriteParams {
701                    write: write.clone(),
702                }),
703            )
704            .await?;
705        Ok(result.id)
706    }
707
708    async fn write_project(&self, write: &ItemWrite<Project>) -> Result<NativeId, SourceError> {
709        self.knows(write.item.status.category)?;
710        let result: WriteResult = self
711            .ask(
712                "write_project",
713                params(&ProjectWriteParams {
714                    write: write.clone(),
715                }),
716            )
717            .await?;
718        Ok(result.id)
719    }
720
721    async fn set_task_status(
722        &self,
723        id: &NativeId,
724        category: StatusCategory,
725    ) -> Result<Option<Status>, SourceError> {
726        self.updates("set_task_status")?;
727        self.knows(category)?;
728        let result: StatusResult = self
729            .ask(
730                "set_task_status",
731                params(&StatusParams {
732                    id: id.clone(),
733                    category,
734                }),
735            )
736            .await?;
737        Ok(result.status)
738    }
739
740    async fn set_task_priority(
741        &self,
742        id: &NativeId,
743        priority: Priority,
744    ) -> Result<Option<Priority>, SourceError> {
745        // §4.19: sent only to a plugin whose handshake declares it holds a priority.
746        if !self.capabilities.priority.is_native() {
747            return Err(unwritable_field(self.kind, "priority"));
748        }
749        let result: PriorityResult = self
750            .ask(
751                "set_task_priority",
752                params(&PriorityParams {
753                    id: id.clone(),
754                    priority,
755                }),
756            )
757            .await?;
758        Ok(result.priority)
759    }
760
761    async fn set_task_content(
762        &self,
763        id: &NativeId,
764        content: &str,
765    ) -> Result<Option<()>, SourceError> {
766        // §3.9: sent only to a plugin whose handshake sets content_updates.
767        if !self.content_updates {
768            return Err(unwritable_field(self.kind, "content"));
769        }
770        let result: ContentResult = self
771            .ask(
772                "set_task_content",
773                params(&ContentParams {
774                    id: id.clone(),
775                    content: content.to_owned(),
776                }),
777            )
778            .await?;
779        match result.id {
780            None => Ok(None),
781            Some(written) if &written == id => Ok(Some(())),
782            Some(written) => Err(SourceError::Malformed {
783                message: format!(
784                    "the plugin answered set_task_content for {id} with the id {written}, which \
785                     is not the task it was asked to write"
786                ),
787            }),
788        }
789    }
790
791    async fn set_delivered_by(
792        &self,
793        id: &NativeId,
794        delivered_by: &[TaskRef],
795    ) -> Result<Option<()>, SourceError> {
796        self.updates("set_delivered_by")?;
797        let result: DeliveredByResult = self
798            .ask(
799                "set_delivered_by",
800                params(&DeliveredByParams {
801                    id: id.clone(),
802                    delivered_by: delivered_by.to_vec(),
803                }),
804            )
805            .await?;
806        Ok(result.delivered_by.map(|_| ()))
807    }
808
809    async fn set_task_metadata(
810        &self,
811        id: &NativeId,
812        key: &MetadataKey,
813        value: &Value,
814    ) -> Result<Option<Task>, SourceError> {
815        self.metadata_updates(MetadataRecord::Task)?;
816        let result: TaskResult = self
817            .ask("set_task_metadata", metadata_params(id, key, value))
818            .await?;
819        Ok(result.task)
820    }
821
822    async fn set_project_metadata(
823        &self,
824        id: &NativeId,
825        key: &MetadataKey,
826        value: &Value,
827    ) -> Result<Option<Project>, SourceError> {
828        self.metadata_updates(MetadataRecord::Project)?;
829        let result: ProjectResult = self
830            .ask("set_project_metadata", metadata_params(id, key, value))
831            .await?;
832        Ok(result.project)
833    }
834
835    async fn set_document_metadata(
836        &self,
837        id: &NativeId,
838        key: &MetadataKey,
839        value: &Value,
840    ) -> Result<Option<Document>, SourceError> {
841        self.metadata_updates(MetadataRecord::Document)?;
842        let result: DocumentResult = self
843            .ask("set_document_metadata", metadata_params(id, key, value))
844            .await?;
845        Ok(result.document)
846    }
847
848    async fn task_template_answers(
849        &self,
850        id: &NativeId,
851    ) -> Result<Option<BTreeMap<String, Value>>, SourceError> {
852        let _ = id;
853        Err(self.unrendered("a task's stored template answers, which `task answers` prints and `task render` regenerates over"))
854    }
855
856    async fn document_template_answers(
857        &self,
858        id: &NativeId,
859    ) -> Result<Option<BTreeMap<String, Value>>, SourceError> {
860        let _ = id;
861        Err(self.unrendered("a document's stored template answers, which `document answers` prints and `document render` regenerates over"))
862    }
863
864    async fn write_task_rendered(
865        &self,
866        write: &ItemWrite<Task>,
867        answers: &BTreeMap<String, Value>,
868    ) -> Result<NativeId, SourceError> {
869        let _ = (write, answers);
870        Err(self.unrendered("a task create from a template"))
871    }
872
873    async fn write_document_rendered(
874        &self,
875        write: &ItemWrite<Document>,
876        answers: &BTreeMap<String, Value>,
877    ) -> Result<NativeId, SourceError> {
878        let _ = (write, answers);
879        Err(self.unrendered("a document create from a template"))
880    }
881
882    async fn set_task_rendering(
883        &self,
884        id: &NativeId,
885        content: &str,
886        provenance: &Value,
887        answers: &BTreeMap<String, Value>,
888    ) -> Result<Option<()>, SourceError> {
889        let _ = (id, content, provenance, answers);
890        Err(self.unrendered("a task's regenerate in place"))
891    }
892
893    async fn set_document_rendering(
894        &self,
895        id: &NativeId,
896        content: &str,
897        provenance: &Value,
898        answers: &BTreeMap<String, Value>,
899    ) -> Result<Option<()>, SourceError> {
900        let _ = (id, content, provenance, answers);
901        Err(self.unrendered("a document's regenerate in place"))
902    }
903
904    async fn delete_task(&self, id: &NativeId) -> Result<(), SourceError> {
905        let _: IgnoredResult = self
906            .ask("delete_task", params(&DeleteParams { id: id.clone() }))
907            .await?;
908        Ok(())
909    }
910
911    async fn delete_project(&self, id: &NativeId) -> Result<(), SourceError> {
912        let _: IgnoredResult = self
913            .ask("delete_project", params(&DeleteParams { id: id.clone() }))
914            .await?;
915        Ok(())
916    }
917
918    async fn get_document(&self, id: &NativeId) -> Result<Option<Document>, SourceError> {
919        let result: DocumentResult = self
920            .ask("get_document", params(&IdParams { id: id.clone() }))
921            .await?;
922        Ok(result.document)
923    }
924
925    async fn query_documents(
926        &self,
927        query: &DocumentQuery,
928        page: &PageRequest,
929    ) -> Result<Page<Document>, SourceError> {
930        self.ask(
931            "query_documents",
932            params(&DocumentQueryParams {
933                query: query.clone(),
934                page: page.clone(),
935            }),
936        )
937        .await
938    }
939
940    async fn write_document(&self, write: &ItemWrite<Document>) -> Result<NativeId, SourceError> {
941        let result: WriteResult = self
942            .ask(
943                "write_document",
944                params(&DocumentWriteParams {
945                    write: write.clone(),
946                }),
947            )
948            .await?;
949        Ok(result.id)
950    }
951
952    async fn delete_document(&self, id: &NativeId) -> Result<(), SourceError> {
953        let _: IgnoredResult = self
954            .ask("delete_document", params(&DeleteParams { id: id.clone() }))
955            .await?;
956        Ok(())
957    }
958
959    async fn task_comments(
960        &self,
961        task: &NativeId,
962        page: &PageRequest,
963    ) -> Result<Option<Page<Comment>>, SourceError> {
964        let result: CommentsResult = self
965            .ask(
966                "task_comments",
967                params(&CommentsParams {
968                    task: task.clone(),
969                    page: page.clone(),
970                }),
971            )
972            .await?;
973        Ok(result.page)
974    }
975
976    async fn add_comment(
977        &self,
978        task: &NativeId,
979        comment: &NewComment,
980    ) -> Result<Option<Comment>, SourceError> {
981        let result: CommentResult = self
982            .ask(
983                "add_comment",
984                params(&AddCommentParams {
985                    task: task.clone(),
986                    comment: comment.clone(),
987                }),
988            )
989            .await?;
990        Ok(result.comment)
991    }
992
993    async fn edit_comment(
994        &self,
995        task: &NativeId,
996        comment: &NativeId,
997        body: &CommentBody,
998    ) -> Result<Option<Comment>, SourceError> {
999        let result: CommentResult = self
1000            .ask(
1001                "edit_comment",
1002                params(&EditCommentParams {
1003                    task: task.clone(),
1004                    comment: comment.clone(),
1005                    body: body.clone(),
1006                }),
1007            )
1008            .await?;
1009        Ok(result.comment)
1010    }
1011
1012    async fn delete_comment(
1013        &self,
1014        task: &NativeId,
1015        comment: &NativeId,
1016    ) -> Result<Option<NativeId>, SourceError> {
1017        let result: DeletedCommentResult = self
1018            .ask(
1019                "delete_comment",
1020                params(&DeleteCommentParams {
1021                    task: task.clone(),
1022                    comment: comment.clone(),
1023                }),
1024            )
1025            .await?;
1026        Ok(result.deleted)
1027    }
1028
1029    async fn metering(&self) -> Result<Option<Metering>, SourceError> {
1030        // Never sent to a plugin that did not declare it (§3.4), which is what lets a
1031        // plugin written before there was metering go on working without an edit.
1032        if !self.meters {
1033            return Ok(None);
1034        }
1035        let result: MeteringResult = self.ask("metering", json!({})).await?;
1036        Ok(result.metering)
1037    }
1038}
1039
1040/// The object §4.10 answers with, decoded so that `ask` has a type to hand back.
1041///
1042/// A named type rather than `serde_json::Value` so a plugin answering with something other
1043/// than an object is still refused where every other method's answer is. It does **not**
1044/// require that object to be empty, and no `deny_unknown_fields` belongs here: §2.1 is that
1045/// a reader ignores members it does not know, at every level, which is what lets a later
1046/// version add an optional one without a version bump. Refusing an unknown member here
1047/// would refuse that plugin outright, and would be the only type of this boundary that did.
1048#[derive(serde::Deserialize)]
1049struct IgnoredResult {}
1050
1051/// One method's parameters as the object the envelope carries.
1052///
1053/// Every parameter type in `wire` is built from contract types that all serialize, so
1054/// this cannot fail for a reason a caller could act on.
1055fn params<T: serde::Serialize>(value: &T) -> Value {
1056    serde_json::to_value(value).expect("method parameters are plain data")
1057}
1058
1059/// The parameters of any of the three narrow metadata writes (§4.18).
1060fn metadata_params(id: &NativeId, key: &MetadataKey, value: &Value) -> Value {
1061    params(&MetadataParams {
1062        id: id.clone(),
1063        key: key.clone(),
1064        value: value.clone(),
1065    })
1066}