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