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