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